import { readFile } from "node:fs/promises"; import { createServer, type Server } from "node:http"; import { createRequire } from "node:module"; import { dirname, join } from "node:path"; import { createInterface, type Interface as ReadlineInterface } from "node:readline"; import type { Readable } from "node:stream"; import { fileURLToPath, URL } from "node:url"; import WebSocket, { WebSocketServer } from "ws"; import type { TraceEnvelope, TraceEventType } from "../src/types.ts"; import { ViewerEventStore } from "./store.ts"; const EVENT_TYPES = new Set([ "run.start", "run.end", "interaction.start", "interaction.end", "turn.start", "turn.end", "llm.request", "llm.response", "tool.request", "tool.response", "trace.warning", ]); export interface StartViewerServerOptions { tracePath?: string; token: string; port: number; assetsDir?: string; input?: Readable; onWarning?: (message: string) => void; } export interface RunningViewerServer { port: number; ingest(event: TraceEnvelope): boolean; close(): Promise; } function childArgument(arguments_: readonly string[], name: string): string | undefined { const index = arguments_.indexOf(name); return index >= 0 ? arguments_[index + 1] : undefined; } function isTraceEnvelope(value: unknown): value is TraceEnvelope { if (typeof value !== "object" || value === null || Array.isArray(value)) return false; const record = value as Record; return record.schemaVersion === 1 && typeof record.seq === "number" && typeof record.eventId === "string" && typeof record.timestamp === "string" && typeof record.sessionId === "string" && typeof record.runId === "string" && typeof record.type === "string" && EVENT_TYPES.has(record.type as TraceEventType) && "payload" in record; } function parseTraceLine(line: string): TraceEnvelope | undefined { try { const parsed: unknown = JSON.parse(line); return isTraceEnvelope(parsed) ? parsed : undefined; } catch { return undefined; } } function tokenMatches(requestUrl: URL, token: string): boolean { return requestUrl.searchParams.get("token") === token; } async function replayTrace(tracePath: string | undefined, store: ViewerEventStore, warn: (message: string) => void): Promise { if (tracePath === undefined) return; let contents: string; try { contents = await readFile(tracePath, "utf8"); } catch (error) { if (error instanceof Error && "code" in error && error.code === "ENOENT") return; warn(`Could not replay trace ${tracePath}: ${error instanceof Error ? error.message : String(error)}`); return; } const lines = contents.split("\n"); if (!contents.endsWith("\n")) lines.pop(); for (const line of lines) { if (line.length === 0) continue; const event = parseTraceLine(line); if (event !== undefined) store.ingest(event); } } function closeHttpServer(server: Server): Promise { return new Promise((resolve) => { server.close(() => resolve()); }); } export async function startViewerServer(options: StartViewerServerOptions): Promise { const warn = options.onWarning ?? (() => undefined); const store = new ViewerEventStore(); await replayTrace(options.tracePath, store, warn); const assetsDir = options.assetsDir ?? dirname(fileURLToPath(import.meta.url)); const markedPath = createRequire(import.meta.url).resolve("marked"); const assetPaths: Readonly> = { "/": { path: join(assetsDir, "index.html"), contentType: "text/html; charset=utf-8" }, "/app.js": { path: join(assetsDir, "app.js"), contentType: "text/javascript; charset=utf-8" }, "/collapse-state.js": { path: join(assetsDir, "collapse-state.js"), contentType: "text/javascript; charset=utf-8" }, "/conversation.js": { path: join(assetsDir, "conversation.js"), contentType: "text/javascript; charset=utf-8" }, "/detail-model.js": { path: join(assetsDir, "detail-model.js"), contentType: "text/javascript; charset=utf-8" }, "/labels.js": { path: join(assetsDir, "labels.js"), contentType: "text/javascript; charset=utf-8" }, "/markdown.js": { path: join(assetsDir, "markdown.js"), contentType: "text/javascript; charset=utf-8" }, "/reconnect.js": { path: join(assetsDir, "reconnect.js"), contentType: "text/javascript; charset=utf-8" }, "/styles.css": { path: join(assetsDir, "styles.css"), contentType: "text/css; charset=utf-8" }, "/tree-state.js": { path: join(assetsDir, "tree-state.js"), contentType: "text/javascript; charset=utf-8" }, "/vendor/marked.js": { path: markedPath, contentType: "text/javascript; charset=utf-8" }, }; const server = createServer((request, response) => { void (async () => { const requestUrl = new URL(request.url ?? "/", "http://127.0.0.1"); if (request.method !== "GET") { response.writeHead(405).end("Method Not Allowed"); return; } if (requestUrl.pathname === "/api/snapshot") { if (!tokenMatches(requestUrl, options.token)) { response.writeHead(401).end("Unauthorized"); return; } response.writeHead(200, { "content-type": "application/json; charset=utf-8", "cache-control": "no-store" }); response.end(JSON.stringify(store.snapshot())); return; } const asset = assetPaths[requestUrl.pathname]; if (asset === undefined) { response.writeHead(404).end("Not Found"); return; } try { const contents = await readFile(asset.path); response.writeHead(200, { "content-type": asset.contentType, "cache-control": "no-store" }); response.end(contents); } catch (error) { warn(`Could not serve Viewer asset ${requestUrl.pathname}: ${error instanceof Error ? error.message : String(error)}`); response.writeHead(500).end("Viewer asset unavailable"); } })().catch((error: unknown) => { warn(`Viewer request failed: ${error instanceof Error ? error.message : String(error)}`); if (!response.headersSent) response.writeHead(500); response.end(); }); }); const webSockets = new WebSocketServer({ noServer: true }); server.on("upgrade", (request, socket, head) => { const requestUrl = new URL(request.url ?? "/", "http://127.0.0.1"); if (requestUrl.pathname !== "/ws" || !tokenMatches(requestUrl, options.token)) { socket.write("HTTP/1.1 401 Unauthorized\r\nConnection: close\r\n\r\n"); socket.destroy(); return; } webSockets.handleUpgrade(request, socket, head, (webSocket) => { webSockets.emit("connection", webSocket, request); }); }); webSockets.on("connection", (webSocket) => { webSocket.on("error", (error) => warn(`Viewer WebSocket failed: ${error.message}`)); webSocket.send(JSON.stringify({ type: "snapshot", snapshot: store.snapshot() })); }); webSockets.on("error", (error) => warn(`Viewer WebSocket server failed: ${error.message}`)); const unsubscribe = store.subscribe((event) => { const message = JSON.stringify({ type: "event", event }); for (const client of webSockets.clients) { if (client.readyState === WebSocket.OPEN) client.send(message); } }); let inputLines: ReadlineInterface | undefined; if (options.input !== undefined) { inputLines = createInterface({ input: options.input, crlfDelay: Number.POSITIVE_INFINITY }); inputLines.on("line", (line) => { const event = parseTraceLine(line); if (event !== undefined) store.ingest(event); }); inputLines.on("error", (error) => warn(`Viewer stdin failed: ${error.message}`)); } await new Promise((resolve, reject) => { const onError = (error: Error): void => reject(error); server.once("error", onError); server.listen(options.port, "127.0.0.1", () => { server.off("error", onError); resolve(); }); }); server.on("error", (error) => warn(`Viewer HTTP server failed: ${error.message}`)); const address = server.address(); if (address === null || typeof address === "string") throw new Error("Viewer server did not bind a TCP port"); let closed = false; return { port: address.port, ingest: (event) => store.ingest(event), close: async () => { if (closed) return; closed = true; unsubscribe(); inputLines?.close(); for (const client of webSockets.clients) client.terminate(); webSockets.close(); await closeHttpServer(server); }, }; } export async function runViewerChild(arguments_: readonly string[]): Promise { const token = childArgument(arguments_, "--token"); const portValue = childArgument(arguments_, "--port") ?? "0"; const port = Number(portValue); if (token === undefined || !Number.isInteger(port) || port < 0 || port > 65535) { throw new Error("Viewer requires --token and a valid --port"); } const server = await startViewerServer({ tracePath: childArgument(arguments_, "--trace-path"), token, port, input: process.stdin, onWarning: (message) => process.stderr.write(`${message}\n`), }); process.stdout.write(`${JSON.stringify({ type: "ready", port: server.port })}\n`); process.stdin.once("end", () => { void server.close().catch((error: unknown) => { process.stderr.write(`Viewer shutdown failed: ${error instanceof Error ? error.message : String(error)}\n`); }); }); }