Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
48 changes: 48 additions & 0 deletions apps/web/src/__tests__/connection.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -234,6 +234,54 @@ describe("4001 auth failure handling", () => {
});
});

describe("liveness heartbeat", () => {
it("closes a silent socket when the heartbeat goes unanswered", async () => {
vi.useFakeTimers();
const statusSpy = vi.fn();
const transport = createWsTransport("ws://localhost:1234", {
onStatusChange: statusSpy,
});
mockWsInstance.simulateOpen();

// 30s interval fires the probe; the mock never responds; the 10s
// heartbeat timeout closes the socket and drives the reconnect path.
await vi.advanceTimersByTimeAsync(30_000);
statusSpy.mockClear();
await vi.advanceTimersByTimeAsync(10_000);

expect(statusSpy).toHaveBeenCalledWith("reconnecting");
transport.close();
vi.useRealTimers();
});

it("keeps the socket open when the heartbeat is answered", async () => {
vi.useFakeTimers();
const statusSpy = vi.fn();
const transport = createWsTransport("ws://localhost:1234", {
onStatusChange: statusSpy,
});
mockWsInstance.simulateOpen();
const sendSpy = vi.spyOn(mockWsInstance, "send");

await vi.advanceTimersByTimeAsync(30_000);
const heartbeat = sendSpy.mock.calls
.map(([raw]) => JSON.parse(raw as string) as { id: string; method?: string })
.find((message) => message.method === "app.version");
expect(heartbeat).toBeDefined();
mockWsInstance.onmessage?.({
data: JSON.stringify({ id: heartbeat!.id, result: "0.0.1" }),
});

await vi.advanceTimersByTimeAsync(10_000);
expect(mockWsInstance.readyState).toBe(1);
statusSpy.mockClear();
await vi.advanceTimersByTimeAsync(10_000);
expect(statusSpy).not.toHaveBeenCalledWith("reconnecting");
transport.close();
vi.useRealTimers();
});
});

/**
* Browser-faithful socket: send() throws while CONNECTING and silently
* discards while CLOSING/CLOSED, matching the WebSocket spec. The silent
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,11 @@ class TimeoutSocket {
const parsed: unknown = JSON.parse(raw);
if (!isRpcRequest(parsed)) throw new Error("Expected an RPC request");
this.requests.push(parsed);
// A socket that ignores the liveness probe is closed by the watchdog;
// these tests exercise RPC timeouts on a live connection instead.
if (parsed.method === "app.version") {
queueMicrotask(() => this.respond(parsed, "0.0.1"));
}
}

open(): void {
Expand Down
43 changes: 42 additions & 1 deletion apps/web/src/transport/ws-transport.ts
Original file line number Diff line number Diff line change
Expand Up @@ -107,6 +107,13 @@ const TERMINAL_LATE_CREATE_CLEANUP_TIMEOUT_MS = 10_000;
/** Maximum Terminal creates whose late responses still need exact cleanup. */
const MAX_PENDING_TERMINAL_CREATE_CLEANUPS = 8;

/** Interval between liveness probes on an open socket. */
const HEARTBEAT_INTERVAL_MS = 30_000;
/** How long a heartbeat may go unanswered before the socket is treated as dead. */
const HEARTBEAT_TIMEOUT_MS = 10_000;
/** Deadline for `agent.send` so a composer submit cannot park forever on a dead socket. */
const SEND_MESSAGE_TIMEOUT_MS = 20_000;

/** Last thread-list refresh timestamp per workspace, triggered on WS reconnect. */
const lastLoadThreadsAtByWorkspace = new Map<string, number>();
/** Minimum interval between reconnect-triggered thread-list fetches to avoid rapid-reconnect storms. */
Expand Down Expand Up @@ -386,6 +393,8 @@ export function createWsTransport(
let closed = false;
let reconnectDelay = MIN_RECONNECT_MS;
let reconnectTimer: ReturnType<typeof setTimeout> | null = null;
let heartbeatTimer: ReturnType<typeof setInterval> | null = null;
let heartbeatInFlight = false;
let terminalSelectionPromise: Promise<TerminalBackendCapabilities> | null = null;
// Track consecutive auth failures so we apply backoff after 3 immediate
// retries, preventing a tight loop when the token is persistently wrong.
Expand Down Expand Up @@ -554,6 +563,7 @@ export function createWsTransport(
setAttachmentTransportWsUrl(url);
resolveReady();
options?.onStatusChange?.("connected");
startHeartbeat();
invalidateLiveTurnDiff();
void selectTerminalClientWithRecovery();

Expand Down Expand Up @@ -614,6 +624,7 @@ export function createWsTransport(

ws.onclose = (event: CloseEvent) => {
freshTurnDiffThreads.clear();
stopHeartbeat();
rejectPending("WebSocket disconnected");
lateResponseHandlers = new Map();
pendingTerminalCreateCleanups.clear();
Expand All @@ -640,6 +651,35 @@ export function createWsTransport(
pending = new Map();
}

// A half-open socket answers nothing and never fires `close`; without a
// probe, reconnect and runtime resync wait forever while the UI shows stale
// streaming state. The probe is a cheap existing RPC bounded by timeoutMs;
// on timeout we close locally so onclose drives the recovery ladder.
function heartbeatTick(): void {
if (heartbeatInFlight || closed || ws.readyState !== WebSocket.OPEN) return;
heartbeatInFlight = true;
rpc<string>("app.version", {}, { timeoutMs: HEARTBEAT_TIMEOUT_MS })
.catch(() => {
if (ws.readyState === WebSocket.OPEN) ws.close();
})
.finally(() => {
heartbeatInFlight = false;
});
}

function startHeartbeat(): void {
stopHeartbeat();
heartbeatTimer = setInterval(heartbeatTick, HEARTBEAT_INTERVAL_MS);
}

function stopHeartbeat(): void {
if (heartbeatTimer) {
clearInterval(heartbeatTimer);
heartbeatTimer = null;
}
heartbeatInFlight = false;
}

function invalidateLiveTurnDiff(): void {
void Promise.all([import("@/features/projects/state/workspaceStore"), import("@/stores/diffStore")]).then(([workspace, diff]) => {
const threadId = workspace.useWorkspaceStore.getState().activeThreadId;
Expand Down Expand Up @@ -1153,7 +1193,7 @@ export function createWsTransport(
...(replyToMessageId && { replyToMessageId }),
...(quotedText && { quotedText }),
...guardrails,
});
}, { timeoutMs: SEND_MESSAGE_TIMEOUT_MS });
},
getRecoveryIncident: () =>
rpc<import("@mcode/contracts").RecoveryIncident | null>("agent.recoveryIncident", {}),
Expand Down Expand Up @@ -1543,6 +1583,7 @@ export function createWsTransport(
// Lifecycle
close: () => {
closed = true;
stopHeartbeat();
if (reconnectTimer) {
clearTimeout(reconnectTimer);
reconnectTimer = null;
Expand Down
Loading