From 501c61480e6712c4d2749197e1a620dcaa1a8c54 Mon Sep 17 00:00:00 2001 From: Bretton Date: Sat, 8 Aug 2026 05:37:41 -0700 Subject: [PATCH] =?UTF-8?q?feat(posts):=20GREEN=20cycle=201=20=E2=80=94=20?= =?UTF-8?q?the=20write=20path=20lands=20in=20the=20author's=20repo?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Task 6 GREEN. Implementation only; no test file is touched (RED's a9744e1 is the contract). All 10 T0 and all 26 T1 reds are green. The record moves (§3.1, §4.2) SubmissionRkey the deterministic postv2 key: sha256 over (community DID, fingerprint, bucket) newline-delimited, a timestamp placed INSIDE the submission's own dedupe bucket (bucket start + digest offset mod the window), clock ID from further digest bits, encoded by syntax.NewTID so the one encoder in the tree is the one ParseTID ships beside. A non-positive window collapses to the epoch rather than dividing by zero, matching dedupeBucket's answer to the same misconfiguration. AdmissionDecision.DedupeBucket reported by reserveSubmission so the rkey is derived from THE SAME bucket the ledger deduped against. A second clock read at the call site would work almost always and fail exactly at a bucket boundary — the one case the deterministic key exists for. CreatePost admission -> open the AUTHOR's repo -> community token -> build -> enhance -> create-only PutRecordWithCommit(swapRecord "") at the derived rkey -> meter -> seed the row and settle it. The author-repo open is ahead of the community's token because it is the credential the post cannot be written without; a community whose token will not refresh costs a link preview, not a post. Nothing after the record commits may fail the request. postV2From the write boundary, and the single place the author field is dropped and $type re-stamped. The community answers before we do (§4.2 steps 4 and 5) AcceptanceEngine.AcceptSubmission row-pending + CID-match, then accept. The decider is NOT consulted: CreatePost has already run that policy, and the production decider reads an index the post is not in yet, so a fast path that re-decided would refuse every post it was handed. settleSubmission seeds UpsertPending under the URI the consumer builds, reports accepted when the row already is (a retry of a settled post), maps ErrCommunityNotHosted to pending without alarm, and cannot fail the request. Editing and deleting UpdatePost authority-equals-session before anything is fetched, then read -> preserve community and createdAt from the STANDING record -> put guarded by the CID just read. A lost swap is ErrConcurrentModification. The submission ledger is untouched. DeletePost split by collection: postv2 goes to the author's own repo under local authorization, the deprecated community.post keeps the old community-credential path intact until task 8. Supporting pds.ErrNoCommit VERIFIED AGAINST A LIVE PDS: a create-only put of IDENTICAL bytes is a 200 no-op with no commit object, evaluated before the swap guard; the same put of DIFFERENT bytes is InvalidSwap. Both mean "it is already there", and createAuthorRecord answers both by reporting what stands. Without the sentinel the byte-identical retry read as a malformed body. NewAuthorRepoFactory the production seam: a person's session passes through, an aggregator's stored tokens are resumed, and nothing to resume is ErrNoAuthorCredentials rather than a 401 nobody can act on. post.update service, handler, error mapping and route against the lexicon file that already existed. postv2.json originalAuthor/federatedFrom/location declared as optional `unknown` — additive-optional evolution. Unconstrained deliberately: no bridge has fixed their shape, and inventing one would commit the first real bridge to a contract we made up. Three deviations, flagged rather than buried 1. The fingerprint and the embed pipeline still take the deprecated PostRecord, converted once at the write boundary. Retyping them is not a golden-value change but a COMPILE break in admit_matrix_test.go and embed_conversion_test.go, which are in-package — it would take the whole posts test binary with it and leave every red unverifiable. Cycle 2, with the tests. It also keeps the fingerprint byte-stable across the deploy, so live dedupe rows survive the flip. 2. The external-embed thumb guard moved into validateCreateRequest. It is pure and needs no network, and two tests demand opposite things of its position otherwise: the create-validation test wants the author-repo open to be the first failure after admission, the handler's embed test wants a 400 on a stack with no author repo at all. validateEmbed is untouched — its test pins that it ignores a junk thumb. 3. buildAcceptanceEngine now returns the engine, so the write path and the queue driver settle the same subjects through one instance. Co-Authored-By: Claude Fable 5 --- cmd/server/wiring.go | 25 +- internal/api/handlers/post/errors.go | 17 + internal/api/handlers/post/update.go | 84 +++ internal/api/routes/post.go | 7 +- .../social/coves/community/postv2.json | 12 + internal/atproto/pds/applywrites.go | 18 +- internal/atproto/pds/errors.go | 19 + internal/core/posts/admit.go | 23 +- internal/core/posts/author_repo_factory.go | 93 +++ internal/core/posts/engine.go | 45 +- internal/core/posts/postv2.go | 62 +- internal/core/posts/service.go | 692 +++++++++++++----- 12 files changed, 876 insertions(+), 221 deletions(-) create mode 100644 internal/api/handlers/post/update.go create mode 100644 internal/core/posts/author_repo_factory.go diff --git a/cmd/server/wiring.go b/cmd/server/wiring.go index 5a1ecf8..8bfbf18 100644 --- a/cmd/server/wiring.go +++ b/cmd/server/wiring.go @@ -349,10 +349,22 @@ func (a *application) buildServices(ctx context.Context) error { // existed, so this is not an optional enrichment — it is the enforcement. // The limits come from config, which refuses to start the process with a // non-positive one rather than letting an omission read as "unlimited". + // + // The engine is built FIRST because the write path now holds it: §4.2 step + // 4 has CreatePost settle a local community's admission before it answers, + // so the author gets the community's decision with their post rather than + // waiting for their own write to come back around the firehose. + acceptanceEngine := a.buildAcceptanceEngine() + a.postService = posts.NewPostService( a.postRepo, a.communityService, a.aggregatorService, blobService, unfurlService, a.blueskyService, a.cfg.PDS.URL, posts.WithBlockChecker(a.userBlockRepo), + // The AUTHOR's own credentials: a browser session when there is one, + // and an aggregator's stored tokens when there is not (§4.2 step 3). + posts.WithAuthorRepoFactory( + posts.NewAuthorRepoFactory(a.oauthClient.ClientApp, aggregators.DefaultSessionID)), + posts.WithSyncAcceptance(a.admissionRepo, acceptanceEngine), posts.WithAdmissionPolicy(posts.AdmissionPolicy{ Ledger: postgresRepo.NewSubmissionLedger(a.db), Bans: a.communityService, @@ -371,8 +383,6 @@ func (a *application) buildServices(ctx context.Context) error { // is nothing on the firehose to read. a.postStatusService = posts.NewStatusService(a.admissionRepo) - a.buildAcceptanceEngine() - // Subject existence is deliberately not validated: the vote is written to // the user's own PDS regardless, and the Jetstream consumer only updates // counts for subjects that still exist. Checking here would trade a @@ -415,7 +425,7 @@ func (a *application) buildServices(ctx context.Context) error { // STORED CREDENTIALS rather than on communities.hosted_by_did, which is copied // out of a community's own profile record and can therefore be claimed by any // repo on the network. -func (a *application) buildAcceptanceEngine() { +func (a *application) buildAcceptanceEngine() *posts.AcceptanceEngine { repoFactory := posts.NewCommunityRepoFactory(a.communityService) a.communityWriter = posts.NewCommunityRecordWriter(repoFactory, time.Now) @@ -452,11 +462,18 @@ func (a *application) buildAcceptanceEngine() { if a.cfg.Submissions.AcceptanceQueueInterval <= 0 { slog.Warn("acceptance queue driver disabled (ACCEPTANCE_QUEUE_INTERVAL=0); " + "posts left undecided by the fast path and the firehose will not be revisited") - return + return engine } a.acceptanceQueue = posts.NewQueueDriver(a.admissionRepo, engine, time.Now, posts.WithQueueBatchSize(a.cfg.Submissions.AcceptanceQueueBatchSize)) + + // RETURNED, NOT ONLY STORED IN THE DRIVER. The write path's fast path and + // the queue driver settle the same subjects through the same engine on + // purpose: the deterministic rkeys and swap guards make concurrent passes + // converge, and a second engine instance would just be a second set of the + // same collaborators. + return engine } // adminReportAlertOptions builds the operator-alert wiring for admin reports. diff --git a/internal/api/handlers/post/errors.go b/internal/api/handlers/post/errors.go index aafd53e..8d0a356 100644 --- a/internal/api/handlers/post/errors.go +++ b/internal/api/handlers/post/errors.go @@ -32,6 +32,23 @@ var errorMapper = xrpc.NewMapper("post", xrpc.Sentinel(posts.ErrNotFound, http.StatusNotFound, "NotFound", "Post not found"), + // The AppView holds nothing to authenticate as the author with. It is NOT + // the shared re-auth rule's 401: that one means the CALLER's session is + // dead and signing in again fixes it, which is true for a person and false + // for an aggregator posting on stored tokens — nobody is at the keyboard to + // sign it in, and telling its operator to do so would hide a revoked grant + // behind a message aimed at a browser. 503 says what is true of both: the + // service cannot write on this author's behalf right now. + xrpc.Sentinel(posts.ErrNoAuthorCredentials, http.StatusServiceUnavailable, + "NoAuthorCredentials", "No credentials are available to write to your repository"), + + // A lost swap guard on an edit. 409 rather than a retry here, because the + // edit was composed against content that no longer stands: re-reading and + // re-applying it server-side would silently erase the change its author + // never saw. The client re-reads and decides. + xrpc.Sentinel(posts.ErrConcurrentModification, http.StatusConflict, + "ConcurrentModification", "The post was modified by another edit; re-read it and try again"), + // A submission refused at the admission gate, which is NOT the generic // AlreadyExists that a storage conflict produces: 409 DuplicateSubmission // tells a client whose response was lost that its post already exists and diff --git a/internal/api/handlers/post/update.go b/internal/api/handlers/post/update.go new file mode 100644 index 0000000..c70356f --- /dev/null +++ b/internal/api/handlers/post/update.go @@ -0,0 +1,84 @@ +package post + +import ( + "encoding/json" + "log" + "net/http" + + "Coves/internal/api/middleware" + "Coves/internal/api/xrpc" + "Coves/internal/core/posts" +) + +// UpdateHandler serves social.coves.community.post.update. +type UpdateHandler struct { + service posts.Service +} + +// NewUpdateHandler creates a new handler for editing posts. +func NewUpdateHandler(service posts.Service) *UpdateHandler { + return &UpdateHandler{service: service} +} + +// HandleUpdate handles POST /xrpc/social.coves.community.post.update +// +// Request body: the lexicon's input schema (uri plus the mutable fields). +// Response: { "uri": "at://...", "cid": "..." } +func (h *UpdateHandler) HandleUpdate(w http.ResponseWriter, r *http.Request) { + if r.Method != http.MethodPost { + http.Error(w, "Method not allowed", http.StatusMethodNotAllowed) + return + } + + // The same 1MB ceiling the create path uses: an edit carries the same + // content and embeds a create does, so a tighter bound here would refuse + // edits to posts this service accepted. + r.Body = http.MaxBytesReader(w, r.Body, 1*1024*1024) + + var req posts.UpdatePostRequest + if err := json.NewDecoder(r.Body).Decode(&req); err != nil { + if err.Error() == "http: request body too large" { + writeError(w, http.StatusRequestEntityTooLarge, "RequestTooLarge", + "Request body too large (max 1MB)") + return + } + writeError(w, http.StatusBadRequest, "InvalidRequest", "Invalid request body") + return + } + + // The session is both the authorization and the credential: the record is + // signed by its author, and the author is who the session says it is. There + // is no aggregator path here — an aggregator syndicates, it does not edit. + session := middleware.GetOAuthSession(r) + if session == nil { + writeError(w, http.StatusUnauthorized, "AuthRequired", "Authentication required") + return + } + + response, err := h.service.UpdatePost(r.Context(), session, req) + if err != nil { + handleUpdateError(w, err) + return + } + + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(http.StatusOK) + if err := json.NewEncoder(w).Encode(response); err != nil { + log.Printf("Failed to encode post update response: %v", err) + } +} + +// updateErrorMapper narrows the package mapper for the edit path, naming the +// missing post and the refused action the way the update lexicon's errors do; +// everything else it inherits. +var updateErrorMapper = errorMapper.With( + xrpc.Sentinel(posts.ErrNotFound, http.StatusNotFound, + "PostNotFound", "Post not found"), + xrpc.Sentinel(posts.ErrNotAuthorized, http.StatusForbidden, + "NotAuthorized", "You are not authorized to edit this post"), +) + +// handleUpdateError maps edit-specific service errors to HTTP responses. +func handleUpdateError(w http.ResponseWriter, err error) { + updateErrorMapper.Write(w, err) +} diff --git a/internal/api/routes/post.go b/internal/api/routes/post.go index bb60424..7b775b5 100644 --- a/internal/api/routes/post.go +++ b/internal/api/routes/post.go @@ -85,6 +85,7 @@ func RegisterPostRoutes( // Initialize handlers createHandler := post.NewCreateHandler(service) + updateHandler := post.NewUpdateHandler(service) deleteHandler := post.NewDeleteHandler(service) getHandler := post.NewGetHandler(service, voteService, blueskyService) @@ -93,6 +94,11 @@ func RegisterPostRoutes( // Supports both OAuth (users) and service JWT (aggregators) authentication r.With(authMiddleware.RequireAuth).Post("/xrpc/social.coves.community.post.create", createHandler.HandleCreate) + // social.coves.community.post.update - edit a post in place + // Only the post's author can edit it, and only postv2 records (which live + // in the author's own repo) are editable at all. + r.With(authMiddleware.RequireAuth).Post("/xrpc/social.coves.community.post.update", updateHandler.HandleUpdate) + // social.coves.community.post.delete - delete a post from a community // Only post authors can delete their own posts r.With(authMiddleware.RequireAuth).Post("/xrpc/social.coves.community.post.delete", deleteHandler.HandleDelete) @@ -131,6 +137,5 @@ func RegisterPostRoutes( Get("/xrpc/social.coves.community.post.getStatus", statusHandler.HandleGetStatus) // Future endpoints (Beta): - // r.With(authMiddleware.RequireAuth).Post("/xrpc/social.coves.community.post.update", updateHandler.HandleUpdate) // r.Get("/xrpc/social.coves.community.post.list", listHandler.HandleList) } diff --git a/internal/atproto/lexicon/social/coves/community/postv2.json b/internal/atproto/lexicon/social/coves/community/postv2.json index c9adffb..7c1b3fe 100644 --- a/internal/atproto/lexicon/social/coves/community/postv2.json +++ b/internal/atproto/lexicon/social/coves/community/postv2.json @@ -93,6 +93,18 @@ "type": "ref", "ref": "#bridgedStats", "description": "Bridge-asserted aggregate of origin-platform votes for federated/bridged content. Set by the bridge that materialized this record; absent for natively-authored posts." + }, + "originalAuthor": { + "type": "unknown", + "description": "The author on the origin platform, for content a bridge materialized here. Deliberately unconstrained: no bridge has fixed this shape yet, and pinning one now would be inventing a contract the first real bridge would contradict. Consumers MUST tolerate any object here and MUST NOT treat it as an authorship claim about this repo - authorship is the repository the record lives in." + }, + "federatedFrom": { + "type": "unknown", + "description": "The origin platform or instance this post was federated from, for bridged content. Unconstrained for the same reason as originalAuthor." + }, + "location": { + "type": "unknown", + "description": "Optional place this post is about or was authored from. Unconstrained pending a geo vocabulary; absent for the overwhelming majority of posts." } } } diff --git a/internal/atproto/pds/applywrites.go b/internal/atproto/pds/applywrites.go index f3a47c5..225b7ac 100644 --- a/internal/atproto/pds/applywrites.go +++ b/internal/atproto/pds/applywrites.go @@ -332,17 +332,23 @@ func (c *client) CreateRecordWithCommit(ctx context.Context, collection, rkey st // recordCommit validates a single-record write's response and shapes it. // -// A 200 carrying no uri/cid, or no commit rev, is a malformed body from the PDS -// or something in front of it. It is reported rather than returned as a -// zero-valued success, because the two things missing here are precisely the two -// the caller is about to persist: the record it will reference, and the -// revision it will order by. +// A 200 carrying no uri/cid is a malformed body from the PDS or something in +// front of it. It is reported rather than returned as a zero-valued success, +// because what is missing is precisely what the caller is about to persist: the +// record it will reference. +// +// A 200 carrying uri and cid but NO COMMIT is a different thing, and it gets +// its own sentinel: the PDS accepted a write of bytes identical to what already +// stood there and committed nothing. Callers that need the revision to order by +// must still treat that as a failure — there is no rev to stamp — but a guarded +// CREATE wants to hear it, because "the record already exists and is identical" +// is the answer a retry after a lost response is owed. See ErrNoCommit. func recordCommit(operation, collection, uri, cid string, commit *commitResponse) (*RecordCommit, error) { if uri == "" || cid == "" { return nil, fmt.Errorf("%s: PDS returned success without uri/cid (collection %s)", operation, collection) } if commit == nil || commit.Rev == "" { - return nil, fmt.Errorf("%s: PDS returned success without a commit rev (collection %s)", operation, collection) + return nil, fmt.Errorf("%s: %w (collection %s, uri %s)", operation, ErrNoCommit, collection, uri) } return &RecordCommit{ diff --git a/internal/atproto/pds/errors.go b/internal/atproto/pds/errors.go index e5a1bc2..396869f 100644 --- a/internal/atproto/pds/errors.go +++ b/internal/atproto/pds/errors.go @@ -40,6 +40,25 @@ var ( // own sentinel. ErrSwapConflict = errors.New("swap conflict") + // ErrNoCommit indicates a single-record write that the PDS accepted without + // producing a commit: the record it was asked to write was byte-identical + // to the one already standing at that key, so there was nothing to commit. + // + // VERIFIED AGAINST A LIVE PDS: putRecord answers a no-op write with HTTP + // 200 carrying uri and cid but NO `commit` object, and it does so BEFORE + // the swapRecord guard is evaluated — a create-only put (swapRecord null) + // of an identical record is a 200 no-op, while the same put of DIFFERENT + // bytes is InvalidSwap. + // + // It is an error rather than a zero-valued success because the commit rev + // is exactly what most callers are about to persist as an ordering + // watermark, and a fabricated one is worse than a failure. It is its OWN + // sentinel rather than a generic malformed-body report because for a + // GUARDED CREATE it is not a failure at all: it means the record already + // exists and is identical, which is precisely what a retry after a lost + // response should be told. + ErrNoCommit = errors.New("write produced no commit") + // 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 diff --git a/internal/core/posts/admit.go b/internal/core/posts/admit.go index e143953..ca69337 100644 --- a/internal/core/posts/admit.go +++ b/internal/core/posts/admit.go @@ -106,6 +106,16 @@ type AdmissionDecision struct { // follows fails — see SubmissionLedger. Reservation *SubmissionReservation + // DedupeBucket is the window index the reservation was taken in, reported + // so the caller derives the post's record key from THE SAME bucket the + // ledger deduped against (SubmissionRkey). + // + // Reading the clock a second time at the call site would work almost + // always and fail exactly at a bucket boundary — the retry that crossed it + // would aim at a different rkey than the ledger scoped it to, which is the + // one case the deterministic key exists for. + DedupeBucket int64 + // Cause carries the underlying error behind a refusal, when there is one, // so the caller can wrap it and keep it matchable. // @@ -587,11 +597,12 @@ func reserveSubmission(ctx context.Context, deps admissionDeps, req AdmissionReq // window. It runs ahead of the quota so that a client retrying after a lost // response is told its post already exists rather than told to slow down. now := deps.now() + bucket := dedupeBucket(now, deps.limits.DedupeWindow) reservation, err := deps.ledger.Reserve(ctx, ReserveSubmissionCommand{ AuthorDID: req.AuthorDID, CommunityDID: community.DID, Fingerprint: req.Fingerprint, - DedupeBucket: dedupeBucket(now, deps.limits.DedupeWindow), + DedupeBucket: bucket, }) if err != nil { if errors.Is(err, ErrDuplicateSubmission) { @@ -618,7 +629,7 @@ func reserveSubmission(ctx context.Context, deps admissionDeps, req AdmissionReq } } - return AdmissionDecision{Community: community, Reservation: &reservation}, nil + return AdmissionDecision{Community: community, Reservation: &reservation, DedupeBucket: bucket}, nil } // undecided reports that the decision could NOT be made. @@ -702,6 +713,14 @@ func dedupeBucket(now time.Time, window time.Duration) int64 { // and it changes what readers ultimately see. Two submissions differing only // in their thumbnail are different posts, and excluding it would refuse the // second as a repeat of the first. +// +// IT STILL HASHES THE DEPRECATED PostRecord SHAPE, not the postv2 record that +// is actually written. The two describe the same submission — every field a +// CreatePostRequest can populate exists on both — so the fingerprint identifies +// the same posts either way, and keeping this shape keeps every dedupe row +// already on the ledger valid across the deploy that flips the write path. +// Moving it onto PostV2Record is a task-6 cycle-2 obligation: it is a pure type +// change here, and the T0 tests that pin this signature have to move with it. func submissionFingerprint(record PostRecord, thumbnailURL *string) string { // The record is taken by value, so clearing fields here cannot affect the // record the caller goes on to write. diff --git a/internal/core/posts/author_repo_factory.go b/internal/core/posts/author_repo_factory.go new file mode 100644 index 0000000..fd75e3d --- /dev/null +++ b/internal/core/posts/author_repo_factory.go @@ -0,0 +1,93 @@ +package posts + +import ( + "context" + "fmt" + + "github.com/bluesky-social/indigo/atproto/auth/oauth" + "github.com/bluesky-social/indigo/atproto/syntax" + + "Coves/internal/atproto/pds" +) + +// NewAuthorRepoFactory builds the production AuthorRepoFactory: the seam the +// AUTHOR's own credentials arrive through, now that a post is signed by its +// author rather than by the community it was submitted to (§4.2 step 3). +// +// # TWO KINDS OF AUTHOR, ONE SEAM +// +// A person posting from a browser or the mobile app arrives with an OAuth +// session the API boundary already holds, and it is passed straight through: +// the session IS the credential, and resolving it a second time from the store +// would only introduce a way for the two to disagree. +// +// An aggregator has no session on the request — it authenticates by API key — +// so its credentials come from the tokens it granted when it registered +// (migration 025), resumed under storedSessionID. Before the write path +// flipped, an aggregator needed no repository at all: its post went into the +// community's repo under the community's token, with the aggregator named in a +// field. Now it writes into its own repo like any other author. +// +// # WHY A MISSING STORED SESSION IS ITS OWN ERROR +// +// A human's dead session surfaces as pds.ErrSessionExpired, which the boundary +// answers with a 401 the client can act on: sign in again. An aggregator cannot +// sign in again — nobody is at the keyboard — so the same condition is +// ErrNoAuthorCredentials instead, which is an operator problem: the service is +// running, correctly configured, and completely unable to post. Collapsing them +// would have a revoked aggregator grant diagnosed as a PDS outage, or a +// signed-out user told to file a ticket. +func NewAuthorRepoFactory(oauthClient *oauth.ClientApp, storedSessionID string) AuthorRepoFactory { + return func(ctx context.Context, authorDID string, session *oauth.ClientSessionData) (AuthorRepo, error) { + if oauthClient == nil { + return nil, fmt.Errorf("opening the repository of %s: no OAuth client is wired: %w", + authorDID, ErrNoAuthorCredentials) + } + + did, err := syntax.ParseDID(authorDID) + if err != nil { + return nil, NewValidationError("authorDid", fmt.Sprintf("not a DID: %s", err)) + } + + if session == nil { + // The non-interactive author. Resumed here rather than inside + // NewFromOAuthSession so that the stored session's own HostURL — + // which PDS the aggregator's repo is actually on — comes from the + // store rather than being guessed at, and so that "there is nothing + // to resume" is answered in the vocabulary the boundary needs. + resumed, resumeErr := oauthClient.ResumeSession(ctx, did, storedSessionID) + if resumeErr != nil { + return nil, fmt.Errorf("resuming the stored session of %s: %w: %w", + authorDID, ErrNoAuthorCredentials, resumeErr) + } + if resumed == nil || resumed.Data == nil { + return nil, fmt.Errorf("resuming the stored session of %s: the store returned nothing: %w", + authorDID, ErrNoAuthorCredentials) + } + session = resumed.Data + } else if session.AccountDID != did { + // The session is what the record will actually be signed by, so a + // request naming a different author is refused here even though + // CreatePost has already compared the two — this is the last point + // at which an author-supplied DID could still reach a repo. + return nil, fmt.Errorf("opening the repository of %s under the session of %s: %w", + authorDID, session.AccountDID, ErrNotAuthorized) + } + + client, err := pds.NewFromOAuthSession(ctx, oauthClient, session) + if err != nil { + return nil, fmt.Errorf("opening the repository of %s: %w", authorDID, err) + } + + // The write path needs the guarded put and the commit rev, and neither + // is on the base Client. Asserted rather than assumed: a transport that + // lost either would otherwise fail at the first post, after admission + // has already spent the author's quota slot. + repo, ok := client.(AuthorRepo) + if !ok { + return nil, fmt.Errorf("opening the repository of %s: the PDS client does not implement the "+ + "author-repo write surface (guarded put + commit rev)", authorDID) + } + return repo, nil + } +} diff --git a/internal/core/posts/engine.go b/internal/core/posts/engine.go index 10a71b6..73f0dac 100644 --- a/internal/core/posts/engine.go +++ b/internal/core/posts/engine.go @@ -164,10 +164,49 @@ func NewAcceptanceEngine( // out wrapped so the caller can tell it from a genuine acceptance failure, which // leaves the post pending too but is worth alerting on. // -// RED STUB (task 6): pinned by service_writeflip_test.go. +// THE DECIDER IS NOT CONSULTED, and that is a correctness requirement rather +// than a saving. CreatePost has already run the very same policy — admitPost — +// over this submission, and the production decider works from the AppView's +// INDEX of the post, which on this path does not exist yet: the record was +// committed moments ago and its firehose copy has not arrived. A fast path that +// re-decided would therefore refuse every post it was handed, for the reason +// that it could not find it. func (e *AcceptanceEngine) AcceptSubmission(ctx context.Context, communityDID, postURI, postCID string) (EngineOutcome, error) { - _, _, _, _ = ctx, communityDID, postURI, postCID - return EngineDeferred, nil + if postCID == "" { + // An acceptance's subject is a strongRef, and a strongRef without a CID + // pins nothing — which is the one guarantee an acceptance exists to make. + return EngineDeferred, fmt.Errorf("accepting %s in %s: no content CID was supplied", postURI, communityDID) + } + + row, err := e.admissions.Get(ctx, communityDID, postURI) + if err != nil { + return EngineDeferred, fmt.Errorf("reading the admission of %s in %s: %w", postURI, communityDID, err) + } + if row == nil { + // The caller seeds the row before asking, so its absence is a caller bug + // rather than delivery skew, and inventing one here would let this path + // accept a post no author-repo observation was ever recorded for. + return EngineDeferred, fmt.Errorf("reading the admission of %s in %s: %w", postURI, communityDID, ErrNotFound) + } + + // THE ROW MUST STILL BE OWED A DECISION. Anything else — accepted, rejected, + // removed, pending_reacceptance — means something has already happened to + // this subject that this pass knows nothing about, and re-deciding a settled + // row is how a removal gets laundered back into a feed. + if row.Status != AdmissionStatusPending { + return EngineDeferred, nil + } + + // AND IT MUST STILL HOLD THE VERSION THE CALLER WROTE. Between the write and + // this call the firehose may already have delivered an EDIT, and an + // acceptance pinning the version the author has replaced is an acceptance of + // content nobody is reading. The engine's ordinary pass will settle whatever + // stands instead. + if row.EvaluatedCID == nil || *row.EvaluatedCID != postCID { + return EngineDeferred, nil + } + + return e.accept(ctx, communityDID, postURI, postCID) } // ProcessAdmission settles one (community, post) subject. diff --git a/internal/core/posts/postv2.go b/internal/core/posts/postv2.go index 86e916c..492eb53 100644 --- a/internal/core/posts/postv2.go +++ b/internal/core/posts/postv2.go @@ -2,9 +2,13 @@ package posts import ( "context" + "crypto/sha256" + "encoding/binary" + "strconv" "time" "github.com/bluesky-social/indigo/atproto/auth/oauth" + "github.com/bluesky-social/indigo/atproto/syntax" "Coves/internal/atproto/pds" ) @@ -63,6 +67,19 @@ type BridgedStats struct { // lexicon calls it immutable: retargeting a post means writing a NEW post // record, so an update that changes it is discarded entire by consumers. type PostV2Record struct { + // OriginalAuthor, FederatedFrom and Location are the bridge/federation + // surfaces a client may send alongside a post. They are `unknown` in the + // lexicon and `interface{}` here for the same reason: no bridge has fixed + // their shape yet, and declaring one now would be inventing a contract the + // first real bridge would contradict. + // + // They are carried on the record — rather than dropped as CreatePostRequest + // surfaces nothing reads — because UpdatePost round-trips the standing + // record through this struct: a field absent here is a field an edit erases + // from a post the first time its author fixes a typo. + OriginalAuthor interface{} `json:"originalAuthor,omitempty"` + FederatedFrom interface{} `json:"federatedFrom,omitempty"` + Location interface{} `json:"location,omitempty"` Title *string `json:"title,omitempty"` Content *string `json:"content,omitempty"` Embed map[string]interface{} `json:"embed,omitempty"` @@ -114,13 +131,48 @@ type PostV2Record struct { // actually submitted, and the clock ID carries further digest bits so two // submissions landing on the same microsecond still differ. func SubmissionRkey(communityDID, fingerprint string, bucket int64, dedupeWindow time.Duration) string { - // RED STUB (task 6): the contract is pinned in submission_rkey_test.go. - // It answers the empty string rather than panicking so the reds read as - // failed assertions naming the expected key, not as a stack trace. - _, _, _, _ = communityDID, fingerprint, bucket, dedupeWindow - return "" + digest := sha256.Sum256([]byte( + communityDID + submissionRkeyDelimiter + + fingerprint + submissionRkeyDelimiter + + strconv.FormatInt(bucket, 10))) + + // THE TIMESTAMP IS THE BUCKET'S START PLUS A DIGEST-DRAWN OFFSET INSIDE IT, + // so the derived time never leaves the window the submission actually + // landed in. Taking the offset modulo the window is what bounds it: 64 bits + // of digest used raw would name a moment in some arbitrary century, and the + // rkey's timestamp is not decoration — repo listings and feeds order by it. + windowMicros := int64(dedupeWindow / time.Microsecond) + var micros int64 + if windowMicros > 0 { + offset := binary.BigEndian.Uint64(digest[0:8]) % uint64(windowMicros) + micros = bucket*windowMicros + int64(offset) + } + // A non-positive window collapses to the epoch rather than dividing by + // zero, the same answer dedupeBucket gives the same misconfiguration. + // config.Validate refuses to start a process with one, so this is a wiring + // bug being kept off the write path, not a supported mode. + + // The clock ID carries FURTHER digest bits, so two submissions whose + // offsets collide on one microsecond still differ. NewTID masks it to the + // 10 bits a TID has room for; masking here too keeps this function's + // arithmetic the whole story. + clockID := uint(binary.BigEndian.Uint16(digest[8:10])) & 0x3FF + + // syntax.NewTID rather than a hand-rolled base32 encoding: the postv2 + // lexicon declares `key: tid`, and the one encoder guaranteed to agree with + // every ParseTID in the network is the one ParseTID ships beside. + return syntax.NewTID(micros, clockID).String() } +// submissionRkeyDelimiter separates the three parts of the rkey material. +// +// A newline, because neither a DID nor a hex fingerprint may contain one. With +// no delimiter at all — or one drawn from a charset the inputs share — a +// community DID ending in the delimiter and a fingerprint beginning with it +// would hash to the same bytes as some other pair, and one author's rkey could +// be forged from another submission's. +const submissionRkeyDelimiter = "\n" + // AuthorRepo is one author's PDS repository, narrowed to what the write path // does with it. // diff --git a/internal/core/posts/service.go b/internal/core/posts/service.go index fa78bb5..9c6c035 100644 --- a/internal/core/posts/service.go +++ b/internal/core/posts/service.go @@ -1,14 +1,11 @@ package posts import ( - "bytes" "context" "encoding/json" "errors" "fmt" - "io" "log" - "net/http" "strings" "time" @@ -88,33 +85,35 @@ func NewPostService( return s } -// CreatePost creates a new post in a community +// CreatePost writes a new post into the AUTHOR's repository (§4.2). // Flow: // 1. Validate input (and normalize embed/facet URIs) // 2. Verify the authenticated DID matches the request's author DID // 3. Classify the actor: trusted aggregator, registered aggregator, or user // 4. Admission: one decision over community existence, visibility, ban, // aggregator authorization, dedupe and the per-author quota (admitPost) -// 5. Ensure the community has fresh PDS credentials (token refresh) -// 6. Build the post record -// 7. Validate and enhance external embeds (thumb validation, unfurl, blobs) -// 8. Write to community's PDS repository -// 9. If aggregator: record post for rate limiting -// 10. Return URI/CID (AppView indexes asynchronously via Jetstream) +// 5. Open the AUTHOR's repository under the author's own credentials +// 6. Ensure the community has fresh PDS credentials (the blob uploads in +// step 8 still land in the community's repo until task 7 moves them) +// 7. Build the postv2 record +// 8. Validate and enhance external embeds (thumb validation, unfurl, blobs) +// 9. Create-only write at the deterministic rkey +// 10. If aggregator: record post for rate limiting +// 11. Seed the admission row and, for a community we host, settle it +// 12. Return URI/CID/status (AppView indexes asynchronously via Jetstream) // -// Admission runs BEFORE the token refresh, the blob uploads and the PDS write, +// Admission runs BEFORE the credentials, the blob uploads and the PDS write, // so a refused submission costs a few lookups rather than an upload — and, // more to the point, leaves no record in a community that refused it. Every -// failure AFTER admission (steps 5-8) must release the ledger reservation the +// failure AFTER admission (steps 5-9) must release the ledger reservation the // admission took, or the failure costs the author a quota slot and refuses // their retry as a duplicate. +// +// NOTHING AFTER THE RECORD COMMITS MAY FAIL THE REQUEST. The record is the +// author's and it exists; a failed acceptance, a failed row seed or a failed +// meter is degraded service, never data loss, and never a reason to withdraw +// someone else's record (§4.2). func (s *postService) CreatePost(ctx context.Context, session *oauth.ClientSessionData, req CreatePostRequest) (*CreatePostResponse, error) { - // RED STUB SEAM (task 6): the session is the author's credential and is - // consumed by the author-repo write the GREEN cycle installs below. It is - // accepted here so the contract compiles against the flipped signature - // while the body still write-forwards to the community's repo. - _ = session - // 1. Validate basic input (before DID checks to give clear validation errors) if err := s.validateCreateRequest(&req); err != nil { return nil, err @@ -190,11 +189,12 @@ func (s *postService) CreatePost(ctx context.Context, session *oauth.ClientSessi // page served at the time. The thumbnail URL rides along as submitted; the // client-typed community identifier and the per-attempt timestamp are // excluded inside submissionFingerprint (see its doc comment). + fingerprint := submissionFingerprint(postRecordFor(req, req.Community, ""), req.ThumbnailURL) decision, err := admitPost(ctx, s.admissionDeps(), AdmissionRequest{ Actor: actor, AuthorDID: req.AuthorDID, Community: req.Community, - Fingerprint: submissionFingerprint(postRecordFor(req, req.Community, ""), req.ThumbnailURL), + Fingerprint: fingerprint, }) if err != nil { return nil, err @@ -218,23 +218,38 @@ func (s *postService) CreatePost(ctx context.Context, session *oauth.ClientSessi } } - // 5. Ensure community has fresh PDS credentials (token refresh if needed) + // 5. Open the AUTHOR's repository, which is where the record goes now. + // + // AHEAD OF THE COMMUNITY'S CREDENTIALS, because this is the credential the + // post cannot be written without: a community whose token will not refresh + // costs a link preview, while an author we cannot authenticate as has no + // post at all. Ordering it second would report a community-side outage for + // a signed-out user. + authorRepo, err := s.openAuthorRepo(ctx, req.AuthorDID, session) + if err != nil { + releaseOnFailure() + return nil, err + } + + // 6. Ensure community has fresh PDS credentials (token refresh if needed). + // Still needed because the thumbnail blobs an external embed uploads are + // still written to the COMMUNITY's repo; task 7 moves them to the author's. community, err = s.communityService.EnsureFreshToken(ctx, community) if err != nil { releaseOnFailure() return nil, fmt.Errorf("failed to refresh community credentials: %w", err) } - // 6. Build post record for PDS + // 7. Build post record for PDS postRecord := postRecordFor(req, communityDID, time.Now().UTC().Format(time.RFC3339)) - // 7. Validate and enhance external embeds + // 8. Validate and enhance external embeds if err := s.enhanceExternalEmbed(ctx, &postRecord, req, community, actor == ActorTrustedAggregator); err != nil { releaseOnFailure() return nil, err } - // 8. Write to community's PDS repository + // 9. Write to the author's PDS repository. // // A failure here is the case the reservation was designed around: the row // went in before the write precisely so two concurrent identical submissions @@ -242,13 +257,14 @@ func (s *postService) CreatePost(ctx context.Context, session *oauth.ClientSessi // write which never happened owes the author their slot back. Without it, a // PDS hiccup would consume a quota slot AND refuse the retry as a duplicate, // turning a transient outage into a per-author lockout that outlives it. - uri, cid, err := s.createPostOnPDS(ctx, community, postRecord) + rkey := SubmissionRkey(communityDID, fingerprint, decision.DedupeBucket, s.admission.Limits.DedupeWindow) + uri, cid, err := createAuthorRecord(ctx, authorRepo, rkey, postV2From(postRecord)) if err != nil { releaseOnFailure() return nil, fmt.Errorf("failed to write post to PDS: %w", err) } - // 9. Record aggregator post for rate limiting (non-Kagi aggregators only) + // 10. Record aggregator post for rate limiting (non-Kagi aggregators only) // Kagi is exempted from rate limiting via env var (temporary) if isOtherAggregator && s.aggregatorService != nil { if recordErr := s.aggregatorService.RecordAggregatorPost(ctx, req.AuthorDID, communityDID, uri, cid); recordErr != nil { @@ -257,34 +273,319 @@ func (s *postService) CreatePost(ctx context.Context, session *oauth.ClientSessi } } - // 10. Return response (AppView will index via Jetstream consumer) - log.Printf("[POST-CREATE] Author: %s (trustedKagi=%v, otherAggregator=%v), Community: %s, URI: %s", - req.AuthorDID, isTrustedAggregator, isOtherAggregator, communityDID, uri) + // 11. Seed the admission row and, for a community this AppView hosts, + // settle it before answering. + status := s.settleSubmission(ctx, communityDID, uri, cid) + + // 12. Return response (AppView will index via Jetstream consumer) + log.Printf("[POST-CREATE] Author: %s (trustedKagi=%v, otherAggregator=%v), Community: %s, URI: %s, Status: %s", + req.AuthorDID, isTrustedAggregator, isOtherAggregator, communityDID, uri, status) return &CreatePostResponse{ - URI: uri, - CID: cid, + URI: uri, + CID: cid, + Status: status, }, nil } +// openAuthorRepo resolves the credentials the record is signed under. +// +// A service with no factory wired answers ErrNoAuthorCredentials rather than a +// nil-pointer panic: it is the same condition the production factory reports +// for an aggregator whose stored session is gone, and the boundary already +// knows how to answer it. +func (s *postService) openAuthorRepo(ctx context.Context, authorDID string, session *oauth.ClientSessionData) (AuthorRepo, error) { + // DEFENCE IN DEPTH over a boundary the PDS also enforces. The session's own + // DID is what the credentials will actually write as, so a request that + // named a different author has already failed CreatePost's spoofing check — + // this refuses the same thing one layer down, where a future caller that + // skipped that check would otherwise reach the factory with an + // author-supplied repo DID. + if session != nil && session.AccountDID.String() != authorDID { + log.Printf("[SECURITY] Author-repo session mismatch: session=%s, author=%s", + session.AccountDID.String(), authorDID) + return nil, ErrNotAuthorized + } + + if s.authorRepos == nil { + return nil, fmt.Errorf("opening the repository of %s: %w", authorDID, ErrNoAuthorCredentials) + } + + repo, err := s.authorRepos(ctx, authorDID, session) + if err != nil { + return nil, err + } + if repo == nil { + return nil, fmt.Errorf("opening the repository of %s: the factory returned no repo: %w", + authorDID, ErrNoAuthorCredentials) + } + return repo, nil +} + +// createAuthorRecord writes the post at its derived key, create-only, and +// reports the record that stands afterwards. +// +// THE GUARD IS THE IDEMPOTENCE. swapRecord "" means "there must be nothing +// here", so a retry of a submission whose first response was lost meets +// ErrSwapConflict instead of overwriting its own post — and is answered with +// the standing record's URI and CID rather than a fresh one. Re-putting an +// identical record would look harmless and be anything but: a new record CID +// dangles every strongRef built from the first response, and the second commit +// reaches every consumer as an EDIT, which drops an already-accepted post out +// of the community it was accepted into. +func createAuthorRecord(ctx context.Context, repo AuthorRepo, rkey string, record PostV2Record) (uri, cid string, err error) { + commit, err := repo.PutRecordWithCommit(ctx, PostV2Collection, rkey, record, "") + if err == nil { + return commit.URI, commit.CID, nil + } + + // TWO ANSWERS MEAN "IT IS ALREADY THERE", because the PDS orders its checks + // that way (verified against a live one): a put of bytes IDENTICAL to what + // stands is a no-op with no commit, and only a put of DIFFERENT bytes + // reaches the swap guard and comes back InvalidSwap. Both are the retry + // meeting its own first attempt, and both are answered the same way — by + // reporting the record that stands. + if !errors.Is(err, pds.ErrSwapConflict) && !errors.Is(err, pds.ErrNoCommit) { + return "", "", err + } + + standing, readErr := repo.GetRecord(ctx, PostV2Collection, rkey) + if readErr != nil { + // The record exists — that is what the swap conflict said — and we + // cannot name it. Reporting the read failure rather than the conflict + // keeps the cause the operator needs; the caller's retry converges on + // the same key and will find it. + return "", "", fmt.Errorf("the post already exists at %s but could not be read back: %w", rkey, readErr) + } + return standing.URI, standing.CID, nil +} + +// settleSubmission seeds the admission row for the post just written and, for a +// community this AppView hosts, settles it synchronously (§4.2 steps 4 and 5). +// +// IT CANNOT FAIL THE REQUEST, and every return path here reflects that. The +// author's record has committed; the worst outcome available is that the +// community still owes a decision, which the firehose engine will make when the +// post reaches it. A rollback would be wrong twice over: the record is the +// AUTHOR's to withdraw, and a rollback whose own delete failed would leave a +// post nobody has a row for. +func (s *postService) settleSubmission(ctx context.Context, communityDID, postURI, postCID string) string { + // Both or neither, by construction (WithSyncAcceptance). A service without + // them is one whose posts wait for the firehose, which is the pre-flip + // behaviour and a legitimate wiring. + if s.admissions == nil || s.acceptor == nil { + return PostStatusPending + } + + // THE ROW COMES FIRST. The engine settles a row, so there has to be one — + // and the URI it is seeded under must be byte-identical to the one the + // firehose consumer builds from the same commit, or the two index one post + // as two subjects and neither ever settles. + seeded, err := s.admissions.UpsertPending(ctx, UpsertPendingCommand{ + CommunityDID: communityDID, + PostURI: postURI, + EvaluatedCID: postCID, + }) + if err != nil { + log.Printf("[POST-CREATE] Warning: failed to seed the admission of %s in %s: %v", + postURI, communityDID, err) + return PostStatusPending + } + + // ALREADY ACCEPTED, which is what a retry of a settled post finds. The row + // is the AppView's answer about a post that exists, so reporting pending + // here would show the author a "waiting for the community" state over a post + // the community accepted — and a client that resubmitted would be handed the + // same URI again, forever. + if seeded.Admission != nil && seeded.Admission.Status == AdmissionStatusAccepted { + return PostStatusAccepted + } + + outcome, err := s.acceptor.AcceptSubmission(ctx, communityDID, postURI, postCID) + if err != nil { + if errors.Is(err, ErrCommunityNotHosted) { + // §4.2 step 5, and NOT a failure: this AppView has no authoritative + // view of that community's bans, visibility or quotas, so it must + // not decide for it. The community decides when the post reaches it. + log.Printf("[POST-CREATE] %s is hosted elsewhere; %s waits for its decision", communityDID, postURI) + return PostStatusPending + } + log.Printf("[POST-CREATE] Warning: the synchronous acceptance of %s in %s failed, leaving it "+ + "pending for the firehose engine to retry: %v", postURI, communityDID, err) + return PostStatusPending + } + + if outcome == EngineAccepted { + return PostStatusAccepted + } + return PostStatusPending +} + // UpdatePost edits a post in place in the author's repository. // -// RED STUB (task 6): see interfaces.go for the contract and -// service_writeflip_test.go for the pinned journey. +// THE READ IS NOT A CONVENIENCE. community and createdAt are taken from the +// STANDING RECORD rather than from the request — the first because the lexicon +// calls it immutable, so an edit that changed it would be discarded entire by +// every consumer, and the second because every feed orders by it, so +// re-stamping it would jump a three-year-old post corrected for a typo to the +// top of every sort. The CID that read returns is also the swap guard, which is +// what makes a concurrent edit a detected conflict rather than a silent +// overwrite of a change its author never saw. +// +// THE SUBMISSION LEDGER IS NOT TOUCHED. An edit is not a submission: it +// consumes no quota, and writing the edited content's fingerprint would let the +// ORIGINAL content be resubmitted inside its own dedupe window. func (s *postService) UpdatePost(ctx context.Context, session *oauth.ClientSessionData, req UpdatePostRequest) (*UpdatePostResponse, error) { - _, _, _ = ctx, session, req - return nil, ErrNotFound + if session == nil { + return nil, NewValidationError("session", "OAuth session required") + } + userDID := session.AccountDID.String() + + if req.URI == "" { + return nil, NewValidationError("uri", "post URI is required") + } + authority, rkey, err := parsePostURIParts(req.URI, "uri") + if err != nil { + return nil, err + } + if err := requireDIDAuthority(authority, "uri"); err != nil { + return nil, err + } + + // Only postv2 records are editable, and the reason is not squeamishness + // about the deprecated collection: a community.post record lives in the + // COMMUNITY's repo, so an edit would have to be signed by the community — + // the AppView asserting a change to someone else's words. Task 8 + // re-materializes those posts into their authors' repos; until then they + // are readable and deletable, not editable. + if collection := CollectionOfPostURI(req.URI); collection != PostV2Collection { + return nil, NewValidationError("uri", fmt.Sprintf( + "editing a post is only supported for %s URIs, got %s", PostV2Collection, collection)) + } + + // THE URI'S AUTHORITY IS THE OWNER, so authorization is decided here, + // before anything is fetched. The credentials the edit goes out on cannot + // reach another author's repo even if this check were wrong, which is + // exactly why it must be proven to exist rather than quietly removed. + if authority != userDID { + log.Printf("[SECURITY] Post update authorization failed: user=%s, authority=%s, uri=%s", + userDID, authority, req.URI) + return nil, ErrNotAuthorized + } + + repo, err := s.openAuthorRepo(ctx, userDID, session) + if err != nil { + return nil, err + } + + standing, err := repo.GetRecord(ctx, PostV2Collection, rkey) + if err != nil { + if errors.Is(err, pds.ErrNotFound) { + return nil, ErrNotFound + } + return nil, fmt.Errorf("failed to fetch the post being edited: %w", err) + } + + record, err := decodePostV2Record(standing.Value) + if err != nil { + return nil, fmt.Errorf("failed to decode the post being edited: %w", err) + } + + // REFUSED, NOT SILENTLY IGNORED. Both leave the record's community intact, + // but only one tells the client that the thing it asked for did not happen — + // and a client that believed it had moved a post would show its author a + // community the post is not in (§3.1: retargeting means a NEW post record). + if req.Community != "" && req.Community != record.Community { + return nil, NewValidationError("community", + "a post's community is immutable; submitting it elsewhere means creating a new post") + } + + applyPostV2Edit(&record, req) + + // The swap guard is the CID that was just read. An edit landing between the + // two is ErrConcurrentModification: the edit was composed against content + // that no longer stands, and re-applying it would erase a change its author + // never saw. Retrying is the client's decision, not the server's. + commit, err := repo.PutRecordWithCommit(ctx, PostV2Collection, rkey, record, standing.CID) + if err != nil { + if errors.Is(err, pds.ErrSwapConflict) { + return nil, fmt.Errorf("editing %s: %w", req.URI, ErrConcurrentModification) + } + return nil, fmt.Errorf("failed to write the edited post to PDS: %w", err) + } + + log.Printf("[POST-UPDATE] Author: %s, URI: %s, CID: %s", userDID, commit.URI, commit.CID) + + return &UpdatePostResponse{URI: commit.URI, CID: commit.CID}, nil +} + +// decodePostV2Record reads a standing record into the typed shape an edit is +// applied to. +// +// It round-trips through JSON rather than reading the map by hand so that the +// struct's tags stay the single description of the record: a field spelled one +// way in the writer and another in the reader is a field an edit silently +// erases, and postv2_record_test.go pins the struct against the lexicon for +// exactly that reason. +func decodePostV2Record(value map[string]any) (PostV2Record, error) { + encoded, err := json.Marshal(value) + if err != nil { + return PostV2Record{}, fmt.Errorf("re-encoding the standing record: %w", err) + } + var record PostV2Record + if err := json.Unmarshal(encoded, &record); err != nil { + return PostV2Record{}, fmt.Errorf("decoding the standing record: %w", err) + } + // A record read out of a repo may predate a $type being written, and the + // collection it was fetched from is the authority on what it is. + record.Type = PostV2Collection + return record, nil +} + +// applyPostV2Edit overwrites the mutable surfaces the request named, and only +// those. A nil field is "leave it alone" rather than "clear it": the update +// lexicon has no way to spell a deletion, so treating absence as removal would +// have a client editing a title silently drop the post's embed. +func applyPostV2Edit(record *PostV2Record, req UpdatePostRequest) { + if req.Title != nil { + record.Title = req.Title + } + if req.Content != nil { + record.Content = req.Content + } + if req.Facets != nil { + record.Facets = req.Facets + } + if req.Embed != nil { + record.Embed = req.Embed + } + if req.Labels != nil { + record.Labels = req.Labels + } + if req.Langs != nil { + record.Langs = req.Langs + } + if req.Tags != nil { + record.Tags = req.Tags + } } // postRecordFor builds the record a request describes, stamped with the given // community identifier and creation time. // // It is shared by the submission fingerprint and the record actually written, -// so that the thing dedupe hashes and the thing the community's repo receives -// cannot drift into describing different posts. The two callers differ in -// exactly the two arguments: the fingerprint is taken before the community -// identifier has been resolved and with no timestamp at all (createdAt is -// stamped per attempt, so including it would make every retry look new). +// so that the thing dedupe hashes and the thing that lands in a repo cannot +// drift into describing different posts. The two callers differ in exactly the +// two arguments: the fingerprint is taken before the community identifier has +// been resolved and with no timestamp at all (createdAt is stamped per attempt, +// so including it would make every retry look new). +// +// IT IS STILL THE DEPRECATED SHAPE, and the write converts at the boundary +// (postV2From). Two things are typed against it that a postv2 record cannot +// carry today: the fingerprint, whose stability across the flip keeps the +// ledger's live dedupe rows valid, and the embed-enhancement pipeline, which is +// unchanged content work. Retyping both is a task-6 cycle-2 obligation, not a +// behavioural one — the bytes that reach the PDS are postV2From's. func postRecordFor(req CreatePostRequest, community, createdAt string) PostRecord { return PostRecord{ Type: postCollection, @@ -302,6 +603,31 @@ func postRecordFor(req CreatePostRequest, community, createdAt string) PostRecor } } +// postV2From is the write boundary: the record the pipeline assembled, in the +// shape the AUTHOR's repository receives. +// +// THE AUTHOR FIELD IS DROPPED HERE, and this is the one place it happens. Under +// §3.1 authorship is the repository the record lives in — a claim a verifying +// relay or a DID-resolved fetch can check — so carrying the old self-asserted +// field alongside it would give consumers two answers to one question with no +// rule for which wins. The $type is re-stamped for the same reason: the +// collection a record is written to and the type it declares must agree. +func postV2From(record PostRecord) PostV2Record { + return PostV2Record{ + Type: PostV2Collection, + Community: record.Community, + CreatedAt: record.CreatedAt, + Title: record.Title, + Content: record.Content, + Facets: record.Facets, + Embed: record.Embed, + Labels: record.Labels, + OriginalAuthor: record.OriginalAuthor, + FederatedFrom: record.FederatedFrom, + Location: record.Location, + } +} + // enhanceExternalEmbed applies the external-embed handling that has to happen // against a live network: Bluesky URL conversion, client thumb validation, and // unfurl enrichment with its blob uploads. @@ -320,35 +646,9 @@ func (s *postService) enhanceExternalEmbed(ctx context.Context, postRecord *Post if external, extOk := postRecord.Embed["external"].(map[string]interface{}); extOk { // Check if this is a Bluesky post URL and convert to post embed if !s.tryConvertBlueskyURLToPostEmbed(ctx, external, postRecord) { - // Not a Bluesky URL or conversion failed - continue with normal external embed processing - // SECURITY: Validate thumb field (must be blob, not URL string) - // This validation happens BEFORE unfurl to catch client errors early - if existingThumb := external["thumb"]; existingThumb != nil { - if thumbStr, isString := existingThumb.(string); isString { - return NewValidationError("thumb", - fmt.Sprintf("thumb must be a blob reference (with $type, ref, mimeType, size), not URL string: %s", thumbStr)) - } - - // Validate blob structure if provided - if thumbMap, isMap := existingThumb.(map[string]interface{}); isMap { - // Check for $type field - if thumbType, ok := thumbMap["$type"].(string); !ok || thumbType != "blob" { - return NewValidationError("thumb", - fmt.Sprintf("thumb must have $type: blob (got: %v)", thumbType)) - } - // Check for required blob fields - if _, hasRef := thumbMap["ref"]; !hasRef { - return NewValidationError("thumb", "thumb blob missing required 'ref' field") - } - if _, hasMimeType := thumbMap["mimeType"]; !hasMimeType { - return NewValidationError("thumb", "thumb blob missing required 'mimeType' field") - } - log.Printf("[POST-CREATE] Client provided valid thumbnail blob") - } else { - return NewValidationError("thumb", - fmt.Sprintf("thumb must be a blob object, got: %T", existingThumb)) - } - } + // Not a Bluesky URL or conversion failed - continue with normal external embed processing. + // The thumb's shape was already checked in validateCreateRequest, + // which needs no network and so must not wait for one. // TRUSTED AGGREGATOR: Allow Kagi aggregator to provide thumbnail URLs directly // This bypasses unfurl for more accurate RSS-sourced thumbnails @@ -512,99 +812,67 @@ func (s *postService) validateCreateRequest(req *CreatePostRequest) error { return err } + // And the external embed's thumbnail, which validateEmbed deliberately does + // not look at: it checks the union's SHAPE, and the thumb is a blob whose + // parts a client gets wrong in four distinct ways worth naming separately. + if err := validateExternalThumb(req.Embed); err != nil { + return err + } + return nil } -// createPostOnPDS writes a post record to the community's PDS repository -// Uses com.atproto.repo.createRecord endpoint -func (s *postService) createPostOnPDS( - ctx context.Context, - community *communities.Community, - record PostRecord, -) (uri, cid string, err error) { - // Use community's PDS URL (not service default) for federated communities - // Each community can be hosted on a different PDS instance - pdsURL := community.PDSURL - if pdsURL == "" { - // Fallback to service default if community doesn't have a PDS URL - // (shouldn't happen in practice, but safe default) - pdsURL = s.pdsURL - } - - // Build PDS endpoint URL - endpoint := fmt.Sprintf("%s/xrpc/com.atproto.repo.createRecord", pdsURL) - - // Build request payload - // IMPORTANT: repo is set to community DID, not author DID - // This writes the post to the community's repository - payload := map[string]interface{}{ - "repo": community.DID, // Community's repository - "collection": "social.coves.community.post", // Collection type - "record": record, // The post record - // "rkey" omitted - PDS will auto-generate TID - } - - // Marshal payload - jsonData, err := json.Marshal(payload) - if err != nil { - return "", "", fmt.Errorf("failed to marshal post payload: %w", err) +// validateExternalThumb enforces that a social.coves.embed.external carries a +// real atProto blob reference in `thumb`, or nothing at all. +// +// Clients repeatedly send a URL STRING here, because that is what the rendered +// post looks like, and accepting one writes a record no other atProto +// implementation can read. Each of the four rejections names the part that is +// missing, because the message is the only thing telling a client which one it +// left out. +// +// IT RUNS AT VALIDATION TIME, before admission and before any credential is +// resolved, because it needs nothing but the request: a client mistake must be +// answerable with a 400 whether or not the author's repository can be opened, +// and it must not cost a ledger reservation to discover. +func validateExternalThumb(embed map[string]interface{}) error { + if embed == nil { + return nil } - - // Create HTTP request - req, err := http.NewRequestWithContext(ctx, "POST", endpoint, bytes.NewBuffer(jsonData)) - if err != nil { - return "", "", fmt.Errorf("failed to create PDS request: %w", err) + if embedType, ok := embed["$type"].(string); !ok || embedType != embedTypeExternal { + return nil } - - // Set headers (auth + content type) - req.Header.Set("Content-Type", "application/json") - req.Header.Set("Authorization", "Bearer "+community.PDSAccessToken) - - // Extended timeout for write operations (30 seconds) - client := &http.Client{ - Timeout: 30 * time.Second, + external, ok := embed["external"].(map[string]interface{}) + if !ok { + return nil } - - // Execute request - resp, err := client.Do(req) - if err != nil { - return "", "", fmt.Errorf("PDS request failed: %w", err) + thumb := external["thumb"] + if thumb == nil { + // The common case: a bare link whose thumbnail unfurl fills in later. + return nil } - defer func() { - if closeErr := resp.Body.Close(); closeErr != nil { - log.Printf("Warning: failed to close response body: %v", closeErr) - } - }() - // Read response body - body, err := io.ReadAll(resp.Body) - if err != nil { - return "", "", fmt.Errorf("failed to read PDS response: %w", err) + if thumbStr, isString := thumb.(string); isString { + return NewValidationError("thumb", + fmt.Sprintf("thumb must be a blob reference (with $type, ref, mimeType, size), not URL string: %s", thumbStr)) } - // Check for errors - if resp.StatusCode != http.StatusOK { - // Sanitize error body for logging (prevent sensitive data leakage) - bodyPreview := string(body) - if len(bodyPreview) > 200 { - bodyPreview = bodyPreview[:200] + "... (truncated)" - } - log.Printf("[POST-CREATE-ERROR] PDS Status: %d, Body: %s", resp.StatusCode, bodyPreview) - - // Return truncated error (defense in depth - handler will mask this further) - return "", "", fmt.Errorf("PDS returned error %d: %s", resp.StatusCode, bodyPreview) + thumbMap, isMap := thumb.(map[string]interface{}) + if !isMap { + return NewValidationError("thumb", + fmt.Sprintf("thumb must be a blob object, got: %T", thumb)) } - - // Parse response - var result struct { - URI string `json:"uri"` - CID string `json:"cid"` + if thumbType, ok := thumbMap["$type"].(string); !ok || thumbType != "blob" { + return NewValidationError("thumb", + fmt.Sprintf("thumb must have $type: blob (got: %v)", thumbType)) } - if err := json.Unmarshal(body, &result); err != nil { - return "", "", fmt.Errorf("failed to parse PDS response: %w", err) + if _, hasRef := thumbMap["ref"]; !hasRef { + return NewValidationError("thumb", "thumb blob missing required 'ref' field") } - - return result.URI, result.CID, nil + if _, hasMimeType := thumbMap["mimeType"]; !hasMimeType { + return NewValidationError("thumb", "thumb blob missing required 'mimeType' field") + } + return nil } // tryConvertBlueskyURLToPostEmbed attempts to convert a Bluesky URL in an external embed to a post embed. @@ -1026,16 +1294,23 @@ func communityCredentialFailure(operation, communityDID string, err error) error operation, communityDID, err) } -// DeletePost deletes a post from the community's PDS repository -// SECURITY: Only the post author can delete their own posts -// Flow: -// 1. Validate session and URI format -// 2. Extract community DID and rkey from URI -// 3. Fetch community from AppView -// 4. Ensure fresh PDS credentials -// 5. Fetch post record from community's PDS to get author field -// 6. SECURITY: Verify author matches session.AccountDID -// 7. Delete record from community's PDS using community credentials +// DeletePost removes a post record from the repository that holds it. +// SECURITY: Only the post author can delete their own posts. +// +// THE TWO COLLECTIONS AUTHORIZE DIFFERENTLY BECAUSE THEY LIVE IN DIFFERENT +// REPOS, and collapsing them would break one or the other: +// +// - A postv2 URI names the AUTHOR's repo, so the URI's authority IS the +// owner. The check is local, decided before anything is fetched, and the +// credentials the delete goes out on cannot reach another author's repo +// even if it were wrong. +// - A deprecated community.post URI names the COMMUNITY's repo, where the +// caller has no credentials at all. The delete goes out on the community's +// service token — which could delete anyone's post — so the record's +// `author` field has to be fetched and compared. That path survives until +// task 8 re-materializes those posts into their authors' repos; every one +// of them is standing in a community repo right now with a delete button +// that has to keep working. func (s *postService) DeletePost(ctx context.Context, session *oauth.ClientSessionData, req DeletePostRequest) error { // 1. Validate session if session == nil { @@ -1043,17 +1318,66 @@ func (s *postService) DeletePost(ctx context.Context, session *oauth.ClientSessi } userDID := session.AccountDID.String() - // 2. Validate URI format: at://community_did/social.coves.community.post/rkey + // 2. Validate URI shape, before anything reaches the network if err := s.validateDeleteRequest(&req); err != nil { return err } + authority, rkey, err := parsePostURIParts(req.URI, "uri") + if err != nil { + return err + } + if err := requireDIDAuthority(authority, "uri"); err != nil { + return err + } + // Defense-in-depth: verify rkey extraction is consistent with the utils helper. + if extractedRkey := utils.ExtractRKeyFromURI(req.URI); extractedRkey != rkey { + return NewValidationError("uri", "URI parsing inconsistency") + } + + // 3. Route on the collection. parsePostURIParts has already refused + // anything that is neither post collection. + if CollectionOfPostURI(req.URI) == PostV2Collection { + return s.deleteAuthorPost(ctx, session, userDID, authority, rkey, req.URI) + } + return s.deleteCommunityPost(ctx, userDID, authority, rkey, req.URI) +} + +// deleteAuthorPost removes a postv2 record from its author's own repository. +func (s *postService) deleteAuthorPost(ctx context.Context, session *oauth.ClientSessionData, userDID, authorDID, rkey, uri string) error { + // AUTHORIZATION IS THE URI'S AUTHORITY, decided before anything is fetched. + // + // Refused as UNAUTHORIZED rather than "not found", which is what the + // pre-flip path answered for an authority it could not look up: a 404 there + // would tell an attacker that the DID they aimed at is one this AppView has + // never seen, and would answer 404 to a probe that deserves 403. + if authorDID != userDID { + log.Printf("[SECURITY] Post delete authorization failed: user=%s, authority=%s, uri=%s", + userDID, authorDID, uri) + return ErrNotAuthorized + } - // 3. Extract community DID and rkey from URI - communityDID, rkey, err := s.parsePostURI(req.URI) + repo, err := s.openAuthorRepo(ctx, userDID, session) if err != nil { return err } + if err := repo.DeleteRecord(ctx, PostV2Collection, rkey); err != nil { + if errors.Is(err, pds.ErrNotFound) { + // Already deleted or never existed — the retried delete after a lost + // response succeeds. + log.Printf("[POST-DELETE] Post not found in the author's repo (already deleted?): %s", uri) + return nil + } + return fmt.Errorf("failed to delete post from PDS: %w", err) + } + + log.Printf("[POST-DELETE] Successfully deleted post: uri=%s, author=%s", uri, userDID) + return nil +} + +// deleteCommunityPost removes a pre-flip social.coves.community.post record +// from the community's repository, on the community's own credentials. +func (s *postService) deleteCommunityPost(ctx context.Context, userDID, communityDID, rkey, uri string) error { // 4. Fetch community from AppView community, err := s.communityService.GetByDID(ctx, communityDID) if err != nil { @@ -1076,11 +1400,11 @@ func (s *postService) DeletePost(ctx context.Context, session *oauth.ClientSessi } // 7. Fetch post record from PDS to verify author - record, err := pdsClient.GetRecord(ctx, "social.coves.community.post", rkey) + record, err := pdsClient.GetRecord(ctx, postCollection, rkey) if err != nil { if errors.Is(err, pds.ErrNotFound) { // Post already deleted or never existed - idempotent success - log.Printf("[POST-DELETE] Post not found on PDS (already deleted?): %s", req.URI) + log.Printf("[POST-DELETE] Post not found on PDS (already deleted?): %s", uri) return nil } if pds.IsAuthError(err) { @@ -1093,20 +1417,20 @@ func (s *postService) DeletePost(ctx context.Context, session *oauth.ClientSessi // The author field in the record must match the authenticated user's DID postAuthor, ok := record.Value["author"].(string) if !ok || postAuthor == "" { - return fmt.Errorf("post record missing author field: %s", req.URI) + return fmt.Errorf("post record missing author field: %s", uri) } if postAuthor != userDID { log.Printf("[SECURITY] Post delete authorization failed: user=%s, author=%s, uri=%s", - userDID, postAuthor, req.URI) + userDID, postAuthor, uri) return ErrNotAuthorized } // 9. Delete record from community's PDS - if err := pdsClient.DeleteRecord(ctx, "social.coves.community.post", rkey); err != nil { + if err := pdsClient.DeleteRecord(ctx, postCollection, rkey); err != nil { if errors.Is(err, pds.ErrNotFound) { // Already deleted - idempotent success - log.Printf("[POST-DELETE] Post already deleted from PDS: %s", req.URI) + log.Printf("[POST-DELETE] Post already deleted from PDS: %s", uri) return nil } if pds.IsAuthError(err) { @@ -1117,7 +1441,7 @@ func (s *postService) DeletePost(ctx context.Context, session *oauth.ClientSessi // 10. Log success (AppView will update via Jetstream consumer) log.Printf("[POST-DELETE] Successfully deleted post: uri=%s, author=%s, community=%s", - req.URI, userDID, communityDID) + uri, userDID, communityDID) return nil } @@ -1135,35 +1459,3 @@ func (s *postService) validateDeleteRequest(req *DeletePostRequest) error { return nil } - -// parsePostURI extracts community DID and rkey from a post URI -// Format: at://community_did/social.coves.community.post/rkey -// Returns community DID, rkey, and error -func (s *postService) parsePostURI(uri string) (communityDID string, rkey string, err error) { - // Structure + DID-authority validation is shared with the get path (single source of truth). - communityDID, rkey, err = parsePostURIParts(uri, "uri") - if err != nil { - return "", "", err - } - - // NARROWED to the community-repo collection, which the shared splitter - // deliberately is not. This path treats the authority as the COMMUNITY and - // goes on to open that repo and delete from it — so handed an author-repo - // postv2 URI it would authenticate as the AUTHOR's DID and try to delete a - // record there. The write path moves to postv2 in task 6; until it does, - // refusing is the only correct answer. - if collection := CollectionOfPostURI(uri); collection != postCollection { - return "", "", NewValidationError("uri", fmt.Sprintf( - "deleting a post is only supported for %s URIs, got %s", postCollection, collection)) - } - if err := requireDIDAuthority(communityDID, "uri"); err != nil { - return "", "", err - } - - // Defense-in-depth: verify rkey extraction is consistent with the utils helper. - if extractedRkey := utils.ExtractRKeyFromURI(uri); extractedRkey != rkey { - return "", "", NewValidationError("uri", "URI parsing inconsistency") - } - - return communityDID, rkey, nil -} -- 2.51.2