PostgreSQL storage implementation¶
The PostgreSQL adapter persists processing state, source progress, output intent, quarantine, publication receipts, and the materialized view. Processing decisions, publication receipts, and view application use their respective transaction boundaries described below. It is exercised against the local PostgreSQL 18.6 instance using pgx v5.11.0.
The broker source worker and relay/view transport are implemented separately. The bounded live verifier and inspection metrics use these APIs; the process-crash showcase observes rollback and recovery across actual worker death. The storage integration tests deliver envelopes through a database test harness; the delivery suite also consumes actual broker output. A passing integration test is not itself a product verifier PASS.
Running the database checks¶
From the repository root:
make test-postgres runs fresh race-enabled tests with the integration build tag. It uses AMENDS_TEST_DATABASE_URL, defaulting to the synthetic local endpoint already configured in Compose:
The tests fail if the selected database is unavailable. Direct tagged test execution also fails if the variable is missing; it does not silently skip the suite. Normal make check still requires no services and runs the core/model tests while compiling the database adapter.
Each test creates a unique namespace and removes only its own rows on cleanup. The schema and migration ledger remain. No existing namespace, broker topic, or Docker volume is reset. An interrupted test can leave its synthetic namespace behind; there is no automatic sweep of unrelated state.
make evidence-postgres writes selected toolchain and module versions plus JSON test events under artifacts/postgres/. The environment test records the server version and requires fsync=on. The adapter sets synchronous_commit=on for its transactions. The CI workflow has a separate PostgreSQL service job running the same evidence target; hosted results must be observed separately.
Schema and representation¶
Migrations 001, 002, 003, and 004 are embedded in order by migrations/embed.go. postgres.Migrate checks the contiguous installed history and SHA-256 checksums, then applies missing migrations transactionally under an advisory transaction lock. Repeated installation is idempotent. An unknown version, gap, or changed applied migration fails. Migration 002 preserves the original migration and existing processing data while adding output identity and raw protocol diagnostics.
| Table | Stored content |
|---|---|
amends_schema_migrations |
Applied migration version, checksum, installation time |
amends_namespaces |
Configuration bytes, configuration digest, storage format version |
amends_partitions |
Ownership epoch, exclusive source frontier, emission sequence, optional last committed processing measurement, cumulative statistics, bounded recent history |
amends_key_state |
Complete per-key revision evidence and computed projection |
amends_output_state |
Latest committed envelope and version for each output ID |
amends_outbox |
Exact committed envelope bytes, partition sequence, creation time, acknowledgement time, and observed output offset |
amends_quarantine |
Raw rejected record, source position, and rejection reason |
amends_output_binding |
Immutable output topic name, broker incarnation, and partition count |
amends_view_progress |
Exclusive output frontier, stopped-partition reason, and exact offending offset/key/value |
amends_view_rows |
Highest retained envelope, version, and presence/deletion marker |
Namespace IDs are JSON-encoded strings, and stock/output IDs are JSON-encoded tuples stored as TEXT. Configuration, state, records, and envelopes are JSON bytes stored as BYTEA. This preserves the exact outbox bytes and handles identifiers containing escaped control characters, including NUL, without relying on PostgreSQL JSONB's representation limits.
For this small implementation, revision candidates, conflicts, sparse balances, and alerts are serialized together in one row per stock key. They are not separately indexed relational tables. Epochs, frontiers, emission sequences, output versions, presence, and publication receipts have explicit SQL columns and constraints. This keeps the first adapter inspectable but makes each key update a replacement of its serialized state.
The serialized Go state format remains version 1; the SQL migration history now reaches version 4. These are separate version spaces. Migration 003 adds a nullable measurement without rewriting historical state or backfilling fabricated costs. Migration 004 adds nullable cumulative statistics and the latest 32 samples; earlier history is unknown. Upgrades from populated schemas 1, 2, and 3 are tested. Incompatible persisted state changes need a deliberate format/migration decision. Evidence cleanup, deletion-version cleanup, and live namespace cutover are not implemented.
Adapter entry points¶
The adapter is internal/store/postgres. Shared snapshots, tokens, committed deliveries, and view input/position values live in internal/store/types.go. Worker and delivery protocols use store.ErrFence, store.ErrCommitUnknown, and store.ErrProtocol from errors.go; PostgreSQL session guards and ErrServiceBusy remain adapter-specific. The memory model retains aliases for its existing snapshot/token callers. There is no general Store interface.
| Entry point | Responsibility |
|---|---|
Migrate(ctx, pool) |
Install or validate the schema |
New(pool, config) |
Validate and copy configuration without database mutation |
Store.Ensure(ctx) |
Initialize a namespace or validate its existing immutable configuration |
BindOutput / OutputBinding |
Install or read immutable output identity; reject adoption of previously delivered unbound state |
AcquireService |
Hold a dedicated session lock for a producer, relay, or view; no automatic failover |
Progress |
Observe epochs, source/view positions, pending count/age, receipts, quarantine count, protocol status, the last processing measurement, cumulative statistics, and bounded history in one snapshot |
Receipt |
Read publication outcome after an unknown acknowledgement commit |
ApplyViewInput / ViewPosition |
Apply decoded output or retain raw protocol failure; observe durable recovery position |
Claim(ctx, partition) |
Install a new epoch under the partition row lock |
Observe(ctx, partition) |
Read installed epoch and source progress |
Begin(ctx, token, expected) |
Lock the partition and validate the processing guard |
ProcessingTx.Apply(record) |
Prepare all database writes for one input decision |
ProcessingTx.Commit() |
Commit the entire decision |
ProcessingTx.Rollback() |
Abort and release the transaction |
Process(ctx, token, expected, record) |
Convenience wrapper around begin/apply/commit |
Pending(ctx, partition) |
Read the earliest unacknowledged envelope and its exact bytes |
Acknowledge(ctx, partition, sequence, outputOffset) |
Record an ordered publication receipt |
ApplyView(ctx, partition, expected, offset, envelope) |
Atomically materialize one output and its progress |
Snapshot(ctx) and View(ctx) |
Read consistent processing and view snapshots separately |
Ensure never overwrites an existing manifest or fills in missing state for an existing namespace. Configuration disagreement fails. Migration and initialization are explicit operations; constructing a Store does not mutate the database.
Processing and ownership¶
Begin obtains SELECT ... FOR UPDATE on the namespace/partition row before reading mutable key or output state. It checks the token's namespace, configuration digest, positive epoch, and expected source frontier against the installed values.
Apply uses the pure fold.Validate, fold.ApplyMeasured, and diff.Reconcile functions. fold.Apply delegates to the same measured implementation; the independent oracle does not. Accepted input writes key evidence/projection, output state, and outbox envelopes. Rejected input writes its exact raw record and reason. Both persist the processing measurement, cumulative statistics, and bounded recent history in the same final UPDATE as the source frontier and emission sequence. Statistics are read while the existing partition row lock is held. Rollback/fencing contribute nothing, and replay preserves the observations. Failure of that write poisons the whole transaction. Commit retains its existing commit-outcome classification without added precommit I/O. See measurement definitions and timing limits.
ProcessingTx.recordProgress owns the final measurement/statistics read and frontier UPDATE inside that same transaction. An Apply error, including a failure in this helper, poisons the transaction: Commit rolls back rather than permitting partial writes. Only one record is accepted per processing transaction. Callers must always defer Rollback; a successful commit makes that cleanup harmless.
Claim locks the same partition row before incrementing the epoch. If processing already owns the lock, it can finish before the replacement claim. If the replacement wins first, old processing sees the newer epoch after waiting and fails its guard. The integration test checks database-observed blocking through pg_blocking_pids, rather than assuming that a sleep proves one connection blocked another.
The adapter bounds operations with a 15-second context, a 3-second lock timeout, a 10-second statement timeout, and a 15-second idle-in-transaction timeout. These are limits for the small local implementation, not a latency promise. A deep correction can fail its budget without advancing progress. Cleanup uses a separate bounded context because pgx does not automatically roll back when the operation context is cancelled.
The store has no knowledge of Kafka assignment validity. Observe reports installed state; it does not authorize a revoked worker to resume. The source worker implements that separate lifecycle guard; the database fence remains necessary.
Unknown commit outcomes¶
A successful commit call returns nil. An unclassified commit failure returns an error matching ErrCommitUnknown; it never reports optimistic success or assumes rollback from a lost connection. Explicit pgx closed/rolled-back transaction results remain distinguishable.
After an unknown processing result, read durable epoch/progress and resume from the installed frontier. After an unknown claim, reobserve under a still-valid local assignment. After an unknown view result, use its durable output frontier. Do not compensate, lower a frontier, or allocate another output version merely because a reply was lost.
The test socket wrapper exercises two actual database outcomes:
- Close the connection before transmitting COMMIT, leaving PostgreSQL to roll back.
- Send COMMIT, read a successful server CommandComplete and ReadyForQuery, then close the socket before pgx receives those reply bytes.
Both paths return unknown to the adapter caller. Separate connections inspect what actually committed. These tests cover processing, quarantine, ownership claims, publication receipts, and view application. The fault wrapper is compiled only into integration tests.
A separate case terminates the test's PostgreSQL backend after a correction's writes but before commit. It checks that the old projection, outbox, and frontier remain together, then installs a replacement owner and successfully applies the correction once.
These are connection and transaction failure tests. They do not kill a deployed worker process, crash the PostgreSQL server, simulate power loss, or exercise broker acknowledgement loss.
Outbox and publication receipts¶
Output versions are allocated from committed output state while holding the processing lock. Each emitted envelope is serialized once; the same bytes enter the latest output row and the outbox row in the transaction.
Pending returns those bytes and the earliest pending sequence. The relay publishes Delivery.Bytes unchanged. A later attempt may obtain another output offset, but it must not replace the output ID, version, payload, or cause.
Acknowledge records an observed output offset only for the earliest pending row. Receipts must advance within the partition. Repeating an identical receipt is idempotent, which supports recovery after its database commit reply is lost. A conflicting receipt or an attempt to acknowledge later pending work fails.
The acknowledgement transaction also locks the processing partition row. This is a small, explicit serialization cost in the first adapter. It does not hold that lock during broker publication. The runner preserves the single-relay ordering assumption with serial steps and a dedicated session guard. That guard does not fence broker requests already in flight; see the operating limits.
Creation and acknowledgement timestamps provide the stored information needed for outbox-age reporting. No running metrics endpoint or lag dashboard is delivered in this increment.
Durable view¶
ApplyView locks the view-progress row before reading an output's stored envelope. It applies the existing version-aware view decision, then commits the retained envelope, version, presence marker, and next output offset together.
A higher version replaces a lower one. An equal identical envelope or a lower valid version leaves the current fact intact while advancing consumption progress. Withdrawals retain their version and envelope with present=false.
An invalid envelope or conflicting equal version commits a stopped-partition diagnostic while leaving rows and frontier unchanged. A stale caller frontier is an operation error and does not poison the output partition. The diagnostic survives reconnection and prevents further application through the adapter.
View reads rows, deletion guards, diagnostics, and frontiers from one repeatable-read snapshot. Snapshot does the same for processing state. These separate reads do not establish a complete common source/output boundary; the bounded verifier coordinates that boundary using writer guards, a clean relay handoff, and repeated broker/database observations before comparing the downstream view with the independent source-history oracle.
Evidence and remaining work¶
The integration suite covers:
- All three hand-stated fixture stages, checked independently against the oracle and durable view.
- Cancellation authority, higher-revision reactivation, historical conflicts, blocked-key accumulation, and repair ordering after reconnection.
- Rollback and atomic quarantine, emission/epoch/output-version exhaustion, and arithmetic failures.
- Actual partition lock ordering and rejection of stale ownership.
- Exact outbox retries, ordered acknowledgements, and idempotent receipt recovery.
- Durable duplicate handling, conflicting equal-version diagnostics, and deletion guards after reconnect.
- Both actual outcomes behind a lost commit reply for each database mutation boundary listed above.
- Per-record conformance with the memory store for the selected valid and conflict/rejection histories from seed 7.
- Upgrade from populated schemas 1, 2, and 3 to schema 4, idempotent reruns, checksum rejection, and no fabricated historical measurements.
- Committed refold/output measurements and aggregates, bounded history across restarts/ownership changes and offset gaps, unchanged observations on rollback/fencing/replay, both unknown-commit outcomes without double counting, counter saturation, and lock acquisition measured across an actual database-observed wait.
- Immutable output binding, refusal of previously delivered unbound state, and dedicated relay/view session guards including actual backend termination.
The database conformance sample is separate from the 256-run memory corpus. It is not 256 executions against PostgreSQL.
The source worker, output transport, and bounded verifier exercise this store through actual broker adapters. The showcase observes selected process deaths, rollback, recovery, and automatic relay handoff for its owned pipeline. Committed measurements, cumulative counts, and the latest 32 samples per partition are exposed through inspection; long-term charting is outside this adapter.
For the API and locking basis, see pgx v5.11.0, its transaction implementation, and PostgreSQL 18 explicit locking. The repository's tests supply implementation evidence for the limited cases above.