Files
enricobuehler 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

80 lines
3.0 KiB
TypeScript

// 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<string> => {
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");
});
});