Skip to content

@tetsujs/sse

@tetsujs/sse turns an async generator into a server-sent events response. It handles the framing, the headers, the keep-alive and the backpressure, which are easy to get quietly wrong by hand. stream() does the same for any other streamed format. For how streaming fits into a handler, see Streaming.

Terminal window
bun add @tetsujs/sse
import { sse } from "@tetsujs/sse";
const feed = route({
method: "GET"

The method this route answers.

Kept as a literal rather than widened to

Method

: the method is half of a route's identity, and an application that remembers its routes — for a generated client, for tooling — needs to know which one this is.

method
: "GET",
path: "/prices"

Route path with :param segments, e.g. "/orders/:id/cancel".

Must start with /, contain no empty segments and no trailing slash; a malformed literal is a compile error.

path
: "/prices",
handler: (ctx: {
readonly req: Request & {
readonly cookies?: Bun.CookieMap;
};
readonly server: Bun.Server<unknown>;
readonly out: Outgoing;
readonly route: RouteInfo;
readonly startedAt: number;
readonly params: {};
}) => Response

The endpoint logic; ctx is fully inferred, never annotate it.

The return type is inferred rather than demanded, and checked twice over. Against the route's own contract, by the

HandlerResult

bound on R: answering with something response never declared is a compile error that says so, instead of a structural diff against Response. And against what the framework can serialize at all, by the intersected

ValidateResult

: a stream handed over bare is refused whether or not the route declared anything, because that is the case no contract covers — without a response schema HandlerResult is unknown and accepts every value there is.

handler
: (
ctx: {
readonly req: Request & {
readonly cookies?: Bun.CookieMap;
};
readonly server: Bun.Server<unknown>;
readonly out: Outgoing;
readonly route: RouteInfo;
readonly startedAt: number;
readonly params: {};
}
ctx
) =>
sse(
ctx: {
readonly req: Request & {
readonly cookies?: Bun.CookieMap;
};
readonly server: Bun.Server<unknown>;
readonly out: Outgoing;
readonly route: RouteInfo;
readonly startedAt: number;
readonly params: {};
}
ctx
, async function* (
signal: AbortSignal
signal
) {
for await (const
const price: Price
price
of
const prices: {
watch(options: {
signal: AbortSignal;
}): AsyncIterable<Price>;
}
prices
.watch({
signal: AbortSignal
signal
})) {
yield {
data: unknown

The payload. A string is sent as it is; anything else is JSON, which is what a browser's EventSource expects to parse. A value JSON has no form for — undefined, a function, a symbol — is refused rather than sent as an empty string: a payload that went missing, a map.get() that found nothing, would reach the page as a valid event with nothing in it.

data
:
const price: Price
price
,
id?: string | number | undefined

Event id. The browser sends the last one back as Last-Event-ID when it reconnects, which is how a stream resumes where it stopped.

id
:
const price: Price
price
.
at: number
at
};
}
}),
});

sse() returns a Response that:

  • sets content-type: text/event-stream and cache-control: no-cache, and x-accel-buffering: no, so nginx, which buffers a proxied response by default, passes each event on as it is written;
  • opens with a comment, : open, so the headers go out at once rather than with the first event;
  • sends a : ping comment every 15 seconds, and sets the request’s idle timeout above that, so neither Bun nor a proxy closes an idle connection (see Idle connections);
  • asks the generator for one event at a time, as the client reads. A client that stops reading pauses the generator instead of filling memory.

A header some other proxy needs, such as no-transform, which asks a proxy not to change the body, goes on ctx.out.headers. The framework lays it over the response sse() returns, and set from a hook, it covers a whole group or application:

const noTransform = hook.beforeHandle((ctx) => {
ctx.
out: Outgoing

Response parameters for the serialized handler result.

out
.
headers: Headers

Headers applied to the outgoing response, whatever produced it — set-cookie values are appended, any other name overwrites.

A live standard Headers, created on first access. It is never assigned, only mutated — set() to own a header, append() to add to it — so two hooks writing headers compose instead of overwriting each other's whole set.

headers
.set("cache-control", "no-cache, no-transform");
});

Each yielded event has these fields:

Field
data the payload: a string is sent as is, anything else as JSON. A value with no JSON form, such as undefined, throws
event the name, read by addEventListener(name) instead of onmessage
id a string or a number; the browser sends the last one back as Last-Event-ID when it reconnects
retry milliseconds the browser waits before reconnecting; a whole number, 0 or more

A multi-line data is framed correctly. event and id must fit on one line: a value with a line break or a NUL throws a TypeError, because sending it would add a field the stream never meant to send.

A browser that lost the connection reconnects with the last id it saw. lastEventId(ctx) reads it, so the stream can continue from there:

import { lastEventId, sse } from "@tetsujs/sse";
const resume = route({
method: "GET"

The method this route answers.

Kept as a literal rather than widened to

Method

: the method is half of a route's identity, and an application that remembers its routes — for a generated client, for tooling — needs to know which one this is.

method
: "GET",
path: "/history"

Route path with :param segments, e.g. "/orders/:id/cancel".

Must start with /, contain no empty segments and no trailing slash; a malformed literal is a compile error.

path
: "/history",
handler: (ctx: {
readonly req: Request & {
readonly cookies?: Bun.CookieMap;
};
readonly server: Bun.Server<unknown>;
readonly out: Outgoing;
readonly route: RouteInfo;
readonly startedAt: number;
readonly params: {};
}) => Response

The endpoint logic; ctx is fully inferred, never annotate it.

The return type is inferred rather than demanded, and checked twice over. Against the route's own contract, by the

HandlerResult

bound on R: answering with something response never declared is a compile error that says so, instead of a structural diff against Response. And against what the framework can serialize at all, by the intersected

ValidateResult

: a stream handed over bare is refused whether or not the route declared anything, because that is the case no contract covers — without a response schema HandlerResult is unknown and accepts every value there is.

handler
: (
ctx: {
readonly req: Request & {
readonly cookies?: Bun.CookieMap;
};
readonly server: Bun.Server<unknown>;
readonly out: Outgoing;
readonly route: RouteInfo;
readonly startedAt: number;
readonly params: {};
}
ctx
) =>
sse(
ctx: {
readonly req: Request & {
readonly cookies?: Bun.CookieMap;
};
readonly server: Bun.Server<unknown>;
readonly out: Outgoing;
readonly route: RouteInfo;
readonly startedAt: number;
readonly params: {};
}
ctx
, async function* () {
for await (const
const item: Item
item
of
const history: {
since(id: string | undefined): AsyncIterable<Item>;
}
history
.since(lastEventId(
ctx: {
readonly req: Request & {
readonly cookies?: Bun.CookieMap;
};
readonly server: Bun.Server<unknown>;
readonly out: Outgoing;
readonly route: RouteInfo;
readonly startedAt: number;
readonly params: {};
}
ctx
))) {
yield {
data: unknown

The payload. A string is sent as it is; anything else is JSON, which is what a browser's EventSource expects to parse. A value JSON has no form for — undefined, a function, a symbol — is refused rather than sent as an empty string: a payload that went missing, a map.get() that found nothing, would reach the page as a valid event with nothing in it.

data
:
const item: Item
item
,
id?: string | number | undefined

Event id. The browser sends the last one back as Last-Event-ID when it reconnects, which is how a stream resumes where it stopped.

id
:
const item: Item
item
.
id: string
id
};
}
}),
});

The generator receives an AbortSignal that fires when the stream is over: the client left, the response was discarded, until fired, or the generator ended. A generator that yields regularly needs nothing more: when the client leaves, its loop ends and its finally runs. A generator that can go quiet must pass the signal on to whatever it waits for, such as a queue with no traffic, or it waits inside an await forever. Streaming shows how.

A generator that throws ends the stream where it stood. The response has already gone, so the error goes to the application’s reportError with source: "stream" (or is printed as [tetsu] stream failed: without one). A finally or an onEnd that throws is reported the same way.

A stopping server waits for every response in flight, and a stream never finishes on its own. until ends it from outside: pass the draining signal of @tetsujs/lifecycle, and clients reconnect to another server.

onShutdownSignals() takes the server, and the server is built from the routes, so the handler reads the signal when a request comes in, after the server and the signal exist:

import { createApp, route } from "@tetsujs/core";
import { onShutdownSignals } from "@tetsujs/lifecycle";
import { sse } from "@tetsujs/sse";
const live = route({
method: "GET"

The method this route answers.

Kept as a literal rather than widened to

Method

: the method is half of a route's identity, and an application that remembers its routes — for a generated client, for tooling — needs to know which one this is.

method
: "GET",
path: "/feed"

Route path with :param segments, e.g. "/orders/:id/cancel".

Must start with /, contain no empty segments and no trailing slash; a malformed literal is a compile error.

path
: "/feed",
handler: (ctx: {
readonly req: Request & {
readonly cookies?: Bun.CookieMap;
};
readonly server: Bun.Server<unknown>;
readonly out: Outgoing;
readonly route: RouteInfo;
readonly startedAt: number;
readonly params: {};
}) => Response

The endpoint logic; ctx is fully inferred, never annotate it.

The return type is inferred rather than demanded, and checked twice over. Against the route's own contract, by the

HandlerResult

bound on R: answering with something response never declared is a compile error that says so, instead of a structural diff against Response. And against what the framework can serialize at all, by the intersected

ValidateResult

: a stream handed over bare is refused whether or not the route declared anything, because that is the case no contract covers — without a response schema HandlerResult is unknown and accepts every value there is.

handler
: (
ctx: {
readonly req: Request & {
readonly cookies?: Bun.CookieMap;
};
readonly server: Bun.Server<unknown>;
readonly out: Outgoing;
readonly route: RouteInfo;
readonly startedAt: number;
readonly params: {};
}
ctx
) => sse(
ctx: {
readonly req: Request & {
readonly cookies?: Bun.CookieMap;
};
readonly server: Bun.Server<unknown>;
readonly out: Outgoing;
readonly route: RouteInfo;
readonly startedAt: number;
readonly params: {};
}
ctx
, feed, {
until?: AbortSignal | undefined

Ends the stream when it fires, as a client leaving would. What a server that is stopping closes its streams on — see

StreamOptions.until

.

draining exists only once the server does, and the server is built from the routes: the handler reads it when a request comes in, by which time it is there.

until
:
const shutdown: ShutdownHandle
shutdown
.
draining: AbortSignal

Aborts when the server starts to stop: after the pre-stop delay, the moment server.stop() is called — or with stopping, when there is no delay.

What a response that never ends on its own closes on. server.stop() waits for every request in flight, and an event stream or a long poll is one that is always in flight: left open, it holds the stop for the whole of graceMs, and the process then exits with 1, its connections cut. Closed here, it ends at once, and its client reconnects to a server the balancer is still sending traffic to.

Not stopping: that one fires while this server is still being sent traffic, so a client that reconnects at once lands here again, to be closed again, for as long as the delay runs.

Read, like stopping, when a request comes in: the routes are built before the server, and the signal only with it.

draining
}),
});
const
const server: Bun.Server<unknown>
server
= Bun.serve({ ...createApp({ routes: { live } }),
port?: string | number | undefined

The port the server listens on

@default ― process.env.PORT || "3000"

port
: 3000 });
const
const shutdown: ShutdownHandle
shutdown
= onShutdownSignals(
const server: Bun.Server<unknown>
server
, {
preStopDelayMs?: number | undefined

How long to keep serving after the stop was asked for, before the server is told to stop at all. Zero by default.

This is the half of a graceful shutdown that stopping gracefully does not cover. In Kubernetes the SIGTERM and the pod's removal from the Service travel in parallel: the signal arrives at once, while the endpoint change has to reach kube-proxy, the ingress and whatever balancer sits in front, which takes hundreds of milliseconds and sometimes seconds. Stopping the moment the signal lands therefore cuts exactly the requests that are still being routed here — the ones this package exists to protect.

So the sequence is: say we are not ready, keep serving for this long while the news travels, and only then drain. A second signal cuts the wait short, the same way it cuts the grace period short.

Zero by default because the delay is pure cost anywhere the caller is not behind a balancer that has to be told — a test, a CLI, a process nobody is routing to. Set it where something has to hear.

A delay on its own changes nothing. Something has to start answering readiness with a failure while it runs, or the balancer goes on sending traffic and the wait only postpones the same cut. See stopping on the return of

onShutdownSignals

.

preStopDelayMs
: 5_000 });

A controller built from its dependencies takes the signal as a function; see Streams and sockets.

Bun closes a connection that sends nothing for its idleTimeout, 10 seconds unless the server sets another, and the first heartbeat of a quiet feed comes after 15. So when the stream starts, sse() sets its own request’s timeout to the heartbeat and ten seconds more, with ctx.server.timeout(). Every other request keeps the server’s setting.

It sets the timeout rather than raising it: Bun does not say what the server’s idleTimeout is, so a longer one, or 0, is replaced for these requests too. A client that stops reading is then let go after the stream’s timeout, and a heartbeat over four minutes is cut at 255 seconds. The timeout is set as the stream is read, after the handler has returned, so it also replaces one the handler set itself, such as the server.timeout(req, 0) of Bun’s own guide. A generator that wants another sets it, since it starts later.

Where the stream cannot set it, the server’s idleTimeout decides. A feed with heartbeatMs: 0 is closed once it stays quiet longer. Over HTTP/3, Bun ignores a request’s own timeout: set idleTimeout above the heartbeat. On a unix socket it ignores it too, and its types do not take idleTimeout with unix: keep the heartbeat under 8 seconds there, since Bun checks idleness every 4 and its default of 10 can end a connection after 8. Bun waits 255 seconds at most, so a heartbeat over four minutes keeps no connection open.

An access log records a stream when it starts, so a long feed shows as a fast 200. onEnd reports how the stream actually ended:

const live = route({
method: "GET"

The method this route answers.

Kept as a literal rather than widened to

Method

: the method is half of a route's identity, and an application that remembers its routes — for a generated client, for tooling — needs to know which one this is.

method
: "GET",
path: "/feed"

Route path with :param segments, e.g. "/orders/:id/cancel".

Must start with /, contain no empty segments and no trailing slash; a malformed literal is a compile error.

path
: "/feed",
hooks: { beforeParse: [requestId()] },
handler: (ctx: {
readonly params: {};
readonly req: Request & {
readonly cookies?: Bun.CookieMap;
};
readonly server: Bun.Server<unknown>;
readonly out: Outgoing;
readonly route: RouteInfo;
readonly startedAt: number;
requestId: string;
}) => Response

The endpoint logic; ctx is fully inferred, never annotate it.

The return type is inferred rather than demanded, and checked twice over. Against the route's own contract, by the

HandlerResult

bound on R: answering with something response never declared is a compile error that says so, instead of a structural diff against Response. And against what the framework can serialize at all, by the intersected

ValidateResult

: a stream handed over bare is refused whether or not the route declared anything, because that is the case no contract covers — without a response schema HandlerResult is unknown and accepts every value there is.

handler
: (
ctx: {
readonly params: {};
readonly req: Request & {
readonly cookies?: Bun.CookieMap;
};
readonly server: Bun.Server<unknown>;
readonly out: Outgoing;
readonly route: RouteInfo;
readonly startedAt: number;
requestId: string;
}
ctx
) =>
sse(
ctx: {
readonly params: {};
readonly req: Request & {
readonly cookies?: Bun.CookieMap;
};
readonly server: Bun.Server<unknown>;
readonly out: Outgoing;
readonly route: RouteInfo;
readonly startedAt: number;
requestId: string;
}
ctx
, feed, {
onEnd?: ((summary: SseSummary) => void) | undefined

Called once when the stream is over, with what it did — for a stream that went out: one made and never sent has nothing to report.

The gap this closes: afterResponse runs as the response goes to Bun, which for a stream is the moment it starts. An access log therefore records a forty-minute feed as a 200 that took microseconds, and a torn connection as a success. Delivery to the client is not observable in the fetch model and stays that way — but the end of generation is, and that is what this reports.

It carries no request id on purpose. This callback is written at the call site, where ctx is already in scope, so the caller adds whatever identifies the request better than this package could guess.

onEnd
: (
summary: SseSummary
summary
) =>
const logger: {
info(fields: object, message: string): void;
}
logger
.info({ ...
summary: SseSummary
summary
,
requestId: string
requestId
:
ctx: {
readonly params: {};
readonly req: Request & {
readonly cookies?: Bun.CookieMap;
};
readonly server: Bun.Server<unknown>;
readonly out: Outgoing;
readonly route: RouteInfo;
readonly startedAt: number;
requestId: string;
}
ctx
.
requestId: string
requestId
}, "stream closed"),
}),
});
// { events: 412, bytes: 38104, durationMs: 2401882.6, reason: "cancelled" }

reason is "ended" (the generator finished), "cancelled" (the client left, the response was discarded, or until fired) or "failed" (the generator threw). events counts the events, not the keep-alives or the opening; bytes counts everything sent. The summary has no request id: add what identifies the request yourself, from ctx. A stream that was never sent reports nothing.

sse() is stream() with SSE framing on top. For anything else, such as NDJSON, CSV or a model’s tokens, stream() sends the chunks as they are, with the same backpressure, signal and summary:

import { stream } from "@tetsujs/sse";
const exportRows = route({
method: "GET"

The method this route answers.

Kept as a literal rather than widened to

Method

: the method is half of a route's identity, and an application that remembers its routes — for a generated client, for tooling — needs to know which one this is.

method
: "GET",
path: "/rows"

Route path with :param segments, e.g. "/orders/:id/cancel".

Must start with /, contain no empty segments and no trailing slash; a malformed literal is a compile error.

path
: "/rows",
handler: (ctx: {
readonly req: Request & {
readonly cookies?: Bun.CookieMap;
};
readonly server: Bun.Server<unknown>;
readonly out: Outgoing;
readonly route: RouteInfo;
readonly startedAt: number;
readonly params: {};
}) => Response

The endpoint logic; ctx is fully inferred, never annotate it.

The return type is inferred rather than demanded, and checked twice over. Against the route's own contract, by the

HandlerResult

bound on R: answering with something response never declared is a compile error that says so, instead of a structural diff against Response. And against what the framework can serialize at all, by the intersected

ValidateResult

: a stream handed over bare is refused whether or not the route declared anything, because that is the case no contract covers — without a response schema HandlerResult is unknown and accepts every value there is.

handler
: (
ctx: {
readonly req: Request & {
readonly cookies?: Bun.CookieMap;
};
readonly server: Bun.Server<unknown>;
readonly out: Outgoing;
readonly route: RouteInfo;
readonly startedAt: number;
readonly params: {};
}
ctx
) =>
stream(
ctx: {
readonly req: Request & {
readonly cookies?: Bun.CookieMap;
};
readonly server: Bun.Server<unknown>;
readonly out: Outgoing;
readonly route: RouteInfo;
readonly startedAt: number;
readonly params: {};
}
ctx
,
async function* (
signal: AbortSignal
signal
) {
for await (const
const row: {
id: number;
}
row
of
const rows: {
watch(options: {
signal: AbortSignal;
}): AsyncIterable<{
id: number;
}>;
}
rows
.watch({
signal: AbortSignal
signal
})) {
yield `${JSON.stringify(
const row: {
id: number;
}
row
)}\n`;
}
},
{
contentType?: string | undefined

The response's content-type. Omitted, none is set.

contentType
: "application/x-ndjson" },
),
});

stream() sends no opening and no keep-alive unless asked, because not every format has a line a client will ignore. The headers go out with the first chunk, so a stream that is slow to produce one keeps its client waiting for the headers too. Yield early, or set a keep-alive the consumer’s parser ignores, such as a blank line:

const
const response: Response
response
= stream(ctx, chunks, {
keepAlive?: KeepAlive | undefined

A filler written while nothing else is, which also sets the request's idle timeout above its interval — see

KeepAlive

. Off by default, and then the server's idleTimeout ends a stream that writes nothing for that long.

keepAlive
: {
everyMs: number

How often, in milliseconds, at most 2³¹ − 1; 0 turns it off. NaN, Infinity, a negative or anything not a number is refused with a TypeError where the stream is made.

everyMs
: 15_000,
chunk: string

What to write. It has to be something the consumer's parser ignores — a comment in a format that has them, a blank line in one that does not, nothing at all in a format where neither is true.

chunk
: "\n" } });

A keep-alive sets the request’s idle timeout, as the heartbeat of sse() does. Without one, Bun closes a stream that writes nothing for the server’s idleTimeout. A stream that should arrive as it is written, behind nginx, sends x-accel-buffering: no in headers, as sse() does on its own.

Without a response map, the generated document says only that the route answers 200, and nothing of what it sends. Name the stream’s type in the map, and the same for stream(), with its own type:

const live = route({
method: "GET"

The method this route answers.

Kept as a literal rather than widened to

Method

: the method is half of a route's identity, and an application that remembers its routes — for a generated client, for tooling — needs to know which one this is.

method
: "GET",
path: "/feed"

Route path with :param segments, e.g. "/orders/:id/cancel".

Must start with /, contain no empty segments and no trailing slash; a malformed literal is a compile error.

path
: "/feed",
schema?: {
readonly response: {
readonly 200: {
readonly contentType: "text/event-stream";
};
};
} | undefined

Validation schemas; an absent part is neither parsed nor typed.

schema
: {
response: {
readonly 200: {
readonly contentType: "text/event-stream";
};
}
response
: { 200: {
contentType: "text/event-stream"
contentType
: "text/event-stream" } } },
handler: (ctx: {
readonly req: Request & {
readonly cookies?: Bun.CookieMap;
};
readonly server: Bun.Server<unknown>;
readonly out: Outgoing & DeclaredOutgoing<200>;
readonly route: RouteInfo;
readonly startedAt: number;
readonly params: {};
}) => Response

The endpoint logic; ctx is fully inferred, never annotate it.

The return type is inferred rather than demanded, and checked twice over. Against the route's own contract, by the

HandlerResult

bound on R: answering with something response never declared is a compile error that says so, instead of a structural diff against Response. And against what the framework can serialize at all, by the intersected

ValidateResult

: a stream handed over bare is refused whether or not the route declared anything, because that is the case no contract covers — without a response schema HandlerResult is unknown and accepts every value there is.

handler
: (
ctx: {
readonly req: Request & {
readonly cookies?: Bun.CookieMap;
};
readonly server: Bun.Server<unknown>;
readonly out: Outgoing & DeclaredOutgoing<200>;
readonly route: RouteInfo;
readonly startedAt: number;
readonly params: {};
}
ctx
) => sse(
ctx: {
readonly req: Request & {
readonly cookies?: Bun.CookieMap;
};
readonly server: Bun.Server<unknown>;
readonly out: Outgoing & DeclaredOutgoing<200>;
readonly route: RouteInfo;
readonly startedAt: number;
readonly params: {};
}
ctx
, feed),
});

OpenAPI 3.1 cannot describe the events one by one, so the stream is documented by its type alone.

sse() Default
heartbeatMs 15000 keep-alive interval, up to 2³¹ − 1; 0 turns it off, and NaN, Infinity or a negative throws a TypeError
status 200
until none a signal that ends the stream, such as draining
onEnd none receives the summary when the stream ends
stream() Default
contentType none the content-type header
status 200
headers none more response headers
keepAlive off { everyMs, chunk }; everyMs follows the rules of heartbeatMs
until none a signal that ends the stream
onEnd none receives the summary; it counts chunks instead of events

The package also exports frame(event), which formats one event as it goes on the wire, and the types ServerSentEvent, SseOptions, SseSummary, SseReason, StreamOptions, StreamSummary, StreamReason and KeepAlive.