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)