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 {