Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495// spindle websocket log stream// replays from line one on connect and closes when all workflows settle
import { decodeFirst } from "@atcute/cbor";import { spindleUrl } from "$lib/api/spindle";import type * as SubscribeLogs from "$lib/api/lexicons/types/sh/tangled/ci/subscribePipelineLogs";
export type LogControl = SubscribeLogs.Control;
export type LogData = SubscribeLogs.Data;
export type LogFrame = | { kind: "control"; control: LogControl } | { kind: "data"; data: LogData } | { kind: "error"; error: string; message: string };
// indigo event header opsconst OP_MESSAGE = 1;const OP_ERROR = -1;
const asRecord = (value: unknown): Record<string, unknown> => typeof value === "object" && value !== null ? (value as Record<string, unknown>) : {};
const str = (value: unknown): string => (typeof value === "string" ? value : "");
// decodes binary websocket frame into dag-cbor header and bodyexport const decodeLogFrame = (bytes: Uint8Array): LogFrame | null => { const [rawHeader, rest] = decodeFirst(bytes); const header = asRecord(rawHeader); if (header.op !== OP_MESSAGE && header.op !== OP_ERROR) return null;
const [rawBody] = decodeFirst(rest); const body = asRecord(rawBody);
if (header.op === OP_ERROR) { return { kind: "error", error: str(body.error), message: str(body.message) }; }
switch (header.t) { case "#control": return { kind: "control", control: body as unknown as LogControl }; case "#data": return { kind: "data", data: body as unknown as LogData }; default: return null; }};
export const pipelineLogsUrl = (host: string, pipeline: string, workflows: string[]): string => { const url = new URL( "/xrpc/sh.tangled.ci.subscribePipelineLogs", `${spindleUrl(host).replace(/^http/, "ws")}/` ); url.searchParams.set("pipeline", pipeline); for (const workflow of workflows) url.searchParams.append("workflows", workflow); return url.toString();};
export interface LogStreamHandlers { onFrame: (frame: LogFrame) => void; onClose?: (clean: boolean) => void;}
// opens log subscription, returns teardown functionexport const subscribePipelineLogs = ( host: string, pipeline: string, workflows: string[], handlers: LogStreamHandlers): (() => void) => { const socket = new WebSocket(pipelineLogsUrl(host, pipeline, workflows)); socket.binaryType = "arraybuffer"; let live = true;
socket.onmessage = (event: MessageEvent<ArrayBuffer | string>) => { if (!live || typeof event.data === "string") return; const frame = decodeLogFrame(new Uint8Array(event.data)); if (frame) handlers.onFrame(frame); };
socket.onclose = (event: CloseEvent) => { if (!live) return; live = false; handlers.onClose?.(event.wasClean); };
// refused upgrades trigger an error then close, close handler handles it socket.onerror = () => {};
return () => { live = false; socket.close(); };};