@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.
: 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: {
readonlyreq:Request& {
readonlycookies?:Bun.CookieMap;
};
readonlyserver:Bun.Server<unknown>;
readonlyout:Outgoing;
readonlyroute:RouteInfo;
readonlystartedAt:number;
readonlyparams: {};
}) => 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, asyncfunction* (
signal: AbortSignal
signal) {
forawait (const
constprice:Price
priceof
constprices: {
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:
constprice: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:
constprice: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:
constnoTransform= 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.
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";
constresume=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: {
readonlyreq:Request& {
readonlycookies?:Bun.CookieMap;
};
readonlyserver:Bun.Server<unknown>;
readonlyout:Outgoing;
readonlyroute:RouteInfo;
readonlystartedAt:number;
readonlyparams: {};
}) => 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, asyncfunction* () {
forawait (const
constitem:Item
itemof
consthistory: {
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:
constitem: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.
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:
: 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: {
readonlyreq:Request& {
readonlycookies?:Bun.CookieMap;
};
readonlyserver:Bun.Server<unknown>;
readonlyout:Outgoing;
readonlyroute:RouteInfo;
readonlystartedAt:number;
readonlyparams: {};
}) => 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:
constshutdown: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
constserver: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
constshutdown:ShutdownHandle
shutdown=onShutdownSignals(
constserver: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:
constlive=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: {
readonlyparams: {};
readonlyreq:Request& {
readonlycookies?:Bun.CookieMap;
};
readonlyserver:Bun.Server<unknown>;
readonlyout:Outgoing;
readonlyroute:RouteInfo;
readonlystartedAt: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.
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";
constexportRows=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: {
readonlyreq:Request& {
readonlycookies?:Bun.CookieMap;
};
readonlyserver:Bun.Server<unknown>;
readonlyout:Outgoing;
readonlyroute:RouteInfo;
readonlystartedAt:number;
readonlyparams: {};
}) => 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,
asyncfunction* (
signal: AbortSignal
signal) {
forawait (const
constrow: {
id:number;
}
rowof
constrows: {
watch(options: {
signal:AbortSignal;
}):AsyncIterable<{
id:number;
}>;
}
rows.watch({
signal: AbortSignal
signal })) {
yield`${JSON.stringify(
constrow: {
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
constresponse: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:
constlive=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: {
readonlyreq:Request& {
readonlycookies?:Bun.CookieMap;
};
readonlyserver:Bun.Server<unknown>;
readonlyout:Outgoing&DeclaredOutgoing<200>;
readonlyroute:RouteInfo;
readonlystartedAt:number;
readonlyparams: {};
}) => 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.
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.