Durable streams
Resumable output for chats, workflows and agents, stored in Postgres or Redis, with a cancel flag the writer reads on its next write.
better-supabase/streams stores the output of a long-running response
(a model generation, a workflow run, an agent step) as ordered text chunks.
A reader that disconnects reconnects with the number of chunks it already
has and gets the rest, live or after the writer finished. A reader can also
ask the writer to stop, and the writer learns on its next write.
Two stores implement the StreamStore interface:
| Store | Where chunks live | Wakes readers with |
|---|---|---|
postgresStreamStore(options) | the streams SQL module | a payload-free Realtime ping, or polling |
redisStreamStore(options) | Redis lists, from better-supabase/streams/redis | a Redis pub/sub message, or polling |
Postgres
pnpm better-supabase sql add streamsThe module adds streams and stream_chunks tables and these functions:
| Function | Who can call it | What it does |
|---|---|---|
stream_open(id, owner, tenant, kind, ttl) | service_role | Creates the stream; returns false when it existed |
stream_append(id, from_idx, chunks) | service_role | Stores the chunks at their indexes, skips indexes already stored, returns the cancel flag |
stream_read(id, from_idx, max) | the owner, service_role | Chunks from an index, through RLS |
stream_status(id) | the owner, service_role | The chunk count and the closed and cancelled flags |
stream_close(id) | service_role | Marks the stream finished |
stream_cancel(id) | the owner, service_role | Sets the cancel flag |
purge_streams(older_than, batch) | service_role | Deletes expired streams and closed ones older than older_than |
stream_append fails with the hint STREAM_NOT_FOUND, STREAM_CLOSED or
STREAM_GAP (an index past the end). A retried batch is harmless, because
indexes already stored are skipped.
import { postgresStreamStore, sqlTransport } from "better-supabase/streams";
export const streams = postgresStreamStore({
transport: sqlTransport(postgres.asService()),
});After each batch the writer pings the stream's private Realtime topic
(stream:<id> by default, sql.modules.streams.options.topic changes the
prefix) without the chunk in the payload. A reader that passes realtime
(a Supabase client) wakes on the ping and reads the new chunks through RLS;
without it, the reader polls every pollMs (250 ms). Set wake: "poll" to
skip the pings for apps near the Realtime message quota.
Redis
pnpm add redisimport { redisStreamStore } from "better-supabase/streams/redis";
export const streams = redisStreamStore({ url: process.env.REDIS_URL });With url, the store loads redis the first time it is used and opens a
second connection for pub/sub. Pass client (and subscriber) instead to
use connections the app already has; any client with the node-redis method
names works. keyPrefix (bs:stream by default) namespaces the keys, and
ttl (one day) sets their expiry. A Redis store has no owner check: the
route that resumes a stream must check the caller first.
Writing and resuming
teeToStore copies a stream into the store while passing it through to the
live response:
import { resumeFromStore, teeToStore } from "better-supabase/streams";
import { after } from "next/server";
const { stream, persisted } = teeToStore(streams, chatId, output, {
owner: userId,
kind: "chat",
onCancel: () => controller.abort(),
});
after(() => persisted);
return new Response(stream.pipeThrough(new TextEncoderStream()));persisted keeps writing after the live client disconnects, so hand it to
after or waitUntil. Chunks are written in batches (every 100 ms or
4 KB). writeToStore does the same for a stream the caller already sends
elsewhere.
A reconnecting client sends the number of chunks it has:
const rest = await resumeFromStore(streams, id, { fromIdx }).orThrow();
if (rest === undefined) return new Response(null, { status: 204 });
return new Response(rest.pipeThrough(new TextEncoderStream()));resumeFromStore returns undefined when there is nothing to resume: no
such stream, or a closed one the client read to the end.
Cancelling
streams.cancel(id) sets the flag. The writer reads it on its next append
and calls onCancel once, so abort the generation there. Call cancel from
a route the owner can reach; with the Postgres store an owner can also call
stream_cancel through RLS.
Cleaning up
Run streams.purge({ olderThan: "1 day" }) from a job
schedule, or schedule purge_streams with pg_cron. It resolves with the
number of streams deleted.
Your own store
StreamStore has apiVersion: 1. Run testStreamStore from
better-supabase/testing against your implementation; see
Interfaces.
Last updated on
Workflow builder
Graph workflows that members edit and publish per tenant, with webhook, schedule and event triggers, credentials by reference, node-level run status and alerts on failed or slow runs.
Notifications
In-app notifications with per-recipient state, subject subscriptions, channel preferences, email and push deliveries, and realtime updates on a private topic.