Lab 03: The event arrives twice
Add a queue and task workers. Implement an outbox and deduplication; test a crash between saving a result and acknowledging a message.
Lab source, setup, and report template
Lab 3. Outbox, redelivery, and one local effect
Add asynchronous processing to the same service. Connect two local transactions without claiming a shared etcd+NATS transaction. Prerequisites: Labs 1–2, transactions, messaging and idempotency. Plan for roughly 90 minutes plus independent completion.
Run
Prepare dependencies using the README. From the repository root:
PYTHONDONTWRITEBYTECODE=1 /tmp/ds-course-venv/bin/python labs/distributed-systems/runner.py --lab 3 --solution starter --report /tmp/lab03.json
Complete the create and complete TODOs in starter/protocol.py. Relay/worker transport is provided in service.py; study its ordering and fault hooks. Select --solution reference for comparison; assertions are shared.
Assignment
- Make task, idempotency and outbox creation atomic. The relay reads the outbox, publishes an event, and marks it published only after a publish ACK. A crash between these actions permits publication again.
- Atomically commit dedup + result + completed task + counter. Revision comparisons protect the counter against lost updates from concurrent changes. A result stored under a unique key does not by itself protect another effect from repetition.
- Send a consumer ACK only after commit. Preserve a stable event_id. Broker deduplication is intentionally not used so it cannot hide an application defect.
The effect increments a local etcd counter by amount. An external HTTP request, email or payment lies outside this transaction and needs a different protocol. Computation before commit may repeat.
Three scenarios
- Committed outbox, no publication yet. The relay crashes immediately before publishing. etcd retains the pending intent while the stream is empty. Restarting the relay must eventually complete the task.
- Lost publish ACK. The relay publishes to real JetStream with a reply inbox that has no subscriber. The broker stores the message, but the application receives no ACK and does not mark the outbox. Retrying creates a second copy; the worker must apply the effect once. Real stream counts, dedup records and the counter verify this.
- Committed result, lost consumer ACK. The worker commits its result and exits before ACK. The durable consumer redelivers after AckWait. The new worker recognizes event_id, does not increment the counter again, and acknowledges the message.
Invariant: every accepted task retains publication intent; after dependencies recover and attempts continue, the task completes and its local effect is applied at most once during retention. Dedup keys remain for the full run. This does not promise unique delivery, a single computation, or completion during indefinite unavailability.
Deliverables
- Protocol diff and a diagram of the two atomic boundaries plus the separate broker ACK.
- JSON for all three scenarios, including stream count, redelivery and counter. The supplied sequence uses effects 2, 3 and 5: total 10 and three results even when more than three messages exist.
- One task’s event chain: publication/retry → delivery → commit → redelivery → suppression → ACK.
- A report and ADR explaining retention, atomic scope, and the limits of the shared counter teaching model.
Reading and discussion
Kafka distinguishes logs and consumers from application transactions; Helland frames local atomicity; Sagas introduces compensation when effects cannot share a transaction. The corresponding research decks support instructor preparation. The actual messaging mechanism follows JetStream consumers.
Why can ACK-before-commit lose a task? Why does commit-before-ACK require deduplication? Why does “the broker supports exactly once” not prove a unique external payment? What happens if dedup expires before redelivery? Extension: design an external effect using the receiver’s idempotency API and state the new responsibility boundary.
Recovery and limitations
One file-store NATS broker supports process restart with its volume preserved; disk destruction is outside scope. The runner restores its processes/networks and removes only its own project as described in README. Never expose the teaching fault hooks on a public service. No scenario substitutes an in-memory fake for etcd or JetStream.
Progress after recovery additionally assumes sufficiently stable and timely communication, a reachable majority, running processes, fair scheduling and finite work. This is not a termination proof for an arbitrarily asynchronous network.