From 78dac39a1cc57c23e255fe3c0fd0cea9961bfdc2 Mon Sep 17 00:00:00 2001 From: Bretton Date: Sat, 8 Aug 2026 01:19:28 -0700 Subject: [PATCH] =?UTF-8?q?feat(ingestion):=20GREEN=20cycle=202=20?= =?UTF-8?q?=E2=80=94=20queue=20driver,=20decider,=20factory,=20sweep?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit All 42 cycle-2 reds pass at T0/T1. The three T2 contracts pass only with a one-line fix to the RED helper they share — see BLOCKER below; the file is committed exactly as RED wrote it. (a) ListPendingSubjects + its index Two INNER JOINs, both exclusions in SQL: posts (deleted_at IS NULL) and communities (pds_refresh_token_encrypted IS NOT NULL — credential presence, never hosted_by_did, which any repo can claim about itself). Migration 036 gains the partial index on (created_at) over the two undecided statuses and is renamed to match its widened scope; a 037 would have broken admission_repo_schema_test's asserted 36→35→34 chain. (b) CommunityRepoFactory + credential refresher Credential presence is the hosting test; absence is ErrCommunityNotHosted (permanent), while an unindexed community stays an ordinary error because it may simply not have arrived. Token renewed before the client is built, since a client is bound to the token it was constructed with. (c) AdmissionEngineDecider Post read first; tombstoned or absent returns UNDECIDED wrapping the new ErrSubjectGone rather than a code — the engine turns a code on a pending_reacceptance row into a REMOVAL record, and an author deleting their own post is not the community removing it. Actor classification checks the trusted set first (no lookup), and every uncertain path falls to ActorUser. TrustedAggregatorDIDs is now one helper shared with CreatePost so the two paths cannot disagree about who is privileged. (d) QueueDriver One goroutine, grouped by community (the partition a future worker pool must shard on), per-subject exponential backoff on deferral keyed by (community, post) — NOT by PendingSubject, whose time.Time field makes map equality fragile. Errors checked before outcomes, since a failing engine returns EngineDeferred alongside them. Snapshot is mutex-guarded: the health handler reads it while the job writes. (e) DeleteAcceptance + the tombstone sweep State-shaped single-op applyWrites, swapCommit-guarded, absence reported as a skip. The consumer sweeps only when the admission row says an acceptance stands, only when the tombstone actually applied (not on every redelivery), and never lets a failed sweep hold the local tombstone hostage. (f) Wiring: engine + driver behind ACCEPTANCE_QUEUE_INTERVAL (0 disables, and a nil driver is what omits the health block); acceptanceQueue in /health/consumers via an option, since an existing test pins the handler's arity. TWO PRODUCTION BUGS the T2 contracts caught, both in the READ path: - post.get refused every postv2 URI ("invalid collection in URI"), making author-owned posts unfetchable by the endpoint that hydrates every feed. parsePostURIParts now accepts both collections; the DELETE path narrows back to the community-repo one, because it uses the authority as a community DID and would otherwise try to delete from the AUTHOR's repo. - the synthesized `record` map hardcoded $type community.post and an `author` field, so every postv2 read handed back exactly the field whose removal makes authorship unforgeable. The shape now follows the URI. Co-Authored-By: Claude Fable 5 --- .env.ci | 15 ++ .env.dev.example | 27 +++ .env.prod.example | 27 +++ cmd/server/consumers.go | 8 +- cmd/server/health.go | 71 +++++- cmd/server/jobs.go | 53 +++++ cmd/server/main.go | 7 + cmd/server/routes.go | 10 +- cmd/server/wiring.go | 66 +++++- internal/atproto/jetstream/authorpost.go | 106 ++++++++- internal/atproto/jetstream/post_consumer.go | 27 ++- internal/config/config.go | 64 ++++++ internal/core/posts/community_repo_factory.go | 92 +++++++- internal/core/posts/community_writer.go | 109 ++++++++- internal/core/posts/decider.go | 143 +++++++++++- internal/core/posts/queue.go | 216 ++++++++++++++++-- internal/core/posts/service.go | 64 ++++-- ... 036_deleted_accounts_and_queue_index.sql} | 30 +++ internal/db/postgres/admission_queue_repo.go | 64 +++++- internal/db/postgres/post_repo.go | 27 ++- 20 files changed, 1160 insertions(+), 66 deletions(-) rename internal/db/migrations/{036_create_deleted_accounts.sql => 036_deleted_accounts_and_queue_index.sql} (67%) diff --git a/.env.ci b/.env.ci index 807e963..ade1ef7 100644 --- a/.env.ci +++ b/.env.ci @@ -167,3 +167,18 @@ IMAGE_PROXY_MAX_SOURCE_SIZE_MB=10 # Observability # ============================================================================= OTEL_ENABLED=false + +# ============================================================================= +# Acceptance queue +# ============================================================================= +# CI: absent from .env.dev, which takes the 1m default. The pipeline tier's +# per-probe budget is 45 seconds, so a minute-long cadence means a contract that +# ever needed the driver would time out before its first pass — and would do it +# as a flake rather than as a failure that names the driver. Two seconds keeps +# the pass inside every probe window. +# +# No contract currently DRIVES the engine — no community in the pipeline tier +# holds credentials in the AppView, so the backlog query correctly returns +# nothing — which makes this a liveness setting rather than a functional one: +# the driver runs, against the real database, on every CI boot. +ACCEPTANCE_QUEUE_INTERVAL=2s diff --git a/.env.dev.example b/.env.dev.example index b4e9567..62762c3 100644 --- a/.env.dev.example +++ b/.env.dev.example @@ -212,3 +212,30 @@ OTEL_ENABLED=false # clears itself; there is no sweeper, and the damage is bounded by this # window. # POST_SUBMISSIONS_DEDUPE_WINDOW=1h + +# ----------------------------------------------------------------------------- +# Acceptance queue (the engine's pull side) +# ----------------------------------------------------------------------------- +# A post claiming a community is not visible in it until the community writes an +# acceptance record (docs/PRD_AUTHOR_OWNED_POSTS.md section 5.6). The write path +# and the firehose consumer both push work at the engine that writes those, and +# neither can see a subject left undecided because a credential expired or a +# lookup blipped. This periodic pass is the only thing that reaches those, so its +# cadence is the worst case delay before a stranded post becomes visible. +# +# It is a no-op on an instance that hosts no communities: both the backlog query +# and the repo factory key on STORED PDS CREDENTIALS, which exist only for +# communities this AppView provisioned itself. +# +# How often the pass runs (default: 1m). Set to 0 to DISABLE the driver +# entirely, which is supported rather than a misconfiguration — an AppView that +# hosts no communities can accept nothing. Disabling it also removes the +# acceptanceQueue block from /health/consumers, so an absent driver and an idle +# one stay distinguishable. +# ACCEPTANCE_QUEUE_INTERVAL=1m +# +# How many subjects one pass takes (default: 50). Unlike the quotas above this +# one cannot fail open — the backlog query substitutes its own page size for a +# non-positive value and clamps an over-large one — so it is not validated at +# startup. +# ACCEPTANCE_QUEUE_BATCH_SIZE=50 diff --git a/.env.prod.example b/.env.prod.example index f986583..c0e1e10 100644 --- a/.env.prod.example +++ b/.env.prod.example @@ -407,6 +407,33 @@ OTEL_ENABLED=false # window. # POST_SUBMISSIONS_DEDUPE_WINDOW=1h +# ----------------------------------------------------------------------------- +# Acceptance queue (the engine's pull side) +# ----------------------------------------------------------------------------- +# A post claiming a community is not visible in it until the community writes an +# acceptance record (docs/PRD_AUTHOR_OWNED_POSTS.md section 5.6). The write path +# and the firehose consumer both push work at the engine that writes those, and +# neither can see a subject left undecided because a credential expired or a +# lookup blipped. This periodic pass is the only thing that reaches those, so its +# cadence is the worst case delay before a stranded post becomes visible. +# +# It is a no-op on an instance that hosts no communities: both the backlog query +# and the repo factory key on STORED PDS CREDENTIALS, which exist only for +# communities this AppView provisioned itself. +# +# How often the pass runs (default: 1m). Set to 0 to DISABLE the driver +# entirely, which is supported rather than a misconfiguration — an AppView that +# hosts no communities can accept nothing. Disabling it also removes the +# acceptanceQueue block from /health/consumers, so an absent driver and an idle +# one stay distinguishable. +# ACCEPTANCE_QUEUE_INTERVAL=1m +# +# How many subjects one pass takes (default: 50). Unlike the quotas above this +# one cannot fail open — the backlog query substitutes its own page size for a +# non-positive value and clamps an over-large one — so it is not validated at +# startup. +# ACCEPTANCE_QUEUE_BATCH_SIZE=50 + # ============================================================================= # Optional: Versioning # ============================================================================= diff --git a/cmd/server/consumers.go b/cmd/server/consumers.go index 668b742..2e29419 100644 --- a/cmd/server/consumers.go +++ b/cmd/server/consumers.go @@ -184,7 +184,13 @@ func (a *application) registerFeedConsumers() []feedConsumer { jetstream.WithPostIdentityResolver(a.identityResolver), jetstream.WithAdmissions(a.admissionRepo), jetstream.WithDeletedAccounts(postgresRepo.NewDeletedAccountRepository(a.db)), - jetstream.WithPostRecordFetcher(postFetcher)), + jetstream.WithPostRecordFetcher(postFetcher), + // The host-side half of an author's own deletion (§5.3): when the + // author tombstones a post this instance's community accepted, the + // acceptance in that community's repo is withdrawn. It refuses + // itself for every community this AppView does not host, which on + // most instances is all of them. + jetstream.WithAcceptanceCleanup(a.communityWriter)), }) // Aggregators: service declarations and authorization records, following diff --git a/cmd/server/health.go b/cmd/server/health.go index b4261e6..48698e4 100644 --- a/cmd/server/health.go +++ b/cmd/server/health.go @@ -66,11 +66,33 @@ type acceptanceQueueHealth struct { // buildAcceptanceQueueHealth renders one driver snapshot. // -// RED STUB (task 5, cycle 2). A separate pure function rather than another -// parameter on buildConsumerHealthResponse: the two have no shared logic, and -// widening that signature would touch every existing call site to say nothing. +// A separate pure function rather than another parameter on +// buildConsumerHealthResponse: the two have no shared logic, and widening that +// signature would touch every existing call site to say nothing. +// +// Both optional fields are OMITTED rather than zeroed when there is nothing to +// report, because a zero here would be read, and read wrongly: an age of 0 says +// "something arrived just now" when in fact nothing is waiting, and a +// zero-valued timestamp renders as the epoch, which looks like a driver that +// has been dead since 1970 rather than one that started a minute ago. func buildAcceptanceQueueHealth(snapshot posts.QueueSnapshot, now time.Time) acceptanceQueueHealth { - return acceptanceQueueHealth{} + queue := acceptanceQueueHealth{ + PendingBacklog: snapshot.PendingBacklog, + LastPassAt: snapshot.LastPassAt, + LastPassDeferred: snapshot.LastPassDeferred, + LastPassFailed: snapshot.LastPassFailed, + } + + // The AGE, not the timestamp. A backlog that is merely big is a busy + // instance; a backlog whose oldest entry keeps getting older is an engine + // that has stopped settling anything — and only the age says which is + // happening without the reader doing arithmetic against their own clock. + if snapshot.OldestPendingAt != nil { + age := int64(now.Sub(*snapshot.OldestPendingAt).Seconds()) + queue.OldestPendingAgeSeconds = &age + } + + return queue } // buildConsumerHealthResponse is the pure decision core of /health/consumers, @@ -127,7 +149,39 @@ func buildConsumerHealthResponse(statuses []jetstream.ConnectorStatus, backlogs // and the dead letter backlog per consumer. Responds 503 when any consumer // has been disconnected longer than consumerStalledThreshold (indexing is // stalled) so monitoring can alert on it. -func consumerHealthHandler(connectors []*jetstream.Connector, deadLetterQueue jetstream.DeadLetterQueue) http.HandlerFunc { +// acceptanceQueueReporter is the driver, narrowed to the one method health +// needs. Nil means no driver runs on this deployment, and the queue block is +// omitted entirely rather than reported as all-zero — an all-zero queue reads +// as a driver that is running and settling nothing, which is precisely the +// failure an operator is watching for. +type acceptanceQueueReporter interface { + Snapshot() posts.QueueSnapshot +} + +// consumerHealthOption adds a surface to /health/consumers that not every +// deployment has. +// +// An option rather than another parameter because the queue is genuinely +// optional — an AppView hosting no communities runs no driver — and because a +// nil third argument at every existing call site would say nothing while +// reading as an omission. +type consumerHealthOption func(*consumerHealthConfig) + +type consumerHealthConfig struct { + acceptanceQueue acceptanceQueueReporter +} + +// withAcceptanceQueue reports the acceptance driver alongside the consumers. +func withAcceptanceQueue(queue acceptanceQueueReporter) consumerHealthOption { + return func(c *consumerHealthConfig) { c.acceptanceQueue = queue } +} + +func consumerHealthHandler(connectors []*jetstream.Connector, deadLetterQueue jetstream.DeadLetterQueue, opts ...consumerHealthOption) http.HandlerFunc { + var cfg consumerHealthConfig + for _, opt := range opts { + opt(&cfg) + } + return func(w http.ResponseWriter, r *http.Request) { backlogs, err := deadLetterQueue.CountDeadLetters(r.Context()) backlogUnknown := err != nil @@ -142,7 +196,12 @@ func consumerHealthHandler(connectors []*jetstream.Connector, deadLetterQueue je statuses = append(statuses, connector.Status()) } - response, httpCode := buildConsumerHealthResponse(statuses, backlogs, backlogUnknown, time.Now()) + now := time.Now() + response, httpCode := buildConsumerHealthResponse(statuses, backlogs, backlogUnknown, now) + if cfg.acceptanceQueue != nil { + acceptance := buildAcceptanceQueueHealth(cfg.acceptanceQueue.Snapshot(), now) + response.AcceptanceQueue = &acceptance + } w.Header().Set("Content-Type", "application/json") w.WriteHeader(httpCode) diff --git a/cmd/server/jobs.go b/cmd/server/jobs.go index e79ef50..55abf43 100644 --- a/cmd/server/jobs.go +++ b/cmd/server/jobs.go @@ -6,6 +6,8 @@ import ( "log/slog" "sync" "time" + + "Coves/internal/core/posts" ) const ( @@ -130,6 +132,57 @@ type expiringTokenRefresher interface { RefreshExpiringTokens(ctx context.Context, expiryBuffer time.Duration) (int, []error) } +// acceptanceQueuePass is one walk of the acceptance engine's backlog. +// Declared as an interface so this job body is testable without an engine, a +// PDS or a database behind it. +type acceptanceQueuePass interface { + RunPass(ctx context.Context) (posts.PassReport, error) +} + +// startAcceptanceQueueJob walks the undecided admission backlog on an interval. +// +// It is the PULL half of admission (docs/PRD_AUTHOR_OWNED_POSTS.md §5.6). The +// synchronous write path and the firehose consumer both push work at the +// engine, and neither can see a subject that was left undecided because a +// credential expired or a lookup blipped — this pass is the only thing that +// eventually reaches those, so a deployment where it stops running is one where +// posts quietly stay invisible. +// +// Nothing here fails the process: a pass that could not read the backlog is +// logged and the next tick tries again, because the reason a backlog is +// unreadable is almost always the database being briefly unavailable, which the +// rest of the AppView is already reporting. +func startAcceptanceQueueJob(ctx context.Context, wg *sync.WaitGroup, queue acceptanceQueuePass, interval time.Duration) { + if queue == nil || interval <= 0 { + return + } + + runTicker(ctx, wg, "acceptance-queue", interval, func(ctx context.Context) { + report, err := queue.RunPass(ctx) + if err != nil { + if !errors.Is(err, context.Canceled) { + slog.Error("acceptance queue pass failed", "error", err) + } + return + } + + // Deferrals and failures are reported apart because they mean opposite + // things to whoever is deciding whether to page: a pass that defers + // everything is usually credentials and will clear, while a pass that + // FAILS everything is a bug. A quiet pass logs nothing, which is the + // common case on an instance that hosts no communities. + if report.Processed > 0 || report.Failed > 0 { + slog.Info("acceptance queue pass completed", + "listed", report.Listed, + "processed", report.Processed, + "settled", report.Settled, + "deferred", report.Deferred, + "failed", report.Failed, + ) + } + }) +} + // startAggregatorTokenRefreshJob proactively refreshes aggregator OAuth tokens // before they expire. // diff --git a/cmd/server/main.go b/cmd/server/main.go index a724579..6d88924 100644 --- a/cmd/server/main.go +++ b/cmd/server/main.go @@ -96,6 +96,13 @@ func run() error { startOAuthCleanupJob(backgroundCtx, &backgroundWG, sessionStore) startAggregatorTokenRefreshJob(backgroundCtx, &backgroundWG, app.apiKeyService) + // Nil when the driver is disabled, and passed as a typed nil would be a + // non-nil interface — so the guard is here rather than inside the job. + if app.acceptanceQueue != nil { + startAcceptanceQueueJob(backgroundCtx, &backgroundWG, + app.acceptanceQueue, cfg.Submissions.AcceptanceQueueInterval) + } + consumers, err := startConsumers(backgroundCtx, &backgroundWG, app) if err != nil { // Some connectors may already be running. Drain them under the same diff --git a/cmd/server/routes.go b/cmd/server/routes.go index 3bca2a3..5c5d21d 100644 --- a/cmd/server/routes.go +++ b/cmd/server/routes.go @@ -170,5 +170,13 @@ func registerWebRoutes(r chi.Router, app *application) { func registerHealthRoutes(r chi.Router, app *application, consumers *consumerSet) { r.Get("/health", livenessHandler) r.Get("/xrpc/_health", livenessHandler) - r.Get("/health/consumers", consumerHealthHandler(consumers.connectors, app.jetstreamState)) + // The option is added only when a driver actually runs. Passing a typed nil + // pointer instead would be non-nil to an interface comparison, and the + // response would carry an all-zero queue — which reads as a driver that is + // running and settling nothing, the exact failure an operator watches for. + var healthOptions []consumerHealthOption + if app.acceptanceQueue != nil { + healthOptions = append(healthOptions, withAcceptanceQueue(app.acceptanceQueue)) + } + r.Get("/health/consumers", consumerHealthHandler(consumers.connectors, app.jetstreamState, healthOptions...)) } diff --git a/cmd/server/wiring.go b/cmd/server/wiring.go index 51a75a5..5a1ecf8 100644 --- a/cmd/server/wiring.go +++ b/cmd/server/wiring.go @@ -109,7 +109,15 @@ type application struct { // because a status query needs the admissions store and nothing else, and // widening the write-path interface to reach it would make every test // double of posts.Service carry a method it has no opinion about. - postStatusService posts.StatusService + postStatusService posts.StatusService + // communityWriter publishes acceptances and removals into the repos of + // communities this AppView hosts. Shared by the acceptance engine, which + // writes verdicts, and the post consumer, which withdraws an acceptance + // when the author deletes the post it covers. + communityWriter posts.CommunityRecordWriter + // acceptanceQueue walks the undecided backlog. nil when the driver is + // disabled (ACCEPTANCE_QUEUE_INTERVAL=0). + acceptanceQueue *posts.QueueDriver voteService votes.Service commentService comments.Service userBlockService userblocks.Service @@ -363,6 +371,8 @@ 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 @@ -395,6 +405,60 @@ func (a *application) buildServices(ctx context.Context) error { return nil } +// buildAcceptanceEngine wires the §5.6 acceptance engine and the driver that +// feeds it. +// +// EVERYTHING HERE IS ABOUT WRITING INTO A COMMUNITY'S OWN REPO, which only that +// community's host can do — so on an instance that hosts nothing, all of it is +// a no-op that costs one query a minute: the repo factory refuses with +// ErrCommunityNotHosted and the backlog query returns nothing. Both are keyed on +// 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() { + repoFactory := posts.NewCommunityRepoFactory(a.communityService) + a.communityWriter = posts.NewCommunityRecordWriter(repoFactory, time.Now) + + decider := posts.NewAdmissionEngineDecider(posts.DeciderDeps{ + Posts: a.postRepo, + Communities: a.communityService, + Authorizer: a.aggregatorService, + Aggregators: a.aggregatorService, + Policy: posts.AdmissionPolicy{ + Ledger: postgresRepo.NewSubmissionLedger(a.db), + Bans: a.communityService, + Limits: posts.SubmissionLimits{ + MaxPerAuthorPerCommunity: a.cfg.Submissions.MaxPerAuthorPerCommunity, + Window: a.cfg.Submissions.Window, + DedupeWindow: a.cfg.Submissions.DedupeWindow, + }, + Now: time.Now, + }, + // Resolved ONCE, here, rather than read per decision — and through the + // same helper the write path uses, so the two cannot drift into + // disagreeing about who is privileged. + TrustedAggregatorDIDs: posts.TrustedAggregatorDIDs(), + }) + + engine := posts.NewAcceptanceEngine( + a.admissionRepo, decider, a.communityWriter, + posts.NewCommunityCredentialRefresher(a.communityService)) + + // Zero DISABLES the driver, and leaving the field nil is what makes + // /health/consumers omit the queue block entirely. An all-zero queue and an + // absent one mean different things: the first reads as a driver that is + // running and settling nothing, which is the exact failure an operator + // watches for. + 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 + } + + a.acceptanceQueue = posts.NewQueueDriver(a.admissionRepo, engine, time.Now, + posts.WithQueueBatchSize(a.cfg.Submissions.AcceptanceQueueBatchSize)) +} + // adminReportAlertOptions builds the operator-alert wiring for admin reports. // // Alerting is opt-in (TELEGRAM_ALERTS_ENABLED): most operators running their diff --git a/internal/atproto/jetstream/authorpost.go b/internal/atproto/jetstream/authorpost.go index 0a11c32..25487e2 100644 --- a/internal/atproto/jetstream/authorpost.go +++ b/internal/atproto/jetstream/authorpost.go @@ -47,7 +47,12 @@ import ( // under the same name would have it index authors as communities. A new NSID // makes a stale consumer ignore the records entirely, which is the correct // failure mode. -const PostV2Collection = "social.coves.community.postv2" +// +// It is an alias for the domain's constant rather than a second spelling of the +// string: the read path resolves post URIs against the same name, and two +// literals would let the indexer and the reader drift into disagreeing about +// what a post record is called. +const PostV2Collection = posts.PostV2Collection // DeletedAccountLookup reports whether a DID names an account this AppView was // asked to erase (migration 036, PRD rev 2.7). @@ -362,11 +367,108 @@ func (c *PostEventConsumer) handleAuthorPostEvent(ctx context.Context, event *Je case "create", "update": return c.upsertAuthorPost(ctx, authorDID, commit, event.TimeUS) case "delete": - return c.tombstoneRecord(ctx, recordURI(authorDID, PostV2Collection, commit.RKey), commit.Rev) + return c.tombstoneAuthorPost(ctx, authorDID, commit) } return nil } +// tombstoneAuthorPost soft-deletes the author's post and then withdraws any +// acceptance a community this AppView HOSTS still holds for it (§5.3). +// +// THE ORDER IS THE CONTRACT. The tombstone is the local truth and lands first: +// the author asked for their post to be gone, and a community PDS that cannot +// be reached must not keep this AppView serving it. The withdrawal is +// best-effort cleanup of a REMOTE repo, and it is deliberately not allowed to +// hold the deletion hostage. +func (c *PostEventConsumer) tombstoneAuthorPost(ctx context.Context, authorDID string, commit *CommitEvent) error { + uri := recordURI(authorDID, PostV2Collection, commit.RKey) + + // Read before the tombstone, because the community is what says WHOSE + // acceptance to withdraw and a delete event carries no record to read it + // from. The soft delete leaves the row in place, so this could equally run + // afterwards; doing it first keeps the sweep off the path when the post was + // never indexed here at all. + stored, indexed, err := c.loadStoredPost(ctx, uri) + if err != nil { + return err + } + + applied, err := c.tombstoneRecordIfRevWins(ctx, uri, commit.Rev) + if err != nil { + return err + } + if !applied || !indexed { + // A gate skip means this deletion was already applied — the sweep ran + // with it — so re-sweeping would put one authenticated PDS round trip + // behind every redelivery of every tombstone on the network. + return nil + } + + c.withdrawAcceptance(ctx, stored.communityDID, uri) + return nil +} + +// withdrawAcceptance asks the community's host to delete its acceptance of a +// post whose author has just deleted it. +// +// It is a NO-OP unless three things are true, and each exclusion removes a +// large class of pointless work: +// +// - a sweep is wired at all (nil on any build without the community writer); +// - THIS AppView holds the community's credentials — the acceptance lives in +// the community's repo and needs its keys, so on any instance that is not +// the community's home this is silently not our job, which is the common +// case; +// - an acceptance actually stands, per the AppView's own admission row. +// Consulting it is what keeps this from being a PDS round trip per delete +// event, since most posts a community sees it never accepted. +// +// A failure is LOGGED AND SWALLOWED. Returning it would dead-letter an event +// whose local half already committed, and the redrive would then be rejected by +// the rev gate — so the retry could never reach this code again anyway. The +// standing acceptance is left pointing at a deleted record until something +// revisits it, which nothing currently does: that gap is real, bounded to +// hosted communities, and named here rather than hidden. +func (c *PostEventConsumer) withdrawAcceptance(ctx context.Context, communityDID, postURI string) { + if c.acceptanceCleanup == nil || communityDID == "" { + return + } + + admission, err := c.admissions.Get(ctx, communityDID, postURI) + if err != nil { + if !errors.Is(err, posts.ErrNotFound) { + log.Printf("[ACCEPTANCE-SWEEP] Warning: could not read the admission of %s in %s: %v", + postURI, communityDID, err) + } + return + } + // The URI, not the status. `accepted` and `pending_reacceptance` both have a + // live acceptance record standing in the community's repo — the second + // merely pins content the author has since edited — and both must be + // withdrawn when the subject itself is deleted. + if admission == nil || admission.AcceptanceURI == nil { + return + } + + result, err := c.acceptanceCleanup.DeleteAcceptance(ctx, posts.CommunityAcceptanceDeleteCommand{ + CommunityDID: communityDID, + PostURI: postURI, + }) + switch { + case errors.Is(err, posts.ErrCommunityNotHosted): + // Not this instance's community. Expected, and not worth a warning: + // every AppView sees the deletions of every community it indexes. + log.Printf("debug: not withdrawing the acceptance of %s — %s is hosted elsewhere", postURI, communityDID) + case err != nil: + log.Printf("[ACCEPTANCE-SWEEP] Warning: could not withdraw the acceptance of %s in %s; "+ + "the record now cites a deleted post: %v", postURI, communityDID, err) + case result.Skipped: + log.Printf("debug: no acceptance of %s stood in %s to withdraw", postURI, communityDID) + default: + log.Printf("✓ Withdrew the acceptance of deleted post %s in %s", postURI, communityDID) + } +} + // canRecordAdmissions reports whether this consumer has somewhere to put a // decision. // diff --git a/internal/atproto/jetstream/post_consumer.go b/internal/atproto/jetstream/post_consumer.go index f9896fc..5ee6dbb 100644 --- a/internal/atproto/jetstream/post_consumer.go +++ b/internal/atproto/jetstream/post_consumer.go @@ -247,17 +247,26 @@ func (c *PostEventConsumer) createPost(ctx context.Context, repoDID string, comm // Soft-deletes the post in AppView database by setting deleted_at timestamp func (c *PostEventConsumer) deletePost(ctx context.Context, repoDID string, commit *CommitEvent) error { // Format: at://community_did/social.coves.community.post/rkey - return c.tombstoneRecord(ctx, fmt.Sprintf("at://%s/social.coves.community.post/%s", repoDID, commit.RKey), commit.Rev) + _, err := c.tombstoneRecordIfRevWins(ctx, + fmt.Sprintf("at://%s/social.coves.community.post/%s", repoDID, commit.RKey), commit.Rev) + return err } -// tombstoneRecord soft-deletes the post at uri under the rev gate. +// tombstoneRecordIfRevWins soft-deletes the post at uri under the rev gate, and +// reports whether the deletion APPLIED. // // SOFT, never hard, whichever repo the record lived in: the row is the rev // gate's tombstone, the comment thread's parent, and what moderation still // reads. It is shared by the community-repo and author-repo delete paths // because a deletion is the one operation where the two are identical — the // URI already says whose repo it was. -func (c *PostEventConsumer) tombstoneRecord(ctx context.Context, uri, rev string) error { +// +// The applied flag exists for the author-repo path's acceptance sweep, which +// must fire once per deletion rather than once per DELIVERY of it: the +// connector rewinds its cursor after every reconnect, so a tombstone that +// re-swept on each redelivery would put an authenticated PDS round trip behind +// every replayed event. +func (c *PostEventConsumer) tombstoneRecordIfRevWins(ctx context.Context, uri, rev string) (bool, error) { // REV GATE + soft delete in one transaction (the repo's SoftDelete is not // transaction-aware, and the delete's rev must be recorded atomically with // the tombstone: it is what rejects a stale cross-feed copy of the CREATE @@ -266,7 +275,7 @@ func (c *PostEventConsumer) tombstoneRecord(ctx context.Context, uri, rev string // record is rejected too. tx, err := c.db.BeginTx(ctx, nil) if err != nil { - return fmt.Errorf("failed to begin transaction: %w", err) + return false, fmt.Errorf("failed to begin transaction: %w", err) } defer func() { if rollbackErr := tx.Rollback(); rollbackErr != nil && rollbackErr != sql.ErrTxDone { @@ -276,11 +285,11 @@ func (c *PostEventConsumer) tombstoneRecord(ctx context.Context, uri, rev string won, err := tryAdvanceRecordRev(ctx, tx, uri, rev) if err != nil { - return err + return false, err } if !won { logSkippedStaleRev(ConsumerPosts, "delete", uri, rev) - return nil + return false, nil } // Same statement as postRepo.SoftDelete, inlined for transactionality. @@ -288,15 +297,15 @@ func (c *PostEventConsumer) tombstoneRecord(ctx context.Context, uri, rev string if _, err := tx.ExecContext(ctx, `UPDATE posts SET deleted_at = NOW() WHERE uri = $1 AND deleted_at IS NULL`, uri, ); err != nil { - return fmt.Errorf("failed to soft delete post: %w", err) + return false, fmt.Errorf("failed to soft delete post: %w", err) } if err := tx.Commit(); err != nil { - return fmt.Errorf("failed to commit post delete transaction: %w", err) + return false, fmt.Errorf("failed to commit post delete transaction: %w", err) } log.Printf("✓ Deleted post: %s", uri) - return nil + return true, nil } // updatePost handles post record update events from Jetstream. diff --git a/internal/config/config.go b/internal/config/config.go index 1ba1f98..eb5d0bb 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -355,6 +355,31 @@ type SubmissionsConfig struct { // repeat. It is separate from Window because the two answer different // questions: one bounds volume, the other catches retries. DedupeWindow time.Duration + + // AcceptanceQueueInterval is how often the acceptance engine's driver walks + // the undecided backlog (docs/PRD_AUTHOR_OWNED_POSTS.md §5.6). + // + // It is the PULL side of admission. The synchronous fast path and the + // firehose consumer both push work at the engine, and neither can see a + // subject that was left undecided because a credential expired or a lookup + // blipped — this pass is what eventually reaches those, so its cadence is + // the worst-case delay before a stranded post becomes visible. + // + // Zero DISABLES the driver, which is a supported deployment rather than a + // misconfiguration: an AppView that hosts no communities can accept nothing + // and has no backlog to walk. It is the one submission setting Validate does + // not require to be positive, for exactly that reason. + AcceptanceQueueInterval time.Duration + + // AcceptanceQueueBatchSize bounds how many subjects one pass lists. The + // backlog table grows with every submission the instance has ever seen, so + // an unbounded pass would hold a transaction open across all of it and then + // try to settle it inside a single cycle. + // + // Unlike the quotas above it is not validated, because it cannot fail open: + // the backlog query substitutes its own page size for a non-positive value + // and clamps an over-large one, so a bound exists whatever is set here. + AcceptanceQueueBatchSize int } // TokenEndpointEnabled reports whether the signup-token endpoint can operate. @@ -641,6 +666,19 @@ const ( defaultMaxSubmissionsPerCommunity = 10 defaultSubmissionWindow = time.Hour defaultSubmissionDedupeWindow = time.Hour + + // defaultAcceptanceQueueInterval is the backlog pass's cadence. A minute is + // the compromise the two failure modes point at from opposite directions: + // the pass is the only thing that reaches a subject nothing else will + // retry, so a long interval is a long wait for an author whose post got + // stuck, while a short one repeatedly scans a backlog that is usually empty. + defaultAcceptanceQueueInterval = time.Minute + + // defaultAcceptanceQueueBatch is how many subjects one pass takes. Small + // enough that a pass fits comfortably inside its own interval even when + // every subject needs a PDS round trip, since the backlog is drained across + // passes rather than in one. + defaultAcceptanceQueueBatch = 50 ) func (c *Config) loadSubmissions() error { @@ -657,10 +695,21 @@ func (c *Config) loadSubmissions() error { return err } + queueInterval, err := durationVar("ACCEPTANCE_QUEUE_INTERVAL", defaultAcceptanceQueueInterval) + if err != nil { + return err + } + queueBatch, err := intVar("ACCEPTANCE_QUEUE_BATCH_SIZE", defaultAcceptanceQueueBatch) + if err != nil { + return err + } + c.Submissions = SubmissionsConfig{ MaxPerAuthorPerCommunity: maxPerCommunity, Window: window, DedupeWindow: dedupeWindow, + AcceptanceQueueInterval: queueInterval, + AcceptanceQueueBatchSize: queueBatch, } return nil } @@ -771,6 +820,21 @@ func (c *Config) Validate() error { "it scopes the ledger's uniqueness bucket, and without a width every repost collides with the original forever", c.Submissions.DedupeWindow)) } + // The interval is checked for being NEGATIVE rather than non-positive, + // unlike the three above: zero is the documented way to disable the driver + // on an instance that hosts no communities, while a negative one is a + // time.Ticker panic waiting for the first boot after a typo. + if c.Submissions.AcceptanceQueueInterval < 0 { + problems = append(problems, fmt.Sprintf( + "ACCEPTANCE_QUEUE_INTERVAL cannot be negative (got %s); use 0 to disable the acceptance queue driver", + c.Submissions.AcceptanceQueueInterval)) + } + // The batch size is deliberately NOT validated, unlike every quota above. + // The quotas fail open when unset — an absent limit is no limit — so an + // omission there has to stop the boot. A batch size does not: the backlog + // query substitutes its own page for a non-positive limit and clamps an + // over-large one, so the bound exists whatever this value is, and the worst + // an omission costs is a different page size. if !c.IsDevEnv { switch { diff --git a/internal/core/posts/community_repo_factory.go b/internal/core/posts/community_repo_factory.go index 8d6b0fa..4a1eca8 100644 --- a/internal/core/posts/community_repo_factory.go +++ b/internal/core/posts/community_repo_factory.go @@ -3,12 +3,12 @@ package posts import ( "context" "errors" + "fmt" + "Coves/internal/atproto/pds" "Coves/internal/core/communities" ) -// RED STUB (task 5, cycle 2). Signatures only; the body is GREEN's. - // ErrCommunityNotHosted reports that this AppView does not hold the community's // PDS credentials, so it cannot write records into that community's repo. // @@ -28,6 +28,39 @@ type CommunityCredentialSource interface { EnsureFreshToken(ctx context.Context, community *communities.Community) (*communities.Community, error) } +// NewCommunityCredentialRefresher is the production CredentialRefresher: the +// engine's one forced renewal after a write comes back 401. +// +// KNOWN LIMITATION, recorded rather than hidden. EnsureFreshToken renews only a +// token that is within its expiry BUFFER, so a token the PDS has rejected for +// any other reason — revoked, invalidated by a password change, rotated out of +// band — is re-fetched unchanged and the retry fails identically. The subject +// then defers and the next pass tries again, which is correct but slower than +// it could be. Repairing it properly means a force-renew path on +// communities.Service, which is a wider change than this task, and the current +// behaviour is at worst the behaviour of having no refresher at all. +func NewCommunityCredentialRefresher(source CommunityCredentialSource) CredentialRefresher { + return credentialRefresher{source: source} +} + +type credentialRefresher struct { + source CommunityCredentialSource +} + +func (r credentialRefresher) RefreshCommunityCredentials(ctx context.Context, communityDID string) error { + community, err := r.source.GetByDID(ctx, communityDID) + if err != nil { + return fmt.Errorf("re-reading the credentials of %s: %w", communityDID, err) + } + if community == nil { + return fmt.Errorf("re-reading the credentials of %s: no such community is indexed", communityDID) + } + if _, err := r.source.EnsureFreshToken(ctx, community); err != nil { + return fmt.Errorf("renewing the credentials of %s: %w", communityDID, err) + } + return nil +} + // NewCommunityRepoFactory builds the production CommunityRepoFactory: a // credential lookup, a token refresh, and a PDS client bound to the community's // own repo. @@ -48,8 +81,59 @@ type CommunityCredentialSource interface { // refresh token — which happens exactly once, when it provisioned the account // through social.coves.community.create — or it does not. That is the honest // question, and it is the only one this factory asks. -func NewCommunityRepoFactory(communities CommunityCredentialSource) CommunityRepoFactory { +func NewCommunityRepoFactory(source CommunityCredentialSource) CommunityRepoFactory { return func(ctx context.Context, communityDID string) (CommunityRepo, error) { - return nil, nil + community, err := source.GetByDID(ctx, communityDID) + if err != nil { + if communities.IsNotFound(err) { + // Deliberately NOT ErrCommunityNotHosted. A community nobody has + // indexed may simply not have arrived yet — cross-repo delivery + // order is not guaranteed — and spelling an ordering artefact as + // a permanent skip would abandon a subject that resolves on its + // own once the profile lands. + return nil, fmt.Errorf("opening the repo of %s: no community with that DID is indexed: %w", + communityDID, err) + } + return nil, fmt.Errorf("opening the repo of %s: %w", communityDID, err) + } + if community == nil { + return nil, fmt.Errorf("opening the repo of %s: the community lookup returned nothing", communityDID) + } + + // THE HOSTING TEST, and the only one this factory is allowed to make. + // A stored refresh token exists exactly when this AppView provisioned + // the account itself; nothing a remote repo can publish creates one. + // See the note above for why hosted_by_did is not consulted. + if community.PDSRefreshToken == "" { + return nil, fmt.Errorf("opening the repo of %s: %w", communityDID, ErrCommunityNotHosted) + } + + // The token is renewed BEFORE the client is built rather than after a + // write fails, because a client is bound to the access token it was + // constructed with: refreshing afterwards would leave this caller + // holding the stale one. + fresh, err := source.EnsureFreshToken(ctx, community) + if err != nil { + return nil, fmt.Errorf("refreshing the credentials of %s: %w", communityDID, err) + } + if fresh == nil || fresh.PDSAccessToken == "" { + return nil, fmt.Errorf("refreshing the credentials of %s: no access token came back", communityDID) + } + + client, err := pds.NewFromAccessToken(fresh.PDSURL, fresh.DID, fresh.PDSAccessToken) + if err != nil { + return nil, fmt.Errorf("building a PDS client for %s: %w", communityDID, err) + } + + // The community-repo writers need the commit rev and applyWrites, and + // neither is on the base Client. Asserted rather than assumed: a + // transport that lost either would otherwise fail at the first + // moderation commit, which is after a verdict has already been reached. + repo, ok := client.(CommunityRepo) + if !ok { + return nil, fmt.Errorf("building a PDS client for %s: the client does not implement the community-repo "+ + "write surface (commit rev + applyWrites)", communityDID) + } + return repo, nil } } diff --git a/internal/core/posts/community_writer.go b/internal/core/posts/community_writer.go index 1b45ebc..d715540 100644 --- a/internal/core/posts/community_writer.go +++ b/internal/core/posts/community_writer.go @@ -21,6 +21,17 @@ import ( // because a re-fire must not mint a new record CID. const ( + // PostV2Collection is the AUTHOR-repo collection a post record lives in + // under author-owned posts (§3.1) — the successor to the deprecated + // community-repo social.coves.community.post. + // + // It lives beside the two community-repo collections because the three are + // one vocabulary: an acceptance's subject is a record in this collection, + // and the ingestion consumer re-exports this constant rather than declaring + // its own so that the reader and the writer cannot come to disagree about + // what a post record is called. + PostV2Collection = "social.coves.community.postv2" + // AcceptanceCollection is the community-repo collection holding a // community's attestation that it accepts a post. AcceptanceCollection = "social.coves.community.acceptance" @@ -353,9 +364,83 @@ func (w *communityRecordWriter) RepinAcceptance(ctx context.Context, cmd Communi return w.pinAcceptance(ctx, cmd, acceptanceMustExist) } -// DeleteAcceptance is a RED STUB (task 5, cycle 2); the body is GREEN's. +// DeleteAcceptance withdraws a standing acceptance and writes nothing in its +// place. +// +// STATE-SHAPED, like every other writer here and for the same reason task 4 +// recorded: the PDS has no tolerant delete, and answers a delete of a missing +// record with a 500. So absence is read first and reported as a skip. That is +// not an edge case — it is the COMMON one, because the sweep fires on every +// tombstone event and most posts a community sees were never accepted by it, +// and because the connector rewinds its cursor after every reconnect so each +// tombstone arrives at least twice. +// +// It is a batch of one rather than a putRecord-shaped call because applyWrites +// is where a delete can be guarded by swapCommit: the pre-read that found the +// acceptance is only true until somebody else writes, and the guard is what +// turns a concurrent restore into a detected conflict instead of a silently +// deleted fresh acceptance. func (w *communityRecordWriter) DeleteAcceptance(ctx context.Context, cmd CommunityAcceptanceDeleteCommand) (CommunityWriteResult, error) { - return CommunityWriteResult{}, nil + if err := validateAcceptanceDeleteCommand(cmd); err != nil { + return CommunityWriteResult{}, err + } + + repo, err := w.openRepo(ctx, cmd.CommunityDID) + if err != nil { + return CommunityWriteResult{}, err + } + + rkey := SubjectRkey(cmd.PostURI) + uri := recordURI(repo.DID(), AcceptanceCollection, rkey) + + for attempt := 0; ; attempt++ { + // The head is read before the record, for the same reason commitPair + // reads it first: a swapCommit read afterwards could be newer than the + // state the batch was shaped from, guarding the commit against a + // revision that already contains the change the shape assumed absent. + head, err := repo.GetLatestCommit(ctx) + if err != nil { + return CommunityWriteResult{}, fmt.Errorf("reading the head of %s: %w", cmd.CommunityDID, err) + } + + standing, err := readStandingRecord(ctx, repo, AcceptanceCollection, rkey) + if err != nil { + return CommunityWriteResult{}, err + } + if standing == nil { + // Nothing to withdraw. The head is still reported as the Rev, on the + // same catch-up reasoning the other writers' skips use: a row + // stranded by an earlier failed stamp can be caught up from it. + return CommunityWriteResult{URI: uri, RKey: rkey, Rev: head.Rev, Skipped: true}, nil + } + + result, err := repo.ApplyWrites(ctx, []pds.Write{{ + Op: pds.WriteOpDelete, + Collection: AcceptanceCollection, + RKey: rkey, + }}, head.CID) + if err == nil { + // No CID: a delete leaves no record to name. The Rev is the §5.2 + // watermark the firehose copy of this same deletion is compared + // against, so it is the one field that must be here. + return CommunityWriteResult{URI: uri, RKey: rkey, Rev: result.CommitRev}, nil + } + + // A lost swapCommit and a 500 are the same fact from two directions: the + // state this batch was shaped from is not the state the PDS is in. Both + // are answered by reading again, never by resending the same shape — + // here that matters most for the 500, which is what a delete of a record + // somebody else removed between the pre-read and the commit looks like. + staleShape := errors.Is(err, pds.ErrSwapConflict) || errors.Is(err, pds.ErrServerError) + if !staleShape || attempt >= swapRetryLimit { + return CommunityWriteResult{}, fmt.Errorf("withdrawing the acceptance of %s in %s: %w", + cmd.PostURI, cmd.CommunityDID, err) + } + if err := w.backoff(ctx, attempt); err != nil { + return CommunityWriteResult{}, fmt.Errorf("withdrawing the acceptance of %s in %s: %w", + cmd.PostURI, cmd.CommunityDID, err) + } + } } // pinAcceptance makes an acceptance of cmd.PostCID stand at the subject's rkey. @@ -776,6 +861,26 @@ func validateWriteCommand(cmd CommunityWriteCommand) error { return validateSubjectURI("acceptance write", cmd.PostURI) } +// validateAcceptanceDeleteCommand refuses a withdrawal that names no repo or no +// subject. +// +// The subject check is not symmetry with the other writers — it is the point. +// The rkey is a DIGEST of the subject URI, so a malformed or empty subject +// hashes to a perfectly well-formed key pointing at something else entirely, +// and a delete aimed at the wrong rkey in a community's own repo is a WRITE. A +// validation that only the create paths performed would leave the one operation +// that destroys data unchecked. +func validateAcceptanceDeleteCommand(cmd CommunityAcceptanceDeleteCommand) error { + switch { + case cmd.CommunityDID == "": + return fmt.Errorf("acceptance withdrawal: %w", NewValidationError("communityDID", "is required")) + case cmd.PostURI == "": + return fmt.Errorf("acceptance withdrawal: %w", NewValidationError("postURI", + "is required — the record key is derived from it, so an empty subject deletes a well-formed key belonging to nothing")) + } + return validateSubjectURI("acceptance withdrawal", cmd.PostURI) +} + // validateRemovalCommand refuses a removal with no reason code. `code` is // required by the lexicon and is what a client renders in #removedPost and what // the author is told. diff --git a/internal/core/posts/decider.go b/internal/core/posts/decider.go index fda07c6..6700d33 100644 --- a/internal/core/posts/decider.go +++ b/internal/core/posts/decider.go @@ -2,10 +2,13 @@ package posts import ( "context" + "errors" + "fmt" + "log" + "os" + "strings" ) -// RED STUB (task 5, cycle 2). Signatures only; the body is GREEN's. - // The production AdmissionDecider: the adapter that turns "decide about this // indexed post" into the AdmissionRequest evaluateAdmissionPolicy already // answers (docs/PRD_AUTHOR_OWNED_POSTS.md §5.6). @@ -35,6 +38,37 @@ import ( // redelivery and then refuse the redecision as a duplicate of the very post it // is redeciding. +// TrustedAggregatorDIDs reads the trusted-actor allowlist out of the process +// environment: TRUSTED_AGGREGATOR_DIDS, comma-separated, falling back to the +// legacy single-DID KAGI_AGGREGATOR_DID. +// +// It exists so the WRITE path and the ENGINE cannot disagree about who is +// trusted. Those two decide about the same post at different moments, and a +// trusted actor skips visibility, ban and authorization entirely — so two +// spellings of "read this variable" drifting apart would mean the same author +// is privileged on one path and not the other, which is the least debuggable +// shape a permission bug can take. +// +// It is called ONCE, at wiring time, and its result handed to whoever needs it. +// Reading the environment inside a decision would hide the most consequential +// input to a security decision from the place that makes it, and would make the +// trusted branch untestable alongside t.Parallel — Go's testing package refuses +// t.Setenv there. +func TrustedAggregatorDIDs() map[string]bool { + raw := os.Getenv("TRUSTED_AGGREGATOR_DIDS") + if raw == "" { + raw = os.Getenv("KAGI_AGGREGATOR_DID") + } + + trusted := map[string]bool{} + for _, did := range strings.Split(raw, ",") { + if did = strings.TrimSpace(did); did != "" { + trusted[did] = true + } + } + return trusted +} + // PostLookup reads the indexed post a decision is about. Satisfied by // Repository. type PostLookup interface { @@ -93,7 +127,110 @@ func NewAdmissionEngineDecider(deps DeciderDeps) *AdmissionEngineDecider { return &AdmissionEngineDecider{deps: deps} } +// ErrSubjectGone reports that the post an admission row names no longer stands: +// it was tombstoned by its author, or was never indexed at all. +// +// It travels as an ERROR rather than as a DecisionCode, and the choice is the +// difference between a correct record and a defamatory one. The engine turns a +// code on a `pending_reacceptance` row into a REMOVAL — a signed, portable +// moderation act published to the firehose — and an author deleting their own +// post is not the community removing it. There is no code that can be minted +// here without risking that, so the decider declines to decide instead: nothing +// is written, and the subject leaves the backlog on its own, because +// ListPendingSubjects excludes tombstoned and unindexed posts. +var ErrSubjectGone = errors.New("the post this admission names no longer stands") + // DecideAdmission implements AdmissionDecider. func (d *AdmissionEngineDecider) DecideAdmission(ctx context.Context, communityDID, postURI string) (AdmissionDecision, error) { - return AdmissionDecision{}, nil + // THE POST FIRST, and everything else after it. Two things come out of this + // lookup and both gate the policy: whether there is any content to judge, + // and who wrote it — and the author is what the actor class is derived + // from, so nothing about privilege can be decided before this returns. + post, err := d.deps.Posts.GetByURI(ctx, postURI) + switch { + case err != nil && IsNotFound(err): + // Absent. An admission row can legitimately exist with no post — an + // acceptance that arrived before its subject (§5.4) — so this is a + // normal state rather than a corruption, and there is simply nothing to + // judge yet. + return undecided(fmt.Errorf("deciding %s for %s: %w: it was never indexed", + postURI, communityDID, ErrSubjectGone)) + case err != nil: + // A lookup that FAILED is not a post that is absent. The two look + // identical from here and mean opposite things: the first clears, the + // second does not, and collapsing them would let a Postgres blip refuse + // somebody's post. + return undecided(fmt.Errorf("deciding %s for %s: reading the post: %w", postURI, communityDID, err)) + case post == nil: + return undecided(fmt.Errorf("deciding %s for %s: %w: it was never indexed", + postURI, communityDID, ErrSubjectGone)) + case post.DeletedAt != nil: + // Tombstoned. The driver already excludes these, but a post can be + // deleted between the listing and the decision, and admitting one would + // write an acceptance for content that no longer exists — which the + // host-side tombstone sweep would then delete, once per pass, forever. + return undecided(fmt.Errorf("deciding %s for %s: %w: its author deleted it", + postURI, communityDID, ErrSubjectGone)) + } + + return evaluateAdmissionPolicy(ctx, admissionDeps{ + communities: d.deps.Communities, + bans: d.deps.Policy.Bans, + aggregators: d.deps.Authorizer, + ledger: d.deps.Policy.Ledger, + limits: d.deps.Policy.Limits, + now: d.deps.Policy.Now, + }, AdmissionRequest{ + Actor: d.classify(ctx, post.AuthorDID), + AuthorDID: post.AuthorDID, + // The community DID, which resolves to itself. The engine's input is an + // admission row, and its key is already the resolved DID — there is no + // client-typed handle anywhere on this path to resolve. + Community: communityDID, + // EMPTY, and it must stay empty. Fingerprint is the dedupe key, read + // only by reserveSubmission, and this path deliberately never reserves: + // the engine re-decides posts that already exist, so a ledger row here + // would charge an author's quota for a firehose redelivery and then + // refuse the redecision as a duplicate of the very post it is + // redeciding. + Fingerprint: "", + }) +} + +// classify decides what class of actor the author is. +// +// EVERY UNCERTAIN PATH FALLS TO ActorUser, the stricter class. A trusted +// aggregator skips visibility, ban and authorization entirely, so resolving a +// failed lookup UPWARD would hand the widest privileges in the system to +// whoever managed to make the lookup fail. Guessing downward costs an +// aggregator some refused posts until the lookup recovers — and CreatePost +// already made exactly this choice (service.go step 3), so the engine agreeing +// with it is also what keeps the write path and the ingestion path from +// disagreeing about who someone is. +func (d *AdmissionEngineDecider) classify(ctx context.Context, authorDID string) ActorClass { + // The trusted set is checked FIRST, which is both the cheaper path and the + // only one that costs nothing: it is an in-memory set resolved at + // construction, so a trusted actor never pays for a database lookup to + // learn what the process already knew. + if d.deps.TrustedAggregatorDIDs[authorDID] { + return ActorTrustedAggregator + } + + // With no aggregator collaborators wired — a deployment with no aggregator + // support at all — nobody can be classified as one, which is the strict + // answer rather than a degraded one. + if d.deps.Aggregators == nil || d.deps.Authorizer == nil { + return ActorUser + } + + registered, err := d.deps.Aggregators.IsAggregator(ctx, authorDID) + if err != nil { + log.Printf("[ADMISSION-DECIDER] Warning: classifying %s fell back to the user class, IsAggregator failed: %v", + authorDID, err) + return ActorUser + } + if registered { + return ActorRegisteredAggregator + } + return ActorUser } diff --git a/internal/core/posts/queue.go b/internal/core/posts/queue.go index f0c539a..7c116ba 100644 --- a/internal/core/posts/queue.go +++ b/internal/core/posts/queue.go @@ -2,11 +2,12 @@ package posts import ( "context" + "fmt" + "log" + "sync" "time" ) -// RED STUB (task 5, cycle 2). Signatures only; every body returns zero values. - // The acceptance engine's driver: the thing that decides WHEN the engine runs // and on what (docs/PRD_AUTHOR_OWNED_POSTS.md §5.6, §8). // @@ -129,22 +130,67 @@ type QueueDriver struct { backoffBase time.Duration backoffMax time.Duration - // deferrals holds the per-subject retry-not-before times. In-memory on - // purpose: it is a politeness hint, not state anything is allowed to depend - // on, so a restart that forgets it costs one extra attempt per subject and - // nothing else. - deferrals map[PendingSubject]time.Time + // deferrals holds the per-subject backoff state. In-memory on purpose: it is + // a politeness hint, not state anything is allowed to depend on, so a + // restart that forgets it costs one extra attempt per subject and nothing + // else. + // + // It is keyed by (community, post) rather than by the whole PendingSubject + // because the third field is a time.Time read back from Postgres, and Go + // compares those by wall clock AND monotonic reading AND location. Two + // reads of one unchanged row can therefore produce values that are equal to + // a human and distinct to a map, which would silently defeat the backoff. + deferrals map[subjectKey]deferral + // mu guards deferrals and snapshot. Snapshot is read by the health handler + // on an HTTP goroutine while RunPass is writing on the job goroutine, so + // this is a genuine race rather than a defensive one. + mu sync.Mutex snapshot QueueSnapshot } +// subjectKey identifies one subject by the two fields that actually name it. +type subjectKey struct { + communityDID string + postURI string +} + +func keyOf(subject PendingSubject) subjectKey { + return subjectKey{communityDID: subject.CommunityDID, postURI: subject.PostURI} +} + +// deferral is how long one subject is held back, and until when. +// +// The delay is carried alongside the deadline so it can GROW: a subject that +// defers repeatedly is one whose community is wedged, and re-offering it on a +// fixed interval would keep a steady trickle of doomed requests pointed at a +// PDS that is already failing. +type deferral struct { + until time.Time + delay time.Duration +} + +// Default queue bounds, applied when the corresponding option is not given. +// +// The batch bound is not a nicety: this query runs on a timer against a table +// that grows with every submission the instance has ever seen, so a driver +// built without one must still not ask for the whole backlog. +const ( + defaultQueueBatchSize = 100 + defaultQueueBackoffBase = time.Minute + defaultQueueBackoffMax = 15 * time.Minute +) + // NewQueueDriver wires the driver. func NewQueueDriver(subjects PendingSubjectLister, engine AdmissionProcessor, now Clock, opts ...QueueDriverOption) *QueueDriver { d := &QueueDriver{ - subjects: subjects, - engine: engine, - now: now, - deferrals: make(map[PendingSubject]time.Time), + subjects: subjects, + engine: engine, + now: now, + batchSize: defaultQueueBatchSize, + backoffBase: defaultQueueBackoffBase, + backoffMax: defaultQueueBackoffMax, + deferrals: make(map[subjectKey]deferral), } for _, opt := range opts { opt(d) @@ -159,10 +205,154 @@ func NewQueueDriver(subjects PendingSubjectLister, engine AdmissionProcessor, no // must not stop every other community's posts from being decided, which is what // an early return would do — and the row is still in the backlog next pass. func (d *QueueDriver) RunPass(ctx context.Context) (PassReport, error) { - return PassReport{}, nil + startedAt := d.now() + report := PassReport{StartedAt: startedAt} + + subjects, err := d.subjects.ListPendingSubjects(ctx, d.batchSize) + if err != nil { + // The one failure a pass has nothing to do about. Every other outcome + // below is per-subject and counted; this one means there is no work + // list to count against. + return report, fmt.Errorf("listing the acceptance backlog: %w", err) + } + report.Listed = len(subjects) + + for _, subject := range groupByCommunity(subjects) { + if d.heldBack(subject) { + continue + } + + outcome, err := d.engine.ProcessAdmission(ctx, subject.CommunityDID, subject.PostURI) + report.Processed++ + + // The ERROR is checked before the outcome, because a failing engine + // returns EngineDeferred alongside it and reading the outcome first + // would file every failure as a deferral — collapsing the two numbers an + // operator uses to decide whether this is credentials or a bug. + switch { + case err != nil: + report.Failed++ + log.Printf("[ACCEPTANCE-QUEUE] Warning: %s in %s could not be settled: %v", + subject.PostURI, subject.CommunityDID, err) + // NOT backed off. A failure is unexplained, so the driver has no + // basis for guessing how long to wait; the next pass re-lists it and + // the row is still there. Deferral is the engine SAYING "later". + case outcome == EngineDeferred: + report.Deferred++ + d.deferSubject(subject, startedAt) + default: + report.Settled++ + d.clearDeferral(subject) + } + } + + d.record(subjects, report, startedAt) + return report, nil } // Snapshot returns the driver's health surface as of the last completed pass. func (d *QueueDriver) Snapshot() QueueSnapshot { - return QueueSnapshot{} + d.mu.Lock() + defer d.mu.Unlock() + return d.snapshot +} + +// heldBack reports whether a subject's backoff has yet to elapse. +func (d *QueueDriver) heldBack(subject PendingSubject) bool { + d.mu.Lock() + defer d.mu.Unlock() + + held, ok := d.deferrals[keyOf(subject)] + return ok && d.now().Before(held.until) +} + +// deferSubject holds a subject back, doubling its wait each consecutive time. +func (d *QueueDriver) deferSubject(subject PendingSubject, at time.Time) { + d.mu.Lock() + defer d.mu.Unlock() + + key := keyOf(subject) + delay := d.backoffBase + if previous, ok := d.deferrals[key]; ok && previous.delay > 0 { + delay = previous.delay * 2 + } + if delay > d.backoffMax { + delay = d.backoffMax + } + d.deferrals[key] = deferral{until: at.Add(delay), delay: delay} +} + +// clearDeferral forgets a settled subject, so the map tracks the backlog rather +// than the history of everything the driver has ever seen. +func (d *QueueDriver) clearDeferral(subject PendingSubject) { + d.mu.Lock() + defer d.mu.Unlock() + delete(d.deferrals, keyOf(subject)) +} + +// record publishes the pass's health surface. +func (d *QueueDriver) record(subjects []PendingSubject, report PassReport, at time.Time) { + d.mu.Lock() + defer d.mu.Unlock() + + snapshot := QueueSnapshot{ + PendingBacklog: report.Listed, + LastPassDeferred: report.Deferred, + LastPassFailed: report.Failed, + } + // Taken as a MINIMUM rather than as subjects[0], even though the query + // orders by age. The oldest entry's age is the queue's only early warning, + // and deriving it from an ordering assumption would make it silently wrong + // the first time anything reorders the list. + for i, subject := range subjects { + if i == 0 || subject.CreatedAt.Before(*snapshot.OldestPendingAt) { + oldest := subject.CreatedAt + snapshot.OldestPendingAt = &oldest + } + } + passedAt := at + snapshot.LastPassAt = &passedAt + + d.snapshot = snapshot +} + +// groupByCommunity returns the subjects with each community's contiguous, and +// with duplicates dropped. +// +// GROUPING. swapCommit is repo-global, so two writers on one community's repo +// starve each other — task 4 recorded that before this driver existed. A single +// goroutine satisfies it today whatever the order; what the grouping protects is +// the NEXT version, where a worker pool has to shard on something, and the only +// safe partition is the community. Output that interleaved communities would be +// output with no partition in it. +// +// DEDUPLICATION. A subject appearing twice in one listing is what a racing edit +// produces, and §8's edit-debounce is exactly this rule: a post edited in a +// storm coalesces into one pending_reacceptance row, and a driver that re-decided +// it per occurrence would re-run the whole policy per keystroke. +// +// The order WITHIN a community is preserved, and so is the order BETWEEN them +// (first appearance wins), so the query's oldest-first discipline survives. +func groupByCommunity(subjects []PendingSubject) []PendingSubject { + var order []string + byCommunity := make(map[string][]PendingSubject) + seen := make(map[subjectKey]bool, len(subjects)) + + for _, subject := range subjects { + if seen[keyOf(subject)] { + continue + } + seen[keyOf(subject)] = true + + if _, known := byCommunity[subject.CommunityDID]; !known { + order = append(order, subject.CommunityDID) + } + byCommunity[subject.CommunityDID] = append(byCommunity[subject.CommunityDID], subject) + } + + grouped := make([]PendingSubject, 0, len(seen)) + for _, communityDID := range order { + grouped = append(grouped, byCommunity[communityDID]...) + } + return grouped } diff --git a/internal/core/posts/service.go b/internal/core/posts/service.go index 8fe70a7..4858bf1 100644 --- a/internal/core/posts/service.go +++ b/internal/core/posts/service.go @@ -9,7 +9,6 @@ import ( "io" "log" "net/http" - "os" "strings" "time" @@ -135,22 +134,11 @@ func (s *postService) CreatePost(ctx context.Context, req CreatePostRequest) (*C return nil, fmt.Errorf("authenticated DID does not match author DID") } - // 3. Determine actor type: trusted aggregator, other aggregator, or regular user - // Check against comma-separated list of trusted aggregator DIDs - trustedDIDs := os.Getenv("TRUSTED_AGGREGATOR_DIDS") - if trustedDIDs == "" { - // Fallback to legacy single DID env var - trustedDIDs = os.Getenv("KAGI_AGGREGATOR_DID") - } - isTrustedAggregator := false - if trustedDIDs != "" { - for _, did := range strings.Split(trustedDIDs, ",") { - if strings.TrimSpace(did) == req.AuthorDID { - isTrustedAggregator = true - break - } - } - } + // 3. Determine actor type: trusted aggregator, other aggregator, or regular user. + // The allowlist is read through the shared helper rather than inline, so + // this path and the acceptance engine's decider cannot drift into disagreeing + // about who is trusted — see TrustedAggregatorDIDs. + isTrustedAggregator := TrustedAggregatorDIDs()[req.AuthorDID] // Check if this is a non-trusted aggregator (requires database lookup) var isOtherAggregator bool @@ -858,8 +846,19 @@ func parsePostURIParts(uri, field string) (authority string, rkey string, err er if authority == "" { return "", "", NewValidationError(field, "invalid post URI: missing authority") } - if collection != postCollection { - return "", "", NewValidationError(field, fmt.Sprintf("invalid collection in URI: expected %s, got %s", postCollection, collection)) + // EITHER post collection is a well-formed post URI. A post now lives in the + // author's repo under social.coves.community.postv2 (§3.1), while every post + // written before the flip is still at the deprecated community-repo NSID, and + // a reader has to be able to name both — refusing postv2 here made the new + // records unfetchable by the endpoint that hydrates every feed. + // + // What the authority MEANS differs between them — the community for the old + // collection, the author for the new — so a caller that goes on to use it as + // one or the other must narrow this itself. parsePostURI does, because it + // writes to the repo the authority names. + if collection != postCollection && collection != PostV2Collection { + return "", "", NewValidationError(field, fmt.Sprintf("invalid collection in URI: expected %s or %s, got %s", + postCollection, PostV2Collection, collection)) } if rkey == "" { return "", "", NewValidationError(field, "invalid post URI: missing rkey") @@ -867,6 +866,22 @@ func parsePostURIParts(uri, field string) (authority string, rkey string, err er return authority, rkey, nil } +// CollectionOfPostURI returns the collection segment of an at:// record URI, or +// "" when the URI is not shaped like one. +// +// It is exported because the two post collections are now indexed into one +// table, so every layer that renders or narrows a post has to ask the same +// question of the same URI — the repository decides which record shape to +// build from it, and this path decides which writes it will accept. A second +// spelling of the split is a second place for the two to disagree. +func CollectionOfPostURI(uri string) string { + parts := strings.Split(strings.TrimPrefix(uri, "at://"), "/") + if len(parts) != 3 { + return "" + } + return parts[1] +} + // requireDIDAuthority enforces that a parsed post-URI authority is a DID (not a handle). // Handles are mutable, so a handle-based URI would break after a community rename, or // mis-resolve if the handle is later reassigned. field names the request parameter for errors. @@ -1107,6 +1122,17 @@ func (s *postService) parsePostURI(uri string) (communityDID string, rkey string 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 } diff --git a/internal/db/migrations/036_create_deleted_accounts.sql b/internal/db/migrations/036_deleted_accounts_and_queue_index.sql similarity index 67% rename from internal/db/migrations/036_create_deleted_accounts.sql rename to internal/db/migrations/036_deleted_accounts_and_queue_index.sql index 153d4c4..fdc5f66 100644 --- a/internal/db/migrations/036_create_deleted_accounts.sql +++ b/internal/db/migrations/036_deleted_accounts_and_queue_index.sql @@ -56,5 +56,35 @@ COMMENT ON TABLE deleted_accounts IS 'Erasure markers: DIDs deleted on purpose, COMMENT ON COLUMN deleted_accounts.deleted_at IS 'When the deletion happened; read by retention and audit, never by the ingestion gate itself'; COMMENT ON COLUMN deleted_accounts.deleted_rev IS 'Repo revision the erasure was observed at, when one is known; NULL for AppView-initiated deletions'; +-- The acceptance engine's backlog scan (PRD_AUTHOR_OWNED_POSTS.md §5.6). +-- +-- It rides along in this migration rather than getting one of its own because +-- both halves are the same task's enablers, and an index is not worth a +-- version of its own on a table one migration away. +-- +-- WHY THE EXISTING INDEX CANNOT SERVE IT. Migration 034's +-- (community_did, status, created_at) index leads with the community, which is +-- exactly right for a moderator paging ONE community's queue and useless for +-- the driver's question — "what is undecided ANYWHERE" — which has no community +-- in hand and would scan every community's rows through it. +-- +-- WHY PARTIAL. The undecided rows are a small and roughly constant slice of a +-- table that grows with every submission the instance has ever seen: a settled +-- admission stays forever, and pending ones drain. Restricting the index to the +-- two undecided statuses keeps it proportional to the BACKLOG rather than to +-- history, so a pass that runs on a timer does not get more expensive every day +-- the instance stays up. +-- +-- Nothing FAILS without this index, which is precisely the danger: the query +-- keeps returning correct answers and quietly costs more every week, and the +-- symptom arrives as general database pressure with nothing pointing here. +CREATE INDEX idx_admissions_pending_queue + ON community_post_admissions (created_at) + WHERE status IN ('pending', 'pending_reacceptance'); + +COMMENT ON INDEX idx_admissions_pending_queue IS 'CRITICAL: the acceptance engine cross-community backlog scan, oldest first (PRD_AUTHOR_OWNED_POSTS 5.6)'; + -- +goose Down +DROP INDEX IF EXISTS idx_admissions_pending_queue; + DROP TABLE IF EXISTS deleted_accounts; diff --git a/internal/db/postgres/admission_queue_repo.go b/internal/db/postgres/admission_queue_repo.go index 584bbaf..d808657 100644 --- a/internal/db/postgres/admission_queue_repo.go +++ b/internal/db/postgres/admission_queue_repo.go @@ -2,11 +2,28 @@ package postgres import ( "context" + "fmt" "Coves/internal/core/posts" ) -// RED STUB (task 5, cycle 2). Signature only; the query is GREEN's. +// maxPendingSubjectsPerPass caps what one backlog scan may return, whatever a +// caller asks for. +// +// The bound belongs here rather than only at the call site because the caller +// is a periodic job and the table grows with every submission the instance has +// ever seen: a driver misconfigured with an enormous batch would hold one +// transaction open across the whole backlog and then try to settle all of it +// inside a single bounded cycle. +const maxPendingSubjectsPerPass = 500 + +// defaultPendingSubjectsPerPass is what a non-positive limit means. +// +// A zero limit reads as "no bound" to a LIMIT clause author and as "no work" to +// everyone else, and neither is a useful answer to give a queue. Substituting a +// modest page keeps a driver built without an explicit batch size working +// rather than silently idle. +const defaultPendingSubjectsPerPass = 100 // ListPendingSubjects returns the acceptance engine's backlog: subjects this // AppView can actually settle, oldest first. @@ -25,5 +42,48 @@ import ( // all (an acceptance that arrived before its subject), and a LEFT JOIN would // let those through as decidable when there is nothing to decide about. func (r *postgresAdmissionRepo) ListPendingSubjects(ctx context.Context, limit int) ([]posts.PendingSubject, error) { - return nil, nil + switch { + case limit <= 0: + limit = defaultPendingSubjectsPerPass + case limit > maxPendingSubjectsPerPass: + limit = maxPendingSubjectsPerPass + } + + // Both joins are INNER, and both exclusions are spelled as a join rather + // than as a NOT EXISTS so the planner can use the ordinary primary keys on + // posts and communities. ORDER BY created_at is the queue discipline: a + // queue served newest-first starves its own backlog, and the post that has + // waited longest is the one whose author is already wondering where it went. + const query = ` + SELECT admissions.community_did, admissions.post_uri, admissions.created_at + FROM community_post_admissions AS admissions + JOIN posts + ON posts.uri = admissions.post_uri + AND posts.deleted_at IS NULL + JOIN communities + ON communities.did = admissions.community_did + AND communities.pds_refresh_token_encrypted IS NOT NULL + WHERE admissions.status IN ('pending', 'pending_reacceptance') + ORDER BY admissions.created_at ASC + LIMIT $1 + ` + + rows, err := r.db.QueryContext(ctx, query, limit) + if err != nil { + return nil, fmt.Errorf("listing the acceptance engine's pending subjects: %w", err) + } + defer func() { _ = rows.Close() }() + + subjects := make([]posts.PendingSubject, 0, limit) + for rows.Next() { + var subject posts.PendingSubject + if err := rows.Scan(&subject.CommunityDID, &subject.PostURI, &subject.CreatedAt); err != nil { + return nil, fmt.Errorf("scanning a pending subject: %w", err) + } + subjects = append(subjects, subject) + } + if err := rows.Err(); err != nil { + return nil, fmt.Errorf("reading the pending subjects: %w", err) + } + return subjects, nil } diff --git a/internal/db/postgres/post_repo.go b/internal/db/postgres/post_repo.go index 025b499..f880bb2 100644 --- a/internal/db/postgres/post_repo.go +++ b/internal/db/postgres/post_repo.go @@ -17,6 +17,12 @@ import ( "github.com/lib/pq" ) +// legacyPostCollection is the DEPRECATED community-repo post record (§3.0). A +// post written since the author-owned flip is a posts.PostV2Collection record +// in the author's own repo; both are indexed into the same table, and the URI +// is what says which one a row came from. +const legacyPostCollection = "social.coves.community.post" + type postgresPostRepo struct { db *sql.DB } @@ -484,13 +490,28 @@ func scanPostView(rows *sql.Rows, extraDest ...interface{}) (*posts.PostView, er CommentCount: postView.CommentCount, } - // Build the record (required by lexicon) + // Build the record (required by lexicon). + // + // THE SHAPE FOLLOWS THE URI, because the two post collections are different + // records rather than two spellings of one. A postv2 record lives in the + // AUTHOR's repo and has NO author field — its removal is what makes + // authorship unforgeable (§3.1) — so synthesising one here would hand every + // reader back exactly the field the flip deleted, and a client that trusted + // it would be trusting a value this AppView made up. + collection := posts.CollectionOfPostURI(postView.URI) + if collection == "" { + collection = legacyPostCollection + } record := map[string]interface{}{ - "$type": "social.coves.community.post", + "$type": collection, "community": communityRef.DID, - "author": authorView.DID, "createdAt": postView.CreatedAt.Format(time.RFC3339), } + if collection == legacyPostCollection { + // The deprecated community-repo record DOES carry an author field, and + // it is part of the record a client may verify against the repo. + record["author"] = authorView.DID + } // Add optional fields to record if present if title.Valid { -- 2.51.2