fix: harden resume and SSE replay
This commit is contained in:
@@ -13,7 +13,7 @@ export type ClientEvent =
|
||||
| { type: "activity_delta"; text: string }
|
||||
| { type: "activity_event"; activity: ToolActivity }
|
||||
| { type: "info"; [k: string]: any }
|
||||
| { type: "system_event"; [k: string]: any };
|
||||
| { type: "system_event"; event: string };
|
||||
|
||||
export type TurnState = "idle" | "running" | "waiting" | "failed";
|
||||
|
||||
@@ -61,7 +61,9 @@ export class SessionBridge {
|
||||
this.emitToolActivity(m, "end");
|
||||
// Tool updates and raw payloads remain intentionally dropped.
|
||||
} else if (m.type === "system_event") {
|
||||
this.fan(m as ClientEvent);
|
||||
if (typeof m.event === "string" && m.event.trim() !== "") {
|
||||
this.fan({ type: "system_event", event: m.event });
|
||||
}
|
||||
} else if (m.type === "agent_end") {
|
||||
if (this.state !== "failed" && !this.pending) this.state = "idle";
|
||||
this.fan({ type: "system_event", event: "agent_end" });
|
||||
|
||||
@@ -8,6 +8,18 @@ import type { ReadinessManager } from "../runtime/readiness-manager.js";
|
||||
|
||||
const BOOTSTRAP_FAILURE_MESSAGE =
|
||||
"Session startup failed. Check configuration and connectivity, then Resume the session.";
|
||||
const READINESS_FAILURE_MESSAGE =
|
||||
"Session services are not ready. Check configuration and connectivity, then try again.";
|
||||
|
||||
function eventCursor(...values: unknown[]): number {
|
||||
let cursor = 0;
|
||||
for (const value of values.flatMap((item) => Array.isArray(item) ? item : [item])) {
|
||||
if (typeof value !== "string" || !/^\d+$/.test(value)) continue;
|
||||
const parsed = Number(value);
|
||||
if (Number.isSafeInteger(parsed)) cursor = Math.max(cursor, parsed);
|
||||
}
|
||||
return cursor;
|
||||
}
|
||||
|
||||
export function sessionRoutes(
|
||||
app: FastifyInstance,
|
||||
@@ -57,7 +69,7 @@ export function sessionRoutes(
|
||||
const b = req.body as { question: string; name?: string };
|
||||
const s = d.getSettings();
|
||||
const ensure = await d.readiness.ensure(s.workspace ?? "");
|
||||
if (!ensure.ok) return reply.code(503).send({ error: ensure.error ?? "Ollama/embeddings non disponibili" });
|
||||
if (!ensure.ok) return reply.code(503).send({ error: READINESS_FAILURE_MESSAGE });
|
||||
// Settings (global) supply workspace/provider/model/thinking. The new-question
|
||||
// form sends only the question text. `workspace` selects the tht `-c <config>`.
|
||||
const { id } = await d.tht.sessionNew({
|
||||
@@ -115,7 +127,7 @@ export function sessionRoutes(
|
||||
}
|
||||
const settings = d.getSettings();
|
||||
const ensure = await d.readiness.ensure(settings.workspace ?? "");
|
||||
if (!ensure.ok) return reply.code(503).send({ error: ensure.error ?? "Ollama/embeddings non disponibili" });
|
||||
if (!ensure.ok) return reply.code(503).send({ error: READINESS_FAILURE_MESSAGE });
|
||||
const saved = manifest as { provider?: string; model?: string; thinking?: string } | null;
|
||||
const options = {
|
||||
provider: saved?.provider,
|
||||
@@ -131,7 +143,7 @@ export function sessionRoutes(
|
||||
bindRuntime(id, rt);
|
||||
info(id, "Resuming session");
|
||||
bootstrap(id, rt, d.mgr.configure(rt, options), null, () => d.mgr.start(id, rt, options));
|
||||
return reply.code(200).send({ id });
|
||||
return reply.code(200).send({ id, alreadyActive: false });
|
||||
});
|
||||
app.post("/sessions/:id/close", async (req) => {
|
||||
const id = (req.params as { id: string }).id;
|
||||
@@ -146,6 +158,10 @@ export function sessionRoutes(
|
||||
app.get("/sessions/:id/events", (req, reply) => {
|
||||
const id = (req.params as any).id;
|
||||
const rt = d.mgr.get(id);
|
||||
const afterId = eventCursor(
|
||||
req.headers["last-event-id"],
|
||||
(req.query as { lastEventId?: unknown }).lastEventId,
|
||||
);
|
||||
// Add CORS headers manually: reply.raw.writeHead bypasses Fastify's onSend hooks
|
||||
// (where @fastify/cors injects headers), so we must set them explicitly here.
|
||||
const origin = (req.headers.origin as string | undefined) ?? "*";
|
||||
@@ -160,8 +176,12 @@ export function sessionRoutes(
|
||||
// Send the handshake immediately. Without this, Node waits for the first event body and
|
||||
// proxies/clients cannot establish an idle SSE subscription or inspect its headers.
|
||||
reply.raw.flushHeaders();
|
||||
const send = (event: string, data: object) => reply.raw.write(`event: ${event}\ndata: ${JSON.stringify(data)}\n\n`);
|
||||
const off = d.hub.subscribe(id, send, rt?.bridge.pendingWidget() ?? null);
|
||||
const send = (event: string, data: object, eventId: number) =>
|
||||
reply.raw.write(`id: ${eventId}\nevent: ${event}\ndata: ${JSON.stringify(data)}\n\n`);
|
||||
const off = d.hub.subscribe(id, send, {
|
||||
afterId,
|
||||
pending: rt?.bridge.pendingWidget() ?? null,
|
||||
});
|
||||
req.raw.on("close", off);
|
||||
});
|
||||
app.post("/sessions/:id/rename", async (req, reply) => {
|
||||
|
||||
+68
-24
@@ -1,48 +1,92 @@
|
||||
type Send = (event: string, data: object) => void;
|
||||
type Send = (event: string, data: object, id: number) => void;
|
||||
|
||||
/**
|
||||
* Per-session ring buffer of recent events.
|
||||
*
|
||||
* When the browser opens the SSE connection slightly after session creation
|
||||
* (React re-render, navigation, etc.), events produced by the Pi process in
|
||||
* that gap would be lost. The buffer replays them to late subscribers.
|
||||
*
|
||||
* `system_event` with `event: "agent_end"` acts as the natural sentinel:
|
||||
* once a subscriber sees it, the buffer for that session is safe to clear.
|
||||
* Event ids are transport identity: they are monotonically increasing for the lifetime of a
|
||||
* session id, including across a cold Resume that clears the old buffer and subscribers.
|
||||
*/
|
||||
const BUFFER_LIMIT = 200;
|
||||
|
||||
interface BufferedEvent { event: string; data: object }
|
||||
interface BufferedEvent {
|
||||
id: number;
|
||||
event: string;
|
||||
data: object;
|
||||
}
|
||||
|
||||
interface SubscribeOptions {
|
||||
afterId?: number;
|
||||
pending?: object | null;
|
||||
}
|
||||
|
||||
function descriptorId(value: unknown): string | null {
|
||||
if (!value || typeof value !== "object") return null;
|
||||
const id = (value as { id?: unknown }).id;
|
||||
return typeof id === "string" && id.trim() !== "" ? id : null;
|
||||
}
|
||||
|
||||
function bufferedGateId(item: BufferedEvent): string | null {
|
||||
if (item.event !== "ui_request") return null;
|
||||
const request = (item.data as { ui_request?: unknown }).ui_request;
|
||||
return descriptorId(request);
|
||||
}
|
||||
|
||||
export class SseHub {
|
||||
private subs = new Map<string, Set<Send>>();
|
||||
private buffers = new Map<string, BufferedEvent[]>();
|
||||
private lastIds = new Map<string, number>();
|
||||
|
||||
subscribe(sessionId: string, send: Send, pending?: object | null): () => void {
|
||||
subscribe(
|
||||
sessionId: string,
|
||||
send: Send,
|
||||
options: SubscribeOptions = {},
|
||||
): () => void {
|
||||
if (!this.subs.has(sessionId)) this.subs.set(sessionId, new Set());
|
||||
this.subs.get(sessionId)!.add(send);
|
||||
|
||||
// Replay buffered events so late subscribers don't miss the session lifecycle.
|
||||
const buf = this.buffers.get(sessionId);
|
||||
if (buf) for (const { event, data } of buf) send(event, data);
|
||||
const afterId = Number.isSafeInteger(options.afterId) && (options.afterId ?? 0) >= 0
|
||||
? options.afterId ?? 0
|
||||
: 0;
|
||||
const buf = this.buffers.get(sessionId) ?? [];
|
||||
for (const item of buf) {
|
||||
if (item.id > afterId) send(item.event, item.data, item.id);
|
||||
}
|
||||
|
||||
// A pending gate normally already exists in the buffer. If ring-buffer eviction removed it,
|
||||
// give the fallback a fresh transport id and buffer it for future cursor-based reconnects.
|
||||
// Matching is by stable descriptor identity only, never by event content.
|
||||
const pendingId = descriptorId(options.pending);
|
||||
const pendingBuffered = pendingId !== null && buf.some((item) => bufferedGateId(item) === pendingId);
|
||||
if (pendingId !== null && !pendingBuffered) {
|
||||
const data = { type: "ui_request", ui_request: options.pending! };
|
||||
const item = this.buffer(sessionId, "ui_request", data);
|
||||
send(item.event, item.data, item.id);
|
||||
}
|
||||
|
||||
// Match the live-event shape (hub.publish sends the full ClientEvent):
|
||||
// both carry { type, ui_request } so re-emit and live widgets are identical.
|
||||
if (pending) send("ui_request", { type: "ui_request", ui_request: pending });
|
||||
return () => this.subs.get(sessionId)?.delete(send);
|
||||
}
|
||||
|
||||
publish(sessionId: string, event: string, data: object): void {
|
||||
// Buffer the event for late subscribers.
|
||||
let buf = this.buffers.get(sessionId);
|
||||
if (!buf) { buf = []; this.buffers.set(sessionId, buf); }
|
||||
buf.push({ event, data });
|
||||
if (buf.length > BUFFER_LIMIT) buf.splice(0, buf.length - BUFFER_LIMIT);
|
||||
|
||||
for (const s of this.subs.get(sessionId) ?? []) s(event, data);
|
||||
publish(sessionId: string, event: string, data: object): number {
|
||||
const item = this.buffer(sessionId, event, data);
|
||||
for (const send of this.subs.get(sessionId) ?? []) send(event, data, item.id);
|
||||
return item.id;
|
||||
}
|
||||
|
||||
/** Clear buffer and subscribers for a finished session. */
|
||||
private buffer(sessionId: string, event: string, data: object): BufferedEvent {
|
||||
const id = (this.lastIds.get(sessionId) ?? 0) + 1;
|
||||
this.lastIds.set(sessionId, id);
|
||||
let buf = this.buffers.get(sessionId);
|
||||
if (!buf) {
|
||||
buf = [];
|
||||
this.buffers.set(sessionId, buf);
|
||||
}
|
||||
const item = { id, event, data };
|
||||
buf.push(item);
|
||||
if (buf.length > BUFFER_LIMIT) buf.splice(0, buf.length - BUFFER_LIMIT);
|
||||
return item;
|
||||
}
|
||||
|
||||
/** Clear buffered/runtime bindings while retaining session event-id monotonicity. */
|
||||
clear(sessionId: string): void {
|
||||
this.buffers.delete(sessionId);
|
||||
this.subs.delete(sessionId);
|
||||
|
||||
@@ -185,7 +185,7 @@ test.each(["idle", "failed"])(
|
||||
});
|
||||
const response = await app.inject({ method: "POST", url: "/sessions/s1/resume" });
|
||||
await new Promise((resolve) => setImmediate(resolve));
|
||||
expect(response.json()).toEqual({ id: "s1" });
|
||||
expect(response.json()).toEqual({ id: "s1", alreadyActive: false });
|
||||
expect(order).toEqual(["teardown:s1", "clear:s1", "reopen", "create", "start"]);
|
||||
},
|
||||
);
|
||||
@@ -223,7 +223,7 @@ test("POST resume without a runtime clears stale SSE state before cold start", a
|
||||
const response = await app.inject({ method: "POST", url: "/sessions/crashed/resume" });
|
||||
await new Promise((resolve) => setImmediate(resolve));
|
||||
|
||||
expect(response.json()).toEqual({ id: "crashed" });
|
||||
expect(response.json()).toEqual({ id: "crashed", alreadyActive: false });
|
||||
expect(order).toEqual(["clear:crashed", "reopen", "create", "start"]);
|
||||
expect(createOptions).toMatchObject({
|
||||
provider: "local-qwen", model: "qwen3.6-35b-a3b", thinking: "low", mode: "resume",
|
||||
@@ -359,11 +359,13 @@ test("POST resume on an archived session is refused with 409", async () => {
|
||||
expect(res.statusCode).toBe(409);
|
||||
});
|
||||
|
||||
test("POST /sessions refuses with 503 when ollamaEnsure fails (no session created)", async () => {
|
||||
test("POST /sessions readiness failure returns one fixed public message without raw diagnostics", async () => {
|
||||
let createdCalled = false;
|
||||
const rawFailure =
|
||||
"connect https://secret.invalid/ready?token=DO_NOT_LEAK using /srv/private/model-key";
|
||||
const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness" }), {
|
||||
thtRunner: {
|
||||
ollamaEnsure: async () => ({ ok: false, stage: "model", error: "modello non installato" }),
|
||||
ollamaEnsure: async () => ({ ok: false, stage: "model", error: rawFailure }),
|
||||
searchPack: async () => {},
|
||||
sessionNew: async () => { createdCalled = true; return { id: "s1" }; },
|
||||
} as any,
|
||||
@@ -372,7 +374,10 @@ test("POST /sessions refuses with 503 when ollamaEnsure fails (no session create
|
||||
});
|
||||
const res = await app.inject({ method: "POST", url: "/sessions", payload: { question: "q" } });
|
||||
expect(res.statusCode).toBe(503);
|
||||
expect(res.json().error).toContain("non installato");
|
||||
expect(res.json()).toEqual({
|
||||
error: "Session services are not ready. Check configuration and connectivity, then try again.",
|
||||
});
|
||||
expect(res.body).not.toMatch(/secret\.invalid|DO_NOT_LEAK|\/srv\/private\/model-key/);
|
||||
expect(createdCalled).toBe(false);
|
||||
});
|
||||
|
||||
@@ -403,10 +408,12 @@ test("POST /sessions/:id/resume returns 409 for a read-only session without call
|
||||
expect(ensureCalled).toBe(false);
|
||||
});
|
||||
|
||||
test("POST /sessions/:id/resume refuses with 503 when ollamaEnsure fails", async () => {
|
||||
test("POST /sessions/:id/resume readiness failure returns the same fixed public message", async () => {
|
||||
const rawFailure =
|
||||
"stderr https://secret.invalid/resume?api_key=DO_NOT_LEAK /srv/private/resume-key";
|
||||
const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness" }), {
|
||||
thtRunner: {
|
||||
ollamaEnsure: async () => ({ ok: false, error: "Ollama down" }),
|
||||
ollamaEnsure: async () => ({ ok: false, error: rawFailure }),
|
||||
sessionShow: async () => ({ status: "open", archived: false }),
|
||||
} as any,
|
||||
getSettings: () => ({ workspace: "psd" }) as any,
|
||||
@@ -414,6 +421,10 @@ test("POST /sessions/:id/resume refuses with 503 when ollamaEnsure fails", async
|
||||
});
|
||||
const res = await app.inject({ method: "POST", url: "/sessions/s1/resume" });
|
||||
expect(res.statusCode).toBe(503);
|
||||
expect(res.json()).toEqual({
|
||||
error: "Session services are not ready. Check configuration and connectivity, then try again.",
|
||||
});
|
||||
expect(res.body).not.toMatch(/secret\.invalid|DO_NOT_LEAK|\/srv\/private\/resume-key/);
|
||||
});
|
||||
|
||||
test("POST /runtime/prewarm returns 202 without awaiting readiness", async () => {
|
||||
|
||||
@@ -164,6 +164,26 @@ test("agent_end di Pi diventa un system_event agent_end per il FE", () => {
|
||||
expect(seen).toEqual([{ type: "system_event", event: "agent_end" }]);
|
||||
});
|
||||
|
||||
test("generic Pi system events expose only an allowlisted non-empty event name", () => {
|
||||
const { rpc, fire } = fakeRpc();
|
||||
const bridge = new SessionBridge(rpc);
|
||||
const seen: any[] = [];
|
||||
bridge.onClientEvent((event) => seen.push(event));
|
||||
|
||||
fire({
|
||||
type: "system_event",
|
||||
event: "session_exit",
|
||||
command: "curl https://secret.invalid/?token=DO_NOT_LEAK",
|
||||
result: { path: "/srv/private/key" },
|
||||
});
|
||||
fire({ type: "system_event", event: " ", stderr: "DO_NOT_LEAK_STDERR" });
|
||||
fire({ type: "system_event", event: 42, args: "DO_NOT_LEAK_ARGS" });
|
||||
|
||||
expect(seen).toEqual([{ type: "system_event", event: "session_exit" }]);
|
||||
expect(Object.keys(seen[0])).toEqual(["type", "event"]);
|
||||
expect(JSON.stringify(seen)).not.toMatch(/DO_NOT_LEAK|command|result|path|stderr|args/);
|
||||
});
|
||||
|
||||
test("steer invia un comando steer e riattiva il turno", () => {
|
||||
const { rpc, sent, fire } = fakeRpc();
|
||||
const bridge = new SessionBridge(rpc);
|
||||
|
||||
@@ -1,16 +1,93 @@
|
||||
import { test, expect } from "vitest";
|
||||
import { SseHub } from "../src/sse/sse-hub.js";
|
||||
|
||||
test("re-emette il widget pendente alla sottoscrizione", () => {
|
||||
const hub = new SseHub(); const sent: any[] = [];
|
||||
hub.subscribe("s1", (ev, data) => sent.push({ ev, data }), { id: "u1", widget: "select" });
|
||||
expect(sent[0]).toEqual({ ev: "ui_request", data: { type: "ui_request", ui_request: { id: "u1", widget: "select" } } });
|
||||
test("publish assigns monotonic ids and a subscriber replays only ids newer than its cursor", () => {
|
||||
const hub = new SseHub();
|
||||
const firstId = hub.publish("s1", "text_delta", { text: "one" });
|
||||
const secondId = hub.publish("s1", "text_delta", { text: "two" });
|
||||
const sent: any[] = [];
|
||||
|
||||
hub.subscribe(
|
||||
"s1",
|
||||
(event, data, id) => sent.push({ event, data, id }),
|
||||
{ afterId: firstId },
|
||||
);
|
||||
const thirdId = hub.publish("s1", "text_delta", { text: "three" });
|
||||
|
||||
expect([firstId, secondId, thirdId]).toEqual([1, 2, 3]);
|
||||
expect(sent).toEqual([
|
||||
{ event: "text_delta", data: { text: "two" }, id: 2 },
|
||||
{ event: "text_delta", data: { text: "three" }, id: 3 },
|
||||
]);
|
||||
});
|
||||
|
||||
test("publish raggiunge i subscriber e unsubscribe li stacca", () => {
|
||||
const hub = new SseHub(); const sent: any[] = [];
|
||||
const off = hub.subscribe("s1", (ev, data) => sent.push({ ev, data }));
|
||||
hub.publish("s1", "text_delta", { text: "x" });
|
||||
off(); hub.publish("s1", "text_delta", { text: "y" });
|
||||
expect(sent).toEqual([{ ev: "text_delta", data: { text: "x" } }]);
|
||||
test("clear drops stale subscribers and buffers but preserves the per-session id sequence", () => {
|
||||
const hub = new SseHub();
|
||||
const stale: any[] = [];
|
||||
hub.subscribe("s1", (event, data, id) => stale.push({ event, data, id }));
|
||||
expect(hub.publish("s1", "info", { text: "before" })).toBe(1);
|
||||
|
||||
hub.clear("s1");
|
||||
expect(hub.publish("s1", "info", { text: "after" })).toBe(2);
|
||||
const resumed: any[] = [];
|
||||
hub.subscribe(
|
||||
"s1",
|
||||
(event, data, id) => resumed.push({ event, data, id }),
|
||||
{ afterId: 1 },
|
||||
);
|
||||
|
||||
expect(stale).toEqual([{ event: "info", data: { text: "before" }, id: 1 }]);
|
||||
expect(resumed).toEqual([{ event: "info", data: { text: "after" }, id: 2 }]);
|
||||
});
|
||||
|
||||
test("a buffered pending gate is emitted exactly once and matched by descriptor id", () => {
|
||||
const hub = new SseHub();
|
||||
const descriptor = { id: "gate-1", widget: "select", title: "Original" };
|
||||
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: 0, pending: { ...descriptor, title: "Runtime copy" } },
|
||||
);
|
||||
|
||||
expect(replayed).toEqual([{
|
||||
event: "ui_request",
|
||||
data: { type: "ui_request", ui_request: descriptor },
|
||||
id: 1,
|
||||
}]);
|
||||
|
||||
const alreadySeen: any[] = [];
|
||||
hub.subscribe(
|
||||
"s1",
|
||||
(event, data, id) => alreadySeen.push({ event, data, id }),
|
||||
{ afterId: 1, pending: descriptor },
|
||||
);
|
||||
expect(alreadySeen).toEqual([]);
|
||||
});
|
||||
|
||||
test("an unbuffered pending gate receives one fresh buffered id", () => {
|
||||
const hub = new SseHub();
|
||||
hub.publish("s1", "info", { text: "ready" });
|
||||
const first: any[] = [];
|
||||
|
||||
hub.subscribe(
|
||||
"s1",
|
||||
(event, data, id) => first.push({ event, data, id }),
|
||||
{ afterId: 1, pending: { id: "gate-1", widget: "select" } },
|
||||
);
|
||||
expect(first).toEqual([{
|
||||
event: "ui_request",
|
||||
data: { type: "ui_request", ui_request: { id: "gate-1", widget: "select" } },
|
||||
id: 2,
|
||||
}]);
|
||||
|
||||
const reconnect: any[] = [];
|
||||
hub.subscribe(
|
||||
"s1",
|
||||
(event, data, id) => reconnect.push({ event, data, id }),
|
||||
{ afterId: 1, pending: { id: "gate-1", widget: "select" } },
|
||||
);
|
||||
expect(reconnect).toEqual(first);
|
||||
});
|
||||
|
||||
@@ -0,0 +1,82 @@
|
||||
import { test, expect } from "vitest";
|
||||
import { buildApp } from "../src/app.js";
|
||||
import { loadConfig } from "../src/config.js";
|
||||
import { SseHub } from "../src/sse/sse-hub.js";
|
||||
|
||||
async function readUntil(
|
||||
reader: ReadableStreamDefaultReader<Uint8Array>,
|
||||
predicate: (text: string) => boolean,
|
||||
): Promise<string> {
|
||||
const decoder = new TextDecoder();
|
||||
let text = "";
|
||||
const deadline = Date.now() + 1_000;
|
||||
while (Date.now() < deadline) {
|
||||
const result = await Promise.race([
|
||||
reader.read(),
|
||||
new Promise<{ timeout: true }>((resolve) =>
|
||||
setTimeout(() => resolve({ timeout: true }), Math.max(0, deadline - Date.now())),
|
||||
),
|
||||
]);
|
||||
if ("timeout" in result || result.done) break;
|
||||
text += decoder.decode(result.value);
|
||||
if (predicate(text)) break;
|
||||
}
|
||||
return text;
|
||||
}
|
||||
|
||||
async function captureReplay(options: { query?: string; lastEventId?: string }): Promise<string> {
|
||||
const hub = new SseHub();
|
||||
hub.publish("s1", "info", { type: "info", text: "one" });
|
||||
hub.publish("s1", "info", { type: "info", text: "two" });
|
||||
hub.publish("s1", "info", { type: "info", text: "three" });
|
||||
const app = buildApp(loadConfig({ THT_HARNESS_DIR: "/tmp/h" }), {
|
||||
hub,
|
||||
thtRunner: {} 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${options.query ?? ""}`,
|
||||
{
|
||||
headers: options.lastEventId ? { "Last-Event-ID": options.lastEventId } : undefined,
|
||||
signal: controller.signal,
|
||||
},
|
||||
);
|
||||
const reader = response.body!.getReader();
|
||||
const body = await readUntil(reader, (text) => text.includes('"three"'));
|
||||
await reader.cancel();
|
||||
return body;
|
||||
} finally {
|
||||
controller.abort();
|
||||
await app.close();
|
||||
}
|
||||
}
|
||||
|
||||
test("SSE emits ids and honors the native Last-Event-ID replay cursor", async () => {
|
||||
const body = await captureReplay({ lastEventId: "1" });
|
||||
|
||||
expect(body).not.toContain('"one"');
|
||||
expect(body).toContain("id: 2\nevent: info\n");
|
||||
expect(body).toContain('data: {"type":"info","text":"two"}');
|
||||
expect(body).toContain("id: 3\nevent: info\n");
|
||||
expect(body).toContain('data: {"type":"info","text":"three"}');
|
||||
});
|
||||
|
||||
test("SSE honors the manual lastEventId query cursor", async () => {
|
||||
const body = await captureReplay({ query: "?lastEventId=2" });
|
||||
|
||||
expect(body).not.toContain('"one"');
|
||||
expect(body).not.toContain('"two"');
|
||||
expect(body).toContain("id: 3\nevent: info\n");
|
||||
expect(body).toContain('data: {"type":"info","text":"three"}');
|
||||
});
|
||||
|
||||
test("SSE uses the newer valid cursor when header and query are both present", async () => {
|
||||
const body = await captureReplay({ query: "?lastEventId=2", lastEventId: "1" });
|
||||
|
||||
expect(body).not.toContain('"two"');
|
||||
expect(body).toContain("id: 3\nevent: info\n");
|
||||
});
|
||||
Reference in New Issue
Block a user