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");
+ });
+});