# 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)` 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`.
- 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). ERROR (with traceback): a stream whose `read` or `start` raised or returned non-`Event` items; the client reconnects.
- 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)`. `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) 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`.
- `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.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. 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. 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: return `None` from `key` 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. 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), 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.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 `read(cursor, wake)` and `hub.ephemeral`
