psycodict.notifications¶
LISTEN/NOTIFY support for psycodict: schema-change notifications and a small
general-purpose publish/subscribe primitive built on PostgreSQL’s asynchronous
notification mechanism (LISTEN/NOTIFY/pg_notify).
Design¶
Emission is transactional. Notifications are sent with
SELECT pg_notify(channel, payload) executed on the database’s main
connection, through the same _execute path (and therefore the same
transaction and DelayCommit bookkeeping) as every other statement. This is
deliberate: PostgreSQL delivers a notification only when the transaction that
sent it commits, and drops it if that transaction rolls back. A schema change
and its notification are thus atomic – a listener never hears about a
create_table that was rolled back, and always hears about one that
committed. PostgresDatabase.notify and the internal schema hooks both rely
on this; when they run inside a DelayCommit block the notification rides the
surrounding transaction, and when they run standalone _execute commits
immediately (delivering at once).
Listening uses a dedicated connection. A listener must not reuse the main
connection: that connection is busy running the application’s transactions and
buffered server-side cursors, and a connection sitting in a transaction does not
see notifications committed by others until it ends that transaction.
NotificationListener therefore opens its own psycopg connection from
the same Configuration options, in autocommit
mode, and issues LISTEN on it. It is pull-based (synchronous): call
poll() to collect whatever has arrived within a
bounded window, or iterate listen(). There are no
background threads and no thread-delivered callbacks by design – a caller (such
as the LMFDB website) can wrap the pull loop in whatever concurrency model it
prefers.
The schema channel contract. psycodict’s schema-changing operations
(create_table(), drop_table,
rename_table, add_column(),
drop_column and the reload swap reload_final_swap) each emit on the
single channel named by SCHEMA_CHANNEL ("psycodict_schema"). The
payload is just the affected table’s name, as a plain string – nothing else.
Keeping the payload to a bare table name keeps the contract simple and easy to
consume; a richer (e.g. JSON) payload describing exactly what changed is a
possible future extension, and is intentionally not attempted here. A rename
emits twice, once for the old name and once for the new one, so that a listener
can drop the stale metadata and pick up the new table.
Reconnection¶
Reconnection is intentionally not handled automatically. If the listening
connection drops, poll() / listen raise; a
long-lived listener should catch the error and build a fresh listener (any
notifications sent while it was disconnected are, per PostgreSQL semantics,
lost – LISTEN only receives notifications sent after it was issued).
Keeping v1 free of silent auto-reconnect magic makes that loss visible to the
caller rather than hiding it.
Forking¶
A listener, like any libpq connection, must not be used across fork():
parent and child would share one socket, so each would receive an
unpredictable subset of the notification stream, and an explicit close()
in either process sends a protocol Terminate message over the shared socket,
killing the other process’s subscription as well. Create each listener in
the process that will poll it – under a pre-forking web server (e.g.
gunicorn --preload) that means each worker builds its own listener after
the fork, for example lazily on first use, with os.getpid() recorded at
creation time to detect an inherited one. A process that does find itself
holding a listener from before a fork should simply abandon the object: drop
the reference without calling close. That is safe – psycopg
deliberately skips the protocol shutdown when a connection object is
garbage-collected in a process other than the one that created it, precisely
to protect the parent’s copy.
Hot standbys¶
Notifications do not traverse replication. NOTIFY is not WAL-logged, so
nothing is delivered on physical replicas (nor by logical replication, which
publishes only data changes); moreover a server in recovery refuses the
subscription itself – LISTEN raises cannot execute LISTEN during
recovery (SQLSTATE 25006), so building a listener against a hot
standby fails outright. Treat that error as permanent for the server rather
than retrying: a process whose queries go to a standby must fall back to
refreshing on a schedule or on error, or subscribe on the primary (bearing
in mind that a notification can then arrive before the corresponding change
has replayed on the standby).
- psycodict.notifications.validate_channel_name(channel)[source]¶
Check that
channelis a plain identifier and return it, else raise.INPUT:
channel– the candidate channel name
OUTPUT:
channelunchanged, if it is a non-empty string of letters, digits and underscores that does not start with a digit.
- class psycodict.notifications.NotificationListener(config, channels=('psycodict_schema',), **connect_kwargs)[source]¶
Bases:
objectA pull-based subscriber to one or more PostgreSQL notification channels.
The listener owns a dedicated
autocommitconnection (separate from the database’s main connection) built from the same configuration options, and issuesLISTENon each requested channel. Retrieve notifications withpoll()(bounded, returns a list) or by iteratinglisten().It is a context manager; leaving the
withblock closes the connection:with db.listener() as listener: for channel, payload in listener.listen(timeout=60): ...
A listener is bound to the process that created it and cannot subscribe on a server in recovery; see Forking and Hot standbys in the module docstring before using one in a forking application or against a replica.
INPUT:
config– aConfiguration; thepostgresqloptions are used to open the dedicated connectionchannels– a channel name, or an iterable of them (default:("psycodict_schema",)); each is validated and subscribed withLISTEN**connect_kwargs– extra keyword arguments passed on topsycopg.connect(e.g. keepalive settings), overriding the configuration where they overlap
- poll(timeout=0.0)[source]¶
Return the notifications available within
timeoutseconds.Waits up to
timeoutseconds for the first notification, then returns it together with any others already buffered, without blocking further. Withtimeout <= 0(the default) it does not wait at all, returning only what is already buffered. A timed-out poll that saw nothing returns an empty list.OUTPUT: a list of
(channel, payload)pairs, in arrival order.
- listen(timeout=None)[source]¶
Yield
(channel, payload)pairs as notifications arrive.With
timeout=None(the default) this blocks indefinitely, yielding each notification as it is received (until the connection is closed). With a numerictimeoutthe iterator stops after that many seconds, whether or not anything arrived.This is a thin wrapper over psycopg’s own notification generator; break out of the loop (or close the listener) to stop early.