From bbab497b8e081425576f8fd7cafd500a559c01b9 Mon Sep 17 00:00:00 2001 From: Bretton Date: Mon, 17 Aug 2026 21:33:49 -0700 Subject: [PATCH] outbound: dedupe phrase pinned, cancelled parents decided, poison gated MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Chunk-4 findings 3 + 1 and chunk-9 finding 4 from the second-opinion re-review. - isDuplicate matched ANY 400 whose body contained "already", so Lemmy's genuine rejections (bans, blocks, already_invalid, duplicate_title) went terminal `delivered` — flipping vote state the scores read and opening the causal gate for objects that never landed, with the justifying body discarded. It now matches the received-activity dedupe phrase ("already received", underscore- folded), and the one branch that turns a 400 into a delivered row Debug-logs status + body so a misfire leaves evidence. - The causal gate had verdicts for accepted and poisoned parents but not the third terminal state, cancelled: a child of a cancelled parent looped claim→park at next_attempt=now() with worked=true, a tight 2-UPDATE/3-SELECT spin against Postgres for the full 6h budget. ParentDeliveryPoisoned becomes ParentDeliveryDisposition (open/poisoned/cancelled — cancelled only when EVERY delivery was cancelled and at least one exists, so a pending sibling keeps the parent open and fediverse-origin parents never read as decided); the child poisons as parent_cancelled (redrivable, reason kept — a cancelled child would be unrecoverable and reasonless). The ordinary causal park gains a 1s delay so no future wait path can respin at DB speed, and a parent-lookup store error now parks on the worker's real backoff instead of the causal cadence. - poison() now mirrors countPark: metricPoisoned increments only when the (exists, applied) fence accepted the mark, and a bounced poison Warns — the counter DEPLOY.md tells operators to redrive on no longer counts stale-lease no-ops. New PoisonClassParentCancelled joins neverReachedTheWireClasses, per that file's own contract. Co-Authored-By: Claude Opus 4.8 --- .../outbound/causal_cancelled_parent_test.go | 297 ++++++++++++++++++ internal/outbound/causal_gating_test.go | 5 + .../outbound/duplicate_classification_test.go | 141 +++++++++ internal/outbound/park_bounce_test.go | 3 +- internal/outbound/poison_bounce_test.go | 62 ++++ internal/outbound/worker.go | 164 ++++++++-- internal/store/divergence.go | 20 +- internal/store/divergence_unknown_test.go | 18 +- internal/store/interfaces.go | 19 +- internal/store/outbound_deliveries.go | 73 ++++- 10 files changed, 738 insertions(+), 64 deletions(-) create mode 100644 internal/outbound/causal_cancelled_parent_test.go create mode 100644 internal/outbound/duplicate_classification_test.go create mode 100644 internal/outbound/poison_bounce_test.go diff --git a/internal/outbound/causal_cancelled_parent_test.go b/internal/outbound/causal_cancelled_parent_test.go new file mode 100644 index 0000000..7c9a126 --- /dev/null +++ b/internal/outbound/causal_cancelled_parent_test.go @@ -0,0 +1,297 @@ +package outbound + +import ( + "context" + "database/sql" + "fmt" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "tidepool/internal/errors" + "tidepool/internal/store" +) + +// THE CAUSAL GATE HAS A THIRD TERMINAL PARENT, and it used to have no verdict +// for it. +// +// The gate reads three outcomes off a bridge-origin parent: accepted → the child +// is eligible, its delivery poisoned → the child is poisoned, anything else → +// wait. But a parent delivery can go terminal a THIRD way — CANCELLED, by the +// consent recheck at claim time, an operator cancel, or an actor/community +// sweep. A cancelled parent leaves the pending index, never gets accepted_at, +// and is not poisoned, so it fell into "anything else": wait forever. +// +// Forever is the cheap word for it. parkCausal released the child at +// next_attempt_at = now(), which makes it INSTANTLY re-claimable, and DeliverNext +// reports worked=true so Run never sleeps — a claim/park loop of two UPDATEs and +// two or three SELECTs per iteration, hammering Postgres at whatever rate the +// round trip allows, for the whole six-hour CausalWaitBudget. The trigger is +// ordinary: author A's comment is consent-cancelled at claim time while B's +// reply sits behind it on the same ordering key. +// +// The child is POISONED rather than cancelled, and the choice is about recovery +// rather than about tone. Cancelled is the queue's terminal state for "this must +// not go out", and it is unreachable afterwards: RedrivePoisoned matches only +// poisoned rows, so a cancelled child could never be revived. But the parent's +// own cancellation may well be redressed — a delivery_paused actor is a +// transient #account state, an operator cancel is reversed by re-enqueueing — +// and when it is, the operator needs both rows back. Poisoning also records WHY +// (last_error_class = parent_cancelled) where a cancel records nothing, and +// store.PoisonClassParentCancelled joins the never-reached-the-wire set so the +// divergence sweep does not report a delivery nobody ever sent as an unknown +// outcome. + +// seedParentDeliveryInState writes the parent's OWN activity + delivery — the activity +// whose object.id maps back to gParentATURI, which is what the causal gate keys +// on — and drives it to the given terminal state. +func seedParentDeliveryInState(t *testing.T, conn *sql.DB, state store.DeliveryState) { + t.Helper() + ctx := context.Background() + parentID := "https://coves.social/ap/activity/" + repeatHex64("cancelledparent") + _, err := store.NewOutboundActivities(conn).Insert(ctx, store.OutboundActivity{ + ActivityID: parentID, + ActorDID: wActorDID, + Kind: "Create", + Payload: []byte(fmt.Sprintf(`{"id":%q,"type":"Create","object":{"type":"Page","id":%q}}`, + parentID, objectURLFor(gParentATURI))), + }) + require.NoError(t, err) + _, err = store.NewOutboundDeliveries(conn).Enqueue(ctx, store.OutboundDelivery{ + ActivityID: parentID, TargetInbox: wInbox, OrderingKey: wCommunityAPID, + }) + require.NoError(t, err) + _, err = conn.ExecContext(ctx, + `UPDATE outbound_deliveries SET state = $2 WHERE activity_id = $1`, parentID, string(state)) + require.NoError(t, err) +} + +func TestCausalGating_CancelledParentDoesNotSpin(t *testing.T) { + conn := workerTestDB(t) + ctx := context.Background() + seedWorkerActor(t, conn, true, false) + seedBridgeParent(t, conn, false) // the parent object exists and was never accepted + seedParentDeliveryInState(t, conn, store.DeliveryStateCancelled) + + id := seedDelivery(t, conn, "Create", gParentATURI, createPayload("child")) + + sender := &fakeSender{} + w := newWorker(t, conn, sender, nil) + + worked, err := w.DeliverNext(ctx) + require.NoError(t, err) + require.True(t, worked) + require.Zero(t, sender.count(), "a child of cancelled content is never POSTed") + + d := getDelivery(t, conn, id) + assert.Equal(t, store.DeliveryStatePoisoned, d.State, + "a parent whose delivery was CANCELLED is terminal-and-never-accepted: the child cannot "+ + "land behind it, so it is decided now rather than waiting for a parent that will "+ + "never arrive") + assert.Equal(t, store.PoisonClassParentCancelled, d.LastErrorClass, + "and it says WHY, distinctly from parent_poisoned (a failure) and parent_unaccepted "+ + "(a deadline) — the three reach the same state from different causes and an "+ + "operator triages them differently") + + // The spin itself: the second claim must find nothing. Before the fix the + // child was released at now(), so this returned worked=true immediately and + // went on doing so for the full six-hour budget. + worked, err = w.DeliverNext(ctx) + require.NoError(t, err) + assert.False(t, worked, + "and the queue is EMPTY afterwards: a child that stays instantly re-claimable turns "+ + "DeliverNext into a hot claim/park loop that Run never sleeps out of") +} + +// The delay is the second half of the fix, and it stands on its own: whatever +// verdict the gate reaches, a child that is genuinely still WAITING must not be +// re-claimable at database speed. parkCausal used to stamp next_attempt_at = +// now(), which is a spin under any residual wait path. +func TestCausalGating_ParkedChildIsNotInstantlyReclaimable(t *testing.T) { + conn := workerTestDB(t) + ctx := context.Background() + seedWorkerActor(t, conn, true, false) + seedBridgeParent(t, conn, false) // merely pending: the child legitimately waits + id := seedDelivery(t, conn, "Create", gParentATURI, createPayload("child")) + + w := newWorker(t, conn, &fakeSender{}, nil) + worked, err := w.DeliverNext(ctx) + require.NoError(t, err) + require.True(t, worked) + + d := getDelivery(t, conn, id) + require.Equal(t, store.DeliveryStatePending, d.State, "it is held, not decided") + assert.True(t, d.NextAttemptAt.After(time.Now()), + "a causal park schedules the child a real interval out: at now() the very next claim "+ + "picks it straight back up, and the loop runs at whatever rate the round trip allows") + + worked, err = w.DeliverNext(ctx) + require.NoError(t, err) + assert.False(t, worked, + "so back-to-back claims find nothing — which is what turns the wait from a busy loop "+ + "into a wait") +} + +// failingObjects is the real store with its parent lookup broken: a connection +// blip, a statement timeout, the transient the gate is meant to hold through. +type failingObjects struct { + store.OutboundObjects +} + +func (failingObjects) GetByATURI(context.Context, string) (*store.OutboundObject, error) { + return nil, fmt.Errorf("dial tcp: connection reset by peer") +} + +// A lookup ERROR takes the same "hold rather than poison" verdict a pending +// parent does — correctly, since a blip must not poison anybody's reply — and so +// it inherited the same zero-delay release. The result is worse than the +// cancelled-parent spin: it is a hot loop against a database that is ALREADY in +// trouble, for every gated delivery at once. +func TestCausalGating_ParentLookupErrorBacksOff(t *testing.T) { + conn := workerTestDB(t) + ctx := context.Background() + seedWorkerActor(t, conn, true, false) + seedBridgeParent(t, conn, false) + id := seedDelivery(t, conn, "Create", gParentATURI, createPayload("child")) + + sender := &fakeSender{} + // A REAL backoff base: the helper's 1ms would put next_attempt_at in the past + // before the row could be read, making the assertion vacuous. + w := newWorker(t, conn, sender, func(o *WorkerOptions) { + o.Objects = failingObjects{OutboundObjects: store.NewOutboundObjects(conn)} + o.BackoffBase = 30 * time.Second + }) + + worked, err := w.DeliverNext(ctx) + require.NoError(t, err) + require.True(t, worked) + require.Zero(t, sender.count(), + "a delivery whose causal eligibility could not be READ is not delivered on a guess") + + d := getDelivery(t, conn, id) + assert.Equal(t, store.DeliveryStatePending, d.State, + "a store blip holds the delivery — it must never poison somebody's reply") + assert.Equal(t, "parent_lookup_failed", d.LastErrorClass, + "and it is held under its own class, so an operator can tell a database problem from a "+ + "parent that is simply a beat behind") + assert.True(t, d.NextAttemptAt.After(time.Now().Add(time.Second)), + "with a REAL backoff: re-asking a database that just failed, as fast as the round trip "+ + "allows, is the loop that turns a blip into an outage") + assert.Zero(t, d.Attempts, + "and the hold is still attempt-neutral — the delivery was never tried") +} + +// The whole reason the cancelled parent had no verdict is that the gate could +// only ask one question of the parent's delivery. This pins the store answer +// underneath the worker: the three dispositions must come apart, and an ordinary +// pending parent must read as neither. +func TestParentDeliveryDisposition_SeparatesPoisonedFromCancelled(t *testing.T) { + for _, tc := range []struct { + name string + state store.DeliveryState + want store.ParentDeliveryDisposition + }{ + {"pending", store.DeliveryStatePending, store.ParentDeliveryOpen}, + {"delivered", store.DeliveryStateDelivered, store.ParentDeliveryOpen}, + {"poisoned", store.DeliveryStatePoisoned, store.ParentDeliveryPoisoned}, + {"cancelled", store.DeliveryStateCancelled, store.ParentDeliveryCancelled}, + } { + t.Run(tc.name, func(t *testing.T) { + conn := workerTestDB(t) + ctx := context.Background() + seedWorkerActor(t, conn, true, false) + seedParentDeliveryInState(t, conn, tc.state) + + got, err := store.NewOutboundDeliveries(conn). + ParentDeliveryDisposition(ctx, gParentATURI, wInbox) + require.NoError(t, err) + assert.Equal(t, tc.want, got) + }) + } +} + +// A parent with NO delivery row at all reads open, not cancelled: "every row is +// cancelled" over an empty set is vacuously true, and taking that as a verdict +// would poison every child of a fediverse-origin parent the moment the worker +// asked. +func TestParentDeliveryDisposition_NoParentRowIsOpen(t *testing.T) { + conn := workerTestDB(t) + got, err := store.NewOutboundDeliveries(conn). + ParentDeliveryDisposition(context.Background(), gParentATURI, wInbox) + require.NoError(t, err) + assert.Equal(t, store.ParentDeliveryOpen, got, + "no parent delivery is not a cancelled parent — an empty set must not answer a "+ + "question about what happened to the parent") +} + +// One object can be federated by more than one activity (a Create and a later +// Update carry the same object.id), so the parent's disposition is read across +// ALL of them. A single cancelled row beside a live one does not condemn the +// child: something can still make the parent land. +func TestParentDeliveryDisposition_ALiveSiblingKeepsTheParentOpen(t *testing.T) { + conn := workerTestDB(t) + ctx := context.Background() + seedWorkerActor(t, conn, true, false) + seedParentDeliveryInState(t, conn, store.DeliveryStateCancelled) + + // The Update of the same object, still pending. + updateID := "https://coves.social/ap/activity/" + repeatHex64("parentupdate") + _, err := store.NewOutboundActivities(conn).Insert(ctx, store.OutboundActivity{ + ActivityID: updateID, + ActorDID: wActorDID, + Kind: "Update", + Payload: []byte(fmt.Sprintf(`{"id":%q,"type":"Update","object":{"type":"Page","id":%q}}`, + updateID, objectURLFor(gParentATURI))), + }) + require.NoError(t, err) + _, err = store.NewOutboundDeliveries(conn).Enqueue(ctx, store.OutboundDelivery{ + ActivityID: updateID, TargetInbox: wInbox, OrderingKey: wCommunityAPID, + }) + require.NoError(t, err) + + got, err := store.NewOutboundDeliveries(conn).ParentDeliveryDisposition(ctx, gParentATURI, wInbox) + require.NoError(t, err) + assert.Equal(t, store.ParentDeliveryOpen, got, + "one cancelled activity for the parent object is not the parent's fate while another is "+ + "still in flight — the child waits rather than being poisoned out from under a "+ + "delivery that may still land") + + // Poison still wins outright, which is the pre-existing contract: a poisoned + // parent delivery poisons the child whatever else is on the object. + _, err = conn.ExecContext(ctx, + `UPDATE outbound_deliveries SET state = 'poisoned' WHERE activity_id = $1`, updateID) + require.NoError(t, err) + got, err = store.NewOutboundDeliveries(conn).ParentDeliveryDisposition(ctx, gParentATURI, wInbox) + require.NoError(t, err) + assert.Equal(t, store.ParentDeliveryPoisoned, got, + "and a poisoned row outranks a cancelled one: parent_poisoned is the louder verdict and "+ + "the pre-existing one") +} + +// A guard on the assumption the fix rests on: errors.IsNotFound is what the gate +// reads to call a parent fediverse-origin, and the disposition query must not +// change that path. +func TestCausalGating_CancelledParentStillEligibleWhenObjectIsFediverseOrigin(t *testing.T) { + conn := workerTestDB(t) + ctx := context.Background() + seedWorkerActor(t, conn, true, false) + // No outbound_objects row: a Lemmy post. Whatever deliveries exist against + // that at-uri, the reply is eligible — the object is already on the peer. + _, err := store.NewOutboundObjects(conn).GetByATURI(ctx, gParentATURI) + require.True(t, errors.IsNotFound(err)) + seedParentDeliveryInState(t, conn, store.DeliveryStateCancelled) + + id := seedDelivery(t, conn, "Create", gParentATURI, createPayload("child")) + sender := &fakeSender{} + w := newWorker(t, conn, sender, nil) + worked, err := w.DeliverNext(ctx) + require.NoError(t, err) + require.True(t, worked) + + assert.Equal(t, 1, sender.count(), + "a reply whose parent has no outbound_objects row delivers immediately — the "+ + "cancelled-parent verdict must not reach past the bridge-origin test") + assert.Equal(t, store.DeliveryStateDelivered, getDelivery(t, conn, id).State) +} diff --git a/internal/outbound/causal_gating_test.go b/internal/outbound/causal_gating_test.go index 5478f00..08fbbb2 100644 --- a/internal/outbound/causal_gating_test.go +++ b/internal/outbound/causal_gating_test.go @@ -56,6 +56,11 @@ func TestCausalGating_BridgeParentUnacceptedIsIneligible(t *testing.T) { // Flip the parent to accepted: the SAME delivery must now go out. This is // what keeps the negative non-vacuous — the gate opens, it delivers. require.NoError(t, store.NewOutboundObjects(conn).SetAccepted(context.Background(), gParentATURI)) + // A causal park schedules the child a real interval out (causalParkDelay), + // so the re-claim is rewound rather than slept out — the house idiom, cf. + // park_budget_test.go. The interval is what keeps the hold from spinning; + // TestCausalGating_ParkedChildIsNotInstantlyReclaimable pins it directly. + clearParkDelay(t, conn, id) worked, err := w.DeliverNext(context.Background()) require.NoError(t, err) assert.True(t, worked) diff --git a/internal/outbound/duplicate_classification_test.go b/internal/outbound/duplicate_classification_test.go new file mode 100644 index 0000000..dd281f8 --- /dev/null +++ b/internal/outbound/duplicate_classification_test.go @@ -0,0 +1,141 @@ +package outbound + +import ( + "context" + "fmt" + "net/http" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "tidepool/internal/store" +) + +// THE DEDUPE CLASSIFICATION IS THE ONE PLACE A REJECTION BECOMES A SUCCESS. +// +// classify() maps Lemmy's received-activity dedupe — a 400 whose body says the +// activity was already received — onto DELIVERED, because a redelivery after a +// crash is expected and safe. Everything else about a 400 is a refusal. +// +// Matching on the bare word "already" cannot tell those apart. Lemmy answers 400 +// with a whole family of slugs carrying it — banned_from_community, +// person_is_blocked, already_invalid, duplicate_title — and reading any of them +// as a delivery is not a missed retry, it is a FABRICATED one: +// +// - a Like the peer refused flips outbound_votes.delivered_state, and that +// column is an input to the score users are served (task 17b); +// - stampAccepted opens the causal gate for an object that never landed, so +// every child delivered behind it is rejected in turn; +// - the delivered path stores no excerpt, so the body that would have shown +// an operator what really happened is discarded. +// +// Each case below is a real 400 body containing "already" that is NOT the +// dedupe, and each must follow the ordinary 4xx path. + +func TestWorker_Non_DedupeFourHundredContainingAlreadyIsNotDelivered(t *testing.T) { + for _, tc := range []struct { + name string + body string + }{ + {"banned from community", `{"error":"person_is_banned_from_community","message":"this person is already banned"}`}, + {"blocked by the recipient", `{"error":"You have already been blocked from this community"}`}, + {"already invalid", `{"error":"already_invalid"}`}, + {"duplicate title", `{"error":"duplicate_title","message":"a post with this title already exists"}`}, + } { + t.Run(tc.name, func(t *testing.T) { + conn := workerTestDB(t) + seedWorkerActor(t, conn, true, false) + id := seedDelivery(t, conn, "Create", "", createPayload("x")) + setAttempts(t, conn, id, 2) // MaxAttempts is 3: the next 4xx decides + + w := newWorker(t, conn, senderReturning(httpErr(http.StatusBadRequest, tc.body)), nil) + worked, err := w.DeliverNext(context.Background()) + require.NoError(t, err) + require.True(t, worked) + + d := getDelivery(t, conn, id) + assert.Equal(t, store.DeliveryStatePoisoned, d.State, + "a 400 that merely CONTAINS \"already\" is a genuine rejection, not the "+ + "received-activity dedupe: it must follow the normal 4xx classification") + assert.Equal(t, "4xx", d.LastErrorClass, + "and it is labelled as the rejection it is") + assert.Contains(t, d.ResponseExcerpt, "already", + "with the peer's own words kept, which is the evidence the delivered path throws away") + }) + } +} + +// The Like case is the one with a number attached: a vote the peer REFUSED must +// not be recorded as delivered, because that column is subtracted from the score +// the API serves and nothing reconciles it afterwards. +func TestWorker_RefusedLikeDoesNotFlipTheVoteLedger(t *testing.T) { + conn := workerTestDB(t) + ctx := context.Background() + seedWorkerActor(t, conn, true, false) + + likeID := "https://coves.social/ap/activity/" + repeatHex64("Like") + voteATURI := "at://" + wActorDID + "/social.coves.feed.vote/3lzvotebanned" + _, err := store.NewOutboundVotes(conn).Upsert(ctx, store.OutboundVote{ + VoteATURI: voteATURI, + ActorDID: wActorDID, + SubjectATURI: "at://" + wCommunityDID + "/social.coves.community.postv2/3lzpost", + SubjectAPID: "https://lemmy.world/post/1", + CommunityDID: wCommunityDID, + Direction: "up", + CurrentActivityID: likeID, + }) + require.NoError(t, err) + + id := seedDelivery(t, conn, "Like", "", []byte(fmt.Sprintf( + `{"id":%q,"type":"Like","actor":%q,"object":"https://lemmy.world/post/1"}`, likeID, wActorID))) + setAttempts(t, conn, id, 2) + + // The author is banned in that community, so the vote is refused. + w := newWorker(t, conn, senderReturning( + httpErr(http.StatusBadRequest, + `{"error":"person_is_banned_from_community","message":"this person is already banned"}`)), nil) + worked, err := w.DeliverNext(ctx) + require.NoError(t, err) + require.True(t, worked) + + assert.Equal(t, store.DeliveryStatePoisoned, getDelivery(t, conn, id).State, + "a refused Like is a refusal") + + vote, err := store.NewOutboundVotes(conn).GetByATURI(ctx, voteATURI) + require.NoError(t, err) + assert.Equal(t, store.DeliveredStatePending, vote.DeliveredState, + "and the vote ledger must NOT read delivered for a vote the peer refused — that column "+ + "is an input to the score users are served, and nothing reconciles it later") +} + +// The dedupe itself still lands, in both spellings the peer stack can produce: +// the prose sentence the fake Lemmy in the outer acceptance test returns, and +// the underscored slug an error enum serializes to. Without this the fix above +// could be "classify nothing as a duplicate", which re-poisons every redelivery +// after a crash. +func TestWorker_ReceivedActivityDedupeStillDelivers(t *testing.T) { + for _, tc := range []struct { + name string + body string + }{ + {"prose", `{"error":"activity was already received"}`}, + {"capitalised prose", `{"error":"Activity was already received"}`}, + {"slug", `{"error":"already_received"}`}, + } { + t.Run(tc.name, func(t *testing.T) { + conn := workerTestDB(t) + seedWorkerActor(t, conn, true, false) + id := seedDelivery(t, conn, "Create", "", createPayload("x")) + + w := newWorker(t, conn, senderReturning(httpErr(http.StatusBadRequest, tc.body)), nil) + worked, err := w.DeliverNext(context.Background()) + require.NoError(t, err) + require.True(t, worked) + + assert.Equal(t, store.DeliveryStateDelivered, getDelivery(t, conn, id).State, + "a body reporting the activity was ALREADY RECEIVED is our own redelivery coming "+ + "back: it is delivered, never poisoned") + }) + } +} diff --git a/internal/outbound/park_bounce_test.go b/internal/outbound/park_bounce_test.go index 7d8b1cb..aa1495d 100644 --- a/internal/outbound/park_bounce_test.go +++ b/internal/outbound/park_bounce_test.go @@ -39,7 +39,8 @@ func TestWorker_BouncedParkIsNotCounted(t *testing.T) { name: "causal park", class: "parent_pending", park: func(w *Worker, ctx context.Context, d *store.OutboundDelivery) error { - return w.parkCausal(ctx, d, "parent_pending", "waiting for bridge-origin parent") + return w.parkCausal(ctx, d, "parent_pending", "waiting for bridge-origin parent", + causalParkDelay) }, }, } { diff --git a/internal/outbound/poison_bounce_test.go b/internal/outbound/poison_bounce_test.go new file mode 100644 index 0000000..8f29b67 --- /dev/null +++ b/internal/outbound/poison_bounce_test.go @@ -0,0 +1,62 @@ +package outbound + +import ( + "context" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "tidepool/internal/store" +) + +// The poison counter was the last outcome recorder left ungated by `applied`. +// +// MarkPoisoned carries the same (exists, applied) fence as every other terminal +// mark — a stale worker whose lease lapsed writes zero rows and is told so — and +// delivered, cancelled and parked all gate their metric on that answer. +// tidepool_outbound_poisoned did not: it bumped whether or not anything was +// poisoned. +// +// It is the counter operators are told to read as the queue's verdict (DEPLOY.md +// walks a redrive off it), and a dead-letter number that climbs while nothing is +// dead-lettered sends an incident to the wrong place — the more so because the +// row itself is fine here, so the miscount is the only trace the bounce leaves. +// Cf. TestWorker_BouncedParkIsNotCounted, which is the same rule one state over. +func TestWorker_BouncedPoisonIsNotCounted(t *testing.T) { + conn := workerTestDB(t) + ctx := context.Background() + seedWorkerActor(t, conn, true, false) + id := seedDelivery(t, conn, "Create", "", createPayload("x")) + deliveries := store.NewOutboundDeliveries(conn) + + // A worker claims the delivery and then wedges. + stale, err := deliveries.ClaimNext(ctx, time.Minute) + require.NoError(t, err) + require.NotNil(t, stale.ClaimedUntil) + + // Its lease lapses and a second worker takes the row. + _, err = conn.ExecContext(ctx, + `UPDATE outbound_deliveries SET claimed_until = now() - interval '1 minute' WHERE activity_id = $1`, id) + require.NoError(t, err) + reclaimed, err := deliveries.ClaimNext(ctx, time.Minute) + require.NoError(t, err) + require.Equal(t, 2, reclaimed.Attempts, "the second worker owns the claim now") + + // The wedged worker wakes up holding a token nobody honours and tries to + // dead-letter what it thinks is still its delivery. + w := newWorker(t, conn, &fakeSender{}, nil) + before := metricPoisoned.Value() + require.NoError(t, w.poison(ctx, stale, "4xx", "rejected", 400)) + + assert.Equal(t, before, metricPoisoned.Value(), + "a poison the fence REFUSED must not be counted: tidepool_outbound_poisoned is the "+ + "number an operator reads as the queue's verdict, and a stale worker bouncing off "+ + "a claim somebody else owns dead-lettered nothing") + + got := getDelivery(t, conn, id) + assert.Equal(t, store.DeliveryStatePending, got.State, + "and the row is untouched — the fence saw to that, which is exactly why the miscount is "+ + "the only trace the bounce leaves") +} diff --git a/internal/outbound/worker.go b/internal/outbound/worker.go index 3862440..3e646f8 100644 --- a/internal/outbound/worker.go +++ b/internal/outbound/worker.go @@ -255,12 +255,21 @@ func (w *Worker) handle(ctx context.Context, delivery *store.OutboundDelivery) e case causalEligible: // fall through to consent + delivery case causalWait: - // Held, NOT failed: keep it immediately re-eligible so it delivers the - // instant its parent is accepted, and do not advance the poison budget - // (the causal wait is wall-clock-bounded in causalStatus). - return w.parkCausal(ctx, delivery, "parent_pending", "waiting for bridge-origin parent to be accepted") + // Held, NOT failed: it becomes re-eligible one short interval from now, + // so it delivers within a beat of its parent being accepted, and the + // poison budget does not advance (the causal wait is wall-clock-bounded + // in causalStatus). + return w.parkCausal(ctx, delivery, "parent_pending", + "waiting for bridge-origin parent to be accepted", causalParkDelay) + case causalLookupFailed: + // The gate could not be READ. Holding is right — a store blip must not + // poison somebody's reply — but on the worker's real backoff, not the + // causal cadence: re-asking a database that just failed as fast as the + // round trip allows is how a blip becomes an outage. + return w.parkCausal(ctx, delivery, "parent_lookup_failed", + "causal parent lookup failed; holding for retry", w.backoff(delivery.Attempts)) // THE CLASS NAMES COME FROM store, and so do cross_authority and signer - // below. The reconciliation sweep excludes exactly these four from its + // below. The reconciliation sweep excludes exactly these from its // unknown-outcome report (store.neverReachedTheWireClasses): they carry // last_status_code 0 like a dial timeout does and mean the opposite — nothing // was sent, so the peer's state is not unknown, they simply do not have it. @@ -272,6 +281,9 @@ func (w *Worker) handle(ctx context.Context, delivery *store.OutboundDelivery) e case causalPoisonParent: return w.poison(ctx, delivery, store.PoisonClassParentPoisoned, "parent delivery poisoned; descendant cannot land", 0) + case causalPoisonParentCancelled: + return w.poison(ctx, delivery, store.PoisonClassParentCancelled, + "every delivery of the parent was cancelled; descendant cannot land", 0) } // Consent recheck (retraction asymmetry): a Delete/Undo always goes out — @@ -330,9 +342,7 @@ func (w *Worker) classify(ctx context.Context, delivery *store.OutboundDelivery, if stderrors.As(err, &he) { switch { case isDuplicate(he): - // Lemmy's received_activity dedupe (400 + "already received") is a - // SUCCESS by our stable id: a redelivery after a crash is expected. - return w.deliverSuccess(ctx, delivery, activity, he.StatusCode) + return w.duplicateDelivered(ctx, delivery, activity, he) case he.StatusCode == http.StatusNotFound || he.StatusCode == http.StatusGone: // 404/410 is an endpoint-GONE signal: re-resolve the inbox once @@ -355,6 +365,23 @@ func (w *Worker) classify(ctx context.Context, delivery *store.OutboundDelivery, return w.releaseOrPoison(ctx, delivery, "transport", err.Error(), 0) } +// duplicateDelivered records the received-activity dedupe as the success it is: +// Lemmy already holds this activity under our stable id, which is exactly what a +// redelivery after a crash is meant to discover. +// +// THE BODY IS LOGGED because this is the one branch that turns a 400 into a +// delivered row, and the delivered path stores no excerpt — so if the match ever +// fires on a rejection that merely resembles the dedupe, this line is the only +// record of what the peer actually said. Debug level: on a healthy bridge it +// fires only behind a crash-redelivery, and an operator chasing a wrong +// `delivered` is already turning the level up. +func (w *Worker) duplicateDelivered(ctx context.Context, delivery *store.OutboundDelivery, activity *store.OutboundActivity, he ap.HTTPError) error { + w.logger.Debug("peer reports the activity was already received; classifying as delivered", + "activity", delivery.ActivityID, "inbox", delivery.TargetInbox, + "status", he.StatusCode, "body", he.Body) + return w.deliverSuccess(ctx, delivery, activity, he.StatusCode) +} + // rotateInbox handles a 401/404/410: re-resolve the community's inbox ONCE // bypassing the cache (an endpoint rotation must not become a poison), retry, // then deliver-or-poison. @@ -374,7 +401,7 @@ func (w *Worker) rotateInbox(ctx context.Context, delivery *store.OutboundDelive } var he ap.HTTPError if stderrors.As(err, &he) && isDuplicate(he) { - return w.deliverSuccess(ctx, delivery, activity, he.StatusCode) + return w.duplicateDelivered(ctx, delivery, activity, he) } // Still bad after the single re-resolve: the endpoint is genuinely gone. status := first.StatusCode @@ -523,12 +550,27 @@ func (w *Worker) releaseOrPoison(ctx context.Context, delivery *store.OutboundDe return nil } -// poison marks the delivery permanently failed under its fencing token. +// poison marks the delivery permanently failed under its fencing token, and +// counts THE POISON THE FENCE ACCEPTED — the same rule countPark states for the +// park counter, one state over. +// +// MarkPoisoned carries the same (exists, applied) fence every other terminal +// mark does, and tidepool_outbound_poisoned is the number DEPLOY.md sends an +// operator to read as the queue's verdict (a redrive is decided off it). A stale +// worker whose lease lapsed poisons nothing — the fence says so — so counting +// its bounce puts a dead letter on the dashboard that does not exist anywhere in +// the table. The row is safe either way, which is precisely why the miscount +// would be the bounce's only trace, and why it is logged rather than passed over. func (w *Worker) poison(ctx context.Context, delivery *store.OutboundDelivery, class, excerpt string, status int) error { - _, _, err := w.deliveries.MarkPoisoned(ctx, delivery.ActivityID, delivery.TargetInbox, class, excerpt, status, *delivery.ClaimedUntil) + _, applied, err := w.deliveries.MarkPoisoned(ctx, delivery.ActivityID, delivery.TargetInbox, class, excerpt, status, *delivery.ClaimedUntil) if err != nil { return fmt.Errorf("poison delivery %s: %w", delivery.ActivityID, err) } + if !applied { + w.logger.Warn("poison did not apply: claim lost or row terminal", + "activity", delivery.ActivityID, "inbox", delivery.TargetInbox, "class", class) + return nil + } metricPoisoned.Add(1) return nil } @@ -601,18 +643,36 @@ func (w *Worker) countPark(delivery *store.OutboundDelivery, class string, appli metricParked.Add(1) } -// 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, 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. +// causalParkDelay is how long a causally-held child waits before it can be +// re-claimed. +// +// SHORT, BECAUSE THE POINT OF THE HOLD IS TO END. In practice the parent is a +// lower-seq delivery on the same serial line and goes out first, so a child +// rarely cycles here at all; a second is small enough that a reply lands within +// a beat of its parent being accepted. +// +// NONZERO, BECAUSE now() IS A SPIN. ReleaseParked leaves the row pending, and a +// row scheduled at now() is claimable by the very next ClaimNext — while +// DeliverNext returns worked=true, so Run never sleeps. Any wait the parent does +// not promptly end therefore becomes a claim/park loop running at whatever rate +// the round trip allows, for as long as the wait lasts: two UPDATEs and two or +// three SELECTs per iteration, against the database the rest of the bridge is +// sharing. The cancelled-parent verdict above removes the case that could run +// that loop for the full six-hour budget; this delay is what keeps any FUTURE +// wait path from reintroducing it. +const causalParkDelay = time.Second + +// parkCausal holds a causally-ineligible delivery for delay. Like park it never +// poisons, and like park it is attempt-neutral — which matters most here, since +// a child cycling through the hold would otherwise spend its whole retry budget +// waiting rather than trying. // // 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 { - _, applied, err := w.deliveries.ReleaseParked(ctx, delivery.ActivityID, delivery.TargetInbox, class, reason, 0, time.Now(), *delivery.ClaimedUntil) +func (w *Worker) parkCausal(ctx context.Context, delivery *store.OutboundDelivery, class, reason string, delay time.Duration) error { + next := time.Now().Add(delay) + _, applied, err := w.deliveries.ReleaseParked(ctx, delivery.ActivityID, delivery.TargetInbox, class, reason, 0, next, *delivery.ClaimedUntil) if err != nil { return fmt.Errorf("park (causal) delivery %s: %w", delivery.ActivityID, err) } @@ -626,8 +686,17 @@ type causalStatus int const ( causalEligible causalStatus = iota causalWait + // causalLookupFailed is a WAIT the store could not answer, kept apart from + // causalWait because the two want different cadences: an ordinary wait is + // ended by the parent landing (so it re-checks briskly), while a failed + // lookup is ended by the database recovering (so it backs off). + causalLookupFailed causalPoisonUnaccepted causalPoisonParent + // causalPoisonParentCancelled is the THIRD way a parent goes terminal, and + // the one the gate used to have no verdict for: cancelled. See + // store.PoisonClassParentCancelled. + causalPoisonParentCancelled ) func (w *Worker) causalStatus(ctx context.Context, delivery *store.OutboundDelivery, activity *store.OutboundActivity) causalStatus { @@ -643,22 +712,30 @@ func (w *Worker) causalStatus(ctx context.Context, delivery *store.OutboundDeliv } if err != nil { w.logger.Error("causal parent lookup failed", "parent", activity.ParentATURI, "error", err) - return causalWait // transient: hold rather than poison on a lookup blip + return causalLookupFailed // transient: hold rather than poison on a lookup blip } if parent.IsAccepted() { return causalEligible } - // Bridge-origin parent, not yet accepted. Poison ONLY if the child's ACTUAL - // parent delivery is poisoned (it will never land) — keyed on parent_at_uri, - // not seq-ancestry, so an unrelated poisoned row on the same line does not + // Bridge-origin parent, not yet accepted. Decide against the child ONLY on + // the child's ACTUAL parent deliveries — keyed on parent_at_uri, not + // seq-ancestry, so an unrelated poisoned row on the same line does not // poison this child. - poisoned, err := w.deliveries.ParentDeliveryPoisoned(ctx, activity.ParentATURI, delivery.TargetInbox) + // + // TWO TERMINAL PARENTS, NOT ONE. A poisoned parent will never land; so will + // one whose every delivery was CANCELLED, and that second case is what used + // to fall through to the wait below — a wait for something already decided + // against, which the delivery then spent its whole causal budget on. + disposition, err := w.deliveries.ParentDeliveryDisposition(ctx, activity.ParentATURI, delivery.TargetInbox) if err != nil { - w.logger.Error("parent-delivery poisoned check failed", "error", err) - return causalWait + w.logger.Error("parent-delivery disposition check failed", "error", err) + return causalLookupFailed } - if poisoned { + switch disposition { + case store.ParentDeliveryPoisoned: return causalPoisonParent + case store.ParentDeliveryCancelled: + return causalPoisonParentCancelled } // Otherwise the parent is merely pending: WAIT, bounded by WALL CLOCK from // the delivery's creation — never by the attempt count, so a parent that @@ -717,10 +794,37 @@ func (w *Worker) backoff(attempts int) time.Duration { // it — or the reverse. func isRetraction(kind string) bool { return slices.Contains(store.RetractionKinds, kind) } -// isDuplicate reports whether an HTTPError is Lemmy's duplicate-activity -// response (a 400 whose body reports the activity was already received). +// duplicateActivityPhrase is the received-activity dedupe, and it is the WHOLE +// phrase on purpose. +// +// This is the one classification that turns a rejection into a success, so it +// has to name the single rejection that IS one. Lemmy answers 400 with a family +// of slugs carrying the bare word "already" — banned_from_community, +// already_invalid, duplicate_title, "you have already been blocked" — and +// matching on that word alone read every one of them as a delivery: +// +// - a Like the peer refused flipped outbound_votes.delivered_state, which is +// an input to the score users are served and which nothing reconciles; +// - stampAccepted opened the causal gate for an object that never landed, so +// every child behind it was delivered into a rejection of its own; +// - and the delivered path stores no excerpt, so the body that would have +// shown an operator what happened was discarded. +// +// The phrase is what the peer actually says — the fake Lemmy in the outer +// acceptance test answers `{"error":"activity was already received"}`, which is +// the shape worker_test.go pins — and matching is done on a body lowercased with +// underscores folded to spaces, so a serialized enum (`already_received`) and +// the prose sentence are the same string here while none of the slugs above come +// near it. +const duplicateActivityPhrase = "already received" + +// isDuplicate reports whether an HTTPError is that dedupe response. func isDuplicate(he ap.HTTPError) bool { - return he.StatusCode == http.StatusBadRequest && strings.Contains(strings.ToLower(he.Body), "already") + if he.StatusCode != http.StatusBadRequest { + return false + } + normalized := strings.ReplaceAll(strings.ToLower(he.Body), "_", " ") + return strings.Contains(normalized, duplicateActivityPhrase) } // classForStatus labels a retryable HTTP status for the retry taxonomy. diff --git a/internal/store/divergence.go b/internal/store/divergence.go index fa729a2..ce78c6e 100644 --- a/internal/store/divergence.go +++ b/internal/store/divergence.go @@ -700,13 +700,13 @@ func (r *postgresDivergences) RecastDivergenceCount(ctx context.Context) (int, e // THE DECLARATION IS SHARED SO THE TWO SIDES CANNOT DRIFT IN SPELLING. The // worker used to write these as string literals at its own call sites while the // denylist below repeated them, and nothing compiled the two lists against each -// other: a respelling on either side would have silently moved four known +// other: a respelling on either side would have silently moved the known // non-deliveries into the unknown-outcome report, whose entire worth is that its // numbers stay small enough to trust. The compiler now refuses that. What it -// still cannot catch is a FIFTH never-wire class introduced as a fresh literal +// still cannot catch is a FURTHER never-wire class introduced as a fresh literal // — the worker's wire classes (transport, 4xx, 5xx …) are literals, so the // surrounding style invites one — which is why a new class belongs in this block -// and in neverReachedTheWireClasses, and why the store test that seeds all four +// and in neverReachedTheWireClasses, and why the store test that seeds them all // by name (divergence_unknown_test.go) is the pin on the set. const ( // PoisonClassParentUnaccepted: the causal wait budget expired and the parent @@ -715,6 +715,12 @@ const ( // PoisonClassParentPoisoned: the parent delivery poisoned; a descendant // cannot land, so it is not attempted. PoisonClassParentPoisoned = "parent_poisoned" + // PoisonClassParentCancelled: every delivery of the parent was CANCELLED — + // a consent recheck, an operator cancel, an actor or community sweep — so + // the parent is terminal, was never accepted, and nothing is left that could + // make it land. The descendant is decided here rather than left waiting out + // the causal budget against a parent that is never coming. + PoisonClassParentCancelled = "parent_cancelled" // PoisonClassCrossAuthority: the stored target inbox is not same-authority // with its community, and the worker refuses to sign a POST to it. PoisonClassCrossAuthority = "cross_authority" @@ -725,10 +731,10 @@ const ( PoisonClassSigner = "signer" ) -// neverReachedTheWireClasses are those four classes as the divergence sweep +// neverReachedTheWireClasses are those classes as the divergence sweep // reads them. // -// All four are written with last_status_code 0 — the same shape a dial timeout +// All of them are written with last_status_code 0 — the same shape a dial timeout // leaves behind, and the opposite meaning. Nothing was sent, so the peer's // state is not unknown at all: they do not have it. That is the rule that // already excludes a cancelled delivery, one step later in the worker. @@ -743,7 +749,7 @@ const ( // runs the other way — a new never-wire poison class must be added to the block // above AND to this list, or it inflates these counts. var neverReachedTheWireClasses = []string{ - PoisonClassParentUnaccepted, PoisonClassParentPoisoned, + PoisonClassParentUnaccepted, PoisonClassParentPoisoned, PoisonClassParentCancelled, PoisonClassCrossAuthority, PoisonClassSigner, } @@ -777,7 +783,7 @@ var neverReachedTheWireClasses = []string{ // selected as a boolean rather than left to the caller to infer, so that no // reader downstream can rediscover the wrong rule from LastStatusCode. // -// AND IT MUST HAVE BEEN SENT. Four poison classes are decided before any +// AND IT MUST HAVE BEEN SENT. Several poison classes are decided before any // POST — see neverReachedTheWireClasses — and they carry status 0 exactly // like a transport failure does. Filing them here would put a KNOWN // non-delivery in the bucket whose whole meaning is that the answer is diff --git a/internal/store/divergence_unknown_test.go b/internal/store/divergence_unknown_test.go index bf4ecff..d60dbb8 100644 --- a/internal/store/divergence_unknown_test.go +++ b/internal/store/divergence_unknown_test.go @@ -244,10 +244,11 @@ func TestUnknownDeliveryOutcomes_ACancelledDeliveryIsNotUnknown(t *testing.T) { // completeness half of this class, and it is the same rule that excludes a // cancelled delivery — one step later in the worker. // -// Four poison classes are decided BEFORE any POST is made (worker.go): a causal -// wait that expired (parent_unaccepted), a parent that poisoned -// (parent_poisoned), an inbox that is not same-authority with its community -// (cross_authority), and a signer that would not resolve (signer). All four +// Several poison classes are decided BEFORE any POST is made (worker.go): a +// causal wait that expired (parent_unaccepted), a parent that poisoned +// (parent_poisoned), a parent whose every delivery was cancelled +// (parent_cancelled), an inbox that is not same-authority with its community +// (cross_authority), and a signer that would not resolve (signer). All of them // record status 0, which is the same shape a dial timeout leaves behind — and // they mean the opposite. Nothing was ever sent, so the peer's state is not // unknown at all: they do not have it. @@ -260,8 +261,10 @@ func TestUnknownDeliveryOutcomes_ACancelledDeliveryIsNotUnknown(t *testing.T) { func TestUnknownDeliveryOutcomes_APoisonThatNeverReachedTheWireIsNotUnknown(t *testing.T) { database := acceptanceTestDB(t) - // The four never-wire classes, exactly as the worker writes them. - neverSent := []string{"parent_unaccepted", "parent_poisoned", "cross_authority", "signer"} + // Every never-wire class, exactly as the worker writes them. + neverSent := []string{ + "parent_unaccepted", "parent_poisoned", "parent_cancelled", "cross_authority", "signer", + } for _, class := range neverSent { seedPoisonedDelivery(t, database, "https://coves.social/ap/activity/dv-neverwire-"+class, @@ -285,7 +288,8 @@ func TestUnknownDeliveryOutcomes_APoisonThatNeverReachedTheWireIsNotUnknown(t *t } assert.ElementsMatch(t, []string{sentSilently, sentRefused}, ids, "only deliveries that REACHED THE WIRE are unknown. %v are all decided before any POST "+ - "is made — a causal wait that expired, a poisoned parent, a cross-authority inbox, "+ + "is made — a causal wait that expired, a poisoned or wholly cancelled parent, a "+ + "cross-authority inbox, "+ "an unresolvable signer — so the peer does not have them and nothing about their "+ "outcome is uncertain. They carry status 0 like a dial timeout does, which is the "+ "whole trap: identical column, opposite meaning. Counting them inflates the one "+ diff --git a/internal/store/interfaces.go b/internal/store/interfaces.go index 8ea77fe..a2e596a 100644 --- a/internal/store/interfaces.go +++ b/internal/store/interfaces.go @@ -817,14 +817,19 @@ type OutboundDeliveries interface { // implementation for why, and for the instances it cannot reach). DistinctInboxesForActor(ctx context.Context, actorDID string) ([]DeliveryTarget, error) - // ParentDeliveryPoisoned reports whether the delivery of the child's ACTUAL - // parent (the activity that federated parentATURI as its object, to the same - // inbox) is poisoned — the causal signal task 15's worker reads to poison a - // child whose bridge-origin parent will NEVER land (parent_poisoned), as - // distinct from one merely waiting for a pending parent (parent_unaccepted). + // ParentDeliveryDisposition reports what has become of the deliveries of the + // child's ACTUAL parent (the activities that federated parentATURI as their + // object, to the same inbox) — the causal signal task 15's worker reads to + // decide a child whose bridge-origin parent will NEVER land, as distinct + // from one merely waiting for a pending parent (parent_unaccepted). + // // Keyed on the parent's object id, NOT on seq-ancestry, so an unrelated - // poisoned row on the same serial line does not poison the child. - ParentDeliveryPoisoned(ctx context.Context, parentATURI, targetInbox string) (bool, error) + // poisoned row on the same serial line does not poison the child. Poisoned + // is reported if ANY delivery poisoned; cancelled only if EVERY delivery was + // cancelled and there is at least one, so a pending sibling keeps the parent + // open and a parent with no deliveries at all (fediverse-origin) never reads + // as decided against. + ParentDeliveryDisposition(ctx context.Context, parentATURI, targetInbox string) (ParentDeliveryDisposition, error) // CancelClaimed cancels a SINGLE claimed delivery under its fencing token // (the consent-block outcome for one create/update), leaving the actor's diff --git a/internal/store/outbound_deliveries.go b/internal/store/outbound_deliveries.go index 5362cbd..69344ee 100644 --- a/internal/store/outbound_deliveries.go +++ b/internal/store/outbound_deliveries.go @@ -595,27 +595,76 @@ func (r *postgresOutboundDeliveries) Get(ctx context.Context, activityID, target return delivery, nil } -func (r *postgresOutboundDeliveries) ParentDeliveryPoisoned(ctx context.Context, parentATURI, targetInbox string) (bool, error) { +// ParentDeliveryDisposition is what has become of the deliveries carrying a +// child's ACTUAL parent to one inbox — the causal gate's read on whether the +// parent can still land. +// +// THREE VALUES BECAUSE THE PARENT HAS THREE ENDINGS, and the third is the one a +// boolean could not express. A delivery goes terminal three ways, and the gate +// used to ask only "is it poisoned?": a CANCELLED parent — the consent recheck +// at claim time, an operator cancel, an actor or community sweep — answered +// false, left the pending index, and never got accepted_at, so its child waited +// on something that was never coming. +type ParentDeliveryDisposition string + +const ( + // ParentDeliveryOpen means nothing has decided against the parent: it is + // still pending, already delivered, or has no delivery row at all. The child + // waits. + ParentDeliveryOpen ParentDeliveryDisposition = "open" + // ParentDeliveryPoisoned means a delivery of the parent poisoned. It + // outranks every other reading — a poisoned parent is the loudest verdict + // available and the one that predates this type. + ParentDeliveryPoisoned ParentDeliveryDisposition = "poisoned" + // ParentDeliveryCancelled means the parent HAS deliveries and every one of + // them was cancelled: nothing is left that could make it land. + ParentDeliveryCancelled ParentDeliveryDisposition = "cancelled" +) + +func (r *postgresOutboundDeliveries) ParentDeliveryDisposition(ctx context.Context, parentATURI, targetInbox string) (ParentDeliveryDisposition, error) { // The parent's delivery is the one whose activity federated parentATURI as // its object: the activity payload's object.id is the served object URL, // which ends in "/ap/object///" — exactly the // at-uri's three parts. Match on that suffix so we need no origin here (and // DIDs/NSIDs/TIDs carry no LIKE metacharacters). + // + // ONE OBJECT CAN HAVE SEVERAL DELIVERIES to the same inbox — a Create and + // every later Update carry the same object.id — so both readings are + // aggregates over the whole set, and they are deliberately asymmetric: + // + // poisoned — ANY row. A poisoned delivery is a failure the peer may + // already have half-seen, and the pre-existing contract is that + // it condemns the descendant. + // cancelled — EVERY row, and at least one. A cancel is a decision about one + // activity, not about the object: while a sibling is still + // pending, something can yet make the parent land, and poisoning + // the child out from under it would be a wrong answer arrived at + // early. The "at least one" guard is what keeps the vacuous + // all-of-nothing from reading as cancelled — a parent with NO + // deliveries is a fediverse-origin object, whose children are + // always eligible. suffix := strings.TrimPrefix(parentATURI, "at://") - var exists bool + var poisoned, allCancelled bool err := r.db.QueryRowContext(ctx, ` - SELECT EXISTS ( - SELECT 1 - FROM outbound_deliveries d - JOIN outbound_activities a ON a.activity_id = d.activity_id - WHERE d.state = 'poisoned' - AND d.target_inbox = $2 - AND a.payload -> 'object' ->> 'id' LIKE '%/ap/object/' || $1)`, - suffix, targetInbox).Scan(&exists) + SELECT + COUNT(*) FILTER (WHERE d.state = 'poisoned') > 0, + COUNT(*) > 0 AND COUNT(*) FILTER (WHERE d.state <> 'cancelled') = 0 + FROM outbound_deliveries d + JOIN outbound_activities a ON a.activity_id = d.activity_id + WHERE d.target_inbox = $2 + AND a.payload -> 'object' ->> 'id' LIKE '%/ap/object/' || $1`, + suffix, targetInbox).Scan(&poisoned, &allCancelled) if err != nil { - return false, fmt.Errorf("check parent delivery poisoned for %q: %w", parentATURI, err) + return "", fmt.Errorf("read parent delivery disposition for %q: %w", parentATURI, err) + } + switch { + case poisoned: + return ParentDeliveryPoisoned, nil + case allCancelled: + return ParentDeliveryCancelled, nil + default: + return ParentDeliveryOpen, nil } - return exists, nil } func (r *postgresOutboundDeliveries) CancelClaimed(ctx context.Context, activityID, targetInbox string, claimToken time.Time) (bool, bool, error) { -- 2.51.2