Sitelet https://github.com/jamarius-fortson/effect-queue
Skip to content

About

Transactional outbox and idempotent inbox over SQLite. Exactly-once is stated honestly as at-least-once delivery plus idempotent application, and proven by killing real processes. Python, zero runtime dependencies.

Topics

Resources

Stars

0 stars

Watchers

0 watching

Forks

Latest commit

 

History

29 Commits

Folders and files

Repository files navigation

effect-queue

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.

Structure

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
Loading

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 honest boundary

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:

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.

The proof: real processes, killed

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
Loading

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 one

That 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.

The checker is proven to check

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.

Live processes, racing

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.

The simulation

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.

Numbers

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.

Simulation results

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.

Reproduce a failure

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_window

Run it

Python 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 harness

The 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

Use it

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 window

The 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.

When you do not need this

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.

How it fits the platform

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.

Failure modes

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.

Layout

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

Decisions

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).

Limitations

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. SqliteBroker is 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 — deliver returns 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 to Database(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=FULL costs 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.pid is really the interpreter (conftest.py probes 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() is TerminateProcess here and SIGKILL on POSIX; the conftest.py probe exists because of a Windows-specific hazard, and the POSIX path has not been run. The CI workflow names ubuntu-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.

Development

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 -q

Local 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.

Every gate in CI, run locally — and shown to fail

.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.py and tests/test_inbox_window.py: 7 failed, 26 passed, the failures including test_the_hard_case_consumer_dies_after_the_effect_and_before_the_inbox_row and both test_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, with duplicated (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 is effects=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.

Licence

MIT — see LICENSE.

About

Transactional outbox and idempotent inbox over SQLite. Exactly-once is stated honestly as at-least-once delivery plus idempotent application, and proven by killing real processes. Python, zero runtime dependencies.

Topics

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages