fix: close delete resume and SSE replay races
This commit is contained in:
@@ -52,9 +52,15 @@ export class SseHub {
|
||||
const subscriber = { send, close: options.close ?? (() => undefined), closed: false };
|
||||
this.subs.get(sessionId)!.add(subscriber);
|
||||
|
||||
const afterId = Number.isSafeInteger(options.afterId) && (options.afterId ?? 0) >= 0
|
||||
const requestedAfterId = Number.isSafeInteger(options.afterId) && (options.afterId ?? 0) >= 0
|
||||
? options.afterId ?? 0
|
||||
: 0;
|
||||
// Native EventSource persists Last-Event-ID across a backend process restart. A cursor newer
|
||||
// than anything this hub generation has produced cannot identify an event in this process,
|
||||
// so replay the fresh generation from its beginning instead of suppressing every low id.
|
||||
const afterId = requestedAfterId > (this.lastIds.get(sessionId) ?? 0)
|
||||
? 0
|
||||
: requestedAfterId;
|
||||
const buf = this.buffers.get(sessionId) ?? [];
|
||||
for (const item of buf) {
|
||||
if (item.id > afterId) send(item.event, item.data, item.id);
|
||||
|
||||
@@ -151,3 +151,30 @@ test("an unbuffered pending gate receives one fresh buffered id", () => {
|
||||
);
|
||||
expect(reconnect).toEqual(first);
|
||||
});
|
||||
|
||||
test("a cursor newer than this hub generation replays low ids and a buffered gate once", () => {
|
||||
const hub = new SseHub();
|
||||
const descriptor = { id: "gate-fresh", widget: "select", title: "Choose" };
|
||||
hub.publish("s1", "info", { type: "info", text: "fresh generation" });
|
||||
hub.publish("s1", "ui_request", { type: "ui_request", ui_request: descriptor });
|
||||
const replayed: any[] = [];
|
||||
|
||||
hub.subscribe(
|
||||
"s1",
|
||||
(event, data, id) => replayed.push({ event, data, id }),
|
||||
{ afterId: 900, pending: { ...descriptor } },
|
||||
);
|
||||
|
||||
expect(replayed).toEqual([
|
||||
{
|
||||
event: "info",
|
||||
data: { type: "info", text: "fresh generation" },
|
||||
id: 1,
|
||||
},
|
||||
{
|
||||
event: "ui_request",
|
||||
data: { type: "ui_request", ui_request: descriptor },
|
||||
id: 2,
|
||||
},
|
||||
]);
|
||||
});
|
||||
|
||||
@@ -92,6 +92,42 @@ test("SSE uses the newer valid cursor when header and query are both present", a
|
||||
expect(body).toContain("id: 3\nevent: info\n");
|
||||
});
|
||||
|
||||
test("SSE resets a cursor from an older process and emits a buffered pending gate once", async () => {
|
||||
const hub = new SseHub();
|
||||
const descriptor = { id: "gate-fresh", widget: "select", title: "Choose" };
|
||||
hub.publish("s1", "info", { type: "info", text: "fresh generation" });
|
||||
hub.publish("s1", "ui_request", { type: "ui_request", ui_request: descriptor });
|
||||
const app = buildApp(loadConfig({ THT_HARNESS_DIR: "/tmp/h" }), {
|
||||
hub,
|
||||
thtRunner: {} as any,
|
||||
mgr: {
|
||||
get: () => ({ bridge: { pendingWidget: () => ({ ...descriptor }) } }),
|
||||
} as any,
|
||||
});
|
||||
await app.listen({ port: 0, host: "127.0.0.1" });
|
||||
const port = (app.server.address() as { port: number }).port;
|
||||
const controller = new AbortController();
|
||||
|
||||
try {
|
||||
const response = await fetch(`http://127.0.0.1:${port}/sessions/s1/events`, {
|
||||
headers: { "Last-Event-ID": "900" },
|
||||
signal: controller.signal,
|
||||
});
|
||||
const reader = response.body!.getReader();
|
||||
const body = await readUntil(reader, (text) => text.includes('"gate-fresh"'));
|
||||
await reader.cancel();
|
||||
|
||||
expect(body).toContain("id: 1\nevent: info\n");
|
||||
expect(body).toContain('data: {"type":"info","text":"fresh generation"}');
|
||||
expect(body.match(/event: ui_request/g)).toHaveLength(1);
|
||||
expect(body).toContain("id: 2\nevent: ui_request\n");
|
||||
expect(body).toContain('"ui_request":{"id":"gate-fresh","widget":"select","title":"Choose"}');
|
||||
} finally {
|
||||
controller.abort();
|
||||
await app.close();
|
||||
}
|
||||
});
|
||||
|
||||
test("clear ends every live SSE response and a cursor reconnect replays post-clear events", async () => {
|
||||
const hub = new SseHub();
|
||||
const app = buildApp(loadConfig({ THT_HARNESS_DIR: "/tmp/h" }), {
|
||||
|
||||
Reference in New Issue
Block a user