# pg-sse

> Python library: live updates (server-sent events) from PostgreSQL LISTEN/NOTIFY, for Starlette/FastAPI (`hub.sse`) or any framework (`hub.stream`). Notifications carry only a key and mean "re-read from your cursor"; the database stays the source of truth. psycopg 3, Python 3.11+, Postgres 14+.

Install: `pip install "pg-sse[starlette]"` (or `[fastapi]`); plain `pip install pg-sse` is enough for `hub.stream`.

Rules when generating code with pg-sse:

- Migration, once: `conn.execute(pg_sse.trigger_sql("messages", key_column="room", channel="room_change", order_column="id"))`. Never call pg_notify from app code; never put row data in a notification.
- Notify only for writes a reader can see: `trigger_sql(..., on=("insert",))` (any of "insert", "update", "delete") on tables with edits, read receipts or soft deletes, or `update_of=("body",)` to wake for updates that name those columns (the key column is always added). Otherwise every update wakes every reader of the key for a read that returns nothing. Run it again with other options to change an installed trigger.
- Alembic: `op.execute(trigger_sql(...))` in `upgrade` after the table is created, `op.execute(drop_trigger_sql(...))` in `downgrade` before it is dropped. The SQL is several statements: run migrations on the psycopg driver (`postgresql+psycopg://`); asyncpg rejects it. With SQLAlchemy, `read` and `start` are async functions that open an `AsyncSession` (`async_sessionmaker(engine, expire_on_commit=False)`); wire them with `functools.partial`.
- `hub = pg_sse.Hub(DIRECT_DATABASE_URL, channel="room_change")` with a direct (not pooled/PgBouncer) URL (or a function returning one, called on every connect).
- `app = FastAPI(lifespan=hub.lifespan)`, or `async with hub.lifespan():` inside an existing lifespan.
- Route (async def, FastAPI/Starlette): `return hub.sse(room, read=lambda after: read_after(room, after), cursor=parsed_last_event_id, start=lambda: latest_id(room))`.
- Other frameworks: `frames = hub.stream(room, read=..., cursor=..., start=...)` yields encoded SSE `str` frames; it raises `TooManySubscribers` when called (answer 503). `hub.sse` is a wrapper over it; its 503 can be replaced with `hub.sse(..., refuse_with=lambda err: YourError())` (e.g. an exception your error tracker ignores).
- `read(after)` returns `[pg_sse.Event(data, id=row_id), ...]` for rows with id > after, oldest first (`ORDER BY id LIMIT n`). Prefer async functions on a psycopg_pool.AsyncConnectionPool (use `functools.partial(read_after, room)`, not a lambda); plain functions run in a limited thread pool.
- `read` returns a list of Events, never a generator (a generator would run its queries on the event loop). `key` must be a `str`, the text form of the key column: pass `str(room_id)` for an integer or UUID column. `hub.stream`, `hub.sse`, `hub.subscribe` and `hub.wake` raise `TypeError` for any other type (an int would subscribe fine and never be woken). `None` means every key, for streams only. A NULL key notifies `""`.
- Optional `read(after, wake)`: a `read` whose second parameter is named `wake`, or is required, gets `pg_sse.Wake(start, keys, signals, resync, reread, signal_keys)` saying why it runs (`signal_keys[label]` is the keys that signal was sent for; not part of equality); a defaulted second parameter under another name (`limit=500`) is not a wake. Use it only for signal streams that tell the client what to refetch (e.g. `Event({}, event="typing")` when `"typing" in wake.signals`); `wake.data_may_have_changed` is False when only `hub.wake` signals arrived.
- `start()` returns the current max id; pass it as `start=`, never compute it before `hub.sse` (that leaves a gap).
- Parse the Last-Event-ID header (client input) before passing it as `cursor`, with `pg_sse.IntCursor`: `IntCursor(parts=1, *, sep=".", max_digits=15)` (`max_digits` 1 to 19; a part above the bigint maximum raises `ValueError` in both `parse` and `format`); in a route, `cursor=cursors.parse_one(last_event_id)` gives the int or `None` (no header, or one that is not a cursor: start from now, as `EventSource` stops for good on a 422) and never raises for client input, text or raw bytes (`ValueError` if `parts != 1`, `TypeError` for an int, already parsed by the code); `.parse(text)` gives `None` for no cursor (`None` or `""`) or a tuple of ints (also for one part), and raises `ValueError` for anything else; never pass the tuple as `cursor` (a tuple cursor or `Event` id raises `TypeError`, since it could not be sent back as an id; use `.format(*parsed)` for several parts); `.format(*values)` makes the text for `Event(id=...)` (plain non-negative ints only; a bool or float raises `ValueError`). ASCII digits only. Never write `int(x) if x.isdigit()`: `"²".isdigit()` is true and `int` raises, a 500 on client input.
- `pg_sse.signal_read(*, start="ready", resync="resync", change="change", change_data=None, signals=None)` returns an async `read(cursor, wake)` for streams that say what to refetch: event names for the first read, for a resync or quiet re-read (these two send only that event), for keys notified (data `{}` or `change_data(keys)`), and `signals={label: event_name}` in mapping order (sent after `change` when both arrive; data `{}`, or a value `(event_name, data_fn)` gives `data_fn(keys_the_signal_was_sent_for)`; those keys are empty for a hand-built `Wake`, so do not index into them). A name of `None` sends nothing. Any other value, a non-str name, or `change_data` with `change=None`, raises `TypeError`, and an empty name `ValueError`, when `signal_read` is called. With `resync=None` a resync wake still answers its keys and signals. Use it instead of hand-writing the `Wake` if-chain.
- Logging: the `pg_sse` logger, no handlers added. WARNING: a stream refused at a cap or during shutdown; the listener losing its connection (with traceback). WARNING, at most once a minute: a client cursor with a line break or NUL (the stream starts from now); an id longer than an `IntCursor`'s `max_digits`. ERROR (with traceback): a stream whose `read` or `start` raised or returned non-`Event` items (the client reconnects); a failed `on_error` callback.
- Browser: `new EventSource(url)`; it reconnects with Last-Event-ID by itself.
- `Event(data, id=None, event=None, final=False)`: str data sent as is, other data as JSON (`Hub(json_default=str)` for datetimes/UUIDs/Decimals); `final=True` ends the stream.
- `Hub(conninfo, channel, *, max_subscribers=500, max_per_owner=5, max_total_per_owner=None, queue_size=64, max_stream_seconds=900, ping_seconds=(15, 25), reread_seconds=60, json_default=None, end_event=None, handle_signals=True, listen_timeout=5.0, resync_spread=2.0, application_name="pg-sse-listener", read_executor=None, on_error=None)`. `stream` and `sse` also take `max_stream_seconds=` to override the hub's age limit for one stream. `max_per_owner` is per key and owner, and applies only when `owner=` is passed: one owner can still hold that many streams on every key it names. `max_total_per_owner` (default None) caps an owner across all keys. Pass `owner` (user id, session, client address, as text) on public endpoints, or one client can hold every slot. `owner` is a `str` or an `int` (turned into text, so `close(owner=7)` finds `owner="7"`), in `stream`, `sse`, `subscribe` and `close` alike; anything else raises `TypeError`. A stream takes its slot when `stream`/`sse` is called and gives it back when it ends, or, if its response is never sent, when it is garbage-collected (at once on CPython) or 30 s after it was made at the latest. `lifespan` waits up to `listen_timeout` for the first LISTEN and starts without a database anyway. Pings cost no read; a stream re-reads after `reread_seconds` of quiet. After the listener connects or reconnects every stream re-reads, each after a random 0 to `resync_spread` seconds (a notification for its key starts the resync early, except for an all-keys stream; any read that re-reads everything in the meantime, which is every read of a `read(cursor)`, counts as the resync; a `read(cursor, wake)` woken by a signal alone reads at once and the resync still follows; `0` reads at once; must be finite). `max_stream_seconds=math.inf` turns the age limit off. `application_name` is the listener connection's name in `pg_stat_activity`. `read_executor` (a `ThreadPoolExecutor`; a process or interpreter pool raises `ValueError`) runs plain `read` and `start` instead of asyncio's shared 32-thread pool. `on_error(key, exc)` is called when a stream fails: `read` or `start` raised, or `read` returned something unsendable (not an `Event`, data `json.dumps` refuses, a tuple id) (as a task, a plain one in asyncio's default pool, not `read_executor`, so a slow tracker holds up neither the response nor the loop, and leaving `lifespan` waits up to 5 s for it before cancelling; after logging; the stream still ends with no goodbye). Durations (here and in `Ephemeral`) take any real number, `Decimal` included, and are stored as `float`; a `bool`, a string or NaN raises `ValueError`.
- Before a planned close (reasons `max_age`, `shutdown`, `behind`, `revoked`) the stream sends `retry:` (1 s; 1-5 s jittered on shutdown), which EventSource applies by itself. Set `Hub(end_event="end")` (or `end_event=` per `stream`/`sse`) when the client is not EventSource: it then also gets `event: end` with `{"reason": ...}`. No goodbye after `Event(final=True)` or a failed `read`.
- `pg_sse.Read` and `pg_sse.Start` are the types of the `read` and `start` arguments, for annotating a helper that wraps `hub.sse`.
- `hub.close(owner=None, key=None)` ends this process's streams opened with that `owner` (pass `owner=` on authenticated streams) and/or on that `key`, with reason `revoked`, so a revoked session does not keep a stream for `max_stream_seconds`. Revoke in the database first (the client reconnects and the route must refuse it); call it in every worker. Called on the loop it takes effect at once; called from another thread it reaches the loop later; either way it spares streams opened after the call. A closed or dropped stream keeps counting against the caps (and in `open_streams`) until it has said goodbye or for 10 s at most, so a client that stops reading cannot pile up connections and a half-open one holds its slot no longer; a reconnect meanwhile can get a 503. Never call it in the stream route to keep one tab per user: the closed tab reconnects and closes the other, back and forth forever. At least one argument; `owner` is a `str` or an `int` (`TypeError` otherwise). `hub.close_all()` ends everything.
- `hub.stats()` returns `HubStats(open_streams, refused_total, dropped_total, listener_connected, listener_reconnects_total, last_notification_at, read_errors_total, ended_total)`; use it for a /healthz route (503 when `listener_connected` is false). `ended_total` is a frozen `EndedTotals(max_age, shutdown, revoked, behind)` of goodbyes sent by reason, e.g. `ended_total.behind` for slow streams; `dataclasses.asdict(hub.stats())` turns it all into dicts. `hub.streams()` lists every stream holding a slot as `StreamInfo(key, owner, age_seconds, ended)`, oldest first, for an admin view (`ended`: closed or dropped, still saying goodbye).
- `drop_trigger_sql(table, key_column, *, order_column=None)` undoes `trigger_sql` in a downgrade migration. `order_column` must be a serial, bigserial or identity column and serialises inserts per key until commit; it can deadlock across keys: keep that transaction short. Its sequence name is resolved once when the SQL runs: run `trigger_sql` again after renaming the sequence.
- `hub.call_later(seconds, callback, *args)`: run a callback on the hub's event loop after a delay, from any thread; does nothing without a running loop. For follow-up on ephemeral state (check a typing signal lapsed, then `hub.wake`).
- `pg_sse.distinct(read, key)` wraps `read` so a snapshot is sent only when it changed: it drops an event whose `key(event)` equals the key of the last one sent; `None` (and a `final` event) always sends, and the event after a `None` is compared with nothing, so it is sent too. The wrapper is async; a sync `read` runs in the hub's `read_executor` (asyncio's pool when called directly), and `key` on the event loop, so keep it cheap. Prefer a `key` that covers everything shown, including the newest row or message id, over `None`. Call it inside the route, once per request (state is per stream; one shared wrapper would starve the second stream). A dropped event's `id` is never applied, so drop only events that would not advance the cursor: make `key` change with every new row, or, failing that, return `None` for any event that carries rows.
- `hub.ephemeral(name, *, ttl, signal=None, grace=0.25, all_keys=False)` returns `pg_sse.Ephemeral`, per-key signals that lapse unless renewed (typing, presence), in this process's memory only: `.set(key, who, on=True)` renews/starts/removes `who` and returns whether what a reader sees changed, waking the key's streams (`wake.signals` contains `signal`, default `name`) only then; `.clear(key, who=None)`; `.live(key)` lists who has not lapsed, in order; `.reset()` forgets everything without waking anyone (tests). Call `set` on every keystroke or heartbeat; a lapse wakes the key by itself. Safe from any thread. Create it once, next to the hub. Not shared between workers.
- Tests: `from pg_sse.testing import next_event, next_frame, goodbye, ended, until, decode`. Run the hub with `async with hub.lifespan():` (leaving it clears subscriptions; a pytest fixture `async with hub.lifespan(): yield hub` needs an async plugin), read `hub.stream(...)` with `await next_event(frames)` (skips pings and the id-only first frame; `event="name"` also skips other events), `goodbye(frames)` for the `retry:` frame, `ended(frames)` to assert the stream is over. Frames decode to `Frame(id, event, data, retry, comment, raw)` with `.json()`, `.is_ping`, `.is_cursor`, `.is_goodbye`. Timeouts raise `TimeoutError` naming the frames seen. `hub.ping_seconds`, `hub.reread_seconds` and an `Ephemeral`'s `ttl` are plain attributes: a test can set them to about 0.2 s after construction (keep `ping_seconds` a `(low, high)` tuple).
- `hub.wake(key, signal="wake", *, all_keys=True)`: wake streams for a non-database change, from any thread, in this process only (other workers' streams are not reached). Use `all_keys=False` for ephemeral per-key state (typing, presence) so the all-keys streams (`key=None`) are left alone; database notifications always reach them.

## Docs

- [README](README.md): usage, SQLAlchemy and Alembic, API reference, limits
- [Example app](examples/chat.py): a runnable FastAPI chat room with async reads on a connection pool and a browser page
- [Signal-stream example](examples/inbox.py): a runnable FastAPI inbox and typing indicator using `signal_read` and `hub.ephemeral`
