From fa3807495a3470328a9ac4b19db105789d8d0d86 Mon Sep 17 00:00:00 2001 From: Bretton Date: Sat, 8 Aug 2026 02:15:30 -0700 Subject: [PATCH] test(ingestion): pin the handle-flood root causes and the fetch taxonomy MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Four reds; T0 fully green and no pre-existing T1 breakage. 1. communities: the profile record CreateCommunity writes carries no `handle`. Read back via getRecord — the failure dumps the actual record, so the evidence is in the output rather than in an argument. Without the field the consumer must resolve one from the DID document, which on an egress-blocked stack yields "handle.invalid" into a UNIQUE column. 2. jetstream: a profile event whose handle is held by a DIFFERENT DID is currently swallowed as an idempotent replay and returns nil. Pinned as a PERMANENT refusal naming the contested handle, with the community absent and the incumbent untouched. Permanent is the load-bearing half: left transient it costs ~4.2s of inline blocking plus ten redrives per delivery, which is the flood. The other half of the narrowing is pinned beside it and PASSES today — a same-DID replay must stay a silent no-op. A fix that widened the refusal to every conflict would dead-letter every community's profile on every cursor rewind, and nothing else in the suite would notice. 3. jetstream: a genuine XRPC RecordNotFound from the author's PDS is classified transient today; pinned as permanent, with the httptest request count asserting one event produces exactly one fetch. Bare 404 and 5xx are pinned as transient and PASS today — they are the guards that stop the fix over-reaching, since a bare 404 usually means the request never reached a PDS at all (users.FetchProfileRecord draws the same distinction). Scope note: driving HandleEvent directly, no dead-letter row is written and no redrive runs — the consumer never touches that table, the connector does. These pin the input the connector switches on (the ErrPermanentEvent wrapping) plus the absence of any retry loop inside the consumer/fetcher. The connector's half is TestConnector_DeadLettersAfterRetryExhaustion. 4. jetstream: a community whose handle resolves to "handle.invalid" must not be stored, classified TRANSIENT (the PLC may be unreachable now and fine in a minute). Mirrors authorpost.go's user-path guard. No conflict with fix 2 — different causes, composing guards: this one refuses before the insert, that one refuses a collision the insert reports. Co-Authored-By: Claude Fable 5 --- .../jetstream/acceptance_consumer_test.go | 141 +++++++++++++ .../community_handle_conflict_test.go | 198 ++++++++++++++++++ .../community_profile_handle_test.go | 90 ++++++++ 3 files changed, 429 insertions(+) create mode 100644 internal/atproto/jetstream/community_handle_conflict_test.go create mode 100644 internal/core/communities/community_profile_handle_test.go diff --git a/internal/atproto/jetstream/acceptance_consumer_test.go b/internal/atproto/jetstream/acceptance_consumer_test.go index 534f7d5..416c2e3 100644 --- a/internal/atproto/jetstream/acceptance_consumer_test.go +++ b/internal/atproto/jetstream/acceptance_consumer_test.go @@ -498,6 +498,147 @@ func TestDirectPostFetcher_RefusesAPrivateHostByDefault(t *testing.T) { assert.Falsef(t, reached, "the guard must refuse before the request is made, not after the server has already answered") } +// --------------------------------------------------------------------------- +// §5.4 direct fetch: classifying what the author's PDS answered +// --------------------------------------------------------------------------- + +// How a failed fetch is CLASSIFIED, which decides what the connector does next. +// +// The consumer returns an error and the connector reads its shape: an error +// wrapping ErrPermanentEvent is dead-lettered with the redrive budget already +// spent, and anything else is retried inline (~4.2 seconds of blocking, per +// event) and then redriven ten times. So the classification is not a label — it +// is the difference between one forensic row and forty pointless refetches +// against somebody else's PDS while this consumer stops indexing. +// +// A non-200 from getRecord is currently one undifferentiated failure, and the +// three cases below need three different answers: +// +// - A GENUINE XRPC RecordNotFound is a definite fact about the repo: the PDS +// was reached, it understood the question, and the record is not there. No +// retry changes that. +// - A BARE 404 — no XRPC error envelope — usually means the request never +// reached a PDS at all: a stale pds_url pointing at a reverse proxy or a +// generic web server, which answers 404 for everything. Treating that as +// proof the record does not exist would permanently discard a post over a +// misconfigured hostname. users.FetchProfileRecord already draws exactly +// this distinction, and for exactly this reason. +// - A 5xx is the PDS saying it is having a bad time. Definitionally transient. +// +// WHAT IS ASSERTED HERE, AND WHAT IS NOT. These drive HandleEvent directly, so +// no dead-letter row is written and no redrive runs — the consumer never touches +// that table; the connector does. What these pin is the input the connector +// switches on (the ErrPermanentEvent wrapping) plus the request count, which +// proves the consumer and fetcher hold no retry loop of their own. The +// connector's half — that a permanent error is dead-lettered exhausted and a +// transient one is retried — is TestConnector_DeadLettersAfterRetryExhaustion. + +// countingAuthorPDS is fakeAuthorPDS with a request tally, so a test can prove +// one event produced exactly one outbound fetch. +func countingAuthorPDS(t *testing.T, expectRepo string, requests *int, handler http.HandlerFunc) *httptest.Server { + t.Helper() + return fakeAuthorPDS(t, expectRepo, func(w http.ResponseWriter, r *http.Request) { + *requests++ + handler(w, r) + }) +} + +func TestAcceptanceConsumer_GenuineRecordNotFound_IsPermanentlyRefusedAfterOneFetch(t *testing.T) { + t.Parallel() + + db := testkit.DB(t) + base := time.Now().UnixMicro() + uri := accPostURI("accgone") + + var requests int + srv := countingAuthorPDS(t, accAuthor, &requests, func(w http.ResponseWriter, r *http.Request) { + // The reference PDS' answer for a repo that exists and a record that + // does not: 400 with an XRPC error envelope naming RecordNotFound. + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(http.StatusBadRequest) + _, _ = w.Write([]byte(`{"error":"RecordNotFound","message":"Could not locate record: ` + uri + `"}`)) + }) + defer srv.Close() + + f := newAccFixture(t, db, WithPostRecordFetcher(newFetcherAt(t, accAuthor, srv.URL))) + + err := f.consumer.HandleEvent(context.Background(), + acceptanceEvent(accCommunity, uri, "bafyreiaccgonepinned", testkit.TID(), base)) + + require.Error(t, err, "an acceptance whose subject the PDS says does not exist cannot be applied") + assert.ErrorIs(t, err, ErrPermanentEvent, + "a genuine RecordNotFound is a definite fact about the repo — the PDS was reached and answered — so no retry can change it. "+ + "Left transient, every one of these costs ~4.2s of inline blocking plus ten redrives, and a community can mint them at will "+ + "by writing acceptances for URIs nobody wrote") + + assert.Equalf(t, 1, requests, + "one event produced %d fetches: the consumer or the fetcher is retrying internally. Retries belong to the connector, "+ + "which can classify and budget them; a loop in here is invisible to it and unbounded", requests) + + assert.Zero(t, countRows(t, db, `SELECT count(*) FROM posts WHERE uri = $1`, uri), + "nothing may be indexed for a subject the PDS says does not exist") +} + +func TestAcceptanceConsumer_BareNotFound_StaysTransient(t *testing.T) { + t.Parallel() + + db := testkit.DB(t) + base := time.Now().UnixMicro() + uri := accPostURI("accbare404") + + var requests int + srv := countingAuthorPDS(t, accAuthor, &requests, func(w http.ResponseWriter, r *http.Request) { + // No XRPC envelope. This is what a reverse proxy, a load balancer or a + // generic web server answers — which is what a stale pds_url in a DID + // document points at. + w.Header().Set("Content-Type", "text/html") + w.WriteHeader(http.StatusNotFound) + _, _ = w.Write([]byte("404 Not Found")) + }) + defer srv.Close() + + f := newAccFixture(t, db, WithPostRecordFetcher(newFetcherAt(t, accAuthor, srv.URL))) + + err := f.consumer.HandleEvent(context.Background(), + acceptanceEvent(accCommunity, uri, "bafyreiaccbarepinned", testkit.TID(), base)) + + require.Error(t, err, "a fetch that did not produce a record cannot apply the acceptance") + assert.NotErrorIs(t, err, ErrPermanentEvent, + "a bare 404 carries no XRPC error envelope, which means the request most likely never reached a PDS at all — a stale pds_url "+ + "pointing at a proxy. Reading it as proof the record does not exist permanently discards a real post over a misconfigured "+ + "hostname; users.FetchProfileRecord draws the same distinction for the same reason") + + assert.Equal(t, 1, requests, "one event, one fetch — the connector owns retries") +} + +func TestAcceptanceConsumer_PDSServerError_StaysTransient(t *testing.T) { + t.Parallel() + + db := testkit.DB(t) + base := time.Now().UnixMicro() + uri := accPostURI("accpds5xx") + + var requests int + srv := countingAuthorPDS(t, accAuthor, &requests, func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(http.StatusServiceUnavailable) + _, _ = w.Write([]byte(`{"error":"InternalServerError","message":"upstream unavailable"}`)) + }) + defer srv.Close() + + f := newAccFixture(t, db, WithPostRecordFetcher(newFetcherAt(t, accAuthor, srv.URL))) + + err := f.consumer.HandleEvent(context.Background(), + acceptanceEvent(accCommunity, uri, "bafyreiacc5xxpinned", testkit.TID(), base)) + + require.Error(t, err) + assert.NotErrorIs(t, err, ErrPermanentEvent, + "a 5xx is the author's PDS saying it is unwell, which is the definition of transient; discarding the acceptance permanently "+ + "would lose a post because somebody else's server restarted") + + assert.Equal(t, 1, requests, "one event, one fetch — the connector owns retries") +} + // readAdmissionRow returns every mutable column of one admission row as a // comparable value, so "the row did not change" can be asserted as a whole // rather than field by field — a new column added later is covered without diff --git a/internal/atproto/jetstream/community_handle_conflict_test.go b/internal/atproto/jetstream/community_handle_conflict_test.go new file mode 100644 index 0000000..9c0c903 --- /dev/null +++ b/internal/atproto/jetstream/community_handle_conflict_test.go @@ -0,0 +1,198 @@ +//go:build integration + +package jetstream + +import ( + "context" + "fmt" + "testing" + "time" + + "Coves/internal/atproto/identity" + "Coves/internal/db/postgres" + "Coves/tests/fixtures" + "Coves/tests/testkit" + + _ "github.com/lib/pq" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// What a handle collision does to a community profile event. +// +// # THE SWALLOW, AND WHY IT LOOKS REASONABLE +// +// createCommunity ends in repo.Create, and treats communities.IsConflict as +// "already indexed — idempotent, nothing to do" (community_consumer.go). That +// reading is correct for exactly one of the two conflicts IsConflict matches. +// +// - ErrCommunityAlreadyExists means the DID is already in the table. That IS +// an idempotent replay: the community this event describes is indexed, the +// event changed nothing, and returning nil is right. Jetstream redelivers +// constantly, so this path is walked all the time. +// - ErrHandleTaken means a DIFFERENT DID already holds this handle. Nothing +// about that is idempotent. The community in the event was NOT indexed, is +// not in the table under any DID, and never will be — and the AppView says +// nothing, logs it as a successful replay, and moves on. +// +// The consequence is not confined to the community. Posts, comments and votes +// naming it are refused as "community not found", which the taxonomy classifies +// TRANSIENT (correctly — it is ordinarily a delivery race), so each one burns +// the connector's full inline retry budget before dead-lettering. A single +// swallowed collision turns into a sustained flood of 4.2-second blocking +// failures in three other consumers, none of which points anywhere near here. +// +// So the pin is in two halves, and both are needed: the collision must surface, +// and the genuine replay must keep NOT surfacing. A fix that widened the error +// into every conflict would dead-letter every redelivered profile event in the +// system. + +const conflictProfileCollection = "social.coves.community.profile" + +// communityProfileEvent builds the commit a community's profile write produces. +// The handle is carried in the record, which is the shape the AppView's own +// CreateCommunity produces once it stops omitting the field +// (internal/core/communities/community_profile_handle_test.go). +func communityProfileEvent(did, handle, name, rev string) *JetstreamEvent { + return &JetstreamEvent{ + Did: did, + Kind: "commit", + TimeUS: time.Now().UnixMicro(), + Commit: &CommitEvent{ + Rev: rev, + Operation: "create", + Collection: conflictProfileCollection, + RKey: "self", + CID: "bafyconflict" + rev, + Record: map[string]interface{}{ + "$type": conflictProfileCollection, + "handle": handle, + "name": name, + "displayName": "Conflict " + name, + "createdBy": "did:plc:conflicttestcreator", + "hostedBy": "did:web:test.local", + "visibility": "public", + "federation": map[string]interface{}{"allowExternalDiscovery": true}, + "createdAt": time.Now().UTC().Format(time.RFC3339), + }, + }, + } +} + +func TestCommunityConsumer_HandleTakenByAnotherDID_IsAPermanentRefusal(t *testing.T) { + t.Parallel() + + db := testkit.DB(t) + ctx := context.Background() + // skipVerification: the hostedBy/handle-domain check is a different security + // property with its own tests, and leaving it on would reject these events + // before the conflict path is ever reached. + consumer := NewCommunityEventConsumer(postgres.NewCommunityRepository(db), "did:web:test.local", true, nil) + + suffix := testkit.UniqueID(t) + contested := fmt.Sprintf("c-first%s.test.local", suffix) + incumbent := fixtures.DID("incumbent" + suffix) + newcomer := fixtures.DID("newcomer" + suffix) + + require.NoError(t, consumer.HandleEvent(ctx, communityProfileEvent(incumbent, contested, "first"+suffix, "3lconflicta")), + "fixture: the first community must index cleanly") + + // A SECOND, DIFFERENT community claiming the same handle. In production this + // is what two communities resolving to "handle.invalid" look like, but it is + // equally what a genuine handle race or a hostile duplicate looks like — and + // none of them is an idempotent replay. + err := consumer.HandleEvent(ctx, communityProfileEvent(newcomer, contested, "second"+suffix, "3lconflictb")) + + require.Errorf(t, err, + "a profile event whose handle is already held by a DIFFERENT community (%s) was accepted as an idempotent replay. "+ + "The community was never indexed, and every post naming it will dead-letter as \"community not found\" with "+ + "nothing in the logs pointing here", incumbent) + assert.ErrorIsf(t, err, ErrPermanentEvent, + "the refusal must be PERMANENT. A handle held by another DID does not resolve itself by waiting, so a transient "+ + "classification spends the connector's full inline retry budget (~4.2s, blocking the consumer) and then ten "+ + "redrives, per delivery, forever") + assert.Containsf(t, err.Error(), contested, + "the error must name the contested handle: it is the only thing that makes this diagnosable from a log line") + + // The newcomer is genuinely absent, and the incumbent is untouched — a + // refusal that had partially applied would be worse than the swallow. + var newcomerRows, incumbentRows int + require.NoError(t, db.QueryRow(`SELECT count(*) FROM communities WHERE did = $1`, newcomer).Scan(&newcomerRows)) + assert.Zero(t, newcomerRows, "the refused community must not be indexed") + + require.NoError(t, db.QueryRow( + `SELECT count(*) FROM communities WHERE did = $1 AND handle = $2`, incumbent, contested).Scan(&incumbentRows)) + assert.Equal(t, 1, incumbentRows, "the community that legitimately holds the handle must be untouched by the refusal") +} + +func TestCommunityConsumer_SameDIDReplay_StaysSilent(t *testing.T) { + t.Parallel() + + db := testkit.DB(t) + ctx := context.Background() + consumer := NewCommunityEventConsumer(postgres.NewCommunityRepository(db), "did:web:test.local", true, nil) + + suffix := testkit.UniqueID(t) + handle := fmt.Sprintf("c-replay%s.test.local", suffix) + did := fixtures.DID("replay" + suffix) + + event := communityProfileEvent(did, handle, "replay"+suffix, "3lreplaya") + require.NoError(t, consumer.HandleEvent(ctx, event)) + + // THE OTHER HALF OF THE NARROWING, and the reason it has to be pinned + // alongside the collision rather than left implied. The connector rewinds + // its cursor five seconds after every reconnect and the AppView consumes + // overlapping feeds, so this exact commit is guaranteed to be redelivered — + // constantly, for every community. A fix that widened the refusal to every + // conflict would dead-letter all of it. + require.NoError(t, consumer.HandleEvent(ctx, event), + "a redelivered profile event for the SAME DID is a genuine idempotent replay and must stay a silent no-op; "+ + "refusing it would dead-letter every community's profile on every cursor rewind") + + var rows int + require.NoError(t, db.QueryRow(`SELECT count(*) FROM communities WHERE did = $1`, did).Scan(&rows)) + assert.Equal(t, 1, rows, "a replay must not produce a second row") +} + +func TestCommunityConsumer_UnverifiableHandle_IsNotStored(t *testing.T) { + t.Parallel() + + db := testkit.DB(t) + ctx := context.Background() + + // A record with NO handle — the federated shape, where resolution is the + // only option — and a resolver that cannot verify the DID's handle. atProto + // identity resolution reports that as the reserved "handle.invalid" rather + // than as an error, so the consumer receives a perfectly well-formed + // identity naming a handle that is not one. + did := fixtures.DID("unverified" + testkit.UniqueID(t)) + resolver := &mockIdentityResolverForUser{identities: map[string]*identity.Identity{ + did: {DID: did, Handle: invalidHandle, PDSURL: "https://pds.example.invalid"}, + }} + consumer := NewCommunityEventConsumer(postgres.NewCommunityRepository(db), "did:web:test.local", true, resolver) + + event := communityProfileEvent(did, "", "unverified", "3lunverified") + delete(event.Commit.Record, "handle") + + err := consumer.HandleEvent(ctx, event) + + // This is the guard authorpost.go already applies on the user path, for the + // identical reason: "handle.invalid" is a PLACEHOLDER, not a handle, and the + // column it would land in is UNIQUE. Store it once and the next unverifiable + // community collides with it — which is the collision the test above pins, + // arriving from a completely different direction and with no attacker + // involved. + require.Errorf(t, err, + "a community whose handle could not be verified was indexed anyway. \"handle.invalid\" is the reserved "+ + "placeholder for exactly this case, and communities.handle is UNIQUE — so the FIRST one indexed takes the "+ + "placeholder and every later one collides with it") + assert.NotErrorIsf(t, err, ErrPermanentEvent, + "an unverifiable handle is a RESOLUTION failure, not a property of the record: the PLC directory may be "+ + "unreachable right now and verifiable in a minute, so the redrive has to be allowed to succeed") + + var stored int + require.NoError(t, db.QueryRow( + `SELECT count(*) FROM communities WHERE handle = $1`, invalidHandle).Scan(&stored)) + assert.Zerof(t, stored, + "%q was written into communities.handle; the next community that cannot resolve will collide with it", invalidHandle) +} diff --git a/internal/core/communities/community_profile_handle_test.go b/internal/core/communities/community_profile_handle_test.go new file mode 100644 index 0000000..5466529 --- /dev/null +++ b/internal/core/communities/community_profile_handle_test.go @@ -0,0 +1,90 @@ +//go:build integration + +package communities_test + +import ( + "context" + "strings" + "testing" + + "Coves/internal/core/communities" + "Coves/tests/testkit" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// The handle in the record CreateCommunity writes to the PDS. +// +// # WHY A MISSING MAP KEY IS A PIPELINE OUTAGE +// +// The community profile record is the ONLY thing that tells the AppView a +// community exists — the community consumer indexes repos it has never seen, and +// there is no signup step for communities. When the record carries no `handle`, +// the consumer falls through to resolving one from the DID document +// (community_consumer.go createCommunity), and on an egress-blocked stack that +// resolution cannot reach the PLC directory. What comes back is the reserved +// "handle.invalid". +// +// communities.handle carries a UNIQUE constraint. So the FIRST community indexed +// that way takes "handle.invalid", every subsequent one collides with it, and +// the collision is currently swallowed as an idempotent replay — the community +// is silently dropped, and every post, comment and vote naming it dead-letters +// as "community not found". One absent map key, and the visible symptom is a +// flood of unrelated transient failures four layers away. +// +// Writing the handle is not a departure from atProto's "handles are mutable, +// resolve them from DIDs" guidance. That guidance is about trusting a stranger's +// self-reported handle; this AppView PROVISIONED the account and asked the PDS +// for this exact handle, so putting it in the record states a fact it already +// holds. The consumer's resolution path stays for federated communities, where +// it is the only option available. +func TestCreateCommunity_ProfileRecordCarriesTheHandle(t *testing.T) { + t.Parallel() + + service, _, pdsServer := newCommunityService(t) + ctx := context.Background() + + // "c-" plus the name must stay inside the PDS' 18-character local-label cap, + // which is why UniqueIDWithPrefix is the only generator allowed here. + name := testkit.UniqueIDWithPrefix(t, "hc") + require.LessOrEqualf(t, len("c-"+name), testkit.MaxIDLength, + "the generated community name %q makes a handle label the PDS will refuse", name) + + community, err := service.CreateCommunity(ctx, communities.CreateCommunityRequest{ + Name: name, + DisplayName: "Handle Carrier", + Description: "a community whose profile record must name its handle", + Visibility: "public", + CreatedByDID: "did:plc:handlecarriertest", + }) + require.NoError(t, err) + require.NotEmpty(t, community.Handle, "fixture: the service must have derived a handle") + + // Read the record back out of the community's own repo — the bytes the + // firehose will carry, not the service's in-memory view of them. The + // consumer parses exactly this. + session := pdsServer.Login(t, community.Handle, community.PDSPassword) + record := session.GetRecord(t, "social.coves.community.profile", "self") + + handle, ok := record.Value["handle"].(string) + require.Truef(t, ok, + "the profile record carries no `handle` field: %#v.\n"+ + "Without it the consumer must resolve one from the DID document, which on an egress-blocked stack yields "+ + "the reserved \"handle.invalid\" — and communities.handle is UNIQUE, so the second community indexed that "+ + "way collides with the first and is dropped, taking every post that names it into the dead-letter queue "+ + "as \"community not found\"", record.Value) + + assert.Equal(t, community.Handle, handle, + "the record's handle must be the one the service provisioned and reports; a record that disagrees with the "+ + "AppView's own row is worse than a missing field, because the index and the firehose would then be "+ + "describing two different communities") + + // The shape is asserted independently of the service's own string building, + // so a change to the derivation has to be a deliberate edit here too. + assert.Truef(t, strings.HasPrefix(handle, "c-"+name+"."), + "the handle %q must be the c-{name}.{domain} form the provisioner builds; a handle in any other shape does not "+ + "resolve to this community's DID, and every client addressing it by handle 404s while the row looks healthy", handle) + assert.NotContains(t, handle, "handle.invalid", + "the reserved unverifiable-identity handle must never reach a record: it is the value that collides on the UNIQUE constraint") +} -- 2.51.2