Chaos Testing Worker Crashes and Redelivery
Every job system claims at-least-once delivery; very few teams have seen what their own jobs do when a worker dies halfway through one. This guide builds crash tests that kill workers at precise points and assert on the side effects left behind, as part of Testing Background Jobs in Backend Frameworks & Worker Scaling. It starts with deterministic crash points in CI and ends with randomised crash drills in staging.
Problem Statement
A payouts service moves money to sellers with a Celery task: it calls the payment provider, records the transfer, and marks the payout paid. The code looks idempotent — it checks payout.status before transferring. During a Kubernetes node drain last quarter, 14 sellers were paid twice. The worker was killed after the provider call returned but before the database update committed; the broker redelivered the task, the status still read pending, and the transfer ran again. Code review had not caught it and no test could have, because no test ever killed a worker. You want tests that reproduce this crash deterministically, fail on the current code, pass after the fix, and keep passing as the code evolves.
Prerequisites
- An integration test setup that runs a real worker process against a real broker (see testing Celery tasks with pytest or the BullMQ equivalent).
- Late acknowledgement configured the way production runs (
task_acks_late=Truefor Celery; BullMQ and Sidekiq acknowledge after completion by default). - A fake of each external side effect that records calls durably (a table or file), not in memory — the process that records them is the one you are about to kill.
- Permission to run disruptive tests in a staging environment for the drill in Step 6.
Step 1 — List the Crash Points That Matter
A handler with side effects has a small number of boundaries where a crash changes the outcome. Write them down; each becomes a test.
payout task:
(a) after receive, before provider call -> redelivered, provider called once: fine
(b) after provider call, before DB commit -> redelivered, provider called AGAIN <- the bug
(c) after DB commit, before ack -> redelivered, status=paid, skipped: fine
(d) during provider call (connection cut) -> unknown outcome at the provider <- needs a key
Points (b) and (d) are where "check the status first" does not help, because the status write has not happened yet. Only an idempotency key the provider honours, or recording intent before calling out, closes them — the approach in preventing duplicate job execution with idempotency.
Step 2 — Add Fault Hooks at the Crash Points
To crash deterministically at point (b), the worker must die at exactly that line. A tiny fault-injection hook, active only when an environment variable names the point, makes that possible without sleeping and hoping.
# faults.py — no-op in production; kills the process at a named point in tests
import os, signal
_ARMED = os.environ.get("CRASH_AT") # e.g. "payout.after_provider"
_MARKER = os.environ.get("CRASH_MARKER", "/tmp/crash-fired")
def crash_point(name: str) -> None:
if _ARMED != name or os.path.exists(_MARKER):
return # fire once only, so the retry can succeed
open(_MARKER, "w").close()
os.kill(os.getpid(), signal.SIGKILL) # no cleanup, no ack: like an OOM kill
# tasks.py
@app.task(bind=True, acks_late=True, reject_on_worker_lost=True)
def pay_out(self, payout_id: int) -> None:
payout = Payout.objects.get(pk=payout_id)
if payout.status == "paid":
return
transfer = provider.transfer(payout.seller_account, payout.amount_cents,
idempotency_key=f"payout-{payout_id}") # the fix
crash_point("payout.after_provider")
Payout.objects.filter(pk=payout_id).update(status="paid", transfer_id=transfer.id)
SIGKILL matters: it gives the process no chance to run finally blocks, send a nack, or flush anything — the same as the OOM killer or a node losing power. reject_on_worker_lost=True makes Celery requeue a task whose worker died, instead of marking it failed. The marker file makes the fault fire once, so the redelivered attempt can complete.
Step 3 — Run the Worker as a Separate Process and Kill It
An embedded worker thread cannot be killed without killing the test. Start the worker as a subprocess with the fault armed, let it die, then start a second, unarmed worker that processes the redelivery.
# test_crash.py
import os, subprocess, time, pytest
def start_worker(env_extra=None):
env = {**os.environ, **(env_extra or {})}
return subprocess.Popen(["celery", "-A", "app", "worker", "-Q", "payouts",
"--concurrency=1", "--without-heartbeat"], env=env)
@pytest.mark.integration
def test_crash_after_provider_call_does_not_pay_twice(db, provider_fake, tmp_path):
payout = Payout.objects.create(seller_account="acct-9", amount_cents=5000, status="pending")
marker = tmp_path / "fired"
w1 = start_worker({"CRASH_AT": "payout.after_provider", "CRASH_MARKER": str(marker)})
pay_out.delay(payout.id)
assert w1.wait(timeout=30) == -9 # died from SIGKILL at the crash point
assert marker.exists()
w2 = start_worker() # healthy worker picks up the redelivery
try:
wait_until(lambda: Payout.objects.get(pk=payout.id).status == "paid", timeout=60)
finally:
w2.terminate(); w2.wait(timeout=30)
assert provider_fake.transfers_for("acct-9") == 1 # the real assertion
Against the original code (no idempotency key), provider_fake records two transfers and the test fails — which is the point. Redelivery timing depends on the broker: RabbitMQ requeues as soon as the channel closes; Redis-backed Celery waits for visibility_timeout, so set it to a few seconds in the test configuration.
Step 4 — Cover Point (d): the Unknown Outcome
At point (d) the connection drops mid-call, so the worker does not know whether the provider performed the transfer. Simulate it in the provider fake: perform the transfer, then raise a timeout to the caller.
class ProviderFake:
def __init__(self, store):
self.store = store # durable: a table the test can read
def transfer(self, account, amount, idempotency_key):
existing = self.store.get(idempotency_key)
if existing:
return existing # provider-side dedup, like real APIs
result = self.store.put(idempotency_key, account, amount)
if os.environ.get("TIMEOUT_AFTER_TRANSFER") and not self.store.flag("timed_out"):
self.store.set_flag("timed_out")
raise ProviderTimeout("read timed out") # transfer happened; caller can't tell
return result
With the key, the retry returns the existing transfer and the test sees one payout. Without the key, it sees two. The same fake lets you test the opposite design mistake — classifying ProviderTimeout as permanent and giving up — which leaves a transfer made but the payout marked failed. The classification rules are in retrying only transient errors by exception type.
Step 5 — Test Broker Restarts and Lost Acks
Workers are not the only thing that crashes. Restart the broker container between the side effect and the ack, and check that the job is neither lost nor duplicated beyond what idempotency absorbs.
@pytest.mark.integration
def test_broker_restart_mid_job(redis_container, db, provider_fake):
payout = Payout.objects.create(seller_account="acct-3", amount_cents=900, status="pending")
w = start_worker({"SLOW_AFTER_PROVIDER": "3"}) # sleep 3 s after the transfer
pay_out.delay(payout.id)
wait_until(lambda: provider_fake.transfers_for("acct-3") == 1, timeout=30)
redis_container.get_wrapped_container().restart() # broker gone for a moment
wait_until(lambda: Payout.objects.get(pk=payout.id).status == "paid", timeout=90)
w.terminate(); w.wait()
assert provider_fake.transfers_for("acct-3") == 1
If Redis runs without persistence in the test (the default for the container), a restart loses unacked messages entirely — a useful demonstration of why production brokers need persistence configured, as covered in Redis persistence: AOF vs RDB for queues. Run the test with --appendonly yes on the container to match production.
Step 6 — Run Randomised Crash Drills in Staging
Deterministic tests cover the points you listed. A drill covers the ones you did not: under realistic load, kill random worker pods repeatedly and reconcile outcomes afterwards.
# drill.sh — 20 minutes of random worker kills under synthetic load
python loadgen.py --queue payouts --rate 20 --duration 1200 &
end=$((SECONDS + 1200))
while [ $SECONDS -lt $end ]; do
pod=$(kubectl -n staging get pods -l app=payout-worker -o name | shuf -n 1)
kubectl -n staging delete "$pod" --grace-period=0 --force # SIGKILL equivalent
sleep $((RANDOM % 60 + 15))
done
wait
python reconcile.py --since "20 minutes ago" # payouts vs provider transfers, one-to-one
reconcile.py compares every payout created during the drill with the provider sandbox's transfer list: each payout should have exactly one transfer and a paid status. Any mismatch is a bug the deterministic tests missed — add a crash point for it. Tools like Chaos Mesh or Litmus can schedule the pod kills declaratively instead of a shell loop.
Verification
The suite is doing its job when it fails on the pre-fix code and passes after:
git stash # remove the idempotency key fix
pytest -m integration -k crash # expect: test_crash_after_provider_call... FAILED
git stash pop
pytest -m integration -k crash # expect: all passed
In production, reconcile daily: payouts marked paid without a transfer, and transfers without a paid payout, should both be zero. A daily reconciliation is the chaos test that runs forever.
Gotchas & Edge Cases
Fakes that live in the killed process. An in-memory fake dies with the worker, so the test cannot count calls. Record to a database table, file, or separate service.
acks_early hides the bug by losing jobs. With early acknowledgement (Celery's default), a killed worker's task is simply lost, and a crash test "passes" with one transfer — for a payout that never gets marked paid. Assert on the final state, not only on the absence of duplicates.
Prefetched messages. A killed worker also releases messages it had prefetched but not started. Run crash tests with the production prefetch settings so those are exercised too.
Leaving faults armed. Guard fault hooks so they cannot activate in production: read the variable only when a test-only setting is present, and alert if it is ever set outside test environments.
FAQ
Is this worth it for every job? No — for jobs whose side effects are expensive, external, or irreversible: payments, emails and SMS, writes to third-party systems. Pure database jobs with idempotent upserts rarely need crash tests.
Why SIGKILL rather than SIGTERM? SIGTERM triggers graceful shutdown, which is a different (and important) test — see graceful shutdown for Go workers. SIGKILL reproduces the crashes you cannot handle: OOM kills, node failures, and forced pod deletion.
Can I run the drill in production? Only with mature reconciliation and alerting, and ideally with the blast radius limited to a canary pool. Most teams get the same value from staging drills plus daily production reconciliation.
Related
- Testing Background Jobs — the layered strategy that ends here.
- Celery acks_late and Worker Crash Safety — the settings these tests exercise.
- Debugging Duplicate Deliveries in SQS — the same failure diagnosed in production.
- Integration Testing BullMQ Workers with Testcontainers — the harness to adapt for Node.