From 66bc5255e5b239908d84723172ae2063335c8466 Mon Sep 17 00:00:00 2001 From: Bretton Date: Sat, 8 Aug 2026 02:22:22 -0700 Subject: [PATCH] fix(ingestion): close the handle-collision flood and sort the fetch taxonomy MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Four fixes for RED's pins at fa38074. Three of them are one causal chain that ended as a dead-letter flood on a lane four layers from the cause. 1. communities: CreateCommunity's profile record now carries `handle`. That record is the ONLY thing that tells an AppView a community exists, so omitting the field left the consumer resolving one from the DID document — which on an egress-blocked stack yields "handle.invalid". Writing it is not a departure from "handles are mutable, resolve from DIDs": that guidance is about trusting a STRANGER's self-reported handle, and this process provisioned the account and asked the PDS for exactly this one. Federated communities still resolve, because there resolution is the only option. 2. jetstream: the community consumer no longer stores an unverifiable handle. Identity resolution reports an unverifiable DID by returning the reserved "handle.invalid" rather than an error, so a well-formed identity naming a non-handle was landing in a UNIQUE column — and the first one to do so made every later unverifiable community collide with it. TRANSIENT, because it is a fact about the resolution and not about the record: the directory may answer fine a minute later. Mirrors the user path's guard in authorpost.go. 3. jetstream: the conflict swallow is narrowed. communities.IsConflict matches two errors that mean opposite things — ErrCommunityAlreadyExists IS an idempotent replay (walked constantly, must stay a silent no-op), while ErrHandleTaken means a DIFFERENT DID holds the handle and the community in the event was never indexed at all. The second is now a PERMANENT refusal naming both DIDs and the handle; an unclassified conflict is reported rather than swallowed, so a future unique constraint cannot inherit the same silence. 4. jetstream: DirectPostFetcher sorts a non-200 getRecord instead of reporting one undifferentiated failure. A genuine XRPC RecordNotFound is permanent — the PDS was reached and answered, and left transient any community could mint unlimited lane-blocking by writing acceptances for URIs nobody wrote. A BARE 404 stays transient: with no envelope the request most likely never reached a PDS, so trusting it would discard a real post over a mistyped hostname. 5xx stays transient by definition. The predicate is EXPORTED from users (IsRecordNotFoundResponse) and FetchProfileRecord now calls it, so the line is drawn once. Why permanence is the load-bearing half of 3 and 4: a transient error costs the connector three inline retries — about 4.2 seconds of a blocked lane that now carries four collections — plus ten redrives, per delivery. ErrPermanentEvent short-circuits all of it. Co-Authored-By: Claude Fable 5 --- internal/atproto/jetstream/authorpost.go | 27 +++++++ .../atproto/jetstream/community_consumer.go | 77 ++++++++++++++++++- internal/core/communities/service.go | 21 ++++- internal/core/users/profile_backfill.go | 44 ++++++++--- 4 files changed, 156 insertions(+), 13 deletions(-) diff --git a/internal/atproto/jetstream/authorpost.go b/internal/atproto/jetstream/authorpost.go index 25487e2..84aaef0 100644 --- a/internal/atproto/jetstream/authorpost.go +++ b/internal/atproto/jetstream/authorpost.go @@ -251,6 +251,33 @@ func (f *DirectPostFetcher) FetchPost(ctx context.Context, postURI string) (*Fet if len(detail) > maxFetchErrorDetailBytes { detail = detail[:maxFetchErrorDetailBytes] } + // THE CLASSIFICATION IS THE EXPENSIVE PART OF THIS FUNCTION. The + // connector reads the returned error's shape: ErrPermanentEvent is + // dead-lettered with its redrive budget already spent, and anything + // else costs three inline retries (~4.2s of a blocked lane that also + // carries posts) plus ten redrives. So a non-200 has to be sorted, not + // merely reported. + // + // A GENUINE XRPC RecordNotFound is a definite fact about the repo: the + // PDS was reached, understood the question, and answered that the + // record is not there. Nothing a retry does changes it — and left + // transient, any community can mint unlimited lane-blocking by writing + // acceptances for URIs nobody ever wrote. + // + // A BARE 404 is the opposite, and the distinction is not pedantry: + // with no XRPC envelope the request most likely never reached a PDS at + // all — a stale pds_url pointing at a reverse proxy or a generic web + // server, both of which 404 everything. Reading that as proof the + // record does not exist would permanently discard a real post over a + // misconfigured hostname. Everything else, 5xx included, is the PDS + // having a bad time and is transient by definition. + // + // users.FetchProfileRecord draws exactly this line for exactly this + // reason; the predicates are shared with it rather than re-derived. + if users.IsRecordNotFoundResponse(resp.StatusCode, body) { + return nil, fmt.Errorf("%w: the PDS serving %s answered getRecord with status %d: %s", + ErrPermanentEvent, postURI, resp.StatusCode, strconv.Quote(detail)) + } // Quoted so control characters and ANSI escapes from a hostile PDS // cannot corrupt log output. return nil, fmt.Errorf("the PDS serving %s answered getRecord with status %d: %s", diff --git a/internal/atproto/jetstream/community_consumer.go b/internal/atproto/jetstream/community_consumer.go index 5a9a517..0e22938 100644 --- a/internal/atproto/jetstream/community_consumer.go +++ b/internal/atproto/jetstream/community_consumer.go @@ -7,6 +7,7 @@ import ( "Coves/internal/core/richtext" "context" "encoding/json" + "errors" "fmt" "log" "net/http" @@ -213,6 +214,25 @@ func (c *CommunityEventConsumer) createCommunity(ctx context.Context, did string if err != nil { return fmt.Errorf("failed to resolve handle from PLC for %s: %w (no fallback - will retry during backfill)", did, err) } + // "handle.invalid" IS NOT A HANDLE. atProto identity resolution + // reports a DID whose handle it could not verify bidirectionally + // by returning that reserved placeholder rather than an error, so + // what arrives here is a perfectly well-formed identity naming a + // non-handle — and communities.handle is UNIQUE. Store it once and + // every subsequent unverifiable community collides with it, which + // the insert reports as a conflict and the swallow below used to + // discard silently. + // + // TRANSIENT, deliberately: this is a fact about the RESOLUTION, not + // about the record. The PLC directory may be unreachable this + // second and answer fine the next, so the redrive has to be allowed + // to succeed. The user path applies the same guard for the same + // reason (authorpost.go, hydrateAuthorOpportunistically). + if identity.Handle == "" || identity.Handle == invalidHandle { + return fmt.Errorf("resolving the handle of community %s: identity resolution returned %q, "+ + "which is the reserved placeholder for an unverifiable handle and must never be stored in a unique column "+ + "(retryable — the directory may verify it later)", did, identity.Handle) + } profile.Handle = identity.Handle // Persist the resolved PDS host: BridgeTrust gates bridgedStats // on the post's community row carrying its repo's PDS URL, and a @@ -300,13 +320,42 @@ func (c *CommunityEventConsumer) createCommunity(ctx context.Context, did string } } - // Index in AppView database + // Index in AppView database. + // + // THE TWO CONFLICTS MEAN OPPOSITE THINGS, and treating them alike is what + // turned one handle collision into a flood of unrelated dead letters. + // communities.IsConflict matches both, so it is too wide to switch on here. _, err = c.repo.Create(ctx, community) if err != nil { - // Check if it already exists (idempotency) - if communities.IsConflict(err) { + switch { + case errors.Is(err, communities.ErrCommunityAlreadyExists): + // The DID is already in the table: a genuine idempotent replay. + // The connector rewinds its cursor after every reconnect and the + // AppView consumes overlapping feeds, so this path is walked + // constantly for every community and must stay a silent no-op. log.Printf("Community already indexed: %s (%s)", community.Handle, community.DID) return nil + + case errors.Is(err, communities.ErrHandleTaken): + // A DIFFERENT DID already holds this handle. Nothing about that is + // idempotent: the community in this event was NOT indexed, is in + // the table under no DID at all, and never will be while the + // incumbent stands. + // + // PERMANENT. A handle held by someone else does not resolve itself + // by waiting, so a transient classification spends the connector's + // full inline retry budget — about 4.2 seconds of blocking, on a + // lane that also carries posts — and then ten redrives, per + // delivery, forever. Both DIDs and the handle go in the message + // because a log line is the only place this is diagnosable from. + return fmt.Errorf("%w: cannot index community %s: handle %q is already held by community %s", + ErrPermanentEvent, community.DID, community.Handle, c.incumbentOfHandle(ctx, community.Handle)) + + case communities.IsConflict(err): + // A conflict this build does not recognise. Reported rather than + // swallowed: the swallow is what hid the handle collision, and a + // new unique constraint would otherwise inherit the same silence. + return fmt.Errorf("failed to index community %s: unclassified conflict: %w", community.DID, err) } return fmt.Errorf("failed to index community: %w", err) } @@ -315,6 +364,28 @@ func (c *CommunityEventConsumer) createCommunity(ctx context.Context, did string return nil } +// incumbentOfHandle names the community that already holds a contested handle, +// for the refusal message. +// +// BEST EFFORT, and never allowed to change the outcome. The refusal is already +// decided by the time this runs; this only fills in the half of the message an +// operator cannot otherwise get — "some other community has it" sends them to +// the database, "did:plc:… has it" sends them to the community. A lookup that +// fails yields a placeholder rather than an error, because replacing a precise +// permanent refusal with a vague transient one would trade the diagnosis for +// the flood this whole change exists to stop. +func (c *CommunityEventConsumer) incumbentOfHandle(ctx context.Context, handle string) string { + const unknown = "an unidentified DID" + if c.repo == nil { + return unknown + } + incumbent, err := c.repo.GetByHandle(ctx, handle) + if err != nil || incumbent == nil { + return unknown + } + return incumbent.DID +} + // updateCommunity updates an existing community from the firehose func (c *CommunityEventConsumer) updateCommunity(ctx context.Context, did string, commit *CommitEvent) error { if commit.Record == nil { diff --git a/internal/core/communities/service.go b/internal/core/communities/service.go index 8c5d629..be227f5 100644 --- a/internal/core/communities/service.go +++ b/internal/core/communities/service.go @@ -195,10 +195,29 @@ func (s *communityService) CreateCommunity(ctx context.Context, req CreateCommun return nil, fmt.Errorf("generated atProto handle is invalid: %w", validateErr) } - // Build community profile record + // Build community profile record. + // + // THE HANDLE IS IN THE RECORD, and its absence used to be a pipeline + // outage. This record is the ONLY thing that tells an AppView a community + // exists — the community consumer indexes repos it has never seen, and + // communities have no signup step — so a record without a handle leaves the + // consumer to resolve one from the DID document. On an egress-blocked stack + // that resolution cannot reach the PLC directory and yields the reserved + // "handle.invalid". communities.handle is UNIQUE, so the first community + // indexed that way takes the placeholder and every later one collides with + // it; the community is dropped, and every post, comment and vote naming it + // dead-letters as "community not found", four layers from the cause. + // + // This 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 exactly this handle, so writing it states a fact this process + // already holds. Consumers still resolve for federated communities, where + // resolution is the only option there is. profile := map[string]interface{}{ "$type": "social.coves.community.profile", "name": req.Name, // Short name for !mentions (e.g., "gaming") + "handle": pdsAccount.Handle, "visibility": req.Visibility, "hostedBy": s.instanceDID, // V2: Instance hosts, community owns "createdBy": req.CreatedByDID, diff --git a/internal/core/users/profile_backfill.go b/internal/core/users/profile_backfill.go index 481fee9..e61bf62 100644 --- a/internal/core/users/profile_backfill.go +++ b/internal/core/users/profile_backfill.go @@ -80,15 +80,9 @@ func FetchProfileRecord(ctx context.Context, client *http.Client, pdsURL, did st switch { case resp.StatusCode == http.StatusOK: // fall through to parse below - case resp.StatusCode == http.StatusNotFound && isXRPCErrorBody(body): - // Some PDS implementations 404 on missing records. Only trust the 404 - // when the body is an XRPC error object — a bare 404 (HTML from a - // reverse proxy or generic web server behind a stale pds_url) means we - // never reached a PDS at all, so it must surface as an error, not be - // silently classified as "user has no profile record". - return nil, nil - case resp.StatusCode == http.StatusBadRequest && isRecordNotFoundBody(body): - // Reference PDS returns 400 RecordNotFound when the repo exists but has no profile + case IsRecordNotFoundResponse(resp.StatusCode, body): + // The PDS was reached and said the record is not there. Absence is a + // normal outcome here, not an error. return nil, nil default: // SECURITY: cap the echoed body so a hostile PDS can't flood our logs, @@ -117,6 +111,38 @@ func FetchProfileRecord(ctx context.Context, client *http.Client, pdsURL, did st return &input, nil } +// IsRecordNotFoundResponse reports whether a getRecord response is a PDS +// saying, definitively, that the record is not there. +// +// It is exported because the answer decides how OTHER callers classify a failed +// fetch, and there must be exactly one line drawn. The firehose consumer's §5.4 +// direct fetch marks a genuine not-found as a PERMANENT event — dead-lettered +// with its redrive budget spent — while everything else stays transient and +// costs the connector three inline retries plus ten redrives. Two +// implementations of "is this a real not-found" would eventually disagree, and +// the disagreement would show up as either discarded posts or a blocked lane. +// +// THE TWO SHAPES IT ACCEPTS, and the one it deliberately does not: +// +// - 400 with an XRPC RecordNotFound (or an InvalidRequest whose message says +// it could not locate the record) — what the reference PDS answers when the +// repo exists and the record does not. +// - 404 WITH an XRPC error envelope — what some other implementations answer. +// - NOT a bare 404. With no envelope the request most likely never reached a +// PDS at all: a stale pds_url pointing at a reverse proxy or a generic web +// server, both of which answer 404 for everything. Trusting that would +// report a record as gone because somebody mistyped a hostname. +func IsRecordNotFoundResponse(statusCode int, body []byte) bool { + switch statusCode { + case http.StatusNotFound: + return isXRPCErrorBody(body) + case http.StatusBadRequest: + return isRecordNotFoundBody(body) + default: + return false + } +} + // isXRPCErrorBody reports whether body is an XRPC error JSON object (has a // non-empty "error" field). Used to distinguish a real PDS not-found response // from a bare 404 served by whatever non-PDS host a stale pds_url points at. -- 2.51.2