---
title: Observe PostgreSQL tables across processes
editUrl: https://github.com/Cratis/Arc.TypeScript/edit/main/Documentation/sql/observing-postgresql.md
description: Install application-owned change triggers and opt into PostgreSQL LISTEN/NOTIFY for Drizzle read models.
---


**Experimental.** This protocol and its versioned trigger names may change before it leaves experimental status. Use [in-process observation](/arc/backend/typescript/sql/observing-tables/) for SQLite, MySQL, or PostgreSQL when only local, explicitly announced changes are needed. PostgreSQL mode observes committed changes from other connections and processes; there is no fallback to in-process observation if the listener fails.

## Install triggers in an application migration

Arc does not create or repair database objects at runtime. For every registered observed PostgreSQL table, generate a drizzle-kit **custom migration**, paste the SQL returned by `postgresqlChangeTrigger(table, { schema: 'app' })` into that migration, and deploy it with your application's migrations before opening observations. For tables declared with `pgSchema('app').table(...)`, omit `schema`; a conflicting override is rejected. Schema-less tables require an explicit migration schema (even if it is `public`). Pass the base `PgTable`, not a Drizzle `alias(...)`; aliases are rejected even if an existing table has the alias name. Run the helper once per table. An existing trigger with the reserved name fails migration installation instead of being overwritten; replace an owned trigger deliberately during an upgrade. Ordinary reads and writes still work without the triggers, but starting an observation fails before its first model read.

```typescript
import { postgresqlChangeTrigger } from '@cratis/arc.drizzle';
import { tasks } from './schema.js';

const migrationSql = postgresqlChangeTrigger(tasks, { schema: 'app' });
// Write migrationSql into an application-owned drizzle-kit custom migration.
```

The helper creates `arc_notify_changes_v1()` and `arc_changes_v1` in the table's schema: an invoker-security, statement-level AFTER INSERT/UPDATE/DELETE/TRUNCATE trigger. It publishes only the canonical, SQL-quoted `schema.table` name to the fixed `arc_changes` channel, not row contents, keys, or tenant identifiers. The channel is database-wide and is **not** an authorization boundary. Duplicate notifications or extra snapshots are possible; NOTIFY delivers no durable history. The reader must see current committed data. Queries with long-lived transaction snapshots or lagging replicas are unsupported.

## Connect one dedicated listener per active tenant

Install `pg` ^8 as an application dependency when using `nodePostgresListener`; `@cratis/arc.drizzle` does not load `pg` at runtime for other modes. The listener factory must return a **new, unconnected, dedicated** `pg.Client` (not a `Pool` client or the connection Drizzle uses). Arc connects and closes it. Use a direct PostgreSQL connection or session-pooled proxy, not PgBouncer transaction/statement pooling. The factory receives a cancellation signal for this tenant listener, not a request's lifetime.

```typescript
import { Client } from 'pg';
import { DrizzleDialect, DrizzleObservation, nodePostgresListener } from '@cratis/arc.drizzle';

builder.withDrizzle({
    dialect: DrizzleDialect.PostgreSQL,
    databaseFactory: (tenant) => queryDatabaseFor(tenant),
    readModels: [{ type: TaskRecord, table: tasks }],
    observation: {
        mode: DrizzleObservation.PostgreSQLNotify,
        listener: async (tenant, { signal }) => {
            const configuration = await listenerConfigurationFor(tenant, signal);
            return nodePostgresListener(new Client(configuration));
        }
    }
});
```

The factory and Drizzle query database must route each tenant to the **same primary database**. A schema-less Drizzle table resolves through the **reader's** `search_path`, not the listener's. Keep every connection in a reader pool on the same stable role and `search_path`; request-dependent `SET LOCAL` is unsupported. Separate schemas in one database route by canonical schema/table payload; tenants reading the same physical table receive the same invalidations. Table filters and RLS remain your responsibility. Temporary tables, views, partitioned tables, partitions, tables with inheritance parents or children, and aliased `PgTable` read-model registrations are unsupported. A statement trigger fires only on the table a write names: a write through a parent never reaches a child's trigger, and a write to a child never reaches the trigger of a parent whose reads include the child's rows. Register the physical base table instead; alias names are not physical trigger targets. Normal writes with `session_replication_role = replica` may skip the trigger and are unsupported.

`observe()`, `observeById()`, and `observePage()` open a lease before reading rows. The first lease for a tenant opens LISTEN; every lease then resolves the table it reads, checks that the listener is connected to the same database, and checks the trigger in `pg_trigger` and its helper in `pg_proc`. Missing, disabled, or incompatible triggers fail that observation; Arc does not attempt DDL or silently degrade. A plain HTTP GET of an observable query calls `current()` and likewise incurs a dedicated per-tenant listener and catalog reads until its scope closes. One tenant with several active subscriptions shares a connection; closing the final lease ends it. Budget one listener backend per active tenant and monitor `pg_stat_activity`. `notifyChanged()` still validates its registered target but publishes **nothing locally** in PostgreSQL mode: the database trigger is the sole invalidation source.

On listener `error`/`end` or heartbeat transport failure, Arc pauses reads, retains each subscriber's last snapshot, and requests a fresh dedicated connection. The internal reconnect delays are 500 ms, 1, 2, 4, and 8 seconds; each connection, listener and catalog operation has a 5-second deadline. After LISTEN is acknowledged, Arc re-resolves each live reader mapping and validates their triggers, then forces a fresh read of **every live subscription, including an unused `current()` prime**. A changed mapping or incompatible trigger fails the affected observation rather than silently switching tables. New subscriptions wait for recovery before their first model read. A successful heartbeat resets the retry budget; if the bounded attempts fail, all remaining tenant subscriptions error, and a subsequent subscription can start a new listener. Arc never falls back to in-process notifications. Trigger validation occurs on acquisition and reconnect, **not** during heartbeat checks; coordinate migrations with observer restart.

Committed changes invalidate matching observations, and recovery catches up to the eventual committed state. PostgreSQL NOTIFY is not a durable event log: intermediate states during a disconnect are not replayed, and duplicate notifications can cause extra snapshots. Pages are not atomic count-and-page snapshots. The installed `@cratis/arc` frontend's direct SSE `EventSource` retries a closed stream automatically; its SSE hub also reconnects and re-subscribes active queries. Custom SSE/HTTP consumers must reconnect after a terminal recovery error; an error delivered as a result frame may require application-level retry. For half-open connections Arc sends `SELECT 1` every 30 seconds. Watch `pg_notification_queue_usage()` as well: a stalled listener in a long-running transaction can fill the notification queue and make writers fail. Arc's listener only executes autocommit statements, never `BEGIN`.
