# 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.
- `hub = pg_sse.Hub(DIRECT_DATABASE_URL, channel="room_change")` with a direct (not pooled/PgBouncer) URL.
- `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.
- `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` is the text form of the key column (`str(room_id)` for an integer column); 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)` saying why it runs; 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`.
- 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, 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)`. `max_per_owner` applies only when `owner=` is passed: pass `owner` (user id, session, client address) on public endpoints, or one client can hold every slot. `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.
- Before a planned close (reasons `max_age`, `shutdown`, `behind`, `busy`) the stream sends `retry:` (1 s; 1-5 s jittered on shutdown; 30 s when busy), 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`.
- `hub.stats()` returns `HubStats(open_streams, refused_total, dropped_total, listener_connected, listener_reconnects_total, last_notification_at)`; use it for a /healthz route (503 when `listener_connected` is false).
- `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.
- `hub.wake(key, signal="wake", *, all_keys=True)`: wake streams for a non-database change, from any thread. 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, API table, limits
- [Example app](examples/chat.py): a runnable FastAPI chat room with async reads on a connection pool and a browser page
