Skip to content

Commit 1b248f6

Browse files
committed
fix(ws): recover terminal replay after silent disconnects
1 parent d52395f commit 1b248f6

10 files changed

Lines changed: 515 additions & 2 deletions

File tree

packages/server/src/__tests__/ws-hub.test.ts

Lines changed: 87 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -13,6 +13,7 @@ import type { FastifyRequest } from "fastify";
1313
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
1414
import WebSocket from "ws";
1515
import { EventBus } from "../bus/event-bus.js";
16+
import "../commands/connection.js";
1617
import "../commands/terminal.js";
1718
import { clearPendingTerminalInput } from "../commands/terminal.js";
1819
import type { ServerConfig } from "../config.js";
@@ -34,6 +35,7 @@ type MockSocket = {
3435

3536
type MessageHandler = ((data: Buffer, isBinary?: boolean) => void) | undefined;
3637
type CloseHandler = (() => void) | undefined;
38+
type PongHandler = (() => void) | undefined;
3739

3840
type ResultMessage = Extract<ServerToClient, { kind: "result" }>;
3941

@@ -90,6 +92,11 @@ const getCloseHandler = (socket: MockSocket): CloseHandler =>
9092
| CloseHandler
9193
| undefined;
9294

95+
const getPongHandler = (socket: MockSocket): PongHandler =>
96+
socket.on.mock.calls.find((call: unknown[]) => call[0] === "pong")?.[1] as
97+
| PongHandler
98+
| undefined;
99+
93100
const subscribeToAllTopics = (socket: MockSocket) => {
94101
const messageHandler = getMessageHandler(socket);
95102

@@ -435,6 +442,86 @@ describe("WsHub", () => {
435442
expect(socket.ping).toHaveBeenCalled();
436443
});
437444

445+
it("handles connection.probe commands", async () => {
446+
const socket = createMockSocket();
447+
hub.handleConnection(socket as never, createMockRequest());
448+
const messageHandler = getMessageHandler(socket);
449+
socket.send.mockClear();
450+
451+
messageHandler?.(
452+
Buffer.from(
453+
JSON.stringify({
454+
kind: "command",
455+
id: "00000000-0000-4000-8000-000000000001",
456+
op: "connection.probe",
457+
args: {},
458+
})
459+
)
460+
);
461+
462+
await Promise.resolve();
463+
await new Promise((resolve) => setImmediate(resolve));
464+
465+
expect(findResultMessage(socket, "00000000-0000-4000-8000-000000000001")).toMatchObject({
466+
ok: true,
467+
data: {
468+
ok: true,
469+
},
470+
});
471+
});
472+
473+
it("closes clients that stay unresponsive across keepalive sweeps", () => {
474+
const socket = createMockSocket();
475+
hub.handleConnection(socket as never, createMockRequest());
476+
477+
hub.pingAll();
478+
expect(socket.ping).toHaveBeenCalledTimes(1);
479+
expect(socket.close).not.toHaveBeenCalled();
480+
481+
hub.pingAll();
482+
expect(socket.close).toHaveBeenCalledTimes(1);
483+
});
484+
485+
it("does not close clients that answer the previous keepalive ping", () => {
486+
const socket = createMockSocket();
487+
hub.handleConnection(socket as never, createMockRequest());
488+
const pongHandler = getPongHandler(socket);
489+
490+
hub.pingAll();
491+
pongHandler?.();
492+
hub.pingAll();
493+
494+
expect(socket.ping).toHaveBeenCalledTimes(2);
495+
expect(socket.close).not.toHaveBeenCalled();
496+
});
497+
498+
it("treats inbound commands as proof of life before the next keepalive sweep", async () => {
499+
const socket = createMockSocket();
500+
hub.handleConnection(socket as never, createMockRequest());
501+
const messageHandler = getMessageHandler(socket);
502+
503+
hub.pingAll();
504+
505+
messageHandler?.(
506+
Buffer.from(
507+
JSON.stringify({
508+
kind: "command",
509+
id: "00000000-0000-4000-8000-000000000002",
510+
op: "connection.probe",
511+
args: {},
512+
})
513+
)
514+
);
515+
516+
await Promise.resolve();
517+
await new Promise((resolve) => setImmediate(resolve));
518+
519+
hub.pingAll();
520+
521+
expect(socket.close).not.toHaveBeenCalled();
522+
expect(socket.ping).toHaveBeenCalledTimes(2);
523+
});
524+
438525
it("should return null for writer (deprecated - use FencingManager)", () => {
439526
const socket = createMockSocket();
440527
hub.handleConnection(socket as never, createMockRequest());
Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,6 @@
1+
import { z } from "zod";
2+
import { registerCommand } from "../ws/dispatch.js";
3+
4+
registerCommand("connection.probe", z.object({}).default({}), async () => {
5+
return { ok: true as const };
6+
});

packages/server/src/commands/index.ts

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,7 @@
66

77
import "./workspace.js";
88
import "./workspace-activity.js";
9+
import "./connection.js";
910
import "./session.js";
1011
import "./terminal.js";
1112
import "./file.js";

packages/server/src/server.ts

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -44,6 +44,8 @@ import { WsHub } from "./ws/hub.js";
4444

4545
import "./commands/index.js";
4646

47+
const WS_KEEPALIVE_INTERVAL_MS = 15_000;
48+
4749
export interface Server {
4850
app: FastifyInstance;
4951
stop: () => Promise<void>;
@@ -235,12 +237,18 @@ export async function createServer(
235237
}, STARTUP_GC_DELAY_MS);
236238
gcTimer.unref();
237239

240+
const wsKeepaliveTimer = setInterval(() => {
241+
wsHub.pingAll();
242+
}, WS_KEEPALIVE_INTERVAL_MS);
243+
wsKeepaliveTimer.unref();
244+
238245
let stopped = false;
239246
const stopServer = async () => {
240247
if (stopped) return;
241248
stopped = true;
242249

243250
clearTimeout(gcTimer);
251+
clearInterval(wsKeepaliveTimer);
244252
await app.close();
245253
autoFetch.stop();
246254
supervisorMgr.stop();

packages/server/src/ws/client.ts

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -64,8 +64,14 @@ export class WsClient {
6464
this.setupSocketHandlers();
6565
}
6666

67+
private markAlive(): void {
68+
this.isAlive = true;
69+
}
70+
6771
private setupSocketHandlers(): void {
6872
this.socket.on("message", (data: Buffer, isBinary: boolean) => {
73+
this.markAlive();
74+
6975
if (isBinary) {
7076
this.messageHandler?.(data);
7177
return;
@@ -87,7 +93,7 @@ export class WsClient {
8793
});
8894

8995
this.socket.on("pong", () => {
90-
this.isAlive = true;
96+
this.markAlive();
9197
});
9298
}
9399

packages/server/src/ws/hub.ts

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -320,6 +320,10 @@ export class WsHub implements Broadcaster {
320320
*/
321321
pingAll(): void {
322322
for (const client of this.clients.values()) {
323+
if (!client.alive) {
324+
client.close(1011, "keepalive_timeout");
325+
continue;
326+
}
323327
client.ping();
324328
}
325329
}

packages/web/src/features/terminal-panel/__tests__/xterm-host.test.tsx

Lines changed: 191 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -4449,6 +4449,78 @@ describe("XtermHost", () => {
44494449
});
44504450
});
44514451

4452+
it("replays terminal history when websocket recovery succeeds without a status transition", async () => {
4453+
const store = createStore();
4454+
const initialReplayChunk = new TextEncoder().encode("initial replay\n");
4455+
const recoveredReplayChunk = new TextEncoder().encode("recovered after probe\n");
4456+
let recoveryHandler:
4457+
| ((trigger: "visibility_resume" | "network_online" | "manual_retry" | "reconnected") => void)
4458+
| undefined;
4459+
let replayCount = 0;
4460+
const sendCommand = vi
4461+
.fn()
4462+
.mockImplementation((op: string, args: { terminalId?: string; lastSeq?: number }) => {
4463+
if (op !== "terminal.replay") {
4464+
return Promise.resolve({ status: "ok" });
4465+
}
4466+
4467+
replayCount += 1;
4468+
if (replayCount === 1) {
4469+
expect(args.lastSeq).toBe(0);
4470+
return Promise.resolve({
4471+
status: "ok",
4472+
transport: "binary",
4473+
streamId: 980,
4474+
size: initialReplayChunk.byteLength,
4475+
seq: 100,
4476+
bytes: initialReplayChunk,
4477+
} satisfies TerminalReplayPayload);
4478+
}
4479+
4480+
expect(args.lastSeq).toBe(100);
4481+
return Promise.resolve({
4482+
status: "ok",
4483+
transport: "binary",
4484+
streamId: 981,
4485+
size: recoveredReplayChunk.byteLength,
4486+
seq: 126,
4487+
bytes: recoveredReplayChunk,
4488+
} satisfies TerminalReplayPayload);
4489+
});
4490+
4491+
store.set(wsClientAtom, {
4492+
sendCommand,
4493+
subscribe: vi.fn(() => vi.fn()),
4494+
getStatus: vi.fn(() => "connected"),
4495+
onStatus: vi.fn(() => () => {}),
4496+
onRecovery: vi.fn((handler: typeof recoveryHandler) => {
4497+
recoveryHandler = handler;
4498+
return () => {};
4499+
}),
4500+
} as never);
4501+
4502+
render(
4503+
<Provider store={store}>
4504+
<XtermHost terminalId="probe-recovery-terminal" workspaceId="test-workspace" />
4505+
</Provider>
4506+
);
4507+
4508+
await waitFor(() => {
4509+
expectTerminalWriteData(initialReplayChunk);
4510+
});
4511+
4512+
await act(async () => {
4513+
recoveryHandler?.("network_online");
4514+
await Promise.resolve();
4515+
await Promise.resolve();
4516+
});
4517+
4518+
await waitFor(() => {
4519+
expect(replayCount).toBe(2);
4520+
expectTerminalWriteData(recoveredReplayChunk);
4521+
});
4522+
});
4523+
44524524
it("waits for reconnect replay bytes to finish rendering before starting another reconnect recovery", async () => {
44534525
const store = createStore();
44544526
const initialReplayChunk = new TextEncoder().encode("initial replay\n");
@@ -4573,6 +4645,125 @@ describe("XtermHost", () => {
45734645
});
45744646
});
45754647

4648+
it("queues probe-triggered recovery that arrives while historical recovery writes are still in flight", async () => {
4649+
const store = createStore();
4650+
const initialReplayChunk = new TextEncoder().encode("initial replay\n");
4651+
const delayedReconnectReplay = new TextEncoder().encode("delayed reconnect replay\n");
4652+
const queuedProbeReplay = new TextEncoder().encode("queued probe recovery\n");
4653+
let recoveryHandler:
4654+
| ((trigger: "visibility_resume" | "network_online" | "manual_retry" | "reconnected") => void)
4655+
| undefined;
4656+
let replayCount = 0;
4657+
let releaseDelayedReconnectWrite: (() => void) | undefined;
4658+
4659+
mockTerminal.write.mockImplementation((data: Uint8Array | string, callback?: () => void) => {
4660+
if (data === delayedReconnectReplay) {
4661+
releaseDelayedReconnectWrite = callback;
4662+
return;
4663+
}
4664+
callback?.();
4665+
});
4666+
4667+
const sendCommand = vi
4668+
.fn()
4669+
.mockImplementation((op: string, args: { terminalId?: string; lastSeq?: number }) => {
4670+
if (op !== "terminal.replay") {
4671+
return Promise.resolve({ status: "ok" });
4672+
}
4673+
4674+
replayCount += 1;
4675+
if (replayCount === 1) {
4676+
expect(args.lastSeq).toBe(0);
4677+
return Promise.resolve({
4678+
status: "ok",
4679+
transport: "binary",
4680+
streamId: 982,
4681+
size: initialReplayChunk.byteLength,
4682+
seq: 100,
4683+
bytes: initialReplayChunk,
4684+
} satisfies TerminalReplayPayload);
4685+
}
4686+
4687+
if (replayCount === 2) {
4688+
expect(args.lastSeq).toBe(100);
4689+
return Promise.resolve({
4690+
status: "ok",
4691+
transport: "binary",
4692+
streamId: 983,
4693+
size: delayedReconnectReplay.byteLength,
4694+
seq: 200,
4695+
bytes: delayedReconnectReplay,
4696+
} satisfies TerminalReplayPayload);
4697+
}
4698+
4699+
expect(args.lastSeq).toBe(200);
4700+
return Promise.resolve({
4701+
status: "ok",
4702+
transport: "binary",
4703+
streamId: 984,
4704+
size: queuedProbeReplay.byteLength,
4705+
seq: 240,
4706+
bytes: queuedProbeReplay,
4707+
} satisfies TerminalReplayPayload);
4708+
});
4709+
4710+
mockTerminal.cols = 132;
4711+
mockTerminal.rows = 36;
4712+
4713+
store.set(wsClientAtom, {
4714+
sendCommand,
4715+
subscribe: vi.fn(() => vi.fn()),
4716+
getStatus: vi.fn(() => "connected"),
4717+
onStatus: vi.fn(() => () => {}),
4718+
onRecovery: vi.fn((handler: typeof recoveryHandler) => {
4719+
recoveryHandler = handler;
4720+
return () => {};
4721+
}),
4722+
} as never);
4723+
4724+
render(
4725+
<Provider store={store}>
4726+
<XtermHost terminalId="queued-probe-recovery-terminal" workspaceId="test-workspace" />
4727+
</Provider>
4728+
);
4729+
4730+
await waitFor(() => {
4731+
expectTerminalWriteData(initialReplayChunk);
4732+
});
4733+
4734+
await act(async () => {
4735+
recoveryHandler?.("network_online");
4736+
await Promise.resolve();
4737+
await Promise.resolve();
4738+
});
4739+
4740+
await waitFor(() => {
4741+
expect(replayCount).toBe(2);
4742+
expectTerminalWriteData(delayedReconnectReplay);
4743+
});
4744+
4745+
expect(typeof releaseDelayedReconnectWrite).toBe("function");
4746+
4747+
await act(async () => {
4748+
recoveryHandler?.("visibility_resume");
4749+
await Promise.resolve();
4750+
await Promise.resolve();
4751+
});
4752+
4753+
expect(replayCount).toBe(2);
4754+
4755+
await act(async () => {
4756+
releaseDelayedReconnectWrite?.();
4757+
await Promise.resolve();
4758+
await Promise.resolve();
4759+
});
4760+
4761+
await waitFor(() => {
4762+
expect(replayCount).toBe(3);
4763+
expectTerminalWriteData(queuedProbeReplay);
4764+
});
4765+
});
4766+
45764767
it("uses the flushed pending chunk seq for a later reconnect after a successful recovery", async () => {
45774768
const store = createStore();
45784769
const initialReplayChunk = new TextEncoder().encode("initial replay\n");

0 commit comments

Comments
 (0)