diff --git a/FOLLOWUPS.md b/FOLLOWUPS.md index 67fab9c..3127a00 100644 --- a/FOLLOWUPS.md +++ b/FOLLOWUPS.md @@ -481,27 +481,71 @@ and the comparisons it deliberately did not build. never act" should not ship a new drop site — on the highest-volume inbound path, which already carries a recorded perf concern. -- **The re-cast divergence sweep drives off a sequential scan of - `outbound_activities`, the highest-volume table in the system.** Confirmed by - `EXPLAIN` on the real schema: every other leg is index-served (the Undo - subquery rides `idx_outbound_activities_actor`, the delivered-delivery joins - ride `idx_outbound_deliveries_activity`, the ledger exclusion rides - `outbound_votes_actor_delivered_idx`), but the outer driver is a full scan - filtered on `kind IN ('Like','Dislike') AND parent_at_uri <> ''`, on every - sweep. Deliberately NOT fixed in 17e — adding an index there is a migration - whose write cost lands on the delivery worker's hot path, and the sweep runs - every 15 minutes against a table that is currently small. The fix when it is - needed: - +- **TWO legs of the divergence sweep scan `outbound_activities` in full, not + one.** This entry previously said the re-cast driver was the only unindexed + leg and that "every other leg is index-served". THAT WAS WRONG, and wrong in + the direction that stops the next person looking: the undelivered-acceptance + leg joins `outbound_activities` on a JSONB EXPRESSION, and no index in + `internal/db/migrations/` can serve it. `outbound_activities` carries exactly + two indexes — `outbound_activities_pkey (activity_id)` and + `idx_outbound_activities_actor (actor_did)` — and neither is on + `payload -> 'object' ->> 'id'`. + + Re-confirmed by `EXPLAIN` against the migrated test schema. The + acceptance leg: + + -> Hash Right Join + Hash Cond: (((a.payload -> 'object'::text) ->> 'id'::text) = o.ap_object_id) + -> Seq Scan on outbound_activities a + -> Hash + -> Seq Scan on outbound_objects o + Filter: (accepted_at IS NULL) + + and the proof that it is unindexed rather than merely cheap here: with + `SET enable_seqscan = off` — which prices a sequential scan at 1e10 — the + planner STILL chooses `Seq Scan on outbound_activities a` for that condition, + because there is nothing else it could use. + + The re-cast leg's original finding stands: its three exclusion subqueries are + index-served (the Undo and superseding-vote subqueries ride + `idx_outbound_activities_actor`, the delivered-delivery joins ride + `idx_outbound_deliveries_activity`, the ledger exclusion rides + `outbound_votes_actor_delivered_idx`), while the outer driver is a full scan + filtered on `kind IN ('Like','Dislike') AND parent_at_uri <> ''`. + + COST NOTE, so the next profile is not a surprise: bounding the sweep put an + exact `COUNT(*)` beside each limited read, so the acceptance leg's full scan + of `outbound_activities` now happens TWICE per sweep rather than once. That is + the price of a count an operator can size an incident from; it is also the + strongest argument for the expression index below, since indexing that join + fixes both statements at once. + + Deliberately NOT fixed here — an index is a migration whose write cost lands + on the delivery worker's hot path, and the sweep runs every 15 minutes against + tables that are currently small. The two candidates, when it is needed: + + -- serves the re-cast driver, and tightens the Undo subquery CREATE INDEX outbound_activities_vote_subject_idx ON outbound_activities (actor_did, parent_at_uri, created_at) WHERE kind IN ('Like','Dislike'); - It serves the driver and tightens the Undo subquery, which today index-scans - by `actor_did` and applies kind/parent_at_uri/created_at as residuals; a - sibling partial index on `kind = 'Undo'` finishes that half. Row estimates - from a near-empty test database are meaningless — the plan SHAPE is the - finding, so re-profile against production volumes before choosing. + -- serves the undelivered-acceptance join, which has nothing today + CREATE INDEX outbound_activities_object_id_idx + ON outbound_activities ((payload -> 'object' ->> 'id')) + WHERE payload -> 'object' ->> 'id' IS NOT NULL; + + The first also tightens the Undo subquery, which today index-scans by + `actor_did` and applies kind/parent_at_uri/created_at as residuals; a sibling + partial index on `kind = 'Undo'` finishes that half. The second is an + EXPRESSION index and must be spelled with exactly the operators the query uses + (`->` then `->>`) or the planner will not match it — and it is the one whose + write cost is least predictable, because every outbound activity has a payload + and the extraction runs on every insert. + + ROW ESTIMATES FROM A NEAR-EMPTY TEST DATABASE ARE MEANINGLESS. Every cost and + row number in the plans above is noise; the plan SHAPE — a full scan of + `outbound_activities` per sweep leg, per statement — is the entire finding. + Re-profile against production volumes before choosing either index. - **An acceptance with NO delivery row at all is not reported by the acceptance-undelivered classes**, and this is stated in the query rather than @@ -511,3 +555,21 @@ and the comparisons it deliberately did not build. `outbound_objects` row, the delivery enqueue and the `admissions` ledger row all land in ONE transaction — so a test asserting it today would be vacuous. If it is ever observed, it wants its own class, not a widened existing one. + +## Found while building 17e's fixtures (NOT a 17e defect) + +- **Re-using a deleted vote record's rkey silently never federates.** The + outbound activity id derives from (at-uri, "create", seq). Delete a vote row + and the seq restarts at 1, so a vote re-created under the SAME rkey + reproduces the FIRST vote's activity id. `InsertTx` then no-ops on the + existing id, `EnqueueTx` reads back the standing **delivered** row (correctly + — that is the idempotency that makes the fan-out safe), and nothing goes out. + Observed on the wire as `["Dislike","Undo"]` with the third vote simply + missing, no error anywhere. + + Whether this is reachable depends on Coves' rkey policy for re-created votes. + If rkeys are ever reused after a delete, a user's vote silently does not + federate and no signal is produced. Worth confirming against Coves before + deciding whether it needs a fix here; the candidate fix is to derive the + activity id from something that does not restart (the row's create timestamp, + or a monotonic per-actor counter that survives deletion). diff --git a/cmd/tidepool/main.go b/cmd/tidepool/main.go index bd81785..1d7bab4 100644 --- a/cmd/tidepool/main.go +++ b/cmd/tidepool/main.go @@ -491,7 +491,13 @@ func run(logger *slog.Logger) error { divergence, err := ingest.NewDivergenceReconciler(ingest.DivergenceOptions{ DB: database, Interval: cfg.DivergenceInterval, - Logger: logger, + // The staleness window is CONFIGURED rather than left to default, because + // it is the knob that decides whether the acceptance classes cry wolf on + // a slow queue or stay quiet through a stopped one, and an operator who + // has to rebuild the binary to tune it will instead learn to ignore the + // report. + AcceptanceStaleAfter: cfg.DivergenceAcceptanceStaleAfter, + Logger: logger, }) if err != nil { return err diff --git a/internal/config/config.go b/internal/config/config.go index cc83ee8..ae46b8a 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -14,6 +14,7 @@ import ( "strings" "time" + "tidepool/internal/ingest" "tidepool/internal/personas" ) @@ -255,6 +256,18 @@ type Config struct { // runs one on demand regardless. It never writes anything (decision 19), // which is what makes an always-on schedule safe. DivergenceInterval time.Duration + // DivergenceAcceptanceStaleAfter is how long a pending delivery may sit + // before the sweep reports its acceptance as STALE + // (DIVERGENCE_ACCEPTANCE_STALE_AFTER, a Go duration, must be positive). + // + // It is the report's one crying-wolf knob: too short and every ordinary + // in-flight post is a finding, too long and a queue that stopped this + // morning is not in the report tonight. The default is derived from the + // retry schedule and the causal wait budget — see + // ingest.DefaultAcceptanceStaleAfter, which is NAMED here rather than + // respelled as a number, so the reasoning and the value an operator + // actually gets cannot come apart. + DivergenceAcceptanceStaleAfter time.Duration } // Load reads configuration from the environment. logger must not be nil; @@ -595,6 +608,11 @@ func Load(logger *slog.Logger) (*Config, error) { if err != nil { return nil, err } + cfg.DivergenceAcceptanceStaleAfter, err = durationVar(logger, + "DIVERGENCE_ACCEPTANCE_STALE_AFTER", ingest.DefaultAcceptanceStaleAfter) + if err != nil { + return nil, err + } defaultUserAgent := fmt.Sprintf("tidepool/0.1 (+https://%s)", cfg.BridgeHostname) cfg.UserAgent = os.Getenv("USER_AGENT") diff --git a/internal/ingest/divergence.go b/internal/ingest/divergence.go index da55958..627059a 100644 --- a/internal/ingest/divergence.go +++ b/internal/ingest/divergence.go @@ -6,6 +6,8 @@ import ( "expvar" "fmt" "log/slog" + "sync" + "sync/atomic" "time" "tidepool/internal/errors" @@ -32,6 +34,34 @@ import ( // defaultDivergenceInterval is the sweep cadence when options leave it zero. const defaultDivergenceInterval = 15 * time.Minute +// divergenceSweepTimeout bounds ONE pass, and it exists because of the pool +// rather than because of the clock. +// +// Every comparison here is a multi-table join taken from the SHARED connection +// pool — the same 25 connections serving inbound ingestion and the delivery +// worker — and there is no statement_timeout anywhere in this codebase. A +// pathological sweep with no deadline therefore does not merely run late: it +// holds pool connections while it runs, and the sweep is slowest exactly when +// the database is struggling, which is when federation most needs those +// connections. GET /admin/divergence makes that reachable from outside on an +// unrated GET, so an operator refreshing the report during an incident can +// deepen it. +// +// GENEROUS, because this is a background comparison and not a scrape: the +// precedent is consume's deadLetterScrapeTimeout (internal/consume/metrics.go), +// three seconds for a read taken on the metrics handler, and the reasoning +// inverts here. Nobody is blocked on a sweep, a real one on a large database +// legitimately takes seconds to a minute, and a deadline that fires on a merely +// slow pass would report failures that are not divergences. Two minutes is well +// past any healthy pass and well inside the 15-minute cadence, so a sweep can +// never queue behind its own predecessor. +// +// WHAT IT BUYS is that the failure is LOUD and BOUNDED: the pass aborts, the +// failure counter moves, the age gauge starts climbing, and the previous gauges +// keep standing — the same treatment any other failed read gets — instead of a +// query pinning a connection for as long as the database will let it. +const divergenceSweepTimeout = 2 * time.Minute + // DefaultAcceptanceStaleAfter is how long a delivery may sit pending before its // absence from the peer is worth an operator's attention. // @@ -62,11 +92,19 @@ const ( DivergencePersonaVoteEvent = "persona-vote-event" // DivergenceVoteRecastUndelivered: a peer is holding a vote this bridge no - // longer claims. Re-casting a delivered vote resets the row to pending under - // a new activity id; when that new delivery poisons, the peer keeps counting - // the OLD vote while our ledger accounts for nothing — permanently, since - // nothing re-drives a poisoned delivery and the reseed subtracts only - // delivered rows. + // longer claims — the peer accepted the activity, and nothing in our + // accounting covers it: no live delivered vote row names it, no delivered + // Undo withdraws it, and no later delivered vote supersedes it. + // + // The NAME says re-cast because that is the common cause, but the class is + // deliberately wider than its name and THREE populations reach it, with + // different remedies: a re-cast whose new delivery poisoned; an Undo that + // poisoned (the withdrawal failed, not the vote); and purge residue, where + // the destructive opt-out tier marks a vote `undone` at decision time while + // keeping current_activity_id. Whichever it is, it is permanent — nothing + // re-drives a poisoned delivery, and the reseed subtracts only delivered + // rows. The Detail states only what the row proves and leaves the cause to + // the operator's reading of the ledger; do not narrow it back to a re-cast. DivergenceVoteRecastUndelivered = "vote-recast-undelivered" // DivergenceDeliveryUnknownRefused / …Unanswered: a delivery we SENT and @@ -114,6 +152,15 @@ const ( MetricDivergenceVoteRecast = "tidepool_divergence_vote_recast_undelivered" MetricDivergenceUnknownRefused = "tidepool_divergence_delivery_unknown_refused" MetricDivergenceUnknownUnanswered = "tidepool_divergence_delivery_unknown_unanswered" + + // The FRESHNESS pair. Every gauge above reports a count that a sweep last + // wrote, and a sweep that has stopped running leaves all seven standing at + // their final values — a low, stable, entirely healthy-looking set. These two + // are how an operator tells "nothing is diverging" from "nothing is + // measuring": the age says when the standing numbers were established, and + // the failure count says whether the sweep has been trying and losing. + MetricDivergenceSweepAgeSeconds = "tidepool_divergence_sweep_age_seconds" + MetricDivergenceSweepFailures = "tidepool_divergence_sweep_failures" ) // divergenceUnswept is what a gauge reads before any sweep has completed. @@ -151,9 +198,20 @@ var ( metricUnknownUnanswered = newDivergenceGauge(MetricDivergenceUnknownUnanswered) ) -// divergenceGauges maps each class to the gauge that reports it, so publish -// cannot set one and forget another: a class added to the report without a -// gauge here fails to compile at the map literal rather than going unwatched. +// divergenceGauges maps each class to the gauge that reports it, and is the +// SINGLE DECLARATION of the sweep's class vocabulary: newDivergenceReport +// derives the report's Counts keys from this map, so the two key sets are one +// set by construction and cannot drift. +// +// They used to be two literals, under a comment claiming a class added without +// a gauge "fails to compile at the map literal". IT DOES NOT — a +// map[string]*expvar.Int literal is not exhaustive over any key set, nothing in +// Go compares two literals against each other, and a missing key read out of +// Counts yields 0, which is HEALTH. The claim was worse than the hole it +// described, because it told every later reader not to check. What actually +// holds the two together now is this derivation, and +// TestEveryDivergenceClassHasAGauge pins it against the report a sweep builds. +// Adding a class means adding it HERE and nowhere else. var divergenceGauges = map[string]*expvar.Int{ DivergencePersonaVoteEvent: metricPersonaVoteEvents, DivergenceAcceptanceCancelled: metricAcceptanceCancelled, @@ -170,6 +228,72 @@ func newDivergenceGauge(name string) *expvar.Int { return gauge } +// metricSweepFailures counts sweeps that aborted on a read error. +// +// It starts at 0 and 0 is an HONEST zero here, unlike the class gauges: a +// process that has never failed a sweep has genuinely failed none. What a +// counter alone cannot say is whether any sweep has RUN — zero failures is also +// what a scheduler that never started reports — which is what the age gauge +// below is for. The two are read together or neither means anything. +var metricSweepFailures = expvar.NewInt(MetricDivergenceSweepFailures) + +// The publication guard. It protects the seven process-global gauges above, NOT +// the database reads (see DivergenceReconciler for why those need nothing). +var ( + // divergenceSweepGeneration numbers sweeps in the order they START, across + // every reconciler in the process, because the gauges they publish to are + // process-global too. + divergenceSweepGeneration atomic.Uint64 + + divergencePublishMu sync.Mutex + // divergencePublishedGeneration is the generation of the newest sweep whose + // numbers are STANDING in the gauges. A sweep whose generation is not above + // it started earlier than the values on display and is dropped rather than + // published. Guarded by divergencePublishMu. + divergencePublishedGeneration uint64 + // divergenceLastPublished is when those standing numbers were established, + // and the zero value means no sweep has ever published. It moves only when a + // publication actually lands: a dropped stale sweep succeeded at reading, but + // its numbers are not the ones being displayed, so counting it as freshness + // would date the gauges to a pass whose results were thrown away. Guarded by + // divergencePublishMu. + divergenceLastPublished time.Time +) + +// The age gauge, and the ONE place expvar.Func is right in this file. +// +// The seven class gauges are deliberately pushed rather than pulled because +// they are multi-table joins and a Func would run them on the scrape — on the +// endpoint an operator reads BECAUSE the system is struggling. This one reads a +// stored timestamp and subtracts, which costs nothing, and it must be pulled: +// an age pushed at publication time would be the one number that stops updating +// at exactly the moment it becomes the number that matters. That is +// consume.PublishMetrics' reasoning about cursor age (internal/consume/ +// metrics.go), applied to the failure it does not cover — there, a value goes +// stale when the consumer stalls; here, ALL SEVEN GAUGES go stale together when +// the sweep stops, and they freeze at plausible, low, healthy-looking values. A +// permissions change or a statement timeout at 04:00 leaves them saying "0 +// divergences" for the life of the process while the queue stops entirely. +// +// It reports divergenceUnswept before the first successful publication, for the +// same reason the class gauges do: a 0 there would say "swept just now". +func init() { + expvar.Publish(MetricDivergenceSweepAgeSeconds, expvar.Func(func() any { + divergencePublishMu.Lock() + published := divergenceLastPublished + divergencePublishMu.Unlock() + if published.IsZero() { + return float64(divergenceUnswept) + } + age := time.Since(published).Seconds() + if age < 0 { + // Clock movement, not a sweep from the future. + return 0.0 + } + return age + })) +} + // DivergenceOptions configures a DivergenceReconciler. type DivergenceOptions struct { // DB is the bridge database. BOTH sides of every comparison are local @@ -194,12 +318,25 @@ type DivergenceOptions struct { // REPORTS what disagrees. It never writes: not to remote instances, and not to // our own tables (PLAN.md decision 19). // -// There is NO MUTEX, unlike FollowReconciler. That one serializes because two -// concurrent sweeps could each see a community as absent and subscribe it -// twice, and the racing EnsureCommunity calls can mint a permanent DID twice. -// Nothing here writes anything, so concurrent sweeps can only read the same rows -// and reach the same answer; a lock would buy nothing and would make the admin -// report queue behind a background pass. +// THE READS ARE UNSERIALIZED, unlike FollowReconciler's. That one holds a lock +// across its whole pass because two concurrent sweeps could each see a community +// as absent and subscribe it twice, and the racing EnsureCommunity calls can +// mint a permanent DID twice. Nothing here writes anything, so concurrent +// passes can only read the same rows and reach the same answer — and a lock +// spanning the reads would make an operator's on-demand report queue behind a +// background pass, which is the request that is most urgent. +// +// THE PUBLICATION IS SERIALIZED, and that is a different hazard entirely: the +// seven gauges are process-global expvar.Ints shared by the Run loop and every +// GET /admin/divergence. Two passes assigning them interleave in randomized map +// order, so a scrape could read a set woven from two sweeps — numbers that never +// described one moment. Worse, sweeps take different times, so a background pass +// that started at T=0 and finished at T=4m could overwrite an on-demand pass +// that started at T=1m and published current state, walking the gauges BACKWARDS +// at precisely the moment an operator is refreshing them during an incident. +// divergencePublishMu makes each publication whole, and the generation compared +// under it drops a pass that started before the one already on display (see +// publish). The lock is held only for the assignment, never across a read. type DivergenceReconciler struct { divergences store.Divergences interval time.Duration @@ -267,51 +404,160 @@ type DivergenceReport struct { // stopped queue — and an unbounded list is then a second outage: a response // nobody can load, on the endpoint an operator reaches for BECAUSE something is // wrong. The counts stay exact; only the examples are bounded. -const MaxDivergenceEntries = 500 +// +// IT IS THE STORE'S CONSTANT, not a second number that happens to match. The +// cap has to be applied as a SQL LIMIT to bound anything (see +// store.MaxDivergenceExamples), so this is the same budget spelled where the +// report talks about it; two constants would drift, and the drift would be +// invisible — a larger budget here would just never be reached, and a smaller +// one would silently throw away rows the database was asked for. +const MaxDivergenceEntries = store.MaxDivergenceExamples + +// divergenceExamples stages one sweep's examples per class, so the report's +// budget can be shared out rather than handed to whoever read first. +// +// WITHOUT IT, CLASS ORDER IS THE ALLOCATION. Entries were appended in the order +// the comparisons run, so one bulk class — a stopped queue, a broken community, +// the very situations this report exists for — consumed the entire budget and +// every later class arrived with a non-zero count and NOT ONE EXAMPLE. That is +// the exact failure DivergenceEntry.Subject exists to prevent: "a class and a +// count alone give an operator a number they cannot investigate", said of the +// classes that need investigating most, at the moment they need it. +// +// The staging is bounded by construction — each comparison is already limited +// to store.MaxDivergenceExamples rows at the database — so this holds at most +// one page per class and never the divergent population. +type divergenceExamples struct { + // order is the classes in the order they were first seen, which is the + // sweep's fixed comparison order. Deterministic on purpose: allocation that + // depended on Go's randomized map iteration would hand different classes + // their examples on every sweep, and an operator refreshing the report + // would watch subjects appear and disappear with nothing changing. + order []string + byClass map[string][]DivergenceEntry +} + +func newDivergenceExamples() *divergenceExamples { + return &divergenceExamples{byClass: make(map[string][]DivergenceEntry)} +} + +func (e *divergenceExamples) add(entry DivergenceEntry) { + if _, seen := e.byClass[entry.Class]; !seen { + e.order = append(e.order, entry.Class) + } + e.byClass[entry.Class] = append(e.byClass[entry.Class], entry) +} + +// fill deals the staged examples into the report, ONE PER CLASS PER ROUND until +// the budget runs out. +// +// Round-robin rather than an equal share computed up front, because an equal +// share wastes the budget: a sweep where one class has thousands and the others +// have three each would cap the big one at a quarter of the page and leave the +// rest of the page empty. Dealing a card at a time gives every class its +// examples first and then spends everything left on whoever still has rows, so a +// single-class incident still fills the page — which is what makes this +// compatible with the bound the report already promises. +func (e *divergenceExamples) fill(report *DivergenceReport) { + for round := 0; ; round++ { + dealt := false + for _, class := range e.order { + entries := e.byClass[class] + if round >= len(entries) { + continue + } + dealt = true + report.addEntry(entries[round]) + } + if !dealt { + return + } + } +} // addEntry records one example, while there is room for one. // -// The cap is applied HERE rather than by truncating a finished list, because a -// list that must first be built in full is not bounded at all: the sweep that -// hits this limit is the sweep running against a database with a stopped queue -// in it, and materialising every row of that before cutting it down would put -// the memory spike in exactly the process an operator is trying to keep alive. +// IT IS THE LAST OF THREE BOUNDS, NOT THE BOUND. The reads are limited at the +// database (store.MaxDivergenceExamples), because a list that must first be +// built in full is not bounded at all: the sweep that reaches this limit is the +// sweep running against a database with a stopped queue in it, and materialising +// every row of that before cutting it down would put the memory spike in exactly +// the process an operator is trying to keep alive. divergenceExamples then +// decides WHICH of those rows get the page. This function only refuses the ones +// that no longer fit. // // COUNTS ARE NOT TOUCHED HERE, deliberately: every caller counts what the -// comparison returned, not what fitted. A cap that also bounded the measurement -// would cap the number an operator escalates on at the size of a page, and -// "500" would then mean both "500" and "a catastrophe". +// comparison MEASURED — an exact COUNT(*) taken beside the limited read — not +// what fitted. A cap that also bounded the measurement would cap the number an +// operator escalates on at the size of a page, and "500" would then mean both +// "500" and "a catastrophe". func (report *DivergenceReport) addEntry(entry DivergenceEntry) { if len(report.Entries) >= MaxDivergenceEntries { // SAY SO. A silently truncated list reads as the whole problem, and an - // operator sizes an incident from what they can see. + // operator sizes an incident from what they can see. Note that this is + // no longer the only way Truncated is set — a comparison whose count + // exceeds the examples it returned truncated at the database, where + // this function never sees the missing rows. See markTruncated. report.Truncated = true return } report.Entries = append(report.Entries, entry) } +// markTruncated says whether the examples are fewer than the findings. +// +// IT IS THE ONLY HONEST TEST NOW THAT THE LIMIT IS IN SQL. addEntry can only +// notice truncation by being handed a row it has no room for, and a comparison +// bounded by `LIMIT 500` never hands over the 501st: the report would carry +// exactly 500 examples for a population of fifty thousand and claim to be +// complete. Comparing the counts — which are exact, by construction — against +// the examples catches every shape of it, including the one addEntry does see. +// +// It only ever SETS the flag. A truncation already noticed must not be cleared +// by a later recount, and the counts and the examples come from separate +// statements, so a row that arrives between them is a reason to be careful in +// one direction only. +func (report *DivergenceReport) markTruncated() { + total := 0 + for _, count := range report.Counts { + total += count + } + if total > len(report.Entries) { + report.Truncated = true + } +} + // newDivergenceReport is an EMPTY report with every class present at zero — and // non-nil throughout, so the JSON is `{"entries":[],"counts":{...}}` rather than // nulls a client has to special-case. +// +// The classes are DERIVED from divergenceGauges rather than listed again here. +// A second literal is a second declaration of the same vocabulary, and two +// literals drift: the class dropped from one of them is the one whose gauge then +// publishes 0 — health — for a comparison nobody is making. Deriving makes the +// two key sets the same set, and leaves exactly one place to edit when a class +// is added. func newDivergenceReport() DivergenceReport { + counts := make(map[string]int, len(divergenceGauges)) + for class := range divergenceGauges { + counts[class] = 0 + } return DivergenceReport{ Entries: []DivergenceEntry{}, - Counts: map[string]int{ - DivergencePersonaVoteEvent: 0, - DivergenceAcceptanceCancelled: 0, - DivergenceAcceptancePoisoned: 0, - DivergenceAcceptanceStale: 0, - DivergenceVoteRecastUndelivered: 0, - DivergenceDeliveryUnknownRefused: 0, - DivergenceDeliveryUnknownUnanswered: 0, - }, + Counts: counts, } } // Run sweeps once immediately, then on every interval tick until ctx is // cancelled. A failed sweep is logged, never fatal: the next tick is the retry, // and there is no state to leave half-applied because there is no state. +// +// The log line is not the only signal any more, and it must not be: this loop is +// the one that runs at 04:00 with nobody reading logs, and a permissions change +// or a statement timeout that fails every pass from then on would otherwise +// leave seven frozen, healthy-looking gauges behind it. Sweep counts each +// failure and the age gauge keeps climbing while they stand (see the metric +// constants), so a stopped sweep is visible on the same endpoint as its results. func (r *DivergenceReconciler) Run(ctx context.Context) { r.logger.Info("divergence reconciler started", "interval", r.interval) if _, err := r.Sweep(ctx); err != nil && ctx.Err() == nil { @@ -340,16 +586,56 @@ func (r *DivergenceReconciler) Run(ctx context.Context) { // happened to succeed would publish a claim about state nobody read — and the // gauges would then be SET to a number that quietly excludes it. The error goes // back instead, the gauges keep their previous values (see publish), and the -// operator sees a failed sweep rather than a clean one. -func (r *DivergenceReconciler) Sweep(ctx context.Context) (DivergenceReport, error) { - report := newDivergenceReport() - - personaVotes, err := r.divergences.PersonaVoteEvents(ctx) +// operator sees a failed sweep rather than a clean one. A DELIVERY STATE WITH NO +// CLASS aborts it the same way, for the same reason: see the acceptance loop. +// +// A FAILED PASS IS COUNTED, once, wherever it aborted. Leaving the previous +// values standing is right, but on its own it is silent: a sweep that fails +// forever leaves seven plausible numbers frozen and a log line nobody is +// watching. The counter and the age gauge are what distinguish those frozen +// numbers from a healthy bridge. The count is deferred rather than written at +// each return so a comparison added later cannot forget it, and it is skipped +// when ctx is already done because a sweep cut short by shutdown is not a +// failing sweep — inflating the counter on every restart would make the number +// mean "restarts plus failures", which is a number an operator cannot act on. +func (r *DivergenceReconciler) Sweep(ctx context.Context) (report DivergenceReport, err error) { + // THE GUARD READS THE PARENT CONTEXT, NOT THE DEADLINED ONE BELOW. A sweep + // that ran out of its own time is a FAILING sweep and must be counted: it is + // the pathological pass this timeout exists to cut short, and the counter is + // how anyone learns it happened. Only a sweep cut short by SHUTDOWN is + // exempt, and that is the parent's cancellation, which is what this reads. + defer func() { + if err != nil && ctx.Err() == nil { + metricSweepFailures.Add(1) + } + }() + + // Taken BEFORE the reads, so generations order sweeps by when they started + // looking at the database — which is what makes a slow pass detectably older + // than a fast one that started after it. Taking it at publication time would + // order them by finish and defeat the whole guard. + generation := divergenceSweepGeneration.Add(1) + report = newDivergenceReport() + // The examples are STAGED per class and dealt out at the end, so a bulk + // class cannot spend the whole page before the later comparisons have run. + // See divergenceExamples. + examples := newDivergenceExamples() + + // ONE DEADLINE FOR THE WHOLE PASS, not one per read: what has to be bounded + // is how long this sweep holds pool connections in total, and eight reads + // with their own two-minute budgets would bound nothing. It is derived from + // the caller's context so a shutdown still cancels immediately, and it + // applies to the on-demand endpoint as well as the loop — a request context + // carries whatever deadline the client chose, which for a curl is none. + sweepCtx, cancel := context.WithTimeout(ctx, divergenceSweepTimeout) + defer cancel() + + personaVotes, err := r.divergences.PersonaVoteEvents(sweepCtx) if err != nil { return DivergenceReport{}, fmt.Errorf("divergence: read persona vote events: %w", err) } for _, vote := range personaVotes { - report.addEntry(DivergenceEntry{ + examples.add(DivergenceEntry{ Class: DivergencePersonaVoteEvent, Subject: vote.ActivityID, Detail: fmt.Sprintf( @@ -359,40 +645,87 @@ func (r *DivergenceReconciler) Sweep(ctx context.Context) (DivergenceReport, err vote.SubjectAPID, vote.VoterAPID, vote.ActorDID), }) } - // The COUNT is the comparison's length, never len(Entries): see addEntry. - report.Counts[DivergencePersonaVoteEvent] = len(personaVotes) + // THE COUNT IS ITS OWN MEASUREMENT, never len(entries) and never + // len(Entries): the list above is bounded by store.MaxDivergenceExamples, + // so its length is the size of a page rather than the size of the problem. + // See addEntry. + personaTotal, err := r.divergences.PersonaVoteEventCount(sweepCtx) + if err != nil { + return DivergenceReport{}, fmt.Errorf("divergence: count persona vote events: %w", err) + } + report.Counts[DivergencePersonaVoteEvent] = personaTotal - undelivered, err := r.divergences.UndeliveredAcceptances(ctx, r.staleAfter) + undelivered, err := r.divergences.UndeliveredAcceptances(sweepCtx, r.staleAfter) if err != nil { return DivergenceReport{}, fmt.Errorf("divergence: read undelivered acceptances: %w", err) } for _, acceptance := range undelivered { class := acceptanceClass(acceptance.DeliveryState) if class == "" { - // A delivery state this sweep has no class for. Logged rather than - // filed under the nearest class: putting it in one would describe it - // wrongly, and dropping it silently would make an unclassifiable - // state look like a healthy one. - r.logger.Warn("divergence: undelivered acceptance in an unclassified delivery state", - "post", acceptance.PostURI, "state", acceptance.DeliveryState) - continue + // A delivery state this sweep has no class for ABORTS THE PASS, the + // same treatment a read failure gets, and for the same reason: it + // means our model of the schema is stale, so every number this sweep + // is about to publish was computed by code that does not know what it + // is looking at. + // + // IT USED TO WARN AND SKIP, which dropped the row from Entries AND + // from every Counts key — the report then said three undelivered + // acceptances when there were four, with the missing one visible only + // in a log line, in a sweep whose entire premise is that logs are not + // where an operator looks. A silent undercount is the one failure this + // file cannot tolerate, because it is indistinguishable from health. + // + // A VISIBLE CLASS WAS THE ALTERNATIVE and it is worse here: a + // catch-all class needs a gauge (divergenceGauges is the vocabulary), + // and a gauge named for "states we do not understand" is a number + // nobody can alert on or act on — while the aborted sweep already has + // an operator-visible surface that says exactly this, the failure + // counter plus the climbing age gauge, and leaves the last real + // numbers standing rather than replacing them with numbers derived + // from a schema we have misread. + // + // Unreachable today: the store's filter admits only cancelled, + // poisoned and pending. It fires the day someone adds a delivery + // state — which is precisely the change that must not land quietly. + return DivergenceReport{}, fmt.Errorf( + "divergence: undelivered acceptance %s is in delivery state %q, which this sweep "+ + "has no class for: the state vocabulary has changed and acceptanceClass has not", + acceptance.PostURI, acceptance.DeliveryState) } - report.addEntry(DivergenceEntry{ + examples.add(DivergenceEntry{ Class: class, Subject: acceptance.PostURI, - Detail: fmt.Sprintf( - "community %s accepted this post and the peer was never told: its delivery is %s", - acceptance.CommunityDID, acceptance.DeliveryState), + Detail: acceptanceDetail(acceptance), }) - report.Counts[class]++ + } + // COUNTED BY DELIVERY STATE, and mapped to classes through the SAME + // function the entries are classified with, so the count and the example + // for one post can never disagree about which class it belongs to. + acceptanceCounts, err := r.divergences.UndeliveredAcceptanceCounts(sweepCtx, r.staleAfter) + if err != nil { + return DivergenceReport{}, fmt.Errorf("divergence: count undelivered acceptances: %w", err) + } + for state, count := range acceptanceCounts { + class := acceptanceClass(state) + if class == "" { + // The same abort as above, at the count, and the count is where the + // undercount would have done its real damage: an operator sizes an + // incident from these numbers, and a class total that quietly omits a + // state reads as a smaller problem rather than an unknown one. + return DivergenceReport{}, fmt.Errorf( + "divergence: %d undelivered acceptances are in delivery state %q, which this "+ + "sweep has no class for: the state vocabulary has changed and acceptanceClass "+ + "has not", count, state) + } + report.Counts[class] += count } - recasts, err := r.divergences.RecastDivergences(ctx) + recasts, err := r.divergences.RecastDivergences(sweepCtx) if err != nil { return DivergenceReport{}, fmt.Errorf("divergence: read recast divergences: %w", err) } for _, recast := range recasts { - report.addEntry(DivergenceEntry{ + examples.add(DivergenceEntry{ Class: DivergenceVoteRecastUndelivered, // The SUBJECT is what an operator checks the tally of; the delivered // activity id is what the peer is actually holding, and the only @@ -400,16 +733,32 @@ func (r *DivergenceReconciler) Sweep(ctx context.Context) (DivergenceReport, err // the subject alone say a number is wrong without saying which vote // is making it wrong. Subject: recast.SubjectATURI, + // THE DETAIL STATES WHAT THE ROW PROVES AND NOTHING ELSE. It used to + // assert a re-cast happened, and the query does not establish that: + // at least three populations reach this class — a re-cast whose new + // delivery poisoned, an Undo that poisoned, and purge residue + // (Purger.undoLiveVotes sets delivered_state='undone' at DECISION + // time while keeping current_activity_id, so the pair is reported + // from that moment until the Undo lands, and forever if it does + // not). Naming one cause sends an operator to check + // the wrong thing on two of the three, and a Detail is a claim in + // exactly the way a metric name is. What the row DOES prove is the + // delivery and the absence, so say that and let the ledger say why. Detail: fmt.Sprintf( - "the peer still holds vote activity %s cast by %s, which this bridge no longer "+ - "claims: the row was re-cast under a new activity id and that delivery never "+ - "landed, so nothing in the ledger accounts for the vote they are counting", + "the peer accepted vote activity %s cast by %s and nothing in this bridge's "+ + "accounting covers it: no live delivered vote row names that activity and no "+ + "delivered Undo or later delivered vote for the pair supersedes it. Read the "+ + "activity's ledger row for the cause", recast.DeliveredActivityID, recast.ActorDID), }) } - report.Counts[DivergenceVoteRecastUndelivered] = len(recasts) + recastTotal, err := r.divergences.RecastDivergenceCount(sweepCtx) + if err != nil { + return DivergenceReport{}, fmt.Errorf("divergence: count recast divergences: %w", err) + } + report.Counts[DivergenceVoteRecastUndelivered] = recastTotal - unknown, err := r.divergences.UnknownDeliveryOutcomes(ctx) + unknown, err := r.divergences.UnknownDeliveryOutcomes(sweepCtx) if err != nil { return DivergenceReport{}, fmt.Errorf("divergence: read unknown delivery outcomes: %w", err) } @@ -424,17 +773,19 @@ func (r *DivergenceReconciler) Sweep(ctx context.Context) (DivergenceReport, err detail := fmt.Sprintf( "we sent this %s to %s and got no answer at all (%s): it may never have arrived, or "+ "it may have been applied and the response lost — %s is the instance to ask, "+ - "because we cannot tell from here", - outcome.Kind, outcome.TargetInbox, outcome.LastErrorClass, outcome.TargetInbox) + "because we cannot tell from here.%s", + outcome.Kind, outcome.TargetInbox, outcome.LastErrorClass, outcome.TargetInbox, + unknownAcceptanceOverlap) if outcome.Refused { class = DivergenceDeliveryUnknownRefused detail = fmt.Sprintf( "we sent this %s to %s and the peer answered %d (%s): that is evidence it was not "+ "applied and not proof of it, since a peer can apply an activity and then "+ - "fail to respond", - outcome.Kind, outcome.TargetInbox, outcome.LastStatusCode, outcome.LastErrorClass) + "fail to respond.%s", + outcome.Kind, outcome.TargetInbox, outcome.LastStatusCode, outcome.LastErrorClass, + unknownAcceptanceOverlap) } - report.addEntry(DivergenceEntry{ + examples.add(DivergenceEntry{ // The SUBJECT is the activity id: what was sent is the only handle // that exists here. There is no post at-uri to point at — these are // every kind of activity, and the acceptance classes are where a post @@ -443,16 +794,35 @@ func (r *DivergenceReconciler) Sweep(ctx context.Context) (DivergenceReport, err Subject: outcome.ActivityID, Detail: detail, }) - report.Counts[class]++ } - - r.publish(report) + // The same split, counted rather than listed. It is read as a pair from one + // grouped statement so the two sub-counts cannot come from two different + // spellings of "the peer answered" — which is precisely how a dial timeout + // would end up in the refused bucket. + unknownCounts, err := r.divergences.UnknownDeliveryOutcomeCounts(sweepCtx) + if err != nil { + return DivergenceReport{}, fmt.Errorf("divergence: count unknown delivery outcomes: %w", err) + } + report.Counts[DivergenceDeliveryUnknownRefused] = unknownCounts.Refused + report.Counts[DivergenceDeliveryUnknownUnanswered] = unknownCounts.Unanswered + + // The staged examples become the report's page here, after every comparison + // has had its say — dealing them out earlier would be the class-order + // allocation this staging exists to remove. + examples.fill(&report) + // And last, the counts against the examples: the reads are bounded in SQL, + // so this is where a report learns it is smaller than the problem it + // describes. + report.markTruncated() + + r.publish(generation, report) return report, nil } // acceptanceClass maps a delivery state to the class an operator triages on. -// An unknown state yields "" — the caller reports that rather than guessing, -// because a state nobody has a response for is itself worth seeing. +// An unknown state yields "" and the caller ABORTS THE SWEEP on it rather than +// guessing: a state nobody has a response for means the schema has moved under +// this sweep, and a report built on that is a set of numbers nobody can trust. func acceptanceClass(state store.DeliveryState) string { switch state { case store.DeliveryStateCancelled: @@ -466,16 +836,160 @@ func acceptanceClass(state store.DeliveryState) string { } } +// unknownAcceptanceOverlap is the sentence that stops an operator adding two +// gauges that count one row. +// +// A poisoned delivery carrying an accepted post is reported TWICE by one sweep: +// under acceptance-undelivered-poisoned keyed by the post at-uri, and here +// keyed by the activity id. Both entries are wanted — see acceptanceDetail for +// why neither side is dropped — but only if the report says they are the same +// row, because two entries with two subjects and two counts otherwise read as +// two problems, and the sum reads as a total. +// +// It is stated CONDITIONALLY because this read cannot tell: the unknown classes +// cover every activity kind (votes, Undos, comments), and only a Create whose +// object is an accepted, unstamped post has a twin. Asserting the twin exists +// would be the same overclaim this whole class avoids. +const unknownAcceptanceOverlap = " If this activity carried a post its community had already " + + "accepted, the same delivery is reported again under acceptance-undelivered-poisoned, keyed " + + "by the post's at-uri: one row asked two questions, so the two counts must not be added." + +// acceptanceDetail is what the report SAYS about an undelivered acceptance, and +// only ONE of the three states supports a definite claim. +// +// A DETAIL IS A CLAIM, exactly as a class name is. The comparison behind all +// three entries is `accepted_at IS NULL`, and that column is stamped only by +// delivery success (worker.stampAccepted), so its absence proves that success +// was never SETTLED HERE — never that the peer lacks the post: +// +// cancelled — definite, and earned. The row left the queue without a POST (a +// ban, an opt-out), so the peer really was never told and nothing will ever +// carry it. +// poisoned — NOT definite. The activity reached the wire and the answer is +// unavailable, which is the entire justification for the delivery-unknown-* +// classes forty lines up. A Detail asserting non-delivery here would +// contradict those classes about the same row, in the same report. +// stale — NOT definite either. A pending delivery past the window has +// usually been ATTEMPTED (attempts is what drives the backoff), and a row +// that has been attempted is silent about whether any attempt landed. What +// is certain is that the queue is not finishing with it. +// +// THE POISONED CASE ALSO NAMES ITS OVERLAP. One poisoned acceptance Create is +// reported twice by this sweep — here, keyed by the post at-uri, and again +// under delivery-unknown-refused/unanswered keyed by the activity id — because +// the two comparisons ask different questions about the same row: "is this +// accepted post confirmed on the peer?" and "what became of this delivery?". +// They are stated as an overlap rather than resolved by dropping one, for three +// reasons. The unknown classes are the ONLY report of a poisoned vote or Undo, +// so excluding acceptance-carrying activities from them would leave those +// counts meaning "everything except posts" — a stranger claim than the overlap. +// Suppressing this entry instead would take the post's at-uri, its community, +// and the redrive action out of the report, leaving an activity id an operator +// cannot map back to a post. And excluding either way requires one query to +// re-derive the other's join — accepted admission, unstamped object, latest +// delivery, staleness window — which is the second-query-for-one-condition +// drift the class vocabulary is already defended against (see the lapsed-ban +// test). What the overlap really costs is an operator ADDING the gauges, so +// both Details say plainly that these two are one row. +func acceptanceDetail(acceptance store.UndeliveredAcceptance) string { + switch acceptance.DeliveryState { + case store.DeliveryStateCancelled: + return fmt.Sprintf( + "community %s accepted this post and the peer was never told: its delivery (%s) was "+ + "CANCELLED before any POST — a decision, a ban or an opt-out, took it out of the "+ + "queue — so the peer does not have it and nothing will ever carry it. This is the "+ + "one class here whose non-delivery is a fact rather than an inference", + acceptance.CommunityDID, acceptance.ActivityID) + case store.DeliveryStatePoisoned: + return fmt.Sprintf( + "community %s accepted this post and its delivery (%s) is POISONED (%s): accepted_at "+ + "was never stamped, which records only that success was never settled HERE. "+ + "Whether the peer applied it is not knowable from this bridge — the POST may have "+ + "been refused, may never have arrived, or may have been applied with the response "+ + "lost. The SAME delivery is reported again under delivery-unknown-refused or "+ + "delivery-unknown-unanswered, keyed by that activity id and carrying the inbox to "+ + "ask: two questions about one row, not two findings, so those counts and this one "+ + "must not be added together", + acceptance.CommunityDID, acceptance.ActivityID, acceptance.LastErrorClass) + case store.DeliveryStatePending: + return fmt.Sprintf( + "community %s accepted this post and its delivery (%s) is still PENDING far past the "+ + "staleness window: the QUEUE has stopped moving, and that investigation starts at "+ + "the worker rather than at the community. accepted_at was never stamped, which "+ + "records only that no delivery has settled HERE — an attempt that went out and was "+ + "never confirmed looks identical on these columns to one that never left — so "+ + "whether the peer holds this post is not knowable from here", + acceptance.CommunityDID, acceptance.ActivityID) + default: + // Unreachable while acceptanceClass gates the caller: a state with no + // class aborts the sweep before a Detail is ever composed. Kept + // non-committal anyway, because the failure this file exists to prevent + // is a sentence that claims more than the row establishes. + return fmt.Sprintf( + "community %s accepted this post and its delivery (%s) is in state %q, which this "+ + "sweep has no reading for: what the peer holds is not knowable from here", + acceptance.CommunityDID, acceptance.ActivityID, acceptance.DeliveryState) + } +} + // publish assigns the gauges from a COMPLETED sweep. // // It runs only on success, and that is the discipline the whole gauge design // rests on: a failed sweep must leave the previous values standing rather than // write a zero, because a zero written by a sweep that read nothing is a claim // that everything is fine, made at the moment nobody could check. -func (r *DivergenceReconciler) publish(report DivergenceReport) { +// +// THE WHOLE ASSIGNMENT IS ONE CRITICAL SECTION, and an older sweep's numbers are +// dropped rather than written over a newer sweep's (see DivergenceReconciler for +// the hazard). generation is the caller's sweep number, taken before its reads. +// +// A CLASS WITH NO COUNT IS A SENTINEL, NEVER A ZERO. `report.Counts[class]` +// alone reads a missing key as 0 — the value that says "this invariant is +// holding" — so a comparison dropped from Sweep, or a gauge whose class the +// report does not carry, would broadcast health from a pass that measured +// nothing. Deriving Counts from this same map makes that unreachable by +// construction; the check stays because the failure it prevents is the one this +// whole file exists to prevent, and a construction can be undone by an edit that +// looks harmless. +func (r *DivergenceReconciler) publish(generation uint64, report DivergenceReport) { + divergencePublishMu.Lock() + defer divergencePublishMu.Unlock() + + if generation <= divergencePublishedGeneration { + // A slower pass that started earlier finishing after a newer one. Its + // report still goes back to its own caller — it was true when it was + // read — but writing it here would move the gauges backwards. + r.logger.Info("divergence: dropping a stale sweep's gauge publication", + "sweep", generation, "published", divergencePublishedGeneration, + "entries", len(report.Entries)) + return + } + divergencePublishedGeneration = generation + divergenceLastPublished = time.Now() + for class, gauge := range divergenceGauges { - gauge.Set(int64(report.Counts[class])) + count, counted := report.Counts[class] + if !counted { + gauge.Set(divergenceUnswept) + r.logger.Error("divergence: sweep produced no count for a gauged class, so its gauge "+ + "reports unswept rather than zero: a comparison for this class is missing from "+ + "Sweep, and zero would claim the invariant is holding", + "class", class, "sweep", generation) + continue + } + gauge.Set(int64(count)) } + for class, count := range report.Counts { + if _, gauged := divergenceGauges[class]; !gauged { + // The mirror failure: a comparison whose findings reach the report and + // no gauge, so an operator alerting on the gauges never sees it. There + // is nothing to set — say so loudly instead. + r.logger.Error("divergence: sweep counted a class with no gauge, so nothing an "+ + "operator alerts on reports it", + "class", class, "count", count, "sweep", generation) + } + } + if len(report.Entries) > 0 { r.logger.Warn("divergence sweep found disagreements", "entries", len(report.Entries), "counts", report.Counts) diff --git a/internal/ingest/divergence_endpoint_test.go b/internal/ingest/divergence_endpoint_test.go index 1796d70..dbf6363 100644 --- a/internal/ingest/divergence_endpoint_test.go +++ b/internal/ingest/divergence_endpoint_test.go @@ -422,3 +422,436 @@ func TestDivergenceRunSweepsImmediately(t *testing.T) { "watching the gauges cannot tell a fresh start from a stalled sweep") } } + +// --------------------------------------------------------------------------- +// 17e review — a class with no gauge publishes health it never measured +// --------------------------------------------------------------------------- + +// TestEveryDivergenceClassHasAGauge pins the two key sets against each other. +// +// publish() reads report.Counts[class] for each gauge, and a MISSING KEY yields +// zero — so a class whose gauge was forgotten does not go unwatched, it +// publishes 0: health, for a comparison nobody wired up. The reverse is as bad: +// a gauge with no class in the report is set from a map miss on every sweep and +// reads 0 forever, which is a number an operator can watch indefinitely while it +// measures nothing. +// +// The old comment on divergenceGauges claimed a missing entry "fails to compile +// at the map literal". IT DOES NOT — Go map literals are not exhaustive over any +// key set, and nothing in the language checks these two against each other. That +// claim is why nobody wrote this test, which is the failure mode worth naming: +// a compile-time guarantee asserted in prose is a runtime hole with a comment +// over it. +func TestEveryDivergenceClassHasAGauge(t *testing.T) { + classes := make([]string, 0, len(newDivergenceReport().Counts)) + for class := range newDivergenceReport().Counts { + classes = append(classes, class) + } + gauged := make([]string, 0, len(divergenceGauges)) + for class := range divergenceGauges { + gauged = append(gauged, class) + } + + assert.ElementsMatch(t, classes, gauged, + "every class the report counts must have a gauge, and every gauge must have a class. A "+ + "class with no gauge is not merely unwatched: publish reads a missing key as 0 and "+ + "broadcasts HEALTH for a comparison nobody wired. A gauge with no class is set from "+ + "a map miss on every sweep and reads 0 forever while measuring nothing. Nothing in "+ + "Go checks these two literals against each other — the map literal does NOT fail to "+ + "compile, whatever the comment above it says") + + for class, gauge := range divergenceGauges { + require.NotNil(t, gauge, "the gauge for %q must exist", class) + } +} + +// TestDivergenceReportIsNotTruncatedWhenItFits is the negative control for the +// cap, and it is cheap for a reason: an implementation that sets Truncated +// unconditionally passes the bounded test, and then every report an operator +// ever reads claims examples are missing. A flag that is always true carries no +// information, and the one it should have carried is "the number you are looking +// at is smaller than the problem". +func TestDivergenceReportIsNotTruncatedWhenItFits(t *testing.T) { + h := newHarness(t) + world := newModerationWorld(t, h) + h.admin.SetDivergenceReconciler(newDivergenceReconciler(t, h)) + deliverEverythingQueued(t, h.db) + _ = world + + // A handful of real divergences — far below the cap. + personaActorID := actorIDOf(t, h.db, mtAuthorDID) + overflowPersonaVotes(t, h.db, personaActorID, 3) + + report := fetchDivergence(t, h) + require.Len(t, report.Entries, 3, "precondition: a report well inside the cap") + assert.False(t, report.Truncated, + "a report that FITS must not claim it was cut short: a flag that is always set tells an "+ + "operator nothing, and the thing it was supposed to tell them — that the list is "+ + "smaller than the problem — is exactly what they would stop believing") +} + +// TestDivergenceStaleAfterOptionIsHonored pins that the configured window is the +// one actually applied. +// +// The reconciler takes AcceptanceStaleAfter and defaults it; replacing the field +// with the default constant at the call site leaves every other test green, +// because they all age their fixtures far past both values. An option nobody +// reads is worse than no option: an operator tunes it, observes no change, and +// concludes the report is broken in some deeper way. +func TestDivergenceStaleAfterOptionIsHonored(t *testing.T) { + h := newHarness(t) + world := newModerationWorld(t, h) + deliverEverythingQueued(t, h.db) + + // One accepted post whose delivery has been pending for two hours: stale + // under a one-minute window, in flight under a one-day one. + stale := admitDivergencePost(t, h, world, "3lzdvopt00001", "3lzdvopt00011", 1_775_000_090_000_001) + setDeliveryState(t, h.db, stale, "pending", "", 2*time.Hour) + + tight, err := NewDivergenceReconciler(DivergenceOptions{ + DB: h.db, AcceptanceStaleAfter: time.Minute, + }) + require.NoError(t, err) + tightReport, err := tight.Sweep(context.Background()) + require.NoError(t, err) + assert.Equal(t, 1, tightReport.Counts[DivergenceAcceptanceStale], + "with a one-minute window a delivery pending for two hours is stale") + + relaxed, err := NewDivergenceReconciler(DivergenceOptions{ + DB: h.db, AcceptanceStaleAfter: 24 * time.Hour, + }) + require.NoError(t, err) + relaxedReport, err := relaxed.Sweep(context.Background()) + require.NoError(t, err) + assert.Zero(t, relaxedReport.Counts[DivergenceAcceptanceStale], + "and with a one-day window the SAME row is still in flight. If both sweeps agree, the "+ + "option is not reaching the query — an operator who widens the window to quiet a "+ + "noisy report would see nothing change and go looking for the fault somewhere else") +} + +// --------------------------------------------------------------------------- +// (d) WHO GETS THE PAGE when the budget runs out +// --------------------------------------------------------------------------- + +// TestDivergenceExamplesAreDealtAcrossClassesRatherThanFirstComeFirstServed is +// the only test in the suite that seeds TWO bulk classes, and that is why it +// exists. +// +// Every other fixture here overflows ONE class, so first-come-first-served and +// round-robin produce identical reports and neither is pinned. The failure the +// staging exists to prevent needs a second class to be visible at all: with +// entries appended in comparison order, the class swept FIRST spends the entire +// budget and every later class arrives with a non-zero count and NOT ONE +// EXAMPLE. That is the exact thing DivergenceEntry.Subject exists to prevent — +// "a class and a count alone give an operator a number they cannot investigate" +// — happening to the classes that need investigating most, at the moment they +// need it, because a bulk class is precisely what an incident looks like. +// +// The fixture is built on the sweep's own order: persona vote events are the +// FIRST comparison and the unknown-delivery classes are the LAST, so the small +// class here is the one a first-come allocation would starve. +func TestDivergenceExamplesAreDealtAcrossClassesRatherThanFirstComeFirstServed(t *testing.T) { + h := newHarness(t) + world := newModerationWorld(t, h) + h.admin.SetDivergenceReconciler(newDivergenceReconciler(t, h)) + deliverEverythingQueued(t, h.db) + _ = world + + // The BULK class, swept first: more persona vote events than the whole page. + personaActorID := actorIDOf(t, h.db, mtAuthorDID) + overflowPersonaVotes(t, h.db, personaActorID, MaxDivergenceEntries+1) + + // The SMALL class, swept last: three deliveries nobody answered. Three is + // deliberately tiny — an operator can act on all of them, and a report that + // shows none of them is the regression. + late := []string{ + seedPoisonedComment(t, h.db, "dv-fair-1", "transport", 0), + seedPoisonedComment(t, h.db, "dv-fair-2", "transport", 0), + seedPoisonedComment(t, h.db, "dv-fair-3", "transport", 0), + } + + report := fetchDivergence(t, h) + + require.Len(t, report.Entries, MaxDivergenceEntries, + "precondition: the budget really is exhausted — with room to spare every class gets "+ + "everything and this test would pass without an allocation policy at all") + require.Equal(t, MaxDivergenceEntries+1, report.Counts[DivergencePersonaVoteEvent], + "precondition: the bulk class overflows the page") + require.Equal(t, len(late), report.Counts[DivergenceDeliveryUnknownUnanswered], + "precondition: the late class found exactly the rows this fixture seeded") + + assert.ElementsMatch(t, late, subjectsOfClass(report, DivergenceDeliveryUnknownUnanswered), + "the class swept LAST still gets its examples. Appending entries in comparison order "+ + "hands the whole budget to whoever read first, and every later class then arrives "+ + "with a count an operator cannot investigate — no subject to look up, no inbox to "+ + "ask, on the endpoint they opened BECAUSE something is wrong") + + // Every class the sweep counted must be able to show its work. + for class, count := range report.Counts { + if count == 0 { + continue + } + assert.NotEmpty(t, subjectsOfClass(report, class), + "class %q has a count of %d and no example: a number with nothing to look at is the "+ + "one thing this report promises never to serve", class, count) + } + + // ...AND THE BUDGET IS NOT SHARED OUT EQUALLY, which is the other way to get + // this wrong. An equal share computed up front would cap the bulk class at a + // fraction of the page and leave most of it empty, so a single-class incident + // would come back with a quarter of the examples it could have had. Dealing a + // card at a time gives every class its rows first and spends everything left + // on whoever still has some. + bulk := subjectsOfClass(report, DivergencePersonaVoteEvent) + assert.Equal(t, MaxDivergenceEntries-len(late), len(bulk), + "the bulk class keeps every slot the small classes did not need: round-robin fills the "+ + "page, an equal share wastes it, and the report's bound is only worth having if the "+ + "page is actually full") + + assert.True(t, report.Truncated, + "and the report still says it is smaller than the problem it describes") +} + +// --------------------------------------------------------------------------- +// (e) The sweep's own deadline, and what a failed pass leaves behind +// --------------------------------------------------------------------------- + +// deadlineDivergences records the context the sweep hands its reads. +type deadlineDivergences struct { + store.Divergences + mu sync.Mutex + deadline time.Time + hasDeadline bool +} + +func (d *deadlineDivergences) PersonaVoteEvents(ctx context.Context) ([]store.PersonaVoteEvent, error) { + deadline, ok := ctx.Deadline() + d.mu.Lock() + d.deadline, d.hasDeadline = deadline, ok + d.mu.Unlock() + return d.Divergences.PersonaVoteEvents(ctx) +} + +func (d *deadlineDivergences) Deadline() (time.Time, bool) { + d.mu.Lock() + defer d.mu.Unlock() + return d.deadline, d.hasDeadline +} + +// TestDivergenceSweepBoundsItsReadsWithADeadline pins divergenceSweepTimeout at +// the only place it is observable: the context the reads actually run under. +// +// The bound exists because of the POOL, not the clock. Every comparison is a +// multi-table join taken from the same 25 connections that serve inbound +// ingestion and the delivery worker, there is no statement_timeout anywhere in +// this codebase, and GET /admin/divergence makes a sweep reachable from outside +// on an unrated GET — so an operator refreshing the report during an incident +// can deepen the incident. A sweep with no deadline does not merely run late; it +// holds pool connections for as long as the database will let it, and it is +// slowest exactly when federation most needs them. +// +// THE CALLER HERE HAS NO DEADLINE OF ITS OWN, which is the production case: a +// curl carries none, and the background Run loop's context carries none either. +// If the sweep did not impose one, these reads would run unbounded. +func TestDivergenceSweepBoundsItsReadsWithADeadline(t *testing.T) { + h := newHarness(t) + + watched := &deadlineDivergences{Divergences: store.NewDivergences(h.db)} + reconciler, err := NewDivergenceReconciler(DivergenceOptions{DB: h.db, Divergences: watched}) + require.NoError(t, err) + + // context.Background(): no deadline, exactly like a curl and like Run. + _, err = reconciler.Sweep(context.Background()) + require.NoError(t, err) + + deadline, ok := watched.Deadline() + require.True(t, ok, + "the reads must run under a DEADLINE the sweep imposed. The caller supplied none — that "+ + "is what a curl and the background loop both look like — so without one here a "+ + "pathological pass pins pool connections indefinitely, on the endpoint an operator "+ + "reaches for because the database is already struggling") + remaining := time.Until(deadline) + assert.LessOrEqual(t, remaining, divergenceSweepTimeout, + "and the budget is divergenceSweepTimeout (%s), never longer: it has to be well inside "+ + "the sweep cadence or a slow pass queues behind its own predecessor", divergenceSweepTimeout) + assert.Greater(t, remaining, divergenceSweepTimeout-time.Minute, + "and generously long: a deadline that fired on a merely slow pass would report failures "+ + "that are not divergences, which is the one thing a report like this cannot afford") +} + +// blockingDivergences is the pathological read: it stops on the context and +// never returns until that context is done. +// +// It is how a sweep is made to abort in a test without waiting out the real +// two-minute budget. The mechanism under test is the same one either way — the +// reads run under a context derived from the caller's, and they end when it ends +// — and driving it by cancellation exercises the derivation, which the deadline +// alone cannot. +type blockingDivergences struct { + store.Divergences + entered chan struct{} +} + +func (b *blockingDivergences) PersonaVoteEvents(ctx context.Context) ([]store.PersonaVoteEvent, error) { + select { + case b.entered <- struct{}{}: + default: + } + <-ctx.Done() + return nil, ctx.Err() +} + +// TestDivergenceSweepAbortsWhenItsContextEndsAndLeavesTheGaugesStanding drives a +// sweep that would never finish on its own. +// +// TWO CLAIMS, and the second is the one an operator lives with. The pass ABORTS +// rather than blocking forever — the reads run under the sweep's own context, +// derived from the caller's, so cancelling the caller reaches a read that is +// already inside the database. And the gauges KEEP THEIR PREVIOUS VALUES: a +// sweep that could not read must never write a zero, because a zero written by a +// pass that measured nothing is a claim of health made at the moment nobody +// could check, and it overwrites the last real number the operator was watching. +// +// THE FAILURE COUNTER DELIBERATELY DOES NOT MOVE HERE. This sweep was cut short +// by its CALLER, which is what shutdown looks like, and counting it would make +// the number mean "restarts plus failures" — a number nobody can act on. The +// counter's positive case is the failing-read test below. +func TestDivergenceSweepAbortsWhenItsContextEndsAndLeavesTheGaugesStanding(t *testing.T) { + h := newHarness(t) + world := newModerationWorld(t, h) + deliverEverythingQueued(t, h.db) + _ = world + + // A real divergence first, so the gauge carries a NON-ZERO value an aborted + // sweep could overwrite. A fixture whose healthy value is 0 cannot tell "left + // standing" from "written as zero". + insertVoteEvent(t, h.db, "https://lemmy.world/activities/like/dv-abort", + actorIDOf(t, h.db, mtAuthorDID), "up") + h.admin.SetDivergenceReconciler(newDivergenceReconciler(t, h)) + require.Equal(t, http.StatusOK, h.adminRequest(http.MethodGet, "/admin/divergence", nil).Code) + require.Equal(t, 1, gaugeValue(t, h, MetricDivergencePersonaVoteEvents), + "precondition: a completed sweep published a real number") + failuresBefore := gaugeValue(t, h, MetricDivergenceSweepFailures) + + blocking := &blockingDivergences{ + Divergences: store.NewDivergences(h.db), + entered: make(chan struct{}, 1), + } + stuck, err := NewDivergenceReconciler(DivergenceOptions{DB: h.db, Divergences: blocking}) + require.NoError(t, err) + + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + swept := make(chan error, 1) + go func() { _, sweepErr := stuck.Sweep(ctx); swept <- sweepErr }() + + select { + case <-blocking.entered: + case <-time.After(5 * time.Second): + t.Fatal("precondition: the sweep must reach the read this test blocks in") + } + cancel() + + select { + case sweepErr := <-swept: + require.Error(t, sweepErr, + "a sweep whose reads were cut short is a FAILURE, not an empty report: returning the "+ + "classes that happened to finish publishes a claim about state nobody read") + assert.ErrorIs(t, sweepErr, context.Canceled) + case <-time.After(5 * time.Second): + t.Fatal("the sweep must END when its context does. The reads run under a context DERIVED " + + "from the caller's precisely so that a shutdown — or the two-minute budget — reaches " + + "a query already inside the database; a sweep built on a fresh background context " + + "holds its pool connections for as long as the database will let it, which is the " + + "failure divergenceSweepTimeout exists to bound") + } + + assert.Equal(t, 1, gaugeValue(t, h, MetricDivergencePersonaVoteEvents), + "THE PREVIOUS VALUE STANDS. An aborted pass measured nothing, and a zero written by a "+ + "pass that measured nothing reads as health at exactly the moment nobody could check") + assert.Equal(t, failuresBefore, gaugeValue(t, h, MetricDivergenceSweepFailures), + "and the failure counter does NOT move for a sweep its own caller cancelled: that is "+ + "shutdown, and a counter inflated on every restart means 'restarts plus failures', "+ + "which is a number nobody can alert on") +} + +// TestDivergenceSweepFailureAndAgeAreVisibleToAnOperator reads the two freshness +// metrics through the surface an operator reads them through. +// +// THEY ARE READ FROM /admin/metrics, NOT FROM expvar. scopedMetrics serves only +// keys carrying the tidepool prefix, so a metric named without it is published, +// behaves perfectly, and appears nowhere — indistinguishable from a check that +// never runs. A test that pokes the Var it just watched proves the sweep can +// call Set; it proves nothing about whether anyone can see the result. +// +// WHY THESE TWO EXIST AT ALL: the seven class gauges are all SET by a successful +// sweep, so a sweep that has stopped running leaves seven low, stable, entirely +// healthy-looking numbers behind it. The age says when those numbers were +// established and the failure count says whether the sweep has been trying and +// losing. Read together they are the difference between "nothing is diverging" +// and "nothing is measuring"; either one alone can be read as health. +func TestDivergenceSweepFailureAndAgeAreVisibleToAnOperator(t *testing.T) { + h := newHarness(t) + world := newModerationWorld(t, h) + deliverEverythingQueued(t, h.db) + _ = world + + insertVoteEvent(t, h.db, "https://lemmy.world/activities/like/dv-age", + actorIDOf(t, h.db, mtAuthorDID), "up") + h.admin.SetDivergenceReconciler(newDivergenceReconciler(t, h)) + require.Equal(t, http.StatusOK, h.adminRequest(http.MethodGet, "/admin/divergence", nil).Code) + require.Equal(t, 1, gaugeValue(t, h, MetricDivergencePersonaVoteEvents), + "precondition: a completed sweep published a real number") + + fresh := floatGaugeValue(t, h, MetricDivergenceSweepAgeSeconds) + assert.GreaterOrEqual(t, fresh, 0.0, + "a sweep has published, so the age is a real measurement rather than the unswept "+ + "sentinel: %s reads negative only before the first successful publication, because a "+ + "0 there would say 'swept just now' about a process that has never swept at all", + MetricDivergenceSweepAgeSeconds) + assert.Less(t, fresh, time.Minute.Seconds(), + "and it is SECONDS old, because the sweep it dates finished during this test") + + // PULLED, NOT PUSHED, and this is the whole point of the age gauge. A value + // written at publication time would be the one number that stops updating at + // exactly the moment it starts to matter — a sweep that has stopped running + // would report the age it had when it last ran, forever. + climbing := floatGaugeValue(t, h, MetricDivergenceSweepAgeSeconds) + assert.Greater(t, climbing, fresh, + "%s must CLIMB between two reads of the same standing numbers. Pushed at publication it "+ + "would freeze with the class gauges it is supposed to date, and seven frozen gauges "+ + "beside a frozen age is exactly what a healthy bridge looks like", + MetricDivergenceSweepAgeSeconds) + + // Now a sweep that fails on a read, with a live caller: this is a FAILING + // pass, not a shutdown. + failuresBefore := gaugeValue(t, h, MetricDivergenceSweepFailures) + faulty := &failingDivergences{ + Divergences: store.NewDivergences(h.db), + err: stderrors.New("connection reset by peer"), + } + failing, err := NewDivergenceReconciler(DivergenceOptions{DB: h.db, Divergences: faulty}) + require.NoError(t, err) + h.admin.SetDivergenceReconciler(failing) + require.Equal(t, http.StatusInternalServerError, + h.adminRequest(http.MethodGet, "/admin/divergence", nil).Code, + "precondition: the sweep really did fail") + require.NotZero(t, faulty.Calls(), "precondition: the sweep really did try to read") + + assert.Equal(t, failuresBefore+1, gaugeValue(t, h, MetricDivergenceSweepFailures), + "%s must move on a failed pass. Leaving the previous gauges standing is right, and on "+ + "its own it is SILENT: a sweep that fails forever leaves seven plausible numbers "+ + "frozen and a log line nobody is watching, and this counter plus the climbing age is "+ + "the only thing that tells those apart from a healthy bridge", + MetricDivergenceSweepFailures) + + afterFailure := floatGaugeValue(t, h, MetricDivergenceSweepAgeSeconds) + assert.Greater(t, afterFailure, climbing, + "and the age KEEPS CLIMBING through the failure: it dates the standing numbers, which a "+ + "failed sweep did not refresh. An age reset by a pass that published nothing would "+ + "report the stale gauges as fresh, which is the single reading that would hide a "+ + "sweep that has stopped working") + assert.Equal(t, 1, gaugeValue(t, h, MetricDivergencePersonaVoteEvents), + "while the class gauge keeps its last real value: the counter and the age exist so that "+ + "standing numbers can be told from measured ones, not so they can be overwritten") +} diff --git a/internal/ingest/divergence_test.go b/internal/ingest/divergence_test.go index 93a27f8..45d8bc6 100644 --- a/internal/ingest/divergence_test.go +++ b/internal/ingest/divergence_test.go @@ -64,23 +64,8 @@ const ( dvGauge = "tidepool_divergence_acceptance_undelivered_cancelled" ) -// dvTables is every table the sweep could conceivably touch, plus the repo -// tables underneath them. It is deliberately wider than the reconciler's -// inputs: the assertion is not "it did not write the rows I expected it to -// leave alone", it is "it did not write". -var dvTables = []string{ - // The outbound queue and ledger — the side an over-eager "fix" would - // re-drive, cancel or re-flip. - "outbound_deliveries", "outbound_activities", "outbound_votes", "outbound_objects", - // The atproto side. A reconciler that "restored" an acceptance would commit - // to a community repo, which lands here as new blocks, a moved head and a - // firehose frame — the loudest possible write, and the easiest to make by - // accident when the report has already computed exactly what is missing. - "blocks", "repo_state", "firehose_events", - // The bridge's own records of what it decided. - "admissions", "ap_objects", "ap_actors", "federation_prefs", - "community_bans", "object_moderation", "vote_events", "vote_aggregates", -} +// snapshotTables enumerates the tables from the database itself. See +// allBaseTables for why the list is not written down here. // TestDivergenceReportNamesAnUndeliveredAcceptanceAndWritesNothing is the OUTER // acceptance test for 17e. @@ -115,7 +100,7 @@ func TestDivergenceReportNamesAnUndeliveredAcceptanceAndWritesNothing(t *testing acceptanceStands(t, h, world.communityADID, postURI) // --- The state of the world, in full, immediately before the sweep. - before := snapshotTables(t, h.db, dvTables...) + before := snapshotTables(t, h.db) // --- WHEN: the operator asks what diverges. rec := h.adminRequest(http.MethodGet, "/admin/divergence", nil) @@ -167,8 +152,11 @@ func TestDivergenceReportNamesAnUndeliveredAcceptanceAndWritesNothing(t *testing "unchanging problem, and would never return to zero when it is resolved", dvGauge) // --- AND, THE POINT: the sweep changed NOTHING. - after := snapshotTables(t, h.db, dvTables...) - for _, table := range dvTables { + after := snapshotTables(t, h.db) + require.Equal(t, len(before), len(after), + "the sweep must not CREATE or DROP a table either — the comparison is over whatever the "+ + "schema holds, so a new one appearing is itself a write") + for table := range before { assert.Equal(t, before[table], after[table], "REPORTING IS THE WHOLE JOB: %s must be byte-identical across the sweep. A "+ "reconciler that writes has no witness — it is the only thing that reads both "+ @@ -392,6 +380,30 @@ func gaugeValue(t *testing.T, h *harness, name string) int { return value } +// floatGaugeValue is gaugeValue for a metric that is not a whole number. +// +// The age gauge is seconds as a float — an expvar.Func over a stored timestamp — +// so reading it through gaugeValue would fail on the decimal rather than on +// anything true about the sweep. It goes through the ENDPOINT for the same +// reason gaugeValue does: the tidepool prefix filter in scopedMetrics is exactly +// the kind of silent drop a direct expvar read cannot catch. +func floatGaugeValue(t *testing.T, h *harness, name string) float64 { + t.Helper() + rec := h.adminRequest(http.MethodGet, "/admin/metrics", nil) + require.Equal(t, http.StatusOK, rec.Code) + var metrics map[string]json.RawMessage + require.NoError(t, json.Unmarshal(rec.Body.Bytes(), &metrics), + "/admin/metrics must stay parseable JSON (body: %s)", rec.Body.String()) + raw, ok := metrics[name] + require.True(t, ok, + "%s must be published on /admin/metrics. A metric whose name does not start with "+ + "'tidepool' is filtered out by scopedMetrics and is invisible in exactly the way a "+ + "check that never runs is invisible. Published keys: %v", name, keysOf(metrics)) + var value float64 + require.NoError(t, json.Unmarshal(raw, &value), "%s must be a number, got %s", name, raw) + return value +} + func keysOf(m map[string]json.RawMessage) []string { keys := make([]string, 0, len(m)) for key := range m { @@ -409,8 +421,13 @@ func keysOf(m map[string]json.RawMessage) []string { // comparison that carries every column can say so. row_to_json also survives a // migration adding a column, which a hand-listed projection would silently stop // covering on the day the schema grows. -func snapshotTables(t *testing.T, db *sql.DB, tables ...string) map[string][]string { +func snapshotTables(t *testing.T, db *sql.DB) map[string][]string { t.Helper() + tables := allBaseTables(t, db) + require.Greater(t, len(tables), 15, + "the enumeration must actually find the schema: a query that returned a handful of "+ + "tables would make this whole assertion a spot-check wearing the clothes of an "+ + "exhaustive one") snapshot := make(map[string][]string, len(tables)) for _, table := range tables { rows, err := db.QueryContext(context.Background(), @@ -674,3 +691,40 @@ func classesIn(report divergenceReport) []string { } return classes } + +// allBaseTables lists every base table in the public schema, minus goose's +// migration bookkeeping. +// +// ENUMERATED, NEVER LISTED. A hand-written list is a claim about the schema +// made at the moment it was typed, and it stops being true the next time +// anything is added — silently, in the direction that weakens the assertion. +// The previous version of this file watched 15 of 24 tables, and the gaps were +// exactly the ones that matter: ap_tombstones, which the admin endpoint +// registered four lines above /divergence writes, and communities, which the +// SIBLING reconciler converges by writing. "While we are here, tombstone what +// the peer never got" is the most plausible accidental repair anyone would add +// to this sweep, and the snapshot would not have seen it. +// +// goose_db_version is excluded because migrations are not the sweep's writes and +// a test-run migration would make every comparison here fail for the wrong +// reason. +func allBaseTables(t *testing.T, db *sql.DB) []string { + t.Helper() + rows, err := db.QueryContext(context.Background(), ` + SELECT table_name + FROM information_schema.tables + WHERE table_schema = 'public' + AND table_type = 'BASE TABLE' + AND table_name <> 'goose_db_version' + ORDER BY table_name`) + require.NoError(t, err) + defer func() { require.NoError(t, rows.Close()) }() + var tables []string + for rows.Next() { + var name string + require.NoError(t, rows.Scan(&name)) + tables = append(tables, name) + } + require.NoError(t, rows.Err()) + return tables +} diff --git a/internal/ingest/follow.go b/internal/ingest/follow.go index 652d5ea..68f9db8 100644 --- a/internal/ingest/follow.go +++ b/internal/ingest/follow.go @@ -535,10 +535,6 @@ func (a *Admin) handleBackfill(w http.ResponseWriter, r *http.Request) { writeJSON(w, http.StatusAccepted, communityJSON(community)) } -// handleReconcile runs one synchronous follow-list sweep on demand, so an -// operator can converge right after editing the file instead of waiting for -// the next tick. 501 when no follow list is configured (the handleBackfill -// nil-dependency pattern). // handleDivergence serves the reconciliation report (task 17e, decision 19). // // GET, not POST, and that is a statement rather than a convention: this sweep @@ -547,8 +543,16 @@ func (a *Admin) handleBackfill(w http.ResponseWriter, r *http.Request) { // still working out what is wrong. The moment it needed POST it would have // stopped being a report. // -// A sweep failure is a 500 carrying the reason: a partial report would be a lie -// about the classes it never reached, and this endpoint is operator-facing. +// A sweep failure is a 500 and NO PARTIAL REPORT: a report missing one class +// reads exactly like a class that found nothing, so the classes that happened to +// succeed must not ride out with the error. +// +// The reason goes to the LOG, not to the body, like every other handler here. +// This is admin-authenticated, so the exposure is small, but every read in the +// sweep is raw SQL and a pq error carries table, column and constraint names +// straight out of the schema — detail an operator can read in the log line one +// scroll away, and the only party the body could ever tell is someone who should +// not be reading it. func (a *Admin) handleDivergence(w http.ResponseWriter, r *http.Request) { if a.divergence == nil { http.Error(w, "divergence reconciliation is not configured", http.StatusNotImplemented) @@ -557,13 +561,17 @@ func (a *Admin) handleDivergence(w http.ResponseWriter, r *http.Request) { report, err := a.divergence.Sweep(r.Context()) if err != nil { a.logger.Error("divergence sweep failed", "error", err) - http.Error(w, "divergence sweep failed: "+err.Error(), http.StatusInternalServerError) + http.Error(w, "divergence sweep failed", http.StatusInternalServerError) return } w.Header().Set("Content-Type", "application/json; charset=utf-8") _ = json.NewEncoder(w).Encode(report) } +// handleReconcile runs one synchronous follow-list sweep on demand, so an +// operator can converge right after editing the file instead of waiting for +// the next tick. 501 when no follow list is configured (the handleBackfill +// nil-dependency pattern). func (a *Admin) handleReconcile(w http.ResponseWriter, r *http.Request) { if a.reconciler == nil { http.Error(w, "follow list reconciliation is not configured", http.StatusNotImplemented) diff --git a/internal/outbound/worker.go b/internal/outbound/worker.go index 8c499fa..1f753f5 100644 --- a/internal/outbound/worker.go +++ b/internal/outbound/worker.go @@ -259,10 +259,19 @@ func (w *Worker) handle(ctx context.Context, delivery *store.OutboundDelivery) e // instant its parent is accepted, and do not advance the poison budget // (the causal wait is wall-clock-bounded in causalStatus). return w.parkCausal(ctx, delivery, "parent_pending", "waiting for bridge-origin parent to be accepted") + // THE CLASS NAMES COME FROM store, and so do cross_authority and signer + // below. The reconciliation sweep excludes exactly these four from its + // unknown-outcome report (store.neverReachedTheWireClasses): they carry + // last_status_code 0 like a dial timeout does and mean the opposite — nothing + // was sent, so the peer's state is not unknown, they simply do not have it. + // Spelling them here as literals made that correspondence a comment; naming + // the constants makes it the compiler's problem. case causalPoisonUnaccepted: - return w.poison(ctx, delivery, "parent_unaccepted", "causal wait budget exhausted; parent never accepted", 0) + return w.poison(ctx, delivery, store.PoisonClassParentUnaccepted, + "causal wait budget exhausted; parent never accepted", 0) case causalPoisonParent: - return w.poison(ctx, delivery, "parent_poisoned", "parent delivery poisoned; descendant cannot land", 0) + return w.poison(ctx, delivery, store.PoisonClassParentPoisoned, + "parent delivery poisoned; descendant cannot land", 0) } // Consent recheck (retraction asymmetry): a Delete/Undo always goes out — @@ -290,7 +299,7 @@ func (w *Worker) handle(ctx context.Context, delivery *store.OutboundDelivery) e // with the target community. The resolver already refuses a cross-authority // inbox at enqueue time; this catches a tampered or legacy stored target. if !ap.SameAuthority(delivery.OrderingKey, delivery.TargetInbox) { - return w.poison(ctx, delivery, "cross_authority", + return w.poison(ctx, delivery, store.PoisonClassCrossAuthority, "target inbox is not same-authority with the community; refusing to deliver", 0) } @@ -304,7 +313,7 @@ func (w *Worker) deliver(ctx context.Context, delivery *store.OutboundDelivery, if err != nil { // A signer that cannot be resolved right now is transient (a KEK blip, // a not-yet-replicated actor): retry rather than poison. - return w.releaseOrPoison(ctx, delivery, "signer", err.Error(), 0) + return w.releaseOrPoison(ctx, delivery, store.PoisonClassSigner, err.Error(), 0) } err = w.sender.SendActivityAs(ctx, signer, delivery.TargetInbox, json.RawMessage(activity.Payload)) diff --git a/internal/store/divergence.go b/internal/store/divergence.go index 7bdb2db..404c895 100644 --- a/internal/store/divergence.go +++ b/internal/store/divergence.go @@ -5,6 +5,8 @@ import ( "database/sql" "fmt" "time" + + "github.com/lib/pq" ) // The reconciliation sweep's reads (task 17e, decision 19). @@ -16,6 +18,10 @@ import ( // an Exec — there is none below, and adding one would make this the only // component that both reads a divergence and can act on it, which is exactly // the shape decision 19 forbids. +// +// AND THAT IS NOT LEFT TO THE COMMENT. Every read below runs inside BEGIN … +// READ ONLY (see readOnlyTx), so Postgres itself refuses a write here — the one +// form of enforcement that outlives everyone who has read this paragraph. // PersonaVoteEvent is one INBOUND vote_events row whose voter resolves to one // of our own native personas — a vote the bridge cast on a user's behalf, @@ -64,26 +70,97 @@ type UndeliveredAcceptance struct { LastErrorClass string } +// MaxDivergenceExamples is the EXAMPLE BUDGET, and it is applied as a SQL +// LIMIT rather than by cutting a finished slice down. +// +// The sweep's report caps the examples it carries (ingest.MaxDivergenceEntries, +// which is this constant), and a cap applied in Go after the rows arrive bounds +// nothing that matters: the sweep that reaches the cap is the sweep running +// against a stopped queue, and every one of those rows would already have been +// scanned into a slice, in the process an operator is trying to keep alive, on +// an endpoint they reached for BECAUSE it is struggling. So the bound goes to +// the database — every list below returns at most this many rows. +// +// THE COUNTS ARE NOT BOUNDED WITH THEM. Each list has a matching exact +// COUNT(*), because the count is what an operator SIZES the incident from: a +// number capped at the length of a page would make "500" mean both "500" and +// "a catastrophe". That is the whole reason these are two statements rather +// than one truncated read. +const MaxDivergenceExamples = 500 + // Divergences reads the comparisons the reconciliation sweep reports on. Every // method is a READ; both sides of every comparison are local. +// +// EACH COMPARISON IS TWO METHODS — bounded examples, exact count — and the two +// are deliberately not folded into one call that returns both. The list +// signatures are what the reconciler's failure-injection doubles wrap, and more +// importantly the pair states the contract in the type: a caller that wants a +// number cannot get a truncated one by accident, and a caller that wants +// examples cannot mistake how many there were. +// +// The count is read on EVERY sweep rather than only when a list comes back +// full. It could be skipped in the common case — a list under the limit is its +// own exact count — but that would leave the counting queries dark until the +// first incident, which is the one moment nobody wants to discover them for the +// first time. Running them always makes every existing test that asserts a +// count a test of the counting query too. The cost is one extra aggregate pass +// per comparison per sweep; see FOLLOWUPS.md for what that pass costs on the +// acceptance leg specifically. type Divergences interface { // PersonaVoteEvents lists inbound vote events attributed to our own - // personas, joining vote_events.voter_ap_id to ap_actors.actor_id. + // personas, joining vote_events.voter_ap_id to ap_actors.actor_id. At most + // MaxDivergenceExamples rows. PersonaVoteEvents(ctx context.Context) ([]PersonaVoteEvent, error) + // PersonaVoteEventCount is how many such rows exist, unbounded by the + // example budget. + PersonaVoteEventCount(ctx context.Context) (int, error) + // UndeliveredAcceptances lists accepted posts whose federation never // reached the peer. A pending delivery counts only once it is older than // staleAfter: until then it is in flight, which is the ordinary state of - // every post between acceptance and delivery. + // every post between acceptance and delivery. At most + // MaxDivergenceExamples rows. UndeliveredAcceptances(ctx context.Context, staleAfter time.Duration) ([]UndeliveredAcceptance, error) + // UndeliveredAcceptanceCounts is how many such posts exist, KEYED BY + // DELIVERY STATE — the same split the sweep's three acceptance classes are + // derived from, so the caller maps states to classes in exactly one place + // and a state this sweep has no class for arrives as itself rather than + // folded into a total. + UndeliveredAcceptanceCounts(ctx context.Context, staleAfter time.Duration) (map[DeliveryState]int, error) + // RecastDivergences lists (actor, subject) pairs where a peer is holding a - // vote this bridge no longer claims — see RecastDivergence. + // vote this bridge no longer claims — see RecastDivergence. At most + // MaxDivergenceExamples rows. RecastDivergences(ctx context.Context) ([]RecastDivergence, error) + // RecastDivergenceCount is how many such pairs exist. + RecastDivergenceCount(ctx context.Context) (int, error) + // UnknownDeliveryOutcomes lists poisoned deliveries: activities we SENT and - // never got confirmation for — see UnknownDeliveryOutcome. + // never got confirmation for — see UnknownDeliveryOutcome. At most + // MaxDivergenceExamples rows. UnknownDeliveryOutcomes(ctx context.Context) ([]UnknownDeliveryOutcome, error) + + // UnknownDeliveryOutcomeCounts is how many such deliveries exist, split the + // same way the entries are — see UnknownDeliveryCounts. + UnknownDeliveryOutcomeCounts(ctx context.Context) (UnknownDeliveryCounts, error) +} + +// UnknownDeliveryCounts is the refused/unanswered split, counted rather than +// listed. +// +// It is a struct rather than two ints returned side by side because the split +// is the only information these rows carry, and the two numbers must never be +// added: a peer that ANSWERED told us something a silent one did not. A single +// value returned here would be a total of things we do not know, which is a +// number nobody can act on. +type UnknownDeliveryCounts struct { + // Refused is deliveries where a status code came back. + Refused int + // Unanswered is deliveries where nothing came back at all. + Unanswered int } // UnknownDeliveryOutcome is a delivery whose result this bridge DOES NOT KNOW. @@ -155,6 +232,44 @@ type postgresDivergences struct{ db *sql.DB } // NewDivergences creates the postgres-backed read-only divergence reader. func NewDivergences(db *sql.DB) Divergences { return &postgresDivergences{db: db} } +// readOnlyTx opens the transaction every read in this file runs in, and it is +// the ONLY enforcement of this file's one rule that survives an edit. +// +// "This reconciler cannot write" was a comment at the top of the file plus a +// test that snapshots every table before and after a sweep. Both are worth +// having and neither stops the write: the comment is advice, and the snapshot +// test only fails AFTER someone has already added an Exec and run it. BEGIN … +// READ ONLY moves the rule into the database, which refuses the statement +// outright (25006, read_only_sql_transaction) — so a self-healing write bolted +// on here fails on the day it is written, in the author's own test run, instead +// of the day it corrupts the state the sweep exists to measure. That is decision +// 19 enforced by something other than good intentions. +// +// ONE TRANSACTION PER READ, not one per sweep. The sweep's eight reads are eight +// independent comparisons that already tolerate each other moving — nothing here +// is cross-checked between two of them — and a pass-wide snapshot would hold a +// single pool connection for the whole two-minute budget, which is the opposite +// of what the sweep timeout was added for. The cost is a BEGIN and a ROLLBACK +// per read, on a job that runs every fifteen minutes. +// +// It always ROLLBACKs, never COMMITs: a read-only transaction has nothing to +// make durable, and an ending that cannot possibly persist anything is the one +// that matches what this file claims about itself. +func (r *postgresDivergences) readOnlyTx(ctx context.Context) (*sql.Tx, func(), error) { + tx, err := r.db.BeginTx(ctx, &sql.TxOptions{ReadOnly: true}) + if err != nil { + return nil, nil, fmt.Errorf("begin read-only transaction: %w", err) + } + return tx, func() { _ = tx.Rollback() }, nil +} + +// personaVoteEventsFrom is the POPULATION, written once and shared by the +// example list and the count. Two spellings of one comparison drift, and the +// drift is silent in exactly the direction that matters: a count taken over a +// slightly different join reports a size for a population nobody listed. +const personaVoteEventsFrom = `FROM vote_events v + JOIN ap_actors a ON a.actor_id = v.voter_ap_id` + func (r *postgresDivergences) PersonaVoteEvents(ctx context.Context) ([]PersonaVoteEvent, error) { // THE JOIN IS ON IDENTITY, NEVER ON ACTIVITY IDS. Decision 16 originally // proposed matching inbound votes against outbound_activities by activity @@ -178,13 +293,22 @@ func (r *postgresDivergences) PersonaVoteEvents(ctx context.Context) ([]PersonaV // that cleanup would have removed. A zero therefore means "no // exactly-spelled persona vote", not "no persona vote" — see the KNOWNNARROW // case in divergence_test.go, which pins that limit rather than papering it. + // + // The LIMIT is the example budget reaching the database; the count below + // reads the same population with the same FROM clause and no bound. query := ` SELECT v.activity_id, v.voter_ap_id, v.subject_ap_id, a.did - FROM vote_events v - JOIN ap_actors a ON a.actor_id = v.voter_ap_id - ORDER BY v.activity_id` + ` + personaVoteEventsFrom + ` + ORDER BY v.activity_id + LIMIT $1` + + tx, done, err := r.readOnlyTx(ctx) + if err != nil { + return nil, fmt.Errorf("list persona vote events: %w", err) + } + defer done() - rows, err := r.db.QueryContext(ctx, query) + rows, err := tx.QueryContext(ctx, query, MaxDivergenceExamples) if err != nil { return nil, fmt.Errorf("list persona vote events: %w", err) } @@ -207,33 +331,68 @@ func (r *postgresDivergences) PersonaVoteEvents(ctx context.Context) ([]PersonaV return found, nil } -func (r *postgresDivergences) UndeliveredAcceptances(ctx context.Context, staleAfter time.Duration) ([]UndeliveredAcceptance, error) { - // BOTH SIDES ARE LOCAL, and the second one is the whole comparison. - // - // admissions.status = 'accepted' — the community's repo carries an - // acceptance record, so Coves renders the post in that community. - // outbound_objects.accepted_at IS NULL — the peer never confirmed it. That - // column is stamped ONLY by delivery success (worker.stampAccepted), so - // it is the one signal in the schema that means "it landed" rather than - // "we tried". - // - // The activity join is the SAME correspondence the worker uses to stamp - // that column: a Create/Update payload names the object it federates in - // object.id, and that id is outbound_objects.ap_object_id. Read as a JSONB - // path rather than a substring search, so an id that merely appears - // somewhere in another activity's payload cannot masquerade as this one's. - // - // LATEST DELIVERY FIRST, then classify. Picking the newest attempt (by seq) - // and asking what state IT is in describes the situation now; filtering - // first and taking the newest survivor would report a post whose older - // attempt was cancelled while a fresh one is still in flight. - // - // KNOWN GAP, stated rather than silently swallowed: an acceptance with NO - // delivery row at all (the LEFT JOINs yield NULL) is not reported, because - // the three classes here are all delivery states and it has none. It is a - // real divergence — nothing will ever carry that post — and it needs its own - // class rather than being folded into one that would misdescribe it. - query := ` +// PersonaVoteEventCount counts the same population PersonaVoteEvents lists, +// with no LIMIT. No ORDER BY either — ordering an aggregate is work whose only +// product is a sort the count throws away. +func (r *postgresDivergences) PersonaVoteEventCount(ctx context.Context) (int, error) { + tx, done, err := r.readOnlyTx(ctx) + if err != nil { + return 0, fmt.Errorf("count persona vote events: %w", err) + } + defer done() + + var total int + if err := tx.QueryRowContext(ctx, `SELECT count(*) `+personaVoteEventsFrom).Scan(&total); err != nil { + return 0, fmt.Errorf("count persona vote events: %w", err) + } + return total, nil +} + +// undeliveredAcceptanceLatest and undeliveredAcceptanceFilter are the +// comparison, written ONCE and shared by the example list and the per-state +// counts. Restating either in a second query is how a count comes to describe a +// population the list does not: the CTE decides which delivery is the current +// one, and the filter decides which of those is a divergence. +// +// BOTH SIDES ARE LOCAL, and the second one is the whole comparison. +// +// admissions.status = 'accepted' — the community's repo carries an +// acceptance record, so Coves renders the post in that community. +// outbound_objects.accepted_at IS NULL — the peer never confirmed it. That +// column is stamped ONLY by delivery success (worker.stampAccepted), so +// it is the one signal in the schema that means "it landed" rather than +// "we tried". +// +// THE ACCEPTANCE JOIN IS ON THE PAIR, not on the post alone. admissions is keyed +// (community_did, post_uri) — one post may be admitted by several communities — +// while outbound_objects is keyed by at_uri and carries its OWN community_did. +// Joining on post_uri alone would let the accepted status come from one +// community's admission and the reported CommunityDID from another community's +// object row, so the report would name a community that did not accept it. No +// writer produces that shape today (each post federates to the community that +// admitted it), which is exactly why the predicate is worth spelling: it costs +// nothing and it states the invariant instead of relying on it. +// +// The activity join is the SAME correspondence the worker uses to stamp that +// column: a Create/Update payload names the object it federates in object.id, +// and that id is outbound_objects.ap_object_id. Read as a JSONB path rather +// than a substring search, so an id that merely appears somewhere in another +// activity's payload cannot masquerade as this one's. NOTHING INDEXES THAT +// EXPRESSION — it is a full pass over outbound_activities per sweep, twice now +// that the count is its own statement; see FOLLOWUPS.md, which carries the +// EXPLAIN and the candidate expression index. +// +// LATEST DELIVERY FIRST, then classify. Picking the newest attempt (by seq) +// and asking what state IT is in describes the situation now; filtering +// first and taking the newest survivor would report a post whose older +// attempt was cancelled while a fresh one is still in flight. +// +// KNOWN GAP, stated rather than silently swallowed: an acceptance with NO +// delivery row at all (the LEFT JOINs yield NULL) is not reported, because +// the three classes here are all delivery states and it has none. It is a +// real divergence — nothing will ever carry that post — and it needs its own +// class rather than being folded into one that would misdescribe it. +const undeliveredAcceptanceLatest = ` WITH latest AS ( SELECT DISTINCT ON (o.at_uri) o.at_uri AS post_uri, @@ -243,39 +402,67 @@ func (r *postgresDivergences) UndeliveredAcceptances(ctx context.Context, staleA d.last_error_class AS last_error_class, d.created_at AS delivery_created_at FROM admissions adm - JOIN outbound_objects o ON o.at_uri = adm.post_uri + JOIN outbound_objects o + ON o.at_uri = adm.post_uri + AND o.community_did = adm.community_did LEFT JOIN outbound_activities a ON a.payload -> 'object' ->> 'id' = o.ap_object_id LEFT JOIN outbound_deliveries d ON d.activity_id = a.activity_id WHERE adm.status = 'accepted' AND o.accepted_at IS NULL ORDER BY o.at_uri, d.seq DESC NULLS LAST - ) - SELECT post_uri, - community_did, - COALESCE(activity_id, ''), - COALESCE(delivery_state, ''), - COALESCE(last_error_class, '') - FROM latest + )` + +// undeliveredAcceptanceFilter selects the divergent rows out of that CTE. +// +// A HELD SETTLEMENT IS NOT A DIVERGENCE, and it is the one case this report +// would otherwise cry wolf on. It is pending, old, and unstamped — identical +// to a stuck delivery on every column above except this one — and it means +// the opposite: the peer ALREADY ACCEPTED the activity, and accepted_at is +// missing precisely because writing it is the step that failed. The worker +// resumes it on its next claim. Reporting it fires the sweep on every +// settlement retry, and an operator who learns to ignore this report also +// ignores the cancelled acceptance beside it, which is the finding that +// never heals itself. +const undeliveredAcceptanceFilter = ` WHERE delivery_state IN ($1, $2) OR (delivery_state = $3 AND last_error_class <> $4 - AND delivery_created_at < now() - $5::interval) - ORDER BY post_uri` - - // A HELD SETTLEMENT IS NOT A DIVERGENCE, and it is the one case this report - // would otherwise cry wolf on. It is pending, old, and unstamped — identical - // to a stuck delivery on every column above except this one — and it means - // the opposite: the peer ALREADY ACCEPTED the activity, and accepted_at is - // missing precisely because writing it is the step that failed. The worker - // resumes it on its next claim. Reporting it fires the sweep on every - // settlement retry, and an operator who learns to ignore this report also - // ignores the cancelled acceptance beside it, which is the finding that - // never heals itself. - rows, err := r.db.QueryContext(ctx, query, + AND delivery_created_at < now() - $5::interval)` + +// undeliveredAcceptanceArgs binds that filter. One helper, so the list and the +// count cannot bind $1..$5 in two different orders — which would leave two +// queries that both run and disagree. +func undeliveredAcceptanceArgs(staleAfter time.Duration) []any { + return []any{ string(DeliveryStateCancelled), string(DeliveryStatePoisoned), string(DeliveryStatePending), DeliveryHeldForSettlement, - fmt.Sprintf("%d seconds", int(staleAfter.Seconds()))) + fmt.Sprintf("%d seconds", int(staleAfter.Seconds())), + } +} + +func (r *postgresDivergences) UndeliveredAcceptances(ctx context.Context, staleAfter time.Duration) ([]UndeliveredAcceptance, error) { + // The comparison itself is undeliveredAcceptanceLatest + + // undeliveredAcceptanceFilter, above: this adds the columns an operator + // reads, a stable order, and the example budget. + query := undeliveredAcceptanceLatest + ` + SELECT post_uri, + community_did, + COALESCE(activity_id, ''), + COALESCE(delivery_state, ''), + COALESCE(last_error_class, '') + FROM latest` + undeliveredAcceptanceFilter + ` + ORDER BY post_uri + LIMIT $6` + + tx, done, err := r.readOnlyTx(ctx) + if err != nil { + return nil, fmt.Errorf("list undelivered acceptances: %w", err) + } + defer done() + + rows, err := tx.QueryContext(ctx, query, + append(undeliveredAcceptanceArgs(staleAfter), MaxDivergenceExamples)...) if err != nil { return nil, fmt.Errorf("list undelivered acceptances: %w", err) } @@ -298,41 +485,113 @@ func (r *postgresDivergences) UndeliveredAcceptances(ctx context.Context, staleA return found, nil } -func (r *postgresDivergences) RecastDivergences(ctx context.Context) ([]RecastDivergence, error) { - // THE EVIDENCE IS THE HISTORY, NEVER THE VOTE ROW. outbound_activities is - // append-only and its parent_at_uri carries the subject at-uri from both - // vote enqueue sites, so a Like/Dislike with a DELIVERED delivery is durable - // proof of what a peer was told — which is exactly what the vote row stops - // being the moment a re-cast rewrites it. - // - // The row appears below only as an EXCLUSION, and the direction matters: it - // can suppress a finding, never create one. If our ledger still names this - // exact activity as the live vote AND still calls it delivered, then we - // account for what the peer holds and there is nothing to reconcile. When - // the row has been reset, retracted or deleted — every shape this bug takes - // — the row cannot answer, and the history stands on its own. - // - // TWO INDEPENDENT EXCLUSIONS, because they answer different questions and - // each is the whole defence against a different way of ruining this report: - // - // the ledger still claims it — without this, EVERY cleanly delivered vote - // is a finding. That is most of the highest-volume table in the system, - // and a class that names all of it hides the one case that matters. - // a delivered Undo followed it — without this, a withdrawal that WORKED is - // reported forever. The history is append-only, so the delivered - // Like/Dislike never goes away; only the Undo beside it says the peer - // holds nothing now. - // - // The Undo comparison is by time rather than by id on purpose: an Undo names - // the activity it withdraws in its payload, but a re-cast mints new ids, so - // id-chasing would miss an Undo that withdrew an earlier incarnation of the - // same (actor, subject) vote. Only one vote per pair may be live at a time, - // which is what makes "any delivered Undo at or after this delivery" the - // right question. - // - // DISTINCT because one activity may have several deliveries (the fan-out - // schema); one delivered copy is one thing the peer holds. - query := ` +// UndeliveredAcceptanceCounts counts the same population, grouped by the +// delivery state the sweep classifies on. +// +// GROUPED RATHER THAN TOTALLED, because the three classes are three different +// jobs — a cancelled acceptance is a note, a stopped queue is a page — and a +// single total could only ever be alerted on at the noise level of whichever +// kind is most common. Grouping here also keeps the state-to-class mapping in +// the one place that already owns it (ingest.acceptanceClass): a state this +// sweep has no class for arrives as its own key rather than inflating a class +// that would misdescribe it. +func (r *postgresDivergences) UndeliveredAcceptanceCounts(ctx context.Context, staleAfter time.Duration) (map[DeliveryState]int, error) { + query := undeliveredAcceptanceLatest + ` + SELECT COALESCE(delivery_state, ''), count(*) + FROM latest` + undeliveredAcceptanceFilter + ` + GROUP BY 1` + + tx, done, err := r.readOnlyTx(ctx) + if err != nil { + return nil, fmt.Errorf("count undelivered acceptances: %w", err) + } + defer done() + + rows, err := tx.QueryContext(ctx, query, undeliveredAcceptanceArgs(staleAfter)...) + if err != nil { + return nil, fmt.Errorf("count undelivered acceptances: %w", err) + } + defer func() { _ = rows.Close() }() + + counts := make(map[DeliveryState]int) + for rows.Next() { + var state string + var count int + if err := rows.Scan(&state, &count); err != nil { + return nil, fmt.Errorf("scan undelivered acceptance count: %w", err) + } + counts[DeliveryState(state)] = count + } + if err := rows.Err(); err != nil { + return nil, fmt.Errorf("count undelivered acceptances: %w", err) + } + return counts, nil +} + +// recastDivergenceRows is the comparison, shared by the example list and the +// count so the two cannot come to describe different populations. +// +// THE EVIDENCE IS THE HISTORY, NEVER THE VOTE ROW. outbound_activities is +// append-only and its parent_at_uri carries the subject at-uri from both +// vote enqueue sites, so a Like/Dislike with a DELIVERED delivery is durable +// proof of what a peer was told — which is exactly what the vote row stops +// being the moment a re-cast rewrites it. +// +// The row appears below only as an EXCLUSION, and the direction matters: it +// can suppress a finding, never create one. If our ledger still names this +// exact activity as the live vote AND still calls it delivered, then we +// account for what the peer holds and there is nothing to reconcile. When +// the row has been reset, retracted or deleted — every shape this bug takes +// — the row cannot answer, and the history stands on its own. +// +// THREE INDEPENDENT EXCLUSIONS, because they answer different questions and +// each is the whole defence against a different way of ruining this report: +// +// the ledger still claims it — without this, EVERY cleanly delivered vote +// is a finding. That is most of the highest-volume table in the system, +// and a class that names all of it hides the one case that matters. +// a delivered Undo followed it — without this, a withdrawal that WORKED is +// reported forever. The history is append-only, so the delivered +// Like/Dislike never goes away; only the Undo beside it says the peer +// holds nothing now. +// a LATER DELIVERED VOTE followed it — without this, every successful vote +// FLIP is a finding, forever. A flip is an in-place upsert +// (consume.applyVoteWrite): current_activity_id moves to the new +// activity, delivered_state resets to pending, and NO Undo is enqueued, +// because Lemmy holds one vote per (person, object) and REPLACES it on a +// bare opposite vote. So once the new vote delivers, the old delivered +// activity satisfies neither exclusion above — the ledger names the new +// id and no Undo will ever join it — and the append-only history keeps it +// forever. A later delivered Like/Dislike for the same pair supersedes an +// earlier one EXACTLY as a delivered Undo does, and that is the only +// reason this is correct rather than merely convenient. +// +// Both time exclusions compare by TIME rather than by id on purpose: an Undo +// names the activity it withdraws in its payload, but a re-cast mints new +// ids, so id-chasing would miss an Undo that withdrew an earlier incarnation +// of the same (actor, subject) vote. Only one vote per pair may be live at a +// time, which is what makes "any delivered Undo at or after this delivery" +// the right question. +// +// The superseding-vote comparison is a TUPLE, (created_at, activity_id), +// because created_at cannot be trusted to separate them: it defaults to +// now(), which is the TRANSACTION timestamp, so any two vote activities +// written in one transaction carry it identically, and across transactions +// the column is only microsecond-resolution. A bare `>` would then exclude neither +// from the other and report BOTH; a bare `>=` would exclude both and report +// NEITHER, silently swallowing the poisoned re-cast this class exists for. +// The tuple gives a total order regardless of clock resolution, and the +// activity id is the tiebreak because it is the primary key — unique by +// construction, so the order is total and stable across sweeps. +// +// IT MUST NOT WEAKEN THE TRUE POSITIVE, and it does not: the later vote must +// itself be DELIVERED. A re-cast whose new delivery poisoned has no later +// delivered vote, so the old activity the peer is still counting is still +// reported — which is the entire point of the class. +// +// DISTINCT because one activity may have several deliveries (the fan-out +// schema); one delivered copy is one thing the peer holds. +const recastDivergenceRows = ` SELECT DISTINCT a.actor_did, a.parent_at_uri, a.activity_id FROM outbound_activities a JOIN outbound_deliveries d @@ -355,11 +614,39 @@ func (r *postgresDivergences) RecastDivergences(ctx context.Context) ([]RecastDi AND u.actor_did = a.actor_did AND u.parent_at_uri = a.parent_at_uri AND u.created_at >= a.created_at) - ORDER BY a.actor_did, a.parent_at_uri, a.activity_id` + AND NOT EXISTS ( + SELECT 1 + FROM outbound_activities n + JOIN outbound_deliveries nd + ON nd.activity_id = n.activity_id AND nd.state = $1 + WHERE n.kind IN ($2, $3) + AND n.actor_did = a.actor_did + AND n.parent_at_uri = a.parent_at_uri + AND (n.created_at, n.activity_id) > (a.created_at, a.activity_id))` - rows, err := r.db.QueryContext(ctx, query, +// recastDivergenceArgs binds that comparison's $1..$5. One helper, so the list +// and the count cannot bind them in two different orders. +func recastDivergenceArgs() []any { + return []any{ string(DeliveryStateDelivered), "Like", "Dislike", - string(DeliveredStateDelivered), "Undo") + string(DeliveredStateDelivered), "Undo", + } +} + +func (r *postgresDivergences) RecastDivergences(ctx context.Context) ([]RecastDivergence, error) { + // The comparison is recastDivergenceRows, above; this adds the operator's + // stable order and the example budget. + query := recastDivergenceRows + ` + ORDER BY a.actor_did, a.parent_at_uri, a.activity_id + LIMIT $6` + + tx, done, err := r.readOnlyTx(ctx) + if err != nil { + return nil, fmt.Errorf("list recast divergences: %w", err) + } + defer done() + + rows, err := tx.QueryContext(ctx, query, append(recastDivergenceArgs(), MaxDivergenceExamples)...) if err != nil { return nil, fmt.Errorf("list recast divergences: %w", err) } @@ -379,46 +666,154 @@ func (r *postgresDivergences) RecastDivergences(ctx context.Context) ([]RecastDi return found, nil } +// RecastDivergenceCount counts the same comparison, unbounded. +// +// It wraps the DISTINCT select rather than counting the join, because the +// DISTINCT is part of the definition: one activity may have several deliveries +// (the fan-out schema), and one delivered copy is ONE thing the peer holds. +// count(*) over the un-deduplicated join would report a number larger than the +// list it labels, on the exact rows an operator is trying to size. +func (r *postgresDivergences) RecastDivergenceCount(ctx context.Context) (int, error) { + query := `SELECT count(*) FROM (` + recastDivergenceRows + `) AS divergent` + tx, done, err := r.readOnlyTx(ctx) + if err != nil { + return 0, fmt.Errorf("count recast divergences: %w", err) + } + defer done() + + var total int + if err := tx.QueryRowContext(ctx, query, recastDivergenceArgs()...).Scan(&total); err != nil { + return 0, fmt.Errorf("count recast divergences: %w", err) + } + return total, nil +} + +// The poison classes the worker decides BEFORE any POST is attempted, declared +// HERE and written by internal/outbound/worker.go through these names. +// +// THE DECLARATION IS SHARED SO THE TWO SIDES CANNOT DRIFT IN SPELLING. The +// worker used to write these as string literals at its own call sites while the +// denylist below repeated them, and nothing compiled the two lists against each +// other: a respelling on either side would have silently moved four known +// non-deliveries into the unknown-outcome report, whose entire worth is that its +// numbers stay small enough to trust. The compiler now refuses that. What it +// still cannot catch is a FIFTH never-wire class introduced as a fresh literal +// — the worker's wire classes (transport, 4xx, 5xx …) are literals, so the +// surrounding style invites one — which is why a new class belongs in this block +// and in neverReachedTheWireClasses, and why the store test that seeds all four +// by name (divergence_unknown_test.go) is the pin on the set. +const ( + // PoisonClassParentUnaccepted: the causal wait budget expired and the parent + // was never accepted, so the reply was never offered to anyone. + PoisonClassParentUnaccepted = "parent_unaccepted" + // PoisonClassParentPoisoned: the parent delivery poisoned; a descendant + // cannot land, so it is not attempted. + PoisonClassParentPoisoned = "parent_poisoned" + // PoisonClassCrossAuthority: the stored target inbox is not same-authority + // with its community, and the worker refuses to sign a POST to it. + PoisonClassCrossAuthority = "cross_authority" + // PoisonClassSigner: SignerFor could not resolve a signing key, so + // SendActivityAs is never called. This one poisons only after the attempt + // budget is exhausted (releaseOrPoison), which makes it LOOK like a retried + // wire failure on every column; it is not one. + PoisonClassSigner = "signer" +) + +// neverReachedTheWireClasses are those four classes as the divergence sweep +// reads them. +// +// All four are written with last_status_code 0 — the same shape a dial timeout +// leaves behind, and the opposite meaning. Nothing was sent, so the peer's +// state is not unknown at all: they do not have it. That is the rule that +// already excludes a cancelled delivery, one step later in the worker. +// +// A DENYLIST RATHER THAN AN ALLOWLIST, deliberately. Every other poison class +// (transport, 4xx, 5xx, unauthorized, timeout, rate_limited, transient, +// inbox_gone, inbox_resolve) is recorded after a POST was attempted, so +// "unknown" is the honest default for a class this list has not heard of: an +// allowlist would silently drop a genuinely unknown delivery out of the report +// the day the worker grows a new wire class, and a divergence that vanishes is +// worse than one that is over-reported with its class name attached. The cost +// runs the other way — a new never-wire poison class must be added to the block +// above AND to this list, or it inflates these counts. +var neverReachedTheWireClasses = []string{ + PoisonClassParentUnaccepted, PoisonClassParentPoisoned, + PoisonClassCrossAuthority, PoisonClassSigner, +} + +// unknownDeliveryOutcomeRows is the population, shared by the example list and +// the counts. +// +// POISONED ONLY, and every other state is excluded for a reason that is not +// symmetry: +// +// delivered — the peer confirmed it. The one outcome here we DO know. +// cancelled — it never reached the wire: a ban or an opt-out took it out of +// the queue, so the peer does not have it and that is a fact, not a +// question. Cycle 2's cancelled class already names it. Folding it in +// would put a KNOWN non-delivery into the bucket whose entire meaning is +// that the answer is unavailable, and these are the two numbers an +// operator has to be able to trust as small. +// pending — still in flight, or held for settlement. Nothing is unresolved +// about a delivery that has not finished trying. +// +// AGE IS NOT PART OF THE DEFINITION, unlike the stale-acceptance class: a +// poisoned delivery is terminal the moment it poisons — nothing re-drives it +// — so a recent one and an old one are the same state and the same unknown. +// +// REFUSED IS COMPUTED FROM WHETHER A PEER ANSWERED, NOT FROM THE NUMBER, and +// this is the one subtle thing in the query. The column is nullable, but the +// worker does not use the NULL: classify() passes a literal 0 for a transport +// failure and for an unresolvable signer (worker.go), and MarkPoisoned writes +// that 0. So "the peer never answered" arrives in this table in TWO +// spellings, NULL and 0, and only a predicate that treats both as silence +// keeps a dial timeout out of the sub-count that says the peer spoke. It is +// selected as a boolean rather than left to the caller to infer, so that no +// reader downstream can rediscover the wrong rule from LastStatusCode. +// +// AND IT MUST HAVE BEEN SENT. Four poison classes are decided before any +// POST — see neverReachedTheWireClasses — and they carry status 0 exactly +// like a transport failure does. Filing them here would put a KNOWN +// non-delivery in the bucket whose whole meaning is that the answer is +// unavailable, inflate the one number whose worth depends on staying small, +// and send an operator to ask a stranger about a request that never left +// this process. last_error_class is NOT NULL DEFAULT ” (migration 020), so +// `<> ALL` cannot go NULL and quietly drop a row. +const unknownDeliveryOutcomeRows = ` + FROM outbound_deliveries d + JOIN outbound_activities a ON a.activity_id = d.activity_id + WHERE d.state = $1 + AND d.last_error_class <> ALL($2)` + +// unknownDeliveryRefused is the refused/unanswered discriminator, written once +// so the list and the counts cannot disagree about which sub-count a row +// belongs to. NULL and 0 are both silence — see above. +const unknownDeliveryRefused = `COALESCE(d.last_status_code, 0) > 0` + +// unknownDeliveryOutcomeArgs binds $1..$2 for both statements. +func unknownDeliveryOutcomeArgs() []any { + return []any{string(DeliveryStatePoisoned), pq.Array(neverReachedTheWireClasses)} +} + func (r *postgresDivergences) UnknownDeliveryOutcomes(ctx context.Context) ([]UnknownDeliveryOutcome, error) { - // POISONED ONLY, and every other state is excluded for a reason that is not - // symmetry: - // - // delivered — the peer confirmed it. The one outcome here we DO know. - // cancelled — it never reached the wire: a ban or an opt-out took it out of - // the queue, so the peer does not have it and that is a fact, not a - // question. Cycle 2's cancelled class already names it. Folding it in - // would put a KNOWN non-delivery into the bucket whose entire meaning is - // that the answer is unavailable, and these are the two numbers an - // operator has to be able to trust as small. - // pending — still in flight, or held for settlement. Nothing is unresolved - // about a delivery that has not finished trying. - // - // AGE IS NOT PART OF THE DEFINITION, unlike the stale-acceptance class: a - // poisoned delivery is terminal the moment it poisons — nothing re-drives it - // — so a recent one and an old one are the same state and the same unknown. - // - // REFUSED IS COMPUTED FROM WHETHER A PEER ANSWERED, NOT FROM THE NUMBER, and - // this is the one subtle thing in the query. The column is nullable, but the - // worker does not use the NULL: classify() passes a literal 0 for a transport - // failure and for an unresolvable signer (worker.go), and MarkPoisoned writes - // that 0. So "the peer never answered" arrives in this table in TWO - // spellings, NULL and 0, and only a predicate that treats both as silence - // keeps a dial timeout out of the sub-count that says the peer spoke. It is - // selected as a boolean rather than left to the caller to infer, so that no - // reader downstream can rediscover the wrong rule from LastStatusCode. query := ` SELECT d.activity_id, d.target_inbox, a.kind, d.last_error_class, COALESCE(d.last_status_code, 0), - COALESCE(d.last_status_code, 0) > 0 - FROM outbound_deliveries d - JOIN outbound_activities a ON a.activity_id = d.activity_id - WHERE d.state = $1 - ORDER BY d.activity_id, d.target_inbox` + ` + unknownDeliveryRefused + unknownDeliveryOutcomeRows + ` + ORDER BY d.activity_id, d.target_inbox + LIMIT $3` + + tx, done, err := r.readOnlyTx(ctx) + if err != nil { + return nil, fmt.Errorf("list unknown delivery outcomes: %w", err) + } + defer done() - rows, err := r.db.QueryContext(ctx, query, string(DeliveryStatePoisoned)) + rows, err := tx.QueryContext(ctx, query, + append(unknownDeliveryOutcomeArgs(), MaxDivergenceExamples)...) if err != nil { return nil, fmt.Errorf("list unknown delivery outcomes: %w", err) } @@ -438,3 +833,46 @@ func (r *postgresDivergences) UnknownDeliveryOutcomes(ctx context.Context) ([]Un } return found, nil } + +// UnknownDeliveryOutcomeCounts counts the same population, split by whether a +// peer answered. +// +// The split is made HERE, by the same expression the list selects Refused with, +// rather than by counting twice with two predicates: two spellings of "the peer +// spoke" is exactly how a dial timeout ends up in the sub-count that says it +// did. GROUP BY over that one boolean cannot produce a third bucket, and a +// missing group is an honest zero — no rows of that kind exist. +func (r *postgresDivergences) UnknownDeliveryOutcomeCounts(ctx context.Context) (UnknownDeliveryCounts, error) { + query := `SELECT ` + unknownDeliveryRefused + `, count(*)` + unknownDeliveryOutcomeRows + ` + GROUP BY 1` + + tx, done, err := r.readOnlyTx(ctx) + if err != nil { + return UnknownDeliveryCounts{}, fmt.Errorf("count unknown delivery outcomes: %w", err) + } + defer done() + + rows, err := tx.QueryContext(ctx, query, unknownDeliveryOutcomeArgs()...) + if err != nil { + return UnknownDeliveryCounts{}, fmt.Errorf("count unknown delivery outcomes: %w", err) + } + defer func() { _ = rows.Close() }() + + var counts UnknownDeliveryCounts + for rows.Next() { + var refused bool + var count int + if err := rows.Scan(&refused, &count); err != nil { + return UnknownDeliveryCounts{}, fmt.Errorf("scan unknown delivery outcome count: %w", err) + } + if refused { + counts.Refused = count + continue + } + counts.Unanswered = count + } + if err := rows.Err(); err != nil { + return UnknownDeliveryCounts{}, fmt.Errorf("count unknown delivery outcomes: %w", err) + } + return counts, nil +} diff --git a/internal/store/divergence_acceptance_test.go b/internal/store/divergence_acceptance_test.go index cf3ad1c..0ad6a4c 100644 --- a/internal/store/divergence_acceptance_test.go +++ b/internal/store/divergence_acceptance_test.go @@ -359,3 +359,220 @@ func seedFollowUpDelivery(t *testing.T, database *sql.DB, post dvPost, kind stri fmt.Sprintf("%d seconds", int(age.Seconds()))) require.NoError(t, err, "seed follow-up delivery for %s", post.rkey) } + +// --------------------------------------------------------------------------- +// The ledger term: only an ACCEPTED post is an acceptance +// --------------------------------------------------------------------------- + +// TestUndeliveredAcceptances_ARejectedPostIsNotAnUndeliveredAcceptance pins +// `adm.status = 'accepted'`. +// +// Every other fixture in this file writes `accepted`, so the term deletes green +// — and the state it excludes is not exotic. admissions.status has five values, +// and a REJECTED post is the ordinary outcome of an opted-out author, a banned +// author, a missing title or a locked parent. Its shape is exactly the reported +// one: no acceptance record was written, so nothing stamped accepted_at, and any +// delivery it had was cancelled. +// +// The difference is the whole meaning of the class. An accepted post that never +// reached the peer is a DISAGREEMENT — Coves shows it in the community, Lemmy +// does not. A rejected post is AGREEMENT: it is invisible in both places, +// exactly as decided, and Coves' post.getStatus already tells its author why. +// Reporting it turns every refusal the admission policy makes into a finding, +// and the volume of refusals is not small. +func TestUndeliveredAcceptances_ARejectedPostIsNotAnUndeliveredAcceptance(t *testing.T) { + database := acceptanceTestDB(t) + + for _, status := range []string{"rejected", "removed", "pending", "pending_reacceptance"} { + t.Run(status, func(t *testing.T) { + database := acceptanceTestDB(t) + + undecided := dvPost{rkey: "3lzdvstat0001", state: DeliveryStateCancelled, age: 2 * time.Hour} + seedAcceptedPost(t, database, undecided) + setAdmissionStatus(t, database, undecided, status) + + // A genuinely accepted one beside it, so an empty result cannot pass + // for the right answer. + accepted := dvPost{rkey: "3lzdvstat0002", state: DeliveryStateCancelled, age: 2 * time.Hour} + seedAcceptedPost(t, database, accepted) + + found, err := NewDivergences(database).UndeliveredAcceptances(context.Background(), dvStaleAfter) + require.NoError(t, err) + assert.Equal(t, []string{accepted.uri()}, urisOf(found), + "a post whose admission is %q is not an ACCEPTANCE that failed to reach the peer: "+ + "no acceptance record was ever written, so the community view agrees with "+ + "Lemmy's — both hide it, which is the decision working. Only 'accepted' means "+ + "the two sides disagree, and without that term every refusal the admission "+ + "policy makes becomes a finding", status) + }) + } + _ = database +} + +// setAdmissionStatus rewrites the ledger decision for one post. +func setAdmissionStatus(t *testing.T, database *sql.DB, post dvPost, status string) { + t.Helper() + result, err := database.ExecContext(context.Background(), + `UPDATE admissions SET status = $2 WHERE post_uri = $1`, post.uri(), status) + require.NoError(t, err) + affected, err := result.RowsAffected() + require.NoError(t, err) + require.EqualValues(t, 1, affected, "there must be an admission row for %s to re-decide", post.rkey) +} + +// --------------------------------------------------------------------------- +// The newest attempt decides +// --------------------------------------------------------------------------- + +// TestUndeliveredAcceptances_AFreshRetryAfterACancelledAttemptIsInFlight pins +// the ORDER of the DISTINCT ON: newest delivery first. +// +// Flip it to oldest-first and every assertion in this file still passes, because +// no other fixture gives one post two deliveries in different states. This one +// does, and it is the ordinary shape of recovery: an attempt was cancelled, and +// a newer one is on its way. +// +// Read oldest-first, the report describes a post as permanently undelivered +// while the delivery that will carry it is in the queue right now — and the +// operator's response to "cancelled" (accept it, or re-enqueue by hand) is +// exactly wrong for a post that needs neither. +func TestUndeliveredAcceptances_AFreshRetryAfterACancelledAttemptIsInFlight(t *testing.T) { + database := acceptanceTestDB(t) + + retried := dvPost{rkey: "3lzdvordr0001", state: DeliveryStateCancelled, age: 3 * time.Hour} + seedAcceptedPost(t, database, retried) + // The newer attempt: pending, created NOW, so it is inside any staleness + // window the sweep could be configured with. + seedFollowUpDelivery(t, database, retried, "Create", DeliveryStatePending, "", 0) + + // And a post with ONLY the cancelled attempt, so "reports nothing" cannot + // pass for the right answer. + abandoned := dvPost{rkey: "3lzdvordr0002", state: DeliveryStateCancelled, age: 3 * time.Hour} + seedAcceptedPost(t, database, abandoned) + + found, err := NewDivergences(database).UndeliveredAcceptances(context.Background(), dvStaleAfter) + require.NoError(t, err) + + assert.Equal(t, []string{abandoned.uri()}, urisOf(found), + "the NEWEST delivery decides. This post has a cancelled attempt and a fresh one in "+ + "flight, so it is not a divergence — it is a retry. Ordering the other way describes "+ + "it as permanently undelivered while the delivery that will carry it sits in the "+ + "queue, and sends an operator to re-enqueue a post that needs nothing") +} + +// --------------------------------------------------------------------------- +// The count is the MEASUREMENT; the list is a page of it +// --------------------------------------------------------------------------- + +// TestUndeliveredAcceptanceCounts_MeasureThePopulationTheExamplesOnlySampleFrom +// is the reason these are two statements rather than one truncated read. +// +// Every other fixture in this file is smaller than the page, so the list IS its +// own count and the two queries are indistinguishable. The state that separates +// them is the only one that matters: a broken community or a stopped queue, +// which is what an outage looks like here. If the count were bounded with the +// examples, "500" would mean both "500" and "a catastrophe" — and an operator +// sizes an incident from exactly this number, then decides whether to page +// anyone. +// +// IT ALSO PINS THAT THE TWO STATEMENTS DESCRIBE ONE POPULATION. They share the +// CTE and the filter as constants (undeliveredAcceptanceLatest, +// undeliveredAcceptanceFilter) precisely so a count cannot come to be taken over +// a slightly different join than the list it labels — a drift that is silent in +// the direction that matters, because nothing about the report would look wrong. +// The three states are seeded at three different sizes so a grouped count cannot +// be satisfied by a total, and the delivered sibling is here so the population +// is a comparison rather than a listing. +func TestUndeliveredAcceptanceCounts_MeasureThePopulationTheExamplesOnlySampleFrom(t *testing.T) { + database := acceptanceTestDB(t) + ctx := context.Background() + + // One class far past the page, and two well inside it — three sizes, so a + // count that reported a total, or one class's number for another's, fails + // here rather than reading plausibly. + const overflowing = MaxDivergenceExamples + 1 + seedAcceptedPostsInBulk(t, database, "cncl", overflowing, DeliveryStateCancelled, "", 2*time.Hour) + seedAcceptedPostsInBulk(t, database, "pois", 2, DeliveryStatePoisoned, "4xx", 2*time.Hour) + seedAcceptedPostsInBulk(t, database, "stal", 3, DeliveryStatePending, "", 2*time.Hour) + // THE FALSE-POSITIVE CONTROL, as everywhere else in this file: a post the + // peer really took. A count is the number an operator escalates on, so an + // inflated one is worse than a missing one. + landed := dvPost{rkey: "3lzdvbulk0000d", state: DeliveryStateDelivered, age: 2 * time.Hour, accepted: true} + seedAcceptedPost(t, database, landed) + + divergences := NewDivergences(database) + + found, err := divergences.UndeliveredAcceptances(ctx, dvStaleAfter) + require.NoError(t, err) + require.Len(t, found, MaxDivergenceExamples, + "the EXAMPLES are bounded at the database, not in Go: the sweep that reaches this limit "+ + "is the sweep running against a stopped queue, and materialising every row before "+ + "cutting it down would put the memory spike in the process an operator is trying to "+ + "keep alive") + + counts, err := divergences.UndeliveredAcceptanceCounts(ctx, dvStaleAfter) + require.NoError(t, err) + + assert.Equal(t, overflowing, counts[DeliveryStateCancelled], + "the COUNT is the true measurement, unbounded by the example budget") + assert.Greater(t, counts[DeliveryStateCancelled], len(found), + "and it EXCEEDS the examples, which is the whole contract: a count capped at the size of "+ + "a page makes %d mean both '%d' and 'a catastrophe', and the number an operator sizes "+ + "the incident from would stop growing exactly when the incident does", + MaxDivergenceExamples, MaxDivergenceExamples) + assert.Equal(t, 2, counts[DeliveryStatePoisoned], + "while the small classes keep their own exact numbers: the count is GROUPED by delivery "+ + "state because the three classes are three different jobs, and a single total could "+ + "only ever be alerted on at the noise level of whichever kind is most common") + assert.Equal(t, 3, counts[DeliveryStatePending]) + assert.NotContains(t, counts, DeliveryStateDelivered, + "and the post the peer TOOK is in no bucket at all: it is the healthy state, and a count "+ + "that includes it is a count of every acceptance the bridge has ever made") +} + +// seedAcceptedPostsInBulk writes n accepted posts sharing one delivery state. +// +// FOUR STATEMENTS, NOT n×4 ROUND TRIPS: the population that matters here is +// larger than the example budget, and five hundred round trips is a slow way to +// say "an outage". The rows are otherwise identical to seedAcceptedPost's — +// including the payload's object id, which is the correspondence the sweep +// joins on and the worker stamps accepted_at through. +func seedAcceptedPostsInBulk(t *testing.T, database *sql.DB, tag string, n int, + state DeliveryState, errorClass string, age time.Duration) { + t.Helper() + ctx := context.Background() + uri := "'at://" + dvAuthorDID + "/social.coves.community.postv2/3lzdvbulk" + tag + "' || i" + objectID := "'" + dvUserOrigin + "/ap/object/" + dvAuthorDID + + "/social.coves.community.postv2/3lzdvbulk" + tag + "' || i" + activityID := "'" + dvUserOrigin + "/ap/activity/3lzdvbulk" + tag + "' || i" + + _, err := database.ExecContext(ctx, ` + INSERT INTO outbound_objects (at_uri, ap_object_id, community_did, community_ap_id, + translated_snapshot, accepted_at) + SELECT `+uri+`, `+objectID+`, $1, $2, '{"type":"Page"}'::jsonb, NULL + FROM generate_series(1, $3) AS i`, dvCommunityDID, dvCommunityAPID, n) + require.NoError(t, err, "seed %d outbound_objects", n) + + _, err = database.ExecContext(ctx, ` + INSERT INTO admissions (community_did, post_uri, author_did, status, decision_code) + SELECT $1, `+uri+`, $2, 'accepted', '' + FROM generate_series(1, $3) AS i`, dvCommunityDID, dvAuthorDID, n) + require.NoError(t, err, "seed %d admissions", n) + + _, err = database.ExecContext(ctx, ` + INSERT INTO outbound_activities (activity_id, actor_did, kind, payload) + SELECT `+activityID+`, $1, 'Create', + jsonb_build_object('id', `+activityID+`, 'type', 'Create', + 'object', jsonb_build_object('type', 'Page', 'id', `+objectID+`)) + FROM generate_series(1, $2) AS i`, dvAuthorDID, n) + require.NoError(t, err, "seed %d outbound_activities", n) + + _, err = database.ExecContext(ctx, ` + INSERT INTO outbound_deliveries (activity_id, target_inbox, ordering_key, state, + last_error_class, created_at, updated_at) + SELECT `+activityID+`, $1, $2, $3, $4, now() - $5::interval, now() - $5::interval + FROM generate_series(1, $6) AS i`, + dvCommunityInbox, dvCommunityAPID, string(state), errorClass, + fmt.Sprintf("%d seconds", int(age.Seconds())), n) + require.NoError(t, err, "seed %d outbound_deliveries", n) +} diff --git a/internal/store/divergence_test.go b/internal/store/divergence_test.go index be324a8..3e678d7 100644 --- a/internal/store/divergence_test.go +++ b/internal/store/divergence_test.go @@ -234,3 +234,173 @@ func TestPersonaVoteEvents_ANonCanonicalSpellingIsNotMatched_KNOWNNARROW(t *test "test ever fails because the join learned to normalize, that is an improvement: "+ "delete this test and say so in the report") } + +// --------------------------------------------------------------------------- +// 17e review — the exclusion's SECOND conjunct, decided rather than inherited +// --------------------------------------------------------------------------- + +// TestRecastDivergences_ALedgerThatNamesTheVoteButCallsItPendingIsReported +// decides a state the query has an opinion about and nobody wrote down. +// +// The re-cast exclusion suppresses a finding only when the ledger BOTH names +// this exact activity as the live vote AND calls it delivered. Drop the second +// conjunct and every fixture in the suite stays green, because no fixture has +// ever produced the in-between: current_activity_id pointing at an activity a +// delivery row says was delivered, while delivered_state reads 'pending'. +// +// THE ANSWER IS THAT IT MUST BE REPORTED, and the reason is the same one that +// makes this class worth having. The delivery is the durable fact — the peer +// accepted that activity — and the ledger is the number the reseed reads. A row +// saying 'pending' subtracts nothing, so the community's served score counts a +// vote of ours as a stranger's, permanently, and the vote row that would +// normally be corrected by voteCallback is the one thing that already failed to +// be. "Our ledger half-claims it" is not an account of what the peer holds; only +// a full claim is. +// +// NO WRITER PRODUCES THIS TODAY — voteCallback flips the row to delivered on the +// same success that marks the delivery, and a re-cast moves the id rather than +// resetting the state under it. That is exactly why it is written here: an +// implementation that dropped the conjunct would look identical on every state +// the system currently reaches, and would then swallow this one silently on the +// day some future path leaves the pair half-written. +func TestRecastDivergences_ALedgerThatNamesTheVoteButCallsItPendingIsReported(t *testing.T) { + database := divergenceTestDB(t) + ctx := context.Background() + + subject := "at://" + dvPersonaDID + "/social.coves.community.postv2/3lzdvrcpost1" + activityID := "https://coves.social/ap/activity/dv-recast-halfclaimed" + seedDeliveredVoteActivity(t, database, activityID, dvPersonaDID, subject, "Like") + + // The ledger names THIS activity as the live vote — and calls it pending. + _, err := database.ExecContext(ctx, ` + INSERT INTO outbound_votes (vote_at_uri, actor_did, subject_at_uri, subject_ap_id, + community_did, direction, current_activity_id, delivered_state) + VALUES ($1, $2, $3, $4, $5, 'up', $6, 'pending')`, + "at://"+dvPersonaDID+"/social.coves.feed.vote/3lzdvrcvote1", dvPersonaDID, subject, + "https://lemmy.world/post/9002", dvCommunityDID, activityID) + require.NoError(t, err) + + found, err := NewDivergences(database).RecastDivergences(ctx) + require.NoError(t, err) + require.Len(t, found, 1, + "a HALF-CLAIM is not a claim. The delivery row is durable proof the peer accepted this "+ + "activity, while the ledger says 'pending' — which the reseed reads as 'subtract "+ + "nothing', so the community's served score counts our own vote as a stranger's "+ + "forever. The exclusion exists for the case where our accounting is COMPLETE; "+ + "anything less is the divergence itself, and dropping the delivered_state conjunct "+ + "makes this state silently disappear from a report that has no other way to see it") + assert.Equal(t, activityID, found[0].DeliveredActivityID) + assert.Equal(t, subject, found[0].SubjectATURI) +} + +// TestRecastDivergenceCount_MeasuresThePopulationTheExamplesOnlySampleFrom is +// this comparison's half of the same contract the acceptance counts carry: +// Counts is the true measurement, Entries is a page of it. +// +// The bulk shape is not hypothetical for this class. Votes are the +// highest-volume thing the bridge sends, and the population here is "a peer is +// holding a vote we no longer claim" — which arrives in bulk exactly when it +// matters: an instance that went away mid-flight poisons every re-cast and Undo +// aimed at it, and the count is how an operator learns whether that is three +// votes or fifty thousand. A count bounded with the examples would report the +// size of a page and stop moving while the problem grew. +// +// THE COUNT WRAPS THE SAME DISTINCT SELECT the list runs, and the DISTINCT is +// part of the definition rather than tidiness: one activity may have several +// deliveries (the fan-out schema) and one delivered copy is ONE thing the peer +// holds. Counting the un-deduplicated join would report a number larger than the +// list it labels, on the rows an operator is sizing. +func TestRecastDivergenceCount_MeasuresThePopulationTheExamplesOnlySampleFrom(t *testing.T) { + database := divergenceTestDB(t) + ctx := context.Background() + + const overflowing = MaxDivergenceExamples + 1 + seedRecastDivergencesInBulk(t, database, overflowing) + + // THE EXCLUSION CONTROL, because a count is only worth pinning if the + // population it counts is a comparison. This vote is delivered AND our ledger + // still names it as the live, delivered vote — so we account for what the peer + // holds and there is nothing to reconcile. Without it, a query that reported + // every delivered vote would satisfy every assertion below while naming most + // of the highest-volume table in the system. + covered := "https://coves.social/ap/activity/dv-count-covered" + coveredSubject := "at://" + dvPersonaDID + "/social.coves.community.postv2/3lzdvcntkeep" + seedDeliveredVoteActivity(t, database, covered, dvPersonaDID, coveredSubject, "Like") + _, err := database.ExecContext(ctx, ` + INSERT INTO outbound_votes (vote_at_uri, actor_did, subject_at_uri, subject_ap_id, + community_did, direction, current_activity_id, delivered_state) + VALUES ($1, $2, $3, $4, $5, 'up', $6, 'delivered')`, + "at://"+dvPersonaDID+"/social.coves.feed.vote/3lzdvcntkeep", dvPersonaDID, coveredSubject, + "https://lemmy.world/post/9003", dvCommunityDID, covered) + require.NoError(t, err) + + divergences := NewDivergences(database) + + found, err := divergences.RecastDivergences(ctx) + require.NoError(t, err) + require.Len(t, found, MaxDivergenceExamples, + "the examples are bounded at the database by the same budget the report spends") + for _, entry := range found { + assert.NotEqual(t, covered, entry.DeliveredActivityID, + "and the fully-accounted vote is not among them: our ledger names this exact activity "+ + "as the live vote AND calls it delivered, which is the one state that says the "+ + "peer holds nothing we have not accounted for") + } + + total, err := divergences.RecastDivergenceCount(ctx) + require.NoError(t, err) + assert.Equal(t, overflowing, total, + "the COUNT is exact and unbounded — the same comparison, with no LIMIT — so it counts the "+ + "population rather than the page") + assert.Greater(t, total, len(found), + "and it EXCEEDS the examples. That is the contract these two statements exist to keep: "+ + "the list is a sample an operator investigates, the count is the size they escalate "+ + "on, and a count capped at %d would stop growing exactly when the incident does", + MaxDivergenceExamples) +} + +// seedRecastDivergencesInBulk writes n delivered vote activities that NOTHING in +// the ledger accounts for — the shape this class reports. +// +// Each carries its OWN subject, so no two supersede each other: the comparison +// excludes an earlier delivered vote when a LATER delivered vote for the same +// (actor, subject) pair follows it, which is how an ordinary vote flip stays out +// of the report. Sharing one subject across the bulk would leave exactly one row +// standing and quietly turn this into a fixture of size 1. +func seedRecastDivergencesInBulk(t *testing.T, database *sql.DB, n int) { + t.Helper() + ctx := context.Background() + activityID := "'https://coves.social/ap/activity/dv-bulk-recast-' || i" + subject := "'at://" + dvPersonaDID + "/social.coves.community.postv2/3lzdvbulkrc' || i" + + _, err := database.ExecContext(ctx, ` + INSERT INTO outbound_activities (activity_id, actor_did, kind, payload, parent_at_uri) + SELECT `+activityID+`, $1, 'Like', '{"type":"Like"}'::jsonb, `+subject+` + FROM generate_series(1, $2) AS i`, dvPersonaDID, n) + require.NoError(t, err, "seed %d vote activities", n) + + _, err = database.ExecContext(ctx, ` + INSERT INTO outbound_deliveries (activity_id, target_inbox, ordering_key, state, + last_status_code, delivered_at) + SELECT `+activityID+`, $1, $2, 'delivered', 202, now() + FROM generate_series(1, $3) AS i`, + "https://lemmy.world/inbox", "https://lemmy.world/c/technology", n) + require.NoError(t, err, "seed %d deliveries", n) +} + +// seedDeliveredVoteActivity writes one vote activity with a delivered delivery — +// the durable evidence a peer was told. +func seedDeliveredVoteActivity(t *testing.T, database *sql.DB, activityID, actorDID, subject, kind string) { + t.Helper() + ctx := context.Background() + _, err := database.ExecContext(ctx, ` + INSERT INTO outbound_activities (activity_id, actor_did, kind, payload, parent_at_uri) + VALUES ($1, $2, $3, '{"type":"Like"}'::jsonb, $4)`, activityID, actorDID, kind, subject) + require.NoError(t, err) + _, err = database.ExecContext(ctx, ` + INSERT INTO outbound_deliveries (activity_id, target_inbox, ordering_key, state, + last_status_code, delivered_at) + VALUES ($1, $2, $3, 'delivered', 202, now())`, + activityID, "https://lemmy.world/inbox", "https://lemmy.world/c/technology") + require.NoError(t, err) +} diff --git a/internal/store/divergence_unknown_test.go b/internal/store/divergence_unknown_test.go index b0d6263..bf4ecff 100644 --- a/internal/store/divergence_unknown_test.go +++ b/internal/store/divergence_unknown_test.go @@ -146,9 +146,8 @@ func TestUnknownDeliveryOutcomes_ATransportFailureIsUnansweredAsProductionWrites // The two shapes the worker really produces. timeoutID := "https://coves.social/ap/activity/dv-unknown-timeout" - signerID := "https://coves.social/ap/activity/dv-unknown-signer" seedPoisonedDelivery(t, database, timeoutID, "Create", dvUnknownInboxA, "transport", 0) - seedPoisonedDelivery(t, database, signerID, "Create", dvUnknownInboxB, "signer", 0) + seedPoisonedDelivery(t, database, timeoutID+"-tls", "Create", dvUnknownInboxB, "transport", 0) // And a genuine refusal beside them, so "everything is unanswered" cannot // pass either: the two must come apart. refusedID := "https://coves.social/ap/activity/dv-unknown-503" @@ -168,10 +167,9 @@ func TestUnknownDeliveryOutcomes_ATransportFailureIsUnansweredAsProductionWrites "passes, not NULL — and it means NO ANSWER CAME BACK. A predicate that tests for "+ "NULL instead of for a real status files every dial timeout on the network under "+ "the count that says the peer replied, which is precisely the answer nobody gave") - assert.False(t, byID[signerID].Refused, - "and so is a delivery that never reached the wire at all: an unresolvable signer is "+ - "written with the same 0, and the peer certainly did not answer a request we never "+ - "sent") + assert.False(t, byID[timeoutID+"-tls"].Refused, + "and so is a TLS failure mid-handshake: same literal 0, same silence about whether the "+ + "peer ever saw the request") assert.True(t, byID[refusedID].Refused, "while a real 503 IS an answer: the two must come apart on the value, or the "+ "distinction the sub-counts exist for does not exist") @@ -237,3 +235,148 @@ func TestUnknownDeliveryOutcomes_ACancelledDeliveryIsNotUnknown(t *testing.T) { "not have it. Counting it as unknown inflates the one number whose value is that "+ "it is small and honest") } + +// --------------------------------------------------------------------------- +// "We do not know" requires that we SENT it +// --------------------------------------------------------------------------- + +// TestUnknownDeliveryOutcomes_APoisonThatNeverReachedTheWireIsNotUnknown is the +// completeness half of this class, and it is the same rule that excludes a +// cancelled delivery — one step later in the worker. +// +// Four poison classes are decided BEFORE any POST is made (worker.go): a causal +// wait that expired (parent_unaccepted), a parent that poisoned +// (parent_poisoned), an inbox that is not same-authority with its community +// (cross_authority), and a signer that would not resolve (signer). All four +// record status 0, which is the same shape a dial timeout leaves behind — and +// they mean the opposite. Nothing was ever sent, so the peer's state is not +// unknown at all: they do not have it. +// +// The cost of getting this wrong is precise. The unanswered bucket's Detail +// names an instance for the operator to go ASK, and its whole value is that the +// number is small and every row in it is a real question. Filing decisions the +// bridge made itself under "we do not know" both inflates that number and sends +// someone to ask a stranger about a request that never left the building. +func TestUnknownDeliveryOutcomes_APoisonThatNeverReachedTheWireIsNotUnknown(t *testing.T) { + database := acceptanceTestDB(t) + + // The four never-wire classes, exactly as the worker writes them. + neverSent := []string{"parent_unaccepted", "parent_poisoned", "cross_authority", "signer"} + for _, class := range neverSent { + seedPoisonedDelivery(t, database, + "https://coves.social/ap/activity/dv-neverwire-"+class, + "Create", dvUnknownInboxA, class, 0) + } + + // Two genuine unknowns beside them: one that got no answer, one refused. + // Without these the assertion would be satisfied by a query that returned + // nothing at all. + sentSilently := "https://coves.social/ap/activity/dv-neverwire-transport" + sentRefused := "https://coves.social/ap/activity/dv-neverwire-4xx" + seedPoisonedDelivery(t, database, sentSilently, "Create", dvUnknownInboxA, "transport", 0) + seedPoisonedDelivery(t, database, sentRefused, "Create", dvUnknownInboxB, "4xx", 422) + + found, err := NewDivergences(database).UnknownDeliveryOutcomes(context.Background()) + require.NoError(t, err) + + ids := make([]string, 0, len(found)) + for _, outcome := range found { + ids = append(ids, outcome.ActivityID) + } + assert.ElementsMatch(t, []string{sentSilently, sentRefused}, ids, + "only deliveries that REACHED THE WIRE are unknown. %v are all decided before any POST "+ + "is made — a causal wait that expired, a poisoned parent, a cross-authority inbox, "+ + "an unresolvable signer — so the peer does not have them and nothing about their "+ + "outcome is uncertain. They carry status 0 like a dial timeout does, which is the "+ + "whole trap: identical column, opposite meaning. Counting them inflates the one "+ + "number whose value is that it is small, and the unanswered bucket's own detail "+ + "sends an operator to ask an instance about a request that never left this process", + neverSent) +} + +// --------------------------------------------------------------------------- +// The counts are the measurement; the list is a page of it +// --------------------------------------------------------------------------- + +// TestUnknownDeliveryOutcomeCounts_MeasureThePopulationTheExamplesOnlySampleFrom +// pins the last of the four count/list pairs, and this class is where the +// bounded read is least obvious and most needed. +// +// A peer that goes away poisons everything aimed at it, so this population +// arrives in bulk by nature: one unreachable instance is thousands of rows, all +// with the same target inbox and the same silence. The two sub-counts are the +// numbers an operator judges that by, and their whole worth is that they are +// small and honest — a count bounded at the page would read as five hundred +// open questions no matter how many there were. +// +// THE SPLIT SURVIVES THE BULK, which is the second thing under test. The counts +// are grouped by the same expression the list selects Refused with +// (unknownDeliveryRefused), so five hundred silent transport failures cannot +// tip into the sub-count that says the peer answered — the exact drift that +// would happen if the two statements spelled "the peer spoke" differently. +func TestUnknownDeliveryOutcomeCounts_MeasureThePopulationTheExamplesOnlySampleFrom(t *testing.T) { + database := acceptanceTestDB(t) + ctx := context.Background() + + // One unreachable instance: more silent deliveries than the page holds. + const overflowing = MaxDivergenceExamples + 1 + seedPoisonedDeliveriesInBulk(t, database, overflowing) + // And two that a peer really answered, so the split cannot be satisfied by a + // total and the refused bucket has something in it to be wrong about. + refusedA := "https://coves.social/ap/activity/dv-count-refused-a" + refusedB := "https://coves.social/ap/activity/dv-count-refused-b" + seedPoisonedDelivery(t, database, refusedA, "Create", dvUnknownInboxA, "4xx", 422) + seedPoisonedDelivery(t, database, refusedB, "Create", dvUnknownInboxB, "5xx", 503) + // The outcome we DO know, which belongs in neither number. + seedPoisonedSibling(t, database) + + divergences := NewDivergences(database) + + found, err := divergences.UnknownDeliveryOutcomes(ctx) + require.NoError(t, err) + require.Len(t, found, MaxDivergenceExamples, + "the examples are bounded at the database: an unreachable instance is thousands of rows, "+ + "and a report that grows with the outage cannot be read during one") + + counts, err := divergences.UnknownDeliveryOutcomeCounts(ctx) + require.NoError(t, err) + + assert.Equal(t, overflowing, counts.Unanswered, + "the COUNT is the true measurement: every one of these got no answer at all, and the "+ + "example budget bounds what an operator can read rather than what the sweep measured") + assert.Greater(t, counts.Unanswered, len(found), + "and it EXCEEDS the examples, which is the contract — a count capped at %d would report "+ + "the size of a page while an instance was silently dropping everything we sent it", + MaxDivergenceExamples) + assert.Equal(t, 2, counts.Refused, + "while the peers that ANSWERED stay their own number. The split is the only information "+ + "these rows carry, and it is computed by the same expression the list selects Refused "+ + "with: two spellings of 'the peer spoke' is precisely how five hundred dial timeouts "+ + "end up in the count that says they replied") +} + +// seedPoisonedDeliveriesInBulk writes n poisoned deliveries that got no answer +// at all — one unreachable instance, as the worker records it. +// +// The status is a literal 0 rather than NULL because that is what production +// writes for a transport failure; a fixture that spelled it NULL here would let +// a NULL-based discriminator file every one of these under REFUSED and still +// pass. +func seedPoisonedDeliveriesInBulk(t *testing.T, database *sql.DB, n int) { + t.Helper() + ctx := context.Background() + activityID := "'https://coves.social/ap/activity/dv-bulk-unknown-' || i" + + _, err := database.ExecContext(ctx, ` + INSERT INTO outbound_activities (activity_id, actor_did, kind, payload) + SELECT `+activityID+`, $1, 'Create', '{"type":"Create"}'::jsonb + FROM generate_series(1, $2) AS i`, dvAuthorDID, n) + require.NoError(t, err, "seed %d activities", n) + + _, err = database.ExecContext(ctx, ` + INSERT INTO outbound_deliveries (activity_id, target_inbox, ordering_key, state, + last_error_class, last_status_code, created_at) + SELECT `+activityID+`, $1, $2, 'poisoned', 'transport', 0, now() - interval '2 hours' + FROM generate_series(1, $3) AS i`, dvUnknownInboxA, dvCommunityAPID, n) + require.NoError(t, err, "seed %d deliveries", n) +} diff --git a/internal/votes/divergence_recast_test.go b/internal/votes/divergence_recast_test.go index 7cd88f3..969f8ad 100644 --- a/internal/votes/divergence_recast_test.go +++ b/internal/votes/divergence_recast_test.go @@ -2,6 +2,7 @@ package votes import ( "context" + "database/sql" "fmt" "testing" @@ -222,3 +223,145 @@ func TestRecastDivergence_ADeliveredThenUndoneVoteIsNotADivergence(t *testing.T) "vote with no live row' reports this pair forever, about a withdrawal that worked") assert.Empty(t, recastEntries(sweep(t, w))) } + +// --------------------------------------------------------------------------- +// 17e review — the case that makes this class usable at all +// --------------------------------------------------------------------------- + +// TestRecastDivergence_ASuccessfulRecastIsNotADivergence is the ORDINARY vote +// flip, and it is most of what this table does. +// +// A flip is an in-place upsert: current_activity_id moves to the new id, +// delivered_state resets to pending, and NO Undo is enqueued — Lemmy takes a +// bare opposite vote as a replacement (17b measured this; a flip is not an Undo +// followed by a vote). So once the new vote delivers, the OLD delivered activity +// matches NEITHER exclusion: the ledger names the new id, and no Undo exists to +// pair with it. +// +// Left unpinned, every vote change the bridge has ever federated becomes a +// permanent entry, and the report grows monotonically with ordinary use. The +// detail on each would claim the peer holds a vote whose delivery never landed — +// both halves false — and an operator acting on it would retract the vote the +// user currently holds. +// +// Driven through the real path for the same reason the poisoned case is: the +// state that matters here is produced by the SECOND step of a two-step history, +// and no hand-written row has a history. +func TestRecastDivergence_ASuccessfulRecastIsNotADivergence(t *testing.T) { + w := newRecastWorld(t, 5) + + // Down, delivered. + w.castVote(t, "3lztprev00001", directionDown) + w.deliver(t) + require.Equal(t, []string{"Dislike"}, w.sender.kinds()) + require.Equal(t, string(store.DeliveredStateDelivered), w.state(t)) + + // Flipped to up — and this time it LANDS. + w.castVote(t, "3lztprev00002", directionUp) + w.deliver(t) + require.Equal(t, []string{"Dislike", "Like"}, w.sender.kinds(), + "precondition: the flip really went out as a bare opposite vote — Lemmy's replacement "+ + "semantics, with no Undo between them") + require.Equal(t, string(store.DeliveredStateDelivered), w.state(t), + "precondition: and the ledger caught up, so our accounting names what the peer holds") + + found, err := store.NewDivergences(w.db).RecastDivergences(context.Background()) + require.NoError(t, err) + assert.Empty(t, found, + "a flip that DELIVERED is not a divergence — it is the system working. The old activity "+ + "stays in the append-only history forever and no Undo will ever join it, because a "+ + "flip does not produce one; only the ledger naming the NEW id says the account is "+ + "settled. Reported, this class would grow with every vote change the bridge has "+ + "ever federated, each entry claiming the peer holds a vote that never landed — "+ + "both halves false — and acting on one would retract the vote the user has right now") + assert.Empty(t, recastEntries(sweep(t, w))) +} + +// TestRecastDivergence_AnUndoOlderThanTheDeliveredVoteDoesNotSuppressIt pins the +// TIME direction of the Undo exclusion. +// +// The exclusion asks whether a delivered Undo followed this vote. Drop the +// `>=` and any Undo for the pair suppresses the finding — including one from an +// earlier incarnation of the same (actor, subject), which the append-only +// history keeps forever. The sequence below is reachable and its consequence is +// permanent: the stale Undo from the FIRST vote silently hides a real divergence +// created two votes later, and nothing else in the system reports it. +// +// The second incarnation is a NEW vote record with its own rkey, which is what a +// client that deleted a vote and voted again produces. (Re-using the first +// record's rkey does not work: the activity id derives from the at-uri and a seq +// that restarts at 1 once the row is gone, so the re-created vote reproduces the +// FIRST vote's activity id, finds the already-delivered delivery standing, and +// never goes out. Recorded for the conductor — it is outside 17e.) +func TestRecastDivergence_AnUndoOlderThanTheDeliveredVoteDoesNotSuppressIt(t *testing.T) { + w := newRecastWorld(t, 1) // the last delivery must poison + ctx := context.Background() + + // 1) A vote, delivered. 2) Withdrawn, and the Undo delivers — so the + // history now holds a delivered Undo for this pair, forever. + w.castVote(t, "3lztprev00001", directionDown) + w.deliver(t) + w.deleteVote(t, "3lztprev00002") + w.deliver(t) + require.Equal(t, []string{"Dislike", "Undo"}, w.sender.kinds()) + require.Equal(t, "", w.state(t), "precondition: the withdrawal completed and cleared the row") + + // 3) The user votes again — a new record — and it lands. This is the vote + // the peer holds from here on. + const secondRKey = "3lztemporalvot2" + castVoteRecord(t, w, secondRKey, "3lztprev00003", directionUp) + w.deliver(t) + require.Equal(t, []string{"Dislike", "Undo", "Like"}, w.sender.kinds(), + "precondition: the second vote really went out") + require.Equal(t, string(store.DeliveredStateDelivered), voteRecordState(t, w, secondRKey)) + held := deliveredVoteActivity(t, w) + + // 4) And they change it once more — this time the delivery poisons. + w.sender.fail(fmt.Errorf("lemmy is unreachable")) + castVoteRecord(t, w, secondRKey, "3lztprev00004", directionDown) + w.deliver(t) + var poisoned int + require.NoError(t, w.db.QueryRow( + `SELECT COUNT(*) FROM outbound_deliveries WHERE state = 'poisoned'`).Scan(&poisoned)) + require.Equal(t, 1, poisoned, "precondition: the last delivery poisoned") + + found, err := store.NewDivergences(w.db).RecastDivergences(ctx) + require.NoError(t, err) + require.Len(t, found, 1, + "the peer is holding the Like from step 3, and the Undo in this history is OLDER than "+ + "it — it withdrew a different incarnation of the same (actor, subject) pair. An "+ + "exclusion that accepted any Undo at all would let that stale row suppress this "+ + "finding forever: the history is append-only, so the old Undo never goes away, and "+ + "the divergence it hides is permanent and reported nowhere else") + assert.Equal(t, held, found[0].DeliveredActivityID, + "and the entry cites the vote the peer actually holds — the one delivered AFTER the "+ + "Undo, not the withdrawn one") +} + +// castVoteRecord drives a vote COMMIT for an arbitrary record key, so a history +// can contain more than one vote RECORD for the same subject — which is what a +// client that deleted a vote and voted again produces. +func castVoteRecord(t *testing.T, w *recastWorld, rkey, rev, direction string) { + t.Helper() + frame := fmt.Sprintf( + `{"did":%q,"time_us":9500,"kind":"commit","commit":{"rev":%q,"operation":"create",`+ + `"collection":"social.coves.feed.vote","rkey":%q,"cid":%q,`+ + `"record":{"$type":"social.coves.feed.vote","subject":{"uri":%q,"cid":%q},`+ + `"direction":%q,"createdAt":"2026-08-14T10:00:00.000Z"}}}`, + tpNativeDID, rev, rkey, testCID, w.subject, testCID, direction) + w.handle(t, frame) +} + +// voteRecordState reads one vote record's delivered_state, or "" when the row +// is gone. +func voteRecordState(t *testing.T, w *recastWorld, rkey string) string { + t.Helper() + var state string + err := w.db.QueryRow(`SELECT delivered_state FROM outbound_votes WHERE vote_at_uri = $1`, + "at://"+tpNativeDID+"/social.coves.feed.vote/"+rkey).Scan(&state) + if err == sql.ErrNoRows { + return "" + } + require.NoError(t, err) + return state +} diff --git a/internal/votes/votes_test.go b/internal/votes/votes_test.go index 6f19202..b2a0f6f 100644 --- a/internal/votes/votes_test.go +++ b/internal/votes/votes_test.go @@ -59,6 +59,13 @@ var voteTablesToTruncate = []string{ // Consumer state, for the tests that drive the real write path. "federation_prefs", "jetstream_record_revs", "jetstream_dead_letters", "consumer_cursors", "repo_state", + // The acceptance ledger. 17e's divergence sweep runs here (the re-cast + // tests drive the real reconciler), and it READS this table — so a row left + // by another package's run would surface as an undelivered-acceptance + // finding in a vote test. Harmless only by coincidence today: the query + // inner-joins through outbound_objects, which IS truncated. That coincidence + // is exactly what this list exists to stop depending on. + "admissions", } // truncateVoteTables clears everything in voteTablesToTruncate.