From 8da918b2088b8205dc445c57adfb0e0357c9464d Mon Sep 17 00:00:00 2001 From: Bretton Date: Sat, 15 Aug 2026 17:25:40 -0700 Subject: [PATCH] fix(outbound): a park hands back its claim's attempt, so holds spend no retries MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit ClaimNext charges an attempt to every claim — right for a delivery that was TRIED, wrong for one that was HELD. park/parkCausal settled through Release, which never gave the increment back, so a kill switch held for ~40s exhausted an 8-attempt budget and the first retryable failure after the switch cleared poisoned instead of retrying (FOLLOWUPS "DEFECT, NOT DONE", 2026-08-14). ReleaseParked is Release's fenced sibling built from one shared template (the fence is spelled once); its single difference is attempts = GREATEST(attempts - 1, 0). The claim-token fence is what makes the give-back safe: only the claim that added an increment can subtract one, a stale park matches zero rows, and a terminal row is untouched. Real failures (releaseOrPoison) and held settlements (settleLater) keep Release — those increments are earned. The causal wait's outcome remains wall-clock-decided; a regression guard pins both facts together. Tests: outer acceptance (switch held past the cap, then one 503 → retry, not poison), store fencing/net-zero suite, all three park sites attempt-neutral. DEPLOY.md §5, README, and FOLLOWUPS rewritten to the fixed behavior; the residual write-load cost note stands. Co-Authored-By: Claude Fable 5 --- DEPLOY.md | 110 +++++----- FOLLOWUPS.md | 95 +++----- README.md | 2 +- internal/outbound/park_budget_test.go | 96 +++++++++ internal/outbound/park_sites_test.go | 140 ++++++++++++ internal/outbound/worker.go | 39 ++-- internal/store/interfaces.go | 12 ++ internal/store/outbound_deliveries.go | 64 +++++- internal/store/outbound_delivery_park_test.go | 204 ++++++++++++++++++ 9 files changed, 613 insertions(+), 149 deletions(-) create mode 100644 internal/outbound/park_budget_test.go create mode 100644 internal/outbound/park_sites_test.go create mode 100644 internal/store/outbound_delivery_park_test.go diff --git a/DEPLOY.md b/DEPLOY.md index 5d3e97d..5e8886b 100644 --- a/DEPLOY.md +++ b/DEPLOY.md @@ -124,38 +124,38 @@ transaction, so intent past the gate is never dropped (`cmd/tidepool/main.go:695-713`). Raising `OUTBOUND_WORKERS` later drains whatever accumulated, immediately. Budget for that. -**`OUTBOUND_DISABLED` parks, it does not fail — but parking is not free.** A -blocked delivery stays `pending` and resumes when the switch clears; `park` +**`OUTBOUND_DISABLED` parks, it does not fail — and a park costs writes, not +retries.** A blocked delivery stays `pending` and resumes when the switch +clears; `park` itself never poisons and never cancels (`internal/outbound/worker.go:242-244`, `internal/outbound/switches.go:5-12`). `OUTBOUND_DRY_RUN` parks the same way, before any signing or POST (`worker.go:245-249`). -⚠️ **A park still consumes the delivery's retry budget, and fast.** `ClaimNext` -does `attempts = attempts + 1` on every claim -(`internal/store/outbound_deliveries.go:143`), and `park` settles through -`Release`, whose `SET` clause updates `claimed_until`, `next_attempt_at`, -`last_error_class`, `response_excerpt`, `last_status_code` and `updated_at` — -and **never resets `attempts`** (`outbound_deliveries.go:218-221`). A parked -delivery is rescheduled `parkDelay = 5 * time.Second` out (`worker.go:560`), so -with workers running, each ordering key's head delivery is re-claimed and -re-parked roughly every five seconds and its -`DefaultMaxDeliveryAttempts = 8` budget (`worker.go:78`) is exhausted in about -**40 seconds**. Nothing poisons *while* parked — `park` never calls `poison` — -but once the switch clears, the **first** retryable failure finds -`Attempts >= maxAttempts` and poisons immediately instead of retrying -(`worker.go:512-518`). An hour-long global kill switch leaves every head -delivery with hundreds of attempts and zero retries left. - -It is also a continuous write load: one `UPDATE` per parked head roughly every -five seconds, per ordering key, for as long as the switch is engaged. - -**The remedy is `POST /admin/outbound/redrive`.** `RedrivePoisoned` sets -`attempts = 0` along with `state = 'pending'` -(`internal/store/outbound_deliveries.go:613-617`), so a redrive restores a full -budget. It only matches `state = 'poisoned'` (`:619`), so it repairs the damage -after a delivery has already fallen over, rather than preventing it. **To park -everything, prefer `OUTBOUND_WORKERS=0`** — see [Rollback](#rollback). The -underlying defect is tracked in `FOLLOWUPS.md`. +**A park costs the delivery no retries.** `ClaimNext` does +`attempts = attempts + 1` on every claim +(`internal/store/outbound_deliveries.go:143`), but `park` and `parkCausal` +settle through `ReleaseParked`, not `Release` (`worker.go:572`, `:592`), whose +one difference is `attempts = GREATEST(attempts - 1, 0)` +(`outbound_deliveries.go:243-244`) applied under the same claim fence. A full +claim/park cycle therefore leaves the ledger where it found it: an hour-long +global kill switch spends nothing, and once it clears the **first** retryable +failure is a retry, not a poison. `DefaultMaxDeliveryAttempts = 8` +(`worker.go:78`) counts attempts a delivery actually made. + +⚠️ **It is a continuous write load, though.** A parked head is re-claimed and +re-parked every `parkDelay = 5 * time.Second` (`worker.go:560`), so a held +switch costs one claim plus one `UPDATE` per parked ordering-key head per ~5s, +for as long as it is engaged — indefinitely, since nothing now ends the cycle on +its own. That churn is the remaining argument for **preferring +`OUTBOUND_WORKERS=0` to park everything** — see [Rollback](#rollback) — and it +is a cost argument, not a safety one. + +`POST /admin/outbound/redrive` is the repair for **poisoned** deliveries, not +for parked ones: `RedrivePoisoned` sets `attempts = 0` with `state = 'pending'` +(`internal/store/outbound_deliveries.go:665-666`) and matches +`state = 'poisoned'` only (`:667`). A parked delivery is still `pending`, so it +is neither reached by a redrive nor in need of one — it resumes on its own when +the switch clears. Scope matching, from `config.go:519-525`: @@ -175,7 +175,7 @@ constructed inside `if cfg.OutboundWorkers > 0` switch and nothing consulting one. Setting `OUTBOUND_DISABLED=true` while workers are 0 changes nothing and confirms nothing; it is not a second belt. Conversely, that is also why `OUTBOUND_WORKERS=0` costs nothing: no worker -exists to claim, park, or spend attempts. +exists to claim, park, or write anything. ### Confirming a switch actually engaged @@ -193,8 +193,8 @@ exists to claim, park, or spend attempts. switch engaged: `by_state` still reads `pending` for a parked delivery, exactly as it does for one merely waiting its turn, so the queue view cannot tell you. Read it as a rate, not a level — a parked head is re-parked every ~5s, so the -counter climbs continuously while a switch is held. That climb *is* the budget -being spent. +counter climbs continuously while a switch is held. That climb is the write load +of the hold, not retries being spent. Two caveats on it. It is shared with `parkCausal`, so a nonzero `parked` with no switch engaged means causal waits, not an operator action. And @@ -542,10 +542,10 @@ the moment `CONSUMER_ENABLED=true`, the translator has already run on everything dry run or not. Step 1 is where you check its output (read `outbound_activities.payload` in psql), not here. -And because dry-run parks, it carries the same budget cost as any other park: -each head delivery is re-claimed every ~5s and burns its 8 attempts in ~40 -seconds (see [§2](#2-the-v2-flag-topology)). A "one cycle" dry run is measured -in seconds, not hours, and wants a `redrive` after it. +And because dry-run parks, it carries the same cost as any other park — which +is write load, not retries: each head delivery is re-claimed and re-parked every +~5s for as long as it is engaged, and the attempt it spends doing so is handed +back (see [§2](#2-the-v2-flag-topology)). Nothing needs redriving afterwards. ### On announcement throttling — what actually exists @@ -614,11 +614,11 @@ What each one means: ORDER BY updated_at DESC LIMIT 50;" ``` - Check `attempts` in that output before concluding the peer rejected anything: - a delivery whose budget was spent by a held kill switch (see - [§2](#2-the-v2-flag-topology)) poisons on its first real failure with an - `attempts` far above 8 and an error class that describes one attempt, not - eight. + `attempts` in that output counts claims that were TRIED — a park hands its own + claim's increment back (see [§2](#2-the-v2-flag-topology)) — so a poisoned row + should read about 8. A materially higher number is not extra peer rejections: + it is claims that never settled, since a lapsed lease or a crashed worker + leaves its increment behind with nobody to return it. - **echo drop counters rising steadily** — expected and healthy: our own content arriving back from Lemmy and being correctly refused. A counter at **zero** while native content is flowing is the alarming case; it means @@ -640,10 +640,10 @@ Rollback is a **kill switch, not an un-deploy**. In escalation order, each step being one `.env` edit plus `up -d tidepool`: 1. **`OUTBOUND_DISABLED_COMMUNITIES=`** — park one community. - Everything else keeps flowing; the parked deliveries resume when you clear - it. First because it is the *narrowest*, not because it is free: it is a - park, so it spends that community's head delivery's retry budget at the same - ~5s cadence as (3). Redrive that community after clearing it. + Everything else keeps flowing; the parked deliveries resume with their retry + budget intact when you clear it, and want no redrive. First because it is the + *narrowest*. Its only cost is the same ~5s re-claim/re-park write cycle as + (3), on that community's head delivery alone. 2. **`OUTBOUND_WORKERS=0`** — stop the workers entirely. This is the big red button for delivery. The consumer keeps running and keeps recording intent; nothing reaches any peer. `NewWorker` is only called inside @@ -651,17 +651,17 @@ step being one `.env` edit plus `up -d tidepool`: no worker to claim anything: nothing is re-claimed, no `attempts` are spent, and no rows are written. It costs nothing and it is genuinely lossless. 3. **`OUTBOUND_DISABLED=true`** — park everything outbound. Same *observable* - effect as (2) — nothing reaches any peer — but **it is not free, and it is - ranked below (2) for that reason.** The workers keep running, so every - ordering key's head delivery is re-claimed and re-parked every ~5s and burns - its 8-attempt budget in ~40 seconds (see [§2](#2-the-v2-flag-topology)). The - deliveries this switch exists to protect are exactly the ones left with no - retries. Reach for it only when you need the *scope* it gives you and a - whole-worker stop is too blunt — and expect to `redrive` afterwards, which - is the only thing that resets `attempts` - (`internal/store/outbound_deliveries.go:613-617`). The one thing (3) buys - over (2) is that each parked row carries a recorded `last_error_class` / - `response_excerpt`; a stopped worker records nothing. + effect as (2) — nothing reaches any peer — and the deliveries keep their + retries: a park hands its claim's increment back (see + [§2](#2-the-v2-flag-topology)), so this is safe to hold for as long as you + need. It is **ranked below (2) on cost, not on risk**: the workers keep + running, so every ordering key's head delivery is re-claimed and re-parked + every ~5s — a claim and an `UPDATE` per key per 5s, indefinitely, against a + database you may be trying to leave alone during an incident. Reach for it + when you need the *scope* it gives you and a whole-worker stop is too blunt. + What (3) buys over (2) is that each parked row carries a recorded + `last_error_class` / `response_excerpt` saying why it is held; a stopped + worker records nothing. 4. **`CONSUMER_ENABLED=false`** — stop consuming. Intent stops being recorded. The consumer resumes from its stored cursor when re-enabled, so this is recoverable, but it is the only step that stops *observing*, and diff --git a/FOLLOWUPS.md b/FOLLOWUPS.md index d01be9d..4e2bb6a 100644 --- a/FOLLOWUPS.md +++ b/FOLLOWUPS.md @@ -63,67 +63,40 @@ task documents and git history rather than this list. ## Outbound delivery (task 15) -- **DEFECT, NOT DONE — a PARK spends the poison budget, so the kill switch - destroys the retries of exactly the deliveries it exists to protect.** - Found 2026-08-14 while fact-checking `DEPLOY.md`; the docs and the two park - doc-comments have been corrected to describe it, and **nothing about the - behaviour was changed.** It needs its own test-first subtask. - - *Mechanism.* `ClaimNext` does `attempts = attempts + 1` on every claim - (`internal/store/outbound_deliveries.go:143`). `park` and `parkCausal` settle - through `Release`, whose `SET` clause updates `claimed_until`, - `next_attempt_at`, `last_error_class`, `response_excerpt`, `last_status_code` - and `updated_at` — and never resets `attempts` - (`outbound_deliveries.go:218-221`). `parkDelay = 5 * time.Second` - (`internal/outbound/worker.go:560`) and `DefaultMaxDeliveryAttempts = 8` - (`worker.go:78`), so with `OUTBOUND_DISABLED=true` and workers running, each - ordering key's head delivery is re-claimed and re-parked about every five - seconds and its entire retry budget is gone in roughly **40 seconds**. - - *Why it is not caught by "park never poisons".* It isn't — `park` genuinely - never calls `poison`, which is what made the old comments read as true. The - damage lands later: `releaseOrPoison` poisons on the FIRST retryable failure - once `Attempts >= maxAttempts` (`worker.go:512-518`). So an hour-long kill - switch leaves head deliveries with hundreds of attempts and zero retries, and - the next transient 5xx or dial timeout poisons them immediately. `parkCausal` - shares the mechanism but is materially safer: `causalStatus` bounds the causal - wait by WALL CLOCK from `delivery.CreatedAt`, never by attempt count, so a - held child's outcome is still decided by elapsed time. It is exposed only for - a genuine delivery failure after the parent lands. - - *Remedy that exists today, and its limit.* `RedrivePoisoned` sets - `attempts = 0` (`outbound_deliveries.go:613-617`), but it matches - `state = 'poisoned'` only (`:619`) — it repairs after the fall, and cannot - pre-empt it. `OUTBOUND_WORKERS=0` avoids the whole problem (no worker is - constructed, `cmd/tidepool/main.go:716`) and `DEPLOY.md` §5 now ranks it above - `OUTBOUND_DISABLED` for that reason. - - *Shape of a real fix — this is the design question, not a settled plan.* - `Release` is a single statement serving two callers that mean opposite - things, and it cannot tell them apart: a **retryable failure** (where holding - the incremented `attempts` is exactly right — that is the backoff working) - from a **park by an operator switch or a causal hold** (where it is not; the - delivery never reached the wire and nothing was learned about the peer). Two - candidate directions, with the trade to be argued in the subtask: - - a **park-specific release** — a sibling statement, or a flag on `Release`, - that writes `attempts = attempts - 1` / `attempts = $n` so a park is - attempt-neutral. Cheapest and most local, but adds a second write path - through the most fencing-sensitive statement in the package, and a park - that "un-counts" must not be able to underflow or to un-count a real - attempt on a re-claim race. - - **reset on unpark** — leave the increment and clear `attempts` when a - delivery next passes the switch gate. Keeps `Release` single-purpose, but - the reset then lives on the hot path and has to distinguish "was parked" - from "was retried", which today is only knowable from `last_error_class` - (`switch_parked` / `dry_run` / the causal classes) — i.e. it would make an - error-class string load-bearing for a correctness decision. - - A RED test should pin the operator-visible fact rather than the column: with - the switch engaged, drive N claim/park cycles well past `maxAttempts`, clear - the switch, then fail the delivery once retryably and assert it is - **rescheduled, not poisoned**. Note that any fix must keep `parkCausal`'s - wall-clock poison reachable — a naive "parks never advance anything" change - must not also disarm the causal budget. +- **Parks are attempt-neutral — the design that was chosen, and the two limits + it leaves open.** `ClaimNext` charges an attempt to every claim + (`internal/store/outbound_deliveries.go:143`), which is right for a delivery + that was TRIED and wrong for one that was HELD. The fix is a SIBLING of + `Release`, not a flag on it: `ReleaseParked` (`:276`) runs the same fenced + statement from the same template (`releaseStatement`, `:224`) with one slot + filled, `attempts = GREATEST(attempts - 1, 0)` (`:243-244`), and `park` / + `parkCausal` settle through it (`internal/outbound/worker.go:572`, `:592`). + The fence — `state = 'pending' AND claimed_until = $token` — is what makes the + give-back safe: only the holder of the claim that added an increment can + subtract one, so a park can never un-count an attempt another claim spent. + + *The rejected candidate, recorded because it reads cheaper than it is.* + "Reset on unpark" — leave the increment, clear `attempts` when the delivery + next passes the switch gate — keeps `Release` single-purpose but cannot tell a + park's attempts from real failures that preceded the park, so it erases + genuine failure history instead of handing back a hold; the only signal + available to it is `last_error_class`, which would make an error-class string + load-bearing for a correctness decision; and it repairs late, so `attempts` + reads inflated on the admin surface for the whole duration of the park. + + *Two residual limits, both pre-existing and neither closed by this.* + - **An abandoned claim still leaks its `+1` forever.** A lapsed lease, or a + crash between `ClaimNext` and any settle, leaves an increment with nobody + holding the fence to hand it back. Shared with the real-failure path and + bounded by the lease, so it is cosmetic — but it is why `attempts` is a + count of claims that settled, not of POSTs. + - **A park writes `last_status_code = 0`,** clobbering a real status a PRIOR + failed attempt recorded. The divergence sweep reads that column as its + answered-or-silent discriminator (`COALESCE(d.last_status_code, 0) > 0`, + `internal/store/divergence.go:791`), where 0 means the peer never answered — + so a 502 overwritten by a later park leaves no trace of the peer having + spoken. Unchanged by this fix (a park has no status to write) and worth + revisiting only if the refused/unanswered split has to be trusted per-row. - **outbound_deliveries.ClaimNext lacks a standalone `seq` index.** The loose-scan CTE builds the head set via the `(ordering_key, seq)` partial diff --git a/README.md b/README.md index ab43d70..4ca109d 100644 --- a/README.md +++ b/README.md @@ -299,7 +299,7 @@ Two classes, and the difference matters at boot: | `CONSUMER_ENABLED` | **off** | turns on the Jetstream consumer (task 14): native users' opt-outs, profiles, posts, comments and votes flowing outward. Default off because it writes durable outbound state, and because a deployment that has not been canaried should not start accumulating it — not because the seams behind it are stubbed. They are wired: with it on, the **real** enqueuer persists outbound intent, the acceptance engine admits postv2 and writes community-signed acceptances, and opt-out `deleteRemote` / confirmed account deletions actually purge at peers. It also gates two other things — the `OUTBOUND_WORKERS` AND, and whether `/admin/admissions*` exists at all | | `JETSTREAM_URL` | *(optional)* | the self-hosted Jetstream the consumer subscribes to (`ws://` or `wss://`); **required** when `CONSUMER_ENABLED`, and validated at boot whenever set so a typo fails fast instead of becoming a reconnect loop. May be staged ahead of the flag | | `OUTBOUND_WORKERS` | `0` (**off**) | how many delivery workers run. **There is no `OUTBOUND_ENABLED`:** delivery starts only when this is `>0` *and* `CONSUMER_ENABLED`. With the consumer on and this at `0`, outbound intent still accumulates durably and nothing is POSTed — which is the intended staging step, not a broken state. Raising it drains the accumulated backlog immediately | -| `OUTBOUND_DISABLED` | off | global delivery kill switch. A blocked delivery is **parked** — it stays `pending` and resumes when the switch clears — and park itself never poisons or cancels. **But it is not free:** `ClaimNext` increments `attempts` on every claim and `Release` never resets it, so with workers running each ordering key's head is re-claimed every ~5s and its 8-attempt budget is gone in ~40s; the first real failure after the switch clears then poisons instead of retrying. Only `redrive` resets `attempts`. To stop everything, **`OUTBOUND_WORKERS=0` is the cheaper switch** — no worker is constructed, so nothing is spent. Also note every switch here is **inert** while `OUTBOUND_WORKERS=0`. See `DEPLOY.md` §2 | +| `OUTBOUND_DISABLED` | off | global delivery kill switch. A blocked delivery is **parked** — it stays `pending` and resumes when the switch clears — and park itself never poisons or cancels. A park is also **attempt-neutral**: `ClaimNext` increments `attempts` on every claim, but a park settles through `ReleaseParked`, which hands that increment back under the claim fence, so a held switch costs no retries and needs no `redrive`. **What it does cost is writes:** with workers running each ordering key's head is re-claimed and re-parked every ~5s for as long as the switch is engaged, so **`OUTBOUND_WORKERS=0` is the cheaper switch** for stopping everything — no worker is constructed. Also note every switch here is **inert** while `OUTBOUND_WORKERS=0`. See `DEPLOY.md` §2 | | `OUTBOUND_DISABLED_HOSTS` | *(empty)* | comma-separated inbox **hosts** to park. Lowercased on load and compared case-insensitively — a kill switch must fail closed on case | | `OUTBOUND_DISABLED_COMMUNITIES` | *(empty)* | comma-separated community **AP ids** to park (`https://lemmy.world/c/comicstrips`), matched **exactly and case-sensitively** against the delivery's ordering key. There is no allowlist form: a one-community canary is spelled by disabling every other community | | `OUTBOUND_DISABLED_ACTORS` | *(empty)* | comma-separated actor **DIDs** to park, exact match | diff --git a/internal/outbound/park_budget_test.go b/internal/outbound/park_budget_test.go new file mode 100644 index 0000000..a26b2c0 --- /dev/null +++ b/internal/outbound/park_budget_test.go @@ -0,0 +1,96 @@ +package outbound + +import ( + "context" + "database/sql" + "net/http" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "tidepool/internal/store" +) + +// The park/poison budget. A kill switch is an OPERATOR PAUSE: it exists to keep +// deliveries safe while something is wrong. It must therefore cost the paused +// delivery nothing — the retries it still has when the switch clears must be +// exactly the retries it had when the switch engaged. +// +// This is the outer, operator-visible framing of that invariant: it drives the +// worker through the full claim → park → re-claim cycle against the real store, +// so it observes the budget the way an operator does (the row's state after the +// switch lifts), not the way the worker's internals count attempts. + +// clearParkDelay makes a parked row immediately re-claimable by rewinding ONLY +// its schedule — the house idiom for time control (cf. setAttempts), so a test +// never sleeps out parkDelay. +// +// It deliberately touches nothing but next_attempt_at: resetting `state` too +// (as resetDeliveryToPending does) would revive a poisoned row and hide the very +// failure this test is here to catch. +func clearParkDelay(t *testing.T, conn *sql.DB, activityID string) { + t.Helper() + _, err := conn.ExecContext(context.Background(), + `UPDATE outbound_deliveries SET next_attempt_at = now() WHERE activity_id = $1`, activityID) + require.NoError(t, err) +} + +func TestWorker_KillSwitchDoesNotSpendPoisonBudget(t *testing.T) { + conn := workerTestDB(t) + seedWorkerActor(t, conn, true, false) + id := seedDelivery(t, conn, "Create", "", createPayload("x")) + + switches := &fakeSwitches{allow: false} + // Exactly one retryable failure after the switch lifts: a 503 is the + // canonical "come back later", so the delivery's fate here is decided + // purely by whether it still has budget to come back with. + sender := &fakeSender{respond: func(call int, _ string) error { + if call == 1 { + return httpErr(http.StatusServiceUnavailable, "peer restarting") + } + return nil + }} + // MaxAttempts 3 is the package default; BackoffBase must be REAL time here + // (the helper's 1ms default would put next_attempt_at in the past by the + // time we read it, making the reschedule assertion vacuous). + w := newWorker(t, conn, sender, func(o *WorkerOptions) { + o.Switches = switches + o.MaxAttempts = 3 + o.BackoffBase = 30 * time.Second + }) + + ctx := context.Background() + + // Well past MaxAttempts worth of parks: an operator block held for a minute + // is an ordinary Tuesday, and the worker re-claims a parked row every + // parkDelay for as long as it lasts. + for i := 0; i < 6; i++ { + worked, err := w.DeliverNext(ctx) + require.NoError(t, err) + require.Truef(t, worked, "cycle %d: the parked delivery is still claimable", i) + require.Equalf(t, store.DeliveryStatePending, getDelivery(t, conn, id).State, + "cycle %d: parking never leaves the pending state", i) + clearParkDelay(t, conn, id) + } + require.Zero(t, sender.count(), "nothing is POSTed while the switch is engaged") + + // The operator lifts the block. The queue should resume exactly where it + // paused. + switches.allow = true + + worked, err := w.DeliverNext(ctx) + require.NoError(t, err) + require.True(t, worked) + require.Equal(t, 1, sender.count(), "the delivery is POSTed once the switch clears") + + d := getDelivery(t, conn, id) + assert.Equal(t, store.DeliveryStatePending, d.State, + "a kill switch must not destroy the retries of the deliveries it exists to protect: "+ + "after the switch lifts, one retryable 503 is a RETRY, not a poisoning — "+ + "time spent parked is not an attempt") + assert.True(t, d.NextAttemptAt.After(time.Now()), + "the delivery is rescheduled into the future for its retry — an operator who lifts a "+ + "kill switch gets their queue back, not a dead-letter pile to redrive by hand") +} diff --git a/internal/outbound/park_sites_test.go b/internal/outbound/park_sites_test.go new file mode 100644 index 0000000..8ec9e97 --- /dev/null +++ b/internal/outbound/park_sites_test.go @@ -0,0 +1,140 @@ +package outbound + +import ( + "context" + "database/sql" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "tidepool/internal/store" +) + +// Every place the worker HOLDS a delivery instead of trying it must cost that +// delivery nothing. There are three such sites — the kill switch, dry-run, and +// the causal wait — and they are one behavior, not three: a hold is not an +// attempt, whatever the reason for holding. +// +// Pinning all three together is deliberate. The switch site is the one an +// operator feels (see the outer acceptance test), but a per-site fix would leave +// dry-run silently burning the budget of a queue nobody is watching, and the +// causal wait burning the budget of every reply whose parent is a beat behind. + +func TestWorker_ParkSitesAreAttemptNeutral(t *testing.T) { + for _, tc := range []struct { + name string + class string + parentATURI string + seed func(t *testing.T, conn *sql.DB) + opts func(*WorkerOptions) + }{ + { + name: "kill switch", + class: "switch_parked", + opts: func(o *WorkerOptions) { o.Switches = &fakeSwitches{allow: false} }, + }, + { + name: "dry run", + class: "dry_run", + opts: func(o *WorkerOptions) { o.Switches = &fakeSwitches{allow: true, dryRun: true} }, + }, + { + name: "causal wait", + class: "parent_pending", + // A bridge-origin parent that has not been accepted yet: the reply + // is held until it lands (see causal_gating_test.go). + parentATURI: gParentATURI, + seed: func(t *testing.T, conn *sql.DB) { seedBridgeParent(t, conn, false) }, + }, + } { + t.Run(tc.name, func(t *testing.T) { + conn := workerTestDB(t) + seedWorkerActor(t, conn, true, false) + if tc.seed != nil { + tc.seed(t, conn) + } + id := seedDelivery(t, conn, "Create", tc.parentATURI, createPayload("x")) + require.Zero(t, getDelivery(t, conn, id).Attempts, + "a freshly enqueued delivery has spent nothing yet — this is the value a park must restore") + + sender := &fakeSender{} + w := newWorker(t, conn, sender, tc.opts) + + worked, err := w.DeliverNext(context.Background()) + require.NoError(t, err) + require.True(t, worked, "the delivery was claimed and handled") + require.Zero(t, sender.count(), "a park POSTs nothing — the delivery was HELD, not tried") + + got := getDelivery(t, conn, id) + assert.Equal(t, 0, got.Attempts, + "a full claim/park cycle leaves the attempt ledger where it started: the hold gives back "+ + "the increment its own claim took, so the delivery arrives at the attempt cap only "+ + "through attempts it actually made") + assert.Equal(t, store.DeliveryStatePending, got.State, + "a parked delivery stays pending — holding is not failing") + assert.Equal(t, tc.class, got.LastErrorClass, + "the reason for the hold is recorded, so the give-back is auditable rather than a silent rewind") + }) + } +} + +// The causal wait is the one hold whose OUTCOME is not "wait forever": it +// poisons on a wall-clock deadline. That makes it the site where the two +// mechanisms could quietly be wired to each other — and they must not be. +// +// This guard holds both facts down at once, because each is a plausible way to +// break the other: an attempt-neutral park built by suppressing the release +// bookkeeping would take the causal deadline's accounting with it (the child +// waits forever on a parent that never lands), and a deadline defended by +// leaning on the attempt counter would put the budget back in the park's hands. +// Neither passes this test. Cf. causal_hardening_test.go, which pins the same +// deadline from the other direction (attempts=99 must NOT poison early). +func TestCausalGating_ParksAreFreeButTheWallClockStillPoisons(t *testing.T) { + conn := workerTestDB(t) + ctx := context.Background() + seedWorkerActor(t, conn, true, false) + seedBridgeParent(t, conn, false) // parent never accepted + id := seedDelivery(t, conn, "Create", gParentATURI, createPayload("x")) + + sender := &fakeSender{} + w := newWorker(t, conn, sender, func(o *WorkerOptions) { o.CausalWaitBudget = time.Hour }) + + // The spin a child does while its parent is in flight. parkCausal stamps + // next_attempt_at from the APP clock while ClaimNext compares against the + // DB's, so on a host running a few hundred microseconds ahead of Postgres a + // back-to-back re-claim can miss by a hair; the nudge re-stamps the schedule + // from the DB's own now() so each cycle here is deterministic. It is the + // schedule column only — a poison would still be plainly visible. + for i := 0; i < 5; i++ { + worked, err := w.DeliverNext(ctx) + require.NoError(t, err) + require.Truef(t, worked, "cycle %d: the held child stays claimable", i) + clearParkDelay(t, conn, id) + } + require.Zero(t, sender.count(), "a held child is never POSTed") + + held := getDelivery(t, conn, id) + assert.Equal(t, 0, held.Attempts, + "five causal parks cost the child nothing: the retries it will need once its parent lands "+ + "are still there") + assert.Equal(t, store.DeliveryStatePending, held.State, "and it is still waiting, not decided") + + // Now the wall clock, and only the wall clock, decides. The attempt ledger + // reads zero — a counter-driven deadline would wait forever here. + _, err := conn.ExecContext(ctx, + `UPDATE outbound_deliveries SET created_at = now() - interval '2 hours' WHERE activity_id = $1`, id) + require.NoError(t, err) + + worked, err := w.DeliverNext(ctx) + require.NoError(t, err) + require.True(t, worked) + + expired := getDelivery(t, conn, id) + assert.Equal(t, store.DeliveryStatePoisoned, expired.State, + "the causal budget is bounded by ELAPSED TIME, not by attempts: a parent that never landed "+ + "still poisons its child on the deadline, however cheap the waiting was") + assert.Equal(t, store.PoisonClassParentUnaccepted, expired.LastErrorClass, + "and it poisons for the causal reason, so the outcome stays queryable and redrivable") +} diff --git a/internal/outbound/worker.go b/internal/outbound/worker.go index 60fca98..c4472c6 100644 --- a/internal/outbound/worker.go +++ b/internal/outbound/worker.go @@ -563,24 +563,15 @@ const parkDelay = 5 * time.Second // REAL delay into the future so it is not instantly re-claimable, and never // poisons — park does not call poison, and no state here is terminal. // -// IT DOES, HOWEVER, ADVANCE THE POISON BUDGET, and this comment used to claim -// the opposite. Release's SET clause does not touch `attempts` — but ClaimNext -// already did `attempts = attempts + 1` to get here (store/outbound_deliveries.go), -// and nothing ever puts it back. So a delivery parked for parkDelay is -// re-claimed ~5s later, +1 again, and a kill switch held for ~40s exhausts -// maxAttempts. The delivery does not poison while parked, but the first -// retryable failure AFTER the switch clears sees Attempts >= maxAttempts in -// releaseOrPoison and poisons instead of retrying. -// -// Only RedrivePoisoned resets attempts to 0, and only for rows already in -// 'poisoned' — i.e. the repair exists but runs after the fall, not instead of -// it. Fixing this means teaching the store to distinguish "retryable failure, -// back off" (increment is correct) from "parked by a switch" (it is not); -// see FOLLOWUPS.md, "Outbound delivery (task 15)". Documented deliberately -// rather than patched here: it needs a failing test first. +// It costs the delivery nothing either. ReleaseParked, not Release, is what +// settles a hold: it hands back the increment this claim's own ClaimNext charged, +// under the same fence, so a switch held for hours leaves the retry budget +// exactly where the operator found it (store/outbound_deliveries.go states why +// the fence is what makes that give-back safe). A held delivery reaches the +// attempt cap only through attempts it actually made. func (w *Worker) park(ctx context.Context, delivery *store.OutboundDelivery, class, reason string) error { next := time.Now().Add(parkDelay) - _, _, err := w.deliveries.Release(ctx, delivery.ActivityID, delivery.TargetInbox, class, reason, 0, next, *delivery.ClaimedUntil) + _, _, err := w.deliveries.ReleaseParked(ctx, delivery.ActivityID, delivery.TargetInbox, class, reason, 0, next, *delivery.ClaimedUntil) if err != nil { return fmt.Errorf("park delivery %s: %w", delivery.ActivityID, err) } @@ -591,15 +582,15 @@ func (w *Worker) park(ctx context.Context, delivery *store.OutboundDelivery, cla // parkCausal holds a causally-ineligible delivery WITHOUT a future delay: a held // child must become claimable the instant its bridge-origin parent is accepted // (in practice the parent, a lower-seq delivery on the same serial line, is -// delivered first, so this rarely re-fires). Like park it never poisons — but, -// like park, it DOES advance the poison budget, because ClaimNext incremented -// attempts and Release never puts it back (see park's doc above; this comment -// previously asserted the opposite). What keeps that from mattering here is the -// wall-clock bound: causalStatus poisons on the CausalWaitBudget deadline, not -// on the attempt count, so a held child's outcome is decided by elapsed time -// even after its retry budget is spent. +// delivered first, so this rarely re-fires). Like park it never poisons, and +// like park it is attempt-neutral — which matters most here, since a child that +// re-claims immediately would otherwise spend its whole budget in seconds. +// +// The wait's OUTCOME is still decided by the wall clock and not by the ledger: +// causalStatus poisons on the CausalWaitBudget deadline, so however many times a +// child cycles through this hold, what ends the wait is elapsed time. func (w *Worker) parkCausal(ctx context.Context, delivery *store.OutboundDelivery, class, reason string) error { - _, _, err := w.deliveries.Release(ctx, delivery.ActivityID, delivery.TargetInbox, class, reason, 0, time.Now(), *delivery.ClaimedUntil) + _, _, err := w.deliveries.ReleaseParked(ctx, delivery.ActivityID, delivery.TargetInbox, class, reason, 0, time.Now(), *delivery.ClaimedUntil) if err != nil { return fmt.Errorf("park (causal) delivery %s: %w", delivery.ActivityID, err) } diff --git a/internal/store/interfaces.go b/internal/store/interfaces.go index 18018e5..1e1fdbf 100644 --- a/internal/store/interfaces.go +++ b/internal/store/interfaces.go @@ -706,6 +706,18 @@ type OutboundDeliveries interface { // Same (exists, applied) split as MarkDelivered. Release(ctx context.Context, activityID, targetInbox, errorClass, excerpt string, lastStatusCode int, nextAttempt, claimToken time.Time) (exists, applied bool, err error) + // ReleaseParked is Release for a delivery that was HELD, not tried: an + // operator kill switch, dry-run, or a causal wait. It is identical in every + // respect but one — it is ATTEMPT-NEUTRAL, handing back the single increment + // its own claim took (ClaimNext does attempts = attempts + 1), so a hold + // costs the delivery none of the retry budget it will need when the hold + // lifts. Real failures still spend it; only park claims are given back. + // + // Same fence and non-terminal guard as Release (state = 'pending' AND + // claimed_until = claimToken), which is what stops a stale park from + // un-counting a newer claim's attempt. Same (exists, applied) split. + ReleaseParked(ctx context.Context, activityID, targetInbox, errorClass, excerpt string, lastStatusCode int, nextAttempt, claimToken time.Time) (exists, applied bool, err error) + // MarkPoisoned permanently fails the delivery (state=poisoned, lease // cleared, error class/status/excerpt stored). Poisoned rows are skipped by // ClaimNext and stop blocking their ordering key. claimToken must equal the diff --git a/internal/store/outbound_deliveries.go b/internal/store/outbound_deliveries.go index 099dbaf..1514df6 100644 --- a/internal/store/outbound_deliveries.go +++ b/internal/store/outbound_deliveries.go @@ -207,18 +207,26 @@ func (r *postgresOutboundDeliveries) MarkDelivered(ctx context.Context, activity activityID, targetInbox, lastStatusCode, claimToken.UTC()) } -func (r *postgresOutboundDeliveries) Release(ctx context.Context, activityID, targetInbox, errorClass, excerpt string, lastStatusCode int, nextAttempt, claimToken time.Time) (bool, bool, error) { - // Fencing + non-terminal guard: only the current claim holder reschedules - // (claimed_until == claimToken, state = 'pending'), so a stale worker's late - // release cannot resurrect a delivery a newer attempt already drove to a - // terminal state. The row stays pending with the lease cleared so a retry - // can re-claim after the backoff. - query := ` +// releaseStatement is the fenced reschedule BOTH releases run, with exactly one +// slot: what the attempt ledger does. Release leaves it alone; ReleaseParked +// hands the claim's increment back. +// +// It is one template rather than two statements because the FENCE is what makes +// the give-back safe. `state = 'pending' AND claimed_until = $7` admits only the +// worker still holding the claim, so the only increment a park can subtract is +// the one its own claim just added. A second copy of this statement is a copy of +// that fence, and the day the two drift is the day a stale worker's late park +// un-counts an attempt the CURRENT claim genuinely spent — a delivery quietly +// gaining retries it already used, visible nowhere until a poison budget that +// should have stopped never does. +// +// The slot is filled from a CONSTANT below and never from input. +const releaseStatement = ` WITH updated AS ( UPDATE outbound_deliveries SET claimed_until = NULL, next_attempt_at = $3, last_error_class = $4, response_excerpt = $5, - last_status_code = $6, updated_at = now() + last_status_code = $6, updated_at = now()%s WHERE activity_id = $1 AND target_inbox = $2 AND state = 'pending' AND claimed_until = $7 @@ -228,10 +236,50 @@ func (r *postgresOutboundDeliveries) Release(ctx context.Context, activityID, ta EXISTS (SELECT 1 FROM outbound_deliveries WHERE activity_id = $1 AND target_inbox = $2), EXISTS (SELECT 1 FROM updated)` +// handBackClaimAttempt is ReleaseParked's ONE difference from Release. GREATEST +// floors the ledger at zero so an unmatched arithmetic edge can never write a +// negative attempt count (attempts is INT NOT NULL DEFAULT 0, so there is no +// NULL to guard). +const handBackClaimAttempt = `, + attempts = GREATEST(attempts - 1, 0)` + +func (r *postgresOutboundDeliveries) Release(ctx context.Context, activityID, targetInbox, errorClass, excerpt string, lastStatusCode int, nextAttempt, claimToken time.Time) (bool, bool, error) { + // Fencing + non-terminal guard: only the current claim holder reschedules + // (claimed_until == claimToken, state = 'pending'), so a stale worker's late + // release cannot resurrect a delivery a newer attempt already drove to a + // terminal state. The row stays pending with the lease cleared so a retry + // can re-claim after the backoff. The attempt ClaimNext charged stays + // charged: this delivery was TRIED. + query := fmt.Sprintf(releaseStatement, "") + return r.markResult(ctx, "release", query, activityID, targetInbox, nextAttempt.UTC(), errorClass, excerpt, lastStatusCode, claimToken.UTC()) } +// ReleaseParked is the ATTEMPT-NEUTRAL release: the same fenced reschedule as +// Release, with the claim's own increment handed back. +// +// THE INVARIANT IS "a park hands back exactly its own claim's increment". A +// delivery that was HELD — by the kill switch, a dry run, a causal wait — was +// never tried, but ClaimNext charges an attempt to every claim alike, so without +// the give-back a hold spends retry budget the delivery no longer has when the +// hold lifts, and a long enough hold poisons a delivery nothing ever attempted. +// The fence is what makes the decrement safe rather than merely convenient: only +// the worker whose claim added the increment can subtract one, so no park can +// reach an attempt another claim spent. A park is therefore net-zero and a real +// failure still costs exactly one. +// +// settleLater's Release deliberately keeps its increment instead of parking: a +// held-for-settlement delivery is already accepted by the peer and handle() +// short-circuits before re-POSTing it, so it can never re-poison and its growing +// attempt count only widens the backoff between bookkeeping retries. +func (r *postgresOutboundDeliveries) ReleaseParked(ctx context.Context, activityID, targetInbox, errorClass, excerpt string, lastStatusCode int, nextAttempt, claimToken time.Time) (bool, bool, error) { + query := fmt.Sprintf(releaseStatement, handBackClaimAttempt) + + return r.markResult(ctx, "release parked", query, + activityID, targetInbox, nextAttempt.UTC(), errorClass, excerpt, lastStatusCode, claimToken.UTC()) +} + func (r *postgresOutboundDeliveries) MarkPoisoned(ctx context.Context, activityID, targetInbox, errorClass, excerpt string, lastStatusCode int, claimToken time.Time) (bool, bool, error) { // Fencing + non-terminal guard: only the current claim holder may poison // (claimed_until == claimToken, state = 'pending'). A poisoned delivery diff --git a/internal/store/outbound_delivery_park_test.go b/internal/store/outbound_delivery_park_test.go new file mode 100644 index 0000000..2c6ee99 --- /dev/null +++ b/internal/store/outbound_delivery_park_test.go @@ -0,0 +1,204 @@ +package store + +import ( + "context" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// ReleaseParked — the attempt-neutral release. ClaimNext charges an attempt to +// every claim, which is right for a delivery that was TRIED and wrong for one +// that was HELD: a kill switch, a dry run and a causal wait all put the delivery +// back untouched, and the retry budget they consume is budget the delivery no +// longer has when the hold lifts. ReleaseParked is Release with that one +// increment handed back, under the identical fence — so a hold is free, and a +// failure still costs exactly one. + +// parkedAttempts reads the attempt ledger of the delivery under test — the one +// column every assertion in this file is really about. +func parkedAttempts(t *testing.T, repo OutboundDeliveries) int { + t.Helper() + got, err := repo.Get(context.Background(), delActivityID, delTargetInbox) + require.NoError(t, err) + return got.Attempts +} + +func TestOutboundDeliveries_ReleaseParkedIsAttemptNeutral(t *testing.T) { + database := deliveryTestDB(t) + activities := NewOutboundActivities(database) + repo := NewOutboundDeliveries(database) + ctx := context.Background() + + seedActivity(t, activities, testActivity()) + _, err := repo.Enqueue(ctx, testDelivery()) + require.NoError(t, err) + + claimed, err := repo.ClaimNext(ctx, time.Minute) + require.NoError(t, err) + require.NotNil(t, claimed.ClaimedUntil) + require.Equal(t, 1, claimed.Attempts, "the claim charged its attempt (that is what a park hands back)") + token := *claimed.ClaimedUntil + + next := time.Now().Add(5 * time.Second).UTC() + exists, applied, err := repo.ReleaseParked(ctx, delActivityID, delTargetInbox, + "switch_parked", "outbound kill switch engaged", 0, next, token) + require.NoError(t, err) + assert.True(t, exists, "the row is there to be parked") + assert.True(t, applied, "the claim holder parks its own delivery") + + got, err := repo.Get(ctx, delActivityID, delTargetInbox) + require.NoError(t, err) + assert.Equal(t, 0, got.Attempts, + "a park is ATTEMPT-NEUTRAL: it hands back the increment its own claim took, so a delivery "+ + "held by an operator switch keeps every retry it had before the hold") + assert.Equal(t, DeliveryStatePending, got.State, "a parked delivery stays pending — a hold is not a failure") + assert.Nil(t, got.ClaimedUntil, "the lease is cleared so the delivery can be re-claimed when the hold lifts") + assert.WithinDuration(t, next, got.NextAttemptAt, time.Second, "the park delay is stored") + assert.Equal(t, "switch_parked", got.LastErrorClass, + "the hold's reason is recorded, so an operator can see WHY the queue is idle") +} + +func TestOutboundDeliveries_ReleaseParkedFencing(t *testing.T) { + // The give-back is the dangerous half of this method: an unfenced decrement + // would let a delivery lose attempts it genuinely spent. Both cases below + // are a park arriving too late to be anyone's business. + + t.Run("stale token cannot un-count a newer claim", func(t *testing.T) { + database := deliveryTestDB(t) + activities := NewOutboundActivities(database) + repo := NewOutboundDeliveries(database) + ctx := context.Background() + + seedActivity(t, activities, testActivity()) + _, err := repo.Enqueue(ctx, testDelivery()) + require.NoError(t, err) + + claimed, err := repo.ClaimNext(ctx, time.Minute) + require.NoError(t, err) + require.NotNil(t, claimed.ClaimedUntil) + staleToken := *claimed.ClaimedUntil + + // That worker wedged; its lease lapses and a second worker takes the + // delivery (attempts = 2). Then the first wakes up and tries to park. + _, err = database.ExecContext(ctx, + `UPDATE outbound_deliveries SET claimed_until = now() - interval '1 minute' WHERE activity_id = $1`, + delActivityID) + require.NoError(t, err) + + reclaimed, err := repo.ClaimNext(ctx, time.Minute) + require.NoError(t, err) + require.Equal(t, 2, reclaimed.Attempts) + require.NotNil(t, reclaimed.ClaimedUntil) + + _, applied, err := repo.ReleaseParked(ctx, delActivityID, delTargetInbox, + "switch_parked", "kill switch", 0, time.Now().Add(5*time.Second), staleToken) + require.NoError(t, err) + assert.False(t, applied, "a stale fencing token must not park a delivery someone else now holds") + + got, err := repo.Get(ctx, delActivityID, delTargetInbox) + require.NoError(t, err) + assert.Equal(t, 2, got.Attempts, + "the give-back is fenced: a stale park must never un-count an attempt the CURRENT claim spent") + require.NotNil(t, got.ClaimedUntil, "the current worker's lease survives a stale park") + assert.WithinDuration(t, *reclaimed.ClaimedUntil, *got.ClaimedUntil, time.Second, + "the lease is the second worker's, untouched") + + // The other side of the fence, so "nothing happened" above is the FENCE + // refusing and not the method declining to work at all: the worker that + // actually holds the claim parks, and gets its own attempt back. + _, applied, err = repo.ReleaseParked(ctx, delActivityID, delTargetInbox, + "switch_parked", "kill switch", 0, time.Now().Add(5*time.Second), *reclaimed.ClaimedUntil) + require.NoError(t, err) + assert.True(t, applied, "the CURRENT claim holder parks — the fence blocks the stale worker, not the real one") + assert.Equal(t, 1, parkedAttempts(t, repo), + "the holder's park hands back its OWN claim only: the earlier, genuinely-spent attempt remains") + }) + + t.Run("terminal row is untouched", func(t *testing.T) { + database := deliveryTestDB(t) + activities := NewOutboundActivities(database) + repo := NewOutboundDeliveries(database) + ctx := context.Background() + + seedActivity(t, activities, testActivity()) + _, err := repo.Enqueue(ctx, testDelivery()) + require.NoError(t, err) + + claimed, err := repo.ClaimNext(ctx, time.Minute) + require.NoError(t, err) + require.NotNil(t, claimed.ClaimedUntil) + token := *claimed.ClaimedUntil + + // The actor opts out mid-flight: the claimed row is cancelled under the + // worker. Its park then arrives against a terminal row. + cancelled, err := repo.CancelForActor(ctx, testDID) + require.NoError(t, err) + require.Equal(t, int64(1), cancelled) + + exists, applied, err := repo.ReleaseParked(ctx, delActivityID, delTargetInbox, + "switch_parked", "kill switch", 0, time.Now().Add(5*time.Second), token) + require.NoError(t, err) + assert.True(t, exists, "the row exists — it is simply no longer parkable") + assert.False(t, applied, "a park must not apply to a TERMINAL delivery") + + got, err := repo.Get(ctx, delActivityID, delTargetInbox) + require.NoError(t, err) + assert.Equal(t, DeliveryStateCancelled, got.State, + "a late park must never resurrect a cancelled delivery back into the queue") + assert.Equal(t, 1, got.Attempts, "and must not rewrite the attempt ledger of a decided row") + }) +} + +func TestOutboundDeliveries_ParksAreNetZeroAndFailuresRetain(t *testing.T) { + database := deliveryTestDB(t) + activities := NewOutboundActivities(database) + repo := NewOutboundDeliveries(database) + ctx := context.Background() + + seedActivity(t, activities, testActivity()) + _, err := repo.Enqueue(ctx, testDelivery()) + require.NoError(t, err) + + // claimNow re-claims the delivery immediately; the releases below schedule + // next_attempt_at in the past so no test ever waits on a backoff. + claimNow := func(wantAttempts int, why string) time.Time { + t.Helper() + claimed, err := repo.ClaimNext(ctx, time.Minute) + require.NoError(t, err, "re-claim: %s", why) + require.NotNil(t, claimed.ClaimedUntil) + require.Equalf(t, wantAttempts, claimed.Attempts, "attempts after the claim: %s", why) + return *claimed.ClaimedUntil + } + past := func() time.Time { return time.Now().Add(-time.Second).UTC() } + attempts := func() int { return parkedAttempts(t, repo) } + + // A REAL failure: the attempt is spent and stays spent. + token := claimNow(1, "first attempt") + _, applied, err := repo.Release(ctx, delActivityID, delTargetInbox, "5xx", "bad gateway", 502, past(), token) + require.NoError(t, err) + require.True(t, applied) + require.Equal(t, 1, attempts(), "a genuine failure keeps its attempt") + + // A PARK on top of that history: net zero. The failure below it is not + // forgotten, and the park itself leaves no trace in the budget. + token = claimNow(2, "the park's own claim") + _, applied, err = repo.ReleaseParked(ctx, delActivityID, delTargetInbox, + "switch_parked", "kill switch", 0, past(), token) + require.NoError(t, err) + assert.True(t, applied, "the claim holder parks") + require.Equal(t, 1, attempts(), + "parks are net-zero and failures retain: after a park the ledger reads exactly the ONE real "+ + "failure that came before it — a hold neither spends budget nor erases history") + + // And the next real failure resumes the count where the failures left it. + token = claimNow(2, "second real attempt") + _, applied, err = repo.Release(ctx, delActivityID, delTargetInbox, "5xx", "bad gateway", 502, past(), token) + require.NoError(t, err) + require.True(t, applied) + assert.Equal(t, 2, attempts(), + "two real failures with a park between them cost exactly two attempts — an attempt cap counts "+ + "what was TRIED, never what was held") +} -- 2.51.2