diff --git a/plugin/lvmh-agent.ts b/plugin/lvmh-agent.ts index bbc89a7..0dd7d70 100644 --- a/plugin/lvmh-agent.ts +++ b/plugin/lvmh-agent.ts @@ -20,6 +20,12 @@ * events are additionally kept in a bounded replay buffer and resent * after welcome.lastSeq on reconnect * - reconnect with exponential backoff + jitter, capped, forever + * - handshake watchdog: if welcome does not arrive within + * LVMH_WELCOME_TIMEOUT_MS (default 15s) of connecting, the socket is + * abandoned and retried (covers a daemon that accepts TCP then hangs) + * - stall watchdog: undici exposes no ping(), so a daemon that dies + * without a TCP close is detected via bufferedAmount never draining + * (two LVMH_STALL_CHECK_MS ticks over LVMH_STALL_BYTES, default 30s/16MB) * - errors are logged to ~/.pi/lvmh-agent.log only; stdout stays clean * - session_shutdown closes the socket and stops all timers immediately */ @@ -27,7 +33,10 @@ import * as fs from "node:fs"; import * as os from "node:os"; import * as path from "node:path"; -import type { ExtensionAPI, ExtensionContext } from "@earendil-works/pi-coding-agent"; +import type { + ExtensionAPI, + ExtensionContext, +} from "@earendil-works/pi-coding-agent"; const PROTOCOL_VERSION: number = 1; const ENV_URL: string = "LVMH_URL"; @@ -43,26 +52,40 @@ const BACKOFF_JITTER_MS: number = 1000; const RESULT_PREVIEW_MAX_CHARS: number = 2000; const LOG_LINE_MAX_CHARS: number = 500; +function envPositiveInt(name: string, fallback: number): number { + const raw: string | undefined = process.env[name]; + if (raw === undefined) return fallback; + const n: number = Number.parseInt(raw, 10); + return Number.isFinite(n) && n > 0 ? n : fallback; +} + +// Test/debug knobs; production never sets them. Read lazily so harnesses +// can reconfigure per scenario after the module is already imported. +const welcomeTimeoutMs = (): number => + envPositiveInt("LVMH_WELCOME_TIMEOUT_MS", 15000); +const stallCheckMs = (): number => envPositiveInt("LVMH_STALL_CHECK_MS", 30000); +const stallBytes = (): number => envPositiveInt("LVMH_STALL_BYTES", 16777216); + const TRANSIENT_TYPES: ReadonlySet = new Set(["message_update"]); interface Frame { - v: number; - sessionId: string; - seq: number; - ts: number; - type: string; - [key: string]: unknown; + v: number; + sessionId: string; + seq: number; + ts: number; + type: string; + [key: string]: unknown; } interface SessionSnapshot { - id: string; - name: string | null; - cwd: string; - model: string | null; - provider: string | null; - agent: boolean; - repo: string | null; - startedAt: number; + id: string; + name: string | null; + cwd: string; + model: string | null; + provider: string | null; + agent: boolean; + repo: string | null; + startedAt: number; } /** Per-session seq counters survive reconnects and instance rebinds in-process. */ @@ -73,545 +96,735 @@ const logPath: string = path.join(os.homedir(), ".pi", "lvmh-agent.log"); /** Best-effort append to the log file. Never throws, never touches stdout. */ function log(message: unknown): void { - try { - const line = `${new Date().toISOString()} ${String(message).slice(0, LOG_LINE_MAX_CHARS)}\n`; - fs.appendFile(logPath, line, () => undefined); - } catch { - // even appendFile argument validation failures must stay silent - } + try { + const line = `${new Date().toISOString()} ${String(message).slice(0, LOG_LINE_MAX_CHARS)}\n`; + fs.appendFile(logPath, line, () => undefined); + } catch { + // even appendFile argument validation failures must stay silent + } } function nextSeq(sessionId: string): number { - const n: number = (seqCounters.get(sessionId) ?? 0) + 1; - seqCounters.set(sessionId, n); - return n; + const n: number = (seqCounters.get(sessionId) ?? 0) + 1; + seqCounters.set(sessionId, n); + return n; } /** Concatenate text blocks of a pi content array (or pass through a string). */ function textOfContent(content: unknown): string { - if (typeof content === "string") return content; - if (!Array.isArray(content)) return ""; - let out = ""; - for (const block of content) { - if (block !== null && typeof block === "object" && (block as { type?: unknown }).type === "text") { - const text = (block as { text?: unknown }).text; - if (typeof text === "string") out += text; - } - } - return out; + if (typeof content === "string") return content; + if (!Array.isArray(content)) return ""; + let out = ""; + for (const block of content) { + if ( + block !== null && + typeof block === "object" && + (block as { type?: unknown }).type === "text" + ) { + const text = (block as { text?: unknown }).text; + if (typeof text === "string") out += text; + } + } + return out; } function safeJsonStringify(value: unknown): string { - try { - return JSON.stringify(value) ?? "null"; - } catch { - return "null"; - } + try { + return JSON.stringify(value) ?? "null"; + } catch { + return "null"; + } } /** Map a pi AgentMessage to the protocol Message shape (message_end payload). */ function mapMessage(message: unknown): Record | null { - if (message === null || typeof message !== "object") return null; - const msg = message as { - role?: unknown; - content?: unknown; - timestamp?: unknown; - toolCallId?: unknown; - }; - const role = typeof msg.role === "string" ? msg.role : "unknown"; - const content = msg.content; - const toolCalls: Array<{ id: string; name: string; argsJson: string }> = []; - let thinking = ""; - if (Array.isArray(content)) { - for (const block of content) { - if (block === null || typeof block !== "object") continue; - const b = block as { type?: unknown; id?: unknown; name?: unknown; arguments?: unknown }; - if (b.type === "toolCall") { - toolCalls.push({ - id: String(b.id ?? ""), - name: String(b.name ?? ""), - argsJson: safeJsonStringify(b.arguments), - }); - } else if (b.type === "thinking" && typeof (b as { thinking?: unknown }).thinking === "string") { - thinking += (b as { thinking: string }).thinking; - } - } - } - return { - role, - id: `${role}-${typeof msg.timestamp === "number" ? msg.timestamp : Date.now()}`, - text: textOfContent(content), - thinking: thinking.length > 0 ? thinking : null, - toolCalls, - toolCallId: role === "toolResult" && typeof msg.toolCallId === "string" ? msg.toolCallId : null, - }; + if (message === null || typeof message !== "object") return null; + const msg = message as { + role?: unknown; + content?: unknown; + timestamp?: unknown; + toolCallId?: unknown; + }; + const role = typeof msg.role === "string" ? msg.role : "unknown"; + const content = msg.content; + const toolCalls: Array<{ id: string; name: string; argsJson: string }> = []; + let thinking = ""; + if (Array.isArray(content)) { + for (const block of content) { + if (block === null || typeof block !== "object") continue; + const b = block as { + type?: unknown; + id?: unknown; + name?: unknown; + arguments?: unknown; + }; + if (b.type === "toolCall") { + toolCalls.push({ + id: String(b.id ?? ""), + name: String(b.name ?? ""), + argsJson: safeJsonStringify(b.arguments), + }); + } else if ( + b.type === "thinking" && + typeof (b as { thinking?: unknown }).thinking === "string" + ) { + thinking += (b as { thinking: string }).thinking; + } + } + } + return { + role, + id: `${role}-${typeof msg.timestamp === "number" ? msg.timestamp : Date.now()}`, + text: textOfContent(content), + thinking: thinking.length > 0 ? thinking : null, + toolCalls, + toolCallId: + role === "toolResult" && typeof msg.toolCallId === "string" + ? msg.toolCallId + : null, + }; } export default function (pi: ExtensionAPI): void { - const rawUrl: string | undefined = process.env[ENV_URL]; - const rawToken: string | undefined = process.env[ENV_TOKEN]; - if (!rawUrl || !rawToken) return; // inert: no daemon configured - const url: string = rawUrl; - const token: string = rawToken; - if (typeof WebSocket === "undefined") { - log("global WebSocket unavailable (need Node >= 22); extension inert"); - return; - } + const rawUrl: string | undefined = process.env[ENV_URL]; + const rawToken: string | undefined = process.env[ENV_TOKEN]; + if (!rawUrl || !rawToken) return; // inert: no daemon configured + const url: string = rawUrl; + const token: string = rawToken; + if (typeof WebSocket === "undefined") { + log("global WebSocket unavailable (need Node >= 22); extension inert"); + return; + } - let stopped: boolean = false; - let ws: WebSocket | null = null; - let wsOpen: boolean = false; - let greeted: boolean = false; - let reconnectTimer: ReturnType | null = null; - let backoffAttempt: number = 0; + let stopped: boolean = false; + let ws: WebSocket | null = null; + let wsOpen: boolean = false; + let greeted: boolean = false; + let reconnectTimer: ReturnType | null = null; + let backoffAttempt: number = 0; + let welcomeTimer: ReturnType | null = null; + let stallTimer: ReturnType | null = null; + let lastStallBytes: number | null = null; - let currentSessionId: string | null = null; - let snapshot: SessionSnapshot | null = null; - let lastStreamText: string = ""; - let lastCtx: ExtensionContext | undefined; + let currentSessionId: string | null = null; + let snapshot: SessionSnapshot | null = null; + let lastStreamText: string = ""; + let lastCtx: ExtensionContext | undefined; - let sendQueue: Frame[] = []; - let replayBuf: Frame[] = []; - let droppedEvents: number = 0; + let sendQueue: Frame[] = []; + let replayBuf: Frame[] = []; + let droppedEvents: number = 0; - /** - * Subscribe with a bulletproof wrapper: the inner handler may be sync or - * async; sync throws are caught here and returned promises that reject are - * swallowed. A failing handler logs and yields undefined — it can never - * throw into or reject on pi. - */ - function sub( - event: string, - handler: (event: any, ctx: ExtensionContext) => unknown, - ): void { - (pi as { on: (e: string, h: (event: any, ctx: ExtensionContext) => unknown) => void }).on( - event, - (evt: any, ctx: ExtensionContext): undefined => { - lastCtx = ctx; - try { - const result = handler(evt, ctx); - if (result !== null && typeof result === "object" && typeof (result as Promise).then === "function") { - (result as Promise).then(undefined, (err: unknown) => { - log(`async handler error (${event}): ${err}`); - }); - } - } catch (err) { - log(`handler error (${event}): ${err}`); - } - return undefined; - }, - ); - } + /** + * Subscribe with a bulletproof wrapper: the inner handler may be sync or + * async; sync throws are caught here and returned promises that reject are + * swallowed. A failing handler logs and yields undefined — it can never + * throw into or reject on pi. + */ + function sub( + event: string, + handler: (event: any, ctx: ExtensionContext) => unknown, + ): void { + ( + pi as { + on: ( + e: string, + h: (event: any, ctx: ExtensionContext) => unknown, + ) => void; + } + ).on(event, (evt: any, ctx: ExtensionContext): undefined => { + lastCtx = ctx; + try { + const result = handler(evt, ctx); + if ( + result !== null && + typeof result === "object" && + typeof (result as Promise).then === "function" + ) { + (result as Promise).then(undefined, (err: unknown) => { + log(`async handler error (${event}): ${err}`); + }); + } + } catch (err) { + log(`handler error (${event}): ${err}`); + } + return undefined; + }); + } - function enqueue(frame: Frame): void { - sendQueue.push(frame); - while (sendQueue.length > SEND_QUEUE_MAX) { - const dropped: Frame | undefined = sendQueue.shift(); - if (dropped !== undefined) { - droppedEvents++; - log(`send queue overflow, dropped ${dropped.type} seq=${dropped.seq}`); - } - } - flush(); - } + function enqueue(frame: Frame): void { + sendQueue.push(frame); + while (sendQueue.length > SEND_QUEUE_MAX) { + const dropped: Frame | undefined = sendQueue.shift(); + if (dropped !== undefined) { + droppedEvents++; + log(`send queue overflow, dropped ${dropped.type} seq=${dropped.seq}`); + } + } + flush(); + } - function flush(): void { - while (!stopped && ws !== null && wsOpen && greeted && sendQueue.length > 0) { - const frame: Frame | undefined = sendQueue.shift(); - if (frame === undefined) return; - let text: string; - try { - text = JSON.stringify(frame); - } catch (err) { - log(`serialize error (${frame.type}): ${err}`); - continue; - } - try { - ws.send(text); - } catch (err) { - log(`send error (${frame.type}): ${err}`); - handleDisconnect(); - scheduleReconnect(); // in case no onclose follows the failed send - return; - } - } - } + function flush(): void { + while ( + !stopped && + ws !== null && + wsOpen && + greeted && + sendQueue.length > 0 + ) { + const frame: Frame | undefined = sendQueue.shift(); + if (frame === undefined) return; + let text: string; + try { + text = JSON.stringify(frame); + } catch (err) { + log(`serialize error (${frame.type}): ${err}`); + continue; + } + try { + ws.send(text); + } catch (err) { + log(`send error (${frame.type}): ${err}`); + handleDisconnect(); + scheduleReconnect(); // in case no onclose follows the failed send + return; + } + } + } - /** Persisted-kind events only (message_update is live-only per protocol). */ - function emit(type: string, payload: Record): void { - if (stopped || currentSessionId === null) return; - const frame: Frame = { - v: PROTOCOL_VERSION, - sessionId: currentSessionId, - seq: nextSeq(currentSessionId), - ts: Date.now(), - type, - ...payload, - }; - if (!TRANSIENT_TYPES.has(type)) { - replayBuf.push(frame); - while (replayBuf.length > REPLAY_BUFFER_MAX) { - replayBuf.shift(); - droppedEvents++; - log(`replay buffer overflow (cap ${REPLAY_BUFFER_MAX}), dropped oldest event`); - } - } - enqueue(frame); - } + /** Persisted-kind events only (message_update is live-only per protocol). */ + function emit(type: string, payload: Record): void { + if (stopped || currentSessionId === null) return; + const frame: Frame = { + v: PROTOCOL_VERSION, + sessionId: currentSessionId, + seq: nextSeq(currentSessionId), + ts: Date.now(), + type, + ...payload, + }; + if (!TRANSIENT_TYPES.has(type)) { + replayBuf.push(frame); + while (replayBuf.length > REPLAY_BUFFER_MAX) { + replayBuf.shift(); + droppedEvents++; + log( + `replay buffer overflow (cap ${REPLAY_BUFFER_MAX}), dropped oldest event`, + ); + } + } + enqueue(frame); + } - function handleDisconnect(): void { - const socket: WebSocket | null = ws; - ws = null; - wsOpen = false; - greeted = false; - try { - socket?.close(); - } catch { - // close failures are meaningless here - } - } + function handleDisconnect(): void { + const socket: WebSocket | null = ws; + ws = null; + wsOpen = false; + greeted = false; + clearWelcomeTimer(); + lastStallBytes = null; + try { + socket?.close(); + } catch { + // close failures are meaningless here + } + } - function scheduleReconnect(): void { - if (stopped || reconnectTimer !== null || ws !== null) return; - const expMs: number = BACKOFF_BASE_MS * 2 ** backoffAttempt; - const waitMs: number = - Math.min(Number.isFinite(expMs) ? expMs : BACKOFF_MAX_MS, BACKOFF_MAX_MS) + - Math.floor(Math.random() * BACKOFF_JITTER_MS); - backoffAttempt++; - reconnectTimer = setTimeout(() => { - reconnectTimer = null; - connect(); - }, waitMs); - (reconnectTimer as { unref?: () => void }).unref?.(); - } + function clearWelcomeTimer(): void { + if (welcomeTimer === null) return; + try { + clearTimeout(welcomeTimer); + } catch { + // clearable timers never throw in practice + } + welcomeTimer = null; + } - function connect(): void { - if (stopped || ws !== null || reconnectTimer !== null) return; - let socket: WebSocket; - try { - socket = new WebSocket(url, { headers: { Authorization: `Bearer ${token}` } }); - } catch (err) { - log(`ws construct error: ${err}`); - scheduleReconnect(); - return; - } - ws = socket; - // Stale-socket guards: after a session switch an old socket's callbacks - // can still fire; they must never touch the state of the new socket. - socket.onopen = () => { - if (ws !== socket) return; - try { - wsOpen = true; - sendHello(); - flush(); - } catch (err) { - log(`onopen error: ${err}`); - handleDisconnect(); - scheduleReconnect(); - } - }; - socket.onmessage = (ev: unknown) => { - if (ws !== socket) return; - try { - onMessage(ev); - } catch (err) { - log(`onmessage error: ${err}`); - } - }; - socket.onerror = () => { - // undici always follows onerror with onclose; logging happens there - }; - socket.onclose = () => { - if (ws !== socket) return; - try { - log(`ws closed after backoff attempt #${backoffAttempt}`); - handleDisconnect(); - scheduleReconnect(); - } catch (err) { - log(`onclose error: ${err}`); - } - }; - } + /** A daemon that accepts TCP but never completes the handshake would hang + * the socket (and all mirroring) forever: undici has no handshake timeout + * of its own. Arm at connect time, clear at welcome/disconnect. */ + function armWelcomeTimeout(socket: WebSocket): void { + clearWelcomeTimer(); + welcomeTimer = setTimeout(() => { + welcomeTimer = null; + if (ws !== socket || greeted) return; + log("welcome timeout; abandoning socket"); + handleDisconnect(); + scheduleReconnect(); + }, welcomeTimeoutMs()); + (welcomeTimer as { unref?: () => void }).unref?.(); + } - function sendHello(): void { - if (ws === null || currentSessionId === null || snapshot === null) return; - const hello: Frame = { - v: PROTOCOL_VERSION, - sessionId: currentSessionId, - seq: 0, - ts: Date.now(), - type: "hello", - session: snapshot, - }; - ws.send(JSON.stringify(hello)); - } + /** undici's WHATWG WebSocket has no ping(). Detect a daemon that died + * without a TCP close by watching bufferedAmount: if a backlog over + * STALL_BYTES makes no progress between two ticks, the peer is not + * reading — kill the socket and reconnect. Draining resets the state. */ + function checkStall(): void { + const socket: WebSocket | null = ws; + if (socket === null || !wsOpen) { + lastStallBytes = null; + return; + } + const buffered: number = socket.bufferedAmount; + if (!Number.isFinite(buffered) || buffered < stallBytes()) { + lastStallBytes = null; + return; + } + if (lastStallBytes !== null && buffered >= lastStallBytes) { + log(`ws stalled: bufferedAmount=${buffered} not draining; reconnecting`); + lastStallBytes = null; + handleDisconnect(); + scheduleReconnect(); + return; + } + lastStallBytes = buffered; + } - function onWelcome(frame: Record): void { - if (typeof frame.sessionId === "string" && frame.sessionId !== currentSessionId) { - log(`welcome for foreign session ${String(frame.sessionId)}, ignoring`); - return; - } - const lastSeq: number = typeof frame.lastSeq === "number" && Number.isFinite(frame.lastSeq) ? frame.lastSeq : 0; - backoffAttempt = 0; - greeted = true; - // Post-restart safety: never reuse seq numbers the daemon already has. - if (currentSessionId !== null) { - const counter: number = seqCounters.get(currentSessionId) ?? 0; - if (counter <= lastSeq) seqCounters.set(currentSessionId, lastSeq); - } - // Replay persisted events the daemon is missing, before queued live - // frames. Persisted-kind frames still sitting in the queue are dropped - // here: they are already in replayBuf, so sending both would duplicate - // them (only transient message_update frames survive the merge). - const replay: Frame[] = replayBuf.filter((f) => f.seq > lastSeq); - sendQueue = replay.concat(sendQueue.filter((f) => TRANSIENT_TYPES.has(f.type))); - while (sendQueue.length > SEND_QUEUE_MAX) { - const dropped: Frame | undefined = sendQueue.shift(); - if (dropped !== undefined) { - droppedEvents++; - log(`send queue overflow during replay, dropped ${dropped.type} seq=${dropped.seq}`); - } - } - if (droppedEvents > 0) { - emit("buffer_overflow", { dropped: droppedEvents }); - droppedEvents = 0; - } - flush(); - } + function scheduleReconnect(): void { + if (stopped || reconnectTimer !== null || ws !== null) return; + const expMs: number = BACKOFF_BASE_MS * 2 ** backoffAttempt; + const waitMs: number = + Math.min( + Number.isFinite(expMs) ? expMs : BACKOFF_MAX_MS, + BACKOFF_MAX_MS, + ) + Math.floor(Math.random() * BACKOFF_JITTER_MS); + backoffAttempt++; + reconnectTimer = setTimeout(() => { + reconnectTimer = null; + connect(); + }, waitMs); + (reconnectTimer as { unref?: () => void }).unref?.(); + } - function onMessage(ev: unknown): void { - const data = (ev as { data?: unknown }).data; - let frame: unknown; - try { - frame = JSON.parse(typeof data === "string" ? data : ""); - } catch { - log("malformed json frame from daemon"); - return; - } - if (frame === null || typeof frame !== "object") return; - const f = frame as { v?: unknown; type?: unknown; message?: unknown; sessionId?: unknown }; - if (f.v !== PROTOCOL_VERSION) return; - if (f.type === "welcome") onWelcome(f as Record); - else if (f.type === "prompt") deliverPrompt(f); - else if (f.type === "abort") doAbort(); - // unknown types are ignored (forward compatibility) - } + function connect(): void { + if (stopped || ws !== null || reconnectTimer !== null) return; + let socket: WebSocket; + try { + socket = new WebSocket(url, { + headers: { Authorization: `Bearer ${token}` }, + }); + } catch (err) { + log(`ws construct error: ${err}`); + scheduleReconnect(); + return; + } + ws = socket; + armWelcomeTimeout(socket); + if (stallTimer === null) { + stallTimer = setInterval(checkStall, stallCheckMs()); + (stallTimer as { unref?: () => void }).unref?.(); + } + // Stale-socket guards: after a session switch an old socket's callbacks + // can still fire; they must never touch the state of the new socket. + socket.onopen = () => { + if (ws !== socket) return; + try { + wsOpen = true; + sendHello(); + flush(); + } catch (err) { + log(`onopen error: ${err}`); + handleDisconnect(); + scheduleReconnect(); + } + }; + socket.onmessage = (ev: unknown) => { + if (ws !== socket) return; + try { + onMessage(ev); + } catch (err) { + log(`onmessage error: ${err}`); + } + }; + socket.onerror = () => { + // undici always follows onerror with onclose; logging happens there + }; + socket.onclose = () => { + if (ws !== socket) return; + try { + log(`ws closed after backoff attempt #${backoffAttempt}`); + handleDisconnect(); + scheduleReconnect(); + } catch (err) { + log(`onclose error: ${err}`); + } + }; + } - function deliverPrompt(frame: { message?: unknown; sessionId?: unknown }): void { - const message = frame.message; - if (typeof message !== "string" || message.length === 0) return; - if (typeof frame.sessionId === "string" && frame.sessionId !== currentSessionId) return; - try { - pi.sendUserMessage(message, { deliverAs: "steer" }); - } catch (err) { - log(`sendUserMessage failed: ${err}`); - } - } + function sendHello(): void { + if (ws === null || currentSessionId === null || snapshot === null) return; + const hello: Frame = { + v: PROTOCOL_VERSION, + sessionId: currentSessionId, + seq: 0, + ts: Date.now(), + type: "hello", + session: snapshot, + }; + ws.send(JSON.stringify(hello)); + } - function doAbort(): void { - try { - const maybeAbort = (pi as unknown as { abort?: () => void }).abort; - if (typeof maybeAbort === "function") maybeAbort.call(pi); - else lastCtx?.abort(); - } catch (err) { - log(`abort failed: ${err}`); - } - } + function onWelcome(frame: Record): void { + if ( + typeof frame.sessionId === "string" && + frame.sessionId !== currentSessionId + ) { + log(`welcome for foreign session ${String(frame.sessionId)}, ignoring`); + return; + } + const lastSeq: number = + typeof frame.lastSeq === "number" && Number.isFinite(frame.lastSeq) + ? frame.lastSeq + : 0; + backoffAttempt = 0; + greeted = true; + clearWelcomeTimer(); + // Post-restart safety: never reuse seq numbers the daemon already has. + if (currentSessionId !== null) { + const counter: number = seqCounters.get(currentSessionId) ?? 0; + if (counter <= lastSeq) seqCounters.set(currentSessionId, lastSeq); + } + // Replay persisted events the daemon is missing, before queued live + // frames. Persisted-kind frames still sitting in the queue are dropped + // here: they are already in replayBuf, so sending both would duplicate + // them (only transient message_update frames survive the merge). + const replay: Frame[] = replayBuf.filter((f) => f.seq > lastSeq); + sendQueue = replay.concat( + sendQueue.filter((f) => TRANSIENT_TYPES.has(f.type)), + ); + while (sendQueue.length > SEND_QUEUE_MAX) { + const dropped: Frame | undefined = sendQueue.shift(); + if (dropped !== undefined) { + droppedEvents++; + log( + `send queue overflow during replay, dropped ${dropped.type} seq=${dropped.seq}`, + ); + } + } + if (droppedEvents > 0) { + emit("buffer_overflow", { dropped: droppedEvents }); + droppedEvents = 0; + } + flush(); + } - function buildSnapshot(ctx: ExtensionContext, sessionId: string): SessionSnapshot { - const sm = ctx.sessionManager as { - getSessionName?: () => string | undefined; - getCwd?: () => string; - getHeader?: () => { timestamp?: unknown } | null; - }; - let name: string | null = null; - let cwd: string = ctx.cwd; - let startedAt: number | undefined = sessionStartTs.get(sessionId); - try { - name = sm.getSessionName?.() ?? null; - cwd = sm.getCwd?.() ?? ctx.cwd; - if (startedAt === undefined) { - const headerTs = sm.getHeader?.()?.timestamp; - startedAt = typeof headerTs === "string" ? Date.parse(headerTs) : NaN; - if (!Number.isFinite(startedAt)) startedAt = Date.now(); - sessionStartTs.set(sessionId, startedAt); - } - } catch (err) { - log(`snapshot fallback (${String(err).slice(0, 80)}): using ctx.cwd/now`); - if (startedAt === undefined) { - startedAt = Date.now(); - sessionStartTs.set(sessionId, startedAt); - } - } - const model = ctx.model as { id?: unknown; provider?: unknown } | undefined; - return { - id: sessionId, - name, - cwd, - model: model !== undefined && typeof model.id === "string" ? model.id : null, - provider: model !== undefined && typeof model.provider === "string" ? model.provider : null, - agent: process.env[ENV_AGENT] === "1", - repo: process.env[ENV_REPO] ?? null, - startedAt, - }; - } + function onMessage(ev: unknown): void { + const data = (ev as { data?: unknown }).data; + let frame: unknown; + try { + frame = JSON.parse(typeof data === "string" ? data : ""); + } catch { + log("malformed json frame from daemon"); + return; + } + if (frame === null || typeof frame !== "object") return; + const f = frame as { + v?: unknown; + type?: unknown; + message?: unknown; + sessionId?: unknown; + }; + if (f.v !== PROTOCOL_VERSION) return; + if (f.type === "welcome") onWelcome(f as Record); + else if (f.type === "prompt") deliverPrompt(f); + else if (f.type === "abort") doAbort(); + // unknown types are ignored (forward compatibility) + } - sub("session_start", (_event: any, ctx: ExtensionContext) => { - stopped = false; - let sessionId: string | null = null; - try { - const id = (ctx.sessionManager as { getSessionId?: () => string }).getSessionId?.(); - sessionId = typeof id === "string" && id.length > 0 ? id : null; - } catch (err) { - log(`getSessionId failed: ${err}`); - } - if (sessionId === null) { - log("no session id; lvmh mirroring disabled for this session"); - return; - } - if (sessionId !== currentSessionId) { - sendQueue = []; - replayBuf = []; - droppedEvents = 0; - lastStreamText = ""; - handleDisconnect(); // drop any socket bound to the previous session - try { - if (reconnectTimer !== null) clearTimeout(reconnectTimer); - } catch { - // clearable timers never throw in practice - } - reconnectTimer = null; - } - currentSessionId = sessionId; - snapshot = buildSnapshot(ctx, sessionId); - connect(); - }); + function deliverPrompt(frame: { + message?: unknown; + sessionId?: unknown; + }): void { + const message = frame.message; + if (typeof message !== "string" || message.length === 0) return; + if ( + typeof frame.sessionId === "string" && + frame.sessionId !== currentSessionId + ) + return; + try { + pi.sendUserMessage(message, { deliverAs: "steer" }); + } catch (err) { + log(`sendUserMessage failed: ${err}`); + } + } - sub("session_shutdown", () => { - stopped = true; - try { - if (reconnectTimer !== null) clearTimeout(reconnectTimer); - } catch { - // ignore - } - reconnectTimer = null; - const socket: WebSocket | null = ws; - ws = null; - wsOpen = false; - greeted = false; - try { - if (socket !== null && socket.readyState === WebSocket.OPEN && currentSessionId !== null) { - socket.send( - JSON.stringify({ - v: PROTOCOL_VERSION, - sessionId: currentSessionId, - seq: 0, - ts: Date.now(), - type: "bye", - reason: "shutdown", - }), - ); - } - socket?.close(); - } catch { - // best-effort goodbye; nothing to flush, exit fast - } - }); + function doAbort(): void { + try { + const maybeAbort = (pi as unknown as { abort?: () => void }).abort; + if (typeof maybeAbort === "function") maybeAbort.call(pi); + else lastCtx?.abort(); + } catch (err) { + log(`abort failed: ${err}`); + } + } - sub("session_info_changed", (event: { name?: unknown }) => { - if (snapshot === null) return; - const name = typeof event.name === "string" ? event.name : null; - if (name === snapshot.name) return; - snapshot.name = name; - emit("session_info", { session: snapshot }); - }); + function buildSnapshot( + ctx: ExtensionContext, + sessionId: string, + ): SessionSnapshot { + const sm = ctx.sessionManager as { + getSessionName?: () => string | undefined; + getCwd?: () => string; + getHeader?: () => { timestamp?: unknown } | null; + }; + let name: string | null = null; + let cwd: string = ctx.cwd; + let startedAt: number | undefined = sessionStartTs.get(sessionId); + try { + name = sm.getSessionName?.() ?? null; + cwd = sm.getCwd?.() ?? ctx.cwd; + if (startedAt === undefined) { + const headerTs = sm.getHeader?.()?.timestamp; + startedAt = typeof headerTs === "string" ? Date.parse(headerTs) : NaN; + if (!Number.isFinite(startedAt)) startedAt = Date.now(); + sessionStartTs.set(sessionId, startedAt); + } + } catch (err) { + log(`snapshot fallback (${String(err).slice(0, 80)}): using ctx.cwd/now`); + if (startedAt === undefined) { + startedAt = Date.now(); + sessionStartTs.set(sessionId, startedAt); + } + } + const model = ctx.model as { id?: unknown; provider?: unknown } | undefined; + return { + id: sessionId, + name, + cwd, + model: + model !== undefined && typeof model.id === "string" ? model.id : null, + provider: + model !== undefined && typeof model.provider === "string" + ? model.provider + : null, + agent: process.env[ENV_AGENT] === "1", + repo: process.env[ENV_REPO] ?? null, + startedAt, + }; + } - sub("model_select", (event: { model?: { id?: unknown; provider?: unknown } }) => { - if (snapshot === null) return; - const model = event.model; - const id = model !== undefined && typeof model.id === "string" ? model.id : null; - const provider = model !== undefined && typeof model.provider === "string" ? model.provider : null; - if (id === snapshot.model && provider === snapshot.provider) return; - snapshot.model = id; - snapshot.provider = provider; - emit("session_info", { session: snapshot }); - }); + sub("session_start", (_event: any, ctx: ExtensionContext) => { + stopped = false; + let sessionId: string | null = null; + try { + const id = ( + ctx.sessionManager as { getSessionId?: () => string } + ).getSessionId?.(); + sessionId = typeof id === "string" && id.length > 0 ? id : null; + } catch (err) { + log(`getSessionId failed: ${err}`); + } + if (sessionId === null) { + log("no session id; lvmh mirroring disabled for this session"); + return; + } + if (sessionId !== currentSessionId) { + sendQueue = []; + replayBuf = []; + droppedEvents = 0; + lastStreamText = ""; + handleDisconnect(); // drop any socket bound to the previous session + try { + if (reconnectTimer !== null) clearTimeout(reconnectTimer); + } catch { + // clearable timers never throw in practice + } + reconnectTimer = null; + } + currentSessionId = sessionId; + snapshot = buildSnapshot(ctx, sessionId); + connect(); + }); - sub("message_start", (event: { message?: unknown }) => { - const mapped = mapMessage(event.message); - if (mapped === null) return; - if (mapped.role === "assistant") lastStreamText = ""; - emit("message_start", { message: { role: mapped.role, id: mapped.id } }); - }); + sub("session_shutdown", () => { + stopped = true; + try { + if (reconnectTimer !== null) clearTimeout(reconnectTimer); + } catch { + // ignore + } + reconnectTimer = null; + clearWelcomeTimer(); + try { + if (stallTimer !== null) clearInterval(stallTimer); + } catch { + // ignore + } + stallTimer = null; + lastStallBytes = null; + const socket: WebSocket | null = ws; + ws = null; + wsOpen = false; + greeted = false; + try { + if ( + socket !== null && + socket.readyState === WebSocket.OPEN && + currentSessionId !== null + ) { + socket.send( + JSON.stringify({ + v: PROTOCOL_VERSION, + sessionId: currentSessionId, + seq: 0, + ts: Date.now(), + type: "bye", + reason: "shutdown", + }), + ); + } + socket?.close(); + } catch { + // best-effort goodbye; nothing to flush, exit fast + } + }); - sub("message_update", (event: { message?: unknown }) => { - const mapped = mapMessage(event.message); - if (mapped === null || mapped.role !== "assistant") return; - const text = typeof mapped.text === "string" ? mapped.text : ""; - // Delta only: pi hands us accumulated text; diff against what we sent. - const delta: string = text.startsWith(lastStreamText) ? text.slice(lastStreamText.length) : text; - lastStreamText = text; - if (delta.length > 0) emit("message_update", { delta }); - }); + sub("session_info_changed", (event: { name?: unknown }) => { + if (snapshot === null) return; + const name = typeof event.name === "string" ? event.name : null; + if (name === snapshot.name) return; + snapshot.name = name; + emit("session_info", { session: snapshot }); + }); - sub("message_end", (event: { message?: unknown }) => { - const mapped = mapMessage(event.message); - if (mapped === null) return; - lastStreamText = ""; - emit("message_end", { message: mapped }); - }); + sub( + "model_select", + (event: { model?: { id?: unknown; provider?: unknown } }) => { + if (snapshot === null) return; + const model = event.model; + const id = + model !== undefined && typeof model.id === "string" ? model.id : null; + const provider = + model !== undefined && typeof model.provider === "string" + ? model.provider + : null; + if (id === snapshot.model && provider === snapshot.provider) return; + snapshot.model = id; + snapshot.provider = provider; + emit("session_info", { session: snapshot }); + }, + ); - sub("tool_execution_start", (event: { toolCallId?: unknown; toolName?: unknown; args?: unknown }) => { - emit("tool_execution_start", { - toolCallId: String(event.toolCallId ?? ""), - toolName: String(event.toolName ?? ""), - args: event.args ?? null, - }); - }); + sub("message_start", (event: { message?: unknown }) => { + const mapped = mapMessage(event.message); + if (mapped === null) return; + if (mapped.role === "assistant") lastStreamText = ""; + emit("message_start", { message: { role: mapped.role, id: mapped.id } }); + }); - sub("tool_execution_update", (event: { toolCallId?: unknown; toolName?: unknown; partialResult?: { content?: unknown } }) => { - const partial: string = textOfContent(event.partialResult?.content).slice(0, RESULT_PREVIEW_MAX_CHARS); - emit("tool_execution_update", { - toolCallId: String(event.toolCallId ?? ""), - toolName: String(event.toolName ?? ""), - partial, - }); - }); + sub("message_update", (event: { message?: unknown }) => { + const mapped = mapMessage(event.message); + if (mapped === null || mapped.role !== "assistant") return; + const text = typeof mapped.text === "string" ? mapped.text : ""; + // Delta only: pi hands us accumulated text; diff against what we sent. + const delta: string = text.startsWith(lastStreamText) + ? text.slice(lastStreamText.length) + : text; + lastStreamText = text; + if (delta.length > 0) emit("message_update", { delta }); + }); - sub("tool_execution_end", (event: { toolCallId?: unknown; toolName?: unknown; result?: { content?: unknown }; isError?: unknown }) => { - const preview: string = textOfContent(event.result?.content).slice(0, RESULT_PREVIEW_MAX_CHARS); - emit("tool_execution_end", { - toolCallId: String(event.toolCallId ?? ""), - toolName: String(event.toolName ?? ""), - isError: event.isError === true, - resultPreview: preview, - }); - }); + sub("message_end", (event: { message?: unknown }) => { + const mapped = mapMessage(event.message); + if (mapped === null) return; + lastStreamText = ""; + emit("message_end", { message: mapped }); + }); - sub("agent_start", () => { - emit("agent_start", {}); - }); + sub( + "tool_execution_start", + (event: { toolCallId?: unknown; toolName?: unknown; args?: unknown }) => { + emit("tool_execution_start", { + toolCallId: String(event.toolCallId ?? ""), + toolName: String(event.toolName ?? ""), + args: event.args ?? null, + }); + }, + ); - sub("agent_end", (event: { messages?: unknown }) => { - let inputTokens = 0; - let outputTokens = 0; - let totalCost = 0; - let seen = false; - if (Array.isArray(event.messages)) { - for (const m of event.messages) { - if (m === null || typeof m !== "object" || (m as { role?: unknown }).role !== "assistant") continue; - const u = (m as { usage?: { input?: unknown; output?: unknown; cost?: { total?: unknown } } }).usage; - if (u === undefined) continue; - seen = true; - inputTokens += typeof u.input === "number" ? u.input : 0; - outputTokens += typeof u.output === "number" ? u.output : 0; - totalCost += u.cost !== undefined && typeof u.cost.total === "number" ? u.cost.total : 0; - } - } - emit("agent_end", { usage: seen ? { inputTokens, outputTokens, totalCost } : {} }); - }); + sub( + "tool_execution_update", + (event: { + toolCallId?: unknown; + toolName?: unknown; + partialResult?: { content?: unknown }; + }) => { + const partial: string = textOfContent(event.partialResult?.content).slice( + 0, + RESULT_PREVIEW_MAX_CHARS, + ); + emit("tool_execution_update", { + toolCallId: String(event.toolCallId ?? ""), + toolName: String(event.toolName ?? ""), + partial, + }); + }, + ); - sub("agent_settled", () => { - emit("agent_settled", {}); - }); + sub( + "tool_execution_end", + (event: { + toolCallId?: unknown; + toolName?: unknown; + result?: { content?: unknown }; + isError?: unknown; + }) => { + const preview: string = textOfContent(event.result?.content).slice( + 0, + RESULT_PREVIEW_MAX_CHARS, + ); + emit("tool_execution_end", { + toolCallId: String(event.toolCallId ?? ""), + toolName: String(event.toolName ?? ""), + isError: event.isError === true, + resultPreview: preview, + }); + }, + ); + + sub("agent_start", () => { + emit("agent_start", {}); + }); + + sub("agent_end", (event: { messages?: unknown }) => { + let inputTokens = 0; + let outputTokens = 0; + let totalCost = 0; + let seen = false; + if (Array.isArray(event.messages)) { + for (const m of event.messages) { + if ( + m === null || + typeof m !== "object" || + (m as { role?: unknown }).role !== "assistant" + ) + continue; + const u = ( + m as { + usage?: { + input?: unknown; + output?: unknown; + cost?: { total?: unknown }; + }; + } + ).usage; + if (u === undefined) continue; + seen = true; + inputTokens += typeof u.input === "number" ? u.input : 0; + outputTokens += typeof u.output === "number" ? u.output : 0; + totalCost += + u.cost !== undefined && typeof u.cost.total === "number" + ? u.cost.total + : 0; + } + } + emit("agent_end", { + usage: seen ? { inputTokens, outputTokens, totalCost } : {}, + }); + }); + + sub("agent_settled", () => { + emit("agent_settled", {}); + }); } diff --git a/plugin/mini-daemon.ts b/plugin/mini-daemon.ts index c46d3db..d5e88a3 100644 --- a/plugin/mini-daemon.ts +++ b/plugin/mini-daemon.ts @@ -12,171 +12,192 @@ import type { Duplex } from "node:stream"; export const WS_GUID = "258EAFA5-E914-47DA-95CA-C5AB0DC85B11"; export interface ClientFrame { - v: number; - sessionId: string; - seq: number; - ts: number; - type: string; - [key: string]: unknown; + v: number; + sessionId: string; + seq: number; + ts: number; + type: string; + [key: string]: unknown; } export interface MiniDaemon { - url: string; - frames: ClientFrame[]; - authHeaders: string[]; - connections(): number; - pushAll(text: string): void; - dropConnections(): void; - close(): void; + url: string; + frames: ClientFrame[]; + authHeaders: string[]; + connections(): number; + pushAll(text: string): void; + dropConnections(): void; + /** Pause reads on all live sockets: simulate a daemon wedged without TCP close. */ + wedge(): void; + unwedge(): void; + close(): void; } export function check(name: string, ok: boolean, detail = ""): void { - if (ok) console.log(`ok ${name}`); - else console.log(`FAIL ${name} ${detail}`); + if (ok) console.log(`ok ${name}`); + else console.log(`FAIL ${name} ${detail}`); } export function sleep(ms: number): Promise { - return new Promise((resolve) => setTimeout(resolve, ms)); + return new Promise((resolve) => setTimeout(resolve, ms)); } -export async function waitFor(cond: () => boolean, timeoutMs: number, stepMs = 25): Promise { - const deadline = Date.now() + timeoutMs; - while (Date.now() < deadline) { - if (cond()) return true; - await sleep(stepMs); - } - return cond(); +export async function waitFor( + cond: () => boolean, + timeoutMs: number, + stepMs = 25, +): Promise { + const deadline = Date.now() + timeoutMs; + while (Date.now() < deadline) { + if (cond()) return true; + await sleep(stepMs); + } + return cond(); } function encodeTextFrame(text: string): Buffer { - const payload = Buffer.from(text, "utf8"); - const len = payload.length; - let header: Buffer; - if (len < 126) header = Buffer.from([0x81, len]); - else if (len < 65536) { - header = Buffer.alloc(4); - header[0] = 0x81; - header[1] = 126; - header.writeUInt16BE(len, 2); - } else { - header = Buffer.alloc(10); - header[0] = 0x81; - header[1] = 127; - header.writeBigUInt64BE(BigInt(len), 2); - } - return Buffer.concat([header, payload]); + const payload = Buffer.from(text, "utf8"); + const len = payload.length; + let header: Buffer; + if (len < 126) header = Buffer.from([0x81, len]); + else if (len < 65536) { + header = Buffer.alloc(4); + header[0] = 0x81; + header[1] = 126; + header.writeUInt16BE(len, 2); + } else { + header = Buffer.alloc(10); + header[0] = 0x81; + header[1] = 127; + header.writeBigUInt64BE(BigInt(len), 2); + } + return Buffer.concat([header, payload]); } -function decodeFrames(chunk: Buffer): { frames: Array<{ opcode: number; data: Buffer }>; consumed: number } { - const frames: Array<{ opcode: number; data: Buffer }> = []; - let offset = 0; - while (offset + 2 <= chunk.length) { - const opcode = chunk[offset] & 0x0f; - const masked = (chunk[offset + 1] & 0x80) !== 0; - let len = chunk[offset + 1] & 0x7f; - let cursor = offset + 2; - if (len === 126) { - if (cursor + 2 > chunk.length) break; - len = chunk.readUInt16BE(cursor); - cursor += 2; - } else if (len === 127) { - if (cursor + 8 > chunk.length) break; - len = Number(chunk.readBigUInt64BE(cursor)); - cursor += 8; - } - let mask: Buffer | null = null; - if (masked) { - if (cursor + 4 > chunk.length) break; - mask = chunk.subarray(cursor, cursor + 4); - cursor += 4; - } - if (cursor + len > chunk.length) break; - let data = chunk.subarray(cursor, cursor + len); - if (mask !== null) { - const unmasked = Buffer.allocUnsafe(len); - for (let i = 0; i < len; i++) unmasked[i] = data[i] ^ mask[i % 4]; - data = unmasked; - } - frames.push({ opcode, data }); - offset = cursor + len; - } - return { frames, consumed: offset }; +function decodeFrames(chunk: Buffer): { + frames: Array<{ opcode: number; data: Buffer }>; + consumed: number; +} { + const frames: Array<{ opcode: number; data: Buffer }> = []; + let offset = 0; + while (offset + 2 <= chunk.length) { + const opcode = chunk[offset] & 0x0f; + const masked = (chunk[offset + 1] & 0x80) !== 0; + let len = chunk[offset + 1] & 0x7f; + let cursor = offset + 2; + if (len === 126) { + if (cursor + 2 > chunk.length) break; + len = chunk.readUInt16BE(cursor); + cursor += 2; + } else if (len === 127) { + if (cursor + 8 > chunk.length) break; + len = Number(chunk.readBigUInt64BE(cursor)); + cursor += 8; + } + let mask: Buffer | null = null; + if (masked) { + if (cursor + 4 > chunk.length) break; + mask = chunk.subarray(cursor, cursor + 4); + cursor += 4; + } + if (cursor + len > chunk.length) break; + let data = chunk.subarray(cursor, cursor + len); + if (mask !== null) { + const unmasked = Buffer.allocUnsafe(len); + for (let i = 0; i < len; i++) unmasked[i] = data[i] ^ mask[i % 4]; + data = unmasked; + } + frames.push({ opcode, data }); + offset = cursor + len; + } + return { frames, consumed: offset }; } -export function startMiniDaemon(getLastSeq: () => number): Promise { - const server = http.createServer(); - const sockets = new Set(); - const daemon: MiniDaemon = { - url: "", - frames: [], - authHeaders: [], - connections: () => sockets.size, - pushAll(text: string) { - const frame = encodeTextFrame(text); - for (const s of sockets) s.write(frame); - }, - dropConnections() { - for (const s of sockets) s.destroy(); - sockets.clear(); - }, - close() { - daemon.dropConnections(); - server.close(); - }, - }; - server.on("upgrade", (req: http.IncomingMessage, socket: Duplex) => { - daemon.authHeaders.push(String(req.headers.authorization ?? "")); - const key = String(req.headers["sec-websocket-key"] ?? ""); - const accept = createHash("sha1").update(key + WS_GUID).digest("base64"); - socket.write( - "HTTP/1.1 101 Switching Protocols\r\nUpgrade: websocket\r\nConnection: Upgrade\r\n" + - `Sec-WebSocket-Accept: ${accept}\r\n\r\n`, - ); - sockets.add(socket); - // Frames can split across TCP chunks: buffer per socket and decode only - // complete frames, retaining the remainder for the next chunk. - let recvBuf = Buffer.alloc(0); - socket.on("data", (chunk: Buffer) => { - recvBuf = Buffer.concat([recvBuf, chunk]); - const { frames: decoded, consumed } = decodeFrames(recvBuf); - recvBuf = recvBuf.subarray(consumed); - for (const f of decoded) { - if (f.opcode === 0x8) { - // proper close handshake: echo close frame, then end the socket - socket.write(Buffer.from([0x88, 0x02, 0x03, 0xe8])); - socket.end(); - continue; - } - if (f.opcode !== 0x1) continue; - try { - const frame = JSON.parse(f.data.toString("utf8")) as ClientFrame; - daemon.frames.push(frame); - if (frame.type === "hello") { - socket.write( - encodeTextFrame( - JSON.stringify({ - v: 1, - type: "welcome", - sessionId: frame.sessionId, - seq: 0, - ts: Date.now(), - lastSeq: getLastSeq(), - }), - ), - ); - } - } catch { - // malformed client frame: ignore - } - } - }); - socket.on("close", () => sockets.delete(socket)); - socket.on("error", () => sockets.delete(socket)); - }); - return new Promise((resolve) => { - server.listen(0, "127.0.0.1", () => { - daemon.url = `ws://127.0.0.1:${(server.address() as AddressInfo).port}/agent/ws`; - resolve(daemon); - }); - }); +export function startMiniDaemon( + getLastSeq: () => number, + opts: { suppressWelcome?: boolean } = {}, +): Promise { + const server = http.createServer(); + const sockets = new Set(); + const daemon: MiniDaemon = { + url: "", + frames: [], + authHeaders: [], + connections: () => sockets.size, + pushAll(text: string) { + const frame = encodeTextFrame(text); + for (const s of sockets) s.write(frame); + }, + dropConnections() { + for (const s of sockets) s.destroy(); + sockets.clear(); + }, + wedge() { + for (const s of sockets) s.pause(); + }, + unwedge() { + for (const s of sockets) s.resume(); + }, + close() { + daemon.dropConnections(); + server.close(); + }, + }; + server.on("upgrade", (req: http.IncomingMessage, socket: Duplex) => { + daemon.authHeaders.push(String(req.headers.authorization ?? "")); + const key = String(req.headers["sec-websocket-key"] ?? ""); + const accept = createHash("sha1") + .update(key + WS_GUID) + .digest("base64"); + socket.write( + "HTTP/1.1 101 Switching Protocols\r\nUpgrade: websocket\r\nConnection: Upgrade\r\n" + + `Sec-WebSocket-Accept: ${accept}\r\n\r\n`, + ); + sockets.add(socket); + // Frames can split across TCP chunks: buffer per socket and decode only + // complete frames, retaining the remainder for the next chunk. + let recvBuf = Buffer.alloc(0); + socket.on("data", (chunk: Buffer) => { + recvBuf = Buffer.concat([recvBuf, chunk]); + const { frames: decoded, consumed } = decodeFrames(recvBuf); + recvBuf = recvBuf.subarray(consumed); + for (const f of decoded) { + if (f.opcode === 0x8) { + // proper close handshake: echo close frame, then end the socket + socket.write(Buffer.from([0x88, 0x02, 0x03, 0xe8])); + socket.end(); + continue; + } + if (f.opcode !== 0x1) continue; + try { + const frame = JSON.parse(f.data.toString("utf8")) as ClientFrame; + daemon.frames.push(frame); + if (frame.type === "hello" && !opts.suppressWelcome) { + socket.write( + encodeTextFrame( + JSON.stringify({ + v: 1, + type: "welcome", + sessionId: frame.sessionId, + seq: 0, + ts: Date.now(), + lastSeq: getLastSeq(), + }), + ), + ); + } + } catch { + // malformed client frame: ignore + } + } + }); + socket.on("close", () => sockets.delete(socket)); + socket.on("error", () => sockets.delete(socket)); + }); + return new Promise((resolve) => { + server.listen(0, "127.0.0.1", () => { + daemon.url = `ws://127.0.0.1:${(server.address() as AddressInfo).port}/agent/ws`; + resolve(daemon); + }); + }); } diff --git a/plugin/smoke.ts b/plugin/smoke.ts index ac836a8..25e4353 100644 --- a/plugin/smoke.ts +++ b/plugin/smoke.ts @@ -17,312 +17,530 @@ * 6. session_shutdown -> socket closed, no reconnect afterwards */ -import { check as rawCheck, sleep, waitFor, startMiniDaemon } from "./mini-daemon.ts"; +import { + check as rawCheck, + sleep, + waitFor, + startMiniDaemon, +} from "./mini-daemon.ts"; const SESSION_ID = "sess-1"; let failures = 0; function check(name: string, ok: boolean, detail = ""): void { - rawCheck(name, ok, detail); - if (!ok) failures++; + rawCheck(name, ok, detail); + if (!ok) failures++; } interface FakePi { - handlers: Map unknown>; - sentMessages: Array<{ message: string; options: unknown }>; - aborted: number; - on(event: string, handler: (event: unknown, ctx: unknown) => unknown): void; - sendUserMessage(message: string, options?: unknown): void; - abort(): void; + handlers: Map unknown>; + sentMessages: Array<{ message: string; options: unknown }>; + aborted: number; + on(event: string, handler: (event: unknown, ctx: unknown) => unknown): void; + sendUserMessage(message: string, options?: unknown): void; + abort(): void; } function makeFakePi(): FakePi { - return { - handlers: new Map(), - sentMessages: [], - aborted: 0, - on(event, handler) { - this.handlers.set(event, handler); - }, - sendUserMessage(message, options) { - this.sentMessages.push({ message, options }); - }, - abort() { - this.aborted++; - }, - }; + return { + handlers: new Map(), + sentMessages: [], + aborted: 0, + on(event, handler) { + this.handlers.set(event, handler); + }, + sendUserMessage(message, options) { + this.sentMessages.push({ message, options }); + }, + abort() { + this.aborted++; + }, + }; } function makeFakeCtx(): unknown { - return { - cwd: "/work/repo", - model: { id: "glm-5.3", provider: "zai-renaud" }, - sessionManager: { - getSessionId: () => SESSION_ID, - getSessionName: () => undefined, - getCwd: () => "/work/repo", - getHeader: () => ({ timestamp: "2024-12-03T14:00:00.000Z", id: SESSION_ID }), - }, - }; + return { + cwd: "/work/repo", + model: { id: "glm-5.3", provider: "zai-renaud" }, + sessionManager: { + getSessionId: () => SESSION_ID, + getSessionName: () => undefined, + getCwd: () => "/work/repo", + getHeader: () => ({ + timestamp: "2024-12-03T14:00:00.000Z", + id: SESSION_ID, + }), + }, + }; } async function loadExtension(): Promise<(pi: unknown) => void> { - const mod = (await import("./lvmh-agent.ts")) as { default: (pi: unknown) => void }; - return mod.default; + const mod = (await import("./lvmh-agent.ts")) as { + default: (pi: unknown) => void; + }; + return mod.default; } async function main(): Promise { - // --- Scenario 1: env unset -> inert -------------------------------- - delete process.env.LVMH_URL; - delete process.env.LVMH_TOKEN; - delete process.env.LVMH_AGENT; - delete process.env.LVMH_REPO; - const factory = await loadExtension(); - const inertPi = makeFakePi(); - inertPi.on = ((event: string) => { - throw new Error(`inert extension registered handler ${event}`); - }) as FakePi["on"]; - factory(inertPi); // must not touch pi at all - check("1 inert: no side effects on pi", true); + // --- Scenario 1: env unset -> inert -------------------------------- + delete process.env.LVMH_URL; + delete process.env.LVMH_TOKEN; + delete process.env.LVMH_AGENT; + delete process.env.LVMH_REPO; + const factory = await loadExtension(); + const inertPi = makeFakePi(); + inertPi.on = ((event: string) => { + throw new Error(`inert extension registered handler ${event}`); + }) as FakePi["on"]; + factory(inertPi); // must not touch pi at all + check("1 inert: no side effects on pi", true); - // --- Scenario 2: live daemon --------------------------------------- - let welcomeLastSeq = 0; - const daemon = await startMiniDaemon(() => welcomeLastSeq); - process.env.LVMH_URL = daemon.url; - process.env.LVMH_TOKEN = "smoke-token"; + // --- Scenario 2: live daemon --------------------------------------- + let welcomeLastSeq = 0; + const daemon = await startMiniDaemon(() => welcomeLastSeq); + process.env.LVMH_URL = daemon.url; + process.env.LVMH_TOKEN = "smoke-token"; - const pi = makeFakePi(); - factory(pi); - const handlerNames = [ - "session_start", - "session_shutdown", - "session_info_changed", - "model_select", - "message_start", - "message_update", - "message_end", - "tool_execution_start", - "tool_execution_update", - "tool_execution_end", - "agent_start", - "agent_end", - "agent_settled", - ]; - check("2 all 13 handlers registered", handlerNames.every((n) => pi.handlers.has(n))); + const pi = makeFakePi(); + factory(pi); + const handlerNames = [ + "session_start", + "session_shutdown", + "session_info_changed", + "model_select", + "message_start", + "message_update", + "message_end", + "tool_execution_start", + "tool_execution_update", + "tool_execution_end", + "agent_start", + "agent_end", + "agent_settled", + ]; + check( + "2 all 13 handlers registered", + handlerNames.every((n) => pi.handlers.has(n)), + ); - let lastReturn: unknown = "sentinel"; - const fire = (name: string, event: unknown): void => { - const h = pi.handlers.get(name); - if (h === undefined) throw new Error(`missing handler ${name}`); - lastReturn = h(event, makeFakeCtx()); - if (lastReturn instanceof Promise) lastReturn.catch(() => undefined); - }; + let lastReturn: unknown = "sentinel"; + const fire = (name: string, event: unknown): void => { + const h = pi.handlers.get(name); + if (h === undefined) throw new Error(`missing handler ${name}`); + lastReturn = h(event, makeFakeCtx()); + if (lastReturn instanceof Promise) lastReturn.catch(() => undefined); + }; - fire("session_start", { reason: "startup" }); - const helloSeen = await waitFor(() => daemon.frames.some((f) => f.type === "hello"), 5000); - check("2 connects and sends hello", helloSeen); - check("2 bearer auth on upgrade", daemon.authHeaders.at(-1) === "Bearer smoke-token", daemon.authHeaders.at(-1)); + fire("session_start", { reason: "startup" }); + const helloSeen = await waitFor( + () => daemon.frames.some((f) => f.type === "hello"), + 5000, + ); + check("2 connects and sends hello", helloSeen); + check( + "2 bearer auth on upgrade", + daemon.authHeaders.at(-1) === "Bearer smoke-token", + daemon.authHeaders.at(-1), + ); - const hello = daemon.frames.find((f) => f.type === "hello"); - const hs = (hello?.session ?? {}) as Record; - check( - "2 hello snapshot", - hs.id === SESSION_ID && - hs.name === null && - hs.cwd === "/work/repo" && - hs.model === "glm-5.3" && - hs.provider === "zai-renaud" && - hs.agent === false && - hs.repo === null && - typeof hs.startedAt === "number", - JSON.stringify(hs), - ); + const hello = daemon.frames.find((f) => f.type === "hello"); + const hs = (hello?.session ?? {}) as Record; + check( + "2 hello snapshot", + hs.id === SESSION_ID && + hs.name === null && + hs.cwd === "/work/repo" && + hs.model === "glm-5.3" && + hs.provider === "zai-renaud" && + hs.agent === false && + hs.repo === null && + typeof hs.startedAt === "number", + JSON.stringify(hs), + ); - fire("message_start", { message: { role: "assistant", content: [], timestamp: 1000 } }); - fire("message_update", { message: { role: "assistant", content: [{ type: "text", text: "Hel" }], timestamp: 1000 } }); - fire("message_update", { message: { role: "assistant", content: [{ type: "text", text: "Hello world" }], timestamp: 1000 } }); - fire("message_end", { - message: { - role: "assistant", - content: [ - { type: "thinking", thinking: "hmm" }, - { type: "text", text: "Hello world" }, - { type: "toolCall", id: "tc1", name: "bash", arguments: { command: "ls" } }, - ], - timestamp: 1000, - }, - }); - await waitFor(() => daemon.frames.some((f) => f.type === "message_end"), 2000); + fire("message_start", { + message: { role: "assistant", content: [], timestamp: 1000 }, + }); + fire("message_update", { + message: { + role: "assistant", + content: [{ type: "text", text: "Hel" }], + timestamp: 1000, + }, + }); + fire("message_update", { + message: { + role: "assistant", + content: [{ type: "text", text: "Hello world" }], + timestamp: 1000, + }, + }); + fire("message_end", { + message: { + role: "assistant", + content: [ + { type: "thinking", thinking: "hmm" }, + { type: "text", text: "Hello world" }, + { + type: "toolCall", + id: "tc1", + name: "bash", + arguments: { command: "ls" }, + }, + ], + timestamp: 1000, + }, + }); + await waitFor( + () => daemon.frames.some((f) => f.type === "message_end"), + 2000, + ); - const n0 = daemon.frames.findIndex((f) => f.type === "hello"); - const mirrored = daemon.frames.slice(n0 + 1).map((f) => f.type); - check( - "2 mirror order", - JSON.stringify(mirrored) === JSON.stringify(["message_start", "message_update", "message_update", "message_end"]), - JSON.stringify(mirrored), - ); - const deltas = daemon.frames.filter((f) => f.type === "message_update").map((f) => f.delta); - check("2 delta-only updates", JSON.stringify(deltas) === JSON.stringify(["Hel", "lo world"]), JSON.stringify(deltas)); + const n0 = daemon.frames.findIndex((f) => f.type === "hello"); + const mirrored = daemon.frames.slice(n0 + 1).map((f) => f.type); + check( + "2 mirror order", + JSON.stringify(mirrored) === + JSON.stringify([ + "message_start", + "message_update", + "message_update", + "message_end", + ]), + JSON.stringify(mirrored), + ); + const deltas = daemon.frames + .filter((f) => f.type === "message_update") + .map((f) => f.delta); + check( + "2 delta-only updates", + JSON.stringify(deltas) === JSON.stringify(["Hel", "lo world"]), + JSON.stringify(deltas), + ); - const msg = daemon.frames.find((f) => f.type === "message_end")?.message as Record; - const tc = (msg?.toolCalls as Array> | undefined)?.[0]; - check( - "2 message_end mapping", - msg?.role === "assistant" && - msg?.text === "Hello world" && - msg?.thinking === "hmm" && - typeof msg?.id === "string" && - tc?.id === "tc1" && - tc?.name === "bash" && - tc?.argsJson === '{"command":"ls"}' && - msg?.toolCallId === null, - JSON.stringify(msg), - ); + const msg = daemon.frames.find((f) => f.type === "message_end") + ?.message as Record; + const tc = ( + msg?.toolCalls as Array> | undefined + )?.[0]; + check( + "2 message_end mapping", + msg?.role === "assistant" && + msg?.text === "Hello world" && + msg?.thinking === "hmm" && + typeof msg?.id === "string" && + tc?.id === "tc1" && + tc?.name === "bash" && + tc?.argsJson === '{"command":"ls"}' && + msg?.toolCallId === null, + JSON.stringify(msg), + ); - const big = "x".repeat(3000); - fire("agent_start", {}); - fire("tool_execution_start", { toolCallId: "tc1", toolName: "bash", args: { command: "ls" } }); - fire("tool_execution_update", { toolCallId: "tc1", toolName: "bash", partialResult: { content: [{ type: "text", text: big }] } }); - fire("tool_execution_end", { toolCallId: "tc1", toolName: "bash", isError: true, result: { content: [{ type: "text", text: big }] } }); - fire("agent_end", { messages: [{ role: "assistant", usage: { input: 10, output: 5, cost: { total: 0.25 } } }] }); - fire("agent_settled", {}); - await waitFor(() => daemon.frames.some((f) => f.type === "agent_settled"), 2000); + const big = "x".repeat(3000); + fire("agent_start", {}); + fire("tool_execution_start", { + toolCallId: "tc1", + toolName: "bash", + args: { command: "ls" }, + }); + fire("tool_execution_update", { + toolCallId: "tc1", + toolName: "bash", + partialResult: { content: [{ type: "text", text: big }] }, + }); + fire("tool_execution_end", { + toolCallId: "tc1", + toolName: "bash", + isError: true, + result: { content: [{ type: "text", text: big }] }, + }); + fire("agent_end", { + messages: [ + { + role: "assistant", + usage: { input: 10, output: 5, cost: { total: 0.25 } }, + }, + ], + }); + fire("agent_settled", {}); + await waitFor( + () => daemon.frames.some((f) => f.type === "agent_settled"), + 2000, + ); - const partial = daemon.frames.find((f) => f.type === "tool_execution_update"); - const end = daemon.frames.find((f) => f.type === "tool_execution_end"); - const partialLen = (partial?.partial as string | undefined)?.length; - const previewLen = (end?.resultPreview as string | undefined)?.length; - check("2 partial truncated to 2000", partialLen === 2000, String(partialLen)); - check("2 resultPreview truncated + isError", previewLen === 2000 && end?.isError === true); - const usage = daemon.frames.find((f) => f.type === "agent_end")?.usage as Record; - check( - "2 agent_end usage", - usage?.inputTokens === 10 && usage?.outputTokens === 5 && usage?.totalCost === 0.25, - JSON.stringify(usage), - ); + const partial = daemon.frames.find((f) => f.type === "tool_execution_update"); + const end = daemon.frames.find((f) => f.type === "tool_execution_end"); + const partialLen = (partial?.partial as string | undefined)?.length; + const previewLen = (end?.resultPreview as string | undefined)?.length; + check("2 partial truncated to 2000", partialLen === 2000, String(partialLen)); + check( + "2 resultPreview truncated + isError", + previewLen === 2000 && end?.isError === true, + ); + const usage = daemon.frames.find((f) => f.type === "agent_end") + ?.usage as Record; + check( + "2 agent_end usage", + usage?.inputTokens === 10 && + usage?.outputTokens === 5 && + usage?.totalCost === 0.25, + JSON.stringify(usage), + ); - const seqs = daemon.frames.filter((f) => f.seq > 0).map((f) => f.seq); - check( - "2 seq strictly monotonic", - seqs.length > 0 && seqs.every((s, i) => i === 0 || s > seqs[i - 1]), - JSON.stringify(seqs), - ); - check("2 handler return values are undefined (never a promise)", lastReturn === undefined); + const seqs = daemon.frames.filter((f) => f.seq > 0).map((f) => f.seq); + check( + "2 seq strictly monotonic", + seqs.length > 0 && seqs.every((s, i) => i === 0 || s > seqs[i - 1]), + JSON.stringify(seqs), + ); + check( + "2 handler return values are undefined (never a promise)", + lastReturn === undefined, + ); - // --- Scenario 3: prompt / abort ------------------------------------ - daemon.pushAll( - JSON.stringify({ v: 1, type: "prompt", sessionId: "other-session", seq: 0, ts: Date.now(), promptId: "p0", message: "wrong session" }), - ); - daemon.pushAll( - JSON.stringify({ v: 1, type: "prompt", sessionId: SESSION_ID, seq: 0, ts: Date.now(), promptId: "p1", message: "run the tests" }), - ); - await waitFor(() => pi.sentMessages.length > 0, 2000); - check( - "3 prompt delivered as steer, foreign session ignored", - pi.sentMessages.length === 1 && - pi.sentMessages[0].message === "run the tests" && - JSON.stringify(pi.sentMessages[0].options) === '{"deliverAs":"steer"}', - JSON.stringify(pi.sentMessages), - ); - daemon.pushAll(JSON.stringify({ v: 1, type: "abort", sessionId: SESSION_ID, seq: 0, ts: Date.now() })); - await waitFor(() => pi.aborted > 0, 2000); - check("3 abort calls pi.abort", pi.aborted === 1); + // --- Scenario 3: prompt / abort ------------------------------------ + daemon.pushAll( + JSON.stringify({ + v: 1, + type: "prompt", + sessionId: "other-session", + seq: 0, + ts: Date.now(), + promptId: "p0", + message: "wrong session", + }), + ); + daemon.pushAll( + JSON.stringify({ + v: 1, + type: "prompt", + sessionId: SESSION_ID, + seq: 0, + ts: Date.now(), + promptId: "p1", + message: "run the tests", + }), + ); + await waitFor(() => pi.sentMessages.length > 0, 2000); + check( + "3 prompt delivered as steer, foreign session ignored", + pi.sentMessages.length === 1 && + pi.sentMessages[0].message === "run the tests" && + JSON.stringify(pi.sentMessages[0].options) === '{"deliverAs":"steer"}', + JSON.stringify(pi.sentMessages), + ); + daemon.pushAll( + JSON.stringify({ + v: 1, + type: "abort", + sessionId: SESSION_ID, + seq: 0, + ts: Date.now(), + }), + ); + await waitFor(() => pi.aborted > 0, 2000); + check("3 abort calls pi.abort", pi.aborted === 1); - // --- Scenario 4: drop + reconnect, replay, overflow ---------------- - const seqBeforeDrop = Math.max(...daemon.frames.map((f) => f.seq)); - welcomeLastSeq = seqBeforeDrop; // daemon has everything up to the drop - daemon.dropConnections(); - await sleep(150); // let the client notice - fire("message_end", { message: { role: "user", content: "offline-1", timestamp: 2000 } }); - fire("message_end", { message: { role: "user", content: "offline-2", timestamp: 2001 } }); - const reconnected = await waitFor( - () => daemon.frames.filter((f) => f.type === "hello").length >= 2, - 10000, - ); - check("4 reconnects after drop", reconnected); - const replayedTexts = daemon.frames - .filter((f) => f.type === "message_end" && f.seq > seqBeforeDrop) - .map((f) => (f.message as { text?: string }).text); - check( - "4 buffered events replayed after welcome.lastSeq", - replayedTexts.includes("offline-1") && replayedTexts.includes("offline-2"), - JSON.stringify(replayedTexts), - ); - const allSeqs = daemon.frames.filter((f) => f.seq > 0).map((f) => f.seq); - const dupCount = allSeqs.length - new Set(allSeqs).size; - check("4 no duplicate seq delivery", dupCount === 0, `duplicates: ${dupCount}`); - const seqsAfter = daemon.frames.filter((f) => f.seq > seqBeforeDrop).map((f) => f.seq); - check( - "4 seq continues monotonically across reconnect", - seqsAfter.length >= 2 && seqsAfter.every((s, i) => i === 0 || s > seqsAfter[i - 1]), - JSON.stringify(seqsAfter), - ); + // --- Scenario 4: drop + reconnect, replay, overflow ---------------- + const seqBeforeDrop = Math.max(...daemon.frames.map((f) => f.seq)); + welcomeLastSeq = seqBeforeDrop; // daemon has everything up to the drop + daemon.dropConnections(); + await sleep(150); // let the client notice + fire("message_end", { + message: { role: "user", content: "offline-1", timestamp: 2000 }, + }); + fire("message_end", { + message: { role: "user", content: "offline-2", timestamp: 2001 }, + }); + const reconnected = await waitFor( + () => daemon.frames.filter((f) => f.type === "hello").length >= 2, + 10000, + ); + check("4 reconnects after drop", reconnected); + const replayedTexts = daemon.frames + .filter((f) => f.type === "message_end" && f.seq > seqBeforeDrop) + .map((f) => (f.message as { text?: string }).text); + check( + "4 buffered events replayed after welcome.lastSeq", + replayedTexts.includes("offline-1") && replayedTexts.includes("offline-2"), + JSON.stringify(replayedTexts), + ); + const allSeqs = daemon.frames.filter((f) => f.seq > 0).map((f) => f.seq); + const dupCount = allSeqs.length - new Set(allSeqs).size; + check( + "4 no duplicate seq delivery", + dupCount === 0, + `duplicates: ${dupCount}`, + ); + const seqsAfter = daemon.frames + .filter((f) => f.seq > seqBeforeDrop) + .map((f) => f.seq); + check( + "4 seq continues monotonically across reconnect", + seqsAfter.length >= 2 && + seqsAfter.every((s, i) => i === 0 || s > seqsAfter[i - 1]), + JSON.stringify(seqsAfter), + ); - // overflow both bounded buffers, then reconnect once more - const seqBeforeOverflow = Math.max(...daemon.frames.map((f) => f.seq)); - welcomeLastSeq = seqBeforeOverflow; - daemon.dropConnections(); - await sleep(150); - for (let i = 0; i < 10050; i++) { - fire("agent_settled", {}); // persisted kind: fills replay buffer + send queue - } - const reconnected2 = await waitFor( - () => daemon.frames.filter((f) => f.type === "hello").length >= 3, - 15000, - ); - check("4 reconnects after overflow window", reconnected2); - const overflowNotice = await waitFor( - () => daemon.frames.some((f) => f.type === "buffer_overflow"), - 5000, - ); - const overflow = daemon.frames.find((f) => f.type === "buffer_overflow"); - check( - "4 buffer_overflow notice after drops", - overflowNotice && typeof (overflow?.dropped as number) === "number" && (overflow?.dropped as number) > 0, - JSON.stringify(overflow ?? null), - ); + // overflow both bounded buffers, then reconnect once more + const seqBeforeOverflow = Math.max(...daemon.frames.map((f) => f.seq)); + welcomeLastSeq = seqBeforeOverflow; + daemon.dropConnections(); + await sleep(150); + for (let i = 0; i < 10050; i++) { + fire("agent_settled", {}); // persisted kind: fills replay buffer + send queue + } + const reconnected2 = await waitFor( + () => daemon.frames.filter((f) => f.type === "hello").length >= 3, + 15000, + ); + check("4 reconnects after overflow window", reconnected2); + const overflowNotice = await waitFor( + () => daemon.frames.some((f) => f.type === "buffer_overflow"), + 5000, + ); + const overflow = daemon.frames.find((f) => f.type === "buffer_overflow"); + check( + "4 buffer_overflow notice after drops", + overflowNotice && + typeof (overflow?.dropped as number) === "number" && + (overflow?.dropped as number) > 0, + JSON.stringify(overflow ?? null), + ); - // --- Scenario 6: session_shutdown ---------------------------------- - const helloCountAtShutdown = daemon.frames.filter((f) => f.type === "hello").length; - fire("session_shutdown", { reason: "quit" }); - const closed = await waitFor(() => daemon.connections() === 0, 3000); - check("6 shutdown closes connection", closed); - await sleep(2600); // longer than one backoff attempt (~1-2s) - const helloCountAfter = daemon.frames.filter((f) => f.type === "hello").length; - check("6 no reconnect after shutdown", helloCountAfter === helloCountAtShutdown); + // --- Scenario 9: daemon hangs mid-handshake (no welcome) ---------- + // Fast watchdog knobs; read lazily by the plugin so late env wins. + process.env.LVMH_WELCOME_TIMEOUT_MS = "400"; + process.env.LVMH_STALL_CHECK_MS = "120"; + process.env.LVMH_STALL_BYTES = "1000000"; + const hungDaemon = await startMiniDaemon(() => 0, { suppressWelcome: true }); + const prevUrl = process.env.LVMH_URL; + process.env.LVMH_URL = hungDaemon.url; + const pi3 = makeFakePi(); + factory(pi3); + pi3.handlers.get("session_start")?.({ reason: "startup" }, makeFakeCtx()); + const hungHellos = await waitFor( + () => hungDaemon.frames.filter((f) => f.type === "hello").length >= 2, + 12000, + ); + check( + "9 welcome timeout: abandons hung socket and retries", + hungHellos, + String(hungDaemon.frames.length), + ); + pi3.handlers.get("session_shutdown")?.({ reason: "quit" }, makeFakeCtx()); + hungDaemon.close(); + process.env.LVMH_URL = prevUrl; - daemon.close(); + // --- Scenario 10: daemon dies without TCP close (wedged socket) ---- + // wedge() pauses daemon-side reads: TCP stays "open", undici keeps + // buffering sends. The plugin must detect bufferedAmount not draining + // and reconnect, or mirroring would be dead forever + memory unbounded. + // Fresh plugin instance: its stall interval inherits the fast knobs set + // in scenario 9 (interval period is fixed at first connect). + const stallDaemon = await startMiniDaemon(() => 0); + const urlBeforeStall: string | undefined = process.env.LVMH_URL; + process.env.LVMH_URL = stallDaemon.url; + const pi4 = makeFakePi(); + factory(pi4); + const fire4 = (name: string, event: unknown): void => { + const h = pi4.handlers.get(name); + if (h === undefined) throw new Error(`missing handler ${name}`); + const r = h(event, makeFakeCtx()); + if (r instanceof Promise) r.catch(() => undefined); + }; + fire4("session_start", { reason: "startup" }); + await waitFor(() => stallDaemon.frames.some((f) => f.type === "hello"), 5000); + stallDaemon.wedge(); + const bigText: string = "w".repeat(1024 * 1024); + for (let i = 0; i < 60; i++) { + fire4("message_end", { + message: { role: "user", content: bigText, timestamp: 3000 + i }, + }); + } + const wedgedRecovered = await waitFor( + () => stallDaemon.frames.filter((f) => f.type === "hello").length >= 2, + 20000, + ); + check( + "10 stall watchdog: wedged socket detected, reconnected", + wedgedRecovered, + String(stallDaemon.frames.length), + ); + stallDaemon.unwedge(); + fire4("message_end", { + message: { role: "user", content: "after-stall", timestamp: 4000 }, + }); + const afterStallDelivered = await waitFor( + () => + stallDaemon.frames.some( + (f) => + f.type === "message_end" && + (f.message as { text?: string })?.text === "after-stall", + ), + 15000, + ); + check("10 mirroring works after stall recovery", afterStallDelivered); + fire4("session_shutdown", { reason: "quit" }); + stallDaemon.close(); + process.env.LVMH_URL = urlBeforeStall; + delete process.env.LVMH_WELCOME_TIMEOUT_MS; + delete process.env.LVMH_STALL_CHECK_MS; + delete process.env.LVMH_STALL_BYTES; - // --- Scenario 5: dead host ----------------------------------------- - const deadDaemon = await startMiniDaemon(() => 0); - const deadUrl = deadDaemon.url; - deadDaemon.close(); // port now closed - await sleep(100); - process.env.LVMH_URL = deadUrl; - const pi2 = makeFakePi(); - factory(pi2); - pi2.handlers.get("session_start")?.({ reason: "startup" }, makeFakeCtx()); - await sleep(400); - const t0 = Date.now(); - await sleep(100); - const loopResponsive = Date.now() - t0 < 500; - check("5 dead host: pi loads, event loop responsive", loopResponsive); - let threw = false; - try { - pi2.handlers.get("message_end")?.({ message: { role: "user", content: "hi", timestamp: 1 } }, makeFakeCtx()); - } catch { - threw = true; - } - check("5 dead host: handlers never throw", !threw); + // --- Scenario 6: session_shutdown ---------------------------------- + const helloCountAtShutdown = daemon.frames.filter( + (f) => f.type === "hello", + ).length; + fire("session_shutdown", { reason: "quit" }); + const closed = await waitFor(() => daemon.connections() === 0, 3000); + check("6 shutdown closes connection", closed); + await sleep(2600); // longer than one backoff attempt (~1-2s) + const helloCountAfter = daemon.frames.filter( + (f) => f.type === "hello", + ).length; + check( + "6 no reconnect after shutdown", + helloCountAfter === helloCountAtShutdown, + ); - delete process.env.LVMH_URL; - delete process.env.LVMH_TOKEN; - console.log(failures === 0 ? "\nALL CHECKS PASSED" : `\n${failures} CHECK(S) FAILED`); - process.exit(failures === 0 ? 0 : 1); + daemon.close(); + + // --- Scenario 5: dead host ----------------------------------------- + const deadDaemon = await startMiniDaemon(() => 0); + const deadUrl = deadDaemon.url; + deadDaemon.close(); // port now closed + await sleep(100); + process.env.LVMH_URL = deadUrl; + const pi2 = makeFakePi(); + factory(pi2); + pi2.handlers.get("session_start")?.({ reason: "startup" }, makeFakeCtx()); + await sleep(400); + const t0 = Date.now(); + await sleep(100); + const loopResponsive = Date.now() - t0 < 500; + check("5 dead host: pi loads, event loop responsive", loopResponsive); + let threw = false; + try { + pi2.handlers.get("message_end")?.( + { message: { role: "user", content: "hi", timestamp: 1 } }, + makeFakeCtx(), + ); + } catch { + threw = true; + } + check("5 dead host: handlers never throw", !threw); + + delete process.env.LVMH_URL; + delete process.env.LVMH_TOKEN; + console.log( + failures === 0 ? "\nALL CHECKS PASSED" : `\n${failures} CHECK(S) FAILED`, + ); + process.exit(failures === 0 ? 0 : 1); } main().catch((err: unknown) => { - console.error("smoke harness crashed:", err); - process.exit(1); + console.error("smoke harness crashed:", err); + process.exit(1); });