diff --git a/packages/collab-web/src/lib/client.ts b/packages/collab-web/src/lib/client.ts index 0fc94430c..6fe5830e3 100644 --- a/packages/collab-web/src/lib/client.ts +++ b/packages/collab-web/src/lib/client.ts @@ -67,6 +67,8 @@ const MAX_NOTICES = 50; const TRANSCRIPT_TIMEOUT_MS = 10_000; /** Mirrors the TUI guest's WELCOME_TIMEOUT_MS: a host that never answers hello ends the join. */ const WELCOME_TIMEOUT_MS = 30_000; +/** Mirrors the TUI guest's SNAPSHOT_PROGRESS_TIMEOUT_MS: every snapshot chunk must make progress. */ +const SNAPSHOT_PROGRESS_TIMEOUT_MS = 30_000; interface PendingTranscript { resolve: (result: { text: string; newSize: number } | null) => void; @@ -85,6 +87,7 @@ export class GuestClient { #everConnected = false; #welcomed = false; #welcomeTimer: Timer | null = null; + #snapshotProgressTimer: Timer | null = null; #phase: ConnectionPhase = "connecting"; #endedReason: string | null = null; @@ -135,6 +138,7 @@ export class GuestClient { close(): void { this.#clearWelcomeTimer(); + this.#clearSnapshotProgressTimer(); this.#socket.close(); } @@ -188,6 +192,7 @@ export class GuestClient { } #handleClose(reason: string, willReconnect: boolean): void { + this.#clearSnapshotProgressTimer(); if (this.#phase === "ended") return; if (willReconnect) { this.#phase = "reconnecting"; @@ -200,6 +205,7 @@ export class GuestClient { #end(reason: string): void { if (this.#phase === "ended") return; this.#clearWelcomeTimer(); + this.#clearSnapshotProgressTimer(); this.#phase = "ended"; this.#endedReason = reason; for (const [, pending] of this.#pendingTranscripts) { @@ -218,6 +224,21 @@ export class GuestClient { } } + #armSnapshotProgressTimer(): void { + this.#clearSnapshotProgressTimer(); + this.#snapshotProgressTimer = setTimeout(() => { + this.#snapshotProgressTimer = null; + this.#end("timed out waiting for the host's session snapshot"); + }, SNAPSHOT_PROGRESS_TIMEOUT_MS); + } + + #clearSnapshotProgressTimer(): void { + if (this.#snapshotProgressTimer !== null) { + clearTimeout(this.#snapshotProgressTimer); + this.#snapshotProgressTimer = null; + } + } + /** Surfaces apply failures instead of letting the socket's recv chain swallow them. */ #applyFrameSafe(frame: HostFrame): void { try { @@ -251,7 +272,12 @@ export class GuestClient { this.#readOnly = frame.readOnly === true; this.#welcomed = true; this.#clearWelcomeTimer(); - if (frame.entryCount === 0) this.#phase = "live"; + if (frame.entryCount === 0) { + this.#clearSnapshotProgressTimer(); + this.#phase = "live"; + } else { + this.#armSnapshotProgressTimer(); + } this.#endedReason = null; break; case "snapshot-chunk": { @@ -259,7 +285,12 @@ export class GuestClient { // always closes the train with `final: true`; that flip is what // moves the guest from "waiting" to "live". this.#entries = [...this.#entries, ...frame.entries]; - if (frame.final) this.#phase = "live"; + if (frame.final) { + this.#clearSnapshotProgressTimer(); + this.#phase = "live"; + } else { + this.#armSnapshotProgressTimer(); + } break; } case "entry": diff --git a/packages/collab-web/test/client.test.ts b/packages/collab-web/test/client.test.ts index 1ef339994..51612918e 100644 --- a/packages/collab-web/test/client.test.ts +++ b/packages/collab-web/test/client.test.ts @@ -1,4 +1,4 @@ -import { describe, expect, it } from "bun:test"; +import { describe, expect, it, vi } from "bun:test"; import type { AgentSnapshot, AssistantMessage, @@ -91,6 +91,37 @@ describe("GuestClient frame apply", () => { expect(client.getSnapshot().readOnly).toBe(true); }); + it("times out stalled snapshot chunks and resets the clock on progress", () => { + vi.useFakeTimers(); + try { + const firstEntry = messageEntry("e1", { role: "user", content: "hi", timestamp: 1 }); + const client = new GuestClient(LINK, "tester"); + client.applyFrameForTest(welcomeFrame(2)); + expect(client.getSnapshot().phase).toBe("connecting"); + + vi.advanceTimersByTime(29_999); + expect(client.getSnapshot().phase).toBe("connecting"); + client.applyFrameForTest(snapshotChunk([firstEntry], false)); + expect(client.getSnapshot().entries).toEqual([firstEntry]); + expect(client.getSnapshot().phase).toBe("connecting"); + + vi.advanceTimersByTime(29_999); + expect(client.getSnapshot().phase).toBe("connecting"); + vi.advanceTimersByTime(1); + const snap = client.getSnapshot(); + expect(snap.phase).toBe("ended"); + expect(snap.endedReason).toBe("timed out waiting for the host's session snapshot"); + + const completeClient = new GuestClient(LINK, "tester"); + completeClient.applyFrameForTest(welcomeFrame(1)); + completeClient.applyFrameForTest(snapshotChunk([firstEntry])); + vi.advanceTimersByTime(30_000); + expect(completeClient.getSnapshot().phase).toBe("live"); + } finally { + vi.useRealTimers(); + } + }); + it("message_update sets the stream ghost (synthesizing a missed start)", () => { const client = liveClient(); const partial = assistantMessage("hel");