Design¶
Semantics defines the normative contract. This document explains component responsibilities, authority, transaction/publication boundaries, independent verification, and architectural tradeoffs. The technical guide maps those decisions to code and walks through the algorithms.
1. Design intent¶
When an earlier stock movement changes, already-published balances and alerts may be wrong. amends repairs those results incrementally and keeps a separately implemented full-history calculation available to check them. Choosing which revision applies is source authority; it follows revision numbers, not arrival order.
The worker recalculates balances and alerts from the earliest affected day forward. That portion of history is the affected suffix, and recalculating it is a refold. Comparing the new results with previously committed results—the diff—identifies updates and withdrawals. Inventory receipts add units, issues remove units, and a daily close is the balance after that day's active movements. Alerts record a transition from a previous close at or above the threshold to a close below it, as worked through in evaluation.
Processing must survive a failure between saving a correction and publishing it. PostgreSQL therefore commits the intended publications with the processing state in a transactional outbox; a separate relay delivers them later. Consumers keep the latest results in a materialized view. Withdrawals retain a deletion version, or tombstone, to stop delayed older messages from recreating removed results.
Workers can overlap during reassignment. The database records an ownership generation (epoch) and refuses old-owner commits after a newer epoch is installed; this is fencing. It also stores the next offset to process, the exclusive source frontier. The view has a separate frontier for published output. These are broker coordinates, not business-record counts. A transaction already holding the ownership row lock may finish before the replacement epoch installs; Kafka reassignment is not that database boundary.
The inventory domain is synthetic. This design explains its own choices without comparing them with a client's implementation. It is intended to be inspectable on a laptop, not to establish a throughput record or a production security posture.
Open the illustration for a full-size view. The arrows distinguish the database commit from later delivery.
The worker-to-Postgres arrow represents one transaction boundary. The Postgres-to-relay-to-view path is asynchronous, with possible repeated delivery; it is not one distributed transaction. The verifier's independent oracle reconstructs applicable revisions, balances, and alerts from the full retained source history and compares them with the downstream view. It does not use the worker's computed answers, although shared decoding and representations can still cause shared defects.
2. Component responsibilities¶
| Component | Owns | Does not own |
|---|---|---|
| Generator | Synthetic configuration, fixture loading, seeded arrival schedules | Production ingestion or business authority policy |
| Source topic | Retained transport history and source coordinates | Choosing authoritative revisions |
| Workers | Partition lifecycle, revision decisions, refold and output diff | Delivery of uncommitted output |
| Postgres processing store | Atomic state, progress, epochs, quarantine, versions, outbox | Atomic downstream application across topics |
| Relay | Repeated delivery of unchanged committed envelopes | Allocating replacement versions on retries |
| Output topic | Transport of versioned business and trust-status messages | Deduplicated current business state |
| View service | Durable highest-version materialization and its input frontier | Undoing an external business action |
| amendsctl | Inspection, bounded verification, controlled replay and isolated rebuild | Inventing source authority or bypassing conflicts |
The implemented inspect -format html dashboard renders live observations as a static local page and does not certify a verification boundary. PostgreSQL-only diagnostic review exposes rejected records and current/superseded canonical conflicts without changing state. The last instrumented processing decision is stored with its source frontier and displayed alongside progress; pure refold work and client-side lock/precommit timings remain distinct from full storage-call attempt logs. Committed cumulative counts and the latest 32 samples per partition are available in the same inspection report. Prometheus/Grafana and long-term charting are not running dependencies. Metric continuity is not itself a correctness dependency; structured verification reports remain usable independently.
3. Storage model¶
These are logical responsibilities. The first SQL schema combines revision evidence, authority, status, balances, and alerts in a serialized per-key state row, with separate relational progress, outputs, outbox, quarantine, and view tables. The source manifest binds a topic name to the broker UUID in immutable configuration and requires retained starts of zero. Migration 002 adds immutable output-topic binding and exact view protocol diagnostics; migrations 003 and 004 add committed processing measurements and bounded aggregate/history collection; the delivery guide documents adoption checks and restart constraints. Favor a small number of clear tables over a generic event-processing framework.
| Logical relation | Required durable content |
|---|---|
namespaces |
Immutable config, ledger origin, schema/code identity, source incarnation and start vector, partition mapping |
partition_state |
Namespace, source partition, current epoch, authoritative next offset, next emission sequence |
revision_evidence |
Key, TxnID, revision, canonical candidate payloads and source references, superseded conflict evidence |
transactions |
Current highest-revision resolution: active, cancelled, or ambiguous |
key_state |
READY/BLOCKED and unresolved authority references |
daily_balances |
Current sparse per-key daily closes used for rewind |
output_state |
Per-ID last derived payload/presence and durable version, including withdrawals |
outbox |
Immutable namespaced envelope, source cause, partition emission order, acknowledgement state |
quarantine |
Raw record or retained reference, exact rejection reason, source coordinates |
view_state |
Highest durable version and visible payload or deletion marker per output ID |
view_progress |
Durable output-consumer next offsets |
All tables are namespaced. Use constraints for uniqueness, nonnegative counters where appropriate, positive revisions/versions, and valid state combinations. Never overload SQL NULL to mean both “not yet seen” and “withdrawn.” Domain errors and concurrency conflicts must remain distinguishable.
Retaining revision candidates is intentionally conservative for a bounded reference example: an equal-revision conflict can be diagnosed even after other input has arrived. Retention cleanup of evidence, output tombstones, or configuration would require its own protocol and is not included.
Processing state and view state may share a local Postgres server in separate schemas. The authority model remains separate: the verifier's oracle consumes source history, and the value being checked is the view. Sharing a server is not a high-availability design.
4. Worker transaction protocol¶
The simplest v1 implementation holds the partition-state row lock across each processing decision. PostgreSQL documents that SELECT ... FOR UPDATE prevents conflicting updates/locks until the transaction ends; this is the database behavior this protocol relies on. See reference R1 in references.
The intended sequence is:
receive record under an active partition assignment
begin database transaction
lock partition_state for this namespace and source partition
validate ownership epoch and expected durable source progress
read required mutable state under the same transaction
classify and persist input evidence
malformed -> durable quarantine
stale / duplicate -> no business changes
conflicting highest authority -> block key and retain evidence
unambiguous change -> update authority and refold/diff when allowed
persist business changes, output versions, and outbox intent
advance authoritative source frontier consistently with the consumed record
commit the whole decision
The implementation may validate immutable record syntax before opening the transaction, but rejection evidence and progress still commit together. Mutable authority or output state must not be read optimistically outside the lock and then committed without revalidation.
Epoch installation takes the same row lock. An old transaction already holding the lock may finish first. Once the new epoch has committed, an older token must fail before committing any part of its decision. Use deterministic lock ordering and bounded database timeouts; do not perform broker publication or wait for a human while holding the transaction open.
A real rebalance adapter must cancel stale assignment work and park revoked workers. A database epoch cannot compensate for an obsolete callback that is incorrectly allowed to claim a new epoch. The bounded client spike below informs that adapter; actual database/lifecycle integration is now exercised in the worker suite.
Pinned client findings¶
A development spike exercised franz-go v1.22.1, Go 1.27.1, and the local Redpanda v26.2.3 broker over its host listener. It used classic consumer groups with the eager range balancer, explicit topic retention, disabled auto-commit, and synthetic records. This is the selected configuration for the source adapter. The application now pins this broker dependency and implements the source adapter described in the worker guide.
| Check | Observed result |
|---|---|
| Assignment and external seek | OnPartitionsAssigned preceded AdjustFetchOffsetsFn. An external frontier of 1 overrode a broker group commit of 4; the first returned record was at offset 1. |
| Unavailable position | Requesting offset 100 in a partition ending at 1, with ConsumeResetOffset(NoResetOffset()), returned OFFSET_OUT_OF_RANGE without records. |
| Processing and rebalance | After PollRecords(ctx, 1) with BlockRebalanceOnPoll, the callback-blocked signal arrived before revocation. AllowRebalance released revocation and the second member fetched. The callback's client context remained live. |
| Delayed offset load | Revocation waited for the offset-adjustment callback; the session context alone could not release the load. Independent cancellation released it. Returning context.Canceled stopped group management, requiring explicit consumer retirement outside the callback. |
| Offset gaps | A read-committed scan returned visible records at offsets 0 and 4 around an aborted transaction whose data was at 2. Two visible records did not imply next offset 2; the coordinate after the second record was 5. |
The pinned source serializes offset-adjustment and revocation callbacks. An initial experiment waiting solely for session cancellation timed out; a second expecting automatic rejoin after returning context.Canceled also timed out. These findings require an independently bounded offset load and an explicit recovery path. Callback contexts and serialization are described in the pinned API/source references (R3); a live client context is not evidence of current partition ownership.
Keep offset adjustment limited to obtaining and validating the external frontier. Do not launch detached epoch installers from it. The adapter must invalidate its assignment token on revoke, loss, or shutdown, cancel/join outstanding assignment work, and serialize token validation with any attempt to install ownership. A context check followed by an unguarded asynchronous claim leaves a race. Unknown claim outcomes require durable re-observation, just like unknown processing commits.
Use an exact frontier with WithEpoch(-1) when no compatible broker leader epoch has been persisted; a broker leader epoch is distinct from the database ownership epoch. Validate topic incarnation and retained boundaries before consumption. The worker suite now also rejects deleted required history, changed topic incarnation/partition count, and a database frontier beyond the end. Poll release must occur on success and failure paths, with processing bounded below the rebalance timeout.
The gap check exercises broker coordinates, not a decision to accept arbitrary transactional input. Source and oracle visibility must agree. Advancing through a trailing run of control/filtered records to a captured end vector still needs a validated adapter rule; neither record count nor a high-watermark alone proves processing completion. The later worker suite exercises eager rebalances, database fencing, guarded cancellation, and recovery from controlled unknown observations. Abrupt member loss, cooperative/server-side groups, and the complete verifier remain outside that suite.
Unknown commit outcome¶
Test both outcomes of a lost commit acknowledgement. If the transaction committed, durable progress and output intent exist and re-reading must not produce another semantic transition. If it aborted, none of that decision exists and the record can be processed normally. Recovery re-reads the database after reconnecting; it must not infer rollback from a timeout.
Store the authoritative source position with the state. Kafka's design documentation discusses coordinating an external consumer position with external output; amends chooses a database-local state/progress transaction and a separate outbox rather than claiming a cross-system atomic transaction (R2).
5. Incremental state and output repair¶
For an ordinary revision, the worker compares old and new active contributions, computes the earliest affected day, restores the latest earlier close, and recomputes the suffix. Deletion of the only transaction on a day removes the balance row. Crossing facts are derived from daily closes, not from transaction-by-transaction intraday motion.
The diff compares the newly derived semantic outputs with committed output_state. It emits only changed facts and required withdrawals. The table represents outputs whose intent was committed, not proof that the view has already received them.
An unresolved highest revision blocks computation for its key. Other arrivals for that key are recorded, but its previous business outputs are frozen and visibly uncertified. Once the key becomes unambiguous, a full key refold avoids having to reason about every suspended incremental rewind. Queue repair messages before the READY status message in that key's partition order.
A blocked status is a control-plane trust signal, not a substitute for displaying the view's progress. Because downstream changes are asynchronous, a live reader can temporarily see an old status or partial repair. The complete verification boundary is the point at which the contract compares the whole business result.
Cost of the deliberately simple approach¶
A deep rewind can touch many rows and hold the partition lock for a long transaction. Other keys on that partition wait behind it; ownership transfer may also wait for that lock. The implemented processing measurements expose the last committed rewind depth, refolded-record count, changed-output counts, and client lock/precommit elapsed time with progress. Source-worker logs time the complete storage call and retain unknown/error observations. Lock acquisition includes more than server wait, and suffix record counts exclude the state-copy/authority-scan/output-read overhead. Cumulative committed counts/timing sums and the latest 32 samples are retained per partition. Earlier collection history is unknown; long-term percentiles and isolated server wait accounting are not implemented.
Do not optimize away the example's inspectability before measuring it. If a larger optional history is slow, report that result with hardware and configuration; do not imply that a faster architecture was necessary or achieved. Embedded state, per-key concurrency, batching, and richer snapshots are later alternatives, each requiring new progress/recovery reasoning.
6. Outbox and relay¶
Allocate a per-partition emission sequence under the worker's partition lock and commit it with the outbox row. The relay operates serially within each partition, selecting the earliest pending committed envelope. Use a single relay service per namespace; competing relay ownership is not another distributed coordination project to add now.
Publish the exact stored envelope, await acknowledgement, and mark it sent. If the acknowledgement is missing, retry that same version and content. A retry does not call the business diff again. A transaction that rolls back must leave no publishable envelope.
Keep acknowledgement information sufficient to explain publication progress and establish the verification drain. Never mark an output delivered merely because a send function was called, and never use an empty in-process queue as evidence that the durable outbox drained.
The implemented relay runner uses exact committed bytes and rereads durable receipts after unknown SQL outcomes. Its session advisory lock prevents duplicate healthy starts but does not fence an already in-flight broker request after session loss. Stop and join an old runner before replacement; the bounded verifier requires a clean stopped-relay handoff and holds its service guard before observing the output end.
There is no need for a separate quarantine topic. The durable quarantine table is the operator-facing source. A future quarantine notification can be another outbox message; it must not reintroduce an uncoordinated write-and-offset problem.
7. Durable view materialization¶
The implemented view is a small consumer with its own Postgres transaction. It compares a received output version with the stored version, conditionally updates the business payload or deletion marker, and advances its output-consumer frontier atomically.
Withdrawn rows remain as invisible tombstones with versions. Equal-version conflicting content is a protocol error; the consumer does not guess which message is right. An older valid envelope is ignored, but its consumption decision is still durable.
The view uses SQL tables queried by the current view-show command; inspect now adds progress, outbox age, offset distances, and source diagnostics. Both are live observations; use the bounded verifier for comparison. It must expose enough information to explain presence, current value, version, trust status, and progress. Version monotonicity alone does not make a group of row updates appear atomically.
8. Independent oracle¶
The oracle implements the source contract independently of the incremental path: read a complete raw source prefix; independently validate and group candidate assertions; choose the unique highest canonical revision per identity; reject or qualify ambiguity; aggregate active quantities by key/day; calculate cumulative closes; derive daily crossings.
The oracle may share domain types and a basic wire decoder, but not the incremental revision resolver, fold, diff, applied-state tables, or stored rewind snapshots. Authority equality and day derivation deserve separate expected-value tests so a shared decoder defect is not invisible to both implementations.
Use hand-calculated expected fixtures to validate the oracle itself. The core property is then a three-way relationship on small named cases: specified expected values, independent oracle, and incremental view. Generated cases broaden coverage without replacing those independently stated examples.
The oracle uses math/big.Int accumulation and full-history grouping, favoring clarity over incremental performance. The incremental fold instead cancels opposite-signed quantities before checked int64 accumulation, so a representable daily close does not fail solely because of an intraday subtotal. Its limitations—shared infrastructure, finite test corpus, and supported source contract—belong in the assessment.
9. Bounded verification protocol¶
Correctness requires comparing the same history on both sides. Input must be stable, processing and delivery must reach the required positions, and the relay must be drained, stopped, and joined. This controlled stillness is quiescence. Bounded verification checks the resulting fixed source/output boundaries within a time budget; zero lag or an empty outbox alone is insufficient.
The implemented verifier acquires the owned source loader guard and the operator-stopped relay guard, captures a fixed source end vector, waits for worker progress, requires a drained outbox, captures the output end vector, and waits for the view. It compares a stable view snapshot with an independent source scan. For independently operated pipelines, the operator drains and stops/joins the relay. The showcase runner automates that handoff for its own finite pipeline; it does not adopt arbitrary processes. The semantics, S10, defines exact statuses.
The report must distinguish unavailable history, incomplete materialization, authority ambiguity, and actual drift. A drift report identifies the first differing key/day or alert and the relevant source coordinates, not merely a total mismatch count.
For the combined incomplete-boundary/conflict case, the overall report is INCOMPLETE with known conflict diagnostics. Affected keys retain their BLOCKED trust status. Once the complete boundary is established, unresolved authority yields overall BLOCKED, as specified in S10.
Use a dedicated namespace and controlled producer for the default demonstration. Verifying an independently changing live environment at historical offsets would require a snapshot/cut mechanism not provided by an unversioned current-state table. Do not approximate that feature by racing two reads and calling their values comparable.
Preserve the source history. Kafka documents retention and compaction as mechanisms that can remove records; their correct configuration is a prerequisite for this example's full-history oracle (R2, R4). Validate equivalent behavior on the pinned local broker; Kafka documentation alone does not establish Redpanda configuration parity.
10. Controllable failure boundaries¶
Pure fold/diff code has no I/O. Worker, relay, and view orchestration expose the effects and observations that recovery decisions need. The delivery scheduler drives production delivery.Stepper; the worker scheduler drives production worker.Lifecycle. Both use controlled memory I/O and logical retry waits. The original sim remains a separate serial transaction model.
Minimal operation boundaries¶
| Actor / step | Inputs and observable result |
|---|---|
| Assignment | Partition and local assignment token; claim/read returns durable epoch and next offset, definite failure, or unknown outcome. An invalidated token cannot start another claim. |
| Worker decision | Record coordinates, expected epoch/frontier, and pure decision; commit atomically changes the S8 state and returns a receipt, definite abort, or unknown outcome. |
| Commit recovery | Re-read durable ownership/progress and decision evidence before choosing resume, retry, or retirement. Never allocate a second semantic transition from an acknowledgement timeout. |
| Relay attempt | Read earliest pending immutable envelope; publication returns acknowledged coordinates, definite non-acceptance, or unknown acceptance. Only acknowledgement permits marking sent. |
| View decision | Envelope and expected view frontier; atomically apply version/tombstone rules and progress with the same three commit observations. |
| Boundary observation | Read retained starts, topic incarnation, broker ends, and durable worker/relay/view progress separately; a fetched position is not a processing receipt. |
Keep durable outcome separate from caller observation. For a lost commit reply the simulator must be able to choose either committed or aborted durable state while returning unknown to the caller. An ordinary I/O error does not establish rollback. Similarly, broker acceptance followed by a lost reply leaves an unchanged outbox envelope eligible for redelivery. These cases require no compensation or new output version just to retry transport.
Production retry loops use bounded attempts, cancellable waits, and real operation deadlines. Schedulers supply logical waits and explicit atomic effects/unknown observations. Exhaustion parks or reports the actor; it never advances a frontier. An obsolete callback cannot acquire fresh ownership, and a lost claim reply requires durable re-observation.
These are bounded serial schedules, not a generic concurrent event engine. Arbitrary concurrent completions and logical deadline-expiry schedules are outside their supported model. Real PostgreSQL lock ordering and actual broker callbacks remain separate tests. The mutation campaign checks whether named semantic assertions detect specific guard removals in isolated copies; it does not add runtime fault modes.
11. Recovery boundaries¶
Ordinary restart seeks to durable progress and preserves output versions. Controlled replay uses intact state and a separate scan cursor over already-processed retained input. It is an idempotency/re-delivery exercise, not state reconstruction.
The implemented operator commands default to dry runs. PostgreSQL replay locks the partition, rechecks epoch/frontier and retained evidence, and commits only a higher ownership epoch on execution. Rebuild creates a fresh namespace and output binding, reads the same complete retained source through production worker lifecycle code, drains its relay/view, and invokes independent bounded verification. Producer session guards use source incarnation so both namespaces exclude the same cooperative writers; relay/view guards remain namespace-specific. Old binaries must be stopped before this guard change. Failed or uncertain destination creation is reported and retained, never silently adopted.
Rebuild reads the full required history into fresh processing and output namespaces. This prevents version collisions and stale-message resurrection across runs. The old namespace remains intact until an explicit local cleanup; no live consumer cutover is implied.
A restored database backup would require compatible source availability, source coordinates, configuration, and output-version history. Backup restore and checkpoint support are identified as future design work rather than hidden inside a broad “replay from any offset” claim.
12. When this pattern does and does not fit¶
The example fits a source that supplies stable transaction identities, ordered revision authority, complete replacement records, and enough retained history to reconstruct derived state. It is useful when consumers can accept explicit revisions to previously published information.
It does not directly solve patch-only feeds with missing revisions, ambiguous source authority, irreversible external actions, global inventory constraints, real-time intraday alerting, or bounded-latency processing under arbitrary deep corrections. Those require different or additional contracts. In some workloads, a well-operationalized periodic batch remains the simpler choice.
The design demonstrates how to ask and test those questions; it does not claim that all batches should be replaced or that logistics is the only application domain.