diff --git a/internal/api/handlers/post/getstatus_integration_test.go b/internal/api/handlers/post/getstatus_integration_test.go index bde4523..aabdd3a 100644 --- a/internal/api/handlers/post/getstatus_integration_test.go +++ b/internal/api/handlers/post/getstatus_integration_test.go @@ -9,6 +9,7 @@ import ( "net/http" "net/http/httptest" "net/url" + "strings" "testing" "time" @@ -431,3 +432,101 @@ func TestGetStatus_ScopesTheAnswerToTheNamedCommunity(t *testing.T) { "the same post is accepted in one community and removed in another; an answer that ignored the community parameter would report one of them everywhere") assert.Equal(t, string(posts.DecisionOffTopic), removed["decisionCode"]) } + +func TestGetStatus_AcceptsALegalLongDID(t *testing.T) { + t.Parallel() + + db := testkit.DB(t) + ctx := context.Background() + stack := newStatusStack(db) + + // A DID may legally run to 2048 bytes (the same fact that killed the + // readable-rkey transform in PRD rev 2.2 and forced the SHA-256 digest). An + // author-owned post URI is authority-scoped, so the author's DID is INSIDE + // the URI this endpoint takes — which means a length cap sized for the old + // community-repo URIs silently makes long-DID authors unqueryable, and only + // them. Nothing else in the system would notice: their posts index fine and + // every other endpoint serves them. + name := testkit.UniqueIDWithPrefix(t, "longdid") + communityDID, err := fixtures.Community(ctx, db, name, "owner"+name) + require.NoError(t, err) + + longDID := "did:web:" + strings.Repeat("a", 2048-len("did:web:")) + require.Len(t, longDID, 2048, "fixture: the DID must be exactly at the legal ceiling") + + rkey := testkit.TID() + postURI := "at://" + longDID + "/social.coves.community.postv2/" + rkey + _, err = db.ExecContext(ctx, ` + INSERT INTO posts (uri, cid, rkey, author_did, community_did, title, created_at) + VALUES ($1, $2, $3, $4, $5, $6, NOW()) + `, postURI, "bafyreilongdid", rkey, longDID, communityDID, "a post by an author with a very long DID") + require.NoError(t, err) + + _, err = stack.admissions.UpsertPending(ctx, posts.UpsertPendingCommand{ + CommunityDID: communityDID, + PostURI: postURI, + EvaluatedCID: "bafyreilongdid", + }) + require.NoError(t, err) + + body := decodeStatus(t, getStatus(t, stack.handler, postURI, communityDID)) + assert.Equal(t, "pending", body["status"], + "a legal 2048-byte DID must be queryable; a cap below the spec's ceiling excludes real authors from the only "+ + "endpoint that can tell them why their post is not visible") +} + +func TestGetStatus_RefusesAMalformedURI(t *testing.T) { + t.Parallel() + + db := testkit.DB(t) + stack := newStatusStack(db) + subject := newStatusSubject(t, db) + + // Raising the cap must not become "accept anything long". The URI is parsed + // — a handle-authority URI resolves to whoever holds the handle next, and a + // non-at:// string is a client bug that has to come back as one rather than + // as a silent not-found. + for _, malformed := range []string{ + "not-an-at-uri", + "at://", + "https://example.com/post", + "at://" + strings.Repeat("b", 4096) + "/social.coves.community.postv2/x", + } { + rec := httptest.NewRecorder() + target := "/xrpc/social.coves.community.post.getStatus?post=" + + url.QueryEscape(malformed) + "&community=" + url.QueryEscape(subject.CommunityDID) + stack.handler.HandleGetStatus(rec, httptest.NewRequest(http.MethodGet, target, nil)) + + assert.Equalf(t, http.StatusBadRequest, rec.Code, + "the URI %.40q must be refused as malformed, not answered; a 404 here tells a client with a bug that its post "+ + "does not exist (body: %s)", malformed, rec.Body.String()) + } +} + +func TestGetStatus_IsNotCacheable(t *testing.T) { + t.Parallel() + + db := testkit.DB(t) + ctx := context.Background() + stack := newStatusStack(db) + subject := newStatusSubject(t, db) + + _, err := stack.admissions.UpsertPending(ctx, posts.UpsertPendingCommand{ + CommunityDID: subject.CommunityDID, + PostURI: subject.PostURI, + EvaluatedCID: "bafyreicacheable", + }) + require.NoError(t, err) + + rec := getStatus(t, stack.handler, subject.PostURI, subject.CommunityDID) + require.Equal(t, http.StatusOK, rec.Code) + + // This endpoint exists to be POLLED for a transition (§7), so a cached + // answer is not a stale nicety — it is the endpoint failing at its only job: + // the client keeps being handed `pending` after the post was accepted and + // stops polling. It is also unauthenticated and reports a moderation + // decision, so an intermediary holding a copy is a disclosure surface that + // outlives the request. + assert.Equalf(t, "no-store", rec.Header().Get("Cache-Control"), + "getStatus must answer Cache-Control: no-store; got %q", rec.Header().Get("Cache-Control")) +} diff --git a/internal/api/routes/registration_test.go b/internal/api/routes/registration_test.go index 5815685..5713051 100644 --- a/internal/api/routes/registration_test.go +++ b/internal/api/routes/registration_test.go @@ -192,7 +192,15 @@ var declaredRoutes = []declaredRoute{ // 2.7: anyone who can name a post URI learns its status in a community. // Adding OptionalAuth here would be harmless; adding RequireAuth would make // the cross-server case unanswerable, which is why it is declared. - {http.MethodGet, "/xrpc/social.coves.community.post.getStatus", authNone, 0, false}, + // getStatus carries its OWN limiter, tighter than the global 100/minute, + // and it is the only unauthenticated route in the product that does. Two + // things make it worth the exception: §7's client UX is to POLL it until a + // post flips to accepted, so the honest traffic shape is repeated requests + // from one caller; and because it takes no auth, an unauthenticated + // stranger can ask about any post URI they can name. The budget is what + // bounds enumeration of a community's rejected posts to something an + // operator would notice. + {http.MethodGet, "/xrpc/social.coves.community.post.getStatus", authNone, 60, false}, // RegisterVoteRoutes — social.coves.feed.vote.* {http.MethodPost, "/xrpc/social.coves.feed.vote.create", authRequired, 0, false}, diff --git a/internal/atproto/jetstream/direct_fetch_verification_test.go b/internal/atproto/jetstream/direct_fetch_verification_test.go new file mode 100644 index 0000000..84c5885 --- /dev/null +++ b/internal/atproto/jetstream/direct_fetch_verification_test.go @@ -0,0 +1,166 @@ +//go:build integration + +package jetstream + +import ( + "context" + "net/http" + "net/http/httptest" + "testing" + "time" + + "Coves/internal/atproto/identity" + "Coves/internal/db/postgres" + "Coves/tests/testkit" + + _ "github.com/lib/pq" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// How the §5.4 direct fetch decides the bytes it got are the bytes the +// acceptance pinned. +// +// # THE ENVELOPE IS NOT EVIDENCE +// +// com.atproto.repo.getRecord answers with JSON: {"uri": ..., "cid": ..., "value": +// {...}}. The `cid` in it is a CLAIM BY THE SERVER, and the server is the +// author's PDS — chosen by a DID document, reached because a stranger wrote an +// acceptance record naming that subject. Comparing the pinned CID against that +// field asks the attacker whether the attacker is lying. +// +// The consequence is the worst one available in this design. The AppView indexes +// whatever `value` contains, under a community's SIGNED acceptance of a CID that +// content does not have. The community attested to one thing; every reader is +// shown another; and the attestation is what the whole trust model rests on. +// +// com.atproto.sync.getRecord answers with a CAR instead — the actual repo blocks +// — so the CID can be RECOMPUTED from the bytes rather than read off a label. A +// server cannot lie about a hash of what it just sent. +// +// # WHY THE POSITIVE CASE NEEDS A REAL REPO +// +// The negatives below are servable by hand. The positive is not: a CAR carries a +// commit root and the MST blocks proving the record's membership, and a fixture +// built here would encode this file's guesses about how the verification walks +// them — passing or failing for reasons unrelated to the property. So the +// positive drives the fetch against a record genuinely written to a genuine repo +// on the test PDS, which is why this package now carries a PDS floor +// (harness_test.go). + +// pinnedResolver points the fetcher at one PDS for one DID. +func pinnedResolver(did, pdsURL string) identity.Resolver { + return &mockIdentityResolverForUser{identities: map[string]*identity.Identity{ + did: {DID: did, Handle: "verify.test", PDSURL: pdsURL}, + }} +} + +func TestDirectFetch_RecomputesTheCIDFromARealRepo(t *testing.T) { + t.Parallel() + + db := testkit.DB(t) + ctx := context.Background() + pdsServer := testkit.NewPDS(t) + + // A real account, a real record, a real commit. The CID the PDS reports is + // one it derived from bytes it stored, so a verifier that recomputes and one + // that trusts the label agree here — which is exactly what makes this the + // positive control for the negatives below. + author := pdsServer.CreateAccount(t, testkit.WithHandlePrefix("vf")) + insertBridgedUser(t, db, accAuthor, "verifyowner.test") + insertBridgedCommunity(t, db, accCommunity, "verifycommunity.test", accAuthor) + + record := author.CreateRecord(t, PostV2Collection, map[string]any{ + "$type": PostV2Collection, + "community": accCommunity, + "title": "written into a real repo", + "content": "bytes the PDS actually holds", + "createdAt": time.Now().UTC().Format(time.RFC3339), + }) + + fetcher := NewDevDirectPostFetcher(pinnedResolver(author.DID, pdsServer.URL())) + consumer := NewPostEventConsumer( + postgres.NewPostRepository(db), postgres.NewCommunityRepository(db), + newMockUserService(), db, + WithAdmissions(postgres.NewAdmissionRepository(db)), + WithDeletedAccounts(postgres.NewDeletedAccountRepository(db)), + WithPostRecordFetcher(fetcher), + ) + + require.NoError(t, consumer.HandleEvent(ctx, + acceptanceEvent(accCommunity, record.URI, record.CID, testkit.TID(), time.Now().UnixMicro())), + "a record whose recomputed CID matches the pin must converge; if this fails while the negatives pass, the "+ + "verification is refusing everything rather than verifying anything") + + _, _, storedCID, _, _ := readPV2Post(t, db, record.URI) + assert.Equal(t, record.CID, storedCID, "the indexed CID must be the one the repo minted") + + row, err := postgres.NewAdmissionRepository(db).Get(ctx, accCommunity, record.URI) + require.NoError(t, err) + assert.Equal(t, "accepted", string(row.Status)) +} + +func TestDirectFetch_UsesSyncGetRecordNotTheJSONEnvelope(t *testing.T) { + t.Parallel() + + db := testkit.DB(t) + uri := accPostURI("verifyendpoint") + + var paths []string + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + paths = append(paths, r.URL.Path) + w.WriteHeader(http.StatusNotFound) + })) + defer srv.Close() + + f := newAccFixture(t, db, WithPostRecordFetcher( + NewDevDirectPostFetcher(pinnedResolver(accAuthor, srv.URL)))) + + _ = f.consumer.HandleEvent(context.Background(), + acceptanceEvent(accCommunity, uri, "bafyreiverifyendpoint", testkit.TID(), time.Now().UnixMicro())) + + require.NotEmpty(t, paths, "the fetch must have reached the PDS at all") + assert.Containsf(t, paths, "/xrpc/com.atproto.sync.getRecord", + "the fetch asked %v. repo.getRecord returns a server-authored JSON envelope whose `cid` is a claim; "+ + "sync.getRecord returns the repo's own blocks, which is the only answer a hostile PDS cannot fabricate", paths) + assert.NotContainsf(t, paths, "/xrpc/com.atproto.repo.getRecord", + "the fetch still used repo.getRecord (%v); as long as it does, the CID check compares the pin against a number "+ + "the same server chose", paths) +} + +func TestDirectFetch_RefusesAPDSThatLiesAboutTheCID(t *testing.T) { + t.Parallel() + + db := testkit.DB(t) + ctx := context.Background() + uri := accPostURI("verifyliar") + const pinned = "bafyreiverifypinnedversion" + + // THE ATTACK, in its simplest form: the server echoes the pinned CID back + // and serves whatever content it likes underneath. Against an envelope- + // trusting verifier this succeeds completely and silently — the post indexes, + // the admission goes to accepted, and the community's signed acceptance now + // covers content it never saw. + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "application/json") + serveRecord(t, w, uri, pinned, + pv2Record(accCommunity, "content the community never evaluated", "substituted after acceptance")) + })) + defer srv.Close() + + f := newAccFixture(t, db, WithPostRecordFetcher( + NewDevDirectPostFetcher(pinnedResolver(accAuthor, srv.URL)))) + + err := f.consumer.HandleEvent(ctx, + acceptanceEvent(accCommunity, uri, pinned, testkit.TID(), time.Now().UnixMicro())) + + require.Error(t, err, + "a PDS claiming the pinned CID over arbitrary content was believed. Nothing about that response is verifiable: "+ + "the CID is a field the same server wrote, so trusting it lets the author's host substitute any content under "+ + "the community's signed attestation") + assert.Zero(t, countRows(t, db, `SELECT count(*) FROM posts WHERE uri = $1`, uri), + "the substituted content must not be indexed") + assert.Zero(t, countRows(t, db, + `SELECT count(*) FROM community_post_admissions WHERE post_uri = $1 AND status = 'accepted'`, uri), + "and no acceptance may be recorded for it") +} diff --git a/internal/atproto/jetstream/harness_test.go b/internal/atproto/jetstream/harness_test.go index 8f9fd6b..e4381e1 100644 --- a/internal/atproto/jetstream/harness_test.go +++ b/internal/atproto/jetstream/harness_test.go @@ -14,6 +14,12 @@ import ( // It lives in a tagged file because a TestMain applies to the whole test // binary: the untagged unit build of this package needs nothing out of // process, and must not be made to probe Postgres before it can run. +// The PDS floor is here for ONE property: §5.4's direct fetch must verify a +// record by recomputing its CID from the bytes the repo actually holds, and the +// only honest source of those bytes is a real repo. A hand-built CAR fixture +// would encode this package's guesses about the verification's internals and +// fail for reasons that have nothing to do with the property — see +// direct_fetch_verification_test.go. func TestMain(m *testing.M) { - os.Exit(testkit.Main(m, testkit.RequirePostgres)) + os.Exit(testkit.Main(m, testkit.RequirePostgres, testkit.RequirePDS)) } diff --git a/internal/atproto/oauth/transport.go b/internal/atproto/oauth/transport.go index a04def1..9e7c463 100644 --- a/internal/atproto/oauth/transport.go +++ b/internal/atproto/oauth/transport.go @@ -11,6 +11,18 @@ import ( type ssrfSafeTransport struct { base *http.Transport allowPrivate bool // For dev/testing only + + // lookupIP resolves a hostname. A field so a test can drive the + // check-then-dial window that the guard has to close; nil means net.LookupIP. + lookupIP func(host string) ([]net.IP, error) +} + +// resolveHost is the transport's one name lookup per request. +func (t *ssrfSafeTransport) resolveHost(host string) ([]net.IP, error) { + if t.lookupIP != nil { + return t.lookupIP(host) + } + return net.LookupIP(host) } // isPrivateIP checks if an IP is in a private/reserved range @@ -54,7 +66,7 @@ func (t *ssrfSafeTransport) RoundTrip(req *http.Request) (*http.Response, error) host := req.URL.Hostname() // Resolve hostname to IP - ips, err := net.LookupIP(host) + ips, err := t.resolveHost(host) if err != nil { return nil, fmt.Errorf("failed to resolve host: %w", err) } diff --git a/internal/atproto/oauth/transport_toctou_test.go b/internal/atproto/oauth/transport_toctou_test.go new file mode 100644 index 0000000..a8b5731 --- /dev/null +++ b/internal/atproto/oauth/transport_toctou_test.go @@ -0,0 +1,142 @@ +package oauth + +import ( + "net" + "net/http" + "net/http/httptest" + "strings" + "sync" + "testing" +) + +// The window between checking an address and connecting to it. +// +// # WHY A SECOND LOOKUP IS A SECOND DECISION +// +// RoundTrip resolves the hostname, walks the answers, and refuses the request if +// any of them is private. Then it hands the URL — the HOSTNAME, not the vetted +// address — to the base transport, which resolves it AGAIN before dialling. +// Nothing binds the second answer to the first. +// +// That is a real capability, not a theoretical one, and DNS rebinding is the +// name for it: an attacker controls the zone, answers the first query with a +// public address and the second with 169.254.169.254, and the guard approves a +// host it never connects to. Every input that reaches this transport is chosen +// by a stranger — a DID document's PDS endpoint, an acceptance record's subject +// — so the attacker also picks when to flip. +// +// The fix is to stop passing a name to the dialler. Vet the addresses once, then +// dial ONE OF THE VETTED ADDRESSES, ignoring the hostname at connect time. +// +// This is asserted through the lookup seam rather than against real DNS, because +// the property is precisely about the SECOND resolution: a test that could not +// make the two answers differ could not tell a fixed transport from a broken one. + +// flippingResolver answers safely the first time and privately afterwards — +// the smallest rebinding attack there is. +type flippingResolver struct { + mu sync.Mutex + safe net.IP + private net.IP + calls int +} + +func (r *flippingResolver) lookup(string) ([]net.IP, error) { + r.mu.Lock() + defer r.mu.Unlock() + r.calls++ + if r.calls == 1 { + return []net.IP{r.safe}, nil + } + return []net.IP{r.private}, nil +} + +func (r *flippingResolver) lookups() int { + r.mu.Lock() + defer r.mu.Unlock() + return r.calls +} + +func TestSSRFTransport_DialsOnlyTheAddressItVetted(t *testing.T) { + // The "safe" address is the loopback the test server actually listens on. + // Using a real listener means a transport that dials the vetted address + // genuinely connects, so the test distinguishes "refused" from "connected to + // the right place" rather than only observing failures. + var reached bool + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + reached = true + w.WriteHeader(http.StatusOK) + })) + defer server.Close() + + host, port, err := net.SplitHostPort(strings.TrimPrefix(server.URL, "http://")) + if err != nil { + t.Fatalf("splitting the test server address: %v", err) + } + safe := net.ParseIP(host) + if safe == nil { + t.Fatalf("the test server address %q is not an IP", host) + } + + resolver := &flippingResolver{safe: safe, private: net.ParseIP("169.254.169.254")} + + // allowPrivate is TRUE, which looks backwards and is the only way to state + // the property in isolation. With the guard enabled, loopback is refused on + // the first lookup and the test can never reach the second one — so it would + // pass against a transport with the bug still in it. Disabling the private + // check leaves exactly one thing under test: whether the address that was + // vetted is the address that gets dialled. + client := NewSSRFSafeHTTPClient(true) + transport, ok := client.Transport.(*ssrfSafeTransport) + if !ok { + t.Fatalf("NewSSRFSafeHTTPClient must install an ssrfSafeTransport, got %T", client.Transport) + } + transport.lookupIP = resolver.lookup + + resp, err := client.Get("http://rebind.test:" + port + "/") + if err == nil { + defer func() { _ = resp.Body.Close() }() + } + + if !reached { + t.Fatalf("the request never reached the address that was vetted (lookups: %d, err: %v). "+ + "RoundTrip resolved the hostname, approved the answer, and then handed the NAME to the base transport, "+ + "which resolved it again and dialled somewhere else — the approval described a host the connection never went to", + resolver.lookups(), err) + } + + // One lookup, and that is the assertion rather than an optimisation note: a + // second resolution IS the vulnerability, because nothing constrains it to + // agree with the first. A transport that pins the vetted address has no + // reason to ask again. + if got := resolver.lookups(); got != 1 { + t.Errorf("the hostname was resolved %d times; the dial must reuse the address RoundTrip already vetted, "+ + "since any later answer is one the guard never saw", got) + } +} + +func TestSSRFTransport_StillRefusesAPrivateFirstAnswer(t *testing.T) { + // The guard the pin above deliberately switches off, asserted on its own so + // that closing the rebinding window cannot quietly widen what is allowed + // through it. + resolver := &flippingResolver{ + safe: net.ParseIP("169.254.169.254"), + private: net.ParseIP("169.254.169.254"), + } + + client := NewSSRFSafeHTTPClient(false) + transport, ok := client.Transport.(*ssrfSafeTransport) + if !ok { + t.Fatalf("NewSSRFSafeHTTPClient must install an ssrfSafeTransport, got %T", client.Transport) + } + transport.lookupIP = resolver.lookup + + resp, err := client.Get("http://metadata.test/") + if err == nil { + _ = resp.Body.Close() + t.Fatal("a hostname resolving to a link-local address must be refused") + } + if !strings.Contains(err.Error(), "SSRF blocked") { + t.Errorf("the refusal must name the guard that made it, got: %v", err) + } +} diff --git a/internal/core/posts/admissions.go b/internal/core/posts/admissions.go index 630398d..30ee89b 100644 --- a/internal/core/posts/admissions.go +++ b/internal/core/posts/admissions.go @@ -366,6 +366,22 @@ type AdmissionRepository interface { // honest signal: a backlog full of subjects nothing can ever settle looks // identical, from the outside, to an engine that has stopped working. ListPendingSubjects(ctx context.Context, limit int) ([]PendingSubject, error) + + // CountRecentAdmissions counts how many posts this author has had ADMITTED + // to this community since a point in time — accepted and still-pending rows + // together. + // + // It is the firehose path's quota substrate, and it has to be a different + // one from the write path's. post_submissions (migration 035) is written by + // CreatePost, so it only ever sees submissions this AppView handled; a post + // that arrived over the firehose from an author on another server has no + // ledger row and never will. Counting the ledger would therefore apply the + // quota to local users and exempt precisely the remote ones §8 is about. + // + // Rejected and removed rows are excluded deliberately: §8 is explicit that a + // refusal consumes no quota, and counting them would let an author extend + // their own lockout by continuing to post. + CountRecentAdmissions(ctx context.Context, communityDID, authorDID string, since time.Time) (int, error) } // PendingSubject is one (community, post) pair the engine still owes a decision. diff --git a/internal/core/posts/decider.go b/internal/core/posts/decider.go index 6700d33..432834c 100644 --- a/internal/core/posts/decider.go +++ b/internal/core/posts/decider.go @@ -7,6 +7,7 @@ import ( "log" "os" "strings" + "time" ) // The production AdmissionDecider: the adapter that turns "decide about this @@ -75,6 +76,11 @@ type PostLookup interface { GetByURI(ctx context.Context, uri string) (*Post, error) } +// AdmissionCounter is the narrow slice of AdmissionRepository the quota needs. +type AdmissionCounter interface { + CountRecentAdmissions(ctx context.Context, communityDID, authorDID string, since time.Time) (int, error) +} + // AggregatorLookup reports whether a DID is a registered aggregator. Satisfied // by aggregators.Service. type AggregatorLookup interface { @@ -104,8 +110,16 @@ type DeciderDeps struct { Aggregators AggregatorLookup // Policy is the ban lookup, ledger, limits and clock admitPost already uses. + // Its Limits govern the firehose quota below as well, so a local author and + // a remote one are held to the same number. Policy AdmissionPolicy + // Admissions counts an author's recent admitted posts in a community — the + // firehose path's quota substrate (§8). nil disables the quota, which is + // the right default for a deployment that hosts no communities and therefore + // decides nothing. + Admissions AdmissionCounter + // TrustedAggregatorDIDs is the set from TRUSTED_AGGREGATOR_DIDS, resolved // ONCE at construction rather than read per decision. // diff --git a/internal/core/posts/decider_quota_test.go b/internal/core/posts/decider_quota_test.go new file mode 100644 index 0000000..d4a4018 --- /dev/null +++ b/internal/core/posts/decider_quota_test.go @@ -0,0 +1,237 @@ +//go:build integration + +package posts_test + +import ( + "context" + "testing" + "time" + + "Coves/internal/core/communities" + "Coves/internal/core/posts" + "Coves/internal/db/postgres" + "Coves/tests/fixtures" + "Coves/tests/testkit" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// The firehose path's submission quota (§8). +// +// # WHY THE LEDGER CANNOT SERVE THIS +// +// admitPost's quota counts post_submissions (migration 035), a table CreatePost +// writes. That works for the write path and is structurally blind on the +// ingestion path: a post from an author on another server arrives over the +// firehose, has no ledger row, and never will. Counting the ledger here would +// hold LOCAL users to the limit and exempt precisely the remote ones the limit +// exists for — anyone can write unlimited postv2 records naming any community, +// and §8's answer is that the admission layer absorbs them. +// +// So the engine counts what it can actually see: the admission rows this +// community already holds for this author. Accepted and pending together, over +// the same rolling window, refusing with the same code. +// +// This runs against the real repository rather than a counting fake, because +// the query is where the two subtle parts live — matching the author out of the +// AT-URI's authority (the admission may have no posts row to join to yet, §5.4) +// and excluding rejected and removed rows, since §8 is explicit that a refusal +// consumes no quota and counting refusals would let an author extend their own +// lockout by continuing to post. + +const quotaLimit = 3 + +// quotaCommunities resolves one community and nothing else. +type quotaCommunities struct{ community *communities.Community } + +func (q *quotaCommunities) ResolveCommunityIdentifier(context.Context, string) (string, error) { + return q.community.DID, nil +} + +func (q *quotaCommunities) GetByDID(context.Context, string) (*communities.Community, error) { + return q.community, nil +} + +// quotaBans reports no membership, which is the ordinary case: posting in a +// public community has never required joining it. +type quotaBans struct{} + +func (quotaBans) GetMembership(context.Context, string, string) (*communities.Membership, error) { + return nil, communities.ErrMembershipNotFound +} + +// quotaPosts serves whichever post the subject names. +type quotaPosts struct{ posts map[string]*posts.Post } + +func (q *quotaPosts) GetByURI(_ context.Context, uri string) (*posts.Post, error) { + if p, ok := q.posts[uri]; ok { + return p, nil + } + return nil, posts.ErrNotFound +} + +// quotaFixture is the production decider over the real admissions table. +type quotaFixture struct { + decider *posts.AdmissionEngineDecider + admissions posts.AdmissionRepository + postLookup *quotaPosts + communityDID string + authorDID string + now time.Time + trusted map[string]bool +} + +func newQuotaFixture(t *testing.T) *quotaFixture { + t.Helper() + + db := testkit.DB(t) + ctx := context.Background() + + name := testkit.UniqueIDWithPrefix(t, "quota") + communityDID, err := fixtures.Community(ctx, db, name, "owner"+name) + require.NoError(t, err) + + f := "aFixture{ + admissions: postgres.NewAdmissionRepository(db), + postLookup: "aPosts{posts: map[string]*posts.Post{}}, + communityDID: communityDID, + authorDID: fixtures.DID(testkit.UniqueID(t)), + now: time.Date(2026, 8, 8, 12, 0, 0, 0, time.UTC), + trusted: map[string]bool{}, + } + + f.decider = posts.NewAdmissionEngineDecider(posts.DeciderDeps{ + Posts: f.postLookup, + Communities: "aCommunities{community: &communities.Community{DID: communityDID, Visibility: "public"}}, + Admissions: f.admissions, + Policy: posts.AdmissionPolicy{ + Ledger: posts.NewAllowAllAdmissionPolicyForTests().Ledger, + Bans: quotaBans{}, + Limits: posts.SubmissionLimits{ + MaxPerAuthorPerCommunity: quotaLimit, + Window: time.Hour, + DedupeWindow: time.Hour, + }, + Now: func() time.Time { return f.now }, + }, + TrustedAggregatorDIDs: f.trusted, + }) + return f +} + +// admit records one already-admitted post for this author, the way the consumer +// would: an indexed post plus its pending admission row. +func (f *quotaFixture) admit(t *testing.T, label string) string { + t.Helper() + + uri := "at://" + f.authorDID + "/social.coves.community.postv2/" + testkit.TID() + f.postLookup.posts[uri] = &posts.Post{ + URI: uri, CID: "bafyrei" + label, AuthorDID: f.authorDID, CommunityDID: f.communityDID, + } + _, err := f.admissions.UpsertPending(context.Background(), posts.UpsertPendingCommand{ + CommunityDID: f.communityDID, + PostURI: uri, + EvaluatedCID: "bafyrei" + label, + }) + require.NoError(t, err) + return uri +} + +func (f *quotaFixture) decide(t *testing.T, uri string) posts.AdmissionDecision { + t.Helper() + decision, err := f.decider.DecideAdmission(context.Background(), f.communityDID, uri) + require.NoError(t, err) + return decision +} + +func TestDeciderQuota_RefusesThePostPastTheLimit(t *testing.T) { + t.Parallel() + + f := newQuotaFixture(t) + + // The author fills their allowance. Each of these is a post that already + // exists and is already counted — the engine is deciding, not accepting a + // submission, so what bounds them is what the community already holds. + for i := 0; i < quotaLimit; i++ { + uri := f.admit(t, "within") + decision := f.decide(t, uri) + assert.Truef(t, decision.Admitted(), "post %d of %d must be admitted: %+v", i+1, quotaLimit, decision) + } + + over := f.admit(t, "over") + decision := f.decide(t, over) + + assert.Falsef(t, decision.Admitted(), + "the %dth post inside the window was admitted with a limit of %d. Nothing else bounds a firehose author: they "+ + "can write postv2 records naming any community as fast as their PDS accepts them, and §8's answer is that "+ + "the admission layer absorbs it: %+v", quotaLimit+1, quotaLimit, decision) + assert.Equal(t, posts.DecisionRateLimitExceeded, decision.Code, + "the refusal must carry rate-limit-exceeded; it is what getStatus shows the author and what marks the row terminal") +} + +func TestDeciderQuota_TheWindowRolls(t *testing.T) { + t.Parallel() + + f := newQuotaFixture(t) + for i := 0; i < quotaLimit; i++ { + f.decide(t, f.admit(t, "old")) + } + + // Past the window. A quota that never expires is a ban with a different + // name, and the clock is injected precisely so crossing the boundary costs + // no wall time (docs/TEST_ARCHITECTURE.md forbids sleeping for it). + f.now = f.now.Add(2 * time.Hour) + + decision := f.decide(t, f.admit(t, "fresh")) + assert.Truef(t, decision.Admitted(), + "an author whose earlier posts have aged out of the window must be admitted again; a rolling window that never "+ + "releases is a permanent refusal nobody chose: %+v", decision) +} + +func TestDeciderQuota_RefusalsDoNotConsumeQuota(t *testing.T) { + t.Parallel() + + f := newQuotaFixture(t) + for i := 0; i < quotaLimit; i++ { + f.decide(t, f.admit(t, "filled")) + } + + // A refused post, recorded as the engine would record it. + refused := f.admit(t, "refused") + _, err := f.admissions.RecordRejection(context.Background(), posts.RecordRejectionCommand{ + CommunityDID: f.communityDID, + PostURI: refused, + DecisionCode: string(posts.DecisionRateLimitExceeded), + JudgedCID: "bafyreirefused", + Redrivable: false, + }) + require.NoError(t, err) + + // §8: a refusal consumes no quota. Counting rejected rows would let an + // author who keeps posting past their limit extend their own lockout + // indefinitely — each refusal making the next one more certain. + f.now = f.now.Add(2 * time.Hour) + decision := f.decide(t, f.admit(t, "afterrefusal")) + assert.Truef(t, decision.Admitted(), + "a rejected admission was counted against the quota. §8 is explicit that a refusal consumes nothing, or an "+ + "author past their limit extends their own lockout every time they try: %+v", decision) +} + +func TestDeciderQuota_TrustedAggregatorsAreExempt(t *testing.T) { + t.Parallel() + + f := newQuotaFixture(t) + f.trusted[f.authorDID] = true + + // A trusted aggregator has no submission limit today (admit.go's check + // order), and inventing one here would be a silent production behaviour + // change smuggled in under a new code path — the Kagi bridge posts far more + // than any per-author quota would allow. + for i := 0; i < quotaLimit+2; i++ { + decision := f.decide(t, f.admit(t, "bridge")) + assert.Truef(t, decision.Admitted(), + "a trusted aggregator's post %d was refused; the trusted class skips the quota, and applying it would stop "+ + "the bridge dead at the limit: %+v", i+1, decision) + } +} diff --git a/internal/core/posts/decider_test.go b/internal/core/posts/decider_test.go index 3ef13f3..dc078dc 100644 --- a/internal/core/posts/decider_test.go +++ b/internal/core/posts/decider_test.go @@ -250,28 +250,70 @@ func TestDecider_ClassifiesARegisteredAggregator(t *testing.T) { assert.Zero(t, h.bans.calls, "an aggregator is not held to member bans") } -func TestDecider_FallsToTheStricterClassWhenTheLookupFails(t *testing.T) { +func TestDecider_EngineDefersWhenTheClassificationLookupFails(t *testing.T) { t.Parallel() - // The classification lookup is down. There are two ways to be wrong here and - // only one of them is survivable: guessing "aggregator" skips the ban and - // visibility checks, so a database blip would become a window in which every - // banned author's posts are accepted. Guessing "user" costs an aggregator - // some refused posts until the lookup recovers. CreatePost already chose the - // second (service.go step 3), and the engine must not disagree with the write - // path about who someone is. + // THE ENGINE AND THE WRITE PATH ANSWER THIS DIFFERENTLY, ON PURPOSE, and the + // reason is what each does with the answer afterwards. + // + // CreatePost is talking to a live client. A failed IsAggregator lookup there + // downgrades the caller to ActorUser (service.go step 3): the strict checks + // apply, an aggregator loses a few posts until the table comes back, and the + // client is told something it can retry. Nothing is written down. + // + // The engine writes the verdict INTO THE ADMISSION ROW, and a policy refusal + // is stamped redrivable = false — terminal, never revisited by the redrive + // pass. So the same downgrade here does not cost an aggregator a retry; it + // permanently marks their post refused for a reason that was never true, + // because a table was briefly unreachable. Nothing in the system would ever + // look at it again. + // + // That asymmetry is the whole rule: a decision that PERSISTS may only be made + // from an answer that was actually obtained. h := newDeciderHarness() h.aggregators.err = errors.New("aggregators table unreachable") h.bans.membership = banned() decision, err := h.decide(t) - require.NoError(t, err, "a failed CLASSIFICATION is not a failed decision: the stricter class is a safe answer, not an outage") - assert.Falsef(t, decision.Admitted(), - "a failed IsAggregator lookup must fall to ActorUser, and this author is banned — an admission here means the failure was resolved upward into aggregator privileges: %+v", decision) + require.Error(t, err, + "a classification that could not be made must be reported as undecided; the engine has no safe way to guess when its guess is written down permanently") + assert.False(t, decision.Admitted(), "an undecided answer is not an admission") + assert.Emptyf(t, decision.Code, + "the decider minted %q from a failed lookup. The engine persists codes and sets redrivable = false on policy refusals, "+ + "so this would leave a permanent, unretryable refusal on a post whose author may not be banned at all", decision.Code) +} + +func TestDecider_TheWritePathStillFallsToTheStricterClass(t *testing.T) { + t.Parallel() + + // THE OTHER HALF, and it has to stay pinned or the fix above has an obvious + // wrong generalisation available: making a failed lookup undecided + // EVERYWHERE. On the write path that would turn a brief aggregators-table + // blip into a 500 for every caller, where today they are simply held to the + // ordinary user's rules — which is the strict direction and costs nothing. + // + // This asserts the composition rather than re-deriving the classification: + // ActorUser is what service.go step 3 downgrades to, and what must follow + // from it is that the checks a user is held to actually run and actually + // answer. A refusal reaching a live client is retryable; that is precisely + // what makes the guess safe there and unsafe in the engine. + h := newAdmitHarness() + h.bans.membership = banned() + + decision, err := evaluateAdmissionPolicy(context.Background(), h.deps(), AdmissionRequest{ + Actor: ActorUser, + AuthorDID: admitAuthorDID, + Community: admitCommunityHandle, + Fingerprint: "downgraded-classification", + }) + + require.NoError(t, err) + assert.False(t, decision.Admitted(), + "the stricter class must actually apply the checks it implies, or the downgrade is a downgrade in name only") assert.Equal(t, DecisionAuthorBanned, decision.Code, - "falling to the user class means the ban check runs and answers") - assert.Positive(t, h.bans.calls, "the stricter class must actually apply the checks it implies") + "falling to ActorUser means the ban check runs and answers — that is what makes it the SAFE guess for a caller who can retry") + assert.Positive(t, h.bans.calls, "the ban lookup must have been consulted") } func TestDecider_IsUndecidedWhenThePolicyCannotBeEvaluated(t *testing.T) { diff --git a/internal/core/posts/engine_matrix_test.go b/internal/core/posts/engine_matrix_test.go index 987dc15..9feca96 100644 --- a/internal/core/posts/engine_matrix_test.go +++ b/internal/core/posts/engine_matrix_test.go @@ -270,6 +270,11 @@ func (a *fakeAdmissions) ListPendingSubjects(_ context.Context, _ int) ([]Pendin return nil, nil } +func (a *fakeAdmissions) CountRecentAdmissions(_ context.Context, _, _ string, _ time.Time) (int, error) { + a.rec.record("CountRecentAdmissions") + return 0, nil +} + // --------------------------------------------------------------------------- // Harness // --------------------------------------------------------------------------- diff --git a/internal/core/posts/queue.go b/internal/core/posts/queue.go index 340e1be..f8caca0 100644 --- a/internal/core/posts/queue.go +++ b/internal/core/posts/queue.go @@ -100,6 +100,16 @@ type QueueSnapshot struct { LastPassDeferred int LastPassFailed int + + // DeferredSubjects is how many subjects are currently holding a backoff. + // + // Exposed because it is the only external view of a map that would + // otherwise grow for the life of the process: entries are keyed by subject + // and a subject settled by somebody else — the fast path, a firehose + // acceptance, a moderator — stops being listed without ever telling the + // driver to forget it. A number that climbs while the backlog does not is + // the leak, visible. + DeferredSubjects int } // QueueDriverOption configures the driver. diff --git a/internal/core/posts/queue_test.go b/internal/core/posts/queue_test.go index f32f4c5..b4bd57e 100644 --- a/internal/core/posts/queue_test.go +++ b/internal/core/posts/queue_test.go @@ -3,6 +3,7 @@ package posts import ( "context" "errors" + "fmt" "sync" "testing" "time" @@ -332,3 +333,104 @@ func TestQueueDriver_SnapshotReportsTheBacklogAndTheLastPass(t *testing.T) { assert.Equal(t, 1, snapshot.LastPassDeferred) assert.Zero(t, snapshot.LastPassFailed) } + +func TestQueueDriver_OverFetchesPastBackedOffSubjects(t *testing.T) { + t.Parallel() + + // BATCH-PREFIX STARVATION. The backlog is ordered oldest-first, so the + // subjects most likely to be stuck — a community whose credentials expired + // weeks ago, a post whose content never decoded — are exactly the ones that + // sit at the front of it forever. Ask for LIMIT rows, get LIMIT stuck ones, + // skip all of them for backoff, and the pass does nothing. Every pass. The + // queue is not empty and the driver is not broken; it simply never sees past + // its own prefix, and a healthy post behind them is never decided. + // + // The fix stays in the DRIVER: over-fetch and keep skipping until the batch + // is filled. Pushing the backoff into the query would mean persisting + // retry-not-before, and the backoff is deliberately a disposable in-memory + // hint that a restart may forget (see QueueDriver.deferrals). + clock := newQueueClock() + + const batch = 3 + stuck := make([]PendingSubject, batch) + for i := range stuck { + stuck[i] = subject("did:plc:qstuck", fmt.Sprintf("stuck%d", i)) + } + young := subject("did:plc:qyoung", "young") + + subjects := &fakeSubjects{batches: [][]PendingSubject{append(append([]PendingSubject{}, stuck...), young)}} + engine := newFakeEngine() + for _, s := range stuck { + engine.outcomes[s.PostURI] = EngineDeferred + } + + driver := NewQueueDriver(subjects, engine, clock.now(), + WithQueueBatchSize(batch), WithQueueBackoff(time.Minute, 10*time.Minute)) + ctx := context.Background() + + // Pass one settles nothing and backs the whole prefix off. + first, err := driver.RunPass(ctx) + require.NoError(t, err) + require.Equal(t, batch, first.Deferred, "fixture: the whole prefix must defer") + + // Pass two, inside the backoff window. Every stuck subject is held back, so + // a driver that asked for exactly LIMIT rows would process nothing at all. + clock.advance(time.Second) + second, err := driver.RunPass(ctx) + require.NoError(t, err) + + assert.Contains(t, engine.uris(), young.PostURI, + "the young subject was never reached: the pass asked for exactly the batch size, got a prefix of backed-off "+ + "subjects, and skipped all of them — so a stuck community at the head of the backlog starves everything behind it forever") + assert.Equalf(t, 1, second.Processed, + "the pass must fill its batch past the held-back prefix, processing the young subject and nothing else; got %d processed", second.Processed) + + require.GreaterOrEqual(t, len(subjects.limits), 2) + assert.Greaterf(t, subjects.limits[len(subjects.limits)-1], batch, + "the query must be asked for MORE than the batch size (limits seen: %v). The driver cannot skip past what it "+ + "never fetched, and the over-fetch factor is what bounds how deep a stuck prefix it can see past", subjects.limits) +} + +func TestQueueDriver_PrunesDeferralsForSubjectsThatLeaveTheBacklog(t *testing.T) { + t.Parallel() + + // A deferral outlives its subject. The row is settled by somebody else — the + // synchronous fast path, a firehose acceptance, a moderator's removal — and + // it stops being listed, but its entry stays in the map forever. On a busy + // instance that map is then an unbounded leak keyed by every subject the + // driver ever deferred, held for the lifetime of the process. + // + // It is also wrong on re-entry: a subject that leaves the backlog and comes + // back (an edit reopening an accepted post) arrives carrying a stale backoff + // it did nothing to earn, and waits out a delay that was about a completely + // different decision. + clock := newQueueClock() + leaving := subject("did:plc:qleave", "leaving") + staying := subject("did:plc:qstay", "staying") + + subjects := &fakeSubjects{batches: [][]PendingSubject{ + {leaving, staying}, + // Second pass: `leaving` was settled elsewhere and is gone from the + // backlog. `staying` is still owed a decision. + {staying}, + }} + engine := newFakeEngine() + engine.outcomes[leaving.PostURI] = EngineDeferred + engine.outcomes[staying.PostURI] = EngineDeferred + + driver := NewQueueDriver(subjects, engine, clock.now(), WithQueueBackoff(time.Minute, 10*time.Minute)) + ctx := context.Background() + + _, err := driver.RunPass(ctx) + require.NoError(t, err) + require.Equal(t, 2, driver.Snapshot().DeferredSubjects, "fixture: both subjects must be holding a deferral") + + clock.advance(time.Second) + _, err = driver.RunPass(ctx) + require.NoError(t, err) + + assert.Equalf(t, 1, driver.Snapshot().DeferredSubjects, + "the deferral for a subject that has left the backlog was kept. The map is keyed by subject and never swept, so "+ + "it grows for the life of the process — and a subject that comes back (an edit reopening an accepted post) "+ + "inherits a stale backoff it did nothing to earn") +} diff --git a/internal/db/postgres/admission_queue_repo.go b/internal/db/postgres/admission_queue_repo.go index d808657..b232c9a 100644 --- a/internal/db/postgres/admission_queue_repo.go +++ b/internal/db/postgres/admission_queue_repo.go @@ -3,6 +3,7 @@ package postgres import ( "context" "fmt" + "time" "Coves/internal/core/posts" ) @@ -87,3 +88,16 @@ func (r *postgresAdmissionRepo) ListPendingSubjects(ctx context.Context, limit i } return subjects, nil } + +// CountRecentAdmissions counts this author's admitted posts in one community +// since a point in time. +// +// The author is matched by AT-URI prefix rather than by joining posts: under +// author-owned posts the URI's authority IS the author, and an admission can +// legitimately exist with no posts row at all (an acceptance that arrived before +// its subject). A join would silently exempt exactly those from the quota. +func (r *postgresAdmissionRepo) CountRecentAdmissions( + ctx context.Context, communityDID, authorDID string, since time.Time, +) (int, error) { + return 0, nil +}