# 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.

Source: https://bettersupabase.com/docs/blocks/streams

`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 [#postgres]

```bash
pnpm better-supabase sql add streams
```

The 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.

```ts title="src/lib/streams.ts"
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 [#redis]

```bash
pnpm add redis
```

```ts title="src/lib/streams.ts"
import { 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 [#writing-and-resuming]

`teeToStore` copies a stream into the store while passing it through to the
live response:

```ts title="app/api/chat/route.ts"
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:

```ts title="app/api/chat/[id]/stream/route.ts"
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 [#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 [#cleaning-up]

Run `streams.purge({ olderThan: "1 day" })` from a [job](/docs/blocks/jobs)
schedule, or schedule `purge_streams` with pg\_cron. It resolves with the
number of streams deleted.

## Your own store [#your-own-store]

`StreamStore` has `apiVersion: 1`. Run `testStreamStore` from
`better-supabase/testing` against your implementation; see
[Interfaces](/docs/extending/interfaces#streamstore).