diff --git a/src/config.ts b/src/config.ts index 88932a5..4e90812 100644 --- a/src/config.ts +++ b/src/config.ts @@ -6,6 +6,11 @@ export const HEARTBEAT_INTERVAL_MS = 10_000; export const LEASE_GRACE_MS = 30_000; export const REAP_INTERVAL_MS = 5_000; export const SUBSCRIBER_HEARTBEAT_MS = 15_000; +// Told to the browser as the SSE `retry:` field. Without it browsers use their +// own default (~3s in Chrome), which makes a server restart feel like an outage; +// a blip should heal before the user notices. Longer waits are the client's job: +// it takes the schedule over once a gap lingers, so this never becomes a hammer. +export const SSE_RETRY_MS = 500; export const MAX_SSE_FIELD_BYTES = 8 * 1024; // Actors with no subscribers and no active run are evicted after this TTL. export const ACTOR_IDLE_TTL_MS = 5 * 60_000; diff --git a/src/http.ts b/src/http.ts index c555126..a87f8fd 100644 --- a/src/http.ts +++ b/src/http.ts @@ -3,7 +3,7 @@ import { brotliCompressSync, gzipSync, constants as zlibConstants } from "node:z import { ConversationActor, type Subscriber, type WireEvent } from "./actor"; import { authEnabled, gateApi, getSession, sessionUser } from "./auth"; import type { BlobStore } from "./blobs"; -import { ACTOR_IDLE_TTL_MS, SUBSCRIBER_HEARTBEAT_MS } from "./config"; +import { ACTOR_IDLE_TTL_MS, SSE_RETRY_MS, SUBSCRIBER_HEARTBEAT_MS } from "./config"; import { Event, parseEventId } from "./events"; import { getExecutor } from "./executor"; import { getRegistry } from "./inference"; @@ -301,6 +301,8 @@ function openStream(conversationId: string, req: Request, store: Store): Respons // connection is what triggers replay via Last-Event-ID. Keepalive comments // keep proxies (and Bun's idle timeout) from closing the idle connection. const sseTransform = new TransformStream({ + // The reconnect delay is ours to set, not the browser's to guess. + start: (controller) => controller.enqueue(`retry: ${SSE_RETRY_MS}\n\n`), transform: (item, controller) => { if (item === KEEPALIVE) controller.enqueue(": keepalive\n\n"); else controller.enqueue(sseBlock(item)); diff --git a/tests/server.test.ts b/tests/server.test.ts index dc41498..6a94ef1 100644 --- a/tests/server.test.ts +++ b/tests/server.test.ts @@ -4,6 +4,7 @@ import { tmpdir } from "node:os"; import { join } from "node:path"; import { FsBlobStore } from "../src/blobs"; import { Catalog } from "../src/catalog"; +import { SSE_RETRY_MS } from "../src/config"; import { JobDriver } from "../src/drive"; import { apiRoutes, evictIdleActors, getActor } from "../src/http"; import { setRegistry } from "../src/inference"; @@ -74,6 +75,7 @@ async function readSse(res: Response, until: (frames: Frame[]) => boolean): Prom else if (line.startsWith("id:")) id = line.slice(3).trim(); else if (line.startsWith("data:")) dataLines.push(line.slice(5).trim()); } + if (!dataLines.length) continue; // control frame (`retry:`), not an event frames.push({ event, id, data: JSON.parse(dataLines.join("\n")) }); if (until(frames)) { await reader.cancel(); @@ -172,6 +174,16 @@ test("prompt → SSE stream emits user-message, message-start, deltas, message-e for (let i = 1; i < seqs.length; i++) expect(seqs[i]!).toBeGreaterThan(seqs[i - 1]!); }); +test("the stream sets its own reconnect interval", async () => { + // Without `retry:` the browser guesses (~3s in Chrome), which makes a server + // restart feel like an outage. It has to be the first thing on the wire. + const res = await fetch(`${base}/api/conversations/retryfield/stream`); + const reader = res.body!.getReader(); + const head = new TextDecoder().decode((await reader.read()).value); + await reader.cancel(); + expect(head.startsWith(`retry: ${SSE_RETRY_MS}\n\n`)).toBe(true); +}); + test("resume via HTTP Last-Event-ID replays only the gap", async () => { const conv = "s2"; const actor = getActor(conv, store);