Skip to content
Open
Show file tree
Hide file tree
Changes from 1 commit
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
Prev Previous commit
Improve websocket handshake buffering and request queue safety
  • Loading branch information
jasonLaster committed Mar 9, 2026
commit ad65ff1f53f5f114c8e9c7447ef8237f61c0f9b8
37 changes: 37 additions & 0 deletions apps/server/src/wsServer.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -591,6 +591,43 @@ describe("WebSocket Server", () => {
});
});

it("delivers request-triggered pushes for requests sent immediately after open", async () => {
server = await createTestServer({ cwd: "/test/project" });
const addr = server.address();
const port = typeof addr === "object" && addr !== null ? addr.port : 0;
expect(port).toBeGreaterThan(0);

const ws = await connectWs(port);
connections.push(ws);

const responsePromise = sendRequest(ws, WS_METHODS.terminalOpen, {
threadId: asThreadId("thread-1"),
cwd: "/test/project",
});

const welcome = await waitForPush(ws, WS_CHANNELS.serverWelcome);
const terminalEvent = await waitForPush(ws, WS_CHANNELS.terminalEvent, (push) => push.data.type === "started");
const response = await responsePromise;

expect(welcome.channel).toBe(WS_CHANNELS.serverWelcome);
expect(terminalEvent.channel).toBe(WS_CHANNELS.terminalEvent);
expect(terminalEvent.sequence).toBeGreaterThan(welcome.sequence);
expect(response.id).toBeDefined();
expect(response.result).toEqual(expect.objectContaining({ threadId: "thread-1" }));
});


it("continues startup when keybindings runtime bootstrap fails", async () => {
server = await createTestServer({ cwd: "/test/project", stateDir: "/dev/null" });
const addr = server.address();
const port = typeof addr === "object" && addr !== null ? addr.port : 0;
expect(port).toBeGreaterThan(0);

const [ws, welcome] = await connectAndAwaitWelcome(port);
connections.push(ws);
expect(welcome.channel).toBe(WS_CHANNELS.serverWelcome);
});

it("serves persisted attachments from stateDir", async () => {
const stateDir = makeTempDir("t3code-state-attachments-");
const attachmentPath = path.join(stateDir, "attachments", "thread-a", "message-a", "0.png");
Expand Down
46 changes: 17 additions & 29 deletions apps/server/src/wsServer.ts
Original file line number Diff line number Diff line change
Expand Up @@ -35,7 +35,6 @@ import {
Layer,
Path,
PubSub,
Ref,
Schema,
Scope,
ServiceMap,
Expand Down Expand Up @@ -261,7 +260,6 @@ export const createServer = Effect.fn(function* (): Effect.fn.Return<

const providerStatuses = yield* providerHealth.getStatuses;

const clients = yield* Ref.make(new Set<WebSocket>());
const logger = createLogger("ws");
const readiness = yield* makeServerReadiness;

Expand All @@ -276,13 +274,14 @@ export const createServer = Effect.fn(function* (): Effect.fn.Return<
}

const pushBus = yield* makeServerPushBus({
clients,
logOutgoingPush,
});
yield* readiness.markPushBusReady;
yield* keybindingsManager.start.pipe(
Effect.mapError(
(cause) => new ServerLifecycleError({ operation: "keybindingsRuntimeStart", cause }),
Effect.catchAllCause((cause) =>
Effect.logWarning("keybindings runtime failed to start; continuing server bootstrap", {
error: Cause.pretty(cause),
}),
),
);
yield* readiness.markKeybindingsReady;
Expand Down Expand Up @@ -582,10 +581,11 @@ export const createServer = Effect.fn(function* (): Effect.fn.Return<
});
});

const closeAllClients = Ref.get(clients).pipe(
Effect.flatMap(Effect.forEach((client) => Effect.sync(() => client.close()))),
Effect.flatMap(() => Ref.set(clients, new Set())),
);
const closeAllClients = Effect.sync(() => {
wss.clients.forEach((client) => {
client.close();
});
});

const listenOptions = host ? { host, port } : { port };

Expand Down Expand Up @@ -961,15 +961,13 @@ export const createServer = Effect.fn(function* (): Effect.fn.Return<
...(welcomeBootstrapProjectId ? { bootstrapProjectId: welcomeBootstrapProjectId } : {}),
...(welcomeBootstrapThreadId ? { bootstrapThreadId: welcomeBootstrapThreadId } : {}),
};
// Send welcome before adding to broadcast set so publishAll calls
// cannot reach this client before the welcome arrives.
void runPromise(
readiness.awaitServerReady.pipe(
Effect.flatMap(() => pushBus.publishClient(ws, WS_CHANNELS.serverWelcome, welcomeData)),
Effect.flatMap((delivered) =>
delivered ? Ref.update(clients, (clients) => clients.add(ws)) : Effect.void,
),
),
Effect.gen(function* () {
yield* pushBus.registerClient(ws);
yield* readiness.awaitServerReady;
yield* pushBus.publishClient(ws, WS_CHANNELS.serverWelcome, welcomeData);
yield* pushBus.activateClient(ws);
}),
);

ws.on("message", (raw) => {
Expand All @@ -979,21 +977,11 @@ export const createServer = Effect.fn(function* (): Effect.fn.Return<
});

ws.on("close", () => {
void runPromise(
Ref.update(clients, (clients) => {
clients.delete(ws);
return clients;
}),
);
void runPromise(pushBus.unregisterClient(ws));
});

ws.on("error", () => {
void runPromise(
Ref.update(clients, (clients) => {
clients.delete(ws);
return clients;
}),
);
void runPromise(pushBus.unregisterClient(ws));
});
});

Expand Down
63 changes: 27 additions & 36 deletions apps/server/src/wsServer/pushBus.test.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
import type { WebSocket } from "ws";
import { afterEach, describe, expect, it } from "vitest";
import { Effect, Exit, Ref, Scope } from "effect";
import { Effect, Exit, Scope } from "effect";
import { WS_CHANNELS } from "@t3tools/contracts";

import { makeServerPushBus } from "./pushBus";
Expand Down Expand Up @@ -49,26 +49,21 @@ describe("makeServerPushBus", () => {
scope = null;
});

it("waits for the welcome push before a new client joins broadcast delivery", async () => {
it("queues publishAll pushes for pre_welcome clients and flushes after activation", async () => {
scope = await Effect.runPromise(Scope.make("sequential"));

const client = new MockWebSocket();
const { clients, pushBus } = await Effect.runPromise(
Effect.gen(function* () {
const clients = yield* Ref.make(new Set<WebSocket>());
const pushBus = yield* makeServerPushBus({
clients,
logOutgoingPush: () => {},
});

return { clients, pushBus };
const pushBus = await Effect.runPromise(
makeServerPushBus({
logOutgoingPush: () => {},
}).pipe(Scope.provide(scope)),
);

await Effect.runPromise(
Effect.gen(function* () {
yield* pushBus.registerClient(client as unknown as WebSocket);
yield* pushBus.publishAll(WS_CHANNELS.serverConfigUpdated, {
issues: [{ kind: "keybindings.malformed-config", message: "queued-before-connect" }],
issues: [{ kind: "keybindings.malformed-config", message: "queued-before-welcome" }],
providers: [],
});

Expand All @@ -82,39 +77,35 @@ describe("makeServerPushBus", () => {
);
expect(delivered).toBe(true);

yield* Ref.update(clients, (current) => current.add(client as unknown as WebSocket));

yield* pushBus.publishAll(WS_CHANNELS.serverConfigUpdated, {
issues: [],
providers: [],
});
yield* pushBus.activateClient(client as unknown as WebSocket);
}),
);

await client.waitForSentCount(2);

const messages = client.sent.map(
(message) => JSON.parse(message) as { channel: string; data: unknown },
(message) => JSON.parse(message) as { channel: string; data: unknown; sequence: number },
);

expect(messages).toHaveLength(2);
expect(messages[0]).toEqual({
type: "push",
sequence: 2,
channel: WS_CHANNELS.serverWelcome,
data: {
cwd: "/tmp/project",
projectName: "project",
expect(messages).toEqual([
{
type: "push",
sequence: 2,
channel: WS_CHANNELS.serverWelcome,
data: {
cwd: "/tmp/project",
projectName: "project",
},
},
});
expect(messages[1]).toEqual({
type: "push",
sequence: 3,
channel: WS_CHANNELS.serverConfigUpdated,
data: {
issues: [],
providers: [],
{
type: "push",
sequence: 1,
channel: WS_CHANNELS.serverConfigUpdated,
data: {
issues: [{ kind: "keybindings.malformed-config", message: "queued-before-welcome" }],
providers: [],
},
},
});
]);
});
});
Loading