From 589d48a975e50894fa7c03943af5023e7353a0bf Mon Sep 17 00:00:00 2001 From: enricobuehler Date: Mon, 20 Jul 2026 19:59:04 +0200 Subject: [PATCH] fix(plugin-kit): SSE keepalive must outpace Bun's idle timeout (0.1.4) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit servePluginUi's Bun.serve leaves idleTimeout at Bun's 10s default, so the old 15s keepalive could never arrive — an idle status feed was closed before its first ping, which is exactly what a live host showed. Default is now 5s. Also adds the regression suite that should have caught it: the original test only covered a finite, self-driving stream. The new one exercises the production shape (PubSub-backed, published after the response opens) plus the keepalive — with a reader that does not re-enter read(), which is what made the first version of this test lie. Co-Authored-By: Claude Fable 5 --- plugin-kit/package.json | 2 +- plugin-kit/src/sse.ts | 8 +++- plugin-kit/test/sse-live.test.ts | 79 ++++++++++++++++++++++++++++++++ 3 files changed, 86 insertions(+), 3 deletions(-) create mode 100644 plugin-kit/test/sse-live.test.ts diff --git a/plugin-kit/package.json b/plugin-kit/package.json index c63aa5e2..e489fbf0 100644 --- a/plugin-kit/package.json +++ b/plugin-kit/package.json @@ -1,6 +1,6 @@ { "name": "@punktfunk/plugin-kit", - "version": "0.1.3", + "version": "0.1.4", "description": "Effect-based framework for punktfunk plugins: lifecycle runtime, config/state, sync engine, UI serving, CLI scaffold, and browser helpers.", "type": "module", "license": "MIT OR Apache-2.0", diff --git a/plugin-kit/src/sse.ts b/plugin-kit/src/sse.ts index 576b0756..2a2d36f0 100644 --- a/plugin-kit/src/sse.ts +++ b/plugin-kit/src/sse.ts @@ -12,7 +12,11 @@ export interface SseRouteOptions { readonly event?: string; /** Serialize one item (default JSON.stringify). */ readonly encode?: (a: A) => string; - /** Keepalive comment cadence in seconds (default 15; 0 disables). */ + /** + * Keepalive comment cadence in seconds (0 disables). Default 5 — deliberately below + * Bun's 10 s default `idleTimeout`, which `servePluginUi`'s server does not raise: + * a slower ping never arrives, because the idle connection is closed first. + */ readonly pingSeconds?: number; } @@ -27,7 +31,7 @@ export const sseRoute = ( ): Layer.Layer => { const event = opts?.event ?? "message"; const encode = opts?.encode ?? ((a: A) => JSON.stringify(a)); - const pingSeconds = opts?.pingSeconds ?? 15; + const pingSeconds = opts?.pingSeconds ?? 5; const frames = stream.pipe( Stream.map((a) => `event: ${event}\ndata: ${encode(a)}\n\n`), diff --git a/plugin-kit/test/sse-live.test.ts b/plugin-kit/test/sse-live.test.ts new file mode 100644 index 00000000..65d08f7e --- /dev/null +++ b/plugin-kit/test/sse-live.test.ts @@ -0,0 +1,79 @@ +// Regression: the production shape of sseRoute — a long-lived PubSub-backed stream +// (the engine's status feed) plus the keepalive. The original suite only covered a +// finite, self-driving stream, which hid the fact that nothing ever reached the wire. +import { describe, expect, test } from "bun:test"; +import { Effect, Layer, PubSub, Stream } from "effect"; +import { HttpRouter } from "effect/unstable/http"; +import { httpApiEnv, sseRoute } from "../src/index.js"; + +/** + * Read until the first bytes arrive or `ms` elapses. One sequential read at a time — + * re-entering read() while a previous read is pending is a spec violation and silently + * swallows data (which is exactly how this harness first lied about the ping path). + */ +const readSome = async (res: Response, ms: number): Promise => { + const reader = res.body?.getReader(); + if (!reader) return ""; + const decoder = new TextDecoder(); + let out = ""; + const timer = setTimeout(() => void reader.cancel().catch(() => {}), ms); + try { + while (true) { + const { done, value } = await reader.read(); + if (done) break; + if (value) out += decoder.decode(value, { stream: true }); + if (out.length > 0) break; + } + } catch { + // cancelled by the deadline + } finally { + clearTimeout(timer); + await reader.cancel().catch(() => {}); + } + return out; +}; + +describe("sseRoute (live, PubSub-backed)", () => { + test("delivers frames published AFTER the request opened", async () => { + const program = Effect.gen(function* () { + const hub = yield* PubSub.unbounded<{ n: number }>(); + const routes = sseRoute("/api/events", Stream.fromPubSub(hub), { + event: "status", + pingSeconds: 0, + }); + const { handler, dispose } = HttpRouter.toWebHandler( + Layer.provide(routes, httpApiEnv), + ); + const res = yield* Effect.promise(() => handler(new Request("http://127.0.0.1/api/events"))); + expect(res.status).toBe(200); + // Publish only once the response is open — the real engine's pattern. + setTimeout(() => { + Effect.runFork(PubSub.publish(hub, { n: 1 })); + }, 50); + const body = yield* Effect.promise(() => readSome(res, 3000)); + yield* Effect.promise(() => dispose()); + return body; + }); + const body = await Effect.runPromise(Effect.scoped(program)); + expect(body).toContain('event: status\ndata: {"n":1}'); + }); + + test("emits a keepalive on an otherwise silent stream", async () => { + const program = Effect.gen(function* () { + const hub = yield* PubSub.unbounded<{ n: number }>(); + const routes = sseRoute("/api/events", Stream.fromPubSub(hub), { + event: "status", + pingSeconds: 1, + }); + const { handler, dispose } = HttpRouter.toWebHandler( + Layer.provide(routes, httpApiEnv), + ); + const res = yield* Effect.promise(() => handler(new Request("http://127.0.0.1/api/events"))); + const body = yield* Effect.promise(() => readSome(res, 4000)); + yield* Effect.promise(() => dispose()); + return body; + }); + const body = await Effect.runPromise(Effect.scoped(program)); + expect(body).toContain(": ping"); + }); +});