A transactional outbox and an idempotent inbox over SQLite, with no runtime dependencies.
"Exactly-once" is not a thing a queue can do. What this repository does is compose two things that are achievable — at-least-once delivery and idempotent application within a bounded window — and then spend most of its test suite trying to break the composition by killing real processes at the worst possible moments.
At a glance. Python 3.12+, zero runtime dependencies, python -m pytest -q runs 146 tests
(the timings and every other figure). The claim is
at-least-once delivery plus idempotent application inside a bounded window,
and the boundary is stated before anything else is claimed. The evidence is three independent
harnesses — 200 seeded schedules of real killed processes,
2,000 seeds of a simulated faulty network, and
three live processes racing on one set of files. The first two have each
been watched to catch deliberately broken implementations rather
than trusted; the third has not — what it caught was
a claim this README used to make and could not support. Start with the
structure, or skip to
the boundary, the numbers,
how to replay a failing seed, or
what is not built.
Three roles, three SQLite files, and a checker that reads all three from outside. Every arrow is a place a process can die, and every one of them is a test.
flowchart LR
app(["your application<br/><i>business writes · handlers</i>"])
subgraph prod["producer process"]
produce["<b>outbox.py</b> · Producer.produce<br/>business row + outbox row<br/><i>ONE transaction</i>"]
end
subgraph rel["relay process"]
relay["<b>relay.py</b> · Relay.run_once<br/>deliver, THEN mark sent<br/><i>backoff is volatile, on purpose</i>"]
end
subgraph cons["consumer process"]
receive["<b>inbox.py</b> · Consumer.receive<br/>inbox lookup · handler · inbox row<br/><i>ONE transaction</i>"]
end
pdb[("producer.sqlite3<br/>business tables · outbox")]
bdb[("broker.sqlite3<br/><b>broker.py</b> — does NOT dedup")]
cdb[("consumer.sqlite3<br/>inbox · ledger, with the pid that wrote it")]
app --> produce
produce --> pdb
pdb --> relay
relay -->|deliver| bdb
bdb -->|pull| receive
receive --> cdb
receive -.->|ack| bdb
checker["<b>verify.py</b> — the checker<br/>missing · duplicated · orphaned<br/>torn · unsent · queued"]
pdb -.-> checker
bdb -.-> checker
cdb -.-> checker
db.py is the one connection those three share the shape of — WAL, synchronous=FULL,
BEGIN IMMEDIATE, nested transactions joining the outer one. messages.py is the key, the kind,
the payload and the digest that makes a key a promise. pipeline.py wires all four over three
files; faults.py names ten seams on the critical path so a test can stop a process at one.
Three design decisions, and each is a crash the tests actually stage:
| Where | The rule | The crash it survives |
|---|---|---|
| producer | the business row and the outbox row are one transaction | a kill between them leaves neither — the torn state cannot be built |
| relay | deliver, then mark sent | a kill between them redelivers; it never loses |
| consumer | the effect and the inbox row are one transaction | a kill after applying and before recording leaves nothing to double-apply |
The inbox remembers a key for window_ms (default seven days) after it applied it, then forgets.
A duplicate that arrives after its inbox record was pruned is applied again. That is not a bug
being tolerated; it is the edge of the guarantee, and every system that claims exactly-once has one,
whether or not it says so.
There is no third option. Remember forever and the inbox grows by one row per message for the life of the system. Forget, and a late enough duplicate is indistinguishable from a new message. This package chooses to forget on a configurable schedule, and to be loud about what that costs:
test_the_boundary_is_exactpins the edge to the millisecond: for windows of 1, 2, 10, 1000 ms and the seven-day default, a duplicate atapplied_at + window - 1is dropped and one atapplied_at + windowis applied again.test_a_duplicate_that_outlives_the_window_is_applied_again_and_the_checker_says_soreaches the edge in simulation on purpose — a 200 ms window against reordered delays of up to 1500 ms — and asserts that every re-application is at least a window after the one before it, which is the only way it can happen.exactly_oncereports such a key asduplicated. The checker never hides it.test_pruning_bounds_the_inbox_by_the_window_and_not_by_the_message_countsays what the window buys, as a number: at a steady five keys per quarter-window the inbox holds exactly one window's arrivals — a high-water mark of 20 rows — and it is the same 20 after 50, 200 and 800 ticks. The bound is a function of the window and the arrival rate, never of how many messages the system has ever handled. Its counterpart,test_without_pruning_the_inbox_grows_by_one_row_per_message, measures the alternative: 1,000 rows for 1,000 messages, forever.
Choosing the window is therefore an operational decision, and the rule is: make it longer than the
longest outage you intend to recover from, because a relay that comes back after the window has
passed will redeliver into an inbox that has forgotten. It is measured entirely on the consumer's
clock, so relay/consumer skew does not affect it; a consumer clock that jumps forward prunes early,
and that direction is unsafe and undetectable here. docs/failure-model.md
§4 has the arithmetic.
An exception is not a crash. It unwinds the stack, runs finally, flushes buffers, and closes
files — none of which a dying process does. Every claim about crash safety in this repository is
therefore made against a real subprocess terminated with Popen.kill() (SIGKILL on POSIX,
TerminateProcess on Windows), and counted afterwards by opening the SQLite files from a different
process that was not alive when the effect happened.
Real kills are slow, so the same property is also asserted in simulation, over thousands of seeds of network faults the kills cannot produce. Two harnesses, one property, and neither is trusted alone — the simulation models a crash as an exception unwinding out of a transaction, and the real kills are what check that model is honest. A third kind of evidence, live processes racing each other, covers the interference neither of those two produces: not one process dying, but several running at once.
flowchart TB
subgraph kills["proof 1 — real processes, killed"]
direction TB
child["<b>tests/crash_child.py</b><br/>a real interpreter doing<br/>produce · relay · consume<br/>over three SQLite files"]
seam["<b>faults.py</b> · HaltAt<br/>10 named seams on the critical path<br/>prints HALT point=… n=… pid=…<br/>then parks on an Event nobody sets"]
parent["the parent test process<br/>reads that line, calls Popen.kill<br/><i>SIGKILL / TerminateProcess —<br/>no finally, no atexit, no flush</i>"]
count["reopens the three files afterwards<br/>and counts with <b>verify.py</b><br/><i>ledger rows carry the pid that wrote them</i>"]
child --> seam --> parent --> count
end
subgraph sim["proof 2 — the simulation harness · src/effect_queue/sim"]
direction TB
clock["<b>clock.py</b><br/>virtual ms · advance-driven<br/>never wall time"]
network["<b>network.py</b><br/>drop · delay · duplicate<br/>reorder · partition · heal"]
world["<b>world.py</b><br/>one event heap · crash · restart<br/>a restarted node keeps only its disk<br/>timers die with their incarnation"]
simulation["<b>simulation.py</b><br/>the real Producer / Relay / Consumer<br/>as two nodes on that world"]
trace["<b>trace.py</b><br/>ordered events · render · digest<br/><i>one seed replays byte-for-byte</i>"]
clock --> world
network --> world
world --> simulation
simulation --> trace
end
kills -. "the same checker, the same property" .-> sim
tests/crash_child.py is a real interpreter doing producer, relay and
consumer work over three SQLite files. It parks itself at a named fault point, prints one line
naming where it is, and the parent kills it there:
HALT point=consumer.after_effect n=6 pid=27184
Ten seams are named in faults.py, and no code between two adjacent
seams writes to disk — so a kill at a seam is the same crash as a kill anywhere up to the next one.
The hard case the spec singles out is the consumer dying after applying the effect and before recording that it did. In a two-table design that is a durable effect and a blind inbox, and the redelivery applies it twice. Here both writes are in one transaction, so the kill leaves nothing and the test asserts exactly that:
proc = ws.spawn(intents=1, halt_point=CONSUMER_AFTER_EFFECT)
assert await_signal(proc).point == CONSUMER_AFTER_EFFECT
assert kill(proc) != 0
assert ws.ledger() == {}, "an uncommitted effect leaked to disk"
assert ws.inbox_keys() == []
final = ws.run(attempt=1, intents=1)
assert ws.ledger() == {key: 1}
assert ws.ledger_pids() == {final.pid} # the restarted process did it, not the dead oneThat last line is the one worth pausing on: the ledger records the pid that wrote each row, so the test can say which process applied the effect, not merely that some process did.
The seeded proof runs 200 schedules. Each draws its own number of intents, its own number of
kills, which of the ten seams each kill lands on, and how the child interleaves produce / relay /
consume work — all from Random(seed), so a failure replays from the seed alone. After the last
restart the pipeline drains and the checker counts:
[crash proof] seeds=200 kills=474 intents=803 redeliveries_dropped>=74
kills_by_point={'consumer.after_ack': 25, 'consumer.after_commit': 40, 'consumer.after_effect': 33,
'consumer.after_inbox': 33, 'consumer.before_effect': 44, 'producer.after_business': 69,
'producer.after_commit': 57, 'producer.after_outbox': 61, 'relay.after_deliver': 53,
'relay.after_mark': 59}
Every intent had exactly one ledger row in all 200. The test also refuses to pass vacuously: it asserts every seed was really killed, that every one of the ten seams was hit by some seed, and that some redelivery really reached the inbox and was dropped.
A checker that passes everything makes every other test here worthless. So
tests/broken.py holds four components, each the honest one with a single change
that would survive code review. Three of them are killed by the real-process harness at exactly the
point their bug opens:
| Broken component | The change | Killed at | The checker says |
|---|---|---|---|
TwoTransactionConsumer |
effect in one transaction, inbox row in a second | consumer.after_effect |
duplicated=((key, 2),) |
MarkBeforeDeliverRelay |
mark sent, then deliver | relay.after_mark |
missing=(key,), and drained — the loss is silent |
TwoTransactionProducer |
business row committed before the outbox row | producer.after_business |
torn=(key,) |
The fourth, ForgetfulConsumer, applies every delivery and records none — at-least-once with no
inbox at all. It needs no crash to fail, only a duplicate, so it is the one the simulation is
pointed at rather than the kills:
test_the_simulation_catches_a_consumer_with_no_inbox runs it over
200 seeds under duplicate=0.5 with every node alive, and 182 of the 200 seeds violate
exactly-once, every failure reported as duplicated.
The framing that matters:
test_every_broken_implementation_passes_when_nothing_crashes
asserts that all three of the killed components are green under a test that never crashes. That
is the argument for killing processes, written as code.
The same evidence in the other direction — the honest implementation was itself broken in a throwaway copy of the tree and watched to go red — is in the CI table below.
Killing one process answers "what survives a death". It says nothing about the deployment anyone reaches for when a backlog is too big: several relays and consumers on the same files, all alive at once. A crash tears a sequence in time; a race tears it between writers, and they break different things.
tests/test_concurrent_processes.py runs three real
processes (tests/race_child.py) released together from a stdin barrier
— the parent spawns all three, checks every one is alive, then writes one line to each and flushes,
so they wake within a scheduler round of each other. No sleep and no clock is involved.
Because the operating system picks the schedule, every assertion is an invariant — true under every interleaving — and anything schedule-dependent is printed instead:
| Three processes do this | Asserted | Test |
|---|---|---|
all produce the same 60 keys |
exactly one claims each key; no conflicts; the business table and the outbox agree | test_racing_producers_write_each_key_exactly_once |
| all relay and consume one backlog | exactly one ledger row per key; nothing unsent, nothing queued; the checker is clean | test_racing_workers_apply_each_effect_exactly_once |
| all run the whole pipeline at once | the same, with all three roles contending on all three files | test_racing_full_pipelines_apply_each_effect_exactly_once |
| deliver and consume more times than there are keys | relayed ≥ keys, applied + duplicates ≥ relayed, applied == keys |
test_no_delivery_is_lost_and_only_the_first_of_each_key_is_applied |
What "schedule-dependent" means in practice is worth showing rather than asserting. Three consecutive runs of the same file, unchanged, on this machine — the counters the tests print and deliberately do not assert:
| Run | Which process claimed the 60 keys | Ledger rows written, by process | Deliveries consumed |
|---|---|---|---|
| 1 | new=[0, 60, 0] |
applied=[0, 22, 38] |
relayed=154 consume_attempts=201 applied=60 dropped_as_duplicate=141 |
| 2 | new=[60, 0, 0] |
applied=[1, 35, 24] |
relayed=160 consume_attempts=260 applied=60 dropped_as_duplicate=200 |
| 3 | new=[26, 0, 34] |
applied=[16, 38, 6] |
relayed=160 consume_attempts=314 applied=60 dropped_as_duplicate=254 |
Nothing in the left three columns is stable: one run had a single producer claim all 60 keys, the
next split them 26/34, and the count of duplicates the inbox had to drop moved from 141 to 254
across three runs of identical code. The only figure that never moves is applied=60, in every
column it appears in, and it is the only one asserted — which is the whole reason the table above
asserts invariants and prints everything else. A test that asserted "a duplicate happened" would be
green on this machine and red on a faster one, and it would be measuring the scheduler rather than
the package.
These tests corrected the documentation rather than confirming it. docs/failure-model.md used
to claim two relays on one outbox "serialise rather than double-deliver". They do not:
Relay.pending reads unsent rows outside the transaction that marks them, so two relays can hold
the same row and both deliver it —
test_two_relays_both_deliver_the_same_row_and_only_one_marks_it_sent
steps through that interleaving by hand. Only mark_sent serialises, and only one relay gets
True. Double delivery is at-least-once working as specified and the inbox makes it harmless, but
the document had claimed a stronger property than the code has, so
§7 of the failure model and
ADR-0010 now say what the tests
show.
What makes the racing consumers safe is BEGIN IMMEDIATE rather than any mutual exclusion the
package adds: the write lock is taken before the first read, so a consumer that waited reads the
inbox row the consumer ahead of it just committed and returns DUPLICATE. A deferred transaction
would read a snapshot from before that commit, decide the key was new, and fail on the upgrade.
Check-then-apply is safe across processes only because the check is inside the write transaction.
Correctness under concurrency is tested here. Throughput under concurrency is not measured, and no figure for it appears in this repository — every write serialises on one file's write lock, so adding processes adds contention as well as parallelism.
Real kills are slow: those 200 schedules meant 474 kills, each with a process spawn before it. The
relay path also needs thousands of runs under network faults, and for that this repository carries its own copy of the layer-3
simulation harness (src/effect_queue/sim/) — a virtual clock, a virtual
network with drop / duplicate / delay / reorder / partition, crash and restart, one seed, one trace.
Nothing under src/ imports time, datetime, socket, http or asyncio, and
tests/test_dependencies.py fails the build if that ever changes.
now is an argument to receive, prune and pending, never a clock the package reads. That is
what makes a run replayable, and test_the_same_seed_replays_the_same_trace asserts two runs of one
seed produce byte-identical traces.
The harness is tested before anything is tested with it — a dropped message is not delivered, a
partitioned pair cannot talk, a message in flight when the partition forms is lost, a restarted node
kept its disk and lost every attribute and timer it held, the same seed replays and different seeds
do not (tests/test_sim_harness.py, 22 tests). If those go red, every
property test above them passes for the wrong reason, which is why CI runs them as their own step.
A restarted relay is where the harness earns its keep: the retry backoff lives in the Relay object,
which a crash discards, so a restarted relay resends everything unacked at once. That is a
redelivery, and redelivery is the contract.
Every figure below came from a command in this repository, and nine of the twelve rows are
re-derived by tests/test_docs.py on every run — the test count by collecting
the suite, the line counts by counting, and the six seed and race counts by reading the constants
the tests loop over. Those nine cannot drift and cannot be invented. The other three are
pyproject.toml, a tuple in faults.py, and a number the crash test prints.
| Tests collected | 146 | pytest --collect-only -q |
| Source lines | 2,065 | *.py under src/ |
| Test lines | 3,203 | *.py under tests/ |
| Runtime dependencies | 0 | pyproject.toml, enforced by test_src_imports_only_the_standard_library |
| Fault points on the critical path | 10 | faults.FAULT_POINTS |
| Real-process crash schedules | 200 | test_exactly_once_across_seeded_crash_schedules |
| Real process kills in that test | 474 | printed by the test |
| Live processes raced against each other | 3 | RACERS in tests/test_concurrent_processes.py |
| Keys contended for in each race | 60 | RACE_KEYS in the same file |
| Simulated relay-path schedules | 2,000 | test_exactly_once_under_every_network_fault_and_crashes_across_2000_seeds |
| Simulated window-boundary schedules | 500 | test_a_duplicate_that_outlives_the_window_is_applied_again_and_the_checker_says_so |
| Simulated broken-consumer schedules | 200 | test_the_simulation_catches_a_consumer_that_commits_the_effect_before_the_inbox_row |
Wall time on this machine (6 cores, CPython 3.14.6). Every run that was timed, not a chosen one or a best one; the last figure in the first row is the run that verified the tree as committed, and the only thing that changed after it was that figure being written into this table:
| Command | Wall time, per run |
|---|---|
pytest -q — all 146 tests |
93.05 s, 82.73 s, 82.43 s, 94.52 s, 96.22 s |
pytest tests/test_crash_exactly_once.py — 15 tests, including the 200 schedules and 474 real kills |
59.82 s, 62.04 s |
pytest tests/test_sim_relay_path.py — the 2,000-, 500- and 200-seed simulations |
23.09 s, 31.96 s |
| CI's "The proof, under attack" step — 63 tests over five files | 101.07 s, 85.19 s, 86.50 s |
pytest tests/test_concurrent_processes.py — 8 tests, three real processes racing |
4.25 s, 4.22 s, 3.82 s |
pytest tests/test_docs.py — the 13 documentation tests |
4.72 s, 4.58 s |
pytest tests/test_dependencies.py — the 8 dependency-audit tests |
0.45 s, 0.64 s, 0.52 s |
pytest tests/test_sim_harness.py — the 22 harness tests |
0.41 s, 0.47 s, 0.96 s |
Earlier shapes of the same suite, kept because they are what was actually measured at the time and because the spread is the point:
| Command | Wall time, per run |
|---|---|
pytest -q — at 145 tests, before every place that states the suite size was gated |
109.30 s, 140.96 s, 131.00 s, 95.37 s, 95.88 s, 96.68 s, 130.42 s |
pytest -q — the suite at 132 tests, before the concurrency proof |
112.57 s, 174.97 s, 191.58 s, 127.70 s |
pytest -q — at 131 tests, before the eleventh documentation test |
117.24 s, 134.51 s, 143.71 s, 150.59 s, 152.67 s |
test_exactly_once_across_seeded_crash_schedules alone, at 132 tests |
90.35 s, 94.21 s |
| the 2,000-seed simulation alone, at 132 tests | 16.66 s |
The spread is twelve real interpreters contending for six cores under whatever else the machine was doing: the same 132 tests took 112.57 s on an idle machine and 191.58 s on a busy one, so the difference between the blocks is the machine and not the suite — the 145-test suite on an idle machine (95.37 s) beat the 132-test one on a busy machine (191.58 s) by more than half. The crash schedules run in a thread pool of 12, because a child spends most of its life blocked on IPC and fsync. Every figure here is a wall-clock measurement of the test suite, never of the protocol — no test reads a clock, and none would assert anything different on a slower machine. In particular there is no throughput number for the package anywhere in this repository, under concurrency or otherwise; measuring that would need a real workload on real hardware.
These come from the simulated world described above — a virtual clock and a virtual network. They are not measurements of any real network, any real broker, or any real deployment, and no number here should be read as one. Every row of this table says "simulation" because every row is one; the row that is not simulation is at the bottom, and it says so too.
| Simulation run | Seeds | Parameters, all virtual | Result |
|---|---|---|---|
Simulation — exactly-once under every injectable fault (test_exactly_once_under_every_network_fault_and_crashes_across_2000_seeds) |
2,000 seeds × 4 intents = 8,000 intents | drop=0.2, duplicate=0.2, delay 1–300 ms, reorder=True, relay/consumer partition 400→1300 ms, supervisor kills each live node with p=0.2 every 300 ms, p=0.1 crash at each consumer fault point, retry_ms=500, tick_ms=100, prune_every_ms=250, settle 10,000 ms |
8,000 intents → 8,000 effects. 35,651 relay sends, 5,578 redeliveries dropped, 8,942 crashes of which 2,816 inside a transaction, 7,264 network drops, 5,729 network duplicates |
Simulation — the documented boundary, reached on purpose (test_a_duplicate_that_outlives_the_window_is_applied_again_and_the_checker_says_so) |
500 seeds × 4 intents = 2,000 intents | window_ms=200, duplicate=0.5, delay 0–1500 ms, reorder=True, retry_ms=400, prune_every_ms=50 |
a late duplicate was re-applied in all 500 seeds: 1,910 of 2,000 keys, 4,371 extra applications, one key applied 6 times. Every re-application at least one full window after the previous one |
| Simulation — the same seeds under the default window | the same 500 seeds | window_ms=604,800,000 (seven days), everything else identical |
0 re-applications |
Simulation — the harness catches a broken consumer (test_the_simulation_catches_a_consumer_that_commits_the_effect_before_the_inbox_row) |
200 seeds | TwoTransactionConsumer — the effect committed before the inbox row — on the fault scenario above |
103 of 200 seeds violated exactly-once; the honest consumer passed all 200 of the same seeds |
Simulation — the harness catches a consumer with no inbox at all (test_the_simulation_catches_a_consumer_with_no_inbox) |
200 seeds | ForgetfulConsumer — applies every delivery, records none — under duplicate=0.5 with no crashes and no partition: at-least-once with nothing to make it idempotent |
182 of 200 seeds violated exactly-once, every failure reported as duplicated. It needs no crash to fail, only a duplicate |
Not simulation — real process kills (test_exactly_once_across_seeded_crash_schedules) |
200 seeded schedules, 803 intents | per seed from Random(seed): 2–6 intents, 1–4 kills, which of the ten seams each lands on, and the child's interleaving. Real Popen.kill() on real interpreters, on this machine |
474 kills; 803 intents → 803 effects; every one of the ten seams hit by some seed |
Three of those five rows are printed by their tests as a line of output. The other two — the
103 of 200 and the 182 of 200 — are counted by their tests but asserted rather than printed
(assert failures and assert all(r.report.duplicated for r in failures)), so the figures above
were re-derived by running the same scenarios over the same seed ranges and counting the results:
python -c "import sys; sys.path[:0] = ['src', '.']
from dataclasses import replace
from effect_queue.simulation import run_scenario
from tests.test_sim_relay_path import CHAOS, BROKEN_SEEDS
from tests.broken import TwoTransactionConsumer
broken = replace(CHAOS, consumer_cls=TwoTransactionConsumer)
print(sum(not run_scenario(s, broken).ok for s in BROKEN_SEEDS), 'of', len(BROKEN_SEEDS))"103 of 200
The raw output of the three rows that do print, from the runs the figures came from:
[simulation] seeds=2000 intents=8000 relay_sends=35651 redeliveries_dropped=5578
crashes=8942 in_transaction_crashes=2816 network_drops=7264 network_duplicates=5729
The test asserts every one of those counts is non-zero before it passes, so a scenario that quietly stopped injecting faults would fail rather than pass easily. Then the boundary, reached on purpose:
[simulation] window=200ms delay<=1500ms duplicate=0.5 retry=400ms seeds=500: a late duplicate was re-applied in 500 seeds (1910 keys); window=604800000ms, same seeds: 0
The checker named every one of those keys duplicated rather than hiding it, and the same seeds
under the seven-day default re-applied nothing at all: the boundary demonstrated in both directions,
which is the only way to show that a bound is a bound and not a hope.
No test reads a clock and no kill is timed by one, so every run above is a seed and replays from the seed alone. Both harnesses replay from the repository root with nothing installed.
A simulated schedule. run_scenario(seed, scenario) returns the world, the checker's report and
the trace; describe(tail=N) prints the last N events:
python -c "import sys; sys.path[:0] = ['src', '.']
from effect_queue.simulation import run_scenario
from tests.test_sim_relay_path import CHAOS
print(run_scenario(1234, CHAOS).describe(tail=8))"seed=1234 scenario=Scenario(intents=4, faults=NetworkFaults(drop=0.2, duplicate=0.2, ...), ...)
intents=4 effects=4 unsent=0 queued=0
trace, last 8 of 99 events:
3922ms #91 duplicate_dropped key='s1234-intent-3' node='consumer'
3922ms #92 send at=4176 copy=0 dst='relay' src='/sitelet?url=https%3A%2F%2Fgithub.com%2Fjamarius-fortson%2Fconsumer'
4138ms #93 deliver copy=0 dst='consumer' src='/sitelet?url=https%3A%2F%2Fgithub.com%2Fjamarius-fortson%2Frelay'
4138ms #94 duplicate_dropped key='s1234-intent-3' node='consumer'
4138ms #95 send at=4382 copy=0 dst='relay' src='/sitelet?url=https%3A%2F%2Fgithub.com%2Fjamarius-fortson%2Fconsumer'
4176ms #96 deliver copy=0 dst='relay' src='/sitelet?url=https%3A%2F%2Fgithub.com%2Fjamarius-fortson%2Fconsumer'
4176ms #97 marked_sent key='s1234-intent-3' node='relay'
4382ms #98 deliver copy=0 dst='relay' src='/sitelet?url=https%3A%2F%2Fgithub.com%2Fjamarius-fortson%2Fconsumer'
The times are virtual milliseconds, not wall clock. A failing property test prints exactly this for
the first three failing seeds, and test_the_same_seed_replays_the_same_trace asserts that two runs
of one seed produce byte-identical traces, so the printout above is the failure and not a
reconstruction of it.
A crash schedule. The same for the real-process proof — one seed decides the intent count, the kill count, which seam each kill lands on, and how the child interleaves its work:
python -c "import sys, pathlib, tempfile; sys.path[:0] = ['src', '.']
from tests.test_crash_exactly_once import run_schedule
with tempfile.TemporaryDirectory() as d:
print(run_schedule(pathlib.Path(d), 4).describe())"seed=4 intents=3 kills=[attempt 0: producer.after_outbox (point #2), attempt 1: relay.after_deliver
(point #12), attempt 2: consumer.after_inbox (point #7)] final={'points': 17, 'produced': 3,
'new': 0, 'relayed': 0, 'applied': 3, 'duplicates': 1}
intents=3 effects=3 unsent=0 queued=0
That is three real interpreters killed at three different seams, and three intents with three
effects. Seed 4 is worth running because it is the seed named in the CI table below: against a
deliberately broken consumer it reported intents=3 effects=4, at the same three seams, because the
schedule comes from the seed and not from what the code does.
Or run one named seam instead of a seeded schedule — each is its own test:
python -m pytest tests/test_crash_exactly_once.py -q -s -k after_effect
python -m pytest tests/test_sim_relay_path.py -q -s -k outlives_the_windowPython 3.12+. Nothing to install to run the example.
git clone <this repository> && cd effect-queue
python examples/order_shipping.py # the worked example, no install needed
python -m pip install -e ".[dev]" # pytest, for the suite
python -m pytest -q # the whole thing, 1-2.5 min here
python -m pytest tests/test_crash_exactly_once.py -q -s # just the kills, with the counts
python -m pytest tests/test_sim_harness.py -q # just the harnessThe example places three orders, has the broker hand one of them out twice, and drops the copy:
placed 3 orders; outbox rows unsent=3
relayed 3; outbox rows unsent=0
the broker handed ord-1002 out again: queued copies=2
consumed: applied=3 duplicates_dropped=1
shipped ord-1001: 1 x lamp
shipped ord-1002: 2 x chair
shipped ord-1003: 1 x desk
intents=3 effects=3 unsent=0 queued=0
every order shipped exactly once
Three objects, and one rule that carries the whole guarantee.
from effect_queue import Consumer, Database, Message, Producer, Relay
orders = Database("orders.sqlite3") # the outbox lives beside the business tables
# --- producer: the business row and the intent, in ONE transaction ------------
producer = Producer(orders)
def place_order(conn): # runs inside the producer's transaction
conn.execute("INSERT INTO orders(id, sku) VALUES (?, ?)", ("ord-1", "lamp"))
producer.produce(Message(key="ord-1", kind="shipping.ship", payload={"sku": "lamp"}), place_order)
# --- relay: deliver, THEN mark sent -------------------------------------------
relay = Relay(orders, retry_ms=1000)
relay.run_once(transport, now=now_ms) # transport.deliver returns => durably accepted
# --- consumer: the effect and the inbox row, in ONE transaction ---------------
def ship(conn, message): # runs inside the consumer's transaction
conn.execute("INSERT INTO shipments(order_id) VALUES (?)", (message.key,))
consumer = Consumer(Database("shipping.sqlite3"), {"shipping.ship": ship}, window_ms=7*24*3600*1000)
consumer.receive(message, now=now_ms) # -> Outcome.APPLIED | Outcome.DUPLICATE
consumer.prune(now=now_ms) # forget keys older than the windowThe rule: a handler's effect must be writes through the connection it is given. A handler that
posts to an HTTP endpoint has stepped outside the transaction, and for that part of its effect this
package offers at-least-once with an idempotency key attached and nothing more — the endpoint has to
dedup on the key itself. This is not advice;
test_an_effect_outside_the_connection_is_outside_the_guarantee
demonstrates the side channel keeping its entry through a rollback that removed the ledger row.
Two supporting pieces: SqliteBroker is a queue that deliberately does not deduplicate (SQS
standard and Kafka do not either — a broker that deduplicated would be doing the inbox's job and
hiding whether the inbox works), and exactly_once(producer_db, consumer_db, broker_db) is the
checker, which opens the files from whatever process calls it and names every missing, duplicated,
orphaned or torn key.
Pipeline wires all four together over three SQLite files. It is what the proof runs and what the
example uses, and it is a usable single-node deployment of the pattern.
If the effect lands in the same database the producer writes to, none of this is necessary: put the business row and the effect in one transaction and the problem dissolves — no delivery, so nothing to deliver at-least-once and nothing to deduplicate. This package is for the case where the effect lands somewhere that transaction cannot reach, which is also the only case where the boundary above is worth paying for. ADR-0009 records that as a rejected alternative rather than leaving it unsaid.
Layer 3 of the agent platform: the substrate under the single-process
components of layers 1 and 2. It imports none of them, they import none of it, and it imports
none of its four layer-3 siblings either — this repository carries its own copy of the simulation
harness for exactly that reason, and
ADR-0011 records the shared
package that was rejected to keep it that way. What ties the layer together is the
deterministic-simulation discipline, not an import graph, and the rule is enforced rather than
promised: test_src_imports_only_the_standard_library subtracts sys.stdlib_module_names from
every import under src/ and fails on whatever is left, so a sibling package would fail the build
the day it appeared in an import line.
The coupling to write down (spec §6): the idempotency key is the control plane's, and the dedup
window is the consumer's. The layer-1 repository this would serve is
agent-control-plane, which owns the run state machine and mints a
step_id per unit of work, and already promises that a run which dies mid-step resumes without
repeating a side effect; that promise ends at the process boundary. A side-effecting step whose
effect lands in another system — the sandbox action that calls someone else's API, the step that
charges an account — is exactly the shape this package covers: pass the step's key straight through
as Message.key, and the window becomes a deployment parameter of whoever consumes it. It is not
the control plane's to choose, and it is the reason this README states the boundary rather than
burying it.
Nothing about that coupling is code. agent-control-plane would use this package the way any
application does — Producer.produce inside the transaction that checkpoints the step — and this
repository has no idea that it exists. Spec §6's four other rows are the same shape: documented,
never imported.
The four the platform spec names — partition, crash, duplicate, clock skew — each answered from a test rather than from intuition:
| Fault | What happens | Where that comes from |
|---|---|---|
| Partition, relay cannot reach consumer | every send while it is in force is lost, including a message already in flight when it forms. The outbox row stays unmarked, the relay retries after retry_ms, and delivery resumes on heal. Nothing is applied twice, and nothing is lost |
the 2,000-seed simulation partitions relay from consumer on every seed, 400→1300 ms; the harness's own test_a_message_in_flight_when_the_partition_forms_is_lost and test_a_partitioned_pair_cannot_talk_until_healed |
| Crash, any role, mid-write | a kill inside any of the five in-transaction seams leaves nothing on disk; a kill between deliver and mark redelivers; a kill between commit and ack redelivers into an inbox that drops it | test_a_kill_inside_any_transaction_leaves_no_trace (parametrised over all five), test_the_hard_case_consumer_dies_after_the_effect_and_before_the_inbox_row, test_consumer_dies_after_commit_before_ack_and_the_redelivery_is_dropped, and 474 real kills across 200 seeded schedules |
| Duplicate inside the window | dropped: Outcome.DUPLICATE, the handler is never called, the ledger does not move |
test_a_duplicate_inside_the_window_is_dropped, test_many_duplicates_inside_the_window_are_all_dropped; 5,578 redeliveries dropped across the 2,000-seed simulation |
| Duplicate after the window | applied again. The boundary, not a bug — see "The honest boundary" above | test_the_boundary_is_exact pins it to the millisecond; the 500-seed simulation reaches it on purpose |
Clock skew on the consumer's now |
backwards or slow is safe: rows stay longer and more duplicates are dropped. A jump forward by d prunes the inbox d early, and duplicates arriving in that gap are applied again — unsafe, and undetectable from inside. Relay/consumer skew is irrelevant: the window is measured entirely on the consumer's clock, and produced_at is never compared with applied_at |
arithmetic on test_the_boundary_is_exact, which fixes the row as prunable at exactly applied_at + window_ms. The skew itself is not simulated; docs/failure-model.md §4 says so rather than implying a test exists |
Full tables, each row naming the test that produced it, in
docs/failure-model.md — the process, the transport, the handler, the
caller's clock, the disk, and the key discipline. The three worth knowing before you read any of it:
- A transport that lies — returns "accepted" and then loses the message — loses the effect. The outbox says sent, the consumer never hears, and nothing in this package can detect it. At-least-once rests entirely on the transport keeping that promise. Stated, not tested, because it is not reproducible here.
- A handler that always raises is retried on every redelivery, forever. There is no dead-letter queue.
- A consumer clock that jumps forward prunes the inbox early, and duplicates arriving in that gap are applied again. Backwards is safe; forwards is not, and it is undetectable from inside.
Adversarial angles — forged keys, replay after the window, colliding key schemes, SQL injection
through the checker's table names — are in docs/threat-model.md, with the
one this package genuinely cannot defend (a forged message bearing a key the producer never issued)
named as such.
src/effect_queue/
db.py one SQLite connection: WAL, synchronous=FULL, BEGIN IMMEDIATE, nesting joins
messages.py key + kind + payload, and the digest that makes a key a promise
outbox.py the producer: business rows and the intent in one transaction
relay.py deliver, then mark sent; volatile backoff
inbox.py the consumer: the effect and the inbox row in one transaction; the window
broker.py a durable queue that does not deduplicate
ledger.py the reference effect the proof counts, with the pid that wrote it
pipeline.py the four wired together over three files
verify.py the checker: missing / duplicated / orphaned / torn
faults.py ten named seams; HaltAt parks for a real kill, RaiseAt only raises
simulation.py the protocol as nodes on the harness
sim/ clock, network, world, trace — the harness itself
tests/ 146 tests; crash_child.py is the process that gets killed,
race_child.py is the one that gets company
docs/ ADRs, failure model, threat model
examples/ order_shipping.py
Eleven ADRs, each with Context / Decision / Alternatives / Consequences and the test that holds it —
docs/adr/. Each records an alternative that was genuinely available and
genuinely rejected: two-phase commit, a Bloom filter, per-producer sequence numbers, a
deduplicating broker, change-data capture, an injected Clock object, Postgres with
LISTEN/NOTIFY, a broker visibility timeout, a claim column on the outbox, a shared simulation
package the five layer-3 repositories would all depend on. The three to read
first are
ADR-0003 (why the handler is
handed the connection rather than called and trusted),
ADR-0004 (why the window is bounded and what that
costs) and
ADR-0009 (three
databases — and the deployment that should not use this package at all).
What was not built, from what is actually absent rather than from a template:
- No dead-letter queue. A handler that always raises is retried forever. Nothing counts attempts and nothing parks a poisoned message.
- No real transport.
SqliteBrokeris a local queue over a file. There is no client for SQS, Kafka, NATS or anything else, and the one contract this package cannot check —deliverreturns only when the message is durable — is the transport's to keep. - More than one consumer is correct but does not scale. This bullet used to say a single
consumer per inbox was the only supported shape. Writing the test showed that was wrong —
three racing processes apply each effect exactly once
(
test_racing_full_pipelines_apply_each_effect_exactly_once) — so the limitation is narrower than it was written: every write serialises on one file's write lock, so three consumers do the work of one plus contention, and a handler that holds the lock makes the others wait up toDatabase(timeout=…)and then raise. There is still no partitioning, no lease and no visibility timeout, and throughput under concurrency is not measured — no figure for it appears anywhere here. - No ordering guarantee. Each key is independent. If a key must be applied after another, this package does nothing for you.
- Pruning is manual.
Consumer.prune(now)is a call, not a background thread — this package never starts a thread and never reads a clock. Scheduling it is the deployment's job, and an inbox nobody prunes grows forever. synchronous=FULLcosts throughput and there is no benchmark here saying how much. Measuring it would need a real workload on real hardware, and inventing the number is exactly what this repository refuses to do.- The crash tests are the slow part of the suite — 59.82 s and 62.04 s for the file on its own,
most of a 146-test suite that ran between 82.43 s and 96.22 s over five runs — and they need an
interpreter whose
Popen.pidis really the interpreter (conftest.pyprobes for one, because a Windows virtualenv shim is not). - It has never run on more than one machine. Every process here talks to another through a SQLite file on one filesystem. The pattern is a distributed one and the tests are about what survives a crash, not about what survives a network, which is what the simulated network is for and why its figures are labelled as simulation.
- Only Windows was exercised.
Popen.kill()isTerminateProcesshere andSIGKILLon POSIX; theconftest.pyprobe exists because of a Windows-specific hazard, and the POSIX path has not been run. The CI workflow namesubuntu-latest, and CI has never run — see below. - The simulation models a crash as an exception (
SimCrash) unwinding out of a transaction. That is a model; the real-process tests are what check the model is honest, and they are the reason the simulation is not the only proof.
python -m pip install -e ".[dev]" ruff mypy
python -m ruff check src tests examples
python -m ruff format --check src tests examples
python -m mypy
python -m pytest -qLocal interpreter: CPython 3.14.6, pytest 9.1.1, ruff 0.16.5, mypy 2.3.1. requires-python is
>=3.12 and the CI matrix names 3.12, 3.13 and 3.14; only 3.14 was executed here, and saying so
is cheaper than implying three. CI has never run: there is no remote, so no workflow has ever
been triggered and there is no badge to show.
.github/workflows/ci.yml has nine steps after the install. Every one
was executed on this machine, and every one that can meaningfully be broken was then watched to
go red against a deliberately damaged copy of the tree, because a gate nobody has seen fail is a
gate nobody knows is wired up. The damage was applied to a throwaway copy outside the repository and
is not in the history.
| CI step | Observed here, passing | Watched to fail by | Failing exit |
|---|---|---|---|
| Lint | All checks passed! |
appending import json to relay.py — F401 and E402, 2 errors |
1 |
| Format | 37 files already formatted |
appending x = 1 to relay.py — 1 file would be reformatted |
1 |
| Types | Success: no issues found in 18 source files |
making Relay.mark_sent return the cursor — Incompatible return value type (got "Cursor", expected "bool") |
1 |
| Dependency audit | 8 tests in 0.64 s | adding import time to db.py — fails test_no_module_under_src_reaches_the_network_or_reads_a_clock as db.py imports time |
1 |
| The simulation harness does what it says | 22 tests in 0.47 s | making Network.send ignore drop — fails test_a_dropped_message_is_not_delivered |
1 |
| Tests | 146 passed |
— (covered by the rows around it) | — |
| The proof, under attack | 63 tests in 85.19 s | splitting the consumer's one transaction in two, so the effect commits before the inbox row | 1 |
| The worked example still runs | exit 0 | deleting the broker's second delivery of ord-1002 — the example checks its own premise and prints the redelivery never happened: copies=1 duplicates=0 |
1 |
| Documentation numbers | 13 tests in 4.72 s | itself, three times: adding the concurrency tests left Tests collected at 140 against 141 collected and Test lines at 3,095 against 3,104, and the gate failed with both numbers side by side before either was corrected; adding the thirteenth check then failed it a third time, naming the five other sentences that still said 145 |
1 |
The row that carries the repository is the seventh. Closing Consumer.receive's transaction after
the handler and opening a second one for the inbox row — one inserted line and an indent — looks like
a tidy-up, passes every test that does not crash, and turns the proof red in two independent ways:
tests/test_crash_exactly_once.pyandtests/test_inbox_window.py: 7 failed, 26 passed, the failures includingtest_the_hard_case_consumer_dies_after_the_effect_and_before_the_inbox_rowand bothtest_a_kill_inside_any_transaction_leaves_no_trace[consumer.after_effect]and[consumer.after_inbox]. Inside that, 54 of the 200 real-process crash schedules violated exactly-once, each naming its seed and the seams its kills landed on. The first of them, verbatim:seed=4 intents=3 kills=[attempt 0: producer.after_outbox (point #2), attempt 1: relay.after_deliver (point #12), attempt 2: consumer.after_inbox (point #7)] ... intents=3 effects=4, withduplicated (effect count > 1): [('s4-intent-0', 2)]. That seed is the one the "Reproduce a failure" section above replays against the honest code, where it iseffects=3;tests/test_sim_relay_path.py: 2 failed, 6 passed, with 1,054 of the 2,000 simulated schedules violating exactly-once. Two independent harnesses, the same verdict, from one break.
The eighth row is in the table because it was false before it was true. Deleting the redelivery from
the example used to leave it printing queued copies=1, duplicates_dropped=0 and still exiting
0: three orders and three shipments is exactly-once whether or not a duplicate ever had to be
dropped, so the checker was satisfied and the step named after the example was not the step catching
the damage — tests/test_examples.py, under the Tests step, was. The example now checks the premise
it exists to demonstrate and exits 1 when the redelivery did not happen, which is what the row above
records.
The last row is the one that keeps this file honest. tests/test_docs.py holds thirteen checks:
three re-run the counting commands, so the test count, the source line count and the test line count
cannot drift or be invented; two read the seed and race counts out of the constants the tests loop
over; one fails if the README stops saying, in so many words, that a duplicate outside the window is
applied again; one fails if a line carrying a simulation figure does not label it as one; one bans
the four marketing adjectives that stand in for evidence from every markdown file here; and the rest
hold the documents' own shape — every relative link resolves, every ADR has its four sections and is
indexed, the threat model argues every threat it lists, and every test named anywhere in the
documentation or in a src/ docstring actually exists. That last one earned its place twice in one
sitting: first by finding inbox.py pointing at a test file whose name had lost its middle word —
the file has been tests/test_inbox_window.py since the day it was written — and then by refusing
this very paragraph, because naming the wrong filename in prose is indistinguishable, to a grep,
from naming it in a table. The banned-adjective check has the same kind of story: it caught this
section's ancestor on its first run, when the sentence quoted two of the four in order to name them.
The thirteenth is the newest, and it exists because the first one was not enough.
test_the_readme_test_count_is_the_real_one checks the table row, which is one of six places this
README states the size of the whole suite; the other five — the "at a glance" line at the top, a
timing row, a layout listing, a limitation, and the CI row above — are prose nobody counted. When
the concurrency tests landed, the table row was caught at 140 against 141 collected and corrected;
the four of those five that existed at the time went on saying 140, and no test minded.
test_every_place_the_readme_states_the_suite_size_states_the_same_one now resolves all six against
one collection, and each of its patterns must still match something, so deleting or rewording a
sentence fails the gate rather than quietly removing it. A number that appears in six places has six
chances to be wrong, and the point of a documentation gate is that none of them is left to
proofreading.
A note on order, because the usual claim would be flattering and untrue here. The git history
shows the simulation harness and its 22 tests landing before any protocol code, which is the part
that was genuinely test-first; the crash proof was written afterwards, against a working
implementation. What makes it adversarial anyway is the row above: it was watched to go red against
three broken components shipped in tests/broken.py and against a fourth break applied to the honest
Consumer itself. A test that has never been seen to fail is a claim, not evidence, and every test
in that row has now been seen to fail.
MIT — see LICENSE.