From 0479e64f79c90b5aae06dd5fda81953e4620773b Mon Sep 17 00:00:00 2001 From: Bretton Date: Thu, 13 Aug 2026 18:51:21 -0700 Subject: [PATCH] fix(echo): adversarial-review fix wave (task 17a) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Five-stream review returned needs-work; two streams independently found the same correctness bug and codex judged the patch incorrect. All four HIGH findings fixed test-first and bite-proofed. H1 — Classify descended into TARGETS, not just payloads. `object` is the payload only for wrapper verbs; for Like/Dislike/Delete it is the activity's TARGET, so the walk reached the target of a GENUINE remote activity and condemned the envelope. A Lemmy user's vote on a post we federated out was dropped before ApplyVote (native tallies would sit at zero forever), and a Lemmy moderator's removal of a native post could never fire — the feature 17c exists to deliver. Reached genuine traffic on 3 of 4 routes; the 4th survived only because it has no guard. The drop counter was scoring discarded genuine votes as successful suppressions. Fixed with a payload-verb ALLOWLIST (ap has no TypeFlag/TypeBlock/TypeRemove constants, and an unknown verb must fail toward not condemning). The suite missed this because every "genuine remote" fixture targeted LEMMY-origin content; the breaking shape is a remote actor acting ON our object. H2 — the community backfill bypassed both guards, reaching the materializer directly. A community outbox holds our own posts once native users participate, and MaterializePost calls EnsureActor BEFORE reading any mapping, so an ordinary admin re-backfill was a mint oracle minting a bridged_actors row for our own persona — plus 3 dereferences of our own origin and an uncounted drop. The item then died writing into a repo we do not host, so every trigger reported "1 of 2 items failed" and the community never looked cleanly backfilled. Guarded before resolveEmbedded (resolving is what dereferences) as a skip, not an error, so runs stay clean. H3 — two guards whose removal the suite could not detect: the authority-spoof check (its only negative cases exited via NotFound and never reached it) and the absence of any POSITIVE vanity-origin case, so an implementation comparing against one configured origin — what decision 10 forbids — passed everything. normalizeHost was an untested hand-copy. All three now bite. H4 — classification weighed a bare id string, and our ids are public. Now corroborated against the activity we actually sent: contradiction disqualifies, absence does not. Verb vs stored Kind, carried object vs the stored payload, disqualifying only when the substituted object is not ours (swapping one of our objects cannot lose anyone's content; a Lemmy human's note is the forgery). Never byte-equality — a real echo returns reserialized, and byte-comparison would refuse every genuine one. Structural (four reviewers converged): Echo is now REQUIRED in both ingest.NewHandler and NewBackfill and is an interface, so the dispatcher's fail-safe is finally exercisable — a classification that cannot be MADE parks and retries, never poisons and never materializes, with no counter moved (a failure is not a suppression). echo.New validates its four stores. expvar names underscored for Prometheus. identifyActor compares scheme. Two load-bearing doc comments corrected where the code contradicted them. Co-Authored-By: Claude Opus 5 (1M context) --- cmd/tidepool/main.go | 1 + internal/echo/classify_test.go | 109 +++++++++ internal/echo/echo.go | 232 +++++++++++++++++-- internal/echo/echo_test.go | 87 ++++++- internal/ingest/backfill.go | 69 +++++- internal/ingest/backfill_test.go | 8 +- internal/ingest/echo_backfill_test.go | 224 ++++++++++++++++++ internal/ingest/echo_failsafe_test.go | 160 +++++++++++++ internal/ingest/echo_target_test.go | 315 ++++++++++++++++++++++++++ internal/ingest/handler.go | 51 +++-- internal/ingest/handler_test.go | 23 +- internal/ingest/ingest_test.go | 6 + internal/votes/e2e_test.go | 21 ++ 13 files changed, 1255 insertions(+), 51 deletions(-) create mode 100644 internal/ingest/echo_backfill_test.go create mode 100644 internal/ingest/echo_failsafe_test.go create mode 100644 internal/ingest/echo_target_test.go diff --git a/cmd/tidepool/main.go b/cmd/tidepool/main.go index 9e51773..fd6fe1e 100644 --- a/cmd/tidepool/main.go +++ b/cmd/tidepool/main.go @@ -359,6 +359,7 @@ func run(logger *slog.Logger) error { Communities: communities, Tombstones: tombstones, Seeder: seeder, + Echo: echoClassifier, MaxPosts: cfg.BackfillMaxPosts, // Async runs derive from the run context so a mid-run backfill stops // pulling remote pages once shutdown starts; the drain below waits for diff --git a/internal/echo/classify_test.go b/internal/echo/classify_test.go index 0f3eca3..36da7a9 100644 --- a/internal/echo/classify_test.go +++ b/internal/echo/classify_test.go @@ -500,3 +500,112 @@ func nestedEnvelope(levels int, deepestID string) string { } return body } + +// H4 — CLASSIFICATION MUST BIND THE ENVELOPE, NOT A BARE ID STRING. +// +// Our activity ids are public and derivable: a peer that has received one can +// copy it onto a node wrapping somebody ELSE's content and have the whole +// envelope dropped as "our echo". The practical severity is bounded (the +// announce path runs after the followed-community gate, so the forger must be a +// community we follow, suppressing content it chose to announce) — but the fix +// is cheap and the residual is real: a followed community can make content +// appear announced to everyone else while we silently drop it. +// +// We persist the canonical payload precisely because it is byte-stable for +// replay, so it is available as corroboration. The rule is CONTRADICTION +// DISQUALIFIES, ABSENCE DOES NOT: a node that says nothing (a bare IRI, W4 +// above) is still ours, while a node that says something INCOMPATIBLE with what +// we stored under that id is not. +func TestClassifyRefusesForgedNodesWearingOurActivityIds(t *testing.T) { + classifier, _, database := newWorld(t) + seedMirroredLemmyUser(t, database) + ctx := context.Background() + + cases := []struct { + name string + body string + why string + }{ + { + name: "a Delete wearing the id of a Create we sent", + body: `{ + "id": "` + ecLemmyAnnounceID + `/forged-kind", + "type": "Announce", + "actor": "` + ecLemmyGroupAPID + `", + "object": { + "id": "` + ecActivityID + `", + "type": "Delete", + "actor": "` + ecLemmyPersonAPID + `", + "object": "` + ecLemmyNoteAPID + `" + } + }`, + why: "we stored that id as a Create; a node calling itself a Delete is not the " + + "activity we sent, and honouring the id alone lets a peer suppress any " + + "delete it likes by wearing one of our ids", + }, + { + name: "a Create wearing our id but carrying somebody else's object", + body: `{ + "id": "` + ecLemmyAnnounceID + `/forged-content", + "type": "Announce", + "actor": "` + ecLemmyGroupAPID + `", + "object": { + "id": "` + ecActivityID + `", + "type": "Create", + "actor": "` + ecLemmyPersonAPID + `", + "object": { + "id": "` + ecLemmyNoteAPID + `", + "type": "Note", + "attributedTo": "` + ecLemmyPersonAPID + `" + } + } + }`, + why: "the type matches, so a Kind check alone is not enough: the stored payload " + + "names OUR object and this one names a Lemmy human's note — dropping it is " + + "content loss dressed as echo suppression", + }, + } + + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + identity, err := classifier.Classify(ctx, envelope(t, tc.body)) + require.NoError(t, err) + assert.Equal(t, ClassNone, identity.Class, tc.why) + }) + } +} + +// TestClassifyStillRecognizesARealEchoAfterReserialization is the control that +// keeps the H4 fix honest. A community does not hand our activity back +// byte-for-byte: it re-serializes, reorders keys, and may add or drop +// addressing fields. Corroboration must therefore be on the IDENTIFYING fields +// — the activity's type and the object it carries — never on the bytes, or +// every genuine echo leaks through and re-materializes as remote content. +func TestClassifyStillRecognizesARealEchoAfterReserialization(t *testing.T) { + classifier, _, _ := newWorld(t) + ctx := context.Background() + + identity, err := classifier.Classify(ctx, envelope(t, `{ + "id": "`+ecLemmyAnnounceID+`/reserialized", + "type": "Announce", + "actor": "`+ecLemmyGroupAPID+`", + "object": { + "audience": "`+ecLemmyGroupAPID+`", + "type": "Create", + "cc": ["`+ecLemmyGroupAPID+`/followers", "`+ecLemmyGroupAPID+`"], + "actor": "`+ecActorID+`", + "id": "`+ecActivityID+`", + "object": { + "attributedTo": "`+ecActorID+`", + "type": "Page", + "id": "`+ecNativeAPID+`", + "name": "a title the community re-rendered" + } + } + }`)) + require.NoError(t, err) + assert.Equal(t, ClassLocalActivity, identity.Class, + "reordered keys, extra addressing and a re-rendered name are what a real Announce "+ + "looks like: the echo is still ours, and a byte-equality check would let it back in") + assert.Equal(t, ecAuthorDID, identity.DID) +} diff --git a/internal/echo/echo.go b/internal/echo/echo.go index 008d373..61d6a18 100644 --- a/internal/echo/echo.go +++ b/internal/echo/echo.go @@ -3,9 +3,10 @@ // community Announce, and re-materializing it would duplicate content, // double-count votes, or — worst — read as MODERATION of our own content. // -// It is a LEAF package on purpose (store + errors only, no ingest/materialize -// import): the ingest dispatcher, the vote aggregator and the ancestor walk all -// have to ask the same question, and none of them may import each other. +// It is a LEAF package on purpose (ap + store + errors only, no +// ingest/materialize import): the ingest dispatcher, the vote aggregator and +// the ancestor walk all have to ask the same question, and none of them may +// import each other. // // An id is "ours" IFF the serving surface (internal/personas) would answer 200 // for it — ENTITY EXISTENCE, never path shape alone, because vanity origins @@ -97,8 +98,23 @@ type Classifier struct { actors store.APActors } -// New wires a Classifier. +// New wires a Classifier. All four stores are REQUIRED: this classifier's whole +// contract is that an unanswerable question surfaces as a retryable error, and +// a missing store answers it with a nil dereference on the first inbound +// activity instead. func New(opts Options) (*Classifier, error) { + if opts.Objects == nil { + return nil, errors.NewValidationError("objects", "must not be nil") + } + if opts.OutboundObjects == nil { + return nil, errors.NewValidationError("outbound_objects", "must not be nil") + } + if opts.Activities == nil { + return nil, errors.NewValidationError("activities", "must not be nil") + } + if opts.Actors == nil { + return nil, errors.NewValidationError("actors", "must not be nil") + } return &Classifier{ objects: opts.Objects, outboundObjects: opts.OutboundObjects, @@ -116,7 +132,18 @@ func New(opts Options) (*Classifier, error) { // unenumerable, so the host test is folded into the row itself — the actor's // stored NormalizedOrigin, or the stored id the object/activity rows are keyed // by. Anything else is remote content and stays remote. +// +// An id ALONE is all this can weigh, and our ids are public and derivable. When +// the id arrives inside a node the walk can read, Classify corroborates it +// against the activity we actually sent; see identifyActivity. func (c *Classifier) Identify(ctx context.Context, apID string) (Identity, error) { + return c.identify(ctx, apID, nil) +} + +// identify is Identify with the NODE the id was read from, when there is one. +// node is corroboration material, never the decision: an id that names nothing +// of ours stays not-ours whatever the node claims. +func (c *Classifier) identify(ctx context.Context, apID string, node *ap.Object) (Identity, error) { if apID == "" { return Identity{Class: ClassNone}, nil } @@ -133,25 +160,39 @@ func (c *Classifier) Identify(ctx context.Context, apID string) (Identity, error case strings.HasPrefix(parsed.Path, objectPathPrefix): return c.identifyObject(ctx, apID, strings.TrimPrefix(parsed.Path, objectPathPrefix)) case strings.HasPrefix(parsed.Path, activityPathPrefix): - return c.identifyActivity(ctx, apID, strings.TrimPrefix(parsed.Path, activityPathPrefix)) + return c.identifyActivity(ctx, apID, strings.TrimPrefix(parsed.Path, activityPathPrefix), node) case strings.HasPrefix(parsed.Path, actorPathPrefix): - return c.identifyActor(ctx, apID, normalizeHost(parsed.Host), + return c.identifyActor(ctx, apID, parsed.Scheme, normalizeHost(parsed.Host), strings.TrimPrefix(parsed.Path, actorPathPrefix)) } return Identity{Class: ClassNone}, nil } // identifyObject mirrors personas.handleObject: rest is did/collection/rkey, -// and the body is served from outbound_objects. Two tables answer for this -// route because two halves of the bridge write it — the v1 write side records -// an ap_objects mapping with origin=bridge, while an author-owned record has no -// ap_objects row at all. +// and the body is served from OUTBOUND_OBJECTS (serving.go:137) — so that table +// is the route's real oracle, and the ap_objects read is defence in depth. +// +// Both are consulted because the two rows are written by different halves of +// the bridge and neither implies the other: the enqueuer records a +// bridge-origin ap_objects mapping (enqueuer.go:136-144, 196-205) alongside the +// outbound row, but a legacy v1 write has only the mapping, and a state where +// the outbound row is missing must not make our own object answer "remote". // // The ap_objects read is keyed by the requested id, but a row alone is NOT // enough: genuine Lemmy content the bridge materialized has an ap_objects row // too (origin=fediverse), and treating that as our own echo would drop a whole // community's content. Only origin=bridge is ours. // +// RESIDUAL (accepted, not closed here): like the actor route, this weighs the +// ID and not the node it was read from — our object ids are public, so a peer +// can paint one onto a node wrapping different content and have that node +// dropped. There is no stored per-object payload to corroborate against the way +// identifyActivity has one, and the reachable harm is bounded: the announce +// path runs behind the followed-community gate, so the peer must be a community +// we follow suppressing content it chose to announce, and the bare paths +// authorize separately. Closing it needs the outbound snapshot, which belongs +// with the fetch-binding belt, not here. +// // outbound_objects is keyed by AT-URI, so a row found under the at-uri this // path spells is only ours if the row's own ap id is the id we were asked // about — otherwise any authority could borrow our path shape and have its @@ -197,7 +238,21 @@ func (c *Classifier) identifyObject(ctx context.Context, apID, rest string) (Ide // identifyActivity mirrors personas.handleActivity: the stored id carries its // own origin, so looking the FULL requested id up is both the existence test // and the authority test in one read. -func (c *Classifier) identifyActivity(ctx context.Context, apID, hash string) (Identity, error) { +// +// When the id arrives inside a node, that node is CORROBORATED against the +// activity we stored. Our activity ids are public and derivable, so the id +// alone lets a peer paint one onto a node wrapping somebody else's content and +// have it dropped as "our echo" — suppression turned into a deletion primitive. +// The stored payload is kept for byte-stable replay, which makes it exactly the +// witness this needs. +// +// The rule is CONTRADICTION DISQUALIFIES, ABSENCE DOES NOT, on IDENTIFYING +// FIELDS only — never on bytes. A community re-serializes what it announces: +// keys are reordered, addressing is added, names are re-rendered. Comparing +// bytes (or any field a re-render may touch) would refuse every genuine echo +// and re-materialize all of them, which is the bug this package exists to +// prevent. So a bare IRI, carrying neither type nor object, still classifies. +func (c *Classifier) identifyActivity(ctx context.Context, apID, hash string, node *ap.Object) (Identity, error) { if hash == "" || strings.Contains(hash, "/") { return Identity{Class: ClassNone}, nil } @@ -208,16 +263,90 @@ func (c *Classifier) identifyActivity(ctx context.Context, apID, hash string) (I } return Identity{Class: ClassNone}, fmt.Errorf("echo: outbound activity for %s: %w", apID, err) } + corroborated, err := c.corroborates(ctx, node, activity) + if err != nil || !corroborated { + return Identity{Class: ClassNone}, err + } return Identity{Class: ClassLocalActivity, DID: activity.ActorDID}, nil } +// corroborates reports whether node can be the activity we stored under that +// id. A nil node (a bare id probe) corroborates trivially — there is nothing to +// contradict. +func (c *Classifier) corroborates(ctx context.Context, node *ap.Object, activity *store.OutboundActivity) (bool, error) { + if node == nil { + return true, nil + } + // The VERB. We recorded what we sent; a node calling itself something else + // is not it — and honouring the id alone would let a peer suppress any + // Delete it likes by wearing the id of a Create we sent. + if node.Type != "" && activity.Kind != "" && !strings.EqualFold(node.Type, activity.Kind) { + return false, nil + } + + // The CARRIED OBJECT. Same id, same verb, different content is the forgery + // the verb check cannot see. The comparison is against the object our + // stored payload names, and it disqualifies only when the substituted + // object is NOT ours: swapping one of our objects for another cannot lose + // anybody's content, while swapping in a remote human's note is precisely + // the content loss dressed as echo suppression. + carried := refObjectID(node) + if carried == "" { + return true, nil + } + stored := refObjectID(parsePayload(activity.Payload)) + if stored == "" || stored == carried { + return true, nil + } + identity, err := c.identify(ctx, carried, nil) + if err != nil { + return false, err + } + return identity.Class != ClassNone, nil +} + +// parsePayload reads a stored activity back into its object form. A payload +// that will not parse yields no corroboration material rather than a verdict — +// absence, like any other missing field. +func parsePayload(payload []byte) *ap.Object { + if len(payload) == 0 { + return nil + } + parsed, err := ap.ParseObject(payload) + if err != nil { + return nil + } + return parsed +} + +// refObjectID is the id of the object an activity carries, or "" when it names +// none (including a nil activity). +func refObjectID(activity *ap.Object) string { + if activity == nil || activity.Object == nil { + return "" + } + return activity.Object.ID +} + // identifyActor mirrors personas.lookupResource: the DID is global but the // actor is not, so a row alone is not the answer — the actor must have been -// minted on the very host the id names. The comparison is against the stored +// minted on the very origin the id names. The comparison is against the stored // NormalizedOrigin in full, which makes it label-boundary-safe by // construction: neither a host we are a suffix of nor one we are a prefix of // can equal it. -func (c *Classifier) identifyActor(ctx context.Context, apID, host, rest string) (Identity, error) { +// +// The SCHEME is compared too, against the one the actor was actually minted +// under. The other two routes match a stored id exactly and so cannot be +// spoofed by re-spelling it; without this, http://our.host/ap/actor/{did} — an +// id we never mint and never serve — would classify as ours. +// +// RESIDUAL (accepted, as for identifyObject): an actor id is public, and a node +// can name one of our personas as its actor without our persona having done +// anything. An actor has no per-activity payload to corroborate against, and a +// forged actor claim is already a signature failure at the inbox for the +// top-level activity; what remains is an inner node inside an announce from a +// community we follow. +func (c *Classifier) identifyActor(ctx context.Context, apID, scheme, host, rest string) (Identity, error) { if rest == "" || strings.Contains(rest, "/") { return Identity{Class: ClassNone}, nil } @@ -228,12 +357,23 @@ func (c *Classifier) identifyActor(ctx context.Context, apID, host, rest string) } return Identity{Class: ClassNone}, fmt.Errorf("echo: actor for %s: %w", apID, err) } - if actor.NormalizedOrigin != host { + if actor.NormalizedOrigin != host || !strings.EqualFold(scheme, schemeOf(actor.ActorID)) { return Identity{Class: ClassNone}, nil } return Identity{Class: ClassLocalActor, DID: actor.DID}, nil } +// schemeOf is the URL scheme a stored actor id was minted under, or "" if the +// stored id will not parse — which fails closed, since no requested id's scheme +// can equal "". +func schemeOf(actorID string) string { + parsed, err := url.Parse(actorID) + if err != nil { + return "" + } + return parsed.Scheme +} + // normalizeHost reduces a URL authority to the authority it names — lowercase, // no trailing dot, no default port — so it can be compared with the // normalized_origin an actor was minted under. It is the read-side twin of @@ -273,9 +413,22 @@ const MaxDepth = 8 // in one is legible as its own bug. // // At each node it asks about the node's own id, then its actor, then descends -// into its object. Actor before descent is what catches an echoed vote: -// Announce{Like}'s object is the LEMMY subject, so the inner ACTOR is the only -// handle on it once the id fails. +// into its object — but ONLY through the verbs whose `object` is a PAYLOAD. +// +// This is the difference between a payload and a TARGET, and getting it wrong +// loses content in the most common interaction the product has. For Announce, +// Create, Update and Undo, `object` is what the activity carries: keep asking. +// For Like, Dislike, Delete, Flag, Block, Remove — and for anything else, +// because an unknown verb must fail toward NOT condemning an envelope — it is +// what somebody else's activity is being done TO. A Lemmy human's Like on a +// post we federated out, or a Lemmy moderator's Delete of a native postv2, has +// OUR id as its target; descending would answer "ours" for the whole envelope +// and discard every vote and every moderation action on native content, with no +// error and a counter that says the bridge is working. +// +// Nothing is lost by stopping there: the echoes these guards exist for are +// identified by the target-bearing node ITSELF — an echoed Delete by its own +// activity id, an echoed vote by its actor — never by what sits below it. // // A VOTE node inverts the first two: the voter is asked about before the // activity id. Lemmy 0.19 reconstructs the inner vote of an Announce{Undo{Like}} @@ -301,21 +454,50 @@ func (c *Classifier) Classify(ctx context.Context, envelope *ap.Object) (Identit // recoverable direction (a duplicate, never a drop). node := envelope for depth := 1; node != nil && depth <= MaxDepth; depth++ { - probes := [2]string{node.ID, actorIDOf(node)} + // The node travels with its OWN id — that id claims to name this very + // activity, so the node is what corroborates it. The actor id names a + // different entity entirely and carries no such claim. + probes := [2]probe{{id: node.ID, node: node}, {id: actorIDOf(node)}} if isVote(node) { probes[0], probes[1] = probes[1], probes[0] } - for _, probe := range probes { - identity, err := c.Identify(ctx, probe) + for _, p := range probes { + identity, err := c.identify(ctx, p.id, p.node) if err != nil || identity.Class != ClassNone { return identity, err } } + if !carriesPayload(node) { + // node.Object is this activity's TARGET, not its payload: whatever + // it names belongs to whoever the activity is being done TO. + return Identity{Class: ClassNone}, nil + } node = node.Object } return Identity{Class: ClassNone}, nil } +// carriesPayload reports whether node.Object is the thing the activity CARRIES +// (walk on) rather than the thing it acts UPON (stop). It is an allowlist, so +// an unrecognized verb stops — the recoverable direction, since descending into +// a target can drop genuine content permanently while declining to descend can +// at worst let a duplicate through. +func carriesPayload(node *ap.Object) bool { + switch node.Type { + case ap.TypeAnnounce, ap.TypeCreate, ap.TypeUpdate, ap.TypeUndo: + return true + default: + return false + } +} + +// probe is one question the walk asks: an id, plus the node that id was read +// FROM when the node is a claim about the id itself. +type probe struct { + id string + node *ap.Object +} + // isVote reports whether the node is a vote, whose VOTER identifies it. func isVote(node *ap.Object) bool { return node.Type == ap.TypeLike || node.Type == ap.TypeDislike @@ -336,6 +518,14 @@ func actorIDOf(node *ap.Object) string { // endpoint — which looks exactly like a drop site that never fires. const dropMetricPrefix = "tidepool_echo_drops_" +// dropMetricName is the metric a class is counted under. The class STRING is +// the wire and log spelling and stays hyphenated; the metric suffix is +// underscored, because a hyphen is not legal in a Prometheus metric name and +// every other counter in this repo is underscored. +func dropMetricName(class Class) string { + return dropMetricPrefix + strings.ReplaceAll(string(class), "-", "_") +} + // dropCounters is one expvar per class that MEANS "ours". ClassNone is absent // on purpose: it is the answer "this is genuine remote content", and remote // content is processed, never dropped — a counter for it could only ever @@ -348,7 +538,7 @@ var dropCounters = func() map[Class]*expvar.Int { ClassLocalActor, ClassAncestorShortCircuit, } { - counters[class] = expvar.NewInt(dropMetricPrefix + string(class)) + counters[class] = expvar.NewInt(dropMetricName(class)) } return counters }() diff --git a/internal/echo/echo_test.go b/internal/echo/echo_test.go index a2dffae..676335f 100644 --- a/internal/echo/echo_test.go +++ b/internal/echo/echo_test.go @@ -175,11 +175,18 @@ func newWorld(t *testing.T) (*Classifier, Options, *sql.DB) { _, err = outboundObjects.Tombstone(ctx, ecGoneATURI) require.NoError(t, err, "tombstone the deleted native post") + // The payload is the FULL canonical activity, the way the enqueuer stores it + // (byte-stable for replay). It is what a node claiming this id can be + // corroborated against: the id alone is public and derivable. _, err = activities.Insert(ctx, store.OutboundActivity{ ActivityID: ecActivityID, ActorDID: ecAuthorDID, Kind: "Create", - Payload: []byte(`{"type":"Create","id":"` + ecActivityID + `"}`), + Payload: []byte(`{"@context":"https://www.w3.org/ns/activitystreams",` + + `"id":"` + ecActivityID + `","type":"Create","actor":"` + ecActorID + `",` + + `"to":["https://www.w3.org/ns/activitystreams#Public"],` + + `"cc":["https://lemmy.world/c/technology"],"audience":"https://lemmy.world/c/technology",` + + `"object":{"id":"` + ecNativeAPID + `","type":"Page","attributedTo":"` + ecActorID + `"}}`), }) require.NoError(t, err, "seed the delivered activity") @@ -252,6 +259,16 @@ func TestIdentifyResolvesOurServingSurface(t *testing.T) { "deleted is still an echo, and treating it as remote content would " + "re-materialize the record we just removed", }, + { + name: "P6 a persona minted under a VANITY origin, on that origin", + apID: ecVanityOrigin + "/ap/actor/" + ecVanityDID, + wantClass: ClassLocalActor, + wantDID: ecVanityDID, + why: "decision 10 stated POSITIVELY: the origin set is unenumerable, so the host " + + "test is the actor's OWN stored NormalizedOrigin. An implementation comparing " + + "against one configured origin passes every negative case and fails this one — " + + "and it would take every vanity-origin user's traffic with it", + }, { name: "P5 soft-deleted bridge mapping is still ours", apID: ecSoftAPID, @@ -333,6 +350,22 @@ func TestIdentifyRefusesWhatWeDoNotServe(t *testing.T) { why: "host comparison is label-boundary-safe, never substring: notcoves.social " + "is a different authority that can host the same path", }, + { + name: "N4 a foreign host carrying an at-uri we DO serve (the spoof that reaches the check)", + apID: "https://notcoves.social/ap/object/" + ecAuthorDID + + "/social.coves.community.postv2/" + ecNativeRKey, + why: "outbound_objects is keyed by AT-URI, so this path finds a REAL row — the " + + "only thing standing between it and a mapped-object verdict is comparing the " + + "row's own ap id against the id we were asked about. Without that comparison " + + "any authority can wear our path shape and have its content dropped as our echo", + }, + { + name: "N4 a foreign host carrying an activity hash we minted", + apID: "https://notcoves.social/ap/activity/" + + "9f2c1e6b6d5b4a3f8c7d0e1a2b3c4d5e6f708192a3b4c5d6e7f8091a2b3c4d5e", + why: "the activity route is keyed by the FULL id for this reason: a lookup keyed " + + "on the hash alone would answer for whoever hosts it", + }, { name: "N4 a host that merely starts with ours", apID: "https://coves.social.evil.example/ap/actor/" + ecAuthorDID, @@ -368,6 +401,58 @@ func TestIdentifyRefusesWhatWeDoNotServe(t *testing.T) { } } +// TestIdentifyNormalizesTheHostTheWayServingDoes pins echo.normalizeHost, the +// read-side twin of personas' own. It is a hand-copy, and a hand-copy that +// drifts is invisible: replacing its body with `return host` leaves every other +// test in this package green while every actor id a peer spells slightly +// differently stops being recognized as ours — and an unrecognized echo is a +// duplicate, an over-normalized one is dropped genuine content. +func TestIdentifyNormalizesTheHostTheWayServingDoes(t *testing.T) { + classifier, _, _ := newWorld(t) + ctx := context.Background() + + cases := []struct { + name string + apID string + wantClass Class + why string + }{ + { + name: "the default port is the same origin", + apID: "https://coves.social:443/ap/actor/" + ecAuthorDID, + wantClass: ClassLocalActor, + why: "https://host:443 and https://host name the same authority", + }, + { + name: "the host is case-insensitive", + apID: "https://COVES.SOCIAL/ap/actor/" + ecAuthorDID, + wantClass: ClassLocalActor, + why: "DNS is case-insensitive and normalized_origin is stored lowercased", + }, + { + name: "a fully-qualified trailing dot is the same origin", + apID: "https://coves.social./ap/actor/" + ecAuthorDID, + wantClass: ClassLocalActor, + why: "the root label is implicit; a peer that spells it out names the same host", + }, + { + name: "a NON-default port is a DIFFERENT origin", + apID: "https://coves.social:8091/ap/actor/" + ecAuthorDID, + wantClass: ClassNone, + why: "the dev origin runs on :8091 and is genuinely another origin — stripping " + + "ports wholesale would let it answer for production's actors", + }, + } + + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + identity, err := classifier.Identify(ctx, tc.apID) + require.NoError(t, err) + assert.Equal(t, tc.wantClass, identity.Class, tc.why) + }) + } +} + // TestIdentifyPropagatesStoreFailures pins the FAIL-SAFE DIRECTION. A database // hiccup must surface as a RETRYABLE error so the event is redelivered: // diff --git a/internal/ingest/backfill.go b/internal/ingest/backfill.go index 4bf1b5f..f81d481 100644 --- a/internal/ingest/backfill.go +++ b/internal/ingest/backfill.go @@ -9,6 +9,7 @@ import ( "time" "tidepool/internal/ap" + "tidepool/internal/echo" "tidepool/internal/errors" "tidepool/internal/materialize" "tidepool/internal/store" @@ -43,12 +44,18 @@ type CountSeeder interface { } // BackfillOptions configures NewBackfill. Fetcher, Materializer, -// Communities, and Tombstones are required. +// Communities, Tombstones and Echo are required. type BackfillOptions struct { Fetcher BackfillFetcher Materializer Materializer Communities store.Communities Tombstones store.Tombstones + // Echo keeps the bridge's own federated content out of the walk. A + // bridged community's outbox holds OUR posts the moment a native user + // participates, and this path reaches the materializer directly — past + // the dispatcher's guard, which never sees an outbox item. Required for + // the same reason the dispatcher's is (task 17a). + Echo EchoClassifier // Seeder, when set, seeds each backfilled post's vote aggregates from // the origin's public API (config SEED_COUNTS_FROM_API). Seeder CountSeeder @@ -75,6 +82,7 @@ type Backfill struct { mat Materializer communities store.Communities tombstones store.Tombstones + classifier EchoClassifier seeder CountSeeder maxPosts int minInterval time.Duration @@ -101,6 +109,9 @@ func NewBackfill(opts BackfillOptions) (*Backfill, error) { if opts.Tombstones == nil { return nil, errors.NewValidationError("tombstones", "must not be nil") } + if opts.Echo == nil { + return nil, errors.NewValidationError("echo", "must not be nil") + } logger := opts.Logger if logger == nil { logger = slog.Default() @@ -114,6 +125,7 @@ func NewBackfill(opts BackfillOptions) (*Backfill, error) { mat: opts.Materializer, communities: opts.Communities, tombstones: opts.Tombstones, + classifier: opts.Echo, seeder: opts.Seeder, maxPosts: opts.MaxPosts, minInterval: opts.MinInterval, @@ -266,6 +278,25 @@ func (b *Backfill) materializeOutboxItem(ctx context.Context, item *ap.Object, c return false, skip(item.ID, "outbox item object has no id") } + // Our own content, walked past. It runs on the UNRESOLVED node, before + // resolveEmbedded: our object id is cross-authority with the outbox host, + // so resolving would dereference our own origin to fetch back a record we + // already hold — and MaterializePost then calls EnsureActor on its + // attributedTo BEFORE reading any mapping, minting a bridged actor for our + // own persona. A skip, never an error: our post in a community's history is + // an expected item, and failing it would leave every run reporting failures + // and the community never cleanly backfilled. + // + // The question is asked of the unwrapped OBJECT, not the outbox envelope: + // the community mints its own Announce/Create ids around our content, so + // the object is the only node in the item that can be ours — and it is what + // this path would materialize. + if ours, err := b.suppressEcho(ctx, obj); err != nil { + return false, err + } else if ours { + return false, skip(obj.ID, "outbox item is our own federated content") + } + // Same funnel rules as live deliveries: never resurrect deleted // content, trust embedded bodies only on the outbox host's authority. // The walk reads markers in the backfilled community's scope — its own, @@ -301,6 +332,31 @@ func (b *Backfill) materializeOutboxItem(ctx context.Context, item *ap.Object, c } } +// suppressEcho reports whether an outbox node is content the bridge itself +// federated, counting and logging the drop when it is. +// +// This is the same guard the dispatcher runs, at the only other place content +// enters: a backfill has no envelope and no signer, so nothing upstream of here +// can ask the question. An error is propagated rather than resolved into a +// verdict — the run is resumable and re-walks, where "not ours" would +// re-materialize our own post and "ours" would drop a community's history. +func (b *Backfill) suppressEcho(ctx context.Context, node *ap.Object) (bool, error) { + identity, err := b.classifier.Classify(ctx, node) + if err != nil { + return false, fmt.Errorf("ingest: backfill echo check for %s: %w", node.ID, err) + } + if identity.Class == echo.ClassNone { + return false, nil + } + echo.CountDrop(identity.Class) + // INFO, like every other echo drop: this is the only record that a piece of + // a community's history was deliberately walked past. + b.logger.Info("backfill item skipped: our own federated content", + "object", node.ID, "class", string(identity.Class), + "did", identity.DID, "at_uri", identity.ATURI) + return true, nil +} + // seedCounts imports a backfilled post's historical vote counts (task 07). // Best-effort: failures are logged and never affect the run — a post with a // zero score is strictly better than no post. Warn (matching the @@ -336,6 +392,17 @@ func (b *Backfill) backfillReplies(ctx context.Context, post *ap.Object, communi } count++ note := *item + // Before resolveEmbedded, for the same reason as the outbox item: a + // reply of ours must not be dereferenced from our own origin, and a + // native comment in a Lemmy thread is exactly what this collection + // holds once a Coves user replies. + if ours, err := b.suppressEcho(ctx, ¬e); err != nil { + b.logger.Warn("backfill reply echo check failed", + "post", post.ID, "reply", note.ID, "error", err) + return nil + } else if ours { + return nil + } resolved, err := b.resolveEmbedded(ctx, ¬e, repliesIRI) if err != nil { b.logger.Info("backfill reply skipped", "post", post.ID, "error", err.Error()) diff --git a/internal/ingest/backfill_test.go b/internal/ingest/backfill_test.go index 094622c..11fae46 100644 --- a/internal/ingest/backfill_test.go +++ b/internal/ingest/backfill_test.go @@ -28,7 +28,12 @@ func newBackfill(t *testing.T, h *harness, maxPosts int) *Backfill { Materializer: h.mat, Communities: h.communities, Tombstones: h.tombstones, - MaxPosts: maxPosts, + // The same guard the dispatcher runs, and required for the same reason: + // a community's outbox carries OUR federated content once native users + // participate, and this path reaches the materializer with no envelope + // in front of it. + Echo: h.classifier, + MaxPosts: maxPosts, }) require.NoError(t, err) return b @@ -134,6 +139,7 @@ func TestBackfillSeedsVoteCounts(t *testing.T) { Communities: h.communities, Tombstones: h.tombstones, Seeder: seeder, + Echo: h.classifier, MaxPosts: 10, }) require.NoError(t, err) diff --git a/internal/ingest/echo_backfill_test.go b/internal/ingest/echo_backfill_test.go new file mode 100644 index 0000000..3363882 --- /dev/null +++ b/internal/ingest/echo_backfill_test.go @@ -0,0 +1,224 @@ +package ingest + +import ( + "context" + "encoding/json" + "net/http" + "sync/atomic" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "tidepool/internal/consume" + "tidepool/internal/echo" + "tidepool/internal/errors" + "tidepool/internal/materialize" + "tidepool/internal/outbound" + "tidepool/internal/personas" + "tidepool/internal/store" +) + +// THE BACKFILL BYPASSES BOTH GUARDS. +// +// Backfill.materializeOutboxItem and backfillReplies call MaterializePost / +// MaterializeComment DIRECTLY — past suppressEcho (no envelope reaches them) +// and past materializeContent's legacy origin=bridge check (that funnel is not +// on this path at all). A bridged community's outbox contains OUR federated +// content the moment native users participate, and any re-trigger re-walks it. +// +// Two things then happen, both silent: +// +// - MaterializePost calls EnsureActor(attributedTo) BEFORE it reads any +// mapping, so our own persona acquires a bridged_actors row — the mint +// oracle echo_bare_test.go asserts can never happen, reached by an ordinary +// admin action; +// - PutMapping rewrites our bridge-origin mapping with the default +// origin=fediverse, which makes the classifier answer ClassNone for that id +// from then on. The echo guard disables itself for the post, permanently, +// and nothing anywhere reports it. +// +// The package doc claims the invariant holds "on any path, however it arrives". +// This is a path. +const ( + bfUserOrigin = "https://coves.social" + bfAuthorDID = "did:plc:bfbackfillauthor01" + bfActorID = bfUserOrigin + "/ap/actor/" + bfAuthorDID + bfPostRKey = "3lzbackfill0001" + bfPostATURI = "at://" + bfAuthorDID + "/social.coves.community.postv2/" + bfPostRKey + bfPostAPID = bfUserOrigin + "/ap/object/" + bfAuthorDID + + "/social.coves.community.postv2/" + bfPostRKey + bfPostCID = "bafyreib2rxk3rybk3aobmv5cjuql3bm2twh4jo5uxgf5kpqrsqxi3jgxte" +) + +// TestBackfillDoesNotReMaterializeOurOwnContent is H2. +// +// The community's outbox carries BOTH kinds of history — a post we federated +// out and a post a Lemmy human wrote — because that is what a bridged +// community's outbox looks like once anyone native has posted. The backfill +// must walk past the first and materialize the second. +func TestBackfillDoesNotReMaterializeOurOwnContent(t *testing.T) { + h := newHarness(t) + ctx := context.Background() + h.subscribeTechnology() + h.serveLemmyWorldContent() + + // Our origin serves the object for real: the outbox item is cross-authority + // with the outbox host, so the walk RE-FETCHES it rather than trusting the + // embedded body — and that fetch succeeds. The bypass is not a fixture + // artifact. + userOrigin, err := personas.New(personas.Options{ + DB: h.db, + Custodian: h.custodian, + UserOrigin: bfUserOrigin, + }) + require.NoError(t, err) + selfFetches := &atomic.Int64{} + h.mux.Handle("/ap/", http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + selfFetches.Add(1) + userOrigin.ServeHTTP(w, r) + })) + _, err = userOrigin.CreateActorForDID(ctx, bfAuthorDID, "bfauthor.coves.social") + require.NoError(t, err) + + communityDID := testDIDFor("technology", "lemmy.world") + _, err = store.NewOutboundObjects(h.db).Upsert(ctx, store.OutboundObject{ + ATURI: bfPostATURI, + APObjectID: bfPostAPID, + LastCID: bfPostCID, + LastRev: "3lzbfrev000001", + CommunityDID: communityDID, + CommunityAPID: groupID, + TranslatedSnapshot: bfSnapshot(t), + }) + require.NoError(t, err) + + enqueuer, err := outbound.NewEnqueuer(outbound.EnqueuerOptions{ + DB: h.db, + Translator: outbound.NewTranslator(bfUserOrigin), + Inboxes: outbound.NewInboxResolver(h.client, time.Minute), + Actors: store.NewAPActors(h.db), + UserOrigin: bfUserOrigin, + }) + require.NoError(t, err) + enqueueAs(t, h.db, enqueuer, bfAuthorDID, consume.PostIntent{ + Op: "create", + ATURI: bfPostATURI, + ID: consume.ActivityID(bfUserOrigin, bfPostATURI, "create", 0), + CommunityAPID: groupID, + Snapshot: bfSnapshot(t), + }) + ourMapping, err := h.objects.GetByAPID(ctx, bfPostAPID) + require.NoError(t, err, "precondition: our post is mapped bridge-origin") + require.Equal(t, store.OriginBridge, ourMapping.Origin) + + // The community's outbox: our federated post, then a genuine Lemmy post. + h.serveObject("/c/technology/outbox", map[string]any{ + "@context": "https://www.w3.org/ns/activitystreams", + "type": "OrderedCollection", + "id": groupID + "/outbox", + "totalItems": 2, + "orderedItems": []any{bfAnnouncedPage(bfPostAPID, bfActorID), bfAnnouncedPage(pageID, personID)}, + }) + + before := dropSnapshot() + opsBefore := len(h.firehoseOps()) + selfFetches.Store(0) + + community, err := h.communities.GetByAPGroupID(ctx, groupID) + require.NoError(t, err) + backfill := newBackfill(t, h, 10) + // Non-fatal: today this run FAILS the item (it tries to write a postv2 into + // the native author's repo, which the bridge does not host), and the state + // assertions below are the ones that say what actually happened. A skipped + // item is not a failed one — after the fix the whole run is clean. + assert.NoError(t, backfill.Run(ctx, community, true), + "walking past our own content is a SKIP, not an error: an errored item leaves the "+ + "community's backfill permanently unresumable-looking and retries forever") + + // --- Our own post: untouched. --- + _, err = h.actors.GetByAPActorID(ctx, bfActorID) + assert.True(t, errors.IsNotFound(err), + "NO bridged actor may be minted for our own persona: MaterializePost calls "+ + "EnsureActor before it reads any mapping, so an ordinary re-backfill is a mint "+ + "oracle for every native user who has posted (err=%v)", err) + + after, err := h.objects.GetByAPID(ctx, bfPostAPID) + require.NoError(t, err) + assert.Equal(t, store.OriginBridge, after.Origin, + "our mapping must stay origin=bridge: re-materializing rewrites it to the default "+ + "origin=fediverse, and the classifier then answers ClassNone for this id forever "+ + "— the echo guard silently disabling itself for the post") + assert.Equal(t, ourMapping.ATURI, after.ATURI, "and it must still name the record we federated") + assert.Equal(t, int64(0), selfFetches.Load(), + "our own object must not even be dereferenced: the backfill has no envelope, so the "+ + "id itself is the only thing that can stop it") + + // --- The genuine Lemmy post in the SAME outbox: backfilled normally. --- + lemmyMapping, err := h.objects.GetByAPID(ctx, pageID) + require.NoError(t, err, + "the community's real history must still backfill: suppression here must key on "+ + "whose content it is, never on 'this outbox contains something of ours'") + assert.Equal(t, store.OriginFediverse, lemmyMapping.Origin) + assert.Equal(t, materialize.CollectionPostV2, lemmyMapping.Collection) + assert.Greater(t, len(h.firehoseOps()), opsBefore, + "and it commits: a backfill that drops everything is the failure mode this control exists for") + + // Attribution: the drop is counted like every other, under the class that + // identified it (the object id is ours). + for _, class := range echoClasses { + want := before[class] + if class == echo.ClassMappedObject { + want++ + } + assert.Equal(t, want, echo.Drops(class), + "counter %q after a backfill drop: an uncounted drop on this path is exactly how "+ + "the bypass stayed invisible", class) + } +} + +// bfAnnouncedPage is one outbox item in Lemmy's shape: Announce{Create{Page}}. +func bfAnnouncedPage(pageAPID, author string) map[string]any { + return map[string]any{ + "id": pageAPID + "/announce", + "type": "Announce", + "actor": groupID, + "to": []any{"https://www.w3.org/ns/activitystreams#Public"}, + "object": map[string]any{ + "id": pageAPID + "/create", + "type": "Create", + "actor": author, + "object": map[string]any{ + "id": pageAPID, + "type": "Page", + "attributedTo": author, + "audience": groupID, + "to": []any{"https://www.w3.org/ns/activitystreams#Public", groupID}, + "name": "a post in the community's history", + "content": "

backfilled

", + "published": "2026-08-13T09:00:00.000000Z", + }, + }, + } +} + +func bfSnapshot(t *testing.T) []byte { + t.Helper() + snap, err := json.Marshal(map[string]any{ + "atUri": bfPostATURI, + "cid": bfPostCID, + "rev": "3lzbfrev000001", + "collection": "social.coves.community.postv2", + "record": map[string]any{ + "$type": "social.coves.community.postv2", + "community": testDIDFor("technology", "lemmy.world"), + "title": "a post in the community's history", + "content": "backfilled", + "createdAt": "2026-08-13T09:00:00.000Z", + }, + "communityApId": groupID, + }) + require.NoError(t, err) + return snap +} diff --git a/internal/ingest/echo_failsafe_test.go b/internal/ingest/echo_failsafe_test.go new file mode 100644 index 0000000..8bc4ba7 --- /dev/null +++ b/internal/ingest/echo_failsafe_test.go @@ -0,0 +1,160 @@ +package ingest + +import ( + "context" + stderrors "errors" + "log/slog" + "net/http" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "tidepool/internal/ap" + "tidepool/internal/echo" + "tidepool/internal/errors" +) + +// THE DISPATCHER'S HALF OF THE FAIL-SAFE DIRECTION. +// +// internal/echo pins that a store failure surfaces as a retryable error, and +// internal/votes pins what its probe does with one. The dispatcher's own +// contract — what happens to the EVENT when the classification cannot be made — +// was the asymmetry: only an injectable classifier can exercise it, and until +// EchoClassifier became an interface there was no way in. +// +// Both wrong answers are available and both are permanent: +// +// - treat the failure as "not ours" and the echo re-materializes (a +// duplicate, recoverable); +// - treat it as "ours" and genuine Lemmy content is dropped for good. +// +// The only honest answer is neither: leave the event alone and try again. So +// the event must NOT be poisoned, must NOT be marked processed, and the failing +// attempt must leave no partial work behind — because a retry that finds +// half-applied state is not a retry. + +// failingClassifier is a classifier that cannot answer. +type failingClassifier struct{ err error } + +func (f failingClassifier) Classify(context.Context, *ap.Object) (echo.Identity, error) { + return echo.Identity{Class: echo.ClassNone}, f.err +} + +// swapClassifier rebuilds the dispatcher (and the queue drain() pumps) with a +// different echo classifier, mirroring swapHandlerStores. +// +// The retry schedule is deliberately SLOW, not fast. drain() pumps until the +// queue is empty, so a short backoff lets one drain re-claim the same event +// until it hits the attempt cap and poisons — which is the queue's correct +// universal rule for a retryable error (exempting the classifier would let a +// permanently failing one spin forever), but it would make this test a race +// between the backoff and the drain loop. An hour of backoff bounds the failing +// phase to exactly ONE attempt, and the redelivery is brought forward +// explicitly below. MaxAttempts is raised as a second, independent guard: even +// if someone later shortens the delay, the cap is out of reach. +func (h *harness) swapClassifier(classifier EchoClassifier) { + h.t.Helper() + handler, err := NewHandler(HandlerOptions{ + Materializer: h.mat, + Fetcher: h.client, + Objects: h.objects, + Actors: h.actors, + Communities: h.communities, + Tombstones: h.tombstones, + Records: h.manager, + Votes: h.votes, + Backfill: h.backfills, + Echo: classifier, + ServiceActorID: h.service.ID, + Logger: slog.New(slog.NewTextHandler(h.logs, &slog.HandlerOptions{Level: slog.LevelDebug})), + }) + require.NoError(h.t, err) + h.handler = handler + queue, err := NewQueue(QueueOptions{ + Events: h.events, + Processor: handler, + Workers: 1, + MaxAttempts: 20, + RetryBaseDelay: time.Hour, + Lease: time.Minute, + }) + require.NoError(h.t, err) + h.queue = queue +} + +// TestEchoClassificationFailureRetriesAndRecovers: an announced activity whose +// classification fails is left for the next attempt, and the next attempt +// completes it. +func TestEchoClassificationFailureRetriesAndRecovers(t *testing.T) { + h := newHarness(t) + ctx := context.Background() + group := h.subscribeTechnology() + h.serveLemmyWorldContent() + + boom := stderrors.New("connection reset by peer") + h.swapClassifier(failingClassifier{err: boom}) + dropsBefore := dropSnapshot() + opsBefore := len(h.firehoseOps()) + + announce := loadFixture(t, "announce_create_page_lemmy_world.json") + announceID := announce["id"].(string) + require.Equal(t, http.StatusAccepted, h.deliver(group, announce)) + h.drain() + + // The event survives, unfinished. + event, err := h.events.GetEvent(ctx, announceID) + require.NoError(t, err) + assert.Nil(t, event.FailedAt, + "a classification that could not be MADE is not a verdict about the activity: "+ + "poisoning the event throws away genuine content over a database hiccup") + assert.Nil(t, event.ProcessedAt, + "and it must not be marked processed either — nothing was decided") + assert.Equal(t, 1, event.Attempts, + "exactly one attempt: the backoff parks the event rather than spinning it against a "+ + "store that is still down") + assert.Contains(t, event.Error, boom.Error(), + "the transient cause is kept on the row, or the retry is undiagnosable") + + // The failing attempt left NOTHING behind. + _, err = h.objects.GetByAPID(ctx, pageID) + assert.True(t, errors.IsNotFound(err), + "no mapping may be written on a failed classification: the guard runs BEFORE "+ + "materialization precisely so a retry starts from a clean slate (err=%v)", err) + tombstoned, err := h.tombstones.ExistsFor(ctx, pageID, groupID) + require.NoError(t, err) + assert.False(t, tombstoned, "and no tombstone marker") + assert.Equal(t, opsBefore, len(h.firehoseOps()), "and no commit") + for _, class := range echoClasses { + assert.Equal(t, dropsBefore[class], echo.Drops(class), + "and no drop counter (%s): a failure is not a suppression, and counting it as "+ + "one would hide the outage inside the metric that reports on the guard", class) + } + + // The retry is the whole point: with the classifier healthy again, the SAME + // event completes. A contract that only says "do not poison" would be + // satisfied by an event that never drains. + // + // The backoff window belongs to the queue, so rather than sleep through it, + // bring the redelivery forward — the clock is not what is under test, the + // retry is. (CURRENT_TIMESTAMP, because the claim predicate reads the + // DATABASE clock, not the process's.) + _, err = h.db.ExecContext(ctx, + `UPDATE inbox_events SET next_attempt_at = CURRENT_TIMESTAMP WHERE activity_id = $1`, + announceID) + require.NoError(t, err) + + h.swapClassifier(h.classifier) + h.drain() + + event, err = h.events.GetEvent(ctx, announceID) + require.NoError(t, err) + assert.NotNil(t, event.ProcessedAt, + "the redelivery completes: the event was PARKED by the failure, not lost to it") + mapping, err := h.objects.GetByAPID(ctx, pageID) + require.NoError(t, err, + "and the genuine Lemmy post finally materializes — which is what makes the retry "+ + "the recoverable direction") + assert.False(t, mapping.IsDeleted()) +} diff --git a/internal/ingest/echo_target_test.go b/internal/ingest/echo_target_test.go new file mode 100644 index 0000000..95f81bc --- /dev/null +++ b/internal/ingest/echo_target_test.go @@ -0,0 +1,315 @@ +package ingest + +import ( + "context" + "encoding/json" + "net/http" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "tidepool/internal/consume" + "tidepool/internal/echo" + "tidepool/internal/errors" + "tidepool/internal/materialize" + "tidepool/internal/outbound" + "tidepool/internal/store" +) + +// TARGET vs PAYLOAD — the shape every other test in this suite missed. +// +// An activity's `object` means two different things depending on the verb: +// +// - for the WRAPPER verbs (Announce, Create, Update, Undo) it is the PAYLOAD: +// the activity or object being carried, and the thing the envelope is +// really about; +// - for Like, Dislike, Delete, Flag, Block and Remove it is the TARGET: what +// somebody else's activity is being done TO. +// +// A walk that descends unconditionally reaches the TARGET of a genuine remote +// activity and answers "ours" for the whole envelope — because the target is +// ours. That is the single most common interaction in the product: a fediverse +// user acting on content a Coves user wrote. +// +// Every "genuine remote traffic survives" control we wrote before this one +// targets LEMMY-origin content (W5 votes on a lemmy post, the mod removal of a +// lemmy page). They prove remote traffic survives only when it never touches our +// content, which is exactly the case that cannot break. These tests construct +// the case that can. +// +// The echoes the guards exist for do NOT need the descent: an echoed Delete is +// ours by its OWN activity id at depth 2, and an echoed vote by its ACTOR at +// depth 2. Nothing below a target-bearing verb may decide the envelope. +const ( + tgUserOrigin = "https://coves.social" + tgAuthorDID = "did:plc:tgtargetauthor0001" + tgActorID = tgUserOrigin + "/ap/actor/" + tgAuthorDID + tgPostRKey = "3lztargetpost01" + tgPostATURI = "at://" + tgAuthorDID + "/social.coves.community.postv2/" + tgPostRKey + tgPostAPID = tgUserOrigin + "/ap/object/" + tgAuthorDID + + "/social.coves.community.postv2/" + tgPostRKey + tgPostCID = "bafyreib2rxk3rybk3aobmv5cjuql3bm2twh4jo5uxgf5kpqrsqxi3jgxte" + + // The fediverse humans acting on it: a voter and a moderator, both real + // people on lemmy.world with no ap_actors row anywhere. + tgVoter = personID +) + +// targetWorld is a native post federated into the bridged community, with the +// acceptance the admission wrote — the state a remote vote or removal arrives +// into. +type targetWorld struct { + group *remoteActor + voter *remoteActor + communityDID string + digestRKey string +} + +func setupTargetWorld(t *testing.T, h *harness) targetWorld { + t.Helper() + ctx := context.Background() + group := h.subscribeTechnology() + h.serveLemmyWorldContent() + voter := h.newRemoteActor(tgVoter, person(tgVoter, "LeftLeaningFreedomFighters", nil)) + communityDID := testDIDFor("technology", "lemmy.world") + + _, err := store.NewAPActors(h.db).Create(ctx, store.APActor{ + DID: tgAuthorDID, + Kind: store.ActorTypePerson, + ActorID: tgActorID, + NormalizedOrigin: "coves.social", + LocalPart: "tgauthor", + RSAKeySealed: []byte{0x01, 0x02, 0x03}, + RSAKeyVersion: 1, + PublicKeyPEM: "-----BEGIN PUBLIC KEY-----\nMIIB\n-----END PUBLIC KEY-----\n", + }) + require.NoError(t, err) + _, err = store.NewOutboundObjects(h.db).Upsert(ctx, store.OutboundObject{ + ATURI: tgPostATURI, + APObjectID: tgPostAPID, + LastCID: tgPostCID, + LastRev: "3lztgrev000001", + CommunityDID: communityDID, + CommunityAPID: groupID, + TranslatedSnapshot: tgSnapshot(t), + }) + require.NoError(t, err) + + enqueuer, err := outbound.NewEnqueuer(outbound.EnqueuerOptions{ + DB: h.db, + Translator: outbound.NewTranslator(tgUserOrigin), + Inboxes: outbound.NewInboxResolver(h.client, time.Minute), + Actors: store.NewAPActors(h.db), + UserOrigin: tgUserOrigin, + }) + require.NoError(t, err) + enqueueAs(t, h.db, enqueuer, tgAuthorDID, consume.PostIntent{ + Op: "create", + ATURI: tgPostATURI, + ID: consume.ActivityID(tgUserOrigin, tgPostATURI, "create", 0), + CommunityAPID: groupID, + Snapshot: tgSnapshot(t), + }) + + // The mapping carries its community, the way 17c will leave it. + // + // SECOND LATENT BUG, worth recording: today's enqueuer writes no + // community_did (enqueuer.go:196-205), and CommunityDIDOf then falls + // through to reading the postv2 out of the AUTHOR's repo — which the bridge + // does not host. So an announced vote on a native post fails + // subjectBelongsToCommunity, and an announced removal fails + // authorizeAnnouncedContentDelete, INDEPENDENTLY of the classifier. Both + // paths need this column, so the fixture supplies it; otherwise these tests + // would be unsatisfiable even after the classifier is fixed. + mapping, err := h.objects.GetByAPID(ctx, tgPostAPID) + require.NoError(t, err) + require.Equal(t, store.OriginBridge, mapping.Origin) + mapping.CommunityDID = communityDID + _, err = h.objects.PutMapping(ctx, *mapping) + require.NoError(t, err) + + digest := testDigestRKey(tgPostATURI) + _, err = h.manager.PutRecord(ctx, communityDID, materialize.CollectionAcceptance, digest, + map[string]any{ + "$type": materialize.CollectionAcceptance, + "subject": map[string]any{"uri": tgPostATURI, "cid": tgPostCID}, + "createdAt": "2026-08-13T09:00:00.000Z", + }) + require.NoError(t, err) + + return targetWorld{group: group, voter: voter, communityDID: communityDID, digestRKey: digest} +} + +func tgSnapshot(t *testing.T) []byte { + t.Helper() + snap, err := json.Marshal(map[string]any{ + "atUri": tgPostATURI, + "cid": tgPostCID, + "rev": "3lztgrev000001", + "collection": "social.coves.community.postv2", + "record": map[string]any{ + "$type": "social.coves.community.postv2", + "community": testDIDFor("technology", "lemmy.world"), + "title": "A native post the fediverse votes on and moderates", + "content": "the target, not the payload", + "createdAt": "2026-08-13T09:00:00.000Z", + }, + "communityApId": groupID, + }) + require.NoError(t, err) + return snap +} + +// liveVoteRows counts the voter's live (non-undone) vote_events rows. +func liveVoteRows(t *testing.T, h *harness, voter string) int { + t.Helper() + var n int + require.NoError(t, h.db.QueryRow( + `SELECT COUNT(*) FROM vote_events WHERE voter_ap_id = $1 AND NOT undone`, voter).Scan(&n)) + return n +} + +// TestFediverseVotesOnOurFederatedPostAreCounted: the highest-volume genuine +// interaction in the product. A Lemmy human upvotes a post a Coves user wrote; +// the vote's TARGET is our object, its actor and activity id are theirs. +// +// If the envelope walk descends into that target, every vote on every native +// post is discarded and their tallies sit at zero forever — with no error, no +// dead letter, and a counter that says the bridge is working. +func TestFediverseVotesOnOurFederatedPostAreCounted(t *testing.T) { + h := newHarness(t) + world := setupTargetWorld(t, h) + withRealAggregator(t, h) + before := dropSnapshot() + + // Each shape is its own subtest: they run in sequence (an Undo needs a live + // vote), but a failure in one must not hide the others — the four shapes + // reach the guard by four different routes. + + t.Run("announced Like (depth 3)", func(t *testing.T) { + require.Equal(t, http.StatusAccepted, h.deliver(world.group, + echoAnnounce("https://lemmy.world/activities/announce/like/tg-1", map[string]any{ + "id": "https://lemmy.world/activities/like/tg-1", + "type": "Like", + "actor": tgVoter, + "object": tgPostAPID, + "audience": groupID, + }))) + h.drain() + + assert.Equal(t, 1, voteRows(t, h, tgVoter), + "a Lemmy human's upvote on a NATIVE post must be recorded: our id is the TARGET "+ + "of their activity, not an activity of ours") + up, _, found := aggregateOf(t, h, tgPostAPID) + assert.True(t, found, "the aggregate must exist — this is the post's only score") + assert.Equal(t, 1, up) + }) + + t.Run("announced Undo{Like} (depth 4)", func(t *testing.T) { + require.Equal(t, http.StatusAccepted, h.deliver(world.group, + echoAnnounce("https://lemmy.world/activities/announce/undo/tg-2", map[string]any{ + "id": "https://lemmy.world/activities/undo/tg-2", + "type": "Undo", + "actor": tgVoter, + "object": map[string]any{ + "id": "https://lemmy.world/activities/like/tg-regenerated", + "type": "Like", + "actor": tgVoter, + "object": tgPostAPID, + "audience": groupID, + }, + }))) + h.drain() + // Asserted on the ROW, not only the total: a zero aggregate is also what + // a vote that was never counted looks like, and this subtest is about + // the retraction landing — not about the sum happening to be right. + assert.GreaterOrEqual(t, voteRows(t, h, tgVoter), 1, + "their vote must be on record before it can be withdrawn") + assert.Equal(t, 0, liveVoteRows(t, h, tgVoter), + "and the Undo must mark it undone — depth 4 (Announce{Undo{Like}}) is where "+ + "Lemmy sends retractions") + up, _, _ := aggregateOf(t, h, tgPostAPID) + assert.Equal(t, 0, up, + "a vote that can be cast but not withdrawn leaves a score nobody can correct") + }) + + t.Run("bare Like", func(t *testing.T) { + require.Equal(t, http.StatusAccepted, h.deliver(world.voter, map[string]any{ + "id": "https://lemmy.world/activities/like/tg-bare", + "type": "Like", + "actor": tgVoter, + "object": tgPostAPID, + })) + h.drain() + up, _, _ := aggregateOf(t, h, tgPostAPID) + assert.Equal(t, 1, up, "a bare Like on our object is still their vote") + }) + + t.Run("bare Undo{Like} through Process's guard", func(t *testing.T) { + require.Equal(t, http.StatusAccepted, h.deliver(world.voter, map[string]any{ + "id": "https://lemmy.world/activities/undo/tg-bare", + "type": "Undo", + "actor": tgVoter, + "object": map[string]any{ + "id": "https://lemmy.world/activities/like/tg-bare", + "type": "Like", + "actor": tgVoter, + "object": tgPostAPID, + }, + })) + h.drain() + up, _, _ := aggregateOf(t, h, tgPostAPID) + assert.Equal(t, 0, up, + "the bare Undo branch runs the classifier BEFORE dispatch, so a walk that "+ + "descends into the vote's target strands every retraction on our own content") + }) + + // The mirror-image pin: none of this is an echo, so no counter may move. + for _, class := range echoClasses { + assert.Equal(t, before[class], echo.Drops(class), + "no echo counter may move for a remote actor acting on OUR content (%s): a drop "+ + "here is invisible — the votes simply never appear", class) + } +} + +// TestModeratorRemovalOfNativePostIsHonored: a Lemmy moderator removing a post +// a Coves user wrote. The Delete's TARGET is our postv2; its actor is the +// moderator and its id is the community's. +// +// This is the transition 17c exists to deliver. If the walk descends into the +// target, the removal is dropped as an echo and no native post can ever be +// moderated by the community hosting it — while the M1 case (our OWN delete +// coming home) must stay dropped. The difference is target versus payload, and +// nothing else. +func TestModeratorRemovalOfNativePostIsHonored(t *testing.T) { + h := newHarness(t) + ctx := context.Background() + world := setupTargetWorld(t, h) + before := dropSnapshot() + + reason := "off topic for this community" + h.announceDeleteWithSummary(world.group, + "https://lemmy.world/activities/announce/delete/tg-removal", tgPostAPID, &reason) + + removal, _, err := h.manager.GetRecord(ctx, + world.communityDID, materialize.CollectionRemoval, world.digestRKey) + require.NoError(t, err, + "a moderator's removal of a NATIVE post must be honored: our id is the TARGET of "+ + "the community's activity, and moderating bridged content is the whole point "+ + "of the acceptance model") + assert.Equal(t, reason, removal["reason"]) + assert.Equal(t, "moderator-discretion", removal["code"]) + + _, _, err = h.manager.GetRecord(ctx, + world.communityDID, materialize.CollectionAcceptance, world.digestRKey) + assert.True(t, errors.IsNotFound(err), + "and the acceptance is withdrawn in the same commit (err=%v)", err) + + for _, class := range echoClasses { + assert.Equal(t, before[class], echo.Drops(class), + "no echo counter may move for genuine moderation (%s): counting it here would "+ + "also be the only evidence anyone ever sees of the content it swallowed", class) + } +} diff --git a/internal/ingest/handler.go b/internal/ingest/handler.go index b09bbf9..b03266d 100644 --- a/internal/ingest/handler.go +++ b/internal/ingest/handler.go @@ -61,6 +61,17 @@ type Fetcher interface { FetchObjectSameAuthority(ctx context.Context, iri string) (*ap.Object, error) } +// EchoClassifier answers whether an inbound envelope is the bridge's own +// traffic coming home (task 17a). *echo.Classifier satisfies it. +// +// It is an INTERFACE, like votes.VoterProbe, for one reason: the fail-safe this +// dispatcher owes — a classification that CANNOT be made must retry, never +// poison the event and never materialize — is only exercisable by injecting a +// classifier that fails, and a concrete type leaves that contract untestable. +type EchoClassifier interface { + Classify(ctx context.Context, envelope *ap.Object) (echo.Identity, error) +} + // Backfiller is notified when a community's Follow is accepted (the // backfill trigger). *Backfill implements it; tests inject recorders. type Backfiller interface { @@ -90,7 +101,7 @@ type HandlerOptions struct { Backfill Backfiller // Echo classifies inbound ids against the bridge's own serving surface so // an activity we sent never re-enters as content (task 17a). - Echo *echo.Classifier + Echo EchoClassifier // ServiceActorID is the bridge's own AP actor id; Accepts must wrap a // Follow issued by it. ServiceActorID string @@ -111,7 +122,7 @@ type Handler struct { records RecordGetter votes VoteAggregator backfill Backfiller - echo *echo.Classifier + classifier EchoClassifier echoLog *ratelimit.Sampler serviceID string logger *slog.Logger @@ -143,6 +154,15 @@ func NewHandler(opts HandlerOptions) (*Handler, error) { if opts.Votes == nil { return nil, errors.NewValidationError("votes", "must not be nil") } + // REQUIRED, exactly as votes.NewAggregator requires its voter probe. The + // same guard cannot be mandatory on one path and optional on another: a + // dispatcher without it re-materializes our own content, mints bridged + // actors for our own personas and self-moderates, silently, in whichever + // binary forgot to pass it — which is how it was left out of production the + // first time. + if opts.Echo == nil { + return nil, errors.NewValidationError("echo", "must not be nil") + } if opts.ServiceActorID == "" { return nil, errors.NewValidationError("service_actor_id", "must not be empty") } @@ -160,7 +180,7 @@ func NewHandler(opts HandlerOptions) (*Handler, error) { records: opts.Records, votes: opts.Votes, backfill: opts.Backfill, - echo: opts.Echo, + classifier: opts.Echo, echoLog: ratelimit.NewSampler(echoDropLogInterval), serviceID: opts.ServiceActorID, logger: logger, @@ -313,25 +333,20 @@ func (h *Handler) handleAnnounce(ctx context.Context, announce *ap.Object, signe } } -// suppressEcho drops an announced activity the bridge itself sent. Returning -// our own content back into materialization duplicates it, double-counts our -// own votes, and — worst — reads as MODERATION of our own records. +// suppressEcho drops an activity the bridge itself sent, whether it arrived +// announced by a community or delivered bare. Letting our own content back into +// materialization duplicates it, double-counts our own votes, and — worst — +// reads as MODERATION of our own records. // // It returns a skip when the envelope resolves to one of our own entities, nil // when it is genuine remote traffic, and the classifier's error otherwise. A // failed lookup is never a verdict: calling it "not ours" re-materializes the // echo, calling it "ours" drops real Lemmy content permanently, and only the // retry the wrapped error buys is honest. -// -// A nil classifier means echo suppression is not configured; the handler then -// behaves exactly as it did before task 17a rather than refusing traffic. -func (h *Handler) suppressEcho(ctx context.Context, announceID string, envelope *ap.Object) error { - if h.echo == nil { - return nil - } - identity, err := h.echo.Classify(ctx, envelope) +func (h *Handler) suppressEcho(ctx context.Context, activityID string, envelope *ap.Object) error { + identity, err := h.classifier.Classify(ctx, envelope) if err != nil { - return fmt.Errorf("ingest: echo classification for %s: %w", announceID, err) + return fmt.Errorf("ingest: echo classification for %s: %w", activityID, err) } if identity.Class == echo.ClassNone { return nil @@ -343,13 +358,13 @@ func (h *Handler) suppressEcho(ctx context.Context, announceID string, envelope // false-positive detector, and genuine community content dropped as an // echo is invisible at Debug in production. if h.echoLog.Allow(time.Now()) { - h.logger.Info("dropped an announced echo of our own activity", - "announce_id", announceID, + h.logger.Info("dropped an echo of our own activity", + "activity_id", activityID, "class", string(identity.Class), "did", identity.DID, "at_uri", identity.ATURI) } - return skip(announceID, "echo of our own "+string(identity.Class)) + return skip(activityID, "echo of our own "+string(identity.Class)) } // handleBareCreateUpdate processes a Create/Update delivered directly by a diff --git a/internal/ingest/handler_test.go b/internal/ingest/handler_test.go index 7098e23..1e2514a 100644 --- a/internal/ingest/handler_test.go +++ b/internal/ingest/handler_test.go @@ -1914,15 +1914,20 @@ func (s *oneShotMissingObjects) GetByAPID(ctx context.Context, apID string) (*st func (h *harness) swapHandlerStores(objects store.APObjects, actors store.BridgedActors) { h.t.Helper() handler, err := NewHandler(HandlerOptions{ - Materializer: h.mat, - Fetcher: h.client, - Objects: objects, - Actors: actors, - Communities: h.communities, - Tombstones: h.tombstones, - Records: h.manager, - Votes: h.votes, - Backfill: h.backfills, + Materializer: h.mat, + Fetcher: h.client, + Objects: objects, + Actors: actors, + Communities: h.communities, + Tombstones: h.tombstones, + Records: h.manager, + Votes: h.votes, + Backfill: h.backfills, + // The SAME classifier, deliberately over the harness's real stores: this + // helper swaps the dispatcher's store VIEWS, and pointing the guard at a + // substitute view would change what is being tested here into a test of + // the guard. + Echo: h.classifier, ServiceActorID: h.service.ID, }) require.NoError(h.t, err) diff --git a/internal/ingest/ingest_test.go b/internal/ingest/ingest_test.go index 6ac2c84..e4e4628 100644 --- a/internal/ingest/ingest_test.go +++ b/internal/ingest/ingest_test.go @@ -182,6 +182,11 @@ type harness struct { // same database and unseals the personas' AP keys). db *sql.DB custodian *identity.Custodian + // classifier is the harness's echo classifier over the real stores. It is + // kept so a rebuilt dispatcher keeps the SAME guard: echo suppression is + // mandatory, and a swap that quietly dropped it would disable it for the + // swapped test only. + classifier EchoClassifier // logs captures the dispatcher's own log output, so a test can assert that // a drop was INTENTIONAL — an incidental drop and a deliberate one are // indistinguishable from state alone. @@ -327,6 +332,7 @@ func newHarness(t *testing.T) *harness { Actors: store.NewAPActors(database), }) require.NoError(t, err) + h.classifier = classifier h.handler, err = NewHandler(HandlerOptions{ Materializer: h.mat, Fetcher: h.client, diff --git a/internal/votes/e2e_test.go b/internal/votes/e2e_test.go index 6e32bb4..30d27f6 100644 --- a/internal/votes/e2e_test.go +++ b/internal/votes/e2e_test.go @@ -2,6 +2,7 @@ package votes import ( "context" + "database/sql" "encoding/json" "net/http" "net/url" @@ -11,12 +12,30 @@ import ( "github.com/stretchr/testify/require" "tidepool/internal/ap" + "tidepool/internal/echo" "tidepool/internal/errors" "tidepool/internal/ingest" "tidepool/internal/materialize" "tidepool/internal/store" ) +// e2eClassifier is the REAL echo classifier over this test's stores. These +// fixtures vote as fediverse users on fediverse content, so it answers +// ClassNone throughout — and would stop doing so the moment the classifier +// started matching genuine remote traffic, which is the failure that takes the +// vote pipeline dark. +func e2eClassifier(t *testing.T, database *sql.DB, objects store.APObjects) *echo.Classifier { + t.Helper() + classifier, err := echo.New(echo.Options{ + Objects: objects, + OutboundObjects: store.NewOutboundObjects(database), + Activities: store.NewOutboundActivities(database), + Actors: store.NewAPActors(database), + }) + require.NoError(t, err) + return classifier +} + // The aggregator is the real implementation behind task 06's seam. var _ ingest.VoteAggregator = (*Aggregator)(nil) @@ -161,6 +180,7 @@ func TestFakeLemmyVoteE2E(t *testing.T) { Tombstones: store.NewTombstones(database), Records: &fakeRecords{records: map[string]map[string]any{}}, Votes: agg, + Echo: e2eClassifier(t, database, objects), ServiceActorID: e2eServiceID, }) require.NoError(t, err) @@ -212,6 +232,7 @@ func TestBareVoteDispatch(t *testing.T) { Tombstones: store.NewTombstones(database), Records: &fakeRecords{records: map[string]map[string]any{}}, Votes: agg, + Echo: e2eClassifier(t, database, objects), ServiceActorID: e2eServiceID, }) require.NoError(t, err) -- 2.51.2