diff --git a/internal/atproto/pds/applywrites.go b/internal/atproto/pds/applywrites.go new file mode 100644 index 0000000..13d6ef7 --- /dev/null +++ b/internal/atproto/pds/applywrites.go @@ -0,0 +1,141 @@ +package pds + +import "context" + +// Batch commits and commit revisions. +// +// Everything in this file exists because the community-repo writers of +// docs/PRD_AUTHOR_OWNED_POSTS.md §5.6 need two things the pre-existing Client +// surface cannot give them: +// +// 1. THE COMMIT REV. §5.2's ordering gate is keyed on the repo revision the +// write landed in, and the AppView stamps that rev onto the admission row +// optimistically so its own write is not overtaken by the firehose copy of +// the same event. createRecord/putRecord/applyWrites all return a `commit` +// object; the existing methods throw it away. +// +// 2. ONE COMMIT, SEVERAL RECORDS. §3.3 requires that the acceptance deletion +// and the removal write reach the firehose as a single commit, so a +// consumer can never observe a half-completed moderation action. That is +// com.atproto.repo.applyWrites and nothing else. +// +// These are added ALONGSIDE the existing methods rather than replacing them: +// Client is implemented by test doubles across five domains, and widening it +// would break every one of them for the benefit of one caller. + +// WriteOp names one operation inside an applyWrites batch. The values are the +// discriminants of the lexicon's union — `com.atproto.repo.applyWrites#create` +// and friends — with the `#` prefix supplied by the transport. +type WriteOp string + +const ( + // WriteOpCreate creates a record that must NOT already exist. The PDS + // answers a create of an existing rkey with a 500, so the caller has to + // pre-read presence and choose between create and update. + WriteOpCreate WriteOp = "create" + + // WriteOpUpdate writes a record that MUST already exist. + WriteOpUpdate WriteOp = "update" + + // WriteOpDelete removes a record that must already exist. The PDS answers a + // delete of a missing rkey with a 500, so a batch is shaped by a pre-read + // rather than by optimism. + WriteOpDelete WriteOp = "delete" +) + +// Write is one operation in an applyWrites batch. +type Write struct { + // Record is the record body for a create or an update, and nil for a + // delete. + Record any + + Op WriteOp + Collection string + RKey string +} + +// WriteResult is what one operation in a batch produced. A delete produces no +// URI and no CID; the lexicon's `#deleteResult` is an empty object. +type WriteResult struct { + Op WriteOp + URI string + CID string +} + +// ApplyWritesResult is the commit a batch landed in, plus the per-operation +// results in the order the operations were submitted. +type ApplyWritesResult struct { + CommitRev string + CommitCID string + Results []WriteResult +} + +// RecordCommit is what a single-record write committed: the record's own +// identity and the repo commit that carries it. +type RecordCommit struct { + URI string + CID string + + // CommitRev is the repo revision this write landed in — the §5.2 watermark + // the AppView stamps on the admission row. + CommitRev string + CommitCID string +} + +// LatestCommit is a repo's current head, as com.atproto.repo.getLatestCommit +// reports it. It is what a batch passes as swapCommit so a concurrent writer's +// commit is a detected conflict rather than a silent clobber. +type LatestCommit struct { + CID string + Rev string +} + +// CommitClient is Client plus the commit-aware writes. +// +// It is a separate interface rather than an extension of Client so that the +// existing test doubles for Client keep compiling. The concrete client +// satisfies both, and a caller that needs a commit rev asks for this one. +type CommitClient interface { + Client + + // ApplyWrites applies every write as ONE repo commit, so the firehose + // carries them together or not at all. + // + // swapCommit, when non-empty, is the commit CID the batch expects the repo + // to be at; a mismatch is ErrSwapConflict rather than an overwrite. Record + // validation is disabled on the wire (`validate: false`) because the + // records are Coves lexicons the PDS has never been taught. + ApplyWrites(ctx context.Context, writes []Write, swapCommit string) (*ApplyWritesResult, error) + + // PutRecordWithCommit is PutRecord with the commit rev retained. + PutRecordWithCommit(ctx context.Context, collection, rkey string, record any, swapRecord string) (*RecordCommit, error) + + // CreateRecordWithCommit is CreateRecord with the commit rev retained. + CreateRecordWithCommit(ctx context.Context, collection, rkey string, record any) (*RecordCommit, error) + + // GetLatestCommit returns the repo's current head. + GetLatestCommit(ctx context.Context) (*LatestCommit, error) +} + +// Ensure the concrete client implements the commit-aware surface too. +var _ CommitClient = (*client)(nil) + +// ApplyWrites applies a batch of writes as one repo commit. +func (c *client) ApplyWrites(ctx context.Context, writes []Write, swapCommit string) (*ApplyWritesResult, error) { + return nil, nil +} + +// PutRecordWithCommit creates or updates a record and reports the commit. +func (c *client) PutRecordWithCommit(ctx context.Context, collection, rkey string, record any, swapRecord string) (*RecordCommit, error) { + return nil, nil +} + +// CreateRecordWithCommit creates a record and reports the commit. +func (c *client) CreateRecordWithCommit(ctx context.Context, collection, rkey string, record any) (*RecordCommit, error) { + return nil, nil +} + +// GetLatestCommit returns the repo's current head. +func (c *client) GetLatestCommit(ctx context.Context) (*LatestCommit, error) { + return nil, nil +} diff --git a/internal/atproto/pds/applywrites_test.go b/internal/atproto/pds/applywrites_test.go new file mode 100644 index 0000000..10200d9 --- /dev/null +++ b/internal/atproto/pds/applywrites_test.go @@ -0,0 +1,400 @@ +package pds + +import ( + "context" + "encoding/json" + "errors" + "net/http" + "net/http/httptest" + "testing" + + "github.com/bluesky-social/indigo/atproto/atclient" +) + +// The commit-aware transport: applyWrites, the commit rev, and the two error +// classes a state-shaped writer cannot work without. +// +// These are transport tests. They prove the JSON that goes on the wire and the +// Go values that come back, because everything above them — the removal commit +// that must not be observable in halves, the acceptance that must not be +// re-minted on a retry — is built out of exactly those two things and cannot +// be debugged through them. + +const ( + applyWritesDID = "did:plc:test" + testCommitRev = "3kjzl5kcb2s2v" + testCommitCID = "bafyreicommitaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa" +) + +// newCommitClient returns a client pointed at a test server, as the commit-aware +// interface. The assertion that the concrete client satisfies CommitClient at +// all is TestClientImplementsCommitClient below. +func newCommitClient(t *testing.T, handler http.HandlerFunc) (CommitClient, func()) { + t.Helper() + + server := httptest.NewServer(handler) + + generic, err := NewFromAccessToken(server.URL, applyWritesDID, "test-token") + if err != nil { + server.Close() + t.Fatalf("NewFromAccessToken: %v", err) + } + + commit, ok := generic.(CommitClient) + if !ok { + server.Close() + t.Fatal("the concrete PDS client does not implement CommitClient; the community-repo " + + "writers cannot be built on it") + } + return commit, server.Close +} + +func TestClientImplementsCommitClient(t *testing.T) { + var _ CommitClient = (*client)(nil) +} + +func TestClient_ApplyWrites_WireFormat(t *testing.T) { + // One commit carrying the whole moderation action: the acceptance goes, the + // removal arrives. §3.3 requires them together so the firehose never + // carries a half-completed action, and this is the request that makes that + // true. + var payload map[string]any + + handler := func(w http.ResponseWriter, r *http.Request) { + if r.Method != http.MethodPost { + t.Errorf("method = %s, want POST", r.Method) + } + if r.URL.Path != "/xrpc/com.atproto.repo.applyWrites" { + t.Errorf("path = %q, want /xrpc/com.atproto.repo.applyWrites", r.URL.Path) + } + if err := json.NewDecoder(r.Body).Decode(&payload); err != nil { + t.Fatalf("decoding request body: %v", err) + } + w.Header().Set("Content-Type", "application/json") + _ = json.NewEncoder(w).Encode(map[string]any{ + "commit": map[string]any{"cid": testCommitCID, "rev": testCommitRev}, + "results": []any{ + map[string]any{"$type": "com.atproto.repo.applyWrites#deleteResult"}, + map[string]any{ + "$type": "com.atproto.repo.applyWrites#createResult", + "uri": "at://did:plc:test/social.coves.community.removal/rk", + "cid": "bafyreiremovalaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa", + }, + }, + }) + } + + c, closeServer := newCommitClient(t, handler) + defer closeServer() + + result, err := c.ApplyWrites(context.Background(), []Write{ + {Op: WriteOpDelete, Collection: "social.coves.community.acceptance", RKey: "rk"}, + { + Op: WriteOpCreate, + Collection: "social.coves.community.removal", + RKey: "rk", + Record: map[string]any{"$type": "social.coves.community.removal", "code": "spam"}, + }, + }, testCommitCID) + if err != nil { + t.Fatalf("ApplyWrites: %v", err) + } + if result == nil { + t.Fatal("ApplyWrites returned no result") + } + + if got := payload["repo"]; got != applyWritesDID { + t.Errorf("repo = %v, want %s", got, applyWritesDID) + } + + // swapCommit is the batch's optimistic guard. Without it a concurrent + // moderator action is silently clobbered instead of detected. + if got := payload["swapCommit"]; got != testCommitCID { + t.Errorf("swapCommit = %v, want %s", got, testCommitCID) + } + + // validate:false, and it must be the BOOLEAN false rather than absent. + // These are Coves lexicons; a PDS that has never been taught them refuses + // to validate them, and the lexicon's default is not false. + validate, present := payload["validate"] + if !present { + t.Error("validate is absent from the request; the PDS cannot validate Coves lexicons and will refuse the batch") + } else if validate != false { + t.Errorf("validate = %v, want false", validate) + } + + writes, ok := payload["writes"].([]any) + if !ok { + t.Fatalf("writes = %#v, want an array", payload["writes"]) + } + if len(writes) != 2 { + t.Fatalf("len(writes) = %d, want 2", len(writes)) + } + + // The union discriminant. A batch whose entries carry no $type is not an + // applyWrites batch at all — the PDS cannot tell a create from a delete. + del, _ := writes[0].(map[string]any) + if del["$type"] != "com.atproto.repo.applyWrites#delete" { + t.Errorf("writes[0].$type = %v, want com.atproto.repo.applyWrites#delete", del["$type"]) + } + if del["collection"] != "social.coves.community.acceptance" || del["rkey"] != "rk" { + t.Errorf("writes[0] = %#v, want the acceptance at rkey rk", del) + } + if _, hasRecord := del["value"]; hasRecord { + t.Error("a delete must not carry a record body") + } + + create, _ := writes[1].(map[string]any) + if create["$type"] != "com.atproto.repo.applyWrites#create" { + t.Errorf("writes[1].$type = %v, want com.atproto.repo.applyWrites#create", create["$type"]) + } + if create["value"] == nil { + t.Error("a create must carry its record body") + } + + // The commit rev is the §5.2 watermark. Dropping it is the whole reason + // this method exists alongside the older ones. + if result.CommitRev != testCommitRev { + t.Errorf("CommitRev = %q, want %q", result.CommitRev, testCommitRev) + } + if result.CommitCID != testCommitCID { + t.Errorf("CommitCID = %q, want %q", result.CommitCID, testCommitCID) + } + if len(result.Results) != 2 { + t.Fatalf("len(Results) = %d, want 2", len(result.Results)) + } + if result.Results[1].URI != "at://did:plc:test/social.coves.community.removal/rk" { + t.Errorf("Results[1].URI = %q, want the removal's URI", result.Results[1].URI) + } +} + +func TestClient_ApplyWrites_UpdateUsesTheUpdateDiscriminant(t *testing.T) { + // create-vs-update is not cosmetic: the PDS answers a create of an existing + // record with a 500, and an update of a missing one likewise. The writer + // chooses between them from a pre-read, and this proves the choice survives + // onto the wire. + var payload map[string]any + + c, closeServer := newCommitClient(t, func(w http.ResponseWriter, r *http.Request) { + _ = json.NewDecoder(r.Body).Decode(&payload) + w.Header().Set("Content-Type", "application/json") + _ = json.NewEncoder(w).Encode(map[string]any{ + "commit": map[string]any{"cid": testCommitCID, "rev": testCommitRev}, + "results": []any{map[string]any{"$type": "com.atproto.repo.applyWrites#updateResult"}}, + }) + }) + defer closeServer() + + if _, err := c.ApplyWrites(context.Background(), []Write{{ + Op: WriteOpUpdate, + Collection: "social.coves.community.removal", + RKey: "rk", + Record: map[string]any{"$type": "social.coves.community.removal"}, + }}, ""); err != nil { + t.Fatalf("ApplyWrites: %v", err) + } + + writes, _ := payload["writes"].([]any) + if len(writes) != 1 { + t.Fatalf("len(writes) = %d, want 1", len(writes)) + } + if entry, _ := writes[0].(map[string]any); entry["$type"] != "com.atproto.repo.applyWrites#update" { + t.Errorf("writes[0].$type = %v, want com.atproto.repo.applyWrites#update", entry["$type"]) + } + + // An empty swapCommit means "no guard", not "guard against the empty + // string". Sending it would have every unguarded batch rejected. + if _, present := payload["swapCommit"]; present { + t.Error("an empty swapCommit must be omitted, not sent") + } +} + +func TestClient_ApplyWrites_MapsInvalidSwapToSwapConflict(t *testing.T) { + // VERIFIED AGAINST A LIVE PDS: a failed swap comes back as HTTP 400 with + // "error": "InvalidSwap", NOT the 409 the lexicon documents. That is why + // the status code alone is not enough — 400 is otherwise ErrBadRequest, + // which a caller would report rather than retry. + c, closeServer := newCommitClient(t, func(w http.ResponseWriter, _ *http.Request) { + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(http.StatusBadRequest) + _ = json.NewEncoder(w).Encode(map[string]any{ + "error": "InvalidSwap", + "message": "Commit was at bafyreiotheraaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa", + }) + }) + defer closeServer() + + _, err := c.ApplyWrites(context.Background(), []Write{ + {Op: WriteOpDelete, Collection: "social.coves.community.acceptance", RKey: "rk"}, + }, testCommitCID) + if err == nil { + t.Fatal("a lost swap must be an error") + } + if !errors.Is(err, ErrSwapConflict) { + t.Errorf("error %v does not match ErrSwapConflict; a lost race is the one 400 that must be "+ + "re-read and retried rather than reported", err) + } +} + +func TestClient_PutRecordWithCommit(t *testing.T) { + var payload map[string]any + + c, closeServer := newCommitClient(t, func(w http.ResponseWriter, r *http.Request) { + if r.URL.Path != "/xrpc/com.atproto.repo.putRecord" { + t.Errorf("path = %q, want /xrpc/com.atproto.repo.putRecord", r.URL.Path) + } + _ = json.NewDecoder(r.Body).Decode(&payload) + w.Header().Set("Content-Type", "application/json") + _ = json.NewEncoder(w).Encode(map[string]any{ + "uri": "at://did:plc:test/social.coves.community.acceptance/rk", + "cid": "bafyreiacceptanceaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa", + "commit": map[string]any{"cid": testCommitCID, "rev": testCommitRev}, + }) + }) + defer closeServer() + + result, err := c.PutRecordWithCommit(context.Background(), + "social.coves.community.acceptance", "rk", + map[string]any{"$type": "social.coves.community.acceptance"}, + "bafyreipreviousaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa") + if err != nil { + t.Fatalf("PutRecordWithCommit: %v", err) + } + if result == nil { + t.Fatal("PutRecordWithCommit returned no result") + } + + if payload["swapRecord"] != "bafyreipreviousaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa" { + t.Errorf("swapRecord = %v, want the CID the caller expected to be replacing", payload["swapRecord"]) + } + if result.URI != "at://did:plc:test/social.coves.community.acceptance/rk" { + t.Errorf("URI = %q", result.URI) + } + if result.CID != "bafyreiacceptanceaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa" { + t.Errorf("CID = %q", result.CID) + } + if result.CommitRev != testCommitRev { + t.Errorf("CommitRev = %q, want %q — without it the admission row has no watermark to stamp", + result.CommitRev, testCommitRev) + } +} + +func TestClient_PutRecordWithCommit_MapsInvalidSwapToSwapConflict(t *testing.T) { + c, closeServer := newCommitClient(t, func(w http.ResponseWriter, _ *http.Request) { + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(http.StatusBadRequest) + _ = json.NewEncoder(w).Encode(map[string]any{ + "error": "InvalidSwap", + "message": "Record was at bafyreiotheraaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa", + }) + }) + defer closeServer() + + _, err := c.PutRecordWithCommit(context.Background(), + "social.coves.community.acceptance", "rk", + map[string]any{"$type": "social.coves.community.acceptance"}, "bafyreistale") + if !errors.Is(err, ErrSwapConflict) { + t.Errorf("error %v does not match ErrSwapConflict", err) + } +} + +func TestClient_CreateRecordWithCommit(t *testing.T) { + c, closeServer := newCommitClient(t, func(w http.ResponseWriter, r *http.Request) { + if r.URL.Path != "/xrpc/com.atproto.repo.createRecord" { + t.Errorf("path = %q, want /xrpc/com.atproto.repo.createRecord", r.URL.Path) + } + w.Header().Set("Content-Type", "application/json") + _ = json.NewEncoder(w).Encode(map[string]any{ + "uri": "at://did:plc:test/social.coves.community.acceptance/rk", + "cid": "bafyreiacceptanceaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa", + "commit": map[string]any{"cid": testCommitCID, "rev": testCommitRev}, + }) + }) + defer closeServer() + + result, err := c.CreateRecordWithCommit(context.Background(), + "social.coves.community.acceptance", "rk", + map[string]any{"$type": "social.coves.community.acceptance"}) + if err != nil { + t.Fatalf("CreateRecordWithCommit: %v", err) + } + if result == nil { + t.Fatal("CreateRecordWithCommit returned no result") + } + if result.CommitRev != testCommitRev { + t.Errorf("CommitRev = %q, want %q", result.CommitRev, testCommitRev) + } +} + +func TestClient_GetLatestCommit(t *testing.T) { + // The swapCommit a batch is guarded by has to come from somewhere, and it + // has to be read immediately before the batch is shaped — the pre-read and + // the commit it is consistent with are the same observation. + c, closeServer := newCommitClient(t, func(w http.ResponseWriter, r *http.Request) { + if r.URL.Path != "/xrpc/com.atproto.repo.getLatestCommit" { + t.Errorf("path = %q, want /xrpc/com.atproto.repo.getLatestCommit", r.URL.Path) + } + if got := r.URL.Query().Get("did"); got != applyWritesDID { + t.Errorf("did = %q, want %q", got, applyWritesDID) + } + w.Header().Set("Content-Type", "application/json") + _ = json.NewEncoder(w).Encode(map[string]any{"cid": testCommitCID, "rev": testCommitRev}) + }) + defer closeServer() + + commit, err := c.GetLatestCommit(context.Background()) + if err != nil { + t.Fatalf("GetLatestCommit: %v", err) + } + if commit == nil { + t.Fatal("GetLatestCommit returned no commit") + } + if commit.CID != testCommitCID || commit.Rev != testCommitRev { + t.Errorf("commit = %+v, want cid %q rev %q", commit, testCommitCID, testCommitRev) + } +} + +func TestWrapAPIError_ServerErrorsAreTheirOwnClass(t *testing.T) { + // applyWrites answers a delete of a missing record — and a create of an + // existing one — with a 500. A state-shaped writer meeting that has to know + // its pre-read went stale and re-shape the batch, which it cannot do if a + // 500 is indistinguishable from a socket that died mid-request. + for _, status := range []int{500, 502, 503, 504} { + err := wrapAPIError(&atclient.APIError{ + StatusCode: status, + Name: "InternalServerError", + Message: "Could not delete record: not found", + }, "applyWrites") + + if !errors.Is(err, ErrServerError) { + t.Errorf("HTTP %d: error %v does not match ErrServerError", status, err) + } + } +} + +func TestWrapAPIError_InvalidSwapIsNotAPlainBadRequest(t *testing.T) { + err := wrapAPIError(&atclient.APIError{ + StatusCode: 400, + Name: "InvalidSwap", + Message: "Record was at bafyreiother", + }, "putRecord") + + if !errors.Is(err, ErrSwapConflict) { + t.Errorf("error %v does not match ErrSwapConflict", err) + } + + // An ordinary 400 must keep meaning what it always meant. The InvalidSwap + // branch is a name test on top of the status, not a replacement for it. + plain := wrapAPIError(&atclient.APIError{ + StatusCode: 400, + Name: "InvalidRequest", + Message: "Bad input", + }, "putRecord") + + if !errors.Is(plain, ErrBadRequest) { + t.Errorf("an ordinary 400 must still be ErrBadRequest, got %v", plain) + } + if errors.Is(plain, ErrSwapConflict) { + t.Error("a malformed request is not a lost race") + } +} diff --git a/internal/atproto/pds/errors.go b/internal/atproto/pds/errors.go index d891736..e5a1bc2 100644 --- a/internal/atproto/pds/errors.go +++ b/internal/atproto/pds/errors.go @@ -27,6 +27,27 @@ var ( // ErrPayloadTooLarge indicates the request payload exceeds PDS limits (HTTP 413). ErrPayloadTooLarge = errors.New("payload too large") + // ErrSwapConflict indicates an optimistic-concurrency guard lost: the + // swapRecord CID or the swapCommit CID the request named is not the one the + // repo is at, so another writer got there first. + // + // It is NOT ErrConflict, and the difference is not cosmetic. A PDS answers + // a failed swap with HTTP 400 and `"error": "InvalidSwap"` — verified + // against a live PDS, not inferred from the lexicon, which documents 409 — + // so the status code alone maps it onto ErrBadRequest, indistinguishable + // from a malformed record. A lost race is the one 400 that must be RETRIED + // (re-read, re-shape, write again) rather than reported, so it needs its + // own sentinel. + ErrSwapConflict = errors.New("swap conflict") + + // ErrServerError indicates the PDS failed to process a well-formed request + // (HTTP 5xx). It is separated from the generic wrap because it is the one + // remote failure class that is worth retrying unchanged: applyWrites + // answers a delete of a missing record, or a create of an existing one, + // with a 500, and a caller that cannot tell that from a transport failure + // cannot decide whether to re-read and re-shape its batch. + ErrServerError = errors.New("server error") + // ErrSessionExpired indicates a stored OAuth session could not be resumed: // the refresh token expired, the session was revoked on the PDS, or the // DPoP key no longer matches. Unlike ErrUnauthorized this is detected diff --git a/internal/core/posts/community_writer.go b/internal/core/posts/community_writer.go new file mode 100644 index 0000000..8f341f9 --- /dev/null +++ b/internal/core/posts/community_writer.go @@ -0,0 +1,194 @@ +package posts + +import ( + "context" + + "Coves/internal/atproto/pds" +) + +// The writers that publish a community's own records about a post. +// +// These are the only things in the post system that write into a COMMUNITY's +// repository (docs/PRD_AUTHOR_OWNED_POSTS.md §5.6), and every one of them is +// STATE-SHAPED rather than optimistic: what it sends is decided by what the +// repo already holds, because the PDS answers the wrong shape with a 500 and +// because a re-fire must not mint a new record CID. + +const ( + // AcceptanceCollection is the community-repo collection holding a + // community's attestation that it accepts a post. + AcceptanceCollection = "social.coves.community.acceptance" + + // RemovalCollection is the community-repo collection holding a community's + // record that a post has been removed from it. + RemovalCollection = "social.coves.community.removal" +) + +// CommunityRepo is one community's PDS repository, narrowed to what the +// writers do with it. +// +// It is declared here rather than taken as pds.CommitClient so that the +// writers' tests can fake four methods instead of a dozen, and so that the +// dependency reads as what it is: a repo we read before we write. +type CommunityRepo interface { + // GetRecord is the PRE-READ. Both writers shape their commit from it: the + // acceptance writer to decide whether there is anything to write at all, + // the removal writer to decide create-vs-update and whether to emit a + // delete. + GetRecord(ctx context.Context, collection, rkey string) (*pds.RecordResponse, error) + + // PutRecordWithCommit writes one record, optionally guarded by the CID it + // expects to be replacing. + PutRecordWithCommit(ctx context.Context, collection, rkey string, record any, swapRecord string) (*pds.RecordCommit, error) + + // ApplyWrites commits several operations together, so the firehose never + // carries a half-completed moderation action. + ApplyWrites(ctx context.Context, writes []pds.Write, swapCommit string) (*pds.ApplyWritesResult, error) + + // GetLatestCommit supplies the swapCommit a batch is guarded by. + GetLatestCommit(ctx context.Context) (*pds.LatestCommit, error) + + // DID is the repo being written — the community's own identity, which is + // the authority half of every record URI this writer produces. + DID() string +} + +// CommunityRepoFactory opens an authenticated client on one community's repo. +// +// The engine works from an admission row, which names a community by DID and +// nothing else, so the credentials have to be fetched per subject rather than +// held. A community this AppView does not host has no credentials and the +// factory says so with an error. +type CommunityRepoFactory func(ctx context.Context, communityDID string) (CommunityRepo, error) + +// CommunityWriteCommand is an acceptance write: this community accepts this +// exact version of this post. +type CommunityWriteCommand struct { + CommunityDID string + PostURI string + + // PostCID is the content CID the acceptance pins. It is required: the + // record's subject is a strongRef, and a strongRef without a CID pins + // nothing, which is exactly the guarantee the acceptance exists to make. + PostCID string +} + +// CommunityRemovalCommand is a removal: this community no longer carries this +// post, for this reason. +type CommunityRemovalCommand struct { + CommunityDID string + PostURI string + + // PostCID is the version present at removal time. It is audit metadata — + // removal is URI-scoped and survives later edits (§5.5). + PostCID string + + // Code is the machine-readable reason, from the open set of §3.3. + Code DecisionCode + + // Reason is the optional human-readable explanation. + Reason string +} + +// CommunityWriteResult describes what reached (or deliberately did not reach) +// the community's repo. +type CommunityWriteResult struct { + // URI, RKey and CID identify the record that now stands — the acceptance + // for an acceptance write, the removal for a removal. + URI string + RKey string + CID string + + // Rev is the repo revision the write committed in: the §5.2 watermark. + // + // It is EMPTY when Skipped is true, and that is not an oversight — nothing + // committed, so there is no revision to report, and getRecord does not + // reveal the revision an existing record was written at. A caller must + // therefore not stamp a watermark from a skipped write; see Skipped. + Rev string + + // Skipped reports that the repo already held exactly this record, so + // nothing was written. + // + // This is the property that makes the three independent acceptance writers + // of §3.2 safe to re-fire: the fast path, the firehose engine and the + // notify endpoint all converge on the same rkey pinning the same CID, and + // a writer that re-put it anyway would mint a fresh record CID, emit a + // pointless commit, and invalidate every reference to the record it just + // rewrote — on every retry, forever. + Skipped bool +} + +// CommunityRecordWriter publishes a community's decisions into its own repo. +// +// Every method is idempotent by construction: the rkey is derived from the +// subject (SubjectRkey), so a concurrent or repeated attempt converges on the +// same record instead of allocating a second one. +type CommunityRecordWriter interface { + // WriteAcceptance makes this community's acceptance of cmd.PostCID stand. + // + // It is one putRecord at the deterministic rkey, guarded by swapRecord, and + // it writes NOTHING when the standing record already pins cmd.PostCID. + // + // On a lost swap it re-reads: if the record another writer committed pins + // the CID we wanted, the work is done and the result is a skip; otherwise + // it retries against the new CID, at most twice, and then defers rather + // than spinning against a livelock. + WriteAcceptance(ctx context.Context, cmd CommunityWriteCommand) (CommunityWriteResult, error) + + // WriteRemoval commits the acceptance's deletion and the removal's write + // TOGETHER, per §3.3. + // + // The batch is shaped by a pre-read of both deterministic rkeys: the delete + // is emitted only when an acceptance is actually there, and the removal is + // a create or an update according to whether one already stands. The PDS + // answers a delete of a missing record — and a create of an existing one — + // with a 500, so the shape is not optional. + WriteRemoval(ctx context.Context, cmd CommunityRemovalCommand) (CommunityWriteResult, error) + + // RestoreAcceptance is WriteRemoval's mirror: one commit deleting the + // removal and writing a fresh acceptance. §5.5 is explicit that there is no + // distinct restore operation on the wire — consumers see ordinary events + // winning the §5.2 tuple CAS — so this is a shape, not a new verb. + RestoreAcceptance(ctx context.Context, cmd CommunityWriteCommand) (CommunityWriteResult, error) + + // RepinAcceptance moves a standing acceptance onto a new content CID + // without re-deciding anything — the bridgedStats exception of §5.5. + // + // It updates the SAME record in place, so the acceptance's createdAt keeps + // meaning "when this community accepted this post" rather than being + // restamped every time a bridge refreshes its vote counts. + RepinAcceptance(ctx context.Context, cmd CommunityWriteCommand) (CommunityWriteResult, error) +} + +// communityRecordWriter is the production writer over real repos. +type communityRecordWriter struct { + repos CommunityRepoFactory + now Clock +} + +// NewCommunityRecordWriter returns the writer that publishes acceptances and +// removals into the repos the factory opens. +// +// The clock is injected for the same reason admitPost's is: createdAt is the +// one field a test cannot otherwise pin, and docs/TEST_ARCHITECTURE.md §3.3 +// forbids sleeping to move time. +func NewCommunityRecordWriter(repos CommunityRepoFactory, now Clock) CommunityRecordWriter { + return &communityRecordWriter{repos: repos, now: now} +} + +func (w *communityRecordWriter) WriteAcceptance(ctx context.Context, cmd CommunityWriteCommand) (CommunityWriteResult, error) { + return CommunityWriteResult{}, nil +} + +func (w *communityRecordWriter) WriteRemoval(ctx context.Context, cmd CommunityRemovalCommand) (CommunityWriteResult, error) { + return CommunityWriteResult{}, nil +} + +func (w *communityRecordWriter) RestoreAcceptance(ctx context.Context, cmd CommunityWriteCommand) (CommunityWriteResult, error) { + return CommunityWriteResult{}, nil +} + +func (w *communityRecordWriter) RepinAcceptance(ctx context.Context, cmd CommunityWriteCommand) (CommunityWriteResult, error) { + return CommunityWriteResult{}, nil +} diff --git a/internal/core/posts/engine.go b/internal/core/posts/engine.go new file mode 100644 index 0000000..a39b7e9 --- /dev/null +++ b/internal/core/posts/engine.go @@ -0,0 +1,131 @@ +package posts + +import "context" + +// The acceptance engine: the single decision point of +// docs/PRD_AUTHOR_OWNED_POSTS.md §5.6. +// +// Its input is one admission row — a (community, post) subject the AppView has +// indexed and this AppView hosts the community for. It runs the admission +// policy over the content the row holds and then makes the community's repo +// agree with the answer: an acceptance record, a removal commit, or an +// AppView-local rejection that writes no record at all. +// +// It is the ONLY writer of community-repo records in the post system, which is +// what makes "every write is idempotent" a property of the system rather than a +// convention each call site has to remember. + +// EngineOutcome reports what one pass over one subject DID. +// +// It is a value rather than an error for the same reason AdmissionOutcome is: +// most of the ways a pass ends without writing anything are the engine working +// — a decision that could not be made yet, a row another writer already +// settled — and routing those into the dead-letter queue would bury the +// failures the queue exists to surface. +type EngineOutcome string + +const ( + // EngineAccepted means a community acceptance now stands, pinning the CID + // the AppView has indexed. + EngineAccepted EngineOutcome = "accepted" + + // EngineRejected means the AppView recorded a local refusal. No community + // record was written: §3.3 is explicit that a submission refused before it + // was ever accepted must not bloat the community's repo. + EngineRejected EngineOutcome = "rejected" + + // EngineRemoved means a removal record now stands and the acceptance is + // gone, committed together. + EngineRemoved EngineOutcome = "removed" + + // EngineRepinned means a standing acceptance moved onto new content with no + // re-decision — the bridgedStats exception of §5.5. + EngineRepinned EngineOutcome = "repinned" + + // EngineDeferred means NOTHING was written anywhere and the subject is + // still owed a decision. It covers an undecided policy answer, a row whose + // content CID is not yet known, a row already in a terminal state, and a + // credential failure. In every one of those cases the correct next step is + // to look again later, never to record a verdict. + EngineDeferred EngineOutcome = "deferred" +) + +// AdmissionDecider is the policy half of admitPost, without its ledger +// reservation. +// +// The split matters. admitPost's dedupe step INSERTS a ledger row, and that +// insert is a submission-time gate: it exists so two concurrent submissions of +// identical content cannot both become posts. The engine is not a submission — +// it is deciding about a post that already exists, often one it has decided +// about before — so reserving a ledger slot here would charge an author quota +// for a firehose redelivery and refuse the redecision as its own duplicate. +type AdmissionDecider interface { + // DecideAdmission evaluates one indexed post against one community's + // policy. A refusal is a code on the decision; a decision that could not be + // made is a non-nil error, and AdmissionDecision.Admitted() reports false + // for both. + DecideAdmission(ctx context.Context, communityDID, postURI string) (AdmissionDecision, error) +} + +// CredentialRefresher forces a community's stored PDS token to be renewed. +// +// The engine holds this so that a stale access token costs one retry rather +// than a verdict. A community's token expires on a schedule that has nothing to +// do with moderation, and an engine that treated the resulting 401 as anything +// other than "try again with fresh credentials" would either lose the decision +// or — far worse — record a refusal for a post it never managed to ask about. +type CredentialRefresher interface { + RefreshCommunityCredentials(ctx context.Context, communityDID string) error +} + +// AcceptanceEngine settles one subject at a time. +type AcceptanceEngine struct { + admissions AdmissionRepository + decider AdmissionDecider + writer CommunityRecordWriter + credentials CredentialRefresher +} + +// NewAcceptanceEngine wires the engine. +func NewAcceptanceEngine( + admissions AdmissionRepository, + decider AdmissionDecider, + writer CommunityRecordWriter, + credentials CredentialRefresher, +) *AcceptanceEngine { + return &AcceptanceEngine{ + admissions: admissions, + decider: decider, + writer: writer, + credentials: credentials, + } +} + +// ProcessAdmission settles one (community, post) subject. +// +// ROUTING IS KEYED ON THE ROW'S STATUS AND THE DECISION TOGETHER, because the +// same verdict means different writes depending on what already stands: +// +// pending + admitted → acceptance created +// pending + refused → AppView-local rejection, NO repo write +// pending_reacceptance + admitted → acceptance updated in place, new CID +// pending_reacceptance + refused → removal commit (§5.5: a failed +// re-acceptance is a removal, never a +// rejection — a rejection would suppress +// the acceptance that is standing) +// anything + undecided → nothing, deferred +// +// A row already in a terminal state — accepted, rejected, removed — is a +// defensive skip. The engine's queue can hand it the same subject twice, and +// re-deciding a settled row is how a removal gets laundered back into a feed. +// +// THE REPOSITORY UPDATE IS OPTIMISTIC AND ITS SKIPS ARE SUCCESS. After a write +// commits, the engine stamps the admission row with the commit rev itself +// rather than waiting for the firehose copy of its own event. If the firehose +// got there first, the repository answers skipped_stale; if the row moved on, +// skipped_terminal. Neither is a failure to report — the firehose is the +// authority on what the repo says, and this write is the AppView catching up +// with itself. +func (e *AcceptanceEngine) ProcessAdmission(ctx context.Context, communityDID, postURI string) (EngineOutcome, error) { + return "", nil +} diff --git a/internal/core/posts/engine_contract_test.go b/internal/core/posts/engine_contract_test.go new file mode 100644 index 0000000..d710c2a --- /dev/null +++ b/internal/core/posts/engine_contract_test.go @@ -0,0 +1,557 @@ +//go:build integration + +package posts_test + +import ( + "context" + "testing" + "time" + + "Coves/internal/atproto/pds" + "Coves/internal/core/posts" + "Coves/internal/db/postgres" + "Coves/tests/testkit" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// The acceptance engine's outer contract: one post's whole life inside one +// community, driven through the engine and read back out of a REAL community +// repository on a REAL PDS (docs/PRD_AUTHOR_OWNED_POSTS.md §5.6). +// +// The routing matrix is proven against fakes in engine_matrix_test.go, at the +// width the matrix has. What is unprovable there — and what this file exists +// for — is everything that only a real repo can contradict: +// +// - that the acceptance lands in the COMMUNITY's repo, at the deterministic +// rkey, pinning the CID the AppView indexed. A writer that computed the key +// differently, or wrote to the wrong repo, passes every fake. +// - that RE-FIRING MINTS NOTHING. This is the assertion with teeth. Three +// independent writers converge on this record (§3.2), and every one of them +// retries; a writer that re-put an identical record would produce a fresh +// record CID on every attempt, invalidating every reference to the +// acceptance it just rewrote. No fake can catch that, because the fake is +// the thing that would have to notice. +// - that the removal is ONE commit. §3.3 requires the acceptance's deletion +// and the removal's write to reach the firehose together, and the PDS is +// the only participant that can refuse a badly shaped batch — it answers a +// delete of a missing record, or a create of an existing one, with a 500. +// - that a lost swapRecord race CONVERGES rather than throws. Forced here by +// committing a competing write between the writer's pre-read and its put, +// which is a real InvalidSwap from a real PDS rather than an injected +// error value. + +const ( + // postv2Collection is the author-repo collection posts live in under + // author-owned posts. Real records are written here so the acceptance pins + // a CID the PDS actually minted, rather than a plausible-looking string. + postv2Collection = "social.coves.community.postv2" +) + +// engineFixture is the acceptance engine over a real community repo, a real +// author repo, and the real admissions table. +type engineFixture struct { + *postFixture + + engine *posts.AcceptanceEngine + writer posts.CommunityRecordWriter + admissions posts.AdmissionRepository + decider *scriptedDecider + refreshes *countingRefresher + communityAt *testkit.Account +} + +// scriptedDecider is the admission policy, scripted per case. The policy itself +// has its own tests (admit_matrix_test.go); what matters here is which repo +// write each verdict produces. +type scriptedDecider struct { + code posts.DecisionCode + err error +} + +func (d *scriptedDecider) DecideAdmission(_ context.Context, _, _ string) (posts.AdmissionDecision, error) { + if d.err != nil { + return posts.AdmissionDecision{Cause: d.err}, d.err + } + return posts.AdmissionDecision{Code: d.code}, nil +} + +// countingRefresher records forced credential renewals. The community's token +// is freshly minted here, so a non-zero count means the writer met a 401 it +// should not have. +type countingRefresher struct{ calls int } + +func (r *countingRefresher) RefreshCommunityCredentials(_ context.Context, _ string) error { + r.calls++ + return nil +} + +// newEngineFixture provisions the community (with its real PDS credentials) and +// points the engine at it. +func newEngineFixture(t *testing.T) *engineFixture { + t.Helper() + + base := newPostFixture(t) + account := base.communityAccount(t) + + generic, err := pds.NewFromAccessToken(base.pds.URL(), account.DID, account.AccessToken) + require.NoError(t, err) + + repo, ok := generic.(pds.CommitClient) + require.Truef(t, ok, "the PDS client must implement pds.CommitClient — the community-repo "+ + "writers need the commit rev and applyWrites, and neither is on the base Client") + + admissions := postgres.NewAdmissionRepository(base.db) + decider := &scriptedDecider{} + refreshes := &countingRefresher{} + + writer := posts.NewCommunityRecordWriter( + func(_ context.Context, communityDID string) (posts.CommunityRepo, error) { + require.Equalf(t, base.community.DID, communityDID, + "the writer asked for credentials to a repo that is not the subject's community") + return repo, nil + }, + func() time.Time { return time.Date(2026, 7, 1, 12, 0, 0, 0, time.UTC) }, + ) + + return &engineFixture{ + postFixture: base, + engine: posts.NewAcceptanceEngine(admissions, decider, writer, refreshes), + writer: writer, + admissions: admissions, + decider: decider, + refreshes: refreshes, + communityAt: account, + } +} + +// publishPost writes a real post record into the AUTHOR's repo and returns it. +// Using a real record means the acceptance pins a CID the PDS minted, so a +// strongRef that fails to round-trip fails here rather than in production. +func (f *engineFixture) publishPost(t *testing.T, title string) testkit.Record { + t.Helper() + return f.author.CreateRecord(t, postv2Collection, map[string]any{ + "$type": postv2Collection, + "community": f.community.DID, + "title": title, + "content": "a body seeking admission", + "createdAt": "2026-07-01T12:00:00Z", + }) +} + +// editPost rewrites the post in place, producing a new content CID. +func (f *engineFixture) editPost(t *testing.T, post testkit.Record, title string) testkit.Record { + t.Helper() + edited := f.author.PutRecord(t, postv2Collection, post.RKey, map[string]any{ + "$type": postv2Collection, + "community": f.community.DID, + "title": title, + "content": "a body the community has not judged", + "createdAt": "2026-07-01T12:00:00Z", + }) + require.NotEqualf(t, post.CID, edited.CID, + "the edit produced the same content CID, so this test would prove nothing about re-acceptance") + return edited +} + +// seedPending records the AppView's observation of the post's content, which is +// what puts a row in the engine's way. +func (f *engineFixture) seedPending(t *testing.T, postURI, contentCID string) *posts.Admission { + t.Helper() + result, err := f.admissions.UpsertPending(context.Background(), posts.UpsertPendingCommand{ + CommunityDID: f.community.DID, + PostURI: postURI, + EvaluatedCID: contentCID, + }) + require.NoError(t, err) + require.NotNil(t, result.Admission) + return result.Admission +} + +func (f *engineFixture) process(t *testing.T, postURI string) (posts.EngineOutcome, error) { + t.Helper() + return f.engine.ProcessAdmission(context.Background(), f.community.DID, postURI) +} + +// acceptanceOf reads the community's acceptance record for a subject. +func (f *engineFixture) acceptanceOf(t *testing.T, postURI string) testkit.RecordValue { + t.Helper() + return f.communityAt.GetRecord(t, posts.AcceptanceCollection, posts.SubjectRkey(postURI)) +} + +// assertSubject checks a record's strongRef points at exactly this version of +// exactly this post. +func assertSubject(t *testing.T, record testkit.RecordValue, wantURI, wantCID string) { + t.Helper() + subject, ok := record.Value["subject"].(map[string]any) + require.Truef(t, ok, "the record has no strongRef subject: %#v", record.Value) + assert.Equal(t, wantURI, subject["uri"]) + assert.Equalf(t, wantCID, subject["cid"], + "the strongRef must pin the exact version that was judged; pinning anything else means "+ + "content nobody evaluated renders under this record") +} + +// assertRecordAbsent asserts a record is not in the community's repo. +func (f *engineFixture) assertRecordAbsent(t *testing.T, collection, rkey, what string) { + t.Helper() + err := getRecordErr(context.Background(), f.communityAt, collection, rkey) + require.Errorf(t, err, "%s is still in the community's repo", what) + assert.Truef(t, testkit.IsNotFound(err), "%s: expected the record to be gone, got: %v", what, err) +} + +// --------------------------------------------------------------------------- + +func TestEngine_AcceptanceLandsInTheCommunityRepoAndSurvivesRefiring(t *testing.T) { + t.Parallel() + + f := newEngineFixture(t) + ctx := context.Background() + post := f.publishPost(t, "a post the community will accept") + f.seedPending(t, post.URI, post.CID) + + outcome, err := f.process(t, post.URI) + require.NoError(t, err) + assert.Equal(t, posts.EngineAccepted, outcome) + + rkey := posts.SubjectRkey(post.URI) + acceptance := f.acceptanceOf(t, post.URI) + + // The authority half of the URI is the COMMUNITY. An acceptance written + // into the author's repo would be an author vouching for themselves. + assert.Equal(t, "at://"+f.community.DID+"/"+posts.AcceptanceCollection+"/"+rkey, acceptance.URI) + assert.Equal(t, posts.AcceptanceCollection, acceptance.Value["$type"]) + assert.NotEmpty(t, acceptance.Value["createdAt"]) + assertSubject(t, acceptance, post.URI, post.CID) + + // The row now agrees: an acceptance pinning the indexed CID is `accepted`. + row, err := f.admissions.Get(ctx, f.community.DID, post.URI) + require.NoError(t, err) + assert.Equal(t, posts.AdmissionStatusAccepted, row.Status) + require.NotNil(t, row.AcceptanceRkey) + assert.Equalf(t, rkey, *row.AcceptanceRkey, + "the rkey the AppView recorded must be the one the record actually lives at, or getStatus "+ + "hands clients a URI that resolves to nothing") + require.NotNil(t, row.LastCommunityEvent) + assert.NotEmptyf(t, row.LastCommunityEvent.Rev, + "the row must carry the rev the acceptance COMMITTED in — that is the §5.2 watermark the "+ + "firehose copy of this same event will be compared against") + + firstRecordCID := acceptance.CID + + // FIRE THE ENGINE AGAIN. A redrive, an overlapping feed, a notify racing + // the firehose — the queue hands the engine the same subject constantly. + outcome, err = f.process(t, post.URI) + require.NoError(t, err) + assert.Equalf(t, posts.EngineDeferred, outcome, + "a settled row is a defensive skip; re-deciding it is how a moderated post gets laundered") + + // FIRE THE WRITER AGAIN, which is where the skip-write rule actually lives: + // the engine's defensive skip above never reaches it, and the fast path and + // the notify endpoint call it with no row check at all. + result, err := f.writer.WriteAcceptance(ctx, posts.CommunityWriteCommand{ + CommunityDID: f.community.DID, + PostURI: post.URI, + PostCID: post.CID, + }) + require.NoError(t, err) + assert.Truef(t, result.Skipped, + "the repo already held this exact acceptance, so the writer must write NOTHING") + assert.Emptyf(t, result.Rev, + "nothing committed, so there is no revision to report — and a caller that stamped one "+ + "would write a watermark no commit ever had") + + // THE ASSERTION WITH TEETH. + assert.Equalf(t, firstRecordCID, f.acceptanceOf(t, post.URI).CID, + "re-firing minted a NEW record CID for an identical acceptance; every reference to the "+ + "acceptance record just became stale, and it will happen again on every retry") + assert.Zerof(t, f.refreshes.calls, "the community's token was fresh; nothing should have forced a renewal") +} + +func TestEngine_FailedReacceptanceRemovesInOneCommit(t *testing.T) { + t.Parallel() + + f := newEngineFixture(t) + ctx := context.Background() + + post := f.publishPost(t, "a post the community will accept, then remove") + f.seedPending(t, post.URI, post.CID) + + outcome, err := f.process(t, post.URI) + require.NoError(t, err) + require.Equal(t, posts.EngineAccepted, outcome) + + // The author edits. The acceptance no longer pins the content the AppView + // holds, so the row moves to pending_reacceptance and §5.5 forbids + // rendering the new content under the old acceptance. + edited := f.editPost(t, post, "a title the community would never have accepted") + row := f.seedPending(t, post.URI, edited.CID) + require.Equalf(t, posts.AdmissionStatusPendingReacceptance, row.Status, + "an edit to an accepted post must leave the row awaiting re-acceptance") + + f.decider.code = posts.DecisionRuleViolation + + outcome, err = f.process(t, post.URI) + require.NoError(t, err) + assert.Equalf(t, posts.EngineRemoved, outcome, + "§5.5: a failed re-acceptance is a REMOVAL — a local rejection would leave the published "+ + "acceptance standing, so federated peers would keep rendering what this AppView hid") + + rkey := posts.SubjectRkey(post.URI) + + // Both halves of the one commit, from the repo's point of view. + f.assertRecordAbsent(t, posts.AcceptanceCollection, rkey, "the acceptance record") + + removal := f.communityAt.GetRecord(t, posts.RemovalCollection, rkey) + assert.Equal(t, "at://"+f.community.DID+"/"+posts.RemovalCollection+"/"+rkey, removal.URI) + assert.Equal(t, posts.RemovalCollection, removal.Value["$type"]) + assert.Equalf(t, string(posts.DecisionRuleViolation), removal.Value["code"], + "the removal's code is what a client renders in #removedPost and what the author is told") + assert.NotEmpty(t, removal.Value["createdAt"]) + assertSubject(t, removal, post.URI, edited.CID) + + // Exactly ONE removal record for the subject, ever. The rkey is derived + // from the subject, so a second removal is an update of this one; a writer + // that allocated a TID instead would leave a trail of them. + assert.Equalf(t, []string{rkey}, listRecordKeys(t, f.communityAt, posts.RemovalCollection), + "the community's removal collection must hold exactly one record, at the deterministic rkey") + + // And the row followed the repo. + stored, err := f.admissions.Get(ctx, f.community.DID, post.URI) + require.NoError(t, err) + assert.Equal(t, posts.AdmissionStatusRemoved, stored.Status) + require.NotNil(t, stored.DecisionCode) + assert.Equal(t, string(posts.DecisionRuleViolation), *stored.DecisionCode) + assert.Nilf(t, stored.AcceptanceURI, + "the acceptance is gone from the repo, so the row must not still advertise it") +} + +func TestEngine_RestoreIsTheRemovalCommitRunBackwards(t *testing.T) { + t.Parallel() + + // §5.5: `removed` is exited only by an explicit restore — one commit + // deleting the removal and writing a fresh acceptance. There is no distinct + // restore operation on the wire; it is the same shape as the removal with + // the two halves swapped, which is exactly why it is worth proving the + // symmetry holds against a real repo. + f := newEngineFixture(t) + ctx := context.Background() + + post := f.publishPost(t, "a post that will be removed and restored") + f.seedPending(t, post.URI, post.CID) + require.NoError(t, firstErr(f.process(t, post.URI))) + + f.decider.code = posts.DecisionModeratorDiscretion + edited := f.editPost(t, post, "an edit that fails re-acceptance") + f.seedPending(t, post.URI, edited.CID) + outcome, err := f.process(t, post.URI) + require.NoError(t, err) + require.Equal(t, posts.EngineRemoved, outcome) + + rkey := posts.SubjectRkey(post.URI) + + result, err := f.writer.RestoreAcceptance(ctx, posts.CommunityWriteCommand{ + CommunityDID: f.community.DID, + PostURI: post.URI, + PostCID: edited.CID, + }) + require.NoError(t, err) + assert.Falsef(t, result.Skipped, "there was no acceptance standing; the restore had work to do") + assert.NotEmptyf(t, result.Rev, "a commit happened, so it has a revision") + + f.assertRecordAbsent(t, posts.RemovalCollection, rkey, "the removal record") + + restored := f.acceptanceOf(t, post.URI) + assertSubject(t, restored, post.URI, edited.CID) + assert.Equal(t, "at://"+f.community.DID+"/"+posts.AcceptanceCollection+"/"+rkey, restored.URI, + "the restored acceptance must reuse the subject's rkey, not allocate a new one") + assert.Equalf(t, []string{rkey}, listRecordKeys(t, f.communityAt, posts.AcceptanceCollection), + "exactly one acceptance record for the subject, at the deterministic rkey") +} + +func TestEngine_ApplyAcceptanceTwiceAtTheSameRevIsASkipThatChangesNothing(t *testing.T) { + t.Parallel() + + // The engine stamps the row optimistically and the firehose delivers the + // same event moments later. Both write the SAME rev, so the second must be + // a skip that leaves the row byte-identical — not a re-stamped decision + // timestamp, and certainly not an error into the dead-letter queue. + f := newEngineFixture(t) + ctx := context.Background() + + post := f.publishPost(t, "a post whose acceptance arrives twice") + f.seedPending(t, post.URI, post.CID) + outcome, err := f.process(t, post.URI) + require.NoError(t, err) + require.Equal(t, posts.EngineAccepted, outcome) + + before, err := f.admissions.Get(ctx, f.community.DID, post.URI) + require.NoError(t, err) + require.NotNil(t, before.LastCommunityEvent) + + replay := posts.ApplyAcceptanceCommand{ + CommunityDID: f.community.DID, + PostURI: post.URI, + AcceptanceURI: *before.AcceptanceURI, + AcceptanceRkey: *before.AcceptanceRkey, + PinnedCID: post.CID, + Watermark: posts.CommunityWatermark{Rev: before.LastCommunityEvent.Rev, OpRank: posts.CommunityOpPut}, + } + + result, err := f.admissions.ApplyAcceptance(ctx, replay) + require.NoErrorf(t, err, "a replayed event is the ordering gate WORKING and must not be an error") + assert.Equalf(t, posts.AdmissionSkippedStale, result.Outcome, + "the watermark must be STRICTLY greater; an equal rev is the same event arriving again") + + after, err := f.admissions.Get(ctx, f.community.DID, post.URI) + require.NoError(t, err) + assert.Equalf(t, before, after, + "a skipped event must leave the row untouched, down to updated_at — a re-stamped row is "+ + "how a replay looks like a fresh decision in the moderation log") +} + +// --------------------------------------------------------------------------- +// Swap conflicts +// --------------------------------------------------------------------------- + +// racingRepo commits a competing write between the writer's pre-read and its +// put, so the swapRecord CID the writer computed is already stale when it +// arrives. The conflict is therefore a real InvalidSwap from a real PDS, not an +// error value a fake handed back. +type racingRepo struct { + posts.CommunityRepo + + race func() + done bool +} + +func (r *racingRepo) PutRecordWithCommit(ctx context.Context, collection, rkey string, record any, swapRecord string) (*pds.RecordCommit, error) { + if !r.done { + r.done = true + r.race() + } + return r.CommunityRepo.PutRecordWithCommit(ctx, collection, rkey, record, swapRecord) +} + +// racingWriter returns a writer whose first put loses a swap race to fn. +func (f *engineFixture) racingWriter(t *testing.T, fn func()) posts.CommunityRecordWriter { + t.Helper() + + generic, err := pds.NewFromAccessToken(f.pds.URL(), f.communityAt.DID, f.communityAt.AccessToken) + require.NoError(t, err) + repo, ok := generic.(pds.CommitClient) + require.True(t, ok) + + return posts.NewCommunityRecordWriter( + func(_ context.Context, _ string) (posts.CommunityRepo, error) { + return &racingRepo{CommunityRepo: repo, race: fn}, nil + }, + func() time.Time { return time.Date(2026, 7, 1, 12, 0, 0, 0, time.UTC) }, + ) +} + +// writeAcceptanceDirectly commits an acceptance through a SECOND client on the +// same repo — another AppView instance, or the synchronous fast path racing the +// firehose engine. +func (f *engineFixture) writeAcceptanceDirectly(t *testing.T, postURI, pinnedCID string) { + t.Helper() + f.communityAt.PutRecord(t, posts.AcceptanceCollection, posts.SubjectRkey(postURI), map[string]any{ + "$type": posts.AcceptanceCollection, + "subject": map[string]any{"uri": postURI, "cid": pinnedCID}, + "createdAt": "2026-07-01T11:59:00Z", + }) +} + +func TestEngine_LostSwapRaceToTheSameCIDIsAlreadyDone(t *testing.T) { + t.Parallel() + + // The common race: two writers admit the same post at the same moment and + // aim at the same CID. The loser re-reads, finds the record it wanted + // already standing, and STOPS. Retrying would mint a new record CID for an + // acceptance that is already correct — the exact churn deterministic rkeys + // exist to prevent. + f := newEngineFixture(t) + post := f.publishPost(t, "a post two writers accept at once") + + writer := f.racingWriter(t, func() { + f.writeAcceptanceDirectly(t, post.URI, post.CID) + }) + + result, err := writer.WriteAcceptance(context.Background(), posts.CommunityWriteCommand{ + CommunityDID: f.community.DID, + PostURI: post.URI, + PostCID: post.CID, + }) + require.NoErrorf(t, err, "a lost race to the same outcome is convergence, not a failure") + assert.Truef(t, result.Skipped, + "the winner pinned the CID we wanted, so there is nothing left to write") + + acceptance := f.acceptanceOf(t, post.URI) + assertSubject(t, acceptance, post.URI, post.CID) +} + +func TestEngine_LostSwapRaceToADifferentCIDRetriesAndConverges(t *testing.T) { + t.Parallel() + + // The harder race: the winner pinned a DIFFERENT version — an acceptance of + // the pre-edit content landing while this writer is accepting the edit. The + // loser must re-read and try again against what is actually there, and the + // repo must end up holding OUR target. + f := newEngineFixture(t) + post := f.publishPost(t, "a post accepted at two different versions") + edited := f.editPost(t, post, "the version this writer is accepting") + + writer := f.racingWriter(t, func() { + f.writeAcceptanceDirectly(t, post.URI, post.CID) + }) + + result, err := writer.WriteAcceptance(context.Background(), posts.CommunityWriteCommand{ + CommunityDID: f.community.DID, + PostURI: post.URI, + PostCID: edited.CID, + }) + require.NoErrorf(t, err, "a bounded retry after a lost swap must converge, not surface the conflict") + assert.Falsef(t, result.Skipped, "the standing record pinned the wrong CID, so there WAS work to do") + assert.NotEmpty(t, result.Rev) + + acceptance := f.acceptanceOf(t, post.URI) + assertSubject(t, acceptance, post.URI, edited.CID) + assert.Equalf(t, []string{posts.SubjectRkey(post.URI)}, + listRecordKeys(t, f.communityAt, posts.AcceptanceCollection), + "the race must leave ONE acceptance record, not one per contender") +} + +// --------------------------------------------------------------------------- +// Helpers +// --------------------------------------------------------------------------- + +// listRecordKeys returns the rkeys in a collection of the account's repo, so a +// test can assert that a collection holds exactly the records it should. +func listRecordKeys(t *testing.T, account *testkit.Account, collection string) []string { + t.Helper() + + var resp struct { + Records []struct { + URI string `json:"uri"` + } `json:"records"` + } + ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second) + defer cancel() + + err := account.XRPC().Query(ctx, "com.atproto.repo.listRecords", map[string][]string{ + "repo": {account.DID}, + "collection": {collection}, + "limit": {"100"}, + }, &resp) + require.NoErrorf(t, err, "listing %s in %s", collection, account.DID) + + keys := make([]string, 0, len(resp.Records)) + for _, record := range resp.Records { + keys = append(keys, rkeyOf(t, record.URI)) + } + return keys +} + +// firstErr drops an outcome and keeps the error, for the setup steps whose +// outcome a test has already proven elsewhere. +func firstErr(_ posts.EngineOutcome, err error) error { return err } diff --git a/internal/core/posts/engine_matrix_test.go b/internal/core/posts/engine_matrix_test.go new file mode 100644 index 0000000..df21d46 --- /dev/null +++ b/internal/core/posts/engine_matrix_test.go @@ -0,0 +1,718 @@ +package posts + +import ( + "context" + "errors" + "fmt" + "testing" + "time" + + "Coves/internal/atproto/pds" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// The acceptance engine's routing matrix (docs/PRD_AUTHOR_OWNED_POSTS.md §5.6). +// +// ProcessAdmission is a two-input decision: the row's STATUS and the policy's +// VERDICT together choose what gets written, and neither alone is enough. The +// same "refused" answer means an AppView-local rejection on a pending row and a +// removal commit on a pending_reacceptance row, because §5.5 is explicit that a +// failed re-acceptance is a removal — a rejection there would suppress an +// acceptance that is currently standing and that the community published. +// +// WHAT THESE TESTS ASSERT IS THE CALL SEQUENCE, NOT JUST THE OUTCOME. Every +// interesting failure of an engine like this is a write that happened when it +// should not have: a rejection recorded for a post whose credentials expired, a +// removal committed for a row that was already removed, a repo stamped before +// the PDS write it claims to describe actually committed. An outcome value +// cannot see any of those. The recorded sequence can, so the fakes below share +// one recorder and the assertions name the exact calls in the exact order. +// +// The outer contract — that this really writes records into a real repo — is +// engine_contract_test.go against a live PDS. + +const ( + engineCommunityDID = "did:plc:cccccccccccccccccccccccc" + enginePostURI = "at://did:plc:aaaaaaaaaaaaaaaaaaaaaaaa/social.coves.community.postv2/3kjzl5kcb2s2v" + + engineIndexedCID = "bafyreiaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa" + engineAcceptedCID = "bafyreibbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb" + + engineCommitRev = "3kjzl5kcb2s2v" + engineRecordCID = "bafyreiddddddddddddddddddddddddddddddddddddddddddddddddd" + engineRemovalCID = "bafyreieeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeee" +) + +// --------------------------------------------------------------------------- +// Fakes +// --------------------------------------------------------------------------- + +// engineRecorder is the shared call log. One log across all four collaborators +// is what makes ORDER assertable — "the row was stamped after the commit +// landed" is a claim about two different fakes. +type engineRecorder struct{ calls []string } + +func (r *engineRecorder) record(name string) { r.calls = append(r.calls, name) } + +// mutatingCalls are every call that changes state somewhere. A deferred pass +// must make none of them. +var mutatingCalls = []string{ + "WriteAcceptance", "WriteRemoval", "RestoreAcceptance", "RepinAcceptance", + "ApplyAcceptance", "ApplyRemoval", "RecordRejection", "UpsertPending", + "ApplyAcceptanceDelete", "ApplyRemovalDelete", "RepinAcceptedCID", +} + +// assertWroteNothing fails if the pass touched anything at all. +func assertWroteNothing(t *testing.T, rec *engineRecorder) { + t.Helper() + for _, call := range rec.calls { + for _, mutation := range mutatingCalls { + assert.NotEqualf(t, mutation, call, + "this pass must write nothing anywhere; it called %s (whole sequence: %v)", + call, rec.calls) + } + } +} + +// fakeDecider is the admission policy. It answers with whatever it was handed: +// an admission, a refusal code, or the undecided pair admitPost returns when a +// lookup failed. +type fakeDecider struct { + rec *engineRecorder + + code DecisionCode + err error + + lastCommunityDID string + lastPostURI string +} + +func (d *fakeDecider) DecideAdmission(_ context.Context, communityDID, postURI string) (AdmissionDecision, error) { + d.rec.record("DecideAdmission") + d.lastCommunityDID = communityDID + d.lastPostURI = postURI + if d.err != nil { + // Exactly admitPost's undecided shape: the error is returned AND + // carried on the decision, so a caller inspecting only the value still + // sees "not admitted". + return AdmissionDecision{Cause: d.err}, d.err + } + return AdmissionDecision{Code: d.code}, nil +} + +// fakeWriter is the community repo. Its results are canned; what matters is +// which method the engine reached for and with what. +type fakeWriter struct { + rec *engineRecorder + + acceptanceResult CommunityWriteResult + removalResult CommunityWriteResult + + // acceptanceErrs is consumed one per WriteAcceptance call, so a test can + // fail the first attempt and let the retry through. + acceptanceErrs []error + removalErr error + + acceptanceCmds []CommunityWriteCommand + removalCmds []CommunityRemovalCommand +} + +func (w *fakeWriter) WriteAcceptance(_ context.Context, cmd CommunityWriteCommand) (CommunityWriteResult, error) { + w.rec.record("WriteAcceptance") + w.acceptanceCmds = append(w.acceptanceCmds, cmd) + if len(w.acceptanceErrs) > 0 { + err := w.acceptanceErrs[0] + w.acceptanceErrs = w.acceptanceErrs[1:] + if err != nil { + return CommunityWriteResult{}, err + } + } + return w.acceptanceResult, nil +} + +func (w *fakeWriter) WriteRemoval(_ context.Context, cmd CommunityRemovalCommand) (CommunityWriteResult, error) { + w.rec.record("WriteRemoval") + w.removalCmds = append(w.removalCmds, cmd) + if w.removalErr != nil { + return CommunityWriteResult{}, w.removalErr + } + return w.removalResult, nil +} + +func (w *fakeWriter) RestoreAcceptance(_ context.Context, _ CommunityWriteCommand) (CommunityWriteResult, error) { + w.rec.record("RestoreAcceptance") + return CommunityWriteResult{}, nil +} + +func (w *fakeWriter) RepinAcceptance(_ context.Context, _ CommunityWriteCommand) (CommunityWriteResult, error) { + w.rec.record("RepinAcceptance") + return CommunityWriteResult{}, nil +} + +// fakeRefresher counts forced credential renewals. +type fakeRefresher struct { + rec *engineRecorder + err error +} + +func (r *fakeRefresher) RefreshCommunityCredentials(_ context.Context, _ string) error { + r.rec.record("RefreshCommunityCredentials") + return r.err +} + +// fakeAdmissions is the admission row store. It answers Get with one row and +// records every mutation with its command, so the assertions can check what the +// engine believed it was recording. +type fakeAdmissions struct { + rec *engineRecorder + + row *Admission + getErr error + + acceptanceResult AdmissionResult + removalResult AdmissionResult + rejectionResult AdmissionResult + + acceptanceErr error + removalErr error + rejectionErr error + + acceptanceCmds []ApplyAcceptanceCommand + removalCmds []ApplyRemovalCommand + rejectionCmds []RecordRejectionCommand +} + +func (a *fakeAdmissions) Get(_ context.Context, _, _ string) (*Admission, error) { + a.rec.record("Get") + if a.getErr != nil { + return nil, a.getErr + } + return a.row, nil +} + +func (a *fakeAdmissions) ApplyAcceptance(_ context.Context, cmd ApplyAcceptanceCommand) (AdmissionResult, error) { + a.rec.record("ApplyAcceptance") + a.acceptanceCmds = append(a.acceptanceCmds, cmd) + return a.acceptanceResult, a.acceptanceErr +} + +func (a *fakeAdmissions) ApplyRemoval(_ context.Context, cmd ApplyRemovalCommand) (AdmissionResult, error) { + a.rec.record("ApplyRemoval") + a.removalCmds = append(a.removalCmds, cmd) + return a.removalResult, a.removalErr +} + +func (a *fakeAdmissions) RecordRejection(_ context.Context, cmd RecordRejectionCommand) (AdmissionResult, error) { + a.rec.record("RecordRejection") + a.rejectionCmds = append(a.rejectionCmds, cmd) + return a.rejectionResult, a.rejectionErr +} + +func (a *fakeAdmissions) UpsertPending(_ context.Context, _ UpsertPendingCommand) (AdmissionResult, error) { + a.rec.record("UpsertPending") + return AdmissionResult{}, nil +} + +func (a *fakeAdmissions) ApplyAcceptanceDelete(_ context.Context, _ CommunityDeleteCommand) (AdmissionResult, error) { + a.rec.record("ApplyAcceptanceDelete") + return AdmissionResult{}, nil +} + +func (a *fakeAdmissions) ApplyRemovalDelete(_ context.Context, _ CommunityDeleteCommand) (AdmissionResult, error) { + a.rec.record("ApplyRemovalDelete") + return AdmissionResult{}, nil +} + +func (a *fakeAdmissions) RepinAcceptedCID(_ context.Context, _ RepinAcceptanceCommand) (AdmissionResult, error) { + a.rec.record("RepinAcceptedCID") + return AdmissionResult{}, nil +} + +func (a *fakeAdmissions) GetByPostURIs(_ context.Context, _ []string) (map[string][]*Admission, error) { + a.rec.record("GetByPostURIs") + return nil, nil +} + +func (a *fakeAdmissions) ListByStatusForCommunity(_ context.Context, _ string, _ AdmissionStatus, _ int, _ *string) ([]*Admission, *string, error) { + a.rec.record("ListByStatusForCommunity") + return nil, nil, nil +} + +// --------------------------------------------------------------------------- +// Harness +// --------------------------------------------------------------------------- + +// engineHarness is the engine plus every collaborator, all sharing one call log. +type engineHarness struct { + engine *AcceptanceEngine + rec *engineRecorder + admissions *fakeAdmissions + decider *fakeDecider + writer *fakeWriter + refresher *fakeRefresher +} + +// newEngineHarness builds an engine over a row in the given status holding the +// given indexed CID. A nil evaluatedCID means the AppView has not yet decoded +// the post's content. +func newEngineHarness(status AdmissionStatus, evaluatedCID *string) *engineHarness { + rec := &engineRecorder{} + + row := &Admission{ + CommunityDID: engineCommunityDID, + PostURI: enginePostURI, + Status: status, + EvaluatedCID: evaluatedCID, + Redrivable: true, + CreatedAt: time.Date(2026, 7, 1, 12, 0, 0, 0, time.UTC), + UpdatedAt: time.Date(2026, 7, 1, 12, 0, 0, 0, time.UTC), + } + if status == AdmissionStatusPendingReacceptance || status == AdmissionStatusAccepted { + acceptanceRkey := SubjectRkey(enginePostURI) + acceptanceURI := "at://" + engineCommunityDID + "/" + AcceptanceCollection + "/" + acceptanceRkey + acceptedCID := engineAcceptedCID + row.AcceptanceURI = &acceptanceURI + row.AcceptanceRkey = &acceptanceRkey + row.AcceptedCID = &acceptedCID + } + + admissions := &fakeAdmissions{ + rec: rec, + row: row, + acceptanceResult: AdmissionResult{Outcome: AdmissionApplied, Admission: row}, + removalResult: AdmissionResult{Outcome: AdmissionApplied, Admission: row}, + rejectionResult: AdmissionResult{Outcome: AdmissionApplied, Admission: row}, + } + + acceptanceRkey := SubjectRkey(enginePostURI) + writer := &fakeWriter{ + rec: rec, + acceptanceResult: CommunityWriteResult{ + URI: "at://" + engineCommunityDID + "/" + AcceptanceCollection + "/" + acceptanceRkey, + RKey: acceptanceRkey, + CID: engineRecordCID, + Rev: engineCommitRev, + }, + removalResult: CommunityWriteResult{ + URI: "at://" + engineCommunityDID + "/" + RemovalCollection + "/" + acceptanceRkey, + RKey: acceptanceRkey, + CID: engineRemovalCID, + Rev: engineCommitRev, + }, + } + + decider := &fakeDecider{rec: rec} + refresher := &fakeRefresher{rec: rec} + + return &engineHarness{ + engine: NewAcceptanceEngine(admissions, decider, writer, refresher), + rec: rec, + admissions: admissions, + decider: decider, + writer: writer, + refresher: refresher, + } +} + +func (h *engineHarness) process(t *testing.T) (EngineOutcome, error) { + t.Helper() + return h.engine.ProcessAdmission(context.Background(), engineCommunityDID, enginePostURI) +} + +func cidPtr(cid string) *string { return &cid } + +// --------------------------------------------------------------------------- +// The matrix +// --------------------------------------------------------------------------- + +func TestEngine_PendingAdmittedCreatesTheAcceptance(t *testing.T) { + t.Parallel() + + h := newEngineHarness(AdmissionStatusPending, cidPtr(engineIndexedCID)) + + outcome, err := h.process(t) + require.NoError(t, err) + assert.Equal(t, EngineAccepted, outcome) + + // The whole sequence, in order. Reading the row FIRST is what makes the + // status half of the routing decision come from the database rather than + // from whoever queued the subject; stamping the row LAST is what keeps the + // AppView from claiming an acceptance the PDS never committed. + assert.Equal(t, []string{"Get", "DecideAdmission", "WriteAcceptance", "ApplyAcceptance"}, h.rec.calls) + + require.Len(t, h.writer.acceptanceCmds, 1) + assert.Equal(t, CommunityWriteCommand{ + CommunityDID: engineCommunityDID, + PostURI: enginePostURI, + PostCID: engineIndexedCID, + }, h.writer.acceptanceCmds[0], + "the acceptance must pin the CID the AppView has INDEXED; pinning anything else is an "+ + "acceptance of content nobody evaluated") + + require.Len(t, h.admissions.acceptanceCmds, 1) + assert.Equal(t, ApplyAcceptanceCommand{ + CommunityDID: engineCommunityDID, + PostURI: enginePostURI, + AcceptanceURI: "at://" + engineCommunityDID + "/" + AcceptanceCollection + "/" + SubjectRkey(enginePostURI), + AcceptanceRkey: SubjectRkey(enginePostURI), + PinnedCID: engineIndexedCID, + Watermark: CommunityWatermark{Rev: engineCommitRev}, + }, h.admissions.acceptanceCmds[0], + "the row is stamped with the rev the write actually COMMITTED in — that is the §5.2 "+ + "watermark, and inventing one would let this optimistic update outrank the firehose "+ + "copy of a later event") +} + +func TestEngine_PendingRefusedRecordsALocalRejectionAndWritesNoRecord(t *testing.T) { + t.Parallel() + + h := newEngineHarness(AdmissionStatusPending, cidPtr(engineIndexedCID)) + h.decider.code = DecisionSpam + + outcome, err := h.process(t) + require.NoError(t, err) + assert.Equal(t, EngineRejected, outcome) + + // NO PDS WRITE. §3.3 is explicit: a submission refused before it was ever + // accepted writes no record, because spam must not bloat the community's + // repository. This is the assertion that keeps that promise. + assert.Equal(t, []string{"Get", "DecideAdmission", "RecordRejection"}, h.rec.calls) + assert.Emptyf(t, h.writer.acceptanceCmds, "a refused submission must not reach the community's repo") + assert.Empty(t, h.writer.removalCmds) + + require.Len(t, h.admissions.rejectionCmds, 1) + assert.Equal(t, RecordRejectionCommand{ + CommunityDID: engineCommunityDID, + PostURI: enginePostURI, + DecisionCode: string(DecisionSpam), + // The CID the verdict JUDGED. The repository lands the rejection only + // on a pending row still holding it, so an author who edited between + // the read and the write gets fresh content judged fresh rather than + // condemned by a verdict about something else. + JudgedCID: engineIndexedCID, + // A policy refusal is terminal. Leaving this true would have the + // dead-letter redrive pass retry a decision that will never change. + Redrivable: false, + }, h.admissions.rejectionCmds[0]) +} + +func TestEngine_PendingReacceptanceAdmittedUpdatesTheAcceptanceInPlace(t *testing.T) { + t.Parallel() + + h := newEngineHarness(AdmissionStatusPendingReacceptance, cidPtr(engineIndexedCID)) + + outcome, err := h.process(t) + require.NoError(t, err) + assert.Equal(t, EngineAccepted, outcome) + + assert.Equal(t, []string{"Get", "DecideAdmission", "WriteAcceptance", "ApplyAcceptance"}, h.rec.calls) + + // SAME RKEY, NEW CID. The record key is derived from the subject, so + // re-acceptance is an update of the record that already exists rather than + // a second acceptance — every reference to the acceptance URI stays valid + // and the community's repo does not accumulate one record per edit. + require.Len(t, h.writer.acceptanceCmds, 1) + assert.Equal(t, engineIndexedCID, h.writer.acceptanceCmds[0].PostCID, + "re-acceptance pins the NEW content; pinning the old CID would leave the post pending forever") + + require.Len(t, h.admissions.acceptanceCmds, 1) + assert.Equal(t, SubjectRkey(enginePostURI), h.admissions.acceptanceCmds[0].AcceptanceRkey) + assert.Equal(t, engineIndexedCID, h.admissions.acceptanceCmds[0].PinnedCID) +} + +func TestEngine_PendingReacceptanceRefusedRemovesRatherThanRejects(t *testing.T) { + t.Parallel() + + h := newEngineHarness(AdmissionStatusPendingReacceptance, cidPtr(engineIndexedCID)) + h.decider.code = DecisionRuleViolation + + outcome, err := h.process(t) + require.NoError(t, err) + assert.Equal(t, EngineRemoved, outcome) + + // §5.5: a failed RE-acceptance is a removal, not a rejection. An acceptance + // is currently standing in the community's repo and was published to the + // firehose; recording an AppView-local rejection would leave that record in + // place, so peers would keep rendering the post while this AppView hid it. + assert.Equal(t, []string{"Get", "DecideAdmission", "WriteRemoval", "ApplyRemoval"}, h.rec.calls) + assert.Emptyf(t, h.admissions.rejectionCmds, + "a standing acceptance must be withdrawn with a removal record, never with a local rejection") + + require.Len(t, h.writer.removalCmds, 1) + assert.Equal(t, CommunityRemovalCommand{ + CommunityDID: engineCommunityDID, + PostURI: enginePostURI, + PostCID: engineIndexedCID, + Code: DecisionRuleViolation, + }, h.writer.removalCmds[0]) + + require.Len(t, h.admissions.removalCmds, 1) + assert.Equal(t, ApplyRemovalCommand{ + CommunityDID: engineCommunityDID, + PostURI: enginePostURI, + DecisionCode: string(DecisionRuleViolation), + Watermark: CommunityWatermark{Rev: engineCommitRev}, + }, h.admissions.removalCmds[0]) +} + +func TestEngine_UndecidedWritesNothingAnywhere(t *testing.T) { + t.Parallel() + + // An undecided answer is an infrastructure failure — a community lookup or + // a ban lookup that could not be reached — dressed as neither an admission + // nor a refusal. Recording ANY verdict from it would turn a Postgres blip + // into a permanent decision about someone's post. + for _, status := range []AdmissionStatus{ + AdmissionStatusPending, + AdmissionStatusPendingReacceptance, + } { + t.Run(string(status), func(t *testing.T) { + t.Parallel() + + lookupFailed := errors.New("failed to look up community membership: connection refused") + h := newEngineHarness(status, cidPtr(engineIndexedCID)) + h.decider.err = lookupFailed + + outcome, err := h.process(t) + assert.Equal(t, EngineDeferred, outcome) + require.Error(t, err, "a decision that could not be made is a genuine failure and must be visible") + assert.ErrorIs(t, err, lookupFailed) + + assertWroteNothing(t, h.rec) + }) + } +} + +func TestEngine_AdmittedWithNoIndexedCIDDefers(t *testing.T) { + t.Parallel() + + // An acceptance's subject is a strongRef, and a strongRef without a CID + // pins nothing — which is precisely the guarantee the acceptance exists to + // make. A row with no evaluated_cid is one whose content the AppView has + // not decoded yet, so there is nothing to accept and the answer is "later". + for _, status := range []AdmissionStatus{ + AdmissionStatusPending, + AdmissionStatusPendingReacceptance, + } { + t.Run(string(status), func(t *testing.T) { + t.Parallel() + + h := newEngineHarness(status, nil) + + outcome, err := h.process(t) + require.NoError(t, err, "an un-indexed subject is the engine working, not a failure") + assert.Equal(t, EngineDeferred, outcome) + + assertWroteNothing(t, h.rec) + }) + } +} + +func TestEngine_SettledRowsAreSkippedWithoutBeingReDecided(t *testing.T) { + t.Parallel() + + // The engine's queue can hand it the same subject twice — a redrive, a + // duplicate from an overlapping feed, a notify racing the firehose. A row + // that is already settled must not be re-decided, and the strongest reason + // is `removed`: §5.5 makes removal terminal against everything except a + // moderator restore at a strictly greater watermark, so an engine that + // re-ran policy on a removed row would launder it straight back into the + // feeds it was removed from. + for _, status := range []AdmissionStatus{ + AdmissionStatusAccepted, + AdmissionStatusRejected, + AdmissionStatusRemoved, + } { + t.Run(string(status), func(t *testing.T) { + t.Parallel() + + h := newEngineHarness(status, cidPtr(engineIndexedCID)) + + outcome, err := h.process(t) + require.NoError(t, err) + assert.Equal(t, EngineDeferred, outcome) + + assert.Equal(t, []string{"Get"}, h.rec.calls, + "a settled row costs one read and nothing else — the policy must not even be consulted") + }) + } +} + +func TestEngine_MissingRowDefersWithoutDeciding(t *testing.T) { + t.Parallel() + + // A subject the engine was handed but that has no row. Whether that is + // reported as an error is the engine's business; what is NOT negotiable is + // that it decides nothing, because there is no evaluated CID to judge and + // no status to route on. + h := newEngineHarness(AdmissionStatusPending, cidPtr(engineIndexedCID)) + h.admissions.row = nil + h.admissions.getErr = ErrNotFound + + outcome, _ := h.process(t) + assert.Equal(t, EngineDeferred, outcome) + assertWroteNothing(t, h.rec) +} + +// --------------------------------------------------------------------------- +// Credentials +// --------------------------------------------------------------------------- + +func TestEngine_RetriesOnceWithFreshCredentials(t *testing.T) { + t.Parallel() + + // A community's PDS access token expires on a schedule that has nothing to + // do with moderation. One forced renewal and one retry is the difference + // between a decision that lands and a decision that has to be redriven. + h := newEngineHarness(AdmissionStatusPending, cidPtr(engineIndexedCID)) + h.writer.acceptanceErrs = []error{fmt.Errorf("putRecord: %w: token expired", pds.ErrUnauthorized)} + + outcome, err := h.process(t) + require.NoError(t, err) + assert.Equal(t, EngineAccepted, outcome) + + assert.Equal(t, []string{ + "Get", "DecideAdmission", + "WriteAcceptance", "RefreshCommunityCredentials", "WriteAcceptance", + "ApplyAcceptance", + }, h.rec.calls) +} + +func TestEngine_CredentialFailureDefersAndNeverRejects(t *testing.T) { + t.Parallel() + + // THE MOST DANGEROUS CONFUSION IN THE WHOLE COMPONENT. "I could not write + // the acceptance" and "this post is not acceptable" are opposite facts, and + // an engine that recorded the second when it meant the first would answer + // the author with a permanent verdict — and set redrivable=false on it, so + // nothing would ever retry. + h := newEngineHarness(AdmissionStatusPending, cidPtr(engineIndexedCID)) + expired := fmt.Errorf("putRecord: %w: token expired", pds.ErrUnauthorized) + h.writer.acceptanceErrs = []error{expired, expired} + + outcome, err := h.process(t) + assert.Equal(t, EngineDeferred, outcome) + require.Error(t, err, "credentials that stay dead after a forced renewal are an operator problem and must surface") + assert.ErrorIs(t, err, pds.ErrUnauthorized) + + assert.Equal(t, []string{ + "Get", "DecideAdmission", + "WriteAcceptance", "RefreshCommunityCredentials", "WriteAcceptance", + }, h.rec.calls) + assert.Emptyf(t, h.admissions.rejectionCmds, + "an authentication failure must NEVER be recorded as a rejection of the post") + assert.Empty(t, h.admissions.acceptanceCmds, + "nothing committed, so nothing may be stamped on the row") +} + +func TestEngine_RetriesTheCredentialFailureAtMostOnce(t *testing.T) { + t.Parallel() + + // Bounded, not persistent. A community whose credentials are genuinely gone + // would otherwise have every pass in the queue spin against its PDS. + h := newEngineHarness(AdmissionStatusPendingReacceptance, cidPtr(engineIndexedCID)) + h.decider.code = DecisionSpam + h.writer.removalErr = fmt.Errorf("applyWrites: %w", pds.ErrUnauthorized) + + outcome, _ := h.process(t) + assert.Equal(t, EngineDeferred, outcome) + + refreshes := 0 + for _, call := range h.rec.calls { + if call == "RefreshCommunityCredentials" { + refreshes++ + } + } + assert.Equalf(t, 1, refreshes, "exactly one forced renewal per pass; got %d (sequence: %v)", + refreshes, h.rec.calls) + assert.Empty(t, h.admissions.removalCmds, "nothing committed, so nothing may be stamped on the row") +} + +// --------------------------------------------------------------------------- +// The optimistic repository update +// --------------------------------------------------------------------------- + +func TestEngine_TreatsRepositorySkipsAsSuccess(t *testing.T) { + t.Parallel() + + // The engine stamps the row itself rather than waiting for the firehose + // copy of its own commit. When the firehose gets there first the repository + // answers skipped_stale; when the row has moved on it answers + // skipped_terminal. Neither is a failure — the record IS in the community's + // repo, the firehose is the authority on what the repo says, and reporting + // an error here would dead-letter a write that succeeded. + for _, outcome := range []AdmissionOutcome{AdmissionSkippedStale, AdmissionSkippedTerminal} { + t.Run(string(outcome), func(t *testing.T) { + t.Parallel() + + h := newEngineHarness(AdmissionStatusPending, cidPtr(engineIndexedCID)) + h.admissions.acceptanceResult = AdmissionResult{Outcome: outcome, Admission: h.admissions.row} + + got, err := h.process(t) + require.NoErrorf(t, err, "a skipped stamp means the firehose won a race, not that the write failed") + assert.Equal(t, EngineAccepted, got, + "the acceptance record is in the community's repo; that is what the outcome reports") + }) + } +} + +func TestEngine_SkippedWriteStampsNothing(t *testing.T) { + t.Parallel() + + // When the community's repo already holds an acceptance pinning this exact + // CID, the writer writes nothing and reports no commit rev — getRecord does + // not reveal the revision an existing record was written at. + // + // So the engine must NOT stamp the row: ApplyAcceptance refuses an empty + // rev with ErrInvalidWatermark (correctly — an empty rev is a fabricated + // clock value), and inventing one to get past that would write a watermark + // no commit ever had. + h := newEngineHarness(AdmissionStatusPending, cidPtr(engineIndexedCID)) + h.writer.acceptanceResult = CommunityWriteResult{ + URI: "at://" + engineCommunityDID + "/" + AcceptanceCollection + "/" + SubjectRkey(enginePostURI), + RKey: SubjectRkey(enginePostURI), + CID: engineRecordCID, + Skipped: true, + } + + outcome, err := h.process(t) + require.NoError(t, err) + assert.Equal(t, EngineAccepted, outcome, + "the acceptance stands — it was already there — so the pass succeeded") + + assert.Equal(t, []string{"Get", "DecideAdmission", "WriteAcceptance"}, h.rec.calls) + assert.Empty(t, h.admissions.acceptanceCmds, + "a skipped write has no commit rev, and the repository refuses an empty one as a fabricated watermark") +} + +func TestEngine_ReportsAFailedStamp(t *testing.T) { + t.Parallel() + + // The opposite of a skip. If the repository genuinely errors, the AppView's + // view of the row now disagrees with the community's repo, and that has to + // be visible: the firehose will reconcile it, but a silent divergence is + // how a post stays invisible for hours with nothing to search for. + h := newEngineHarness(AdmissionStatusPending, cidPtr(engineIndexedCID)) + dbDown := errors.New("apply acceptance: connection refused") + h.admissions.acceptanceErr = dbDown + + _, err := h.process(t) + require.Error(t, err) + assert.ErrorIs(t, err, dbDown) +} + +func TestEngine_PassesTheSubjectToThePolicyUnchanged(t *testing.T) { + t.Parallel() + + // Small, but it is the join between the queue and the decision: a policy + // asked about the wrong subject answers confidently about someone else's + // post. + h := newEngineHarness(AdmissionStatusPending, cidPtr(engineIndexedCID)) + + _, err := h.process(t) + require.NoError(t, err) + assert.Equal(t, engineCommunityDID, h.decider.lastCommunityDID) + assert.Equal(t, enginePostURI, h.decider.lastPostURI) +} diff --git a/internal/core/posts/record_diff.go b/internal/core/posts/record_diff.go new file mode 100644 index 0000000..eb9de7f --- /dev/null +++ b/internal/core/posts/record_diff.go @@ -0,0 +1,44 @@ +package posts + +// RecordDiffClass says what kind of change an author made to a post record — +// the classification the bridgedStats exception of §5.5 turns on. +type RecordDiffClass string + +const ( + // RecordDiffNone means the two records are the same content. + RecordDiffNone RecordDiffClass = "none" + + // RecordDiffBridgedStatsOnly means the ONLY field that differs is + // bridgedStats: a bridge refreshing origin-platform vote counts. Such an + // edit is repinned in place — the acceptance moves onto the new CID with no + // status transition, no feed removal and no re-decision. + RecordDiffBridgedStatsOnly RecordDiffClass = "bridged-stats-only" + + // RecordDiffPolicyRelevant means something a moderator would judge changed, + // so the post needs full re-admission. + // + // It is also the answer for any difference this function does not + // recognise. See classifyRecordDiff. + RecordDiffPolicyRelevant RecordDiffClass = "policy-relevant" +) + +// classifyRecordDiff reports whether the change between two versions of a post +// record is the bridgedStats refresh of §5.5 or an edit needing re-admission. +// +// IT TAKES DECODED RECORDS, NOT PostRecord VALUES, AND THAT IS THE WHOLE POINT. +// A typed struct silently drops every field it does not know about, so an +// author who added an unmodelled field — or a bridge running a newer lexicon +// than this build — would produce two structs that compare equal and a diff +// classified as "nothing changed". The classification would then wave through +// an edit nobody looked at. The raw maps keep unknown fields visible. +// +// IT FAILS CLOSED. Only one specific shape of difference earns the exception: +// bridgedStats differs and NOTHING else does. Every other answer — a known +// policy field changed, a field this function has never heard of changed, a +// field appeared or vanished — is RecordDiffPolicyRelevant. The cost of +// getting that wrong in the safe direction is one unnecessary re-admission of a +// bridge post; in the unsafe direction it is edited content rendering under an +// acceptance granted to different content. +func classifyRecordDiff(oldRecord, newRecord map[string]any) RecordDiffClass { + return "" +} diff --git a/internal/core/posts/record_diff_test.go b/internal/core/posts/record_diff_test.go new file mode 100644 index 0000000..4016177 --- /dev/null +++ b/internal/core/posts/record_diff_test.go @@ -0,0 +1,249 @@ +package posts + +import ( + "testing" + + "github.com/stretchr/testify/assert" +) + +// The bridgedStats exception of docs/PRD_AUTHOR_OWNED_POSTS.md §5.5. +// +// A bridge rewrites its records constantly and specifically to refresh +// origin-platform vote counts. Every rewrite changes the record's CID, so +// without an exception every accepted bridge post would drop out of feeds into +// pending_reacceptance several times an hour and the community's repo would +// fill with re-acceptance commits that decided nothing. +// +// The exception is narrow and the narrowness is the security property. What is +// waved through is a re-decision, so anything that slips through unexamined is +// content a moderator accepted once and never looked at again. The function +// therefore answers RecordDiffBridgedStatsOnly for exactly one shape — +// bridgedStats differs, nothing else does — and RecordDiffPolicyRelevant for +// everything it does not fully understand, INCLUDING fields it has never heard +// of. That is why it reads decoded maps rather than PostRecord values: a typed +// struct discards unknown fields, so an author who added one would produce two +// structs comparing equal and an edit classified as "nothing changed". +// +// classifyRecordDiff is not called by anything yet. It ships with the engine so +// that the repin path has a decision procedure to call when it lands, and it is +// specified here so that path cannot be written against a guess. + +// basePostRecord is the fully-populated record every case below mutates. It +// carries every property social.coves.community.postv2 declares, so a case that +// changes one field is changing it in the presence of all the others. +func basePostRecord() map[string]any { + return map[string]any{ + "$type": "social.coves.community.postv2", + "community": "did:plc:cccccccccccccccccccccccc", + "title": "a post with every field", + "content": "the body", + "facets": []any{ + map[string]any{ + "index": map[string]any{"byteStart": float64(0), "byteEnd": float64(1)}, + "features": []any{map[string]any{"$type": "social.coves.richtext.facet#bold"}}, + }, + }, + "embed": map[string]any{ + "$type": "social.coves.embed.external", + "uri": "https://example.com/a", + }, + "langs": []any{"en"}, + "labels": map[string]any{"$type": "com.atproto.label.defs#selfLabels", "values": []any{map[string]any{"val": "spoiler"}}}, + "tags": []any{"gardening"}, + "crosspostOf": map[string]any{ + "uri": "at://did:plc:abc123/social.coves.community.postv2/3kjzl5kcb2s2v", + "cid": "bafyreiaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa", + }, + "crosspostChain": []any{"did:plc:cccccccccccccccccccccccc"}, + "createdAt": "2026-07-01T12:00:00Z", + "bridgedStats": map[string]any{ + "upvotes": float64(10), + "downvotes": float64(2), + "asOf": "2026-07-01T12:00:00Z", + }, + } +} + +// withField returns the base record with one field set, or removed when value +// is nil. +func withField(field string, value any) map[string]any { + record := basePostRecord() + if value == nil { + delete(record, field) + return record + } + record[field] = value + return record +} + +func TestClassifyRecordDiff(t *testing.T) { + t.Parallel() + + for _, tc := range []struct { + name string + change map[string]any + want RecordDiffClass + why string + }{ + { + name: "identical records", + change: basePostRecord(), + want: RecordDiffNone, + why: "a redelivery of the same record is not an edit and must not cost a re-decision", + }, + { + name: "bridgedStats counts refreshed", + change: withField("bridgedStats", map[string]any{ + "upvotes": float64(97), + "downvotes": float64(3), + "asOf": "2026-07-01T13:00:00Z", + }), + want: RecordDiffBridgedStatsOnly, + why: "this is the exception itself: the vote counts a bridge exists to refresh", + }, + { + name: "bridgedStats added where there was none", + change: withField("bridgedStats", map[string]any{"upvotes": float64(1), "downvotes": float64(0)}), + want: RecordDiffBridgedStatsOnly, + why: "the first refresh after a bridge starts reporting is still only a stats change", + }, + { + name: "bridgedStats removed", + change: withField("bridgedStats", nil), + want: RecordDiffBridgedStatsOnly, + why: "a bridge that stops reporting counts has still changed nothing a moderator judges", + }, + { + name: "title changed", + change: withField("title", "a completely different post"), + want: RecordDiffPolicyRelevant, + why: "the headline is the single most re-decidable thing about a post", + }, + { + name: "content changed", + change: withField("content", "something a moderator never saw"), + want: RecordDiffPolicyRelevant, + why: "the body is what acceptance accepted", + }, + { + name: "content removed", + change: withField("content", nil), + want: RecordDiffPolicyRelevant, + why: "a field vanishing is as much an edit as a field changing", + }, + { + name: "facets changed", + change: withField("facets", []any{ + map[string]any{ + "index": map[string]any{"byteStart": float64(0), "byteEnd": float64(5)}, + "features": []any{map[string]any{"$type": "social.coves.richtext.facet#link", "uri": "https://elsewhere.example/"}}, + }, + }), + want: RecordDiffPolicyRelevant, + why: "facets carry links; repointing one changes where readers land without touching the text", + }, + { + name: "embed changed", + change: withField("embed", map[string]any{"$type": "social.coves.embed.external", "uri": "https://elsewhere.example/"}), + want: RecordDiffPolicyRelevant, + why: "swapping the embedded media is the classic bait-and-switch against a stale acceptance", + }, + { + name: "labels changed", + change: withField("labels", map[string]any{"$type": "com.atproto.label.defs#selfLabels", "values": []any{}}), + want: RecordDiffPolicyRelevant, + why: "self-labels are how content warnings are declared; dropping one un-warns the post", + }, + { + name: "tags changed", + change: withField("tags", []any{"politics"}), + want: RecordDiffPolicyRelevant, + why: "tags steer distribution", + }, + { + name: "community changed", + change: withField("community", "did:plc:dddddddddddddddddddddddd"), + want: RecordDiffPolicyRelevant, + why: "the post claims to belong somewhere else entirely", + }, + { + name: "langs changed", + change: withField("langs", []any{"ru"}), + want: RecordDiffPolicyRelevant, + why: "not on §5.5's enumerated list, and everything not on the list is re-decided", + }, + { + name: "crosspostOf changed", + change: withField("crosspostOf", map[string]any{"uri": "at://did:plc:abc123/social.coves.community.postv2/other", "cid": "bafyreibbbb"}), + want: RecordDiffPolicyRelevant, + why: "the post now claims to mirror something else", + }, + { + name: "createdAt changed", + change: withField("createdAt", "2026-07-02T12:00:00Z"), + want: RecordDiffPolicyRelevant, + why: "a record rewritten with a new timestamp is not the record that was accepted", + }, + { + name: "an unknown field appeared", + change: withField("summary", "a field this build has never heard of"), + want: RecordDiffPolicyRelevant, + why: "FAIL CLOSED. A bridge or a client running a newer lexicon than this AppView will " + + "send fields this function cannot name, and treating an unrecognised difference as " + + "harmless would wave through exactly the edits nobody has looked at yet", + }, + { + name: "an unknown field changed value", + change: withField("$type", "social.coves.community.postv3"), + want: RecordDiffPolicyRelevant, + why: "even the record's own type is a difference that has to be re-decided, not assumed benign", + }, + { + name: "bridgedStats changed AND a policy field changed", + change: func() map[string]any { + record := withField("title", "a completely different post") + record["bridgedStats"] = map[string]any{"upvotes": float64(99), "downvotes": float64(0)} + return record + }(), + want: RecordDiffPolicyRelevant, + why: "the exception is bridgedStats AND NOTHING ELSE; a real edit smuggled alongside a " + + "stats refresh is the obvious way to launder one past re-admission", + }, + } { + t.Run(tc.name, func(t *testing.T) { + t.Parallel() + assert.Equalf(t, tc.want, classifyRecordDiff(basePostRecord(), tc.change), "%s", tc.why) + }) + } +} + +func TestClassifyRecordDiff_TreatsAMissingRecordAsPolicyRelevant(t *testing.T) { + t.Parallel() + + // The engine reaches this function with whatever it managed to decode. A + // nil old record means the AppView has no prior version to compare against + // — it cannot know the change is only stats, so it must not say so. + assert.Equal(t, RecordDiffPolicyRelevant, classifyRecordDiff(nil, basePostRecord()), + "with no prior version there is nothing to prove the change is stats-only") + + assert.Equal(t, RecordDiffPolicyRelevant, classifyRecordDiff(basePostRecord(), nil), + "a record that decoded to nothing is not a stats refresh") +} + +func TestClassifyRecordDiff_IsSymmetricAboutBridgedStats(t *testing.T) { + t.Parallel() + + // Direction must not matter. Records arrive out of order from overlapping + // feeds, and a classifier that answered differently depending on which copy + // it called "old" would repin in one delivery order and re-decide in the + // other. + withStats := basePostRecord() + withoutStats := withField("bridgedStats", nil) + + assert.Equal(t, RecordDiffBridgedStatsOnly, classifyRecordDiff(withStats, withoutStats)) + assert.Equal(t, RecordDiffBridgedStatsOnly, classifyRecordDiff(withoutStats, withStats)) + + edited := withField("title", "edited") + assert.Equal(t, RecordDiffPolicyRelevant, classifyRecordDiff(withStats, edited)) + assert.Equal(t, RecordDiffPolicyRelevant, classifyRecordDiff(edited, withStats)) +} diff --git a/internal/core/posts/rkey.go b/internal/core/posts/rkey.go new file mode 100644 index 0000000..141091a --- /dev/null +++ b/internal/core/posts/rkey.go @@ -0,0 +1,30 @@ +package posts + +// SubjectRkey is the record key a community's records about one post use. +// +// It is the unpadded lowercase base32 encoding of the SHA-256 digest of the +// post's AT-URI: a fixed 52 characters drawn entirely from the rkey-safe +// charset, well inside the 512-byte limit, for a subject of any length +// (docs/PRD_AUTHOR_OWNED_POSTS.md §3.2). +// +// WHY A DIGEST RATHER THAN A READABLE TRANSFORM. The obvious scheme — strip +// `at://`, swap `/` for `:` — is not total over the legal subject space. DIDs +// may run to 2048 bytes and may carry percent-escapes, so the transform can +// produce keys that exceed the rkey limit or leave the rkey charset, and a +// non-total key function fails on exactly the identifiers an attacker gets to +// choose. A fixed-size digest is total and just as deterministic. +// +// WHY ONE FUNCTION FOR BOTH RECORD TYPES. The acceptance and the removal for a +// subject share this key. Record keys are scoped to their collection, so there +// is no collision to avoid, and a per-collection salt would buy nothing while +// costing the property that makes the removal commit shapeable: both records +// for a subject are found by ONE key derivation, so a pre-read of both is two +// lookups of one computed value rather than a search. +// +// THE ARGUMENT IS BYTES, NOT A PARSED URI. Whatever the admission row holds is +// what gets hashed — no normalization, no percent-decoding, no case folding. +// The row's bytes are the identity the AppView indexes under, so a writer that +// normalized would key its records to a URI the reader never looks up. +func SubjectRkey(postURI string) string { + return "" +} diff --git a/internal/core/posts/rkey_test.go b/internal/core/posts/rkey_test.go new file mode 100644 index 0000000..ee24346 --- /dev/null +++ b/internal/core/posts/rkey_test.go @@ -0,0 +1,213 @@ +package posts + +import ( + "strings" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// The deterministic record key of docs/PRD_AUTHOR_OWNED_POSTS.md §3.2. +// +// Every idempotency claim in the acceptance design rests on this function and +// nothing else. Three independent writers — the synchronous fast path, the +// firehose engine, and the notify endpoint — are safe to race only because they +// all compute the SAME rkey for the same post, so a race is a putRecord of one +// record rather than two records nobody can reconcile. Re-acceptance after an +// edit is an update rather than a second acceptance for the same reason. +// +// THE GOLDEN VALUES ARE HARD-CODED, NOT RECOMPUTED. A test that hashed the URI +// itself and compared would pass against base32-Hex, against uppercase, against +// padded output, and against a completely different digest — every variant that +// would silently repartition every community repo in the network the day it +// shipped. The constants below were produced OUTSIDE Go, with: +// +// python3 -c 'import hashlib,base64;print(base64.b32encode( +// hashlib.sha256(URI.encode()).digest()).decode().lower().rstrip("="))' +// +// so they encode the spec rather than the implementation. + +const ( + // goldenSubjectURI is the canonical vector: an ordinary post in an ordinary + // author repo. + goldenSubjectURI = "at://did:plc:abc123/social.coves.community.postv2/3kjzl5kcb2s2v" + goldenSubjectRkey = "xxdmibjaexx43drostplutjbp7g4oaw3uriugf5twafpldfkupca" + + // goldenSiblingURI differs from the canonical vector in its final character + // only. It is here so that "distinct URIs get distinct keys" is proven at + // the smallest possible difference rather than at a comfortable one. + goldenSiblingURI = "at://did:plc:abc123/social.coves.community.postv2/3kjzl5kcb2s2w" + goldenSiblingRkey = "iyhgczhg7xsbrayzrrs2qa4fks6amctx7ghyjakqyrlhxztbbl5a" +) + +// subjectRkeyLength is what §3.2 promises: SHA-256 is 256 bits, base32 packs 5 +// bits per character, and 256/5 rounds up to 52 characters (with the padding +// stripped, which is why 52 and not 56). +const subjectRkeyLength = 52 + +// rkeyCharset is the set base32 draws from once lowercased. It is also a subset +// of the atProto record-key charset, which is the property that makes a digest +// safe to use as a key at all. +const rkeyCharset = "abcdefghijklmnopqrstuvwxyz234567" + +func TestSubjectRkey_GoldenVector(t *testing.T) { + t.Parallel() + + assert.Equal(t, goldenSubjectRkey, SubjectRkey(goldenSubjectURI), + "the rkey for %s is fixed by §3.2 and by every acceptance record already written under it; "+ + "a different value here means the encoding changed (uppercase, padded, base32-Hex, "+ + "or a different digest) and every community repo in the network just repartitioned", + goldenSubjectURI) + + assert.Equal(t, goldenSiblingRkey, SubjectRkey(goldenSiblingURI)) +} + +func TestSubjectRkey_ShapeIsFixedRegardlessOfSubjectLength(t *testing.T) { + t.Parallel() + + // The point of a digest is that it is TOTAL over the legal subject space. + // A readable transform of the URI is not: DIDs may run to 2048 bytes, so + // the transform can exceed the 512-byte rkey limit, and they may carry + // percent-escapes, so it can leave the rkey charset. Both of those are + // attacker-chosen inputs. + for _, tc := range []struct { + name string + uri string + }{ + {name: "ordinary did:plc subject", uri: goldenSubjectURI}, + { + // 580 bytes: a did:web whose authority alone is over the 512-byte + // rkey limit, so a readable transform would produce a key the PDS + // refuses outright. + name: "did:web subject longer than the 512-byte rkey limit", + uri: "at://" + longDIDWeb() + "/social.coves.community.postv2/3kjzl5kcb2s2v", + }, + { + // The maximum legal DID length. Nothing about the answer's shape + // may vary with it. + name: "2048-byte DID, the legal maximum", + uri: "at://" + didOfLength(2048) + "/social.coves.community.postv2/3kjzl5kcb2s2v", + }, + { + name: "percent-escaped authority", + uri: "at://did:web:example.com%3A8443/social.coves.community.postv2/3kjzl5kcb2s2v", + }, + { + name: "empty subject", + uri: "", + }, + } { + t.Run(tc.name, func(t *testing.T) { + t.Parallel() + + rkey := SubjectRkey(tc.uri) + + assert.Lenf(t, rkey, subjectRkeyLength, + "§3.2 promises a fixed %d characters for a subject of any length; this one was %d bytes", + subjectRkeyLength, len(tc.uri)) + + for i, r := range rkey { + assert.Containsf(t, rkeyCharset, string(r), + "rkey %q holds %q at index %d, which is outside the unpadded lowercase base32 "+ + "alphabet — uppercase and '=' are the two slips that produce a key the PDS rejects", + rkey, string(r), i) + } + }) + } +} + +func TestSubjectRkey_LongSubjectGoldenVectors(t *testing.T) { + t.Parallel() + + // The long vectors get golden values too, not merely a shape check. A + // truncating implementation — hashing only the first N bytes, or a + // readable transform that clipped to 512 — passes the shape assertions + // above and collides two different long subjects onto one key. + assert.Equal(t, + "fktiwazbhqfsqjg5e7yypm3ukbalk72bfplcr2d2ukm7iuwqeyqa", + SubjectRkey("at://"+longDIDWeb()+"/social.coves.community.postv2/3kjzl5kcb2s2v")) + + assert.Equal(t, + "4pri7ptezxzyzpzzysj5for6shd54yt5vvch6v3nj4e7kmfg62qq", + SubjectRkey("at://"+didOfLength(2048)+"/social.coves.community.postv2/3kjzl5kcb2s2v")) +} + +func TestSubjectRkey_IsStableAcrossCalls(t *testing.T) { + t.Parallel() + + // Stability is the whole contract. If this ever stops holding — a map + // iteration folded into the derivation, a timestamp, a random salt — the + // three independent writers stop converging and every re-fire mints a + // second acceptance record for the same post. + first := SubjectRkey(goldenSubjectURI) + for i := 0; i < 32; i++ { + require.Equalf(t, first, SubjectRkey(goldenSubjectURI), + "call %d returned a different key for the same subject", i) + } +} + +func TestSubjectRkey_DistinguishesSubjectsThatDifferInBytesOnly(t *testing.T) { + t.Parallel() + + // The row's bytes are the identity. No normalization, no percent-decoding, + // no case folding — because the AppView indexes and looks records up under + // exactly the bytes it stored, so a writer that normalized would key its + // records to a URI the reader never asks for. + escaped := "at://did:web:example.com%3A8443/social.coves.community.postv2/3kjzl5kcb2s2v" + decoded := "at://did:web:example.com:8443/social.coves.community.postv2/3kjzl5kcb2s2v" + + assert.NotEqual(t, SubjectRkey(escaped), SubjectRkey(decoded), + "the escaped and decoded spellings are different byte strings and must key differently; "+ + "a function that percent-decoded first would let two rows that the AppView keeps apart "+ + "collide onto one acceptance record") + + assert.Equal(t, "ip6rb5ppturmcux6r54fglc57je3cgfofvi3eajodrxqkgcwjrna", SubjectRkey(escaped)) + assert.Equal(t, "jusonlequb2pkc34y5uwa65ng22cwz5qef3gorg5huqebglag37a", SubjectRkey(decoded)) + + assert.NotEqual(t, SubjectRkey(goldenSubjectURI), SubjectRkey(goldenSiblingURI), + "two posts differing in one character of their rkey must not share an acceptance record") + + // Case is a byte difference like any other. + assert.NotEqual(t, SubjectRkey(goldenSubjectURI), SubjectRkey(strings.ToUpper(goldenSubjectURI))) +} + +func TestSubjectRkey_IsSharedByTheAcceptanceAndTheRemoval(t *testing.T) { + t.Parallel() + + // PINNED SO NOBODY SALTS IT PER COLLECTION. There is exactly one derivation + // for a subject, and the acceptance and the removal both use it. That is + // what makes the removal commit of §3.3 shapeable: WriteRemoval pre-reads + // BOTH records to decide whether to emit a delete and whether the removal + // is a create or an update, and it can only do that with two lookups of one + // computed value. + // + // Record keys are scoped to their collection, so there is no collision to + // avoid and a per-collection salt would buy nothing at all. + rkey := SubjectRkey(goldenSubjectURI) + require.NotEmpty(t, rkey) + + acceptanceURI := "at://did:plc:community/" + AcceptanceCollection + "/" + rkey + removalURI := "at://did:plc:community/" + RemovalCollection + "/" + rkey + + assert.True(t, strings.HasSuffix(acceptanceURI, "/"+rkey)) + assert.True(t, strings.HasSuffix(removalURI, "/"+rkey), + "the removal record for a subject must live at the SAME rkey as its acceptance") +} + +// longDIDWeb is a did:web whose authority alone exceeds the 512-byte rkey +// limit: eight maximum-length DNS labels and a TLD. +func longDIDWeb() string { + label := strings.Repeat("a", 63) + labels := make([]string, 8) + for i := range labels { + labels[i] = label + } + return "did:web:" + strings.Join(labels, ".") + ".example.com" +} + +// didOfLength returns a syntactically DID-shaped identifier of exactly n bytes. +func didOfLength(n int) string { + const prefix = "did:web:" + return prefix + strings.Repeat("b", n-len(prefix)) +}