Skip to content

Streaming

A streamed response sends its body as it is produced: server-sent events, a large export, a model’s tokens. Most streams are easiest to write with @tetsujs/sse; a stream written by hand has a few rules to follow.

A handler that streams returns a Response whose body is the stream. A bare ReadableStream or generator is a compile error and a 500 at runtime: as JSON it would be {}. The Response is also where the content type is stated.

Like any Response, it is not checked by a response schema, and ctx.out.headers is applied to it.

The status and headers leave with the first chunk of the body. A stream that has nothing to send for a while keeps the client waiting for the headers too, so write something the format allows as soon as you can.

@tetsujs/sse builds the Response from an async generator, and handles what a hand-written stream tends to get wrong: backpressure, stopping when the client leaves, ending the generator, and reporting a failure.

sse() formats each yielded event for EventSource, sends keep-alive comments so neither Bun nor a proxy closes an idle connection, and sends a first comment so the headers leave at once:

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: "/clock"

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
: "/clock",
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
) {
while (!
signal: AbortSignal
signal
.
aborted: boolean

The aborted read-only property returns a value that indicates whether the asynchronous operations the signal is communicating with are aborted (true) or not (false).

MDN Reference

aborted
) {
yield {
event?: string | undefined

Event name, read by addEventListener(name) rather than onmessage.

event
: "tick",
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
: {
at: string
at
: new Date().toISOString() } };
await Bun.sleep(1_000);
}
}),
});

stream() is the same without the event format, for NDJSON, CSV or anything else:

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: "/orders.ndjson"

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
: "/orders.ndjson",
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* () {
for await (const
const order: Order
order
of
const orders: {
cursor(): AsyncGenerator<Order, void, undefined>;
}
orders
.cursor()) yield `${JSON.stringify(
const order: Order
order
)}\n`;
},
{
contentType?: string | undefined

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

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

The package page covers event ids and Last-Event-ID, keep-alives, onEnd summaries and every option.

A generator must yield or wait on the signal

Section titled “A generator must yield or wait on the signal”

The generator behind sse() or stream() 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 keeps yielding stops on its own: at the next value the stream sees the abort and calls return(), which runs its finally. A generator stuck in an await is not woken by anything. If it waits on a source that has gone quiet, such as a queue with no traffic, it never resumes, and the subscription inside it lives as long as the process.

So a stream ends with the connection only if its generator keeps yielding, or waits on the signal. A generator that can go quiet passes the signal to whatever it waits on:

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
) {
const
const queue: Subscription
queue
= await
const broker: {
subscribe(topic: string, options: {
signal: AbortSignal;
}): Promise<Subscription>;
}
broker
.subscribe("prices", {
signal: AbortSignal
signal
});
try {
for await (const
const price: Price
price
of
const queue: Subscription
queue
) 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: string
at
};
} finally {
await
const queue: Subscription
queue
.close();
}
}),
});

A source that takes the signal, such as fetch, events.on or the timers of node:timers/promises, rejects with an AbortError when the client leaves. That is not reported as a failure.

A generator that throws ends the stream where it stood. The headers have already gone, so the error goes to reportError with source: "stream", and the client sees an ordinary end of stream.

An access log in afterResponse sees a stream when it starts, so it records an hour-long feed as a fast 200. onEnd on sse() and stream() reports how the stream really ended, with its size and duration.

A stopping server waits for every response in flight, and a stream may never finish on its own. until on sse() and stream() ends it when a signal fires. With @tetsujs/lifecycle that signal is draining, and the client reconnects to a server that stays; see Health checks and shutdown.

A client reads at its own pace. enqueue never blocks, so a stream fed from a loop in start() builds every chunk at once, whether the client is reading or gone. A pull() source applies backpressure: the platform calls it only while the stream wants more, so the source advances at the rate the client reads.

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: "/orders.ndjson"

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
: "/orders.ndjson",
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
: () => {
const
const rows: AsyncGenerator<Order, void, undefined>
rows
=
const orders: {
cursor(): AsyncGenerator<Order, void, undefined>;
}
orders
.cursor();
const
const encoder: TextEncoder
encoder
= new TextEncoder();
const
const body: ReadableStream<Uint8Array<ArrayBufferLike>>
body
= new ReadableStream<Uint8Array>(
{
async
pull?: ((controller: ReadableStreamDefaultController<Uint8Array<ArrayBufferLike>>) => void | PromiseLike<void>) | undefined
pull
(
controller: ReadableStreamDefaultController<Uint8Array<ArrayBufferLike>>
controller
) {
const
const next: IteratorResult<Order, void>
next
= await
const rows: AsyncGenerator<Order, void, undefined>
rows
.next();
if (
const next: IteratorResult<Order, void>
next
.
done?: boolean | undefined
done
)
controller: ReadableStreamDefaultController<Uint8Array<ArrayBufferLike>>
controller
.close();
else
controller: ReadableStreamDefaultController<Uint8Array<ArrayBufferLike>>
controller
.enqueue(
const encoder: TextEncoder
encoder
.encode(`${JSON.stringify(
const next: IteratorYieldResult<Order>
next
.
value: Order
value
)}\n`));
},
async
cancel?: UnderlyingSourceCancelCallback | undefined
cancel
() {
await
const rows: AsyncGenerator<Order, void, undefined>
rows
.return();
},
},
{
highWaterMark?: number | undefined
highWaterMark
: 0 },
);
return new Response(
const body: ReadableStream<Uint8Array<ArrayBufferLike>>
body
, {
headers?: HeadersInit | undefined
headers
: { "content-type": "application/x-ndjson" } });
},
});
  • highWaterMark: 0. With the default of 1, the platform calls pull() as soon as the stream is built. A response can be built and never sent (a hook replaced it, an error displaced it, a HEAD request has no body), and that first pull() would start work for nobody.
  • cancel(). A response the pipeline builds and does not send has its body cancelled, and so does a client that leaves mid-stream. cancel() is where a cursor is closed or a subscription dropped.
  • An idle connection. Bun closes a connection that sends nothing for its idleTimeout, 10 seconds unless set. A stream that can stay quiet longer writes something its format ignores on a timer, and sets its request’s timeout above that with ctx.server.timeout(), as sse() does.

To end such a stream on shutdown too, give its source AbortSignal.any([ctx.req.signal, shutdown.draining]), as Cancellation and timeouts combines a deadline with the request’s signal.