diff --git a/.env.dev.example b/.env.dev.example index d36418a..b4e9567 100644 --- a/.env.dev.example +++ b/.env.dev.example @@ -200,4 +200,15 @@ OTEL_ENABLED=false # the other recognises the retry a client sends after a lost response. Raising # it makes reposting the same content take longer to become possible again # (default: 1h). +# +# Two deliberate edges of this window. First, dedupe is bucketed against the +# epoch rather than against each submission, so the effective protection +# ranges from just above zero up to the full window depending on where in the +# bucket a submission lands — content posted just before a bucket edge can be +# reposted right after it. Accepted so the dedupe key expires on its own with +# no cleanup process. Second, a crash between reserving a submission slot and +# the PDS write leaves an orphaned reservation that burns one quota slot and +# refuses identical content as a duplicate until the bucket rolls, then +# clears itself; there is no sweeper, and the damage is bounded by this +# window. # POST_SUBMISSIONS_DEDUPE_WINDOW=1h diff --git a/.env.prod.example b/.env.prod.example index 18b97d5..f986583 100644 --- a/.env.prod.example +++ b/.env.prod.example @@ -394,6 +394,17 @@ OTEL_ENABLED=false # the other recognises the retry a client sends after a lost response. Raising # it makes reposting the same content take longer to become possible again # (default: 1h). +# +# Two deliberate edges of this window. First, dedupe is bucketed against the +# epoch rather than against each submission, so the effective protection +# ranges from just above zero up to the full window depending on where in the +# bucket a submission lands — content posted just before a bucket edge can be +# reposted right after it. Accepted so the dedupe key expires on its own with +# no cleanup process. Second, a crash between reserving a submission slot and +# the PDS write leaves an orphaned reservation that burns one quota slot and +# refuses identical content as a duplicate until the bucket rolls, then +# clears itself; there is no sweeper, and the damage is bounded by this +# window. # POST_SUBMISSIONS_DEDUPE_WINDOW=1h # ============================================================================= diff --git a/docs/PRD_AUTHOR_OWNED_POSTS.md b/docs/PRD_AUTHOR_OWNED_POSTS.md index 0a349a9..df9c759 100644 --- a/docs/PRD_AUTHOR_OWNED_POSTS.md +++ b/docs/PRD_AUTHOR_OWNED_POSTS.md @@ -35,7 +35,10 @@ moderation.ban ingestion exists; no production ban writer yet); rate limits/dedupe get a synchronous post_submissions ledger (migration 035) — the posts table is unusable as a limiter substrate (ingestion lag, author-supplied created_at, delete-to-evade); per-origin-PDS quota -explicitly deferred to Beta.** +explicitly deferred to Beta. +Rev 2.6 (2026-08-08): task-3 second-opinion — fingerprint normalized to +resolved-DID scope, release decoupled from request context, admission wiring +fail-loud, ActorClass fail-closed.** **Supersedes** the write-path architecture in `docs/federation-prd.md`: that document solves cross-instance posting by service-auth-forwarding the write to @@ -335,6 +338,11 @@ fiction. Failure mode: author-repo write succeeds, acceptance write fails → post stays `pending`; the firehose engine (§5.6) retries idempotently (same rkey). Degraded latency, not data loss. Never roll back the author's record. +There is a lost-response asymmetry here: when the PDS write's outcome is +ambiguous (the record may or may not exist) and the submission reservation is +released, a client retry can produce a duplicate post — the remedy, noted for +task 6, is to derive the record rkey deterministically from the submission +fingerprint so retries become idempotent at the PDS layer. `post.delete` likewise flips to an author-session delete. @@ -596,11 +604,14 @@ Anyone can write unlimited posts naming any community; nothing stops the - Per-author, per-community, and per-origin-PDS submission quotas in the acceptance engine (new policy, §4.1), with `rejected` + `rate-limit-exceeded` decision codes, `redrivable = false`. -- Dedupe identical submissions by (author, community, content CID). +- Dedupe identical submissions by (author, community, canonical-record + fingerprint) — the hash of the canonical record with `createdAt` removed + (§4.1, rev 2.5), bucketed by the dedupe window. - Debounce edit re-evaluation per post (a rapid edit storm collapses to the latest CID). - Retention caps on `pending`/`rejected` admission rows for never-accepted - posts. + posts, and on the `post_submissions` ledger (migration 035), whose + confirmed rows are otherwise never deleted and grow one per admitted post. - Notify endpoint: per-caller and per-PDS quotas on top of service-auth. - All outbound fetches (identity bootstrap §5.3, record fetch §5.4/§7) behind SSRF guards, response-size caps, and timeouts. diff --git a/internal/api/handlers/post/harness_test.go b/internal/api/handlers/post/harness_test.go index 11c067f..b1956e4 100644 --- a/internal/api/handlers/post/harness_test.go +++ b/internal/api/handlers/post/harness_test.go @@ -82,6 +82,8 @@ func newCreateStack(t *testing.T, db *sql.DB) createStack { communityService, nil, nil, nil, nil, pdsURL, + // Handler translation is the subject here, not admission policy. + posts.WithAdmissionPolicy(posts.NewAllowAllAdmissionPolicyForTests()), ) return createStack{ diff --git a/internal/config/submissions_test.go b/internal/config/submissions_test.go index 39b62db..13c3581 100644 --- a/internal/config/submissions_test.go +++ b/internal/config/submissions_test.go @@ -63,6 +63,39 @@ func TestLoad_SubmissionQuotaIsReadFromTheEnvironment(t *testing.T) { } } +// The two tests above prove Load reads the variables and TestValidate below +// proves a zero quota is refused at the struct level. Neither proves the +// COMPOSITION: that a bad value set in the environment actually stops Load() +// itself, the call startup makes. These two close that gap. + +func TestLoad_RejectsAnUnparseableSubmissionWindow(t *testing.T) { + clearEnv(t) + t.Setenv("IS_DEV_ENV", "true") + t.Setenv("POST_SUBMISSIONS_WINDOW", "banana") + + _, err := Load() + if err == nil { + t.Fatal("Load() accepted POST_SUBMISSIONS_WINDOW=banana; an unparseable window must stop startup, not fall back silently") + } + if !strings.Contains(err.Error(), "POST_SUBMISSIONS_WINDOW") { + t.Errorf("error should name POST_SUBMISSIONS_WINDOW so an operator can fix it; got:\n%s", err.Error()) + } +} + +func TestLoad_RejectsAZeroSubmissionQuotaFromTheEnvironment(t *testing.T) { + clearEnv(t) + t.Setenv("IS_DEV_ENV", "true") + t.Setenv("POST_SUBMISSIONS_MAX_PER_COMMUNITY", "0") + + _, err := Load() + if err == nil { + t.Fatal("Load() accepted POST_SUBMISSIONS_MAX_PER_COMMUNITY=0; the process would start with the abuse limit inverted or disabled") + } + if !strings.Contains(err.Error(), "POST_SUBMISSIONS_MAX_PER_COMMUNITY") { + t.Errorf("error should name POST_SUBMISSIONS_MAX_PER_COMMUNITY so an operator can fix it; got:\n%s", err.Error()) + } +} + // A config assembled with the quota left at its zero value must not validate. // This is the assertion that makes "unset means unlimited" unrepresentable // rather than merely discouraged. diff --git a/internal/core/posts/admit.go b/internal/core/posts/admit.go index b03789c..9c84492 100644 --- a/internal/core/posts/admit.go +++ b/internal/core/posts/admit.go @@ -10,6 +10,7 @@ import ( "log" "time" + "Coves/internal/core/aggregators" "Coves/internal/core/communities" ) @@ -76,9 +77,10 @@ type AdmissionRequest struct { Community string // Fingerprint identifies WHAT is being submitted: the hash of the canonical - // record with createdAt removed (see submissionFingerprint). It is the - // dedupe key, and it must exclude the timestamp or every resubmission of - // identical content would look new. + // record with createdAt and the client-typed community identifier removed, + // and the supplied thumbnail URL folded in (see submissionFingerprint). It + // is the dedupe key, and it must exclude the timestamp or every + // resubmission of identical content would look new. Fingerprint string } @@ -259,48 +261,55 @@ type AdmissionPolicy struct { Now Clock } -// WithAdmissionPolicy enables the ban check, dedupe and per-author rate limit -// on CreatePost. +// WithAdmissionPolicy supplies the ban check, dedupe and per-author rate limit +// on CreatePost. It is not optional: NewPostService refuses to construct a +// service without a complete policy (see mustCompleteAdmissionPolicy), because +// a post service whose admission policy silently defaulted to no-ops would be +// one whose ban check and quota do not exist and nothing says so. func WithAdmissionPolicy(policy AdmissionPolicy) PostServiceOption { return func(s *postService) { s.admission = &policy } } -// completeAdmissionPolicy fills in the collaborators a policy did not name, so -// that CreatePost has exactly ONE decision path to run. -// -// The alternative — branching on whether a policy was supplied, and keeping the -// pre-policy checks inline for the other branch — would leave two copies of the -// community/visibility/authorization sequence, and §4.1 of the PRD exists -// because the one copy we had already drifted from what its docstring claimed. -// -// The substitutes are named for what they are. A service constructed without a -// policy enforces exactly what CreatePost enforced before this decision existed: -// community existence, private visibility, and aggregator authorization. It is -// the shape every test fixture that predates §8 uses, and cmd/server always -// supplies the real policy — which is what makes the ban lookup and the quota -// live in production. -func completeAdmissionPolicy(policy *AdmissionPolicy) *AdmissionPolicy { - complete := AdmissionPolicy{} - if policy != nil { - complete = *policy - } - if complete.Ledger == nil { - complete.Ledger = unmeteredLedger{} - } - if complete.Bans == nil { - complete.Bans = unenforcedBans{} - } - if complete.Now == nil { - complete.Now = time.Now +// mustCompleteAdmissionPolicy is NewPostService's guard: a service may not be +// constructed without a complete admission policy. +// +// It panics rather than returning an error, matching how this codebase treats +// every other mandatory collaborator (aggregators.NewAPIKeyService, +// blueskypost.NewService): a missing policy is a wiring bug that must stop the +// process at startup, not a runtime condition to handle. The old alternative — +// silently substituting a no-op ledger and ban lookup — is exactly how the +// pre-§4.1 docstring came to claim "membership/ban validation" that had never +// existed on the write path. +// +// Every limit must be positive for the same reason config.Validate enforces +// it: a quota that silently disappears when a field is left zero is not a +// quota. Tests that are not about admission opt out EXPLICITLY with +// NewAllowAllAdmissionPolicyForTests. +func mustCompleteAdmissionPolicy(policy *AdmissionPolicy) { + switch { + case policy == nil: + panic("posts.NewPostService: an admission policy is required — wire posts.WithAdmissionPolicy " + + "(cmd/server) or posts.NewAllowAllAdmissionPolicyForTests (fixtures that are not about admission)") + case policy.Ledger == nil: + panic("posts.NewPostService: AdmissionPolicy.Ledger cannot be nil") + case policy.Bans == nil: + panic("posts.NewPostService: AdmissionPolicy.Bans cannot be nil") + case policy.Now == nil: + panic("posts.NewPostService: AdmissionPolicy.Now cannot be nil") + case policy.Limits.MaxPerAuthorPerCommunity <= 0: + panic("posts.NewPostService: AdmissionPolicy.Limits.MaxPerAuthorPerCommunity must be positive") + case policy.Limits.Window <= 0: + panic("posts.NewPostService: AdmissionPolicy.Limits.Window must be positive") + case policy.Limits.DedupeWindow <= 0: + panic("posts.NewPostService: AdmissionPolicy.Limits.DedupeWindow must be positive") } - return &complete } -// unmeteredLedger stands in when no submission ledger was wired: it reserves -// nothing, so neither dedupe nor the per-author quota applies. +// unmeteredLedger is the allow-all test policy's ledger: it reserves nothing, +// so neither dedupe nor the per-author quota applies. // // It cannot silently disable a configured limiter — it is only ever reachable -// when AdmissionPolicy.Ledger is nil, which cmd/server never leaves so. +// through NewAllowAllAdmissionPolicyForTests, whose name is the warning. type unmeteredLedger struct{} func (unmeteredLedger) Reserve(context.Context, ReserveSubmissionCommand) (SubmissionReservation, error) { @@ -316,19 +325,44 @@ func (unmeteredLedger) CountSince(context.Context, string, string, time.Time) (i return 0, nil } -// unenforcedBans stands in when no ban lookup was wired, answering the way an -// author with no membership row does. It returns the sentinel rather than a nil -// membership so that it travels the same branch a real absent row does — the -// "no membership means not banned" translation stays in one place. +// unenforcedBans is the allow-all test policy's ban lookup, answering the way +// an author with no membership row does. It returns the sentinel rather than a +// nil membership so that it travels the same branch a real absent row does — +// the "no membership means not banned" translation stays in one place. type unenforcedBans struct{} func (unenforcedBans) GetMembership(context.Context, string, string) (*communities.Membership, error) { return nil, communities.ErrMembershipNotFound } +// NewAllowAllAdmissionPolicyForTests is the explicit opt-out for TEST fixtures +// whose subject is not admission: it admits everything an unconfigured service +// used to — no ban rows to find, no dedupe, no per-author quota — while +// community existence, visibility and aggregator authorization stay enforced. +// +// THE NAME IS THE CONTRACT: this must never be wired in production code. +// cmd/server wires the real policy, and mustCompleteAdmissionPolicy exists +// precisely so that forgetting to do so fails at startup instead of shipping a +// post service whose §8 enforcement quietly does not exist. The limits are +// real (and enormous) only because construction refuses non-positive ones; the +// unmetered ledger never counts against them anyway. +func NewAllowAllAdmissionPolicyForTests() AdmissionPolicy { + return AdmissionPolicy{ + Ledger: unmeteredLedger{}, + Bans: unenforcedBans{}, + Limits: SubmissionLimits{ + MaxPerAuthorPerCommunity: 1 << 30, + Window: time.Hour, + DedupeWindow: time.Hour, + }, + Now: time.Now, + } +} + // admissionDeps assembles the decision's inputs from the service's -// collaborators. s.admission is never nil — NewPostService completes it — so -// this cannot silently hand admitPost a missing ledger or clock. +// collaborators. s.admission is never nil or incomplete — NewPostService +// refuses to construct without a complete policy — so this cannot silently +// hand admitPost a missing ledger or clock. func (s *postService) admissionDeps() admissionDeps { return admissionDeps{ communities: s.communityService, @@ -435,6 +469,17 @@ type admissionDeps struct { // A non-nil error means the decision could NOT be made — a lookup failed — and // is distinct from a refusal, which is a decision. func admitPost(ctx context.Context, deps admissionDeps, req AdmissionRequest) (AdmissionDecision, error) { + // 0. The actor class must be one this decision knows. It gates everything + // below — including the trusted skip of visibility, ban and authorization — + // so an unknown value must fail CLOSED before any lookup runs. Falling + // through would hand the zero value (a caller that forgot to classify) the + // widest privileges in the system. + switch req.Actor { + case ActorUser, ActorRegisteredAggregator, ActorTrustedAggregator: + default: + return undecided(fmt.Errorf("unknown actor class %q: the submission cannot be evaluated", req.Actor)) + } + // 1. Community resolution. Two lookups — the at-identifier to a DID, then // the DID to the indexed row — and either failing to find it is the same // answer to the client. @@ -490,9 +535,18 @@ func admitPost(ctx context.Context, deps admissionDeps, req AdmissionRequest) (A // sentinel that caused it: the boundary tells 403 from 429 by matching on // it, and a bare code would have a well-behaved aggregator retry a // permanent refusal forever. + // + // Only the package's POLICY sentinels are refusals. ValidateAggregatorPost + // also fails when its own lookups do (a wrapped driver error carrying no + // sentinel), and that is an undecided infrastructure failure like any + // other — dressing it as DecisionAggregatorNotAuthorized would mint a + // permanent-sounding 403 out of a Postgres blip. if req.Actor == ActorRegisteredAggregator { if err := deps.aggregators.ValidateAggregatorPost(ctx, req.AuthorDID, community.DID); err != nil { - return AdmissionDecision{Code: DecisionAggregatorNotAuthorized, Cause: err}, nil + if errors.Is(err, aggregators.ErrNotAuthorized) || errors.Is(err, aggregators.ErrRateLimitExceeded) { + return AdmissionDecision{Code: DecisionAggregatorNotAuthorized, Cause: err}, nil + } + return undecided(fmt.Errorf("failed to validate aggregator post: %w", err)) } } @@ -550,8 +604,21 @@ func undecided(err error) (AdmissionDecision, error) { // use. The error is logged rather than returned: every caller reaches this // while already reporting a refusal or a failure, and replacing that answer // with a second one would hide the reason the submission was actually stopped. +// +// The release runs DETACHED from the caller's cancellation (precedent: +// adminreports.raiseAlert), because the most common reason to be here at all +// is that the caller's context is already dead — a client that disconnected +// mid-write is exactly a failed PDS write. A release issued on that context +// would be refused by Postgres as canceled too, and the reservation would +// leak: one quota slot burned and the author's retry refused as a duplicate, +// with nothing but a warning line to say why. Context values (trace IDs) +// survive; only the cancellation signal is dropped, and the fresh timeout +// keeps a wedged database from pinning the goroutine. func releaseReservation(ctx context.Context, ledger SubmissionLedger, reservation SubmissionReservation) { - if err := ledger.Release(ctx, reservation); err != nil { + releaseCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), 5*time.Second) + defer cancel() + + if err := ledger.Release(releaseCtx, reservation); err != nil { log.Printf("[POST-ADMIT] Warning: failed to release submission reservation %d: %v", reservation.ID, err) } } @@ -559,6 +626,16 @@ func releaseReservation(ctx context.Context, ledger SubmissionLedger, reservatio // dedupeBucket is the index of the window `now` falls in, so that two // submissions in the same window collide on the ledger's unique key and two // submissions a window apart do not. +// +// Buckets are aligned to the epoch, not to the submission, so the effective +// dedupe protection ranges over (0, window] depending on where in the bucket +// a submission lands: content submitted just before a bucket edge can be +// resubmitted the moment the edge passes. That tradeoff is deliberate — the +// epoch-aligned key self-expires without a sweeper, where a per-submission +// window would need a range predicate or a cleanup process to expire. The +// same boundary bounds a leaked reservation: a crash between Reserve and the +// PDS write leaves a row that burns one quota slot and refuses identical +// content as a duplicate until the bucket rolls, then heals on its own. func dedupeBucket(now time.Time, window time.Duration) int64 { // A non-positive window would divide by zero. config.Validate refuses to // start a process with one, so reaching this is a wiring bug rather than an @@ -571,19 +648,40 @@ func dedupeBucket(now time.Time, window time.Duration) int64 { return now.UnixNano() / int64(window) } -// submissionFingerprint hashes what a moderator would judge about a record: -// everything except createdAt. +// submissionFingerprint hashes what a moderator would judge about a +// submission: everything on the record except createdAt and community, plus +// the thumbnail URL that rides alongside the record. // // The timestamp has to go. It is stamped by the server at submission time -// (service.go step 9), so it differs on every attempt — including the retry +// (service.go step 6), so it differs on every attempt — including the retry // after a lost response, which is the case dedupe exists to catch. A // fingerprint that included it would never match anything. -func submissionFingerprint(record PostRecord) string { - // The record is taken by value, so clearing the timestamp here cannot - // affect the record the caller goes on to write. +// +// The community field has to go too, for the opposite failure. It holds the +// at-identifier as the CLIENT typed it — a handle one time, a DID the next — +// while the ledger's unique key already scopes the fingerprint by the +// RESOLVED community DID. Hashing the client-typed identifier would let the +// same submission to the same community bypass dedupe simply by switching +// spelling between attempts; leaving it out cannot collide submissions to +// DIFFERENT communities, because the ledger key keeps them apart. +// +// The thumbnail URL is IN, even though it is not a record field: a trusted +// aggregator supplies it alongside the record (CreatePostRequest.ThumbnailURL) +// and it changes what readers ultimately see. Two submissions differing only +// in their thumbnail are different posts, and excluding it would refuse the +// second as a repeat of the first. +func submissionFingerprint(record PostRecord, thumbnailURL *string) string { + // The record is taken by value, so clearing fields here cannot affect the + // record the caller goes on to write. record.CreatedAt = "" + record.Community = "" + + material := struct { + Record PostRecord `json:"record"` + ThumbnailURL *string `json:"thumbnailUrl,omitempty"` + }{Record: record, ThumbnailURL: thumbnailURL} - canonical, err := json.Marshal(record) + canonical, err := json.Marshal(material) if err != nil { // Unreachable in practice: every field of a PostRecord either has a // concrete marshalable type or holds a value decoded from JSON. Hashing @@ -591,7 +689,8 @@ func submissionFingerprint(record PostRecord) string { // a constant fingerprint would collide every submission with every // other, and the second post the instance ever received would be // refused as a repeat of the first. - canonical = []byte(fmt.Sprintf("%#v", record)) + log.Printf("[POST-ADMIT] Warning: submission fingerprint fell back to a Go rendering, canonical JSON marshal failed: %v", err) + canonical = []byte(fmt.Sprintf("%#v", material)) } sum := sha256.Sum256(canonical) diff --git a/internal/core/posts/admit_matrix_test.go b/internal/core/posts/admit_matrix_test.go index 0f77e17..0cbade7 100644 --- a/internal/core/posts/admit_matrix_test.go +++ b/internal/core/posts/admit_matrix_test.go @@ -556,6 +556,62 @@ func TestAdmitPost_AnAggregatorRefusalKeepsItsSentinel(t *testing.T) { // Failing closed // --------------------------------------------------------------------------- +// An actor class the decision does not recognise must fail CLOSED, before any +// lookup runs. The zero value is the dangerous one: a caller that forgot to +// classify the actor would otherwise sail past every check that switches on +// req.Actor — which is exactly the trusted-aggregator skip path — and a +// database outage would be the least of it. +func TestAdmitPost_AnUnknownActorClassFailsClosed(t *testing.T) { + t.Parallel() + + for _, tc := range []struct { + name string + actor ActorClass + }{ + {"the zero value", ActorClass("")}, + {"an unrecognised class", ActorClass("99")}, + } { + t.Run(tc.name, func(t *testing.T) { + t.Parallel() + + h := newAdmitHarness() + + decision, err := h.admit(t, tc.actor, "probe") + require.Error(t, err, + "an unclassifiable actor must fail the request, never fall through to the trusted-skip path") + assert.False(t, decision.Admitted()) + assert.Emptyf(t, decision.Code, + "an unclassifiable actor is a caller bug, not a policy refusal (%q)", decision.Code) + + assert.Zero(t, h.communities.resolveCalls, "no lookup may run for an actor the decision cannot classify") + assert.Zero(t, h.communities.getCalls, "no lookup may run for an actor the decision cannot classify") + assert.Zero(t, h.bans.calls, "no lookup may run for an actor the decision cannot classify") + assert.Zero(t, h.aggregators.calls, "no lookup may run for an actor the decision cannot classify") + assert.Empty(t, h.ledger.reserveCalls, "no reservation may be taken for an actor the decision cannot classify") + }) + } +} + +// A ValidateAggregatorPost failure that is NOT one of the aggregators package's +// policy sentinels is infrastructure, not a refusal. Mapping a database error +// to DecisionAggregatorNotAuthorized would tell a perfectly authorized +// aggregator to stop asking — a 403 minted out of a Postgres blip. +func TestAdmitPost_AnAggregatorLookupFailureIsAnErrorNotARefusal(t *testing.T) { + t.Parallel() + + h := newAdmitHarness() + h.aggregators.err = errors.New("driver: bad connection") + + decision, err := h.admit(t, ActorRegisteredAggregator, "item") + require.Error(t, err, + "an authorization check that could not be evaluated must fail the request, not refuse it") + assert.False(t, decision.Admitted()) + assert.Emptyf(t, decision.Code, + "an infrastructure failure must not be dressed up as an authorization refusal (%q)", decision.Code) + assert.Empty(t, h.ledger.reserveCalls, + "a submission we could not evaluate must not reserve quota") +} + // A ban lookup that fails for any reason OTHER than "no such membership" must // fail the request. // @@ -587,6 +643,11 @@ func TestAdmitPost_InfrastructureFailuresAreErrorsNotRefusals(t *testing.T) { for _, tc := range []struct { name string setup func(*admitHarness) + + // wantReserveCalls is how many times the failing path was expected to + // reach the ledger before the failure stopped it — the precondition + // that makes the liveRows assertion below meaningful. + wantReserveCalls int }{ { name: "the community index is unreachable", @@ -597,12 +658,14 @@ func TestAdmitPost_InfrastructureFailuresAreErrorsNotRefusals(t *testing.T) { setup: func(h *admitHarness) { h.communities.getErr = errors.New("connection reset by peer") }, }, { - name: "the ledger insert fails for a reason that is not a duplicate", - setup: func(h *admitHarness) { h.ledger.reserveErr = errors.New("deadlock detected") }, + name: "the ledger insert fails for a reason that is not a duplicate", + setup: func(h *admitHarness) { h.ledger.reserveErr = errors.New("deadlock detected") }, + wantReserveCalls: 1, }, { - name: "the quota count fails", - setup: func(h *admitHarness) { h.ledger.countErr = errors.New("statement timeout") }, + name: "the quota count fails", + setup: func(h *admitHarness) { h.ledger.countErr = errors.New("statement timeout") }, + wantReserveCalls: 1, }, } { t.Run(tc.name, func(t *testing.T) { @@ -611,12 +674,19 @@ func TestAdmitPost_InfrastructureFailuresAreErrorsNotRefusals(t *testing.T) { h := newAdmitHarness() tc.setup(h) + require.Zero(t, h.ledger.liveRows(), "the ledger must start empty for the release assertion to mean anything") + decision, err := h.admit(t, ActorUser, "probe") require.Error(t, err) assert.False(t, decision.Admitted(), "a decision that could not be made must not read as an admission") assert.Emptyf(t, decision.Code, "an infrastructure failure must not be dressed up as a policy code (%q); the client would be told to stop retrying something that will work in a second", decision.Code) + + assert.Len(t, h.ledger.reserveCalls, tc.wantReserveCalls, + "the failure was injected at a different point in the flow than this case describes") + assert.Zero(t, h.ledger.liveRows(), + "an undecided submission left a reservation on the ledger: it burned quota and will refuse the client's retry as a duplicate") }) } } @@ -901,6 +971,117 @@ func TestAdmitPost_QuotaIsScopedToOneCommunity(t *testing.T) { "the quota is per community; being at the limit in one must not silence the author in another") } +// --------------------------------------------------------------------------- +// Construction +// --------------------------------------------------------------------------- + +// A post service without a complete admission policy is not a lighter post +// service — it is one whose ban check, dedupe and quota silently do not exist. +// Construction must therefore fail loudly, the way this codebase treats every +// other mandatory collaborator (aggregators.NewAPIKeyService, blueskypost), +// rather than substituting no-op defaults a production wiring mistake would +// never notice. +func TestNewPostService_RefusesConstructionWithoutACompleteAdmissionPolicy(t *testing.T) { + t.Parallel() + + validLimits := SubmissionLimits{ + MaxPerAuthorPerCommunity: 3, + Window: time.Hour, + DedupeWindow: time.Hour, + } + complete := func() AdmissionPolicy { + return AdmissionPolicy{ + Ledger: &stubLedger{now: time.Now}, + Bans: &stubBans{}, + Limits: validLimits, + Now: time.Now, + } + } + + for _, tc := range []struct { + name string + opts []PostServiceOption + }{ + { + name: "no admission policy at all", + opts: nil, + }, + { + name: "a policy with no ledger", + opts: []PostServiceOption{WithAdmissionPolicy(func() AdmissionPolicy { + p := complete() + p.Ledger = nil + return p + }())}, + }, + { + name: "a policy with no ban lookup", + opts: []PostServiceOption{WithAdmissionPolicy(func() AdmissionPolicy { + p := complete() + p.Bans = nil + return p + }())}, + }, + { + name: "a policy with no clock", + opts: []PostServiceOption{WithAdmissionPolicy(func() AdmissionPolicy { + p := complete() + p.Now = nil + return p + }())}, + }, + { + name: "a policy with an unset quota", + opts: []PostServiceOption{WithAdmissionPolicy(func() AdmissionPolicy { + p := complete() + p.Limits.MaxPerAuthorPerCommunity = 0 + return p + }())}, + }, + { + name: "a policy with an unset window", + opts: []PostServiceOption{WithAdmissionPolicy(func() AdmissionPolicy { + p := complete() + p.Limits.Window = 0 + return p + }())}, + }, + { + name: "a policy with an unset dedupe window", + opts: []PostServiceOption{WithAdmissionPolicy(func() AdmissionPolicy { + p := complete() + p.Limits.DedupeWindow = 0 + return p + }())}, + }, + } { + t.Run(tc.name, func(t *testing.T) { + t.Parallel() + + require.Panics(t, func() { + NewPostService(nil, nil, nil, nil, nil, nil, "", tc.opts...) + }, "a service constructed without a complete admission policy would enforce nothing and say nothing about it") + }) + } + + t.Run("a complete policy constructs", func(t *testing.T) { + t.Parallel() + + require.NotPanics(t, func() { + NewPostService(nil, nil, nil, nil, nil, nil, "", WithAdmissionPolicy(complete())) + }) + }) + + t.Run("the test-only allow-all policy constructs", func(t *testing.T) { + t.Parallel() + + require.NotPanics(t, func() { + NewPostService(nil, nil, nil, nil, nil, nil, "", + WithAdmissionPolicy(NewAllowAllAdmissionPolicyForTests())) + }, "fixtures that are not about admission need an explicit, honestly-named way to opt out") + }) +} + // --------------------------------------------------------------------------- // The sentinel's wording // --------------------------------------------------------------------------- @@ -924,9 +1105,12 @@ func TestErrDuplicateSubmissionIsNotAStorageConflict(t *testing.T) { // The fingerprint is what makes two submissions "identical". createdAt is // stamped per attempt, so including it would make every retry look new and -// dedupe would never fire; everything a moderator would judge must be included, -// or two genuinely different posts would collide and the second would be -// refused as a repeat of the first. +// dedupe would never fire; the community field is the identifier as the CLIENT +// typed it, so including it would let a handle-vs-DID resubmission bypass +// dedupe; everything a moderator would judge must be included — including the +// thumbnail an aggregator supplies alongside the record — or two genuinely +// different posts would collide and the second would be refused as a repeat of +// the first. func TestSubmissionFingerprint(t *testing.T) { t.Parallel() @@ -948,24 +1132,46 @@ func TestSubmissionFingerprint(t *testing.T) { later := base() later.CreatedAt = "2026-08-01T12:00:09Z" - assert.Equal(t, submissionFingerprint(base()), submissionFingerprint(later), + assert.Equal(t, submissionFingerprint(base(), nil), submissionFingerprint(later, nil), "the server stamps createdAt per attempt, so a fingerprint that included it would never match a retry") }) + t.Run("the community identifier is excluded", func(t *testing.T) { + t.Parallel() + + byHandle := base() + byHandle.Community = admitCommunityHandle + + assert.Equal(t, submissionFingerprint(base(), nil), submissionFingerprint(byHandle, nil), + "the community field holds whatever identifier the client typed; hashing it would let the same "+ + "submission dodge dedupe by naming the community by handle once and by DID the next time — "+ + "the ledger's unique key already scopes the fingerprint to the RESOLVED community DID") + }) + t.Run("a non-empty fingerprint", func(t *testing.T) { t.Parallel() - assert.NotEmpty(t, submissionFingerprint(base()), + assert.NotEmpty(t, submissionFingerprint(base(), nil), "an empty fingerprint would make every submission collide with every other") }) + t.Run("a different thumbnail is a different submission", func(t *testing.T) { + t.Parallel() + + one, two := "https://example.com/thumb-1.jpg", "https://example.com/thumb-2.jpg" + assert.NotEqual(t, submissionFingerprint(base(), &one), submissionFingerprint(base(), &two), + "the thumbnail is submission material an aggregator supplies alongside the record; "+ + "excluding it would refuse a post differing only in its thumbnail as a repeat") + assert.NotEqual(t, submissionFingerprint(base(), nil), submissionFingerprint(base(), &one), + "a submission with a thumbnail is not a repeat of the same submission without one") + }) + for _, tc := range []struct { field string mutate func(*PostRecord) }{ {"title", func(r *PostRecord) { title := "A different title"; r.Title = &title }}, {"content", func(r *PostRecord) { content := "Different body text"; r.Content = &content }}, - {"community", func(r *PostRecord) { r.Community = "did:plc:dddddddddddddddddddddddd" }}, {"author", func(r *PostRecord) { r.Author = "did:plc:eeeeeeeeeeeeeeeeeeeeeeee" }}, {"embed", func(r *PostRecord) { r.Embed = map[string]interface{}{"$type": "social.coves.embed.external"} @@ -976,8 +1182,70 @@ func TestSubmissionFingerprint(t *testing.T) { changed := base() tc.mutate(&changed) - assert.NotEqual(t, submissionFingerprint(base()), submissionFingerprint(changed), + assert.NotEqual(t, submissionFingerprint(base(), nil), submissionFingerprint(changed, nil), "two posts differing in %s would collide, and the second would be refused as a repeat of the first", tc.field) }) } } + +// The dedupe gate must recognise a resubmission no matter which at-identifier +// the client used to name the community. The ledger's unique key scopes the +// fingerprint by the RESOLVED community DID, so the fingerprint itself must not +// re-introduce the client-typed identifier — a fingerprint that hashed it would +// admit the same post twice for anyone who typed the handle once and the DID +// the second time. +func TestAdmitPost_ResubmissionByDIDAfterHandleIsADuplicate(t *testing.T) { + t.Parallel() + + h := newAdmitHarness() + + title, content := "The same post", "the same body" + record := PostRecord{ + Type: postCollection, + Author: admitAuthorDID, + Title: &title, + Content: &content, + } + + byHandle := record + byHandle.Community = admitCommunityHandle + first, err := h.admit(t, ActorUser, submissionFingerprint(byHandle, nil)) + require.NoError(t, err) + require.True(t, first.Admitted()) + + byDID := record + byDID.Community = admitCommunityDID + second, err := h.admit(t, ActorUser, submissionFingerprint(byDID, nil)) + require.NoError(t, err) + assert.Equal(t, DecisionDuplicateSubmission, second.Code, + "naming the community by DID instead of by handle must not turn a resubmission into a new post") + assert.Equal(t, 1, h.ledger.liveRows()) +} + +// The other direction of the same property: two submissions differing ONLY in +// their thumbnail are different posts, and both must be admitted. +func TestAdmitPost_AThumbnailOnlyDifferenceIsNotADuplicate(t *testing.T) { + t.Parallel() + + h := newAdmitHarness() + + title := "The same link, a different thumbnail" + record := PostRecord{ + Type: postCollection, + Community: admitCommunityHandle, + Author: admitAuthorDID, + Title: &title, + } + + one, two := "https://example.com/thumb-1.jpg", "https://example.com/thumb-2.jpg" + + first, err := h.admit(t, ActorUser, submissionFingerprint(record, &one)) + require.NoError(t, err) + require.True(t, first.Admitted()) + + second, err := h.admit(t, ActorUser, submissionFingerprint(record, &two)) + require.NoError(t, err) + assert.Truef(t, second.Admitted(), + "a thumbnail-only difference is a different post, refused with %q", second.Code) + assert.Equal(t, 2, h.ledger.liveRows()) +} diff --git a/internal/core/posts/errors.go b/internal/core/posts/errors.go index 713d9de..6a7286c 100644 --- a/internal/core/posts/errors.go +++ b/internal/core/posts/errors.go @@ -26,7 +26,12 @@ var ( // ErrNotFound is returned when a post is not found by URI ErrNotFound = errors.New("post not found") - // ErrRateLimitExceeded is returned when an aggregator exceeds rate limits + // ErrRateLimitExceeded is returned when a submission is refused for being + // over quota — primarily the per-author, per-community submission limit of + // PRD_AUTHOR_OWNED_POSTS.md §8 (DecisionRateLimitExceeded). The handler + // maps it to a 429. (An aggregator over its OWN hourly quota is refused + // through the aggregators package's sentinel instead, so the boundary can + // tell the two apart.) ErrRateLimitExceeded = errors.New("rate limit exceeded") // ErrInvalidCursor is returned when a pagination cursor is malformed diff --git a/internal/core/posts/service.go b/internal/core/posts/service.go index 20d6b76..8fe70a7 100644 --- a/internal/core/posts/service.go +++ b/internal/core/posts/service.go @@ -74,24 +74,33 @@ func NewPostService( for _, opt := range opts { opt(s) } - s.admission = completeAdmissionPolicy(s.admission) + // The admission policy is mandatory, and a missing or partial one panics + // here rather than defaulting to no-ops: a post service whose ban check and + // quota silently do not exist is a wiring bug, not a configuration. + mustCompleteAdmissionPolicy(s.admission) return s } // CreatePost creates a new post in a community // Flow: -// 1. Validate input -// 2. Check if author is an aggregator (server-side validation using DID from JWT) -// 3. Admission: one decision over community existence, visibility, ban, +// 1. Validate input (and normalize embed/facet URIs) +// 2. Verify the authenticated DID matches the request's author DID +// 3. Classify the actor: trusted aggregator, registered aggregator, or user +// 4. Admission: one decision over community existence, visibility, ban, // aggregator authorization, dedupe and the per-author quota (admitPost) -// 4. Build post record -// 5. Write to community's PDS repository -// 6. If aggregator: record post for rate limiting -// 7. Return URI/CID (AppView indexes asynchronously via Jetstream) +// 5. Ensure the community has fresh PDS credentials (token refresh) +// 6. Build the post record +// 7. Validate and enhance external embeds (thumb validation, unfurl, blobs) +// 8. Write to community's PDS repository +// 9. If aggregator: record post for rate limiting +// 10. Return URI/CID (AppView indexes asynchronously via Jetstream) // // Admission runs BEFORE the token refresh, the blob uploads and the PDS write, // so a refused submission costs a few lookups rather than an upload — and, -// more to the point, leaves no record in a community that refused it. +// more to the point, leaves no record in a community that refused it. Every +// failure AFTER admission (steps 5-8) must release the ledger reservation the +// admission took, or the failure costs the author a quota slot and refuses +// their retry as a duplicate. func (s *postService) CreatePost(ctx context.Context, req CreatePostRequest) (*CreatePostResponse, error) { // 1. Validate basic input (before DID checks to give clear validation errors) if err := s.validateCreateRequest(&req); err != nil { @@ -146,13 +155,13 @@ func (s *postService) CreatePost(ctx context.Context, req CreatePostRequest) (*C // Check if this is a non-trusted aggregator (requires database lookup) var isOtherAggregator bool if !isTrustedAggregator && s.aggregatorService != nil { - aggregator, err := s.aggregatorService.IsAggregator(ctx, req.AuthorDID) + isAggregator, err := s.aggregatorService.IsAggregator(ctx, req.AuthorDID) if err != nil { log.Printf("[POST-CREATE] Warning: failed to check if DID is aggregator: %v", err) // Don't fail the request - treat as regular user if check fails isOtherAggregator = false } else { - isOtherAggregator = aggregator + isOtherAggregator = isAggregator } } @@ -176,12 +185,14 @@ func (s *postService) CreatePost(ctx context.Context, req CreatePostRequest) (*C // The fingerprint is taken from the record as the CLIENT sent it, before // unfurl enhancement rewrites the embed: two submissions of the same content // must hash the same, and an enriched embed varies with whatever the remote - // page served at the time. + // page served at the time. The thumbnail URL rides along as submitted; the + // client-typed community identifier and the per-attempt timestamp are + // excluded inside submissionFingerprint (see its doc comment). decision, err := admitPost(ctx, s.admissionDeps(), AdmissionRequest{ Actor: actor, AuthorDID: req.AuthorDID, Community: req.Community, - Fingerprint: submissionFingerprint(postRecordFor(req, req.Community, "")), + Fingerprint: submissionFingerprint(postRecordFor(req, req.Community, ""), req.ThumbnailURL), }) if err != nil { return nil, err diff --git a/internal/core/posts/service_admission_test.go b/internal/core/posts/service_admission_test.go index 6b43841..f572c17 100644 --- a/internal/core/posts/service_admission_test.go +++ b/internal/core/posts/service_admission_test.go @@ -5,6 +5,7 @@ package posts_test import ( "context" "database/sql" + "encoding/base64" "fmt" "net/http" "net/http/httptest" @@ -283,20 +284,23 @@ func TestService_TheAuthorQuotaStopsTheNextSubmission(t *testing.T) { assert.NoError(t, err, "the quota is per (author, community); being at the limit in one must not close the others") } -// An identical resubmission is a repeat, not a new post. +// An identical resubmission is a repeat, not a new post — even when the retry +// names the community by DID and the original named it by handle. // // The canonical case is a client that retried after a lost response, and the // answer has to be distinguishable from a quota breach: 409 tells the client its // post already exists, 429 tells it to wait. A submission refused as a duplicate // must also not be billed, or a flaky connection would rate-limit a user who -// posted once. +// posted once. Submitting first by HANDLE and retrying by DID is the identifier +// dodge the fingerprint must not fall for: the ledger scopes dedupe by the +// RESOLVED community DID, so the client-typed spelling must not enter the key. func TestService_AnIdenticalResubmissionIsRefusedAsADuplicate(t *testing.T) { t.Parallel() f := newAdmissionFixture(t) - _, err := f.submit(t, f.base.community.DID, "the very same post") - require.NoError(t, err) + _, err := f.submit(t, f.base.community.Handle, "the very same post") + require.NoError(t, err, "submitting by handle must resolve and admit like submitting by DID") _, err = f.submit(t, f.base.community.DID, "the very same post") require.Error(t, err) @@ -363,6 +367,210 @@ func TestService_AFailedPDSWriteReleasesTheReservation(t *testing.T) { assert.Equal(t, f.base.author.DID, record.Value["author"]) } +// A client that goes away MID-WRITE must still get its reservation back. +// +// The failure path runs on the same context the request came in on, and by the +// time the release runs that context is already dead — the canceled write is +// exactly why the path was taken. A release issued on the caller's context +// would be refused by Postgres as canceled too, and the leak would be +// invisible: the request already failed, the log line is a warning, and the +// author discovers it as a duplicate refusal of a post that does not exist. +// The release must therefore run detached from the caller's cancellation +// (precedent: adminreports raiseAlert), bounded by its own timeout. +func TestService_ACancellationDuringThePDSWriteStillReleasesTheReservation(t *testing.T) { + t.Parallel() + + f := newAdmissionFixture(t) + + // The client's context, canceled by the "PDS" at the exact moment the + // write is in flight — the request-scoped context is dead by the time + // CreatePost's failure path runs, which is the shape of a client + // disconnecting mid-request. + ctx, cancel := context.WithCancel(middleware.SetTestUserDID(context.Background(), f.base.author.DID)) + defer cancel() + + canceling := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + // Kill the caller's context while its write is in flight, then refuse + // the write. Whether CreatePost's failure surfaces as the canceled + // context or as the 500 is an interleaving detail; cancel() happens + // before the response is written, so by the time the failure path runs + // the request's context is dead either way. + cancel() + http.Error(w, `{"error":"InternalServerError"}`, http.StatusInternalServerError) + })) + t.Cleanup(canceling.Close) + + healthyURL := communityPDSURL(t, f.base.db, f.base.community.DID) + setCommunityPDSURL(t, f.base.db, f.base.community.DID, canceling.URL) + + const repeatable = "a post whose client disconnects mid-write" + content := "a body that makes this a complete post" + _, err := f.service.CreatePost(ctx, posts.CreatePostRequest{ + Community: f.base.community.DID, + Title: func() *string { s := repeatable; return &s }(), + Content: &content, + AuthorDID: f.base.author.DID, + }) + require.Error(t, err, "the write ran against a dead context and must fail") + + assert.Zerof(t, f.ledgerRows(t, f.base.community.DID), + "the release ran on the caller's canceled context and was refused with it: the reservation leaked, burning a quota slot and blocking the retry as a duplicate") + + setCommunityPDSURL(t, f.base.db, f.base.community.DID, healthyURL) + + // The retry a reconnected client sends: byte-identical content on a live + // context. Admissible only if the canceled attempt released its row. + resp, err := f.submit(t, f.base.community.DID, repeatable) + require.NoError(t, err, "the identical retry was refused, so the canceled attempt leaked its reservation") + require.NotEmpty(t, resp.URI) + assert.Equal(t, 1, f.ledgerRows(t, f.base.community.DID)) +} + +// A token-refresh failure (step 5) happens with the reservation already on the +// ledger, and must give it back for the same reason a failed PDS write must: +// the community's credentials failing is not the author's fault, and must not +// cost them a quota slot or refuse their retry as a duplicate. +func TestService_ATokenRefreshFailureReleasesTheReservation(t *testing.T) { + t.Parallel() + + f := newAdmissionFixture(t) + ctx := context.Background() + + original, err := f.repo.GetByDID(ctx, f.base.community.DID) + require.NoError(t, err) + require.NotEmpty(t, original.PDSAccessToken, "the fixture community must hold real credentials to restore") + + // An expired access token forces EnsureFreshToken down the refresh path, + // and the community's PDS — repointed at a server that 500s everything — + // refuses the refresh. Same seam as the failed-write test: the pds_url and + // credentials are read fresh off the community row on every write. + broken := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + http.Error(w, `{"error":"InternalServerError"}`, http.StatusInternalServerError) + })) + t.Cleanup(broken.Close) + + healthyURL := communityPDSURL(t, f.base.db, f.base.community.DID) + setCommunityPDSURL(t, f.base.db, f.base.community.DID, broken.URL) + require.NoError(t, f.repo.UpdateCredentials(ctx, f.base.community.DID, expiredJWT(t), original.PDSRefreshToken)) + + const repeatable = "a post whose community credentials fail to refresh" + _, err = f.submit(t, f.base.community.DID, repeatable) + require.Error(t, err, "the token refresh failed, so CreatePost must report a failure") + + assert.Zerof(t, f.ledgerRows(t, f.base.community.DID), + "the reservation for a submission that failed at token refresh is still on the ledger") + + setCommunityPDSURL(t, f.base.db, f.base.community.DID, healthyURL) + require.NoError(t, f.repo.UpdateCredentials(ctx, f.base.community.DID, original.PDSAccessToken, original.PDSRefreshToken)) + + resp, err := f.submit(t, f.base.community.DID, repeatable) + require.NoError(t, err, "the identical retry after the credentials recovered was refused, so the refresh-failure path leaked its reservation") + require.NotEmpty(t, resp.URI) + assert.Equal(t, 1, f.ledgerRows(t, f.base.community.DID)) +} + +// expiredJWT builds a structurally valid, long-expired JWT: enough for +// communities.NeedsRefresh (which parses the exp claim without verifying the +// signature) to answer "refresh this now". +func expiredJWT(t *testing.T) string { + t.Helper() + + header := base64.RawURLEncoding.EncodeToString([]byte(`{"alg":"none","typ":"JWT"}`)) + payload := base64.RawURLEncoding.EncodeToString( + []byte(fmt.Sprintf(`{"exp":%d}`, time.Now().Add(-time.Hour).Unix()))) + return header + "." + payload + ".unverified" +} + +// An embed-enhancement failure (step 7) also runs with the reservation held. +// The thumb-must-be-a-blob guard is the reachable failure in that step without +// a network: it refuses the submission after admission, so the refusal must +// hand the slot back or the author's corrected retry meets a quota they never +// spent. +func TestService_AnEmbedEnhancementFailureReleasesTheReservation(t *testing.T) { + t.Parallel() + + f := newAdmissionFixture(t) + + title := "a link post with a malformed thumbnail" + content := "a body that makes this a complete post" + submitWithThumb := func(thumb interface{}) (*posts.CreatePostResponse, error) { + external := map[string]interface{}{ + "uri": "https://example.com/article", + "title": "An article", + "description": "worth reading", + } + if thumb != nil { + external["thumb"] = thumb + } + return f.service.CreatePost( + middleware.SetTestUserDID(context.Background(), f.base.author.DID), + posts.CreatePostRequest{ + Community: f.base.community.DID, + Title: &title, + Content: &content, + AuthorDID: f.base.author.DID, + Embed: map[string]interface{}{ + "$type": "social.coves.embed.external", + "external": external, + }, + }) + } + + // A thumb sent as a URL string passes the lexicon-shape validation of step + // 1 and is refused by the blob guard in step 7 — after admission. + _, err := submitWithThumb("https://example.com/thumb.jpg") + require.Error(t, err) + require.True(t, posts.IsValidationError(err), "the thumb guard reports a validation error, got: %v", err) + + assert.Zerof(t, f.ledgerRows(t, f.base.community.DID), + "the reservation for a submission refused by the embed guard is still on the ledger") + + // The corrected retry — same post, thumb omitted — must be admitted. + resp, err := submitWithThumb(nil) + require.NoError(t, err, "the corrected retry was refused, so the embed-guard path leaked its reservation") + require.NotEmpty(t, resp.URI) + assert.Equal(t, 1, f.ledgerRows(t, f.base.community.DID)) +} + +// The concurrent double-tap the reserve-then-confirm ordering exists to stop: +// two byte-identical submissions racing through CreatePost. The ledger's +// unique key is the only arbiter both goroutines share, so exactly one may be +// admitted; the loser must hear the DUPLICATE sentinel (its post exists — a +// 409), not a generic failure, and must not leave a second row behind. +func TestService_ConcurrentIdenticalSubmissionsAdmitExactlyOne(t *testing.T) { + t.Parallel() + + f := newAdmissionFixture(t) + + const doubleTap = "the same post, submitted twice at once" + var wg sync.WaitGroup + errs := make([]error, 2) + for i := range errs { + wg.Add(1) + go func(slot int) { + defer wg.Done() + _, errs[slot] = f.submit(t, f.base.community.DID, doubleTap) + }(i) + } + wg.Wait() + + winners, losers := 0, 0 + for _, err := range errs { + if err == nil { + winners++ + continue + } + losers++ + assert.ErrorIsf(t, err, posts.ErrDuplicateSubmission, + "the racing loser must hear the duplicate sentinel — its post exists — not %v", err) + } + assert.Equal(t, 1, winners, "exactly one of two identical concurrent submissions may be admitted") + assert.Equal(t, 1, losers) + + assert.Equal(t, 1, f.ledgerRows(t, f.base.community.DID), + "the race must leave exactly the winner's row: the unique key is the arbiter, not a second insert") +} + // communityPDSURL reads a community's stored PDS, so a test that repoints it // can put back what was actually there rather than what it assumed. func communityPDSURL(t *testing.T, db *sql.DB, communityDID string) string { diff --git a/internal/core/posts/service_aggregator_test.go b/internal/core/posts/service_aggregator_test.go index 3080b37..c94e052 100644 --- a/internal/core/posts/service_aggregator_test.go +++ b/internal/core/posts/service_aggregator_test.go @@ -22,8 +22,8 @@ import ( // // An aggregator is a service that writes into communities it does not belong // to, so the membership and visibility rules a human is held to say nothing -// useful about it. CreatePost swaps them for two others (service.go steps 3, 5 -// and 12): the community must have published an authorization record naming +// useful about it. CreatePost swaps them for two others (service.go steps 3, 4 +// and 9): the community must have published an authorization record naming // this aggregator, and the aggregator must be inside its hourly quota. Both are // checked BEFORE anything reaches the community's repository, and a successful // post is then recorded against the aggregator — which is what makes the next @@ -85,7 +85,10 @@ func newAggregatorFixture(t *testing.T) *aggregatorFixture { service: posts.NewPostService( postgres.NewPostRepository(base.db), base.communityService, aggregators.NewAggregatorService(index, base.communityService), - nil, nil, nil, base.pds.URL()), + nil, nil, nil, base.pds.URL(), + // The aggregator's OWN hourly quota is the subject here; the §8 + // per-author policy is opted out of explicitly. + posts.WithAdmissionPolicy(posts.NewAllowAllAdmissionPolicyForTests())), index: index, aggregatorDID: aggregatorDID, authorizationURI: "at://" + base.community.DID + diff --git a/internal/core/posts/service_author_posts_query_test.go b/internal/core/posts/service_author_posts_query_test.go index 0ac97a4..7309f5c 100644 --- a/internal/core/posts/service_author_posts_query_test.go +++ b/internal/core/posts/service_author_posts_query_test.go @@ -96,7 +96,8 @@ func newAuthorPostsFixture(t *testing.T) *authorPostsFixture { communityRepo, pdsURL, fixtures.InstanceDID(), "", nil, nil, nil) // The optional post collaborators (aggregators, blobs, unfurl, bluesky) are // all write-path concerns and stay nil. - postService := posts.NewPostService(postRepo, communityService, nil, nil, nil, nil, pdsURL) + postService := posts.NewPostService(postRepo, communityService, nil, nil, nil, nil, pdsURL, + posts.WithAdmissionPolicy(posts.NewAllowAllAdmissionPolicyForTests())) voteService := votes.NewServiceWithPDSFactory(voteRepo, nil, nil, fixtures.PasswordAuthPDSClientFactory()) auth := fixtures.NewOAuthMiddleware() diff --git a/internal/core/posts/service_create_validation_test.go b/internal/core/posts/service_create_validation_test.go index c7ef6c2..20d828f 100644 --- a/internal/core/posts/service_create_validation_test.go +++ b/internal/core/posts/service_create_validation_test.go @@ -61,7 +61,8 @@ func TestService_CreateResolvesTheCommunityAndValidatesTheRequest(t *testing.T) // in the community's handle. communityService := communities.NewCommunityServiceWithPDSFactory( communityRepo, pdsURL, instanceDID, instanceDomain, nil, nil, nil) - postService := posts.NewPostService(postRepo, communityService, nil, nil, nil, nil, pdsURL) + postService := posts.NewPostService(postRepo, communityService, nil, nil, nil, nil, pdsURL, + posts.WithAdmissionPolicy(posts.NewAllowAllAdmissionPolicyForTests())) authorDID := fixtures.DID("postauthor") _, err := userService.CreateUser(ctx, users.CreateUserRequest{ diff --git a/internal/core/posts/service_writeforward_test.go b/internal/core/posts/service_writeforward_test.go index 01a2880..b30fa03 100644 --- a/internal/core/posts/service_writeforward_test.go +++ b/internal/core/posts/service_writeforward_test.go @@ -39,7 +39,7 @@ import ( // A post record does not live in its author's repo. It lives in the COMMUNITY's // repo, written with the community's own PDS credentials, carrying an `author` // field that names the human who wrote it (internal/core/posts/service.go step -// 9, and the reason the Jetstream consumer's first security check is +// 8, and the reason the Jetstream consumer's first security check is // repoDID == record.community). // // That makes two things testable only from the PDS side. First, that the @@ -125,7 +125,9 @@ func newPostFixture(t *testing.T) *postFixture { return &postFixture{ service: posts.NewPostService( postgres.NewPostRepository(db), communityService, - nil, nil, nil, nil, pdsServer.URL()), + nil, nil, nil, nil, pdsServer.URL(), + // Write-forward is not about admission; the opt-out is explicit. + posts.WithAdmissionPolicy(posts.NewAllowAllAdmissionPolicyForTests())), pds: pdsServer, db: db, communityService: communityService, diff --git a/internal/core/unfurl/post_unfurl_integration_test.go b/internal/core/unfurl/post_unfurl_integration_test.go index d338cb0..d1b13dd 100644 --- a/internal/core/unfurl/post_unfurl_integration_test.go +++ b/internal/core/unfurl/post_unfurl_integration_test.go @@ -61,6 +61,7 @@ func TestPostUnfurl_UnsupportedURL(t *testing.T) { nil, // unfurlService - intentionally nil to test graceful handling nil, // blueskyService testkit.Endpoints().PDS.BaseURL, + posts.WithAdmissionPolicy(posts.NewAllowAllAdmissionPolicyForTests()), ) // Create test user @@ -156,6 +157,7 @@ func TestPostUnfurl_MissingEmbedType(t *testing.T) { unfurlService, nil, // blueskyService testkit.Endpoints().PDS.BaseURL, + posts.WithAdmissionPolicy(posts.NewAllowAllAdmissionPolicyForTests()), ) // Create test user and community diff --git a/internal/db/migrations/035_create_post_submissions.sql b/internal/db/migrations/035_create_post_submissions.sql index 07d77a9..8213073 100644 --- a/internal/db/migrations/035_create_post_submissions.sql +++ b/internal/db/migrations/035_create_post_submissions.sql @@ -24,6 +24,29 @@ -- The bucket is the index of the application's dedupe window, derived from the -- injected clock, so the key expires on its own without a sweeper. -- +-- THE BUCKET BOUNDARY IS A TRADEOFF, taken deliberately. Buckets are aligned +-- to the epoch, not to the submission, so a submission landing just before a +-- bucket edge is protected against an identical resubmission only until that +-- edge: the effective dedupe protection ranges over (0, window] depending on +-- where in the bucket the submission falls. A per-submission window would +-- protect for the full width every time, but expiring it would take a range +-- predicate or a sweeper; the epoch-aligned bucket lets the unique key expire +-- entirely on its own. +-- +-- A LEAKED RESERVATION IS BOUNDED BY THE SAME MECHANISM. If the process dies +-- between reserving the row and completing the PDS write, the orphaned row +-- burns one quota slot until it ages out of the rolling window, and refuses +-- identical content as a duplicate until the bucket rolls. There is no +-- sweeper to reclaim it — deliberate: the damage is bounded and self-healing, +-- and a reaper would be one more process able to disagree with the gate it +-- cleans. +-- +-- THIS TABLE GROWS WITHOUT BOUND. Confirmed rows are never deleted — only +-- releasing a reservation removes a row — so the ledger accumulates one row +-- per admitted post forever. Retention is a known, deliberate deferral: +-- docs/PRD_AUTHOR_OWNED_POSTS.md §8's retention item names post_submissions +-- alongside the pending/rejected admissions rows. +-- -- WHY NO FOREIGN KEYS. Migration 034 dropped posts.fk_author because a -- federated author has no `users` row (§5.3) — the AppView only bootstraps -- authors from trusted bridge PDSs today — and a community named by a diff --git a/tests/e2e/post_admission_contract_test.go b/tests/e2e/post_admission_contract_test.go new file mode 100644 index 0000000..194b850 --- /dev/null +++ b/tests/e2e/post_admission_contract_test.go @@ -0,0 +1,268 @@ +//go:build e2e + +package e2e + +import ( + "context" + "net/http" + "net/url" + "testing" + + "Coves/tests/testkit" + + "github.com/stretchr/testify/require" +) + +// The admission-wiring proof: cmd/server's PRODUCTION-constructed post service +// routes an authenticated social.coves.community.post.create through the §8 +// admission decision (PRD_AUTHOR_OWNED_POSTS, internal/core/posts/admit.go), +// observed end-to-end through the real XRPC surface. +// +// # THE CREDENTIAL, AND WHY IT IS AN AGGREGATOR'S +// +// This contract holds the first real write credential the tier has ever held. +// §3.4b's standing limitation still stands for USERS — nothing but the browser +// OAuth callback mints a sealed session token RequireAuth accepts — but +// post.create is the one Coves route behind DualAuth, and DualAuth's second +// path takes a PDS-signed service JWT from a REGISTERED AGGREGATOR. Every link +// of that chain is mintable inside the hermetic stack: +// +// - the aggregator is a PDS account (provisionAggregatorRepo), whose DID the +// hermetic PLC can resolve to a signing key; +// - it becomes REGISTERED by declaring social.coves.aggregator.service in +// its own repo, which only the firehose can index (aggregator contract's +// opening note) — so holding a working credential at all already proves +// pipeline delivery; +// - the JWT itself comes from the PDS' own com.atproto.server.getServiceAuth, +// signed with the account's repo key, audience'd to the AppView's instance +// DID — exactly what a production bot does. +// +// internal/api/routes/post_aggregator_test.go names this seam as the one it +// cannot reach ("needs the running stack and a token") and injects the +// principal instead. This file is that missing half: the shipped binary's +// DualAuthMiddleware validating a real signature against the hermetic PLC and +// gating on the firehose-fed aggregators table. +// +// # WHAT THIS PROVES ABOUT THE ADMISSION POLICY — AND WHAT IT CANNOT +// +// admitPost classifies this principal ActorRegisteredAggregator, so the checks +// it walks through the wire are community resolution (step 1), aggregator +// authorization (step 4) and the dedupe ledger (step 5). Three properties of +// the NEW decision are pinned here: +// +// - the decision is LIVE: the 403 flips to admitted when the community's +// authorization record arrives over the firehose, with no redeploy; +// - the check ORDER holds: an identical resubmission of a REFUSED submission +// answers the same refusal again, never 409 DuplicateSubmission — with the +// ledger wired, a decision that consulted dedupe ahead of authorization +// would answer 409 the second time (§8: a refusal consumes no quota); +// - reserve-then-release holds: a submission that is ADMITTED but whose PDS +// write then fails must hand its ledger slot back, so retrying it answers +// the write failure again, never 409 — a leaked reservation would turn one +// failed write into a lockout until the dedupe window rolls. +// +// The two USER-classified refusals — 403 Banned (step 3) and the per-author +// 429 RateLimitExceeded (step 6) — are structurally out of this tier's reach: +// they require an ActorUser principal, which requires the sealed-session mint +// that does not exist (§3.4b), and no aggregator credential is ever classified +// ActorUser. They are proven where they can be honestly: the decision matrix +// at T0 (internal/core/posts/admit_matrix_test.go, service_admission_test.go), +// the ledger against real Postgres at T1 (internal/db/postgres), and the +// refusal-to-status mapping at T0 (internal/api/handlers/post/errors_test.go). +// What none of those can see — the production construction in +// cmd/server/wiring.go actually enforcing the decision on the wire — is what +// this file adds. It carries NO ingestion marker: markers are for pipeline +// proofs (§3.4a), and this asserts the client path. +// +// # THE ADMITTED PATH'S KNOWN CEILING, STATED PLAINLY +// +// A community indexed from the firehose carries no PDS credentials in the +// AppView's store — only social.coves.community.create provisions those, and +// it sits behind the OAuth-only middleware this tier cannot satisfy. So an +// ADMITTED submission proceeds past every gate and then fails at the +// community-credential refresh (posts/service.go step 5, EnsureFreshToken on +// an empty token), which the mapper reports as a 500. The assertions below are +// written for the seam under test — refusal vs. admission — and the moment a +// credentialed community becomes reachable at T2, the same test upgrades +// itself to the full dedupe proof (the branch is written out below). +func TestPostAdmissionAPIContract(t *testing.T) { + p := newPipeline(t) + + moderator := p.IndexedAccount(t, "nm") + community := indexedCommunity(t, p, "n", moderator.DID) + aggregator, _ := indexedAggregator(t, p, "na") + + botToken := mintServiceJWT(t, aggregator) + asAggregator := p.AppView.As(botToken) + + // One submission, byte-identical on every attempt: the dedupe fingerprint + // hashes the record as the client sent it, so proving what repeats DON'T + // trigger requires the repeats to be genuine. + title := "admission " + testkit.UniqueID(t) + + t.Run("a service JWT from a DID that is no aggregator stops at the middleware", func(t *testing.T) { + // The security property internal/api/routes/post_aggregator_test.go + // documents as unprovable there: a VALID signature from a real, + // PLC-resolvable identity is still refused when the DID is not in the + // aggregators table. The moderator is exactly that — an indexed USER + // whose PDS mints service JWTs as willingly as anyone's. + // + // The message is asserted as well as the code, and it is load-bearing: + // the middleware answers 401 AuthenticationRequired for a broken + // signature too, and only the message tells "refused by the aggregator + // gate" from "the validator could not resolve the issuer" — the second + // would mean the stack's PLC plumbing is broken, not that the gate held. + err := submitPost(p.AppView.As(mintServiceJWT(t, moderator)), community.DID, title) + refusal := requireXRPCRefusal(t, err, http.StatusUnauthorized, "AuthenticationRequired", + "a non-aggregator's service JWT") + require.Equal(t, "Not a registered aggregator", refusal.XRPCMessage, + "the 401 must come from the aggregator gate, not from signature validation: %v", err) + }) + + t.Run("an unknown community is the decision's first refusal", func(t *testing.T) { + // Answering a POST-mapper refusal at all — not 401 — is the positive + // half of the credential proof: DualAuth validated the aggregator's JWT + // against the hermetic PLC, found the DID in the firehose-fed + // aggregators table, and let the request through to the service, where + // admitPost step 1 refused it. + // + // The DID literal is spelled at 24 base32 characters (a-z, 2-7) for the + // reason TestPostAPIContract gives: UniqueID does not promise that + // alphabet, and a malformed identifier would take the 400 validation + // path instead of the resolution path under test. Nothing indexes this + // DID, on a fresh stack or a kept one. + err := submitPost(asAggregator, "did:plc:aaaaaaaaaanevercommunity", title) + requireXRPCRefusal(t, err, http.StatusNotFound, "CommunityNotFound", + "a submission to a community nobody has indexed") + }) + + t.Run("an unauthorized aggregator is refused, and the refusal consumes nothing", func(t *testing.T) { + // The community exists and is indexed, but has written no authorization + // record for this aggregator — admitPost step 4, carried through the + // mapper as the aggregators-package 403. + err := submitPost(asAggregator, community.DID, title) + requireXRPCRefusal(t, err, http.StatusForbidden, "NotAuthorized", + "a submission from an aggregator the community never authorized") + + // The SAME submission again. §8's check order is observable right here: + // authorization runs AHEAD of dedupe, and a refusal reserves nothing — + // so the identical resubmission meets the identical 403. A decision + // that consulted the ledger first, or leaked a reservation on refusal, + // would answer 409 DuplicateSubmission instead, and this is the only + // tier that can catch the production wiring doing that. + err = submitPost(asAggregator, community.DID, title) + requireXRPCRefusal(t, err, http.StatusForbidden, "NotAuthorized", + "the identical resubmission of a refused submission — a 409 here means a refusal "+ + "consumed a dedupe slot, which §8 forbids") + }) + + // The community lets the aggregator in, the way production does: an + // authorization record in the COMMUNITY's own repo, delivered over the + // firehose. The wait observes the same table ValidateAggregatorPost reads. + community.PutRecord(t, aggregatorAuthorizationCollection, testkit.TID(), + aggregatorAuthorizationRecord(aggregator.DID, community.DID, moderator.DID, true)) + p.Await(t, "the authorization to reach the index the admission decision reads", func() (bool, error) { + enabled, err := p.Authorizations(context.Background(), aggregator.DID, true) + if err != nil { + return false, err + } + return len(enabled) == 1, nil + }) + + t.Run("the authorization's arrival flips the decision without a redeploy", func(t *testing.T) { + // Byte-identical to the submission refused twice above — so everything + // that changed between that 403 and this answer is the firehose-fed + // authorization row, which is the liveness of the decision in one + // assertion. Its two prior refusals reserved nothing, so this attempt's + // own reservation cannot collide with them. + err := submitPost(asAggregator, community.DID, title) + + if err == nil { + // The stack can complete a community-credentialed write — the + // admitted path ran to the PDS and back. The reservation is now + // CONFIRMED on the ledger, so the identical resubmission is the + // full dedupe proof. + err = submitPost(asAggregator, community.DID, title) + requireXRPCRefusal(t, err, http.StatusConflict, "DuplicateSubmission", + "an identical resubmission of an admitted post inside the dedupe window") + return + } + + // Today's ceiling (see the file comment): admission PASSED and the + // write then failed at the community-credential refresh, which no + // firehose-indexed community can satisfy. The mapper reports that + // unclassified failure as exactly one thing, and pinning it keeps this + // branch honest — any 4xx here would mean the admission gate refused, + // which is the regression this contract exists to catch. + requireXRPCRefusal(t, err, http.StatusInternalServerError, "InternalServerError", + "an ADMITTED submission failing at the community-credential refresh — any 4xx here "+ + "means the admission decision refused a submission the community has authorized") + + // And the failed write handed its ledger slot back: the identical + // retry meets the same write failure, never 409. A leaked reservation + // would refuse the retry as a duplicate of a post that does not exist — + // the §8 failure mode where a transient outage becomes a lockout. + err = submitPost(asAggregator, community.DID, title) + requireXRPCRefusal(t, err, http.StatusInternalServerError, "InternalServerError", + "the retry of a failed write — a 409 means the failed write's reservation was never "+ + "released, turning one PDS failure into a lockout until the dedupe window rolls") + }) +} + +// mintServiceJWT asks the stack's PDS to sign a service JWT for account, +// audience'd to the AppView's instance identity — the credential a production +// aggregator bot presents to post.create. +// +// The audience is communityInstanceDID for the reason that constant documents: +// INSTANCE_DID is unset in .env.ci, so the AppView's DualAuth validator was +// built with internal/config's compiled-in default, and a JWT for any other +// audience is refused before the signature is even consulted. lxm is pinned to +// the one route the token is spent on; the AppView validates service JWTs +// endpoint-agnostically (auth.go: lexMethod nil), so this is defence on the +// MINTING side — a leaked test token authorizes nothing else. +func mintServiceJWT(t *testing.T, account *testkit.Account) string { + t.Helper() + + var minted struct { + Token string `json:"token"` + } + err := account.XRPC().Query(context.Background(), "com.atproto.server.getServiceAuth", url.Values{ + "aud": {communityInstanceDID}, + "lxm": {"social.coves.community.post.create"}, + }, &minted) + if err != nil { + t.Fatalf("minting a service JWT for %s via com.atproto.server.getServiceAuth: %v", account.DID, err) + } + if minted.Token == "" { + t.Fatalf("com.atproto.server.getServiceAuth answered 200 with no token for %s", account.DID) + } + return minted.Token +} + +// submitPost drives social.coves.community.post.create as the holder of +// client's credential, with the minimal well-formed body the lexicon requires. +// The transport error is returned rather than asserted: half of this contract +// is about which refusal comes back. +func submitPost(client *testkit.AppView, community, title string) error { + return client.Procedure(context.Background(), "social.coves.community.post.create", map[string]any{ + "community": community, + "title": title, + "content": "submitted through the admission gate", + }, nil) +} + +// requireXRPCRefusal asserts err is an XRPC error envelope with exactly this +// status and error name, and returns it for callers that assert further. The +// name is asserted as well as the status because it is the machine-readable +// half clients switch on: a 403 NotAuthorized tells an aggregator to stop, a +// 403 with any other name tells it nothing. +func requireXRPCRefusal(t *testing.T, err error, status int, code, what string) *testkit.StatusError { + t.Helper() + + var se *testkit.StatusError + require.ErrorAsf(t, err, &se, + "%s must be refused with an XRPC error envelope, got: %v", what, err) + require.Equalf(t, status, se.StatusCode, "%s: answered %v", what, err) + require.Equalf(t, code, se.XRPCError, "%s: answered %v", what, err) + return se +} diff --git a/tests/live/post_unfurl_test.go b/tests/live/post_unfurl_test.go index 7b8ed2f..5cebc3f 100644 --- a/tests/live/post_unfurl_test.go +++ b/tests/live/post_unfurl_test.go @@ -223,6 +223,7 @@ func TestPostUnfurl_UserProvidedMetadata(t *testing.T) { unfurlService, nil, // blueskyService pdsURL, + posts.WithAdmissionPolicy(posts.NewAllowAllAdmissionPolicyForTests()), ) // Create test user and community