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") +}