// 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 ops const OP_MESSAGE = 1; const OP_ERROR = -1; const asRecord = (value: unknown): Record => typeof value === "object" && value !== null ? (value as Record) : {}; const str = (value: unknown): string => (typeof value === "string" ? value : ""); // decodes binary websocket frame into dag-cbor header and body export 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 function export 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) => { 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(); }; };