Files
punktfunk/plugin-kit/src/sse.ts
T
enricobuehlerandClaude Fable 5 589d48a975
rpm / build-publish (43, bazzite, punktfunk-fedora-rpm) (push) Has been cancelled
rpm / build-publish (44, fedora-44, punktfunk-fedora44-rpm) (push) Has been cancelled
docker / build-push (--build-arg FEDORA_VERSION=44, ci, ci/fedora-rpm.Dockerfile, punktfunk-fedora44-rpm) (push) Has been cancelled
docker / build-push (., web/Dockerfile, punktfunk-web) (push) Has been cancelled
docker / build-push (ci, ci/fedora-rpm.Dockerfile, punktfunk-fedora-rpm) (push) Has been cancelled
docker / build-push (ci, ci/rust-ci-noble.Dockerfile, punktfunk-rust-ci-noble) (push) Has been cancelled
docker / build-push (ci, ci/rust-ci.Dockerfile, punktfunk-rust-ci) (push) Has been cancelled
docker / build-push (docs-site, docs-site/Dockerfile, punktfunk-docs) (push) Has been cancelled
docker / deploy-docs (push) Has been cancelled
decky / build-publish (push) Has been cancelled
deb / build-publish (push) Has been cancelled
deb / build-publish-host (push) Has been cancelled
ci / rust (push) Has been cancelled
ci / web (push) Has been cancelled
ci / docs-site (push) Has been cancelled
ci / bench (push) Has been cancelled
arch / build-publish (push) Has been cancelled
apple / swift (push) Has been cancelled
apple / screenshots (push) Has been cancelled
android / android (push) Has been cancelled
fix(plugin-kit): SSE keepalive must outpace Bun's idle timeout (0.1.4)
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 <noreply@anthropic.com>
2026-07-20 19:59:04 +02:00

65 lines
2.1 KiB
TypeScript

// SSE support. `effect/unstable/httpapi` has no event-stream media type (verified at
// beta.99), so the status feed is a raw HttpRouter route beside the HttpApi contract —
// same wire shape the first-generation plugins used (`event: <name>` frames + comment
// pings), which is already proven through the console's reverse proxy.
import { Effect, Layer, Schedule, Stream } from "effect";
import { HttpRouter, HttpServerResponse } from "effect/unstable/http";
const encoder = new TextEncoder();
export interface SseRouteOptions<A> {
/** SSE `event:` name (default "message"). */
readonly event?: string;
/** Serialize one item (default JSON.stringify). */
readonly encode?: (a: A) => string;
/**
* 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;
}
/**
* Register `GET <path>` streaming `stream` as server-sent events. Each subscriber gets
* an independent subscription; disconnects tear the stream down via scope closure.
*/
export const sseRoute = <A>(
path: `/${string}`,
stream: Stream.Stream<A>,
opts?: SseRouteOptions<A>,
): Layer.Layer<never, never, HttpRouter.HttpRouter> => {
const event = opts?.event ?? "message";
const encode = opts?.encode ?? ((a: A) => JSON.stringify(a));
const pingSeconds = opts?.pingSeconds ?? 5;
const frames = stream.pipe(
Stream.map((a) => `event: ${event}\ndata: ${encode(a)}\n\n`),
);
const pings =
pingSeconds > 0
? Stream.fromSchedule(Schedule.spaced(`${pingSeconds} seconds`)).pipe(
Stream.map(() => `: ping\n\n`),
)
: Stream.empty;
const body = Stream.merge(frames, pings).pipe(
Stream.map((s) => encoder.encode(s)),
);
return HttpRouter.add(
"GET",
path,
Effect.succeed(
HttpServerResponse.stream(body, {
contentType: "text/event-stream",
headers: {
"cache-control": "no-cache",
connection: "keep-alive",
"x-accel-buffering": "no",
},
}),
),
);
};