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
2 changes: 2 additions & 0 deletions .changeset/quiet-lifecycles-characterize.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,2 @@
---
---
151 changes: 151 additions & 0 deletions packages/session-broker/src/connection.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -76,6 +76,36 @@ function createSnapshot(): SessionSnapshot<TestSessionState> {
};
}

/** Await one lifecycle signal while turning a lost wakeup into a bounded test failure. */
async function settleWithinTestTimeout<T>(promise: PromiseLike<T> | T, timeoutMs = 500) {
let timeout: ReturnType<typeof setTimeout> | undefined;
try {
return await Promise.race([
Promise.resolve(promise),
new Promise<never>((_resolve, reject) => {
timeout = setTimeout(
() => reject(new Error(`Promise did not settle within ${timeoutMs}ms.`)),
timeoutMs,
);
timeout.unref?.();
}),
]);
} finally {
if (timeout) clearTimeout(timeout);
}
}

/** Create a manually released promise for reconnect race characterization. */
function createDeferredTest<T = void>() {
let resolve!: (value: T | PromiseLike<T>) => void;
let reject!: (reason?: unknown) => void;
const promise = new Promise<T>((resolvePromise, rejectPromise) => {
resolve = resolvePromise;
reject = rejectPromise;
});
return { promise, resolve, reject };
}

const protocolParsers = createSessionBrokerProtocolParsers<
TestSessionInfo,
TestSessionState,
Expand Down Expand Up @@ -153,6 +183,41 @@ describe("session broker connection", () => {
});
});

test("preserves a synchronous socket factory throw and permits a later direct start", () => {
const sockets: TestSocket[] = [];
let attempts = 0;
const connection = createSessionBrokerConnection<
TestSessionInfo,
TestSessionState,
TestSocket,
TestServerMessage,
{ ok: true }
>({
url: "ws://broker.test/session",
createSocket: () => {
attempts += 1;
if (attempts === 1) throw new Error("socket factory exploded");
const socket = new TestSocket();
sockets.push(socket);
return socket;
},
registration: createRegistration(),
snapshot: createSnapshot(),
protocolParsers,
});

expect(() => connection.start()).toThrow("socket factory exploded");
expect(attempts).toBe(1);
expect(sockets).toHaveLength(0);

connection.start();
sockets[0]!.emitOpen();
expect(attempts).toBe(2);
expect(sockets).toHaveLength(1);
expect(JSON.parse(sockets[0]!.sent[0]!)).toMatchObject({ type: "register" });
connection.stop();
});

test("withholds registration and replacement updates until producer authentication completes", async () => {
const socket = new TestSocket();
const pair = (await crypto.subtle.generateKey("Ed25519", false, [
Expand Down Expand Up @@ -1010,6 +1075,92 @@ describe("session broker connection", () => {
expect(sockets).toHaveLength(2);
});

test("explicitly starts a fresh generation after a natural no-reconnect close", async () => {
const sockets: TestSocket[] = [];
const connection = createSessionBrokerConnection<
TestSessionInfo,
TestSessionState,
TestSocket,
TestServerMessage,
{ ok: true }
>({
url: "ws://broker.test/session",
createSocket: () => {
const socket = new TestSocket();
sockets.push(socket);
return socket;
},
registration: createRegistration(),
snapshot: createSnapshot(),
protocolParsers,
resolveClose: () => ({ reconnect: false }),
});

connection.start();
sockets[0]!.emitOpen();
sockets[0]!.emitClose(1000, "complete");
await Bun.sleep(0);
expect(sockets).toHaveLength(1);

connection.start();
sockets[1]!.emitOpen();
expect(sockets).toHaveLength(2);
expect(sockets.map((socket) => (JSON.parse(socket.sent[0]!) as { type: string }).type)).toEqual(
["register", "register"],
);
connection.stop();
});

for (const outcome of ["resolve", "reject"] as const) {
test(`stop fences late reconnect preparation ${outcome} without warning or another socket`, async () => {
const sockets: TestSocket[] = [];
const warnings: string[] = [];
const preparationStarted = createDeferredTest();
const preparationGate = createDeferredTest();
const preparationSettled = createDeferredTest();
const connection = createSessionBrokerConnection<
TestSessionInfo,
TestSessionState,
TestSocket,
TestServerMessage,
{ ok: true }
>({
url: "ws://broker.test/session",
createSocket: () => {
const socket = new TestSocket();
sockets.push(socket);
return socket;
},
registration: createRegistration(),
snapshot: createSnapshot(),
protocolParsers,
reconnectDelayMs: 1,
prepareReconnect: async () => {
preparationStarted.resolve();
try {
await preparationGate.promise;
if (outcome === "reject") throw new Error("late preparation failure");
} finally {
preparationSettled.resolve();
}
},
onWarning: (message) => warnings.push(message),
});

connection.start();
sockets[0]!.emitOpen();
sockets[0]!.emitClose();
await settleWithinTestTimeout(preparationStarted.promise);
connection.stop();
preparationGate.resolve();
await settleWithinTestTimeout(preparationSettled.promise);
await Bun.sleep(5);

expect(warnings).toEqual([]);
expect(sockets).toHaveLength(1);
});
}

test("reconnects after socket close unless a close directive disables it", async () => {
const sockets: TestSocket[] = [];
const warnings: string[] = [];
Expand Down
Loading
Loading