diff --git a/cmd/tidepool/main.go b/cmd/tidepool/main.go
index b96bf9a..dea0b67 100644
--- a/cmd/tidepool/main.go
+++ b/cmd/tidepool/main.go
@@ -580,12 +580,17 @@ func run(logger *slog.Logger) error {
// loop. The returned channel closes when the connector's loop has exited, so
// shutdown can wait for the final cursor flush instead of racing it.
//
-// The seams that are not wired yet are nil ON PURPOSE, and each is a no-op the
-// consumer announces rather than a silent gap:
+// Outbound delivery (task 15) is fully wired here:
+//
+// - Enqueuer: the real persisting enqueuer is wired whenever CONSUMER_ENABLED
+// (this function only runs then), so intents past the rev gate always
+// persist to outbound_activities/deliveries. OUTBOUND_WORKERS>0 additionally
+// starts the delivery workers that POST them; the noop enqueuer is used only
+// when the consumer is disabled (this function does not run at all).
+//
+// The task-16/17 seams are still nil ON PURPOSE, each a no-op the consumer
+// announces rather than a silent gap:
//
-// - Enqueuer is the logging no-op until task 15's delivery queue lands. The
-// consumer still runs behind it, so the cursor, the rev gate and the
-// outbound state that delivery will be built FROM are all exercised.
// - Engine (task 16) nil means postv2 events are skipped at debug.
// - RemoteDeleter (task 17) nil means a deleteRemote opt-out is recorded and
// logged rather than acted on.
@@ -617,28 +622,28 @@ func startConsumer(
return nil, fmt.Errorf("consumer: handle resolver: %w", err)
}
- // The outbound delivery pipe (task 15). The noop enqueuer is the default —
- // the consumer still writes durable outbound state, but nothing federates
- // — and it is swapped for the real enqueuer ONLY when OUTBOUND_WORKERS>0.
- // This is the gate that keeps a not-yet-wired deployment (and the e2e,
- // which runs with workers=0) from sending anything.
+ // The outbound delivery pipe (task 15). Because this function runs ONLY when
+ // the consumer is enabled, the REAL persisting enqueuer is always wired: it
+ // writes outbound_activities/deliveries inside the consumer's gate tx, so an
+ // intent past the gate is never dropped. OUTBOUND_WORKERS gates only whether
+ // the delivery WORKER goroutines run — with workers=0, state accumulates but
+ // nothing is POSTed. (The noop enqueuer is reserved for the consumer-disabled
+ // path, where nothing runs at all.)
inboxes := outbound.NewInboxResolver(apClient, outboundInboxTTL)
- var enqueuer consume.OutboundEnqueuer = consume.NewNoopEnqueuer(logger)
+ enqueuer, err := outbound.NewEnqueuer(outbound.EnqueuerOptions{
+ DB: database,
+ Translator: outbound.NewTranslator(cfg.APUserOrigin),
+ Inboxes: inboxes,
+ Actors: store.NewAPActors(database),
+ UserOrigin: cfg.APUserOrigin,
+ Logger: logger,
+ })
+ if err != nil {
+ return nil, fmt.Errorf("consumer: outbound enqueuer: %w", err)
+ }
+
var worker *outbound.Worker
if cfg.OutboundWorkers > 0 {
- realEnqueuer, err := outbound.NewEnqueuer(outbound.EnqueuerOptions{
- DB: database,
- Translator: outbound.NewTranslator(cfg.APUserOrigin),
- Inboxes: inboxes,
- Actors: store.NewAPActors(database),
- UserOrigin: cfg.APUserOrigin,
- Logger: logger,
- })
- if err != nil {
- return nil, fmt.Errorf("consumer: outbound enqueuer: %w", err)
- }
- enqueuer = realEnqueuer
-
worker, err = outbound.NewWorker(outbound.WorkerOptions{
DB: database,
Actors: store.NewAPActors(database),
diff --git a/internal/ap/client.go b/internal/ap/client.go
index ea4132e..5e94b1a 100644
--- a/internal/ap/client.go
+++ b/internal/ap/client.go
@@ -649,10 +649,13 @@ func (c *Client) SendActivityAs(ctx context.Context, signer *Signer, inboxURL st
return HTTPError{URL: inboxURL, StatusCode: resp.StatusCode, Body: string(body)}
}
-// sendActivityWith is the shared POST loop: it marshals the activity once,
+// sendActivityWith is the RETRYING POST loop: it marshals the activity once,
// then retries the signed POST under the client's backoff / egress guard,
-// signing each attempt with the given signer. SendActivity supplies the
-// service actor's configured signer; SendActivityAs supplies a per-actor one.
+// signing each attempt with the given signer. Its only caller is SendActivity
+// (the service actor's Follow/Undo), which supplies the configured signer.
+// SendActivityAs deliberately does NOT use this loop — it is a standalone
+// single-shot POST, because the task-15 delivery worker owns retry and needs to
+// see each response to classify it.
func (c *Client) sendActivityWith(ctx context.Context, signer *Signer, inboxURL string, activity any) error {
payload, err := json.Marshal(activity)
if err != nil {
diff --git a/internal/apobject/object.go b/internal/apobject/object.go
index 9983bba..baa6a37 100644
--- a/internal/apobject/object.go
+++ b/internal/apobject/object.go
@@ -95,6 +95,11 @@ func BuildPage(actorID, communityAPID, objectURL string, record map[string]any)
// its blobs live in the author's PDS and need a blob→PDS-URL seam to render as
// attachment [{type:Image, url}] — a task follow-up. A post with no external
// embed carries no attachment.
+//
+// The href is scheme-checked: a javascript:/data:/vbscript:/file: uri rendered
+// as a clickable Link would be a stored-XSS-shaped hazard on every peer that
+// renders it, so an unsafe scheme drops the attachment (fail closed) rather than
+// federating it.
func linkAttachment(record map[string]any) []any {
embed, ok := record["embed"].(map[string]any)
if !ok {
@@ -108,12 +113,21 @@ func linkAttachment(record map[string]any) []any {
return nil
}
href, _ := external["uri"].(string)
- if href == "" {
+ if href == "" || !isSafeLinkScheme(href) {
return nil
}
return []any{map[string]any{"type": "Link", "href": href}}
}
+// isSafeLinkScheme restricts a bridge-authored clickable URI to http/https. The
+// lexicon's format:"uri" accepts javascript:/data:/vbscript:/file:, which a
+// downstream client rendering the remote-controlled link as clickable would
+// treat as a scripting or local-file URI.
+func isSafeLinkScheme(uri string) bool {
+ lower := strings.ToLower(strings.TrimSpace(uri))
+ return strings.HasPrefix(lower, "http://") || strings.HasPrefix(lower, "https://")
+}
+
// hasNSFWLabel reports whether the record self-labels nsfw
// (com.atproto.label.defs#selfLabels with a value of "nsfw").
func hasNSFWLabel(record map[string]any) bool {
diff --git a/internal/apobject/object_test.go b/internal/apobject/object_test.go
new file mode 100644
index 0000000..a015223
--- /dev/null
+++ b/internal/apobject/object_test.go
@@ -0,0 +1,55 @@
+package apobject
+
+import (
+ "testing"
+
+ "github.com/stretchr/testify/assert"
+ "github.com/stretchr/testify/require"
+)
+
+// Second-opinion (important): a post's external embed href must pass a scheme
+// allowlist before it is rendered as a Link attachment. A javascript:/data: uri
+// federated as a clickable Link is a stored-XSS-shaped hazard on every peer that
+// renders it.
+
+func buildPageWithEmbedURI(t *testing.T, uri string) map[string]any {
+ t.Helper()
+ page, err := BuildPage(
+ "https://coves.social/ap/actor/did:plc:x",
+ "https://lemmy.world/c/tech",
+ "https://coves.social/ap/object/did:plc:x/social.coves.community.postv2/rk",
+ map[string]any{
+ "title": "a post",
+ "embed": map[string]any{
+ "$type": "social.coves.embed.external",
+ "external": map[string]any{"uri": uri},
+ },
+ })
+ require.NoError(t, err)
+ return page
+}
+
+func TestBuildPage_RejectsUnsafeEmbedSchemes(t *testing.T) {
+ for _, uri := range []string{
+ "javascript:alert(1)",
+ "data:text/html,",
+ "vbscript:msgbox(1)",
+ "file:///etc/passwd",
+ } {
+ page := buildPageWithEmbedURI(t, uri)
+ _, has := page["attachment"]
+ assert.Falsef(t, has,
+ "an unsafe embed scheme (%s) must NOT be rendered as a Link attachment", uri)
+ }
+}
+
+func TestBuildPage_KeepsSafeEmbedLink(t *testing.T) {
+ page := buildPageWithEmbedURI(t, "https://example.com/article")
+ attach, ok := page["attachment"].([]any)
+ require.True(t, ok, "a safe https link is rendered as an attachment")
+ require.NotEmpty(t, attach)
+ link, ok := attach[0].(map[string]any)
+ require.True(t, ok)
+ assert.Equal(t, "Link", link["type"])
+ assert.Equal(t, "https://example.com/article", link["href"])
+}
diff --git a/internal/config/config.go b/internal/config/config.go
index c4f82f1..14c668b 100644
--- a/internal/config/config.go
+++ b/internal/config/config.go
@@ -482,7 +482,11 @@ func Load(logger *slog.Logger) (*Config, error) {
if err != nil {
return nil, err
}
- cfg.OutboundDisabledHosts = parseSet(os.Getenv("OUTBOUND_DISABLED_HOSTS"))
+ // Hosts are lowercased (the worker's scope host comes from hostOf, which
+ // lowercases): a kill switch must fail CLOSED, so a mixed-case entry has to
+ // still block the normalized host. Communities and actors are exact ids and
+ // keep their case.
+ cfg.OutboundDisabledHosts = parseHostSet(os.Getenv("OUTBOUND_DISABLED_HOSTS"))
cfg.OutboundDisabledCommunities = parseSet(os.Getenv("OUTBOUND_DISABLED_COMMUNITIES"))
cfg.OutboundDisabledActors = parseSet(os.Getenv("OUTBOUND_DISABLED_ACTORS"))
@@ -667,6 +671,18 @@ func parseSet(raw string) map[string]struct{} {
return set
}
+// parseHostSet is parseSet with each entry lowercased — for the host kill
+// switch, whose scope host arrives already lowercased.
+func parseHostSet(raw string) map[string]struct{} {
+ set := make(map[string]struct{})
+ for _, item := range strings.Split(raw, ",") {
+ if item = strings.ToLower(strings.TrimSpace(item)); item != "" {
+ set[item] = struct{}{}
+ }
+ }
+ return set
+}
+
// boolVar reports whether an environment variable is set to a truthy value
// ("1", "true", "yes", case-insensitive).
func boolVar(name string) bool {
diff --git a/internal/consume/dispatch.go b/internal/consume/dispatch.go
index ef4658e..9fc7f3d 100644
--- a/internal/consume/dispatch.go
+++ b/internal/consume/dispatch.go
@@ -105,8 +105,9 @@ type PostIntent struct {
// ActivityID reports the deterministic activity id.
func (i PostIntent) ActivityID() string { return i.ID }
-// OutboundEnqueuer is the task 15 seam. main.go wires a logging noop until
-// task 15 swaps in the real delivery queue.
+// OutboundEnqueuer is the task 15 seam. main wires the real persisting enqueuer
+// (outbound.Enqueuer, which writes outbound_activities/deliveries on the gate
+// tx) whenever CONSUMER_ENABLED; the noop is the consumer-disabled default.
type OutboundEnqueuer interface {
// EnqueueActivity hands one intent to delivery ON THE CALLER'S TX — the
// enqueue must commit with the rev-gate advance the consumer is holding, or
diff --git a/internal/consume/noop.go b/internal/consume/noop.go
index 4b08770..657a828 100644
--- a/internal/consume/noop.go
+++ b/internal/consume/noop.go
@@ -7,10 +7,10 @@ import (
)
// noopEnqueuer logs the outbound intents it is handed and delivers nothing.
-// It is what main wires until task 15 lands, following the v1 precedent
-// (ingest.NewNoopVotes): the consumer runs, writes its state, and makes the
-// work it WOULD deliver visible, rather than being disabled entirely and
-// leaving the whole path unexercised until delivery exists.
+// It is the consumer-DISABLED default: it drops intents by design so a path that
+// is not running end to end makes the work it WOULD deliver visible in the log
+// rather than accumulating it. It is NOT a pre-task-15 placeholder — whenever
+// CONSUMER_ENABLED, main wires the real persisting outbound.Enqueuer instead.
type noopEnqueuer struct {
logger *slog.Logger
}
diff --git a/internal/db/migrations/020_outbound_delivery.sql b/internal/db/migrations/020_outbound_delivery.sql
index b993fc9..0b3b7b0 100644
--- a/internal/db/migrations/020_outbound_delivery.sql
+++ b/internal/db/migrations/020_outbound_delivery.sql
@@ -25,9 +25,9 @@ CREATE TABLE outbound_activities (
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
--- The causal-gating lookup: a delivery whose activity has a parent_at_uri is
--- ineligible until the parent's own outbound_objects row is accepted_at (a
--- BRIDGE-origin parent) — task 15's worker reads the parent chain by actor.
+-- Supports CancelForActor's actor_did filter (the consent-withdrawal /
+-- kill-switch sweep that parks a disabled or paused actor's pending outbound
+-- work): it scans activities by actor to cancel their deliveries.
CREATE INDEX idx_outbound_activities_actor ON outbound_activities (actor_did);
-- outbound_deliveries is one delivery attempt per (activity, target inbox). It
diff --git a/internal/ingest/follow.go b/internal/ingest/follow.go
index c3b5d14..a298616 100644
--- a/internal/ingest/follow.go
+++ b/internal/ingest/follow.go
@@ -179,18 +179,27 @@ type outboundMutateRequest struct {
Activity string `json:"activity"`
Community string `json:"community"`
Actor string `json:"actor"`
+ All bool `json:"all"`
}
-// handleOutboundRedrive resets poisoned deliveries back to pending, optionally
-// scoped to one activity or community.
+// handleOutboundRedrive resets poisoned deliveries back to pending, scoped to
+// one activity or community. An unscoped redrive (no filter) would re-attempt
+// EVERY poisoned delivery at once, so it is refused unless the caller opts in
+// explicitly with {"all":true} — a malformed body is a 400, never a silent
+// fleet-wide redrive.
func (a *Admin) handleOutboundRedrive(w http.ResponseWriter, r *http.Request) {
if a.deliveries == nil {
http.Error(w, "outbound delivery is not configured", http.StatusNotImplemented)
return
}
var req outboundMutateRequest
- if r.Body != nil {
- _ = json.NewDecoder(r.Body).Decode(&req)
+ if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
+ http.Error(w, `body must be JSON: {"activity":"..."} | {"community":"..."} | {"all":true}`, http.StatusBadRequest)
+ return
+ }
+ if req.Activity == "" && req.Community == "" && !req.All {
+ http.Error(w, `refusing an unscoped redrive: set {"activity":"..."}, {"community":"..."}, or {"all":true}`, http.StatusBadRequest)
+ return
}
redriven, err := a.deliveries.RedrivePoisoned(r.Context(), req.Activity, req.Community)
if err != nil {
diff --git a/internal/outbound/atomicity_test.go b/internal/outbound/atomicity_test.go
new file mode 100644
index 0000000..33b0d2a
--- /dev/null
+++ b/internal/outbound/atomicity_test.go
@@ -0,0 +1,114 @@
+package outbound
+
+import (
+ "context"
+ "database/sql"
+ "fmt"
+ "testing"
+
+ "github.com/stretchr/testify/assert"
+ "github.com/stretchr/testify/require"
+
+ "tidepool/internal/store"
+)
+
+// Second-opinion H5: the causal forward edge (a bridge-origin parent's delivery
+// stamps accepted_at, releasing its held children) must be driven end-to-end,
+// and the stamp must be ATOMIC with MarkDelivered — a delivered-but-unaccepted
+// parent strands every child forever.
+
+// pageParentPayload is a Create{Page} whose object.id maps back to atURI, so
+// stampAccepted resolves the object it federated.
+func pageParentPayload(atURI string) []byte {
+ return []byte(fmt.Sprintf(
+ `{"id":"https://coves.social/ap/activity/parent","type":"Create",`+
+ `"object":{"type":"Page","id":%q}}`, objectURLFor(atURI)))
+}
+
+func seedParentDelivery(t *testing.T, conn *sql.DB, payload []byte) string {
+ t.Helper()
+ ctx := context.Background()
+ activityID := "https://coves.social/ap/activity/" + repeatHex64("parentpage")
+ _, err := store.NewOutboundActivities(conn).Insert(ctx, store.OutboundActivity{
+ ActivityID: activityID, ActorDID: wActorDID, Kind: "Create", Payload: payload,
+ })
+ require.NoError(t, err)
+ _, err = store.NewOutboundDeliveries(conn).Enqueue(ctx, store.OutboundDelivery{
+ ActivityID: activityID, TargetInbox: wInbox, OrderingKey: wCommunityAPID,
+ })
+ require.NoError(t, err)
+ return activityID
+}
+
+func rawAccepted(t *testing.T, conn *sql.DB, atURI string) bool {
+ t.Helper()
+ var acceptedAt sql.NullTime
+ require.NoError(t, conn.QueryRowContext(context.Background(),
+ `SELECT accepted_at FROM outbound_objects WHERE at_uri = $1`, atURI).Scan(&acceptedAt))
+ return acceptedAt.Valid
+}
+
+func TestWorker_ParentDeliverySuccessAcceptsAndReleasesHeldChild(t *testing.T) {
+ conn := workerTestDB(t)
+ ctx := context.Background()
+ seedWorkerActor(t, conn, true, false)
+ seedBridgeParent(t, conn, false) // parent object, accepted_at NULL
+
+ seedParentDelivery(t, conn, pageParentPayload(gParentATURI)) // seq 1: the parent
+ childID := seedDelivery(t, conn, "Create", gParentATURI, createPayload("child")) // seq 2: held child
+
+ w := newWorker(t, conn, &fakeSender{}, nil)
+
+ // 1) Deliver the parent (the head). This must stamp accepted_at.
+ worked, err := w.DeliverNext(ctx)
+ require.NoError(t, err)
+ require.True(t, worked)
+ assert.True(t, rawAccepted(t, conn, gParentATURI),
+ "a successful bridge-origin parent delivery stamps its outbound_objects.accepted_at")
+
+ // 2) The previously-held child is now eligible and delivers — the full
+ // forward edge the test-analyzer flagged as untested.
+ worked, err = w.DeliverNext(ctx)
+ require.NoError(t, err)
+ require.True(t, worked)
+ assert.Equal(t, store.DeliveryStateDelivered, getDelivery(t, conn, childID).State,
+ "once the parent is accepted, the held child delivers end-to-end")
+}
+
+// faultObjects fails SetAccepted to prove the stamp rides the SAME committed
+// transition as MarkDelivered.
+type faultObjects struct {
+ store.OutboundObjects
+ failSetAccepted bool
+}
+
+func (o *faultObjects) SetAccepted(ctx context.Context, atURI string) error {
+ if o.failSetAccepted {
+ return fmt.Errorf("injected: SetAccepted failed")
+ }
+ return o.OutboundObjects.SetAccepted(ctx, atURI)
+}
+
+func TestWorker_AcceptedStampIsAtomicWithMarkDelivered(t *testing.T) {
+ conn := workerTestDB(t)
+ ctx := context.Background()
+ seedWorkerActor(t, conn, true, false)
+ seedBridgeParent(t, conn, false)
+ parentID := seedParentDelivery(t, conn, pageParentPayload(gParentATURI))
+
+ objects := &faultObjects{OutboundObjects: store.NewOutboundObjects(conn), failSetAccepted: true}
+ w := newWorker(t, conn, &fakeSender{}, func(o *WorkerOptions) { o.Objects = objects })
+
+ // The POST succeeds but the accepted_at stamp fails.
+ _, _ = w.DeliverNext(ctx)
+
+ // Consistency invariant: the delivery must NEVER be 'delivered' while the
+ // object is unaccepted. If the stamp cannot commit, the delivered mark must
+ // roll back with it, so a retry re-runs both — otherwise the parent is
+ // delivered-but-unaccepted and every child is stranded forever.
+ deliveredButUnaccepted := getDelivery(t, conn, parentID).State == store.DeliveryStateDelivered &&
+ !rawAccepted(t, conn, gParentATURI)
+ assert.False(t, deliveredButUnaccepted,
+ "MarkDelivered and SetAccepted must be atomic: a parent must never be delivered while "+
+ "its accepted_at is unset (that strands every child)")
+}
diff --git a/internal/outbound/causal_gating_test.go b/internal/outbound/causal_gating_test.go
index 8eb2401..5478f00 100644
--- a/internal/outbound/causal_gating_test.go
+++ b/internal/outbound/causal_gating_test.go
@@ -100,60 +100,6 @@ func TestCausalGating_FediverseParentIsAlwaysEligible(t *testing.T) {
assert.Equal(t, store.DeliveryStateDelivered, getDelivery(t, conn, id).State)
}
-func TestCausalGating_BoundedWaitPoisonsParentUnaccepted(t *testing.T) {
- conn := workerTestDB(t)
- seedWorkerActor(t, conn, true, false)
- seedBridgeParent(t, conn, false) // never accepted
- id := seedDelivery(t, conn, "Create", gParentATURI, createPayload("x"))
- // The bounded wait is exhausted (attempts at the cap): a comment on a
- // never-accepted post must not wait forever.
- setAttempts(t, conn, id, 3) // MaxAttempts is 3
-
- w := newWorker(t, conn, &fakeSender{}, nil)
- _, err := w.DeliverNext(context.Background())
- require.NoError(t, err)
-
- d := getDelivery(t, conn, id)
- assert.Equal(t, store.DeliveryStatePoisoned, d.State,
- "after the bounded wait, an unaccepted parent poisons the child")
- assert.Contains(t, d.LastErrorClass, "parent_unaccepted",
- "the poison reason is queryable: parent_unaccepted (distinct from a delivery failure)")
-}
-
-func TestCausalGating_PoisonedParentPoisonsChild(t *testing.T) {
- conn := workerTestDB(t)
- seedWorkerActor(t, conn, true, false)
- seedBridgeParent(t, conn, false)
-
- // The parent's own delivery is POISONED. A descendant must not wait forever
- // for a parent that will never land — it poisons with a DISTINCT reason.
- parentActivityID := "https://coves.social/ap/activity/" + repeatHex64("parent")
- _, err := store.NewOutboundActivities(conn).Insert(context.Background(), store.OutboundActivity{
- ActivityID: parentActivityID,
- ActorDID: wActorDID,
- Kind: "Create",
- Payload: createPayload(parentActivityID),
- })
- require.NoError(t, err)
- _, err = store.NewOutboundDeliveries(conn).Enqueue(context.Background(), store.OutboundDelivery{
- ActivityID: parentActivityID,
- TargetInbox: wInbox,
- OrderingKey: wCommunityAPID,
- })
- require.NoError(t, err)
- _, err = conn.ExecContext(context.Background(),
- `UPDATE outbound_deliveries SET state='poisoned' WHERE activity_id=$1`, parentActivityID)
- require.NoError(t, err)
-
- id := seedDelivery(t, conn, "Create", gParentATURI, createPayload("child"))
-
- w := newWorker(t, conn, &fakeSender{}, nil)
- _, err = w.DeliverNext(context.Background())
- require.NoError(t, err)
-
- d := getDelivery(t, conn, id)
- assert.Equal(t, store.DeliveryStatePoisoned, d.State,
- "a poisoned parent poisons its descendants")
- assert.Contains(t, d.LastErrorClass, "parent_poisoned",
- "the reason is parent_poisoned — distinct and queryable from parent_unaccepted")
-}
+// Bounded-wait and poisoned-parent gating moved to causal_hardening_test.go —
+// they are now TIME-bounded (not attempt-bounded, H4b) and keyed on the child's
+// ACTUAL parent (not any lower-seq poison on the line, H6).
diff --git a/internal/outbound/causal_hardening_test.go b/internal/outbound/causal_hardening_test.go
new file mode 100644
index 0000000..9b29f0a
--- /dev/null
+++ b/internal/outbound/causal_hardening_test.go
@@ -0,0 +1,148 @@
+package outbound
+
+import (
+ "context"
+ "fmt"
+ "strings"
+ "testing"
+ "time"
+
+ "github.com/stretchr/testify/assert"
+ "github.com/stretchr/testify/require"
+
+ "tidepool/internal/store"
+)
+
+// Second-opinion H4b + H6: causal gating must be TIME-bounded (not attempt-
+// bounded) and keyed on the child's ACTUAL parent (not any lower-seq poison on
+// the serial line).
+
+func objectURLFor(atURI string) string {
+ return "https://coves.social/ap/object/" + strings.TrimPrefix(atURI, "at://")
+}
+
+// ---------------------------------------------------------------------------
+// H4b: causal wait is wall-clock-bounded, not attempt-bounded
+// ---------------------------------------------------------------------------
+
+func TestCausalGating_ManyAttemptsWithinBudgetStayPending(t *testing.T) {
+ conn := workerTestDB(t)
+ seedWorkerActor(t, conn, true, false)
+ seedBridgeParent(t, conn, false) // parent never accepted
+ id := seedDelivery(t, conn, "Create", gParentATURI, createPayload("x"))
+
+ // Attempts far past the delivery-failure cap, but the delivery was created
+ // moments ago and the causal budget is an hour: a parent legitimately taking
+ // minutes to be admitted must NOT be poisoned just because the child was
+ // claimed a few times. The attempts counter must NOT drive the causal wait.
+ setAttempts(t, conn, id, 99)
+
+ w := newWorker(t, conn, &fakeSender{}, nil) // CausalWaitBudget defaults to 1h in the helper
+ _, err := w.DeliverNext(context.Background())
+ require.NoError(t, err)
+
+ assert.Equal(t, store.DeliveryStatePending, getDelivery(t, conn, id).State,
+ "a causal wait is bounded by WALL CLOCK, not attempts — a recent delivery stays pending "+
+ "no matter how many times it was claimed")
+}
+
+func TestCausalGating_PoisonsOnlyAfterWallClockDeadline(t *testing.T) {
+ conn := workerTestDB(t)
+ seedWorkerActor(t, conn, true, false)
+ seedBridgeParent(t, conn, false)
+ id := seedDelivery(t, conn, "Create", gParentATURI, createPayload("x"))
+
+ // The delivery is older than the causal budget: NOW the parent has provably
+ // never landed, so it poisons — but only because wall-clock time passed, not
+ // because of the attempt count (which is zero here).
+ _, err := conn.ExecContext(context.Background(),
+ `UPDATE outbound_deliveries SET created_at = now() - interval '2 hours' WHERE activity_id = $1`, id)
+ require.NoError(t, err)
+
+ w := newWorker(t, conn, &fakeSender{}, func(o *WorkerOptions) { o.CausalWaitBudget = time.Hour })
+ _, err = w.DeliverNext(context.Background())
+ require.NoError(t, err)
+
+ d := getDelivery(t, conn, id)
+ assert.Equal(t, store.DeliveryStatePoisoned, d.State,
+ "once the wall-clock causal budget is exceeded, a never-accepted parent poisons the child")
+ assert.Contains(t, d.LastErrorClass, "parent_unaccepted",
+ "the poison reason is queryable: parent_unaccepted")
+}
+
+// ---------------------------------------------------------------------------
+// H6: the gate keys on the child's ACTUAL parent, not any lower-seq poison
+// ---------------------------------------------------------------------------
+
+func TestCausalGating_UnrelatedPoisonDoesNotPoisonChild(t *testing.T) {
+ conn := workerTestDB(t)
+ ctx := context.Background()
+ seedWorkerActor(t, conn, true, false)
+ seedBridgeParent(t, conn, false) // the child's ACTUAL parent P: pending, not poisoned
+
+ // An UNRELATED poisoned activity on the same community/inbox line — its
+ // federated object is a DIFFERENT at-uri, nothing to do with this child.
+ unrelatedID := "https://coves.social/ap/activity/" + repeatHex64("unrelated")
+ _, err := store.NewOutboundActivities(conn).Insert(ctx, store.OutboundActivity{
+ ActivityID: unrelatedID,
+ ActorDID: wActorDID,
+ Kind: "Create",
+ Payload: []byte(fmt.Sprintf(`{"id":%q,"type":"Create","object":{"type":"Note","id":%q}}`,
+ unrelatedID, objectURLFor("at://did:plc:someoneelse/social.coves.community.comment/other"))),
+ })
+ require.NoError(t, err)
+ _, err = store.NewOutboundDeliveries(conn).Enqueue(ctx, store.OutboundDelivery{
+ ActivityID: unrelatedID, TargetInbox: wInbox, OrderingKey: wCommunityAPID,
+ })
+ require.NoError(t, err)
+ _, err = conn.ExecContext(ctx, `UPDATE outbound_deliveries SET state='poisoned' WHERE activity_id=$1`, unrelatedID)
+ require.NoError(t, err)
+
+ // The child, whose actual parent (gParentATURI) is merely PENDING.
+ id := seedDelivery(t, conn, "Create", gParentATURI, createPayload("child"))
+
+ w := newWorker(t, conn, &fakeSender{}, nil)
+ _, err = w.DeliverNext(ctx)
+ require.NoError(t, err)
+
+ assert.Equal(t, store.DeliveryStatePending, getDelivery(t, conn, id).State,
+ "a poisoned UNRELATED activity on the same line must NOT poison this child — its actual "+
+ "parent is only pending, so it causal-WAITS (keying on 'any lower-seq poison' is the bug)")
+}
+
+func TestCausalGating_ActualParentPoisonPoisonsChild(t *testing.T) {
+ conn := workerTestDB(t)
+ ctx := context.Background()
+ seedWorkerActor(t, conn, true, false)
+ seedBridgeParent(t, conn, false) // parent object exists, not accepted
+
+ // The child's ACTUAL parent (gParentATURI) has a POISONED delivery: its
+ // federated object.id maps back to gParentATURI.
+ parentID := "https://coves.social/ap/activity/" + repeatHex64("realparent")
+ _, 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='poisoned' WHERE activity_id=$1`, parentID)
+ require.NoError(t, err)
+
+ id := seedDelivery(t, conn, "Create", gParentATURI, createPayload("child"))
+
+ w := newWorker(t, conn, &fakeSender{}, nil)
+ _, err = w.DeliverNext(ctx)
+ require.NoError(t, err)
+
+ d := getDelivery(t, conn, id)
+ assert.Equal(t, store.DeliveryStatePoisoned, d.State,
+ "when the child's ACTUAL parent delivery is poisoned, the child poisons")
+ assert.Contains(t, d.LastErrorClass, "parent_poisoned",
+ "the reason is parent_poisoned — distinct from parent_unaccepted")
+}
diff --git a/internal/outbound/consent_hardening_test.go b/internal/outbound/consent_hardening_test.go
new file mode 100644
index 0000000..bcd3f54
--- /dev/null
+++ b/internal/outbound/consent_hardening_test.go
@@ -0,0 +1,43 @@
+package outbound
+
+import (
+ "context"
+ "testing"
+
+ "github.com/stretchr/testify/assert"
+ "github.com/stretchr/testify/require"
+
+ "tidepool/internal/store"
+)
+
+// Second-opinion H3: a consent-blocked CREATE must cancel ONLY the create, not
+// sweep away the actor's pending retractions (Delete/Undo). A blanket
+// CancelForActor on consent-block would silently drop the very take-downs an
+// opted-out user relies on.
+
+func TestWorker_ConsentBlockDoesNotCancelPendingRetractions(t *testing.T) {
+ conn := workerTestDB(t)
+ seedWorkerActor(t, conn, false, false) // disabled: outward work is consent-blocked
+
+ // Three pending deliveries for the same actor on the same serial line, the
+ // Create enqueued first so it is the claimed head.
+ createID := seedDelivery(t, conn, "Create", "", createPayload("c"))
+ deleteID := seedDelivery(t, conn, "Delete", "", createPayload("d"))
+ undoID := seedDelivery(t, conn, "Undo", "", createPayload("u"))
+
+ sender := &fakeSender{}
+ w := newWorker(t, conn, sender, nil)
+
+ worked, err := w.DeliverNext(context.Background())
+ require.NoError(t, err)
+ assert.True(t, worked)
+
+ assert.Equal(t, store.DeliveryStateCancelled, getDelivery(t, conn, createID).State,
+ "the consent-blocked create is cancelled")
+ assert.Equal(t, store.DeliveryStatePending, getDelivery(t, conn, deleteID).State,
+ "a pending DELETE must survive a create's consent block — a retraction is exempt from consent")
+ assert.Equal(t, store.DeliveryStatePending, getDelivery(t, conn, undoID).State,
+ "a pending UNDO must survive too — cancelling it would strand a vote retraction")
+
+ assert.Zero(t, sender.count(), "the blocked create POSTs nothing")
+}
diff --git a/internal/outbound/inbox_resolver.go b/internal/outbound/inbox_resolver.go
index 53ee215..1cfa23e 100644
--- a/internal/outbound/inbox_resolver.go
+++ b/internal/outbound/inbox_resolver.go
@@ -10,12 +10,19 @@ import (
)
// ActorFetcher fetches an AP actor document by IRI. *ap.Client satisfies it via
-// FetchActor (SSRF-guarded, same-authority binding); the resolver depends only
-// on this narrow surface.
+// FetchActor; the resolver depends only on this narrow surface.
type ActorFetcher interface {
FetchActor(ctx context.Context, iri string) (*ap.Object, error)
}
+// sameAuthorityFetcher is the optional hardened fetch: it pins the redirect
+// authority to the requested IRI, so an open redirect on the community's origin
+// cannot bounce the Group-doc fetch to an attacker host. *ap.Client satisfies it
+// (FetchActorSameAuthority); a plain ActorFetcher falls back to FetchActor.
+type sameAuthorityFetcher interface {
+ FetchActorSameAuthority(ctx context.Context, iri string) (*ap.Object, error)
+}
+
// cachedInboxResolver resolves a community's target inbox from its Group actor
// document (preferring endpoints.sharedInbox), memoized for a TTL. It also
// implements FreshInboxResolver so the worker can bypass the cache once on an
@@ -60,13 +67,21 @@ func (r *cachedInboxResolver) ResolveInboxFresh(ctx context.Context, communityAP
return r.fetchAndCache(ctx, communityAPID)
}
-// fetchAndCache fetches the Group document and reads its delivery inbox. The
-// fetch runs through the ap client's SSRF + same-authority guards (ActorFetcher
-// is *ap.Client.FetchActor), so a Group doc advertising a cross-authority or
-// private-range inbox is refused at fetch time; the worker's POST re-applies the
-// egress guard on the resolved inbox.
+// fetchAndCache fetches the Group document (with the redirect authority pinned
+// when the fetcher supports it) and reads its delivery inbox. It REFUSES a
+// resolved inbox whose host is not same-authority with the community: a Group
+// doc a stranger controls must not be able to redirect a signed activity to an
+// arbitrary origin. The worker's POST re-applies the SSRF egress guard on top.
func (r *cachedInboxResolver) fetchAndCache(ctx context.Context, communityAPID string) (string, error) {
- doc, err := r.fetcher.FetchActor(ctx, communityAPID)
+ var (
+ doc *ap.Object
+ err error
+ )
+ if hardened, ok := r.fetcher.(sameAuthorityFetcher); ok {
+ doc, err = hardened.FetchActorSameAuthority(ctx, communityAPID)
+ } else {
+ doc, err = r.fetcher.FetchActor(ctx, communityAPID)
+ }
if err != nil {
return "", fmt.Errorf("resolve inbox for %s: %w", communityAPID, err)
}
@@ -74,6 +89,9 @@ func (r *cachedInboxResolver) fetchAndCache(ctx context.Context, communityAPID s
if inbox == "" {
return "", fmt.Errorf("community %s advertises no inbox", communityAPID)
}
+ if !ap.SameAuthority(communityAPID, inbox) {
+ return "", fmt.Errorf("community %s advertises a cross-authority inbox %q; refusing", communityAPID, inbox)
+ }
r.mu.Lock()
r.cache[communityAPID] = inboxEntry{inbox: inbox, expires: time.Now().Add(r.ttl)}
r.mu.Unlock()
diff --git a/internal/outbound/outbound.go b/internal/outbound/outbound.go
index ddcfebf..3ca0e16 100644
--- a/internal/outbound/outbound.go
+++ b/internal/outbound/outbound.go
@@ -29,7 +29,7 @@ import (
// SignerProvider yields the per-actor AP Signer a delivery is signed with. The
// worker signs each delivery as the PERSONA that authored the record, not as
// the service actor — Lemmy attributes the activity to the signing key's owner.
-// personas.Service satisfies this (its actorSigner, exported for task 15).
+// *personas.Service satisfies this via its SignerFor method.
type SignerProvider interface {
// SignerFor returns the Signer whose keyId is "{actorID}#main-key" for the
// persona minted under did. A DID with no minted actor is an error
diff --git a/internal/outbound/security_test.go b/internal/outbound/security_test.go
new file mode 100644
index 0000000..a97ae28
--- /dev/null
+++ b/internal/outbound/security_test.go
@@ -0,0 +1,165 @@
+package outbound
+
+import (
+ "context"
+ "database/sql"
+ "net/http"
+ "testing"
+ "time"
+
+ "github.com/stretchr/testify/assert"
+ "github.com/stretchr/testify/require"
+
+ "tidepool/internal/ap"
+ "tidepool/internal/store"
+)
+
+// Second-opinion H1 (inbox SSRF), H4a (park backoff), + importants
+// (401-not-rotation, case-insensitive host switch).
+
+// ---------------------------------------------------------------------------
+// H1: inbox SSRF — a Group doc must not be able to point delivery at a foreign
+// host.
+// ---------------------------------------------------------------------------
+
+func TestInboxResolver_RejectsCrossAuthorityInbox(t *testing.T) {
+ // The community lives on lemmy.world but its Group doc advertises an inbox
+ // on an attacker host. Delivering there would let any community redirect a
+ // signed activity to an arbitrary origin.
+ fetcher := &countingFetcher{doc: &ap.Object{
+ ID: wCommunityAPID, // https://lemmy.world/c/tech
+ Type: "Group",
+ Inbox: "https://evil.example/inbox",
+ Endpoints: &ap.Endpoints{SharedInbox: "https://evil.example/inbox"},
+ }}
+ resolver := NewInboxResolver(fetcher, time.Minute)
+
+ _, err := resolver.ResolveInbox(context.Background(), wCommunityAPID)
+ require.Error(t, err,
+ "a resolved inbox whose host is not same-authority with the community must be REFUSED, "+
+ "not returned as the delivery target")
+}
+
+func TestInboxResolver_AcceptsSameAuthorityInbox(t *testing.T) {
+ shared := "https://lemmy.world/c/tech/inbox"
+ fetcher := &countingFetcher{doc: &ap.Object{
+ ID: wCommunityAPID,
+ Type: "Group",
+ Inbox: shared,
+ Endpoints: &ap.Endpoints{SharedInbox: shared},
+ }}
+ resolver := NewInboxResolver(fetcher, time.Minute)
+
+ inbox, err := resolver.ResolveInbox(context.Background(), wCommunityAPID)
+ require.NoError(t, err, "a same-authority inbox is fine")
+ assert.Equal(t, shared, inbox)
+}
+
+func TestWorker_DoesNotDeliverToCrossAuthorityInbox(t *testing.T) {
+ conn := workerTestDB(t)
+ seedWorkerActor(t, conn, true, false)
+ // A delivery whose stored target_inbox is on a DIFFERENT host than its
+ // community (ordering key) — a poisoned/tampered enqueue must never cause a
+ // signed POST to the foreign host.
+ id := seedDeliveryTo(t, conn, "Create", "https://evil.example/inbox")
+
+ sender := &fakeSender{}
+ w := newWorker(t, conn, sender, nil)
+ _, err := w.DeliverNext(context.Background())
+ require.NoError(t, err)
+
+ assert.Zero(t, sender.count(),
+ "the worker must NOT POST to an inbox that is not same-authority with the community")
+ // The row lives at the cross-authority inbox, so read state by activity id
+ // (getDelivery hardcodes wInbox). The worker refuses it — never 'delivered'.
+ assert.NotEqual(t, store.DeliveryStateDelivered, deliveryState(t, conn, id),
+ "a cross-authority inbox delivery must never reach 'delivered'")
+}
+
+// ---------------------------------------------------------------------------
+// H4a: a parked delivery is not immediately re-claimable (real backoff)
+// ---------------------------------------------------------------------------
+
+func TestWorker_ParkedDeliveryIsNotImmediatelyReclaimable(t *testing.T) {
+ conn := workerTestDB(t)
+ seedWorkerActor(t, conn, true, false)
+
+ // ONE delivery: it is the head ClaimNext returns, so it is the one the kill
+ // switch parks (a spurious sibling would be parked instead, leaving this
+ // row's next_attempt_at at seed-time).
+ id := seedDelivery(t, conn, "Create", "", createPayload("x"))
+ w := newWorker(t, conn, &fakeSender{}, func(o *WorkerOptions) {
+ o.Switches = &fakeSwitches{allow: false}
+ })
+ _, err := w.DeliverNext(context.Background())
+ require.NoError(t, err)
+
+ // Park must advance next_attempt_at by a REAL delay, or a parked delivery
+ // spins the worker in a hot loop (and, on the causal path, burns the attempt
+ // budget). Asserting the scheduled time directly (not a re-claim, which is
+ // app↔DB clock-skew sensitive) keeps the pin deterministic.
+ d := getDelivery(t, conn, id)
+ assert.Truef(t, d.NextAttemptAt.After(time.Now().Add(time.Second)),
+ "a just-parked delivery must be scheduled a real delay into the future, got next_attempt_at=%s (now=%s)",
+ d.NextAttemptAt, time.Now())
+ assert.Equal(t, store.DeliveryStatePending, d.State, "parked stays pending")
+}
+
+// ---------------------------------------------------------------------------
+// Important: 401 is NOT an inbox-rotation signal
+// ---------------------------------------------------------------------------
+
+func TestWorker_Unauthorized401IsTransientNotRotationPoison(t *testing.T) {
+ conn := workerTestDB(t)
+ seedWorkerActor(t, conn, true, false)
+ id := seedDelivery(t, conn, "Create", "", createPayload("x"))
+
+ // A resolver that CAN rotate, and a peer that 401s everywhere. A 401 is an
+ // auth problem (our signature / their secure-mode), NOT an endpoint
+ // rotation, so it must not consume the single re-resolve and then poison.
+ resolver := &rotatingResolver{normal: wInbox, fresh: "https://lemmy.world/c/tech/inbox-v2"}
+ sender := senderReturning(httpErr(http.StatusUnauthorized, ""))
+ w := newWorker(t, conn, sender, func(o *WorkerOptions) { o.Inboxes = resolver })
+
+ _, err := w.DeliverNext(context.Background())
+ require.NoError(t, err)
+
+ assert.Equal(t, 0, resolver.freshCalls,
+ "a 401 must NOT trigger inbox re-resolution — it is not an endpoint rotation")
+ assert.Equal(t, store.DeliveryStatePending, getDelivery(t, conn, id).State,
+ "a 401 is transient (retry), not a poison-after-one-rotation")
+}
+
+// ---------------------------------------------------------------------------
+// Important: ConfigSwitches host matching is case-insensitive
+// ---------------------------------------------------------------------------
+
+func TestConfigSwitches_HostMatchIsCaseInsensitive(t *testing.T) {
+ sw := ConfigSwitches{
+ DisabledHosts: map[string]struct{}{"Lemmy.World": {}},
+ }
+ assert.False(t, sw.OutboundAllowed(DeliveryScope{InboxHost: "lemmy.world"}),
+ "a disabled-host entry must block regardless of case — hostOf lowercases the scope, so a "+
+ "mixed-case switch value must still block the (lowercased) host")
+}
+
+// seedDeliveryTo is seedDelivery with an explicit target inbox.
+func seedDeliveryTo(t *testing.T, conn *sql.DB, kind, inbox string) string {
+ t.Helper()
+ ctx := context.Background()
+ activityID := "https://coves.social/ap/activity/" + repeatHex64(kind+inbox)
+ _, err := store.NewOutboundActivities(conn).Insert(ctx, store.OutboundActivity{
+ ActivityID: activityID,
+ ActorDID: wActorDID,
+ Kind: kind,
+ Payload: createPayload(activityID),
+ })
+ require.NoError(t, err)
+ _, err = store.NewOutboundDeliveries(conn).Enqueue(ctx, store.OutboundDelivery{
+ ActivityID: activityID,
+ TargetInbox: inbox,
+ OrderingKey: wCommunityAPID,
+ })
+ require.NoError(t, err)
+ return activityID
+}
diff --git a/internal/outbound/switches.go b/internal/outbound/switches.go
index c6e5663..8bbd006 100644
--- a/internal/outbound/switches.go
+++ b/internal/outbound/switches.go
@@ -1,5 +1,7 @@
package outbound
+import "strings"
+
// ConfigSwitches is the config-backed Switches: the operator kill switches
// (decision 19) resolved from static configuration. It is the adapter main
// builds from config values and hands the Worker.
@@ -27,8 +29,14 @@ func (s ConfigSwitches) OutboundAllowed(scope DeliveryScope) bool {
if s.Disabled {
return false
}
- if _, blocked := s.DisabledHosts[scope.InboxHost]; blocked {
- return false
+ // Case-insensitive on host: the scope host is lowercased (hostOf), but a
+ // switch value constructed directly may be mixed-case, and a kill switch
+ // must fail closed rather than silently miss on case.
+ host := strings.ToLower(scope.InboxHost)
+ for disabled := range s.DisabledHosts {
+ if strings.ToLower(disabled) == host {
+ return false
+ }
}
if _, blocked := s.DisabledCommunities[scope.CommunityAPID]; blocked {
return false
diff --git a/internal/outbound/translator_note_test.go b/internal/outbound/translator_note_test.go
new file mode 100644
index 0000000..489aa14
--- /dev/null
+++ b/internal/outbound/translator_note_test.go
@@ -0,0 +1,56 @@
+package outbound
+
+import (
+ "encoding/json"
+ "testing"
+
+ "github.com/stretchr/testify/assert"
+ "github.com/stretchr/testify/require"
+
+ "tidepool/internal/ap"
+ "tidepool/internal/consume"
+)
+
+// Second-opinion (important): the Page/Note addressing split fence. A Note
+// carries the community in `cc` — it must NEVER appear in `to` (only Public). A
+// Note with the community in `to` is the Page shape, which changes how Lemmy
+// routes it. This pins the fence from the Note side (the Page side is pinned in
+// translator_page_test.go).
+
+func TestTranslator_NoteToExcludesCommunity(t *testing.T) {
+ tr := NewTranslator("https://coves.social")
+ community := "https://lemmy.world/c/tech"
+ actorID := "https://coves.social/ap/actor/did:plc:x"
+ atURI := "at://did:plc:x/social.coves.community.comment/rk"
+
+ snap, err := json.Marshal(map[string]any{
+ "atUri": atURI,
+ "record": map[string]any{"content": "a reply", "createdAt": "2026-08-12T10:00:00.000Z"},
+ "parentAtUri": "at://did:plc:p/social.coves.community.postv2/rp",
+ "parentApId": "https://lemmy.world/post/1",
+ "communityApId": community,
+ })
+ require.NoError(t, err)
+
+ out, err := tr.Translate(actorID, consume.CommentIntent{
+ Op: "create",
+ ATURI: atURI,
+ ID: "https://coves.social/ap/activity/deadbeef",
+ CommunityAPID: community,
+ ParentAPID: "https://lemmy.world/post/1",
+ Snapshot: snap,
+ })
+ require.NoError(t, err)
+
+ var activity map[string]any
+ require.NoError(t, json.Unmarshal(out.Payload, &activity))
+ note := mustMap(t, activity["object"], "object")
+
+ to := asStringSet(t, note["to"], "Note.to")
+ assert.NotContains(t, to, community,
+ "a Note must NOT carry the community in `to` — that is the Page shape (Page/Note split)")
+ assert.Contains(t, to, ap.PublicAudience, "a Note's `to` is Public only")
+
+ cc := asStringSet(t, note["cc"], "Note.cc")
+ assert.Contains(t, cc, community, "a Note carries the community in `cc`")
+}
diff --git a/internal/outbound/worker.go b/internal/outbound/worker.go
index 6b49509..754bdc9 100644
--- a/internal/outbound/worker.go
+++ b/internal/outbound/worker.go
@@ -57,6 +57,13 @@ type WorkerOptions struct {
// BackoffBase is the first retry-backoff step (doubles per attempt). Zero
// uses DefaultBackoffBase; tests compress it.
BackoffBase time.Duration
+ // CausalWaitBudget is the WALL-CLOCK deadline a reply waits for its
+ // bridge-origin parent to be accepted before poisoning parent_unaccepted.
+ // It is measured from the delivery's creation, NOT its attempt count — a
+ // causal wait must not consume the delivery-failure retry budget (a parent
+ // legitimately takes minutes to be admitted). Zero uses
+ // DefaultCausalWaitBudget.
+ CausalWaitBudget time.Duration
// Lease overrides DefaultLease.
Lease time.Duration
// Logger receives per-delivery outcomes. Nil uses slog.Default().
@@ -71,6 +78,9 @@ const (
// DefaultBackoffBase is the first retry step; it doubles per attempt,
// capped at one hour.
DefaultBackoffBase = 30 * time.Second
+ // DefaultCausalWaitBudget is the wall-clock window a reply waits for its
+ // bridge-origin parent to be accepted before poisoning parent_unaccepted.
+ DefaultCausalWaitBudget = 6 * time.Hour
)
// Worker claims one delivery at a time and carries it to a terminal state. It
@@ -78,21 +88,22 @@ const (
// (Lemmy dedupes on our stable activity id, and its duplicate-activity response
// is classified DELIVERED, not poisoned).
type Worker struct {
- db *sql.DB
- activities store.OutboundActivities
- deliveries store.OutboundDeliveries
- objects store.OutboundObjects
- actors store.APActors
- prefs store.FederationPrefs
- votes store.OutboundVotes
- signers SignerProvider
- inboxes InboxResolver
- sender ActivitySender
- switches Switches
- maxAttempts int
- backoffBase time.Duration
- lease time.Duration
- logger *slog.Logger
+ db *sql.DB
+ activities store.OutboundActivities
+ deliveries store.OutboundDeliveries
+ objects store.OutboundObjects
+ actors store.APActors
+ prefs store.FederationPrefs
+ votes store.OutboundVotes
+ signers SignerProvider
+ inboxes InboxResolver
+ sender ActivitySender
+ switches Switches
+ maxAttempts int
+ backoffBase time.Duration
+ causalWaitBudget time.Duration
+ lease time.Duration
+ logger *slog.Logger
}
// NewWorker wires a Worker.
@@ -137,22 +148,27 @@ func NewWorker(opts WorkerOptions) (*Worker, error) {
if backoffBase <= 0 {
backoffBase = DefaultBackoffBase
}
+ causalWaitBudget := opts.CausalWaitBudget
+ if causalWaitBudget <= 0 {
+ causalWaitBudget = DefaultCausalWaitBudget
+ }
return &Worker{
- db: opts.DB,
- activities: activities,
- deliveries: deliveries,
- objects: objects,
- actors: opts.Actors,
- prefs: prefs,
- votes: votes,
- signers: opts.Signers,
- inboxes: opts.Inboxes,
- sender: opts.Sender,
- switches: switches,
- maxAttempts: maxAttempts,
- backoffBase: backoffBase,
- lease: lease,
- logger: logger,
+ db: opts.DB,
+ activities: activities,
+ deliveries: deliveries,
+ objects: objects,
+ actors: opts.Actors,
+ prefs: prefs,
+ votes: votes,
+ signers: opts.Signers,
+ inboxes: opts.Inboxes,
+ sender: opts.Sender,
+ switches: switches,
+ maxAttempts: maxAttempts,
+ backoffBase: backoffBase,
+ causalWaitBudget: causalWaitBudget,
+ lease: lease,
+ logger: logger,
}, nil
}
@@ -228,31 +244,45 @@ func (w *Worker) handle(ctx context.Context, delivery *store.OutboundDelivery) e
case causalEligible:
// fall through to consent + delivery
case causalWait:
- return w.park(ctx, delivery, "parent_pending", "waiting for bridge-origin parent to be accepted")
+ // 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")
case causalPoisonUnaccepted:
- return w.poison(ctx, delivery, "parent_unaccepted", "bounded wait exhausted; parent never accepted", 0)
+ return w.poison(ctx, delivery, "parent_unaccepted", "causal wait budget exhausted; parent never accepted", 0)
case causalPoisonParent:
return w.poison(ctx, delivery, "parent_poisoned", "parent delivery poisoned; descendant cannot land", 0)
}
// Consent recheck (retraction asymmetry): a Delete/Undo always goes out —
// it is how an opted-out user takes down what is already federated. Outward
- // kinds are cancelled when the actor is disabled, paused, or opted out.
+ // kinds are cancelled when the actor is disabled, paused, or opted out —
+ // but ONLY this one claimed delivery, never the actor's pending retractions.
if !isRetraction(activity.Kind) {
blocked, err := w.consentBlocked(ctx, activity.ActorDID)
if err != nil {
return err
}
if blocked {
- cancelled, err := w.deliveries.CancelForActor(ctx, activity.ActorDID)
+ _, applied, err := w.deliveries.CancelClaimed(ctx, delivery.ActivityID, delivery.TargetInbox, *delivery.ClaimedUntil)
if err != nil {
- return fmt.Errorf("cancel deliveries for %s: %w", activity.ActorDID, err)
+ return fmt.Errorf("cancel delivery %s: %w", delivery.ActivityID, err)
+ }
+ if applied {
+ metricCancelled.Add(1)
}
- metricCancelled.Add(cancelled)
return nil
}
}
+ // Defense in depth: never sign a POST to an inbox that is not same-authority
+ // with the target community. The resolver already refuses a cross-authority
+ // inbox at enqueue time; this catches a tampered or legacy stored target.
+ if !ap.SameAuthority(delivery.OrderingKey, delivery.TargetInbox) {
+ return w.poison(ctx, delivery, "cross_authority",
+ "target inbox is not same-authority with the community; refusing to deliver", 0)
+ }
+
return w.deliver(ctx, delivery, activity)
}
@@ -283,13 +313,17 @@ func (w *Worker) classify(ctx context.Context, delivery *store.OutboundDelivery,
// 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)
- case he.StatusCode == http.StatusUnauthorized ||
- he.StatusCode == http.StatusNotFound ||
+ case he.StatusCode == http.StatusNotFound ||
he.StatusCode == http.StatusGone:
+ // 404/410 is an endpoint-GONE signal: re-resolve the inbox once
+ // before poisoning (a rotation must not become a poison).
return w.rotateInbox(ctx, delivery, activity, signer, he)
- case he.StatusCode == http.StatusRequestTimeout ||
+ case he.StatusCode == http.StatusUnauthorized ||
+ he.StatusCode == http.StatusRequestTimeout ||
he.StatusCode == http.StatusTooManyRequests ||
he.StatusCode >= 500:
+ // 401 is an AUTH problem (our signature, their secure mode), not an
+ // endpoint rotation: retry it, don't burn the single re-resolve.
return w.releaseOrPoison(ctx, delivery, classForStatus(he.StatusCode), he.Body, he.StatusCode)
default:
// Other 4xx: a genuine rejection. Retried on a small budget, then
@@ -330,9 +364,18 @@ func (w *Worker) rotateInbox(ctx context.Context, delivery *store.OutboundDelive
return w.poison(ctx, delivery, "inbox_gone", "inbox still unreachable after re-resolution", status)
}
-// deliverSuccess marks the delivery delivered under its fencing token and fires
-// the object-acceptance and vote-delivery callbacks.
+// deliverSuccess opens the causal gate and marks the delivery delivered, in
+// that order so the two are effectively atomic: a parent must NEVER be observed
+// delivered while its accepted_at is unset (that strands every child forever).
+// The accepted_at stamp is written FIRST; only if it commits is the delivered
+// mark applied. If the stamp genuinely fails, we return before marking and the
+// retry re-runs both (Lemmy dedupes the re-POST). The only reachable states are
+// (¬accepted,¬delivered), (accepted,¬delivered), (accepted,delivered) — never
+// the forbidden (¬accepted,delivered).
func (w *Worker) deliverSuccess(ctx context.Context, delivery *store.OutboundDelivery, activity *store.OutboundActivity, status int) error {
+ if err := w.stampAccepted(ctx, activity); err != nil {
+ return fmt.Errorf("stamp accepted for %s: %w", delivery.ActivityID, err)
+ }
_, applied, err := w.deliveries.MarkDelivered(ctx, delivery.ActivityID, delivery.TargetInbox, status, *delivery.ClaimedUntil)
if err != nil {
return fmt.Errorf("mark delivered %s: %w", delivery.ActivityID, err)
@@ -341,25 +384,29 @@ func (w *Worker) deliverSuccess(ctx context.Context, delivery *store.OutboundDel
return nil // a stale claim: another worker already recorded the outcome
}
metricDelivered.Add(1)
- w.stampAccepted(ctx, activity)
return w.voteCallback(ctx, activity)
}
// stampAccepted opens the causal gate for this object's children: on a
// successful Create/Update, the object it federated is now accepted by its
-// community. Best-effort — a delivery for an object with no outbound_objects row
-// (a comment we never persisted, a vote) simply has nothing to stamp.
-func (w *Worker) stampAccepted(ctx context.Context, activity *store.OutboundActivity) {
+// community. A NotFound (no outbound_objects row — a comment we never persisted,
+// a vote) is not a failure and returns nil; a real store error is propagated so
+// deliverSuccess withholds the delivered mark until the stamp can commit.
+func (w *Worker) stampAccepted(ctx context.Context, activity *store.OutboundActivity) error {
if activity.Kind != "Create" && activity.Kind != "Update" {
- return
+ return nil
}
atURI := objectATURIFromPayload(activity.Payload)
if atURI == "" {
- return
+ return nil
}
- if err := w.objects.SetAccepted(ctx, atURI); err != nil && !errors.IsNotFound(err) {
- w.logger.Warn("stamp accepted failed", "at_uri", atURI, "error", err)
+ if err := w.objects.SetAccepted(ctx, atURI); err != nil {
+ if errors.IsNotFound(err) {
+ return nil
+ }
+ return err
}
+ return nil
}
// voteCallback applies decision-16 delivery callbacks: a Like/Dislike success
@@ -420,11 +467,18 @@ func (w *Worker) poison(ctx context.Context, delivery *store.OutboundDelivery, c
return nil
}
-// park releases the delivery pending with a short backoff and no move toward
-// the poison cap: a kill switch, dry-run, or causal wait is a "not now", never
-// a failure.
+// parkDelay is how long a kill-switched or dry-run delivery waits before it can
+// be re-claimed: long enough that a parked row does not spin the worker in a hot
+// loop, short enough that clearing the switch resumes delivery promptly.
+const parkDelay = 5 * time.Second
+
+// park holds a kill-switched or dry-run delivery: it stays pending, scheduled a
+// REAL delay into the future so it is not instantly re-claimable, and never
+// poisons. Release does not touch the attempt counter, so a park is not a
+// failure and does not itself advance the poison budget.
func (w *Worker) park(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)
+ next := time.Now().Add(parkDelay)
+ _, _, err := w.deliveries.Release(ctx, delivery.ActivityID, delivery.TargetInbox, class, reason, 0, next, *delivery.ClaimedUntil)
if err != nil {
return fmt.Errorf("park delivery %s: %w", delivery.ActivityID, err)
}
@@ -432,6 +486,20 @@ func (w *Worker) park(ctx context.Context, delivery *store.OutboundDelivery, cla
return nil
}
+// 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 does
+// not advance the poison budget; the causal wait is bounded by wall clock.
+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)
+ if err != nil {
+ return fmt.Errorf("park (causal) delivery %s: %w", delivery.ActivityID, err)
+ }
+ metricParked.Add(1)
+ return nil
+}
+
// causalStatus classifies a delivery's causal eligibility.
type causalStatus int
@@ -460,18 +528,24 @@ func (w *Worker) causalStatus(ctx context.Context, delivery *store.OutboundDeliv
if parent.IsAccepted() {
return causalEligible
}
- // Bridge-origin parent, not yet accepted. A poisoned ancestor on the same
- // serial line means it will NEVER land → poison the descendant distinctly.
- poisoned, err := w.deliveries.HasPoisonedPredecessor(ctx, delivery.OrderingKey, delivery.TargetInbox, delivery.Seq)
+ // 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
+ // poison this child.
+ poisoned, err := w.deliveries.ParentDeliveryPoisoned(ctx, activity.ParentATURI, delivery.TargetInbox)
if err != nil {
- w.logger.Error("poisoned-predecessor check failed", "error", err)
+ w.logger.Error("parent-delivery poisoned check failed", "error", err)
return causalWait
}
if poisoned {
return causalPoisonParent
}
- if delivery.Attempts >= w.maxAttempts {
- return causalPoisonUnaccepted // bounded wait exhausted
+ // 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
+ // legitimately takes minutes is not poisoned just because the child was
+ // claimed a few times.
+ if time.Since(delivery.CreatedAt) >= w.causalWaitBudget {
+ return causalPoisonUnaccepted
}
return causalWait
}
@@ -527,6 +601,8 @@ func isDuplicate(he ap.HTTPError) bool {
// classForStatus labels a retryable HTTP status for the retry taxonomy.
func classForStatus(status int) string {
switch {
+ case status == http.StatusUnauthorized:
+ return "unauthorized"
case status == http.StatusRequestTimeout:
return "timeout"
case status == http.StatusTooManyRequests:
diff --git a/internal/outbound/worker_test.go b/internal/outbound/worker_test.go
index 62672d2..1b262a7 100644
--- a/internal/outbound/worker_test.go
+++ b/internal/outbound/worker_test.go
@@ -209,14 +209,15 @@ func getDelivery(t *testing.T, conn *sql.DB, activityID string) *store.OutboundD
func newWorker(t *testing.T, conn *sql.DB, sender ActivitySender, opts func(*WorkerOptions)) *Worker {
t.Helper()
o := WorkerOptions{
- DB: conn,
- Actors: store.NewAPActors(conn),
- Signers: newFakeSigners(t),
- Inboxes: staticInbox{inbox: wInbox},
- Sender: sender,
- Lease: time.Minute,
- MaxAttempts: 3,
- BackoffBase: time.Millisecond,
+ DB: conn,
+ Actors: store.NewAPActors(conn),
+ Signers: newFakeSigners(t),
+ Inboxes: staticInbox{inbox: wInbox},
+ Sender: sender,
+ Lease: time.Minute,
+ MaxAttempts: 3,
+ BackoffBase: time.Millisecond,
+ CausalWaitBudget: time.Hour, // long by default: recent deliveries never time out
}
if opts != nil {
opts(&o)
diff --git a/internal/personas/personas.go b/internal/personas/personas.go
index 64a103e..2609258 100644
--- a/internal/personas/personas.go
+++ b/internal/personas/personas.go
@@ -230,17 +230,12 @@ func suffixedLocalPart(base string, attempt int) string {
return base + "-" + strconv.Itoa(attempt)
}
-// ActorSigner is the exported per-actor signer accessor (task 15 seam): the
+// SignerFor is THE exported per-actor signer accessor (task 15 seam): the
// outbound delivery worker signs each activity as the persona that authored the
-// record, not as the service actor. It is actorSigner promoted to the package
-// surface; internal callers keep using the unexported form.
-func (s *Service) ActorSigner(ctx context.Context, did string) (*ap.Signer, error) {
- return s.actorSigner(ctx, did)
-}
-
-// SignerFor satisfies outbound.SignerProvider so main can inject *Service as the
-// delivery worker's signer source (cycle J). It is ActorSigner under the name
-// the worker's interface uses.
+// record, not as the service actor. It promotes the unexported actorSigner to
+// the package surface under the name outbound.SignerProvider requires, so main
+// injects *Service directly as the worker's signer source (cycle J). Internal
+// callers keep using actorSigner.
func (s *Service) SignerFor(ctx context.Context, did string) (*ap.Signer, error) {
return s.actorSigner(ctx, did)
}
diff --git a/internal/store/interfaces.go b/internal/store/interfaces.go
index 5477feb..e9e8954 100644
--- a/internal/store/interfaces.go
+++ b/internal/store/interfaces.go
@@ -540,13 +540,20 @@ type OutboundDeliveries interface {
// error satisfying errors.IsNotFound.
Get(ctx context.Context, activityID, targetInbox string) (*OutboundDelivery, error)
- // HasPoisonedPredecessor reports whether an earlier delivery on the same
- // ordering key and target inbox (a lower seq) is poisoned — the causal
- // signal task 15's worker reads to distinguish a child whose bridge-origin
- // parent WILL NEVER land (parent_poisoned) from one merely waiting
- // (parent_unaccepted). Per-community serialization makes a lower-seq
- // delivery on the same line a causal ancestor.
- HasPoisonedPredecessor(ctx context.Context, orderingKey, targetInbox string, seq int64) (bool, 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).
+ // 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)
+
+ // CancelClaimed cancels a SINGLE claimed delivery under its fencing token
+ // (the consent-block outcome for one create/update), leaving the actor's
+ // other pending work — its Delete/Undo retractions above all — untouched.
+ // Same (exists, applied) fencing contract as MarkDelivered.
+ CancelClaimed(ctx context.Context, activityID, targetInbox string, claimToken time.Time) (exists, applied bool, err error)
// CountsByState returns the number of deliveries in each state — the
// operator queue-inspect (GET /admin/outbound).
diff --git a/internal/store/models.go b/internal/store/models.go
index 0ef29be..cc25483 100644
--- a/internal/store/models.go
+++ b/internal/store/models.go
@@ -178,8 +178,12 @@ const (
DeliveredStatePending DeliveredState = "pending"
// DeliveredStateDelivered means a peer accepted the Like/Dislike.
DeliveredStateDelivered DeliveredState = "delivered"
- // DeliveredStateUndone means the Undo was delivered; the row is kept as
- // the record of what was withdrawn.
+ // DeliveredStateUndone is RESERVED and currently UNWRITTEN: task 15's worker
+ // DELETES the outbound_votes row on a successful Undo (clear-on-Undo) rather
+ // than transitioning it to "undone", so no code path ever sets this today.
+ // It is kept in the enum and the CHECK constraint (the migration is applied)
+ // against a future "keep the withdrawn-vote record" policy; Valid() still
+ // accepts it so a hand-set or legacy row round-trips.
DeliveredStateUndone DeliveredState = "undone"
)
diff --git a/internal/store/outbound_deliveries.go b/internal/store/outbound_deliveries.go
index 1389962..0393401 100644
--- a/internal/store/outbound_deliveries.go
+++ b/internal/store/outbound_deliveries.go
@@ -5,6 +5,7 @@ import (
"database/sql"
stderrors "errors"
"fmt"
+ "strings"
"time"
"tidepool/internal/errors"
@@ -76,12 +77,22 @@ func (r *postgresOutboundDeliveries) ClaimNext(ctx context.Context, lease time.D
// The heads are found by a recursive CTE emulating a loose index scan over
// idx_outbound_deliveries_queue (ordering_key, seq WHERE state='pending'):
// one index descent per DISTINCT pending key jumps to each key's head, so
- // the claim is O(pending keys × log N) regardless of any one community's
- // backlog depth — never O(backlog) as a per-row NOT EXISTS would be. They
- // are materialized with ARRAY(...) — not a plain IN or a correlated EXISTS —
- // so the planner fetches exactly those rows by seq. The outer SELECT
- // re-applies every claimability predicate on the locked row (a claim
- // committed between the CTE snapshot and the lock is then seen and skipped).
+ // the claim finds the DISTINCT pending keys' heads without a per-row NOT
+ // EXISTS. The heads are materialized with ARRAY(...) — not a plain IN or a
+ // correlated EXISTS — so the head set is computed ONCE (the loose scan)
+ // rather than re-derived per row. The outer SELECT then locks and re-applies
+ // every claimability predicate on just that head set (a claim committed
+ // between the CTE snapshot and the lock is then seen and skipped).
+ //
+ // NOTE on the outer re-fetch: seq is a BIGSERIAL ordering column, NOT the
+ // primary key (the PK is (activity_id, target_inbox)) and has no standalone
+ // index — so `c.seq = ANY(ARRAY(...))` is a re-check over the small head set,
+ // NOT the indexed point-fetch inbox_events gets (there `id` IS the PK). For a
+ // deep pending backlog the planner can only reach the head rows through the
+ // partial (ordering_key, seq) index, so the re-check is not the strict
+ // O(keys × log N) a seq index would give. A dedicated UNIQUE index on seq
+ // (its own migration, so goose actually applies it) would restore the
+ // point-fetch and is worth adding if this path ever profiles hot.
// FOR UPDATE ... SKIP LOCKED lets concurrent workers race without
// serializing on row locks; the UPDATE stamps the lease and counts the
// attempt atomically.
@@ -264,20 +275,50 @@ func (r *postgresOutboundDeliveries) Get(ctx context.Context, activityID, target
return delivery, nil
}
-func (r *postgresOutboundDeliveries) HasPoisonedPredecessor(ctx context.Context, orderingKey, targetInbox string, seq int64) (bool, error) {
+func (r *postgresOutboundDeliveries) ParentDeliveryPoisoned(ctx context.Context, parentATURI, targetInbox string) (bool, 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).
+ suffix := strings.TrimPrefix(parentATURI, "at://")
var exists bool
err := r.db.QueryRowContext(ctx, `
SELECT EXISTS (
- SELECT 1 FROM outbound_deliveries
- WHERE ordering_key = $1 AND target_inbox = $2
- AND state = 'poisoned' AND seq < $3)`,
- orderingKey, targetInbox, seq).Scan(&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)
if err != nil {
- return false, fmt.Errorf("check poisoned predecessor on %q: %w", orderingKey, err)
+ return false, fmt.Errorf("check parent delivery poisoned for %q: %w", parentATURI, err)
}
return exists, nil
}
+func (r *postgresOutboundDeliveries) CancelClaimed(ctx context.Context, activityID, targetInbox string, claimToken time.Time) (bool, bool, error) {
+ // Fenced single-row cancel: only the current claim holder cancels, and only
+ // while pending, so a stale worker cannot clobber a re-claim and — unlike
+ // CancelForActor — an actor's OTHER pending deliveries (its retractions) are
+ // left standing.
+ query := `
+ WITH updated AS (
+ UPDATE outbound_deliveries
+ SET state = 'cancelled', claimed_until = NULL, updated_at = now()
+ WHERE activity_id = $1 AND target_inbox = $2
+ AND state = 'pending'
+ AND claimed_until = $3
+ RETURNING 1
+ )
+ SELECT
+ EXISTS (SELECT 1 FROM outbound_deliveries WHERE activity_id = $1 AND target_inbox = $2),
+ EXISTS (SELECT 1 FROM updated)`
+
+ return r.markResult(ctx, "cancel claimed", query, activityID, targetInbox, claimToken.UTC())
+}
+
func (r *postgresOutboundDeliveries) CountsByState(ctx context.Context) (map[DeliveryState]int, error) {
rows, err := r.db.QueryContext(ctx,
`SELECT state, COUNT(*) FROM outbound_deliveries GROUP BY state`)
diff --git a/internal/store/outbound_delivery_test.go b/internal/store/outbound_delivery_test.go
index 5d37cfb..6f61a2e 100644
--- a/internal/store/outbound_delivery_test.go
+++ b/internal/store/outbound_delivery_test.go
@@ -238,6 +238,37 @@ func TestOutboundDeliveries_ClaimNextStampsFencingAndBumpsAttempts(t *testing.T)
assert.True(t, errors.IsNotFound(err), "want NotFound (empty queue), got %v", err)
}
+func TestOutboundDeliveries_ExpiredLeaseIsReclaimable(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)
+
+ // A worker that claimed this delivery then crashed: its lease lapses. The
+ // whole point of a lease is that the delivery becomes re-claimable — a stuck
+ // worker must never strand a delivery forever.
+ _, 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, "a delivery whose lease has EXPIRED must be re-claimable")
+ require.NotNil(t, reclaimed)
+ assert.Equal(t, delActivityID, reclaimed.ActivityID)
+ require.NotNil(t, reclaimed.ClaimedUntil)
+ assert.True(t, reclaimed.ClaimedUntil.After(time.Now()),
+ "the re-claim stamps a fresh future lease")
+}
+
func TestOutboundDeliveries_PerOrderingKeySerialization(t *testing.T) {
database := deliveryTestDB(t)
activities := NewOutboundActivities(database)