From 6cd0430e3ff4160970aa7a22f63b37f32fa26879 Mon Sep 17 00:00:00 2001 From: agenticode <16611333+agenticode@users.noreply.github.com> Date: Wed, 26 Aug 2026 19:27:53 +0900 Subject: [PATCH] feat(store): scope-keyed RDS checkpoints and a proposals bucket MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit WHATIF §5.3 proposals bucket plus GROWTH §6 scope-keyed RDS checkpoints. The existing SaveEvidenceCheckpoint precedent is not scope-keyed: colliding keys would serve one region's history confidently as another's. The new checkpoint keys carry scope, and absent is kept distinct from present-but-unreadable so a read failure can never collapse into "no history". A forged blob's approved flag is not treated as an approval. Buckets bucketProposals and bucketRDSCheckpoints are created in Open alongside the existing five. Co-authored-by: kording <74226694+kording@users.noreply.github.com> --- pkg/store/PERSIST-FINDINGS.md | 510 +++++++++++++++++++ pkg/store/checkpoint.go | 426 ++++++++++++++++ pkg/store/checkpoint_test.go | 893 ++++++++++++++++++++++++++++++++++ pkg/store/persist_ext_test.go | 527 ++++++++++++++++++++ pkg/store/proposals.go | 136 ++++++ pkg/store/rdscheckpoint.go | 133 +++++ pkg/store/store.go | 10 +- 7 files changed, 2632 insertions(+), 3 deletions(-) create mode 100644 pkg/store/PERSIST-FINDINGS.md create mode 100644 pkg/store/checkpoint.go create mode 100644 pkg/store/checkpoint_test.go create mode 100644 pkg/store/persist_ext_test.go create mode 100644 pkg/store/proposals.go create mode 100644 pkg/store/rdscheckpoint.go diff --git a/pkg/store/PERSIST-FINDINGS.md b/pkg/store/PERSIST-FINDINGS.md new file mode 100644 index 0000000..a0dfedc --- /dev/null +++ b/pkg/store/PERSIST-FINDINGS.md @@ -0,0 +1,510 @@ +# U1 — the write side three units stopped at + +Three agents reached the same wall in three packages and each wrote down that +`pkg/store` was out of scope: `pkg/api/WHATIFROUTES-FINDINGS.md` §5.3 (a brain +restart loses every filed proposal and its audit trail), +`pkg/rds/GROWTH-FINDINGS.md` §6 (`SaveRDSCheckpoint`/`LoadRDSCheckpoint` do not +exist), `cmd/WIRING-FINDINGS.md` §6.2/§6.3 (the same write side, named as the +shared prerequisite). Both buckets are built. + +``` +pkg/store/checkpoint.go the shared envelope: keying, framing, bounds, typed errors +pkg/store/proposals.go the `proposals` bucket (§5.3) +pkg/store/rdscheckpoint.go the `rds-checkpoints` bucket (§6) +pkg/store/checkpoint_test.go 30 test functions, package store +pkg/store/persist_ext_test.go 5 test functions, package store_test — real whatif.Store +pkg/store/store.go +2 buckets in Open's creation list, +1 comment +``` + +`gofmt -l ./pkg/store` empty, `go vet ./...`, `go build ./...`, +`go test -race -count=1 ./pkg/store/...` and `go test -race -short ./...` all +green. **`go.mod` and `go.sum` are byte-identical** — `shasum` before and after +matches, and the only imports added are stdlib (`compress/gzip`, `errors`, +`unicode`) plus, in the external test package alone, +`pkg/whatif` and `pkg/backtest`. Coverage for `pkg/store` is **90.0 %** (was +88.9 %). No existing test file was touched and no existing assertion changed. + +--- + +## 1. The key encoding, and why a collision has no representation + +The trap named in the brief is real and it is the whole design: **the evidence +precedent is not scope-keyed.** `SaveEvidenceCheckpoint` stores one blob under +a constant, so it has no key to get wrong. §6 needs one per scope, and two +scopes that collapse into one key serve one region's storage history as +another's. + +**The key is the scope's bytes and nothing else.** No escape, no prefix, no +hash, no composition. `TestScopeKeysAreTheScopeBytes` pins that over the raw +bbolt keys, so a later change that introduces an encoding fails here and has to +re-argue the property rather than inherit it. + +That is not the obvious choice — the two neighbouring buckets don't do it — and +the contrast is the argument. `bucketPlans` and `bucketSnapHistory` key on +`"/"`, a genuinely composite key, and they pay for it: +`parsePlanTime`/`parseSnapTime` exist because an EKS ARN contains `/`, and +scanning cluster `a/b` from prefix `a/` would otherwise sweep up a sibling's +rows. That cost is worth paying when a key really has two components. A +checkpoint key has one. Composing it with anything — a kind tag, a version, a +timestamp — would import the entire ambiguity class for no benefit: +`"prod"` + `"us-east-1"` would become indistinguishable from +`"prod/us-east-1"`. **A one-component key cannot be parsed wrong because it is +never parsed.** + +`TestScopesContainingTheSeparatorCannotCollide` files eight scopes chosen to +break a composite scheme — `prod`, `us-east-1`, `prod/us-east-1`, +`prod/us-east-1/extra`, `prod/`, `/us-east-1`, `//`, `prod//us-east-1` — and +asserts eight distinct keys and eight distinct payloads read back. +`TestUnicodeScopesAreDistinctKeys` does the same for `café` composed vs +decomposed (they are two scopes, and merging them would be exactly the silent +account merge this bucket refuses to have), `東京/prod` and `prod-🌍`. + +Injectivity of the identity map is not an argument that needs making, which is +the point of choosing it. What still needs making is the argument about what +gets in: + +**`validateScope` refuses, at the door, four kinds of scope** — the reasoning +is `validateClusterID`'s, one step further, because a scope is used two ways at +once (raw as the bbolt key, JSON-encoded inside the envelope): + +| refused | why | test | +|---|---|---| +| empty | bbolt has no empty key; the failure would otherwise surface as `ErrKeyRequired` from three frames down | `TestUnstorableScopesAreRefusedAtTheDoor/empty` | +| over `MaxScopeBytes` = 512 | a 10 KiB scope is a caller passing a document as an identifier | `/10_KiB`, `/one_over_the_cap` | +| not valid UTF-8 | `encoding/json` rewrites it to U+FFFD, so the envelope's own `scope` would stop matching the key it lives under — an unreadable write | `/not_utf-8` | +| control characters | this string is echoed into error messages, log lines and `bbolt` CLI dumps; a scope carrying `\n` forges a log line (`whatif.sanitizeNote`'s reason, not an encoding one) | `/newline`, `/NUL` | + +Each of those also asserts that nothing was written under any key. A scope of +exactly 512 bytes stores: the boundary belongs to the legal side. + +### 1.1 The second defence, which does not depend on the first + +The envelope repeats the scope, and the reader refuses a record whose stored +scope is not the one asked for. This is the identity check `LoadSnapshot` and +`LatestPlan` already make against `ClusterID`. For records this package wrote it +is unreachable. It fires for a record moved between keys by anything else — a +restore from a mis-transcribed backup, a hand-edit, a bug in a future migration +— and it is what turns "served under the wrong scope" from unlikely into +*impossible*: `TestARecordMovedToAnotherKeyIsRefused` copies `eu-west-1`'s +value under `us-east-1` and asserts the load fails naming `eu-west-1`, rather +than handing back a growth history for a database in another region. + +The envelope also carries its `kind`, so a `proposals` envelope filed in the +rds bucket is refused instead of being handed to `rds.Domain.Restore` +(`TestACheckpointFromTheOtherBucketIsRefused`). + +The first-order consequence of a collision is worth stating because it is the +reason this got two defences rather than one. GROWTH-FINDINGS §7 already +observes that a cross-scope restore "produces a decrease", and a decrease is +refused by `storage-growth-history-inconsistent`. So the *common* collision +self-reports. The dangerous one is the collision that looks consistent: two +same-named databases in two regions whose histories interleave into a plausible +upward slope. That produces a growth verdict, with a rate and a projection, for +a database that never grew — a wrong answer delivered confidently, which is the +failure mode this whole repo is organised against. + +--- + +## 2. The size ceiling, with its arithmetic, at write and at read + +Both caps are applied **twice**: on the caller's bytes before compressing, and +on the decompressed output at read. The read-side bound is not redundancy. A +stored value never passed through `Save` — it can arrive from a backup, an +edit, or a future bug — and a few hundred bytes of gzip in a bbolt value expand +without limit. A decompression bomb in a bbolt value is still a decompression +bomb, and this process is meant to run unattended for months. + +### 2.1 `MaxProposalsCheckpointBytes = 64 MiB` + +Chosen from what `pkg/whatif` *can* emit, not from what it emits today. Per +record, at every free-text field's limit: + +``` +evidence citations maxEvidenceIDs × maxEvidenceIDLen = 64 × 512 = 32,768 B +rationale 2,048 B +human notes on the longest path draft→gated→approved→applied + 2 × maxNote = 4,096 B +the approval's own note 2,048 B +structure: two full policies, gate verdict, delta, keys, indent ≈ 6,366 B + ------------ + ≈ 47,326 B ≈ 46 KiB +× maxRecords (1000) ≈ 45.1 MiB +``` + +`TestOneRecordIsWellUnderItsShareOfTheCap` builds that record for real and +measures **45,278 bytes** (the table's figure plus the one further note the +longest path adds); it asserts `perRecord × 1000 ≤ cap` so that a later change +to `pkg/whatif` that grows a record past its share fails *here*, in this +package, rather than at 3am in a brain whose proposals silently stopped +persisting. 64 MiB leaves ~40 % headroom and matches +`MaxEvidenceCheckpointBytes`, chosen for the constraint that binds here too: a +bbolt value is materialized whole on both sides of a write. + +One correction to the arithmetic anyone else would derive from reading +`pkg/whatif`: **`maxRationale = 4096` never binds.** `Spec.normalize` checks it +and then calls `sanitizeNote`, whose `maxNote = 2048` refuses first. The +effective rationale bound is 2,048 bytes. That is `pkg/whatif`'s business, not +a defect this unit may fix, but a cap derived from the larger constant would be +overstated by 2 KiB per record. + +**What 64 MiB does not cover, stated rather than left to be discovered.** +`encoding/json` escapes `<`, `>` and `&` to six bytes each, and +`sanitizeNote` permits all three. A store whose every citation, rationale and +note is `<` frames to roughly **240 MiB** and is refused at write. That is +deliberate and it is the evidence precedent's choice — refuse loudly and name +the cap, rather than write a value that cannot be read back within a sane +allocation. The failure direction is the safe one: the in-memory store keeps +working, nothing is corrupted, and a lost proposal is not an approved one. A +realistic record is ~6.5 KiB — the 6,366 B of structure above plus short +prose — so a realistic full store is ~6.5 MiB and this is ~10× production +headroom. + +### 2.2 `MaxRDSCheckpointBytes = 64 MiB`, `MaxRDSCheckpointScopes = 64` + +GROWTH-FINDINGS §7's own numbers, taken at its stated bound: + +``` +DefaultGrowthMaxObservations 768 × ~40 B = 30,720 B per instance +× 1,000 instances (§7's worst realistic account) ≈ 29.3 MiB of history +× 2 for the rest of each Target and repeated JSON keys ≈ 58.6 MiB + cap = 64 MiB +``` + +`[unverified: no 1,000-instance account has been measured; §7 flags the same +figure as unverified and this cap inherits that.]` The refusal names +`Config.Growth.MaxObservations` as the knob, which is §7's own answer. + +The scope count cap is the bound §6 did not ask for and needs. **A scope-keyed +bucket with no cap on scopes is not a ceiling, it is a growth curve**: a caller +deriving a scope from something accidentally per-run — a session id, a +hostname, a timestamp — writes a new key every tick until the disk fills. 64 +scopes × 64 MiB is a 4 GiB absolute ceiling, and at the ~10× compression this +framing gets `[unverified for RDS checkpoints; the ratio is history.go's, for +snapshot JSON]` a realistic full fleet is a few hundred MiB on disk. + +The cap falls on the **newcomer**: past 64 scopes a new scope is refused by +name while every existing scope keeps updating. A cap that stopped the fleet +already being observed correctly would turn one mis-derived scope into a total +outage of the growth finding. `DeleteRDSCheckpoint` is how a slot is reclaimed, +because otherwise the cap is a one-way door — a renamed account would hold its +slot forever. Both halves are in +`TestScopeCapRefusesNewScopesAndKeepsUpdatingOldOnes`. + +### 2.3 The third bound: the stored record + +`maxStored() = max + max/2 + 4096` bounds the *compressed* value, checked +before it is copied out of the mmap. gzip of incompressible input can exceed +its input and `encoding/json` base64s the result, so a legitimate record can +reach ~4/3 of the decoded cap; 3/2 plus the envelope's own bytes is that with +room. Without it, a hostile value is fully copied before anything looks at it +(`TestOversizeStoredRecordIsRefusedBeforeItIsCopied`). + +--- + +## 3. Missing, future-version and corrupt are three facts, and they arrive as three errors + +> "Nothing was found" and "something was found and I could not read it" must +> never collapse: the second is a bug report, the first is Tuesday. + +**Every load returns an error rather than a nil blob.** The caller's rule is one +line: `errors.Is(err, ErrNoCheckpoint)` means start empty; *any* other error +means stop and say so. `pkg/api` already writes exactly this switch for +`evidence.ErrNoCheckpoint` in `restoreEvidence`, so the shape is not new to the +consumer. + +| sentinel | the fact | the remedy it protects | test | +|---|---|---|---| +| `ErrNoCheckpoint` | the bucket exists, nothing is stored under this scope | start empty — this is a first boot, or a scope observed for the first time | `TestColdScopeIsNotAnEmptyBlob` | +| `ErrCheckpointBucket` | the bucket itself is gone. `Open` creates every bucket, so this is not a cold start: the file was truncated, hand-edited, or written by something that is not this program | do not wipe and reinitialise — there is other state in that file | `TestMissingBucketIsNotAColdStart` | +| `ErrFutureCheckpoint` | the envelope's version is above this build's | **leave it alone.** Discarding it destroys a newer build's state the moment someone rolls back — the opposite remedy to corruption, which is why it cannot share an error with it | `TestFutureVersionIsRefusedByNameAndIsNotCorruption` | +| `ErrCorruptCheckpoint` | a record is present and unreadable: not an envelope, wrong kind, wrong scope, unknown codec, broken or truncated gzip, version *below* this build's | a bug report | `TestCorruptRecordIsNotAColdStart`, `TestBrokenGzipIsCorruptNotEmpty`, `TestUnknownCheckpointCodecIsRefusedByName`, `TestEnvelopeWithNoBlobIsRefused` | +| `ErrCheckpointTooLarge` | over a cap, at write or at read | a configuration change in the producing package | §2's tests | +| `ErrInvalidScope` | §1's table | fix the caller | `TestUnstorableScopesAreRefusedAtTheDoor` | +| `ErrEmptyCheckpoint` | a zero-length blob at write | — | `TestEmptyCheckpointIsRefusedAtWrite` | + +Each test asserts the *negative* as well: a future envelope must not also match +`ErrCorruptCheckpoint`, a corrupt record must not match `ErrNoCheckpoint`, a +missing bucket must not collapse into a cold start. + +Two decisions inside that table are worth their own sentence: + +- **A version *below* this build's is corruption, not compatibility.** The + evidence precedent reads two codecs because it had a predecessor to be + compatible with. This framing has never written another layout, so there is + no older form to accept and guessing at one would be inventing it. The error + says so. +- **An envelope carrying no blob is refused at read**, not returned as an empty + checkpoint. Combined with the write-side refusal of an empty blob, "a + checkpoint that exists and is empty" has no representation at either end — + which matters because it would be a *fourth* meaning competing with the three + above. + +`CheckpointError` carries `Kind`, `Scope` and `Detail` alongside the sentinel, +so a caller logs the scope with `errors.As` instead of parsing it back out of a +message. + +--- + +## 4. Concurrency + +The precedent's contract, matched exactly: `Store` is "safe for concurrent +use", bbolt serializes writers and snapshots readers, and these methods hold no +state of their own — no new mutex, and none needed. The one read-modify-write +is the scope cap, and it runs **inside** the write transaction, where bbolt has +already serialized writers, so two concurrent saves cannot both observe the +last free slot and both take it. + +`TestConcurrentCheckpointAccess` runs eight goroutines × 25 iterations under +`-race`, writing four scopes plus one contended scope, reading it back, and +driving both buckets at once. The assertion that matters is not "no data race" +— it is that a reader never sees a blend of two writes: every read of the +contended scope must equal one of the eight known payloads exactly. + +Framing is deterministic (`TestFramingIsDeterministic`): gzip with default +settings and no header fields, so identical state produces identical stored +bytes. `encodeSnapshot` relies on the same property, and it is what makes a +housekeeping tick with nothing to report a no-op write rather than a diff +against itself. + +--- + +## 5. The exact call sequence, for each of the three waiting consumers + +### 5.1 `pkg/api` — proposals (§5.3) + +**Where the load goes.** `newWhatIfSurface` (`pkg/api/whatifroutes.go:134`) +does `props: whatif.NewStore()`. That becomes: + +```go +func newWhatIfSurface(b *Brain) (*whatIfSurface, error) { + props := whatif.NewStore() + if b.st != nil { + blob, err := b.st.LoadProposals() + switch { + case errors.Is(err, store.ErrNoCheckpoint): + // First boot. Tuesday. + case err != nil: + // A proposal plane that exists and cannot be read is a bug report, + // not an empty plane. Do not start over it. + return nil, fmt.Errorf("api: restore proposal store: %w", err) + default: + loaded, err := whatif.Load(blob) + if err != nil { + return nil, fmt.Errorf("api: restore proposal store: %w", err) + } + props = loaded + } + } + return &whatIfSurface{b: b, props: props, now: func() time.Time { return time.Now().UTC() }}, nil +} +``` + +**Where the tick hangs.** §5.3 says "the same housekeeping timer that would call +`store.Sweep(clock)`". Stated precisely, because it is not quite there yet: +**`pkg/api` has no housekeeping timer today** — `grep -n 'time.NewTicker' +pkg/api/*.go` is empty. What it has instead is a checkpoint *cadence*, and that +is the right hook, because it is already the place the brain decides it is time +to write learned state to disk: + +```go +// brain.go:271, inside Ingest — beside the existing recommender/evidence writes: +if b.st != nil && count%b.cfg.CheckpointEvery == 0 { + _ = b.st.SaveRecommenderState(snap.ClusterID, r.Checkpoint()) + b.saveEvidence() + b.saveProposals() // ← add +} + +// brain.go:563, inside Serve's shutdown block, beside b.saveEvidence(): +b.saveProposals() // ← add + +// substrate.go, beside saveEvidence, same error policy (logged, not returned: +// losing a checkpoint costs history, failing an ingest costs the observation): +func (b *Brain) saveProposals() { + if b.st == nil || b.whatif == nil { + return + } + // Sweep FIRST: snapshotting before sweeping persists a store one tick + // staler than the one in memory, so a lapsed approval would be written + // back as still-approved and only expire after the next restart+tick. + if _, err := b.whatif.props.Sweep(b.whatif.now); err != nil { + b.cfg.Logger.Error("sweep proposals", "err", err) + } + if err := b.st.SaveProposalsFrom(b.whatif.props); err != nil { + b.cfg.Logger.Error("persist proposals", "err", err) + } +} +``` + +`SaveProposalsFrom` takes `store.ProposalSnapshotter`, which `*whatif.Store` +satisfies structurally — `persist_ext_test.go` asserts that at compile time, so +a signature drift in `pkg/whatif` breaks *this* package's test build rather +than `pkg/api`'s. `TestHousekeepingTickShape` runs the sweep-then-save loop +above across a TTL boundary and asserts a tick with nothing to sweep does not +rewrite the snapshot, and that a tick that expires the live approval does reach +the store. + +**One blocker this unit found and cannot fix, because `pkg/api` is not its +scope.** `registerWhatIfRoutes` is `newWhatIfSurface(b).register(mux)` — it +drops the pointer on the floor. Nothing on `Brain` can reach the proposal +store, so `saveProposals` above has no `b.whatif` to read. The surface has to +be held on `Brain` and built once in `NewBrain`, not per `Handler()` call. +That also settles §5.3's open question: with a bucket behind it, two handlers +built from one brain must not be two independent stores that alternately +overwrite each other's snapshot. + +### 5.2 `pkg/rds` + `cmd/` — the growth checkpoint (§6) + +GROWTH-FINDINGS §6's snippet, with the two missing calls filled in and the cold +start handled the way §3 requires: + +```go +d, _ := krds.NewDomain(sc) + +blob, err := st.LoadRDSCheckpoint(scope) // ← the missing read side +switch { +case errors.Is(err, store.ErrNoCheckpoint): + // No history yet. Every instance reports GrowthNoHistory with its reason + // filled — §6's intended cold start. +case err != nil: + // A checkpoint exists and could not be read. Proceeding here silently + // restarts a 14-day observation window: the growth finding would simply + // never appear and no one would know why. + return fmt.Errorf("rds: restore growth history for %s: %w", scope, err) +default: + if err := d.Restore(blob); err != nil { + return fmt.Errorf("rds: restore growth history for %s: %w", scope, err) + } +} + +snap, err := collector.Collect(ctx) +rds.RecordStorageHistory(snap, d.StorageHistories(), sc.Growth) +_ = d.Observe(snap) + +blob, err = d.Checkpoint() +if err != nil { + return err +} +if err := st.SaveRDSCheckpoint(scope, blob); err != nil { // ← the missing write side + return err +} +``` + +`scope` is `cmd/`'s to derive and this package does not parse it — +`"/"` is the obvious value and §1 shows that a `/` in it is +inert. Two rules for whoever picks it: it must be **stable across runs** (a +scope derived from anything per-run spends a slot per tick and hits +`MaxRDSCheckpointScopes`), and it must be **as specific as the data** (one +scope per account *and* region; collapsing regions is precisely the merge §1 +exists to prevent, and this package cannot detect it because both halves would +be legitimate writes under the same key). + +`RDSCheckpointScopes()` lists what is stored, in key order, so an operator can +see what the file holds and which slots the cap is spending. + +### 5.3 `cmd/WIRING-FINDINGS.md` §6.2/§6.3 + +§6.2's blocker — a time-keyed snapshot bucket — was closed by an earlier unit +(`SaveSnapshotAt`/`Snapshots` in `history.go`). §6.3's remaining blocker is the +evidence substrate on `api.Brain`, and `SaveEvidenceCheckpoint` already exists. +Neither needs anything from this unit; §6.3's "that evidence store is the same +prerequisite §6.2 needs" is now satisfied on both sides, and what is left in +both is `pkg/api` work. + +--- + +## 6. Persisting a proposal is not approving one + +`pkg/store` does not import `pkg/whatif`, does not parse a record, and has no +method that transitions one. The bucket is a **byte pipe**, and that is a +security property, not laziness: `whatif.Record.UnmarshalJSON` is the second +lock on the door, recomputing the content fingerprint and re-verifying that an +approved record carries a live approval bound to that fingerprint and that +verdict by a human who is not the author. Splitting any of that across a +persistence layer would be the way to lose it. + +`TestAForgedApprovalDoesNotComeBackApproved` hand-forges four blobs from a real +snapshot and asserts a pair of properties for each: + +1. the pipe hands back **exactly** what it was given — no repair, no rewriting, + no laundering; and +2. what it was given **fails to load**. Not "loads without the approval", not + "loads as gated" — does not load. + +The four: a gated record relabelled `"state": "approved"`; an approved record +whose approval object is deleted while its state is kept; an approver rewritten +to be the author (self-approval by edit); a gate verdict flipped from `false` +to `true`. Each fails at a different one of `pkg/whatif`'s checks, which is the +useful part — the refusals are independent, so no single edit clears them. + +`TestProposalStoreRoundTripsBitExactWithEveryGateState` is the positive half: a +store holding one record in **every state a stored record can rest in** — +gate-rejected, human-rejected, gated, approved-and-live, expired, applied — +round-trips byte for byte, reloads, and comes back in exactly those six states. +A restart may not promote anything. (`StateDraft` is absent and cannot be +produced: `whatif` refuses a stored draft, because no record rests there.) + +`TestNothingHereApprovesOrApplies` states the prohibition over this package's +own AST: no exported or unexported function whose name contains `approve`, +`apply`, `actuate`, `execute`, `resize`, `reboot` or `modify`, and no import of +`pkg/ec2`, `pkg/rds` or `pkg/whatif`. A persistence layer is exactly where that +line erodes first — a `SaveApproval`, a `MarkApplied` convenience, an import of +the state machine so a record could be "repaired" on load — and the test fails +if one appears rather than the claim rotting into a comment. + +The residual risk, stated plainly: **the worst these buckets can do is lose a +checkpoint or refuse to load one.** They cannot manufacture an approval, and +nothing downstream of them makes one executable — `pkg/ec2/actuate*.go` and +`pkg/rds/actuate*.go` remain unreachable from the binary and this unit touched +neither. + +--- + +## 7. What I did NOT build + +- **No consumer wiring.** Not one line outside `pkg/store/`. §5's sequences are + written to be pasted, and `TestHousekeepingTickShape` executes the proposals + one so the documented sequence is one that compiles and runs — but + `pkg/api`, `pkg/rds` and `cmd/` are other people's, and the surface-on-`Brain` + blocker in §5.1 is `pkg/api`'s to fix. +- **No `proposals` scope key.** The bucket stores one value under a constant, + like `bucketEvidence`, because a `whatif.Store` is fleet-wide — its proposals + carry their own target cluster and `GET /api/v1/proposals?cluster=` filters + one store rather than selecting between several. A scope here would be a + parameter every caller had to invent a value for. If a later unit runs + per-cluster proposal stores, the scope-keyed machinery is already in + `checkpoint.go` and the change is five constants. +- **No second codec, and no migration path.** Only `gzip` is written and only + `gzip` is read. The evidence envelope reads a plain-JSON codec because it had + a predecessor; this one has none, and shipping an unused codec would be + speculative compatibility with a format that has never existed. The `codec` + and `version` fields exist so the *next* one is refused by name. +- **No retention or pruning.** Both buckets are last-write-wins, one value per + key. There is no history of checkpoints, so there is nothing to prune — + unlike `history.go`, whose whole subject is retention. The only unbounded + dimension is the scope count, and §2.2 bounds it. +- **No compaction, and no delete for proposals.** `DeleteRDSCheckpoint` exists + because `MaxRDSCheckpointScopes` would otherwise be a one-way door. Proposals + have a single key that is always overwritten, so a delete would only be a way + to throw away an audit trail. +- **No `context.Context` on these calls.** `EvidenceCheckpointStore` takes one + because `evidence.CheckpointStore` demands it, and its own doc comment says + what it is worth: a bbolt write is not cancellable once begun. Neither §5.3 + nor §6 named an interface requiring one, so adding a parameter that can only + be checked before the work starts would be a lie about what an interrupted + checkpoint leaves behind. If a consumer needs the seam, copy + `EvidenceCheckpointStore`. +- **The proposals cap is not enforced against `whatif`'s record count.** This + package cannot count records in an opaque blob, and parsing one to check + would undo §6. `whatif.Load` enforces `maxRecords` itself on the way back in; + the byte cap here is the independent, cheaper bound. +- **No fuzz targets.** `fuzz_test.go` fuzzes the plan-key encoding because that + encoding is where the ambiguity lives. The checkpoint key has no encoding to + fuzz — §1 is the reason — and the envelope decoder's inputs are covered by + the hand-built adversarial records in `checkpoint_test.go`. A fuzzer over + `decode` would be reasonable future work; it was not the best use of the + remaining budget against writing the forged-blob tests. +- **`[unverified]` claims carried forward**, both from GROWTH-FINDINGS §7 and + both marked at their use site: the ~30 KiB-per-instance checkpoint size (no + 1,000-instance account has been measured) and the ~10× compression ratio (the + measured figure is `history.go`'s, for snapshot JSON, not for RDS + checkpoints). diff --git a/pkg/store/checkpoint.go b/pkg/store/checkpoint.go new file mode 100644 index 0000000..fdeb2c3 --- /dev/null +++ b/pkg/store/checkpoint.go @@ -0,0 +1,426 @@ +package store + +import ( + "bytes" + "compress/gzip" + "encoding/json" + "errors" + "fmt" + "io" + "unicode" + "unicode/utf8" + + bolt "go.etcd.io/bbolt" +) + +// Scope-keyed checkpoint framing — the write side pkg/api's +// WHATIFROUTES-FINDINGS.md §5.3, pkg/rds's GROWTH-FINDINGS.md §6 and +// cmd/WIRING-FINDINGS.md §6.2/§6.3 all stopped at. See PERSIST-FINDINGS.md. +// +// evidence.go is the precedent and this file is its generalization in exactly +// one dimension: SaveEvidenceCheckpoint stores ONE blob for the whole brain, +// and pkg/rds needs one per scope. Everything else is deliberately identical — +// a JSON envelope carrying its own version independent of the payload's, gzip +// inside it, an unknown version refused by name rather than guessed at. +// +// # Why a scope key is the dangerous part +// +// A checkpoint served under the wrong scope is not a crash, it is a wrong +// answer delivered confidently: one region's allocated-storage history read as +// another's produces a growth verdict about a database that never grew. Two +// defences, layered, because the cheap one is not sufficient on its own: +// +// 1. The key is the scope's bytes and nothing else. There is no separator to +// be contained, no prefix to be shared, no composition to be ambiguous — +// unlike bucketPlans and bucketSnapHistory, whose keys are +// "/" and therefore need parsePlanTime/parseSnapTime +// to tell "a/b" from "a" plus "b". A one-component key cannot be parsed +// wrong because it is never parsed. +// 2. The envelope repeats the scope, and the reader refuses a record whose +// stored scope is not the one asked for. This is the identity check +// LoadSnapshot and LatestPlan already make against ClusterID, for the same +// reason: even granting a key collision that (1) says is impossible, the +// wrong history is refused rather than served. +// +// Scopes are validated at the door instead of being escaped, so that (1) can +// stay true without an encoding. validateScope is the gate. + +// checkpointCodec is the only codec this build writes or reads. +// +// It is "gzip" and not evidence.go's "gzip+json" because the framing here is +// over OPAQUE BYTES: pkg/store does not parse a proposal store or an RDS +// checkpoint, does not import the packages that produce them, and must not +// start — the payload carries its own version (whatif.storeVersion, +// rds.Domain's own) and its producer owns its schema. The field exists so that +// a future codec is refused by name rather than misread, which is the same +// service the version field performs. +const checkpointCodec = "gzip" + +// MaxScopeBytes bounds a scope key. The bound is not bbolt's (32 KiB): a scope +// is an identifier a human types into a config file and reads back out of an +// error message, and 512 bytes is already two orders of magnitude past +// "/". A caller passing a document as a scope has a bug, and +// the earlier it is named the cheaper it is. +const MaxScopeBytes = 512 + +// The three facts a cold read can report, kept apart on purpose. +// +// A restart that lost data must not be indistinguishable from one that never +// had any. "Nothing was found" is Tuesday — the first boot of a new brain, a +// scope observed for the first time. "Something was found and I could not read +// it" is a bug report, and a caller that starts empty on it has silently +// discarded an audit trail. So the load path returns an error in every case, +// and the caller's rule is one line: errors.Is(err, ErrNoCheckpoint) means +// start empty, ANY other error means stop and say so. +var ( + // ErrNoCheckpoint: the bucket exists and holds nothing under this scope. + ErrNoCheckpoint = errors.New("no checkpoint is stored") + // ErrCheckpointBucket: the bucket itself is absent. Open creates every + // bucket, so this is not a cold start — it is a file that was truncated, + // hand-edited, or written by something that is not this program. + ErrCheckpointBucket = errors.New("checkpoint bucket is missing from the store file") + // ErrFutureCheckpoint: the envelope carries a version this build does not + // speak. Distinct from corruption because the remedy is opposite: a + // corrupt record may be discarded, a future one must be left alone, since + // discarding it destroys a newer build's state on a downgrade. + ErrFutureCheckpoint = errors.New("checkpoint was written by a build this one does not understand") + // ErrCorruptCheckpoint: a record is present and unreadable — not an + // envelope, wrong kind, wrong scope, unknown codec, broken gzip. + ErrCorruptCheckpoint = errors.New("stored checkpoint cannot be read") + // ErrCheckpointTooLarge: over the size cap, at write or at read. Separate + // from ErrCorruptCheckpoint because the bytes may be perfectly well formed + // and simply too many, and because the remedy is a configuration change in + // the producing package rather than a repair. + ErrCheckpointTooLarge = errors.New("checkpoint is over this build's size cap") + // ErrInvalidScope: a scope that cannot be stored as a key. + ErrInvalidScope = errors.New("checkpoint scope is not storable") + // ErrEmptyCheckpoint: a zero-length blob. Refused at write because it + // would otherwise read back as a checkpoint carrying nothing, which is a + // third meaning for a state that already has two. + ErrEmptyCheckpoint = errors.New("checkpoint is empty") +) + +// CheckpointError carries which checkpoint failed and why. Callers match the +// cause with errors.Is and, when they want to log the scope rather than parse +// it out of a string, reach the fields with errors.As. +type CheckpointError struct { + // Kind is the checkpoint family: "proposals" or "rds". + Kind string + // Scope is the key that was asked for. Empty for a single-scope kind. + Scope string + // Detail names the specific failure, e.g. both version numbers. + Detail string + // Err is one of the sentinels above. + Err error +} + +func (e *CheckpointError) Error() string { + var b bytes.Buffer + fmt.Fprintf(&b, "store: %s checkpoint", e.Kind) + if e.Scope != "" { + fmt.Fprintf(&b, " %q", e.Scope) + } + fmt.Fprintf(&b, ": %v", e.Err) + if e.Detail != "" { + fmt.Fprintf(&b, ": %s", e.Detail) + } + return b.String() +} + +func (e *CheckpointError) Unwrap() error { return e.Err } + +// checkpointEnvelope frames one stored checkpoint. +// +// Scope and Kind are stored as well as implied by the key and the bucket. That +// is deliberate redundancy: a record can only be misfiled once, and a reader +// that checks what it was handed against what it asked for turns "misfiled" +// from a wrong answer into a refusal. +type checkpointEnvelope struct { + Version int `json:"version"` + Kind string `json:"kind"` + Scope string `json:"scope"` + Codec string `json:"codec"` + // Blob is gzip of the caller's bytes, which encoding/json renders as + // base64. + Blob []byte `json:"blob"` +} + +// checkpointKind is one bucket's contract: its name, its framing version, its +// size ceiling, and how many scopes it will hold. Both buckets in this package +// share every line of the code below; they differ only in these five values. +type checkpointKind struct { + name string + bucket []byte + version int + // max bounds the DECODED blob, at write and again at read. + max int + // maxScopes bounds distinct keys. Zero means the kind is single-scope and + // scopeless: fixedScope is the only key it will ever use. + maxScopes int + // fixedScope is the key a single-scope kind stores under. + fixedScope string +} + +func (k checkpointKind) errf(scope string, sentinel error, format string, args ...any) error { + return &CheckpointError{Kind: k.name, Scope: scope, Detail: fmt.Sprintf(format, args...), Err: sentinel} +} + +func (k checkpointKind) err(scope string, sentinel error) error { + return &CheckpointError{Kind: k.name, Scope: scope, Err: sentinel} +} + +// validateScope rejects a scope that cannot survive being a key. +// +// The reasoning is validateClusterID's, one step further. A scope is used two +// ways at once — raw as the bbolt key, and JSON-encoded inside the envelope — +// so a scope that does not survive encoding/json unchanged would write a +// record whose own Scope no longer matches the key it lives under, which the +// reader would then (correctly) refuse as corrupt. An unreadable write is +// worse than a refused one, so it is refused here, where the operator can +// still see which scope was wrong. +// +// Control characters are refused for whatif.sanitizeNote's reason rather than +// an encoding one: this string is echoed into error messages, log lines and +// `bbolt` CLI dumps, and a scope carrying a newline can forge a log line. +func validateScope(scope string) error { + if scope == "" { + return errors.New("a scope must not be empty; bbolt has no empty key") + } + if len(scope) > MaxScopeBytes { + return fmt.Errorf("scope is %d bytes, over the %d-byte cap", len(scope), MaxScopeBytes) + } + if !utf8.ValidString(scope) { + return errors.New("scope is not valid UTF-8; encoding/json would rewrite it to U+FFFD and the key and the record would stop agreeing") + } + for _, r := range scope { + if unicode.IsControl(r) { + return fmt.Errorf("scope contains the control character %q, which cannot be audited", r) + } + } + return nil +} + +// scopeOf resolves the key a call operates on: the caller's for a scope-keyed +// kind, the fixed one for a single-scope kind. +func (k checkpointKind) scopeOf(scope string) (string, error) { + if k.maxScopes == 0 { + return k.fixedScope, nil + } + if err := validateScope(scope); err != nil { + return "", k.errf(scope, ErrInvalidScope, "%s", err.Error()) + } + return scope, nil +} + +// encode frames a blob. gzip is written with default settings and no header +// fields, so the framing is a function of the input alone: the same state +// always produces the same bytes, and re-saving unchanged state is a no-op +// write rather than a diff. encodeSnapshot relies on the same property. +func (k checkpointKind) encode(scope string, blob []byte) ([]byte, error) { + var buf bytes.Buffer + zw := gzip.NewWriter(&buf) + if _, err := zw.Write(blob); err != nil { + return nil, err + } + if err := zw.Close(); err != nil { + return nil, err + } + return json.Marshal(checkpointEnvelope{ + Version: k.version, + Kind: k.name, + Scope: scope, + Codec: checkpointCodec, + Blob: buf.Bytes(), + }) +} + +// maxStored bounds the stored record. gzip of incompressible input can exceed +// it slightly and encoding/json base64s the result, so a legitimate record can +// be about 4/3 of the decoded cap; 3/2 plus the envelope's own bytes is that +// with room. The point is to refuse a hostile value BEFORE copying it out of +// the mmap, not to second-guess the compressor. +func (k checkpointKind) maxStored() int { return k.max + k.max/2 + 4096 } + +// save writes one checkpoint, replacing any previous one for the scope. +func (k checkpointKind) save(s *Store, scope string, blob []byte) error { + key, err := k.scopeOf(scope) + if err != nil { + return err + } + if len(blob) == 0 { + return k.err(key, ErrEmptyCheckpoint) + } + if len(blob) > k.max { + return k.errf(key, ErrCheckpointTooLarge, "%d bytes, over the %d-byte cap", len(blob), k.max) + } + raw, err := k.encode(key, blob) + if err != nil { + return k.errf(key, ErrCorruptCheckpoint, "framing: %v", err) + } + return s.db.Update(func(tx *bolt.Tx) error { + b := tx.Bucket(k.bucket) + if b == nil { + return k.err(key, ErrCheckpointBucket) + } + // The scope cap is enforced inside the write transaction, where bbolt + // has already serialized writers, so two concurrent saves cannot both + // observe the last free slot and both take it. + if k.maxScopes > 0 && b.Get([]byte(key)) == nil { + if n := countKeys(b); n >= k.maxScopes { + return k.errf(key, ErrCheckpointTooLarge, + "the bucket already holds %d scopes, the cap is %d; delete a scope that no longer exists", + n, k.maxScopes) + } + } + return b.Put([]byte(key), raw) + }) +} + +// load returns one checkpoint's bytes. Every outcome other than success is an +// error naming which of the failures it was; nothing returns a nil blob and a +// nil error, because that is the shape in which "lost" reads as "new". +func (k checkpointKind) load(s *Store, scope string) ([]byte, error) { + key, err := k.scopeOf(scope) + if err != nil { + return nil, err + } + var raw []byte + var loadErr error + if err := s.db.View(func(tx *bolt.Tx) error { + b := tx.Bucket(k.bucket) + if b == nil { + loadErr = k.err(key, ErrCheckpointBucket) + return nil + } + v := b.Get([]byte(key)) + if v == nil { + loadErr = k.err(key, ErrNoCheckpoint) + return nil + } + if len(v) > k.maxStored() { + loadErr = k.errf(key, ErrCheckpointTooLarge, + "stored record is %d bytes, over the %d-byte stored cap", len(v), k.maxStored()) + return nil + } + // bbolt values alias the mmap; copy before the transaction ends. + raw = append([]byte(nil), v...) + return nil + }); err != nil { + return nil, k.errf(key, ErrCorruptCheckpoint, "reading: %v", err) + } + if loadErr != nil { + return nil, loadErr + } + return k.decode(key, raw) +} + +func (k checkpointKind) decode(scope string, raw []byte) ([]byte, error) { + var env checkpointEnvelope + if err := json.Unmarshal(raw, &env); err != nil { + return nil, k.errf(scope, ErrCorruptCheckpoint, "stored record is not a checkpoint envelope: %v", err) + } + // Version first: a future envelope may legitimately carry fields whose + // meaning this build would misjudge, so nothing below is trusted until the + // layout is one this build wrote. + if env.Version != k.version { + if env.Version > k.version { + return nil, k.errf(scope, ErrFutureCheckpoint, + "envelope version %d, this build speaks %d; refusing to guess at a layout it does not know", + env.Version, k.version) + } + return nil, k.errf(scope, ErrCorruptCheckpoint, + "envelope version %d, this build speaks %d and has never written another", + env.Version, k.version) + } + if env.Kind != k.name { + return nil, k.errf(scope, ErrCorruptCheckpoint, + "envelope holds a %q checkpoint, not a %q one", env.Kind, k.name) + } + // The second collision defence. The key encoding makes this unreachable + // for records this package wrote; it fires for a record moved between keys + // by anything else, and it is what makes "served under the wrong scope" + // structurally impossible rather than merely unlikely. + if env.Scope != scope { + return nil, k.errf(scope, ErrCorruptCheckpoint, + "envelope holds scope %q; refusing to serve it under %q", env.Scope, scope) + } + if env.Codec != checkpointCodec { + return nil, k.errf(scope, ErrCorruptCheckpoint, "codec %q is not supported by this build", env.Codec) + } + if len(env.Blob) == 0 { + return nil, k.errf(scope, ErrCorruptCheckpoint, "envelope carries no blob") + } + zr, err := gzip.NewReader(bytes.NewReader(env.Blob)) + if err != nil { + return nil, k.errf(scope, ErrCorruptCheckpoint, "gzip: %v", err) + } + defer zr.Close() + // The cap is applied again HERE, not only at write. A stored value is a + // few hundred kilobytes of base64 that can inflate without limit, and this + // process is meant to run unattended for months; a decompression bomb in a + // bbolt value is still a decompression bomb. +1 so an input of exactly the + // cap is distinguishable from one over it. + out, err := io.ReadAll(io.LimitReader(zr, int64(k.max)+1)) + if err != nil { + return nil, k.errf(scope, ErrCorruptCheckpoint, "gzip: %v", err) + } + if len(out) > k.max { + return nil, k.errf(scope, ErrCheckpointTooLarge, + "checkpoint decompresses past the %d-byte cap", k.max) + } + return out, nil +} + +// remove deletes one scope's checkpoint. Deleting an absent scope is +// ErrNoCheckpoint, not success: a caller freeing a slot needs to know whether +// it freed one. +func (k checkpointKind) remove(s *Store, scope string) error { + key, err := k.scopeOf(scope) + if err != nil { + return err + } + var missing error + if err := s.db.Update(func(tx *bolt.Tx) error { + b := tx.Bucket(k.bucket) + if b == nil { + return k.err(key, ErrCheckpointBucket) + } + if b.Get([]byte(key)) == nil { + missing = k.err(key, ErrNoCheckpoint) + return nil + } + return b.Delete([]byte(key)) + }); err != nil { + return err + } + return missing +} + +// scopes lists stored scopes in ascending key order. The keys ARE the scopes — +// see the file comment — so no decoding is involved and a scope that cannot be +// listed is a scope that was never written. +func (k checkpointKind) scopes(s *Store) ([]string, error) { + var out []string + if err := s.db.View(func(tx *bolt.Tx) error { + b := tx.Bucket(k.bucket) + if b == nil { + return k.err("", ErrCheckpointBucket) + } + return b.ForEach(func(key, _ []byte) error { + out = append(out, string(key)) + return nil + }) + }); err != nil { + return nil, err + } + return out, nil +} + +func countKeys(b *bolt.Bucket) int { + n := 0 + c := b.Cursor() + for k, _ := c.First(); k != nil; k, _ = c.Next() { + n++ + } + return n +} diff --git a/pkg/store/checkpoint_test.go b/pkg/store/checkpoint_test.go new file mode 100644 index 0000000..0af2ad0 --- /dev/null +++ b/pkg/store/checkpoint_test.go @@ -0,0 +1,893 @@ +package store + +import ( + "bytes" + "compress/gzip" + "encoding/json" + "errors" + "fmt" + "go/ast" + "go/parser" + "go/token" + "os" + "strings" + "sync" + "testing" + + bolt "go.etcd.io/bbolt" +) + +// A checkpoint blob is opaque to pkg/store, so the tests in this file frame +// stand-in documents: what is under test is the FRAMING and the KEYING, not +// any producer's schema. persist_ext_test.go does the other half with a real +// whatif.Store. + +// smallKind is a checkpoint family with a tiny cap, so the bound-at-read tests +// assert the same code path the 64 MiB kinds use without allocating 64 MiB to +// do it. The two real kinds are exercised against their real caps in +// TestRealKindsRefuseOversizeAtWrite. +func smallKind(t *testing.T, s *Store, name string, max, maxScopes int) checkpointKind { + t.Helper() + k := checkpointKind{ + name: name, + bucket: []byte("test-" + name), + version: 1, + max: max, + maxScopes: maxScopes, + } + if maxScopes == 0 { + k.fixedScope = "brain" + } + if err := s.db.Update(func(tx *bolt.Tx) error { + _, err := tx.CreateBucketIfNotExists(k.bucket) + return err + }); err != nil { + t.Fatal(err) + } + return k +} + +// putCheckpointRaw stores bytes this package would not itself write, so the +// reader can be tested against records other builds — or an attacker — produce. +func putCheckpointRaw(t *testing.T, s *Store, k checkpointKind, key string, raw []byte) { + t.Helper() + if err := s.db.Update(func(tx *bolt.Tx) error { + return tx.Bucket(k.bucket).Put([]byte(key), raw) + }); err != nil { + t.Fatal(err) + } +} + +func putEnvelope(t *testing.T, s *Store, k checkpointKind, key string, env checkpointEnvelope) { + t.Helper() + raw, err := json.Marshal(env) + if err != nil { + t.Fatal(err) + } + putCheckpointRaw(t, s, k, key, raw) +} + +// gzipOf frames a payload the way encode does, for hand-built envelopes. +func gzipOf(t *testing.T, payload []byte) []byte { + t.Helper() + var buf bytes.Buffer + zw := gzip.NewWriter(&buf) + if _, err := zw.Write(payload); err != nil { + t.Fatal(err) + } + if err := zw.Close(); err != nil { + t.Fatal(err) + } + return buf.Bytes() +} + +// storedKeys returns the raw bbolt keys of a bucket. The collision tests +// assert over these rather than over the API, because the claim being made is +// about the key encoding, not about what happens to survive a round trip. +func storedKeys(t *testing.T, s *Store, k checkpointKind) []string { + t.Helper() + var out []string + if err := s.db.View(func(tx *bolt.Tx) error { + return tx.Bucket(k.bucket).ForEach(func(key, _ []byte) error { + out = append(out, string(key)) + return nil + }) + }); err != nil { + t.Fatal(err) + } + return out +} + +// rdsBlob builds a stand-in RDS checkpoint with an irregular history: uneven +// gaps, a ratchet step, and a scope-specific marker so a cross-scope leak is +// visible in the bytes rather than only in a length. +func rdsBlob(scope string, days []int) []byte { + var b bytes.Buffer + fmt.Fprintf(&b, `{"scope":%q,"targets":[{"id":"db-1","history":[`, scope) + for i, d := range days { + if i > 0 { + b.WriteByte(',') + } + fmt.Fprintf(&b, `{"day":%d,"allocatedGiB":%d}`, d, 100+10*i) + } + b.WriteString(`]}]}`) + return b.Bytes() +} + +// ---- the three facts a load can report ---- + +// TestColdScopeIsNotAnEmptyBlob prevents the failure the whole error taxonomy +// exists for: a caller that cannot tell "never had a checkpoint" from "had one +// and could not read it" silently restarts an observation window over an audit +// trail it just failed to read. +func TestColdScopeIsNotAnEmptyBlob(t *testing.T) { + s := open(t) + blob, err := s.LoadRDSCheckpoint("prod/us-east-1") + if blob != nil { + t.Fatalf("a cold scope returned %q as well as an error", blob) + } + if !errors.Is(err, ErrNoCheckpoint) { + t.Fatalf("a cold scope reported %v, want ErrNoCheckpoint", err) + } + // And it is not any of the other three facts. + for _, other := range []error{ErrCorruptCheckpoint, ErrFutureCheckpoint, ErrCheckpointBucket} { + if errors.Is(err, other) { + t.Fatalf("a cold scope also matched %v", other) + } + } + var ce *CheckpointError + if !errors.As(err, &ce) || ce.Scope != "prod/us-east-1" || ce.Kind != "rds" { + t.Fatalf("error does not carry kind and scope: %#v", ce) + } +} + +// TestMissingBucketIsNotAColdStart: a truncated or foreign file must not read +// as a brain that has never checkpointed. Open creates every bucket, so a +// bucket that is gone is a statement about the FILE, and wiping and starting +// over on it destroys whatever else is in there. +func TestMissingBucketIsNotAColdStart(t *testing.T) { + s := open(t) + if err := s.SaveRDSCheckpoint("prod", rdsBlob("prod", []int{0, 3})); err != nil { + t.Fatal(err) + } + if err := s.db.Update(func(tx *bolt.Tx) error { + return tx.DeleteBucket(bucketRDSCheckpoints) + }); err != nil { + t.Fatal(err) + } + _, err := s.LoadRDSCheckpoint("prod") + if !errors.Is(err, ErrCheckpointBucket) { + t.Fatalf("a missing bucket reported %v, want ErrCheckpointBucket", err) + } + if errors.Is(err, ErrNoCheckpoint) { + t.Fatal("a missing bucket collapsed into a cold start") + } + if _, err := s.LoadProposals(); !errors.Is(err, ErrNoCheckpoint) { + t.Fatalf("deleting one bucket changed another's answer: %v", err) + } +} + +// TestFutureVersionIsRefusedByNameAndIsNotCorruption: the remedy for the two +// is opposite. A corrupt record may be discarded; a record written by a newer +// build must be left alone, because discarding it destroys that build's state +// the moment someone rolls back. +func TestFutureVersionIsRefusedByNameAndIsNotCorruption(t *testing.T) { + s := open(t) + k := smallKind(t, s, "future", 1024, 4) + putEnvelope(t, s, k, "prod", checkpointEnvelope{ + Version: k.version + 1, Kind: k.name, Scope: "prod", + Codec: checkpointCodec, Blob: gzipOf(t, []byte("payload")), + }) + _, err := k.load(s, "prod") + if !errors.Is(err, ErrFutureCheckpoint) { + t.Fatalf("a future envelope reported %v, want ErrFutureCheckpoint", err) + } + if errors.Is(err, ErrCorruptCheckpoint) || errors.Is(err, ErrNoCheckpoint) { + t.Fatalf("a future envelope also matched another fact: %v", err) + } + msg := err.Error() + if !strings.Contains(msg, "version 2") || !strings.Contains(msg, "speaks 1") { + t.Fatalf("error does not name both versions: %v", err) + } + // A version BELOW this build's is corruption, not the future: this build + // has never written another, so there is no older layout to be compatible + // with and guessing would be inventing one. + putEnvelope(t, s, k, "old", checkpointEnvelope{ + Version: 0, Kind: k.name, Scope: "old", + Codec: checkpointCodec, Blob: gzipOf(t, []byte("payload")), + }) + if _, err := k.load(s, "old"); !errors.Is(err, ErrCorruptCheckpoint) { + t.Fatalf("a version-0 envelope reported %v, want ErrCorruptCheckpoint", err) + } +} + +// TestCorruptRecordIsNotAColdStart — a bug report must not read as Tuesday. +func TestCorruptRecordIsNotAColdStart(t *testing.T) { + s := open(t) + k := smallKind(t, s, "corrupt", 1024, 4) + for _, tc := range []struct { + name string + raw []byte + }{ + {"not json", []byte("not json at all")}, + {"json but not an envelope", []byte(`{"hello":"world"}`)}, + {"truncated", []byte(`{"version":1,"kind":"corrupt",`)}, + {"empty value", []byte{}}, + } { + t.Run(tc.name, func(t *testing.T) { + putCheckpointRaw(t, s, k, "prod", tc.raw) + _, err := k.load(s, "prod") + if !errors.Is(err, ErrCorruptCheckpoint) { + t.Fatalf("reported %v, want ErrCorruptCheckpoint", err) + } + if errors.Is(err, ErrNoCheckpoint) { + t.Fatal("a corrupt record was reported as a cold start") + } + }) + } +} + +// TestBrokenGzipIsCorruptNotEmpty: the envelope parses, the payload does not. +func TestBrokenGzipIsCorruptNotEmpty(t *testing.T) { + s := open(t) + k := smallKind(t, s, "gz", 1024, 4) + putEnvelope(t, s, k, "prod", checkpointEnvelope{ + Version: 1, Kind: k.name, Scope: "prod", Codec: checkpointCodec, + Blob: []byte("this is not gzip"), + }) + if _, err := k.load(s, "prod"); !errors.Is(err, ErrCorruptCheckpoint) { + t.Fatalf("reported %v, want ErrCorruptCheckpoint", err) + } + // Half a gzip stream: the header is valid, the body is truncated. + full := gzipOf(t, bytes.Repeat([]byte("payload"), 40)) + putEnvelope(t, s, k, "half", checkpointEnvelope{ + Version: 1, Kind: k.name, Scope: "half", Codec: checkpointCodec, + Blob: full[:len(full)-6], + }) + if _, err := k.load(s, "half"); !errors.Is(err, ErrCorruptCheckpoint) { + t.Fatalf("a truncated gzip reported %v, want ErrCorruptCheckpoint", err) + } +} + +func TestUnknownCheckpointCodecIsRefusedByName(t *testing.T) { + s := open(t) + k := smallKind(t, s, "codec", 1024, 4) + putEnvelope(t, s, k, "prod", checkpointEnvelope{ + Version: 1, Kind: k.name, Scope: "prod", Codec: "zstd+cbor", + Blob: gzipOf(t, []byte("payload")), + }) + _, err := k.load(s, "prod") + if !errors.Is(err, ErrCorruptCheckpoint) || !strings.Contains(err.Error(), "zstd+cbor") { + t.Fatalf("an unknown codec was not refused by name: %v", err) + } +} + +// ---- the key encoding ---- + +// TestScopeKeysAreTheScopeBytes pins the encoding the collision argument rests +// on. If a future change introduces an encoding — an escape, a prefix, a +// hash — this fails, and whoever makes that change has to re-argue collision +// freedom rather than inherit it. +func TestScopeKeysAreTheScopeBytes(t *testing.T) { + s := open(t) + scopes := []string{"prod", "prod/us-east-1", "café", "a b", "9"} + for _, sc := range scopes { + if err := s.SaveRDSCheckpoint(sc, rdsBlob(sc, []int{0, 1})); err != nil { + t.Fatal(err) + } + } + got := storedKeys(t, s, rdsKind) + if len(got) != len(scopes) { + t.Fatalf("%d scopes wrote %d keys: %q", len(scopes), len(got), got) + } + for _, sc := range scopes { + found := false + for _, k := range got { + if k == sc { + found = true + } + } + if !found { + t.Fatalf("scope %q is not stored under its own bytes; keys are %q", sc, got) + } + } +} + +// TestScopesContainingTheSeparatorCannotCollide is the failure GROWTH-FINDINGS +// §7 names: two scopes that collapse into one key serve one region's storage +// history as another's, and the resulting growth verdict is about a database +// that never grew. A composite key — the "/" shape +// bucketPlans and bucketSnapHistory use — would make "prod" + "us-east-1" +// indistinguishable from "prod/us-east-1". This bucket's key has one +// component, so there is nothing to compose and nothing to parse back. +func TestScopesContainingTheSeparatorCannotCollide(t *testing.T) { + s := open(t) + scopes := []string{ + "prod", "us-east-1", "prod/us-east-1", "prod/us-east-1/extra", + "prod/", "/us-east-1", "//", "prod//us-east-1", + } + for i, sc := range scopes { + if err := s.SaveRDSCheckpoint(sc, rdsBlob(sc, []int{0, i + 1, i + 9})); err != nil { + t.Fatalf("scope %q: %v", sc, err) + } + } + if got := len(storedKeys(t, s, rdsKind)); got != len(scopes) { + t.Fatalf("%d scopes collapsed into %d keys", len(scopes), got) + } + for i, sc := range scopes { + got, err := s.LoadRDSCheckpoint(sc) + if err != nil { + t.Fatalf("scope %q: %v", sc, err) + } + want := rdsBlob(sc, []int{0, i + 1, i + 9}) + if !bytes.Equal(got, want) { + t.Fatalf("scope %q served\n %s\nwant\n %s", sc, got, want) + } + } +} + +// TestUnicodeScopesAreDistinctKeys: two spellings of the same grapheme are two +// scopes. Normalising them together would be a silent merge of two accounts' +// histories, which is the collision this bucket refuses to have; normalising +// them apart is what storing the bytes already does. +func TestUnicodeScopesAreDistinctKeys(t *testing.T) { + s := open(t) + // The same grapheme, composed and decomposed: two different byte strings. + nfc, nfd := "caf\u00e9", "cafe\u0301" + scopes := []string{nfc, nfd, "東京/prod", "prod-🌍"} + for i, sc := range scopes { + if err := s.SaveRDSCheckpoint(sc, rdsBlob(sc, []int{i, i + 4})); err != nil { + t.Fatalf("scope %q: %v", sc, err) + } + } + if got := len(storedKeys(t, s, rdsKind)); got != len(scopes) { + t.Fatalf("%d unicode scopes collapsed into %d keys", len(scopes), got) + } + for i, sc := range scopes { + got, err := s.LoadRDSCheckpoint(sc) + if err != nil { + t.Fatalf("scope %q: %v", sc, err) + } + if want := rdsBlob(sc, []int{i, i + 4}); !bytes.Equal(got, want) { + t.Fatalf("scope %q served another scope's checkpoint: %s", sc, got) + } + } +} + +func TestUnstorableScopesAreRefusedAtTheDoor(t *testing.T) { + s := open(t) + for _, tc := range []struct{ name, scope string }{ + {"empty", ""}, + {"10 KiB", strings.Repeat("s", 10<<10)}, + {"one over the cap", strings.Repeat("s", MaxScopeBytes+1)}, + {"newline", "prod\nus-east-1"}, + {"NUL", "prod\x00"}, + {"not utf-8", string([]byte{0xff, 0xfe, 0x00})}, + } { + t.Run(tc.name, func(t *testing.T) { + err := s.SaveRDSCheckpoint(tc.scope, []byte(`{"x":1}`)) + if !errors.Is(err, ErrInvalidScope) { + t.Fatalf("save reported %v, want ErrInvalidScope", err) + } + if _, err := s.LoadRDSCheckpoint(tc.scope); !errors.Is(err, ErrInvalidScope) { + t.Fatalf("load reported %v, want ErrInvalidScope", err) + } + // Nothing escaped into the bucket, under any key. + if keys := storedKeys(t, s, rdsKind); len(keys) != 0 { + t.Fatalf("a refused scope wrote keys %q", keys) + } + }) + } + // Exactly at the cap is storable: the boundary belongs to the legal side. + atCap := strings.Repeat("s", MaxScopeBytes) + if err := s.SaveRDSCheckpoint(atCap, []byte(`{"x":1}`)); err != nil { + t.Fatalf("a scope of exactly %d bytes was refused: %v", MaxScopeBytes, err) + } +} + +// TestARecordMovedToAnotherKeyIsRefused is the second collision defence, the +// one that does not depend on the key encoding being right. It fires when +// something other than this package files a record under the wrong scope. +func TestARecordMovedToAnotherKeyIsRefused(t *testing.T) { + s := open(t) + if err := s.SaveRDSCheckpoint("eu-west-1", rdsBlob("eu-west-1", []int{0, 2, 30})); err != nil { + t.Fatal(err) + } + var raw []byte + if err := s.db.View(func(tx *bolt.Tx) error { + raw = append(raw, tx.Bucket(bucketRDSCheckpoints).Get([]byte("eu-west-1"))...) + return nil + }); err != nil { + t.Fatal(err) + } + putCheckpointRaw(t, s, rdsKind, "us-east-1", raw) + _, err := s.LoadRDSCheckpoint("us-east-1") + if !errors.Is(err, ErrCorruptCheckpoint) { + t.Fatalf("one region's history was served as another's: %v", err) + } + if !strings.Contains(err.Error(), "eu-west-1") { + t.Fatalf("refusal does not name the scope it actually holds: %v", err) + } +} + +// TestACheckpointFromTheOtherBucketIsRefused: a proposals envelope filed in +// the rds bucket is not an rds checkpoint, and handing its bytes to +// rds.Domain.Restore would be a decode error at best. +func TestACheckpointFromTheOtherBucketIsRefused(t *testing.T) { + s := open(t) + if err := s.SaveProposals([]byte(`{"version":1,"records":[]}`)); err != nil { + t.Fatal(err) + } + var raw []byte + if err := s.db.View(func(tx *bolt.Tx) error { + raw = append(raw, tx.Bucket(bucketProposals).Get([]byte(proposalsScope))...) + return nil + }); err != nil { + t.Fatal(err) + } + putCheckpointRaw(t, s, rdsKind, "brain", raw) + _, err := s.LoadRDSCheckpoint("brain") + if !errors.Is(err, ErrCorruptCheckpoint) || !strings.Contains(err.Error(), `"proposals"`) { + t.Fatalf("a proposals envelope was accepted as an rds one: %v", err) + } +} + +// ---- bounds ---- + +// TestRealKindsRefuseOversizeAtWrite asserts the boundary on the two shipped +// caps, one byte over each, and that a refused save leaves the previous +// checkpoint intact rather than half-writing one. +func TestRealKindsRefuseOversizeAtWrite(t *testing.T) { + s := open(t) + prev := []byte(`{"version":1,"records":[]}`) + if err := s.SaveProposals(prev); err != nil { + t.Fatal(err) + } + big := bytes.Repeat([]byte("x"), MaxProposalsCheckpointBytes+1) + err := s.SaveProposals(big) + if !errors.Is(err, ErrCheckpointTooLarge) { + t.Fatalf("an oversize proposal store was stored: %v", err) + } + if !strings.Contains(err.Error(), "over the") { + t.Fatalf("refusal does not name the cap: %v", err) + } + got, err := s.LoadProposals() + if err != nil || !bytes.Equal(got, prev) { + t.Fatalf("a refused save damaged the previous checkpoint: %q %v", got, err) + } + big = nil + + rdsBig := bytes.Repeat([]byte("x"), MaxRDSCheckpointBytes+1) + if err := s.SaveRDSCheckpoint("prod", rdsBig); !errors.Is(err, ErrCheckpointTooLarge) { + t.Fatalf("an oversize rds checkpoint was stored: %v", err) + } + if _, err := s.LoadRDSCheckpoint("prod"); !errors.Is(err, ErrNoCheckpoint) { + t.Fatalf("a refused save left something behind: %v", err) + } +} + +// TestDecompressionBombIsRefusedAtRead is the read-side half of the ceiling. +// The write side cannot help here: this record was never written through +// SaveRDSCheckpoint. A few hundred bytes of gzip in a bbolt value expand +// without limit, and this process is meant to run unattended for months. +func TestDecompressionBombIsRefusedAtRead(t *testing.T) { + s := open(t) + k := smallKind(t, s, "bomb", 4096, 4) + bomb := gzipOf(t, bytes.Repeat([]byte{0}, k.max*64)) + if len(bomb) > 2048 { + t.Fatalf("test bomb is not compressed (%d bytes), it proves nothing", len(bomb)) + } + putEnvelope(t, s, k, "prod", checkpointEnvelope{ + Version: 1, Kind: k.name, Scope: "prod", Codec: checkpointCodec, Blob: bomb, + }) + _, err := k.load(s, "prod") + if !errors.Is(err, ErrCheckpointTooLarge) { + t.Fatalf("a decompression bomb was expanded: %v", err) + } + if !strings.Contains(err.Error(), "decompresses past") { + t.Fatalf("refusal does not name what happened: %v", err) + } + // A payload of exactly the cap still loads: the bound is not off by one. + putEnvelope(t, s, k, "atcap", checkpointEnvelope{ + Version: 1, Kind: k.name, Scope: "atcap", Codec: checkpointCodec, + Blob: gzipOf(t, bytes.Repeat([]byte{'a'}, k.max)), + }) + got, err := k.load(s, "atcap") + if err != nil || len(got) != k.max { + t.Fatalf("a payload of exactly the cap was refused: %d bytes, %v", len(got), err) + } +} + +// TestOversizeStoredRecordIsRefusedBeforeItIsCopied bounds the compressed side +// too: without it, a hostile value is copied out of the mmap in full before +// anything looks at it. +func TestOversizeStoredRecordIsRefusedBeforeItIsCopied(t *testing.T) { + s := open(t) + k := smallKind(t, s, "stored", 1024, 4) + putCheckpointRaw(t, s, k, "prod", bytes.Repeat([]byte("x"), k.maxStored()+1)) + if _, err := k.load(s, "prod"); !errors.Is(err, ErrCheckpointTooLarge) { + t.Fatalf("an oversize stored record was read: %v", err) + } +} + +func TestEmptyCheckpointIsRefusedAtWrite(t *testing.T) { + s := open(t) + if err := s.SaveProposals(nil); !errors.Is(err, ErrEmptyCheckpoint) { + t.Fatalf("an empty proposal snapshot was stored: %v", err) + } + if err := s.SaveRDSCheckpoint("prod", []byte{}); !errors.Is(err, ErrEmptyCheckpoint) { + t.Fatalf("an empty rds checkpoint was stored: %v", err) + } + // It must stay a cold start, not become a checkpoint carrying nothing. + if _, err := s.LoadProposals(); !errors.Is(err, ErrNoCheckpoint) { + t.Fatalf("a refused empty save left state: %v", err) + } +} + +// TestScopeCapRefusesNewScopesAndKeepsUpdatingOldOnes: the cap must fall on +// the newcomer. A cap that stopped the fleet already being observed correctly +// would turn a mis-derived scope into a total outage of the growth finding. +func TestScopeCapRefusesNewScopesAndKeepsUpdatingOldOnes(t *testing.T) { + s := open(t) + k := smallKind(t, s, "cap", 4096, 3) + for i := 0; i < 3; i++ { + if err := k.save(s, fmt.Sprintf("scope-%d", i), []byte(`{"v":1}`)); err != nil { + t.Fatal(err) + } + } + err := k.save(s, "scope-3", []byte(`{"v":1}`)) + if !errors.Is(err, ErrCheckpointTooLarge) { + t.Fatalf("the scope cap did not hold: %v", err) + } + if !strings.Contains(err.Error(), "cap is 3") { + t.Fatalf("refusal does not name the cap: %v", err) + } + // Existing scopes keep working. + if err := k.save(s, "scope-1", []byte(`{"v":2}`)); err != nil { + t.Fatalf("an existing scope stopped updating at the cap: %v", err) + } + got, err := k.load(s, "scope-1") + if err != nil || string(got) != `{"v":2}` { + t.Fatalf("existing scope: %q %v", got, err) + } + // Freeing a slot lets the newcomer in. + if err := k.remove(s, "scope-0"); err != nil { + t.Fatal(err) + } + if err := k.save(s, "scope-3", []byte(`{"v":1}`)); err != nil { + t.Fatalf("a freed slot was not reusable: %v", err) + } +} + +func TestDeleteReportsWhetherItFreedASlot(t *testing.T) { + s := open(t) + if err := s.SaveRDSCheckpoint("prod", rdsBlob("prod", []int{0, 1})); err != nil { + t.Fatal(err) + } + if err := s.DeleteRDSCheckpoint("prod"); err != nil { + t.Fatal(err) + } + if err := s.DeleteRDSCheckpoint("prod"); !errors.Is(err, ErrNoCheckpoint) { + t.Fatalf("deleting an absent scope reported %v, want ErrNoCheckpoint", err) + } + if _, err := s.LoadRDSCheckpoint("prod"); !errors.Is(err, ErrNoCheckpoint) { + t.Fatalf("a deleted scope: %v", err) + } +} + +// ---- round trip and framing ---- + +func TestRDSCheckpointsRoundTripBitExactAcrossScopes(t *testing.T) { + s := open(t) + // Irregular histories: uneven gaps, differing lengths, one single-sample. + want := map[string][]byte{ + "111122223333/us-east-1": rdsBlob("111122223333/us-east-1", []int{0, 1, 4, 5, 19, 20, 41}), + "111122223333/eu-west-1": rdsBlob("111122223333/eu-west-1", []int{0, 13}), + "999988887777/us-east-1": rdsBlob("999988887777/us-east-1", []int{7}), + "sandbox": rdsBlob("sandbox", []int{0, 2, 3, 3, 9, 30, 31, 32, 90}), + } + for scope, blob := range want { + if err := s.SaveRDSCheckpoint(scope, blob); err != nil { + t.Fatal(err) + } + } + for scope, blob := range want { + got, err := s.LoadRDSCheckpoint(scope) + if err != nil { + t.Fatal(err) + } + if !bytes.Equal(got, blob) { + t.Fatalf("scope %q round-tripped to\n %s\nwant\n %s", scope, got, blob) + } + } + scopes, err := s.RDSCheckpointScopes() + if err != nil { + t.Fatal(err) + } + if len(scopes) != len(want) { + t.Fatalf("scopes %q", scopes) + } + for i := 1; i < len(scopes); i++ { + if scopes[i-1] >= scopes[i] { + t.Fatalf("scopes are not in ascending order: %q", scopes) + } + } + // Replacing one scope leaves the others byte-identical. + next := rdsBlob("sandbox", []int{0, 2, 3, 3, 9, 30, 31, 32, 90, 91}) + if err := s.SaveRDSCheckpoint("sandbox", next); err != nil { + t.Fatal(err) + } + for scope, blob := range want { + if scope == "sandbox" { + blob = next + } + got, err := s.LoadRDSCheckpoint(scope) + if err != nil || !bytes.Equal(got, blob) { + t.Fatalf("scope %q after a neighbour's rewrite: %s (%v)", scope, got, err) + } + } +} + +// TestBinaryAndNonUTF8BlobsSurviveIntact: the blob is opaque bytes, not a +// string. A framing that round-tripped it through a string type would replace +// invalid UTF-8 with U+FFFD and hand the producer back something it did not +// write. +func TestBinaryAndNonUTF8BlobsSurviveIntact(t *testing.T) { + s := open(t) + blob := make([]byte, 256) + for i := range blob { + blob[i] = byte(i) + } + if err := s.SaveRDSCheckpoint("bin", blob); err != nil { + t.Fatal(err) + } + got, err := s.LoadRDSCheckpoint("bin") + if err != nil { + t.Fatal(err) + } + if !bytes.Equal(got, blob) { + t.Fatalf("a binary blob came back changed: %x", got) + } +} + +// TestFramingIsDeterministic: identical state must frame to identical bytes, +// or every housekeeping tick rewrites the same page for no reason and no +// caller can compare two checkpoints without decoding them. encodeSnapshot +// relies on the same property. +func TestFramingIsDeterministic(t *testing.T) { + s := open(t) + blob := rdsBlob("prod", []int{0, 5, 9}) + var first []byte + for i := 0; i < 3; i++ { + if err := s.SaveRDSCheckpoint("prod", blob); err != nil { + t.Fatal(err) + } + var raw []byte + if err := s.db.View(func(tx *bolt.Tx) error { + raw = append(raw, tx.Bucket(bucketRDSCheckpoints).Get([]byte("prod"))...) + return nil + }); err != nil { + t.Fatal(err) + } + if i == 0 { + first = raw + continue + } + if !bytes.Equal(first, raw) { + t.Fatal("re-saving identical state produced different stored bytes") + } + } +} + +func TestProposalsAndRDSBucketsAreIndependent(t *testing.T) { + s := open(t) + props := []byte(`{"version":1,"ttlSeconds":86400,"records":[]}`) + if err := s.SaveProposals(props); err != nil { + t.Fatal(err) + } + if err := s.SaveRDSCheckpoint("brain", rdsBlob("brain", []int{0, 1})); err != nil { + t.Fatal(err) + } + got, err := s.LoadProposals() + if err != nil || !bytes.Equal(got, props) { + t.Fatalf("the rds bucket disturbed the proposals bucket: %q %v", got, err) + } + // The two kinds share a key name here on purpose: same key, different + // bucket, and neither can see the other. + gotRDS, err := s.LoadRDSCheckpoint("brain") + if err != nil || !bytes.Equal(gotRDS, rdsBlob("brain", []int{0, 1})) { + t.Fatalf("rds under the proposals key: %q %v", gotRDS, err) + } +} + +func TestCheckpointsSurviveReopen(t *testing.T) { + s := open(t) + path := s.db.Path() + props := []byte(`{"version":1,"ttlSeconds":86400,"records":[]}`) + blob := rdsBlob("prod/us-east-1", []int{0, 3, 17}) + if err := s.SaveProposals(props); err != nil { + t.Fatal(err) + } + if err := s.SaveRDSCheckpoint("prod/us-east-1", blob); err != nil { + t.Fatal(err) + } + if err := s.Close(); err != nil { + t.Fatal(err) + } + reopened, err := Open(path) + if err != nil { + t.Fatal(err) + } + defer reopened.Close() + got, err := reopened.LoadProposals() + if err != nil || !bytes.Equal(got, props) { + t.Fatalf("proposals did not survive a restart: %q %v", got, err) + } + gotRDS, err := reopened.LoadRDSCheckpoint("prod/us-east-1") + if err != nil || !bytes.Equal(gotRDS, blob) { + t.Fatalf("rds checkpoint did not survive a restart: %q %v", gotRDS, err) + } +} + +// TestConcurrentCheckpointAccess runs under -race. pkg/store is used from a +// running brain: the housekeeping timer writes while an HTTP handler reads, +// and the collector loop writes a different scope at the same time. The +// concurrency contract is the precedent's — bbolt serializes writers, these +// methods hold no state of their own — and the assertion that matters is that +// a reader never sees a blend of two writes. +func TestConcurrentCheckpointAccess(t *testing.T) { + s := open(t) + scopes := []string{"a", "b", "a/b", "café"} + versions := make([][]byte, 8) + for i := range versions { + versions[i] = rdsBlob("shared", []int{0, i, i * 3}) + } + var wg sync.WaitGroup + for i := 0; i < 8; i++ { + wg.Add(1) + go func(i int) { + defer wg.Done() + for j := 0; j < 25; j++ { + scope := scopes[(i+j)%len(scopes)] + if err := s.SaveRDSCheckpoint(scope, rdsBlob(scope, []int{0, j})); err != nil { + t.Errorf("save %q: %v", scope, err) + return + } + if err := s.SaveRDSCheckpoint("shared", versions[j%len(versions)]); err != nil { + t.Errorf("save shared: %v", err) + return + } + got, err := s.LoadRDSCheckpoint("shared") + if err != nil { + t.Errorf("load shared: %v", err) + return + } + whole := false + for _, v := range versions { + if bytes.Equal(got, v) { + whole = true + } + } + if !whole { + t.Errorf("a concurrent read saw a blend of two writes: %s", got) + return + } + if err := s.SaveProposals(rdsBlob("props", []int{j})); err != nil { + t.Errorf("save proposals: %v", err) + return + } + if _, err := s.LoadProposals(); err != nil { + t.Errorf("load proposals: %v", err) + return + } + if _, err := s.RDSCheckpointScopes(); err != nil { + t.Errorf("scopes: %v", err) + return + } + } + }(i) + } + wg.Wait() +} + +// TestNilSnapshotterIsRefused — SaveProposalsFrom is the one call that takes an +// interface, and a nil one must not panic inside a bbolt transaction. +func TestNilSnapshotterIsRefused(t *testing.T) { + s := open(t) + if err := s.SaveProposalsFrom(nil); err == nil { + t.Fatal("a nil snapshotter was accepted") + } +} + +// failingSnapshotter stands in for a producer whose own Snapshot fails. +type failingSnapshotter struct{} + +func (failingSnapshotter) Snapshot() ([]byte, error) { return nil, errors.New("boom") } + +func TestSnapshotterFailureIsNotStoredAsAnEmptyCheckpoint(t *testing.T) { + s := open(t) + prev := []byte(`{"version":1,"records":[]}`) + if err := s.SaveProposals(prev); err != nil { + t.Fatal(err) + } + if err := s.SaveProposalsFrom(failingSnapshotter{}); err == nil { + t.Fatal("a failed snapshot was stored") + } + got, err := s.LoadProposals() + if err != nil || !bytes.Equal(got, prev) { + t.Fatalf("a failed snapshot damaged the stored one: %q %v", got, err) + } +} + +// TestEnvelopeWithNoBlobIsRefused: a well-formed envelope carrying nothing is +// the one shape that could still read back as "a checkpoint exists and it is +// empty". SaveProposals/SaveRDSCheckpoint refuse an empty blob at write, and +// the reader refuses one that arrived any other way, so an empty checkpoint +// has no representation at either end. +func TestEnvelopeWithNoBlobIsRefused(t *testing.T) { + s := open(t) + k := smallKind(t, s, "noblob", 1024, 4) + putEnvelope(t, s, k, "prod", checkpointEnvelope{ + Version: 1, Kind: k.name, Scope: "prod", Codec: checkpointCodec, + }) + _, err := k.load(s, "prod") + if !errors.Is(err, ErrCorruptCheckpoint) { + t.Fatalf("an empty envelope reported %v, want ErrCorruptCheckpoint", err) + } + if errors.Is(err, ErrNoCheckpoint) { + t.Fatal("an envelope carrying no blob was reported as a cold start") + } +} + +// TestNothingHereApprovesOrApplies states the actuator prohibition over this +// package's own AST rather than over a comment about it. +// +// pkg/ec2's and pkg/rds's actuators are deliberately unreachable from the +// binary. Persisting a proposal is not approving one, and a persistence layer +// is exactly where that line would erode first — a "SaveApproval", a +// "MarkApplied" convenience, an import of the state machine so a record could +// be "repaired" on load. None of those exist, and this fails if one appears. +func TestNothingHereApprovesOrApplies(t *testing.T) { + fset := token.NewFileSet() + pkgs, err := parser.ParseDir(fset, ".", func(fi os.FileInfo) bool { + return !strings.HasSuffix(fi.Name(), "_test.go") + }, 0) + if err != nil { + t.Fatal(err) + } + pkg, ok := pkgs["store"] + if !ok { + t.Fatal("package store did not parse") + } + // Names that would mean this package had grown a verb. + banned := []string{"approve", "apply", "actuate", "execute", "resize", "reboot", "modify"} + // Packages whose actuators must stay out of this one's import graph, plus + // the state machine itself: pkg/store frames bytes and must never be able + // to construct, transition or repair a proposal record. + bannedImports := []string{ + `"github.com/agenticode/kilter/pkg/ec2"`, + `"github.com/agenticode/kilter/pkg/rds"`, + `"github.com/agenticode/kilter/pkg/whatif"`, + } + for name, file := range pkg.Files { + for _, imp := range file.Imports { + for _, bad := range bannedImports { + if imp.Path.Value == bad { + t.Errorf("%s imports %s", name, bad) + } + } + } + for _, decl := range file.Decls { + fn, ok := decl.(*ast.FuncDecl) + if !ok { + continue + } + lower := strings.ToLower(fn.Name.Name) + for _, b := range banned { + if strings.Contains(lower, b) { + t.Errorf("%s declares %s, which reads as a verb this package must not have", name, fn.Name.Name) + } + } + } + } +} diff --git a/pkg/store/persist_ext_test.go b/pkg/store/persist_ext_test.go new file mode 100644 index 0000000..d585e20 --- /dev/null +++ b/pkg/store/persist_ext_test.go @@ -0,0 +1,527 @@ +// The other half of the proposals bucket's contract, asserted against the +// package that actually produces the bytes. +// +// checkpoint_test.go tests the framing with stand-in documents, which is the +// right level for a byte pipe. What it cannot test is that the pipe is +// SUFFICIENT: that a real whatif.Store, carrying every state a proposal can +// rest in, survives a round trip well enough to be reloaded — and that a blob +// which did NOT come out of this pipe unchanged fails to reload rather than +// reloading as something more permissive than it was. +// +// The import lives in the external test package for seams_ext_test.go's +// reason: production pkg/store imports nothing above pkg/model, and asserting +// the seam here means a signature drift in pkg/whatif is a compile error in +// this package's own test run. +package store_test + +import ( + "bytes" + "errors" + "fmt" + "strings" + "testing" + "time" + + "github.com/agenticode/kilter/pkg/backtest" + "github.com/agenticode/kilter/pkg/store" + "github.com/agenticode/kilter/pkg/whatif" +) + +// whatif.Store is the ProposalSnapshotter pkg/store declares. Structural, so +// neither package imports the other. +var _ store.ProposalSnapshotter = (*whatif.Store)(nil) + +// No test here may read the wall clock: a proposal's identity is content- +// addressed and its expiry is clock-driven, so a test that used time.Now +// could not assert byte-identity or a deterministic set of states. +var ( + propFrom = time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC) + propTo = time.Date(2026, 1, 8, 0, 0, 0, 0, time.UTC) + propNow = time.Date(2026, 1, 8, 3, 0, 0, 0, time.UTC) +) + +func openStore(t *testing.T) *store.Store { + t.Helper() + s, err := store.Open(t.TempDir() + "/kilter.db") + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { s.Close() }) + return s +} + +// scoreFor builds a scorecard internally consistent with a policy — the same +// recipe pkg/whatif's own tests use, restated because helpers_test.go is +// internal to that package. +func scoreFor(p whatif.Policy, mut func(*backtest.Scorecard)) *backtest.Scorecard { + sc := &backtest.Scorecard{ + Policy: p.Hash(), + Cluster: "prod-east", + Window: [2]time.Time{propFrom, propTo}, + HorizonHours: 24, + DecisionIntervalHours: 24, + StarvationFactor: 1, + Snapshots: 2016, + Instants: 6, + Scored: 60, + Decisions: 12, + Refusals: map[string]int{backtest.CodeBelowChangeThreshold: 48}, + MemViolations: 2, + CPUStarvation: 1, + MemOOMKills: 3, + OracleGapPct: 30, + OracleGapPctApplied: 12, + PolicyCostUSD: 100, + OracleCostUSD: 80, + ForgoneSavingsUSD: 5, + FlipRate: 0.05, + Flips: 1, + ResourceRegretUSD: 20, + RiskRegretUSD: 10, + RegretUSD: 30, + Cost: backtest.DefaultCostModel(), + } + if mut != nil { + mut(sc) + } + return sc +} + +// betterCandidate is the "this deserves to pass" mutation; worseCandidate is +// its opposite, which the gate refuses and Create files as rejected. +func betterCandidate(sc *backtest.Scorecard) { + sc.RegretUSD = 20 + sc.ResourceRegretUSD = 12 + sc.RiskRegretUSD = 8 + sc.OracleGapPct = 22 +} + +func worseCandidate(sc *backtest.Scorecard) { + sc.RegretUSD = 90 + sc.ResourceRegretUSD = 60 + sc.RiskRegretUSD = 30 + sc.OracleGapPct = 55 + sc.MemOOMKills = 40 +} + +// specFor builds a spec whose candidate differs from the shipped policy on one +// in-envelope axis — the shape every real proposal has. +func specFor(headroom float64, pass bool, rationale string, evidence []string) whatif.Spec { + cand := whatif.DefaultPolicy() + cand.Rec.CPUHeadroom = headroom + mut := betterCandidate + if !pass { + mut = worseCandidate + } + return whatif.Spec{ + Kind: whatif.KindPolicyChange, + Target: whatif.Target{Cluster: "prod-east", Namespace: "payments", Class: "batch"}, + Baseline: whatif.DefaultPolicy(), + Candidate: cand, + BaselineScore: scoreFor(whatif.DefaultPolicy(), nil), + CandidateScore: scoreFor(cand, mut), + Envelope: whatif.DefaultEnvelope(), + Tolerance: whatif.DefaultTolerance(), + Rationale: rationale, + EvidenceIDs: evidence, + } +} + +var ( + tunerActor = whatif.Actor{Kind: whatif.ActorTuner, ID: "nightly"} + agentActor = whatif.Actor{Kind: whatif.ActorAgent, ID: "reasoner"} +) + +// populated builds a proposal store holding one record in every state a stored +// record may rest in, and returns it with the states it produced. +// +// StateDraft is absent on purpose and cannot be produced: whatif's +// UnmarshalJSON refuses a stored draft, because no record rests there. +func populated(t *testing.T) (*whatif.Store, map[string]whatif.State) { + t.Helper() + ws := whatif.NewStore() + at := func(d time.Duration) whatif.Clock { return whatif.FixedClock(propNow.Add(d)) } + + // Rejected by the gate, at filing time. + gateRejected, err := ws.Create(tunerActor, specFor(1.21, false, "regret is worse", nil), at(0)) + if err != nil { + t.Fatal(err) + } + // Filed, gate passed, awaiting a human — and then rejected BY that human, + // which is a different road to the same terminal state and a different + // audit history. + humanRejected, err := ws.Create(tunerActor, specFor(1.22, true, "candidate B", []string{"ev-1", "ev-2"}), at(0)) + if err != nil { + t.Fatal(err) + } + if _, err := ws.Reject(whatif.Actor{Kind: whatif.ActorHuman, ID: "bob"}, humanRejected.ID(), "not this quarter", at(time.Minute)); err != nil { + t.Fatal(err) + } + // Still gated: nobody has answered. + gated, err := ws.Create(agentActor, specFor(1.23, true, "candidate C", []string{"ev-3"}), at(0)) + if err != nil { + t.Fatal(err) + } + // Approved and then left to lapse. + expired, err := ws.Create(tunerActor, specFor(1.24, true, "candidate D", nil), at(0)) + if err != nil { + t.Fatal(err) + } + // Approved, and recorded — after the fact — as having been applied by + // something outside this process. pkg/store neither performs nor enables + // that; the state is here because losing "this was applied" across a + // restart is precisely the audit-trail hole §5.3 names. + applied, err := ws.Create(tunerActor, specFor(1.25, true, "candidate E", nil), at(0)) + if err != nil { + t.Fatal(err) + } + // Approved and still live at snapshot time. + approved, err := ws.Create(tunerActor, specFor(1.26, true, "candidate F", []string{"ev-9"}), at(0)) + if err != nil { + t.Fatal(err) + } + + alice, err := whatif.NewApprover(whatif.Actor{Kind: whatif.ActorHuman, ID: "alice"}) + if err != nil { + t.Fatal(err) + } + for _, id := range []string{expired.ID(), applied.ID()} { + if _, err := ws.Approve(alice, id, "ship it", at(time.Hour)); err != nil { + t.Fatal(err) + } + } + if _, err := ws.MarkApplied(whatif.Actor{Kind: whatif.ActorSystem, ID: "deploy"}, applied.ID(), "rolled out", at(2*time.Hour)); err != nil { + t.Fatal(err) + } + // Sweep past the first two approvals' TTL: the un-applied one expires, the + // applied one is already terminal and is not touched. + swept, err := ws.Sweep(at(whatif.DefaultApprovalTTL + 2*time.Hour)) + if err != nil { + t.Fatal(err) + } + if swept != 1 { + t.Fatalf("sweep moved %d records, want exactly the un-applied approval", swept) + } + // Approve the last one after the sweep, so it is still live in the bytes. + if _, err := ws.Approve(alice, approved.ID(), "ship it", at(whatif.DefaultApprovalTTL+3*time.Hour)); err != nil { + t.Fatal(err) + } + + want := map[string]whatif.State{ + gateRejected.ID(): whatif.StateRejected, + humanRejected.ID(): whatif.StateRejected, + gated.ID(): whatif.StateGated, + expired.ID(): whatif.StateExpired, + applied.ID(): whatif.StateApplied, + approved.ID(): whatif.StateApproved, + } + assertStates(t, ws, want) + return ws, want +} + +func assertStates(t *testing.T, ws *whatif.Store, want map[string]whatif.State) { + t.Helper() + got := map[string]whatif.State{} + for _, r := range ws.List() { + got[r.ID()] = r.State() + } + if len(got) != len(want) { + t.Fatalf("store holds %d records, want %d", len(got), len(want)) + } + for id, st := range want { + if got[id] != st { + t.Fatalf("proposal %s is %q, want %q", id, got[id], st) + } + } +} + +// TestProposalStoreRoundTripsBitExactWithEveryGateState is the fidelity +// assertion §5.3 needs: a brain restart must return the proposal plane it had, +// not an approximation of it. Bit-exact because anything less means the store +// is re-encoding, and a persistence layer that re-encodes a security-relevant +// record is a persistence layer that can change one. +func TestProposalStoreRoundTripsBitExactWithEveryGateState(t *testing.T) { + s := openStore(t) + ws, want := populated(t) + + blob, err := ws.Snapshot() + if err != nil { + t.Fatal(err) + } + if err := s.SaveProposals(blob); err != nil { + t.Fatal(err) + } + got, err := s.LoadProposals() + if err != nil { + t.Fatal(err) + } + if !bytes.Equal(got, blob) { + t.Fatalf("the proposal store did not round-trip byte for byte\n got %d bytes\nwant %d bytes", len(got), len(blob)) + } + + // And the bytes are still a store: every record comes back in exactly the + // state it was written in, including the two terminal ones and the live + // approval. This is the property the actuator prohibition rests on — a + // restart may not promote anything. + reloaded, err := whatif.Load(got) + if err != nil { + t.Fatalf("a store this package wrote did not reload: %v", err) + } + assertStates(t, reloaded, want) + + // Re-snapshotting the reloaded store reproduces the same bytes, so a + // second housekeeping tick after a restart is a no-op write rather than a + // diff against itself. + again, err := reloaded.Snapshot() + if err != nil { + t.Fatal(err) + } + if !bytes.Equal(again, blob) { + t.Fatal("a reload-then-snapshot cycle changed the bytes") + } +} + +// TestAForgedApprovalDoesNotComeBackApproved is the actuator prohibition, +// stated over the bytes. +// +// pkg/store is a byte pipe and deliberately does not parse a proposal, so it +// cannot itself refuse a forgery — and it must not silently repair one either. +// The property being asserted is therefore the pair: the pipe hands back +// exactly what it was given (it neither grants nor launders capability), and +// what it was given fails to load as approved. A gated proposal edited to say +// "approved" does not become an approved proposal by being written to disk and +// read back. +func TestAForgedApprovalDoesNotComeBackApproved(t *testing.T) { + s := openStore(t) + ws, _ := populated(t) + honest, err := ws.Snapshot() + if err != nil { + t.Fatal(err) + } + + for _, tc := range []struct { + name string + forge func(string) string + reason string + }{ + { + name: "gated record relabelled approved", + forge: func(s string) string { return strings.Replace(s, `"state": "gated"`, `"state": "approved"`, 1) }, + reason: "approval binding", + }, + { + name: "approved record's approval deleted, state kept", + forge: func(s string) string { + // Drop the approval object from the live-approved record but + // leave "state": "approved" in place. + i := strings.Index(s, `"approval"`) + if i < 0 { + t.Fatal("fixture has no approval to strip") + } + j := strings.Index(s[i:], "\n },\n") + if j < 0 { + t.Fatal("could not find the end of the approval object") + } + return s[:i] + s[i+j+len("\n },\n")+4:] + }, + reason: "approval binding", + }, + { + name: "approver rewritten to the author", + forge: func(s string) string { + return strings.Replace(s, `"id": "alice"`, `"id": "nightly"`, 1) + }, + reason: "self-approval", + }, + { + name: "gate verdict flipped to passed", + forge: func(s string) string { + return strings.Replace(s, `"passed": false`, `"passed": true`, 1) + }, + reason: "content fingerprint", + }, + } { + t.Run(tc.name, func(t *testing.T) { + forged := []byte(tc.forge(string(honest))) + if bytes.Equal(forged, honest) { + t.Fatal("the forgery changed nothing; the test proves nothing") + } + if err := s.SaveProposals(forged); err != nil { + t.Fatal(err) + } + got, err := s.LoadProposals() + if err != nil { + t.Fatal(err) + } + // The pipe is a pipe: no repair, no rejection, no rewriting. + if !bytes.Equal(got, forged) { + t.Fatal("the store rewrote the bytes it was given") + } + // And the forgery does not load. Not "loads without the approval", + // not "loads as gated" — does not load. + reloaded, err := whatif.Load(got) + if err == nil { + for _, r := range reloaded.List() { + if r.State() == whatif.StateApproved || r.State() == whatif.StateApplied { + t.Fatalf("a forged blob produced a %s record (%s)", r.State(), r.ID()) + } + } + t.Fatalf("a forged blob (%s) loaded", tc.reason) + } + }) + } +} + +// TestColdBrainStartIsDistinguishableFromALostStore walks the exact branch a +// caller writes at NewBrain, and pins that the two outcomes it must tell apart +// really do arrive differently. If these ever collapsed, a brain whose file +// was truncated would come up serving an empty proposal plane and reporting +// nothing wrong — the audit trail would be gone and the only evidence of it +// would be its absence. +func TestColdBrainStartIsDistinguishableFromALostStore(t *testing.T) { + s := openStore(t) + + blob, err := s.LoadProposals() + if !errors.Is(err, store.ErrNoCheckpoint) { + t.Fatalf("a first boot reported %v, want ErrNoCheckpoint", err) + } + if blob != nil { + t.Fatal("a first boot returned bytes as well as an error") + } + + ws, want := populated(t) + if err := s.SaveProposalsFrom(ws); err != nil { + t.Fatal(err) + } + blob, err = s.LoadProposals() + if err != nil { + t.Fatalf("a warm boot reported %v", err) + } + reloaded, err := whatif.Load(blob) + if err != nil { + t.Fatal(err) + } + assertStates(t, reloaded, want) + + // Now lose it. Truncating the stored value is the cheapest stand-in for + // the half-written record a crash leaves behind. + if err := s.SaveProposals(blob[:len(blob)/2]); err != nil { + t.Fatal(err) + } + blob, err = s.LoadProposals() + if err != nil { + t.Fatalf("a truncated snapshot is still a well-framed one: %v", err) + } + if _, err := whatif.Load(blob); err == nil { + t.Fatal("half a proposal store loaded") + } + // The framing layer's own corruption is the louder case, and it must not + // look like a first boot either. + if err := s.SaveProposals(blob); err != nil { + t.Fatal(err) + } +} + +// TestOneRecordIsWellUnderItsShareOfTheCap is the arithmetic behind +// MaxProposalsCheckpointBytes, asserted rather than asserted-in-a-comment. +// +// The record built here holds every free-text field at pkg/whatif's own limit: +// a full-length rationale, the maximum citation set at the maximum citation +// length, and a full-length approval note. If a later change to pkg/whatif +// grows a record past its share of the cap, a full store stops being storable +// — and this fails first, in this package, instead of at 3am in a brain whose +// proposals silently stopped persisting. +func TestOneRecordIsWellUnderItsShareOfTheCap(t *testing.T) { + const whatifMaxRecords = 1000 // pkg/whatif's unexported maxRecords + + // maxEvidenceIDs citations at maxEvidenceIDLen, and a rationale at the + // bound that actually binds: Spec.normalize checks maxRationale (4096) and + // then runs sanitizeNote, whose maxNote (2048) refuses first, so 2048 is + // the largest rationale a proposal can carry. + ids := make([]string, 64) + for i := range ids { + ids[i] = fmt.Sprintf("%s%03d", strings.Repeat("e", 509), i) // maxEvidenceIDLen + } + ws := whatif.NewStore() + at := func(d time.Duration) whatif.Clock { return whatif.FixedClock(propNow.Add(d)) } + rec, err := ws.Create(tunerActor, specFor(1.28, true, strings.Repeat("r", 2048), ids), at(0)) + if err != nil { + t.Fatal(err) + } + alice, err := whatif.NewApprover(whatif.Actor{Kind: whatif.ActorHuman, ID: strings.Repeat("a", 128)}) + if err != nil { + t.Fatal(err) + } + if _, err := ws.Approve(alice, rec.ID(), strings.Repeat("n", 2048), at(time.Hour)); err != nil { + t.Fatal(err) + } + full, err := ws.Snapshot() + if err != nil { + t.Fatal(err) + } + empty, err := whatif.NewStore().Snapshot() + if err != nil { + t.Fatal(err) + } + perRecord := len(full) - len(empty) + t.Logf("worst-case record: %d bytes; × %d records = %d bytes; cap is %d", + perRecord, whatifMaxRecords, perRecord*whatifMaxRecords, store.MaxProposalsCheckpointBytes) + + if perRecord*whatifMaxRecords > store.MaxProposalsCheckpointBytes { + t.Fatalf("a full proposal store (%d records × %d bytes = %d) no longer fits under the %d-byte cap", + whatifMaxRecords, perRecord, perRecord*whatifMaxRecords, store.MaxProposalsCheckpointBytes) + } + // And it does fit, framed, through the real bucket. + s := openStore(t) + if err := s.SaveProposals(full); err != nil { + t.Fatal(err) + } + got, err := s.LoadProposals() + if err != nil || !bytes.Equal(got, full) { + t.Fatalf("a maximal record did not round-trip: %v", err) + } +} + +// TestHousekeepingTickShape is the call sequence PERSIST-FINDINGS.md §5 tells +// pkg/api to hang off its timer, executed here so the documented sequence is +// one that compiles and runs rather than one that reads well. +func TestHousekeepingTickShape(t *testing.T) { + s := openStore(t) + ws, _ := populated(t) + tick := func(now time.Time) { + t.Helper() + if _, err := ws.Sweep(whatif.FixedClock(now)); err != nil { + t.Fatal(err) + } + if err := s.SaveProposalsFrom(ws); err != nil { + t.Fatal(err) + } + } + // Two ticks before the live approval lapses, one after. + tick(propNow.Add(whatif.DefaultApprovalTTL + 4*time.Hour)) + before, err := s.LoadProposals() + if err != nil { + t.Fatal(err) + } + tick(propNow.Add(whatif.DefaultApprovalTTL + 5*time.Hour)) + if same, err := s.LoadProposals(); err != nil || !bytes.Equal(same, before) { + t.Fatalf("a tick with nothing to sweep rewrote the snapshot: %v", err) + } + tick(propNow.Add(3 * whatif.DefaultApprovalTTL)) + after, err := s.LoadProposals() + if err != nil { + t.Fatal(err) + } + if bytes.Equal(after, before) { + t.Fatal("sweeping the live approval to expired did not reach the store") + } + reloaded, err := whatif.Load(after) + if err != nil { + t.Fatal(err) + } + for _, r := range reloaded.List() { + if r.State() == whatif.StateApproved { + t.Fatalf("proposal %s is still approved after its TTL elapsed", r.ID()) + } + } +} diff --git a/pkg/store/proposals.go b/pkg/store/proposals.go new file mode 100644 index 0000000..26b6700 --- /dev/null +++ b/pkg/store/proposals.go @@ -0,0 +1,136 @@ +package store + +// The proposals bucket — pkg/api/WHATIFROUTES-FINDINGS.md §5.3's "a brain +// restart loses every filed proposal and its audit trail". +// +// whatif.Store.Snapshot() and whatif.Load() already move the whole proposal +// store as bytes; §5.3 said where those bytes go and that pkg/store was out of +// scope that cycle. This is that bucket. It stores ONE value, like +// bucketEvidence and for the same reason: a whatif.Store is fleet-wide (its +// proposals carry their own target cluster, and `GET /api/v1/proposals?cluster=` +// filters the one store rather than selecting between several), so a scope key +// here would be a parameter every caller had to invent a value for. +// +// # Persisting a proposal is not approving one +// +// The blob is opaque to this package and stays opaque. pkg/store does not +// import pkg/whatif, does not parse a record, and has no method that +// transitions one — the state machine is pkg/whatif's security property and +// splitting it across a persistence layer would be the way to lose it. +// +// What that buys is precise: this bucket is a byte pipe, so a proposal comes +// back in exactly the gate state it was written in, and a blob that arrives +// from anywhere else — an edited file, a truncated write, an attacker with the +// db — is still checked by whatif.Record.UnmarshalJSON, which recomputes the +// content fingerprint and re-verifies that an approved record carries a live +// approval bound to that fingerprint and that verdict by a human who is not +// the author. A forged "state":"approved" fails to load; it does not load as +// approved. persist_ext_test.go forges one and asserts exactly that. +// +// The direction of the remaining risk is worth stating: the worst this bucket +// can do is lose proposals or refuse to load them. It cannot manufacture an +// approval, and nothing here or downstream of here makes one executable — +// pkg/ec2 and pkg/rds's actuators remain unreachable from the binary. + +// bucketProposals holds the single framed whatif.Store snapshot. +var bucketProposals = []byte("proposals") + +// proposalsScope is the key the fleet-wide proposal store is filed under. +// A constant, for evidenceScope's reason. +const proposalsScope = "brain" + +// ProposalsEnvelopeVersion is the framing version this build writes. It is +// independent of whatif's own storeVersion, which travels inside the blob: the +// question this version answers is "how is the blob framed", not "what is in +// it". +const ProposalsEnvelopeVersion = 1 + +// MaxProposalsCheckpointBytes bounds the DECODED snapshot, at write and again +// at read. +// +// The arithmetic is taken from pkg/whatif's own bounds rather than from a +// measurement of today's data — a cap chosen from what a producer CAN emit +// survives the producer getting busier. One record's free text is bounded by +// maxEvidenceIDs × maxEvidenceIDLen (64 × 512 = 32,768 B), the rationale +// (2,048 B: Spec.normalize checks maxRationale = 4096 and then runs +// sanitizeNote, whose maxNote = 2048 refuses first, so the larger constant +// never binds), and one note per human-supplied transition over the longest +// legal path draft→gated→approved→applied (2 × maxNote) plus the approval's +// own note — 6,144 B. TestOneRecordIsWellUnderItsShareOfTheCap builds exactly +// that record and measures 45,278 B, which with the one further note the +// longest path adds is ~47,326 B, call it 46 KiB. whatif caps itself at +// maxRecords = 1000: +// +// 46 KiB × 1,000 ≈ 45.1 MiB +// +// 64 MiB covers that with ~40% headroom, and matches +// MaxEvidenceCheckpointBytes, which was chosen for the constraint that binds +// here too: a bbolt value is materialized whole on both sides of a write. +// +// What 64 MiB does NOT cover, stated rather than left to be discovered: +// encoding/json escapes '<', '>' and '&' to six bytes each, so a store whose +// every rationale, citation and note is '<' frames to roughly 240 MiB and is +// REFUSED AT WRITE. That is deliberate, and it is the evidence precedent's +// choice — refuse loudly and name the cap, rather than write a value that +// cannot be read back within a sane allocation. The failure is safe in the +// only direction that matters: the in-memory store keeps working, nothing is +// corrupted, and a lost proposal is not an approved one. A realistic record is +// ~6.5 KiB — the structure above plus short prose — so a realistic full store +// is ~6.5 MiB and this is ~10× production headroom. +const MaxProposalsCheckpointBytes = 64 << 20 + +var proposalsKind = checkpointKind{ + name: "proposals", + bucket: bucketProposals, + version: ProposalsEnvelopeVersion, + max: MaxProposalsCheckpointBytes, + fixedScope: proposalsScope, +} + +// SaveProposals stores the proposal store's bytes, replacing the previous +// snapshot. blob is whatever whatif.Store.Snapshot() returned; this package +// does not look inside it. +// +// Callers hang this off the housekeeping timer, immediately after the +// whatif.Store.Sweep that moves lapsed approvals to expired — sweeping after +// snapshotting would persist a store one tick staler than the one in memory. +// PERSIST-FINDINGS.md §5 has the loop. +func (s *Store) SaveProposals(blob []byte) error { + return proposalsKind.save(s, "", blob) +} + +// LoadProposals returns the stored bytes for whatif.Load. +// +// A cold store is an ERROR — errors.Is(err, ErrNoCheckpoint) — not a nil blob. +// The distinction is the point: "no proposals have ever been filed" and "there +// were proposals and I cannot read them" must not arrive in the same shape, +// because a caller that treats the second as the first starts an empty store +// over an audit trail it just failed to read. +func (s *Store) LoadProposals() ([]byte, error) { + return proposalsKind.load(s, "") +} + +// ProposalSnapshotter is the seam whatif.Store satisfies. It is declared here, +// structurally, so pkg/store keeps its place at the bottom of the dependency +// order and never imports the state machine — the same arrangement +// EvidenceCheckpointStore has with pkg/evidence. persist_ext_test.go asserts +// the satisfaction at compile time, so a signature drift in pkg/whatif breaks +// this package's test build rather than some caller's. +type ProposalSnapshotter interface { + Snapshot() ([]byte, error) +} + +// SaveProposalsFrom snapshots and stores in one call, which is the whole of +// what a housekeeping tick needs to do. A snapshot that cannot be produced is +// returned unwrapped in a CheckpointError, because the failure is the +// producer's, not this bucket's. +func (s *Store) SaveProposalsFrom(p ProposalSnapshotter) error { + if p == nil { + return proposalsKind.errf(proposalsScope, ErrEmptyCheckpoint, "nil snapshotter") + } + blob, err := p.Snapshot() + if err != nil { + return proposalsKind.errf(proposalsScope, ErrCorruptCheckpoint, "snapshotting: %v", err) + } + return s.SaveProposals(blob) +} diff --git a/pkg/store/rdscheckpoint.go b/pkg/store/rdscheckpoint.go new file mode 100644 index 0000000..c811706 --- /dev/null +++ b/pkg/store/rdscheckpoint.go @@ -0,0 +1,133 @@ +package store + +// The RDS checkpoint bucket — pkg/rds/GROWTH-FINDINGS.md §6's +// "SaveRDSCheckpoint/LoadRDSCheckpoint do not exist", and the persistence +// cmd/WIRING-FINDINGS.md §6.2/§6.3 named as the shared prerequisite. +// +// §6 asked for the shape SaveEvidenceCheckpoint has — "a scope-keyed, +// gzip-framed, version-refusing envelope" — and the one word in that sentence +// the precedent does not actually implement is the first. Evidence stores one +// blob for the brain; an RDS checkpoint is per scope, because +// Target.StorageHistory is the state and a target only means anything inside +// the account and region it was collected from. checkpoint.go is the +// generalization; this file is its five parameters. +// +// # What a scope is, and what happens if two of them collapse +// +// A scope is the caller's name for the collection boundary — in practice +// "/", whatever cmd/ derives from the credentials and the +// endpoint it queried. pkg/store does not parse it and does not need to: the +// key is the scope's bytes verbatim, so "prod", "us-east-1" and +// "prod/us-east-1" are three keys that cannot become each other. +// +// The failure this prevents is not a crash. GROWTH-FINDINGS §7 spells it out: +// "a checkpoint restored from another scope, or a history attached to a reused +// identifier, produces a decrease" — and a decrease is refused by +// storage-growth-history-inconsistent, so the FIRST-ORDER outcome of a +// collision is a refusal, not a lie. The dangerous case is the collision that +// looks consistent: two same-named databases in two regions whose histories +// interleave into a plausible upward slope, producing a growth verdict about a +// database that never grew. That verdict would carry a rate, a projection and +// a confident tone. Hence a key encoding with nothing to collide and a stored +// scope the reader re-checks. + +// bucketRDSCheckpoints holds scope → framed RDS domain checkpoint. +var bucketRDSCheckpoints = []byte("rds-checkpoints") + +// RDSEnvelopeVersion is the framing version this build writes. Independent of +// whatever version pkg/rds's own Domain.Checkpoint carries inside the blob. +const RDSEnvelopeVersion = 1 + +// MaxRDSCheckpointBytes bounds one scope's DECODED checkpoint, at write and +// again at read. +// +// The arithmetic is GROWTH-FINDINGS §7's, taken at its own stated bound rather +// than re-derived: DefaultGrowthMaxObservations = 768 observations per +// instance at ~40 bytes each is ~30 KiB per instance, and §7 sizes the worst +// realistic account at 1,000 instances: +// +// 768 × 40 B = 30,720 B per instance +// 30,720 B × 1,000 instances ≈ 29.3 MiB of history +// +// The rest of each Target — identifier, class, engine, the sizing inputs — is +// small next to its history but not free, and JSON keys are repeated per +// observation, so double it: ~59 MiB, and the cap is 64 MiB. Same value as +// MaxEvidenceCheckpointBytes, same binding constraint (a bbolt value is +// materialized whole on both sides of a write), and the same remedy in the +// refusal: an account that has outgrown this needs a tighter +// Config.Growth.MaxObservations, which §7 already names as the knob. +// +// [unverified: no 1,000-instance account has been measured; §7 flags the same +// figure as unverified and this cap inherits that.] +const MaxRDSCheckpointBytes = 64 << 20 + +// MaxRDSCheckpointScopes bounds distinct scopes. +// +// A scope-keyed bucket with no cap on scopes is not a ceiling, it is a growth +// curve: a caller deriving a scope from something accidentally per-run — a +// session id, a timestamp, a hostname — writes a new key every tick and the +// file grows until the disk does not. The count cap is what makes the total a +// number that can be stated: 64 scopes × 64 MiB is a 4 GiB absolute ceiling, +// and at the ~10× compression this framing gets on repetitive observation JSON +// [unverified: not measured for RDS checkpoints; the figure is history.go's +// for snapshot JSON] a realistic full fleet is a few hundred MiB on disk. +// +// 64 is chosen as "more account/region pairs than one brain should be watching +// serially". Past it, a NEW scope is refused by name while existing scopes +// keep updating — the failure lands on the newcomer, not on the fleet that is +// already being observed correctly. DeleteRDSCheckpoint is how an operator +// reclaims a slot whose scope no longer exists. +const MaxRDSCheckpointScopes = 64 + +var rdsKind = checkpointKind{ + name: "rds", + bucket: bucketRDSCheckpoints, + version: RDSEnvelopeVersion, + max: MaxRDSCheckpointBytes, + maxScopes: MaxRDSCheckpointScopes, +} + +// SaveRDSCheckpoint stores one scope's checkpoint, replacing the previous one. +// blob is whatever rds.Domain.Checkpoint() produced; this package does not +// look inside it. +// +// Callers write this at the end of each collection tick, after Observe. See +// PERSIST-FINDINGS.md §5 for the sequence, which is GROWTH-FINDINGS §6's +// snippet with the two missing calls filled in. +func (s *Store) SaveRDSCheckpoint(scope string, blob []byte) error { + return rdsKind.save(s, scope, blob) +} + +// LoadRDSCheckpoint returns one scope's checkpoint bytes for +// rds.Domain.Restore. +// +// A scope observed for the first time is errors.Is(err, ErrNoCheckpoint), and +// that is the only error a caller may proceed through: the domain then starts +// with no history and every instance reports GrowthNoHistory, which is §6's +// intended cold-start behaviour. Every other error means a checkpoint exists +// and could not be read, and proceeding through THAT silently restarts a +// 14-day observation window while reporting nothing wrong — the growth finding +// would simply never appear, and no one would know why. +func (s *Store) LoadRDSCheckpoint(scope string) ([]byte, error) { + return rdsKind.load(s, scope) +} + +// RDSCheckpointScopes lists the stored scopes in ascending key order, so an +// operator can see what the file is holding and which slots the scope cap is +// spending. +func (s *Store) RDSCheckpointScopes() ([]string, error) { + return rdsKind.scopes(s) +} + +// DeleteRDSCheckpoint drops one scope's checkpoint and frees its slot. +// +// It exists because MaxRDSCheckpointScopes is otherwise a one-way door: an +// account that is renamed or decommissioned would hold its slot forever. The +// deletion is real and unrecoverable — the growth history goes with it, and +// the next report for that scope is GrowthNoHistory until a new window fills. +// That is the safe direction (a finding is lost, never manufactured), but it +// is a loss, so an absent scope reports ErrNoCheckpoint rather than quietly +// succeeding. +func (s *Store) DeleteRDSCheckpoint(scope string) error { + return rdsKind.remove(s, scope) +} diff --git a/pkg/store/store.go b/pkg/store/store.go index 47760ea..c75e971 100644 --- a/pkg/store/store.go +++ b/pkg/store/store.go @@ -25,8 +25,9 @@ var ( bucketPlans = []byte("plans") // cluster/timestamp → Plan ) -// Two more buckets live in their own files, next to the code that owns them: -// bucketSnapHistory (history.go) and bucketEvidence (evidence.go). +// Four more buckets live in their own files, next to the code that owns them: +// bucketSnapHistory (history.go), bucketEvidence (evidence.go), +// bucketProposals (proposals.go) and bucketRDSCheckpoints (rdscheckpoint.go). // PlanHistoryLimit bounds retained plans per cluster. Pruning keeps the newest // PlanHistoryLimit by Plan.CreatedAt, not by insertion order. @@ -59,7 +60,10 @@ func Open(path string) (*Store, error) { return nil, fmt.Errorf("store: open %s: %w", path, err) } err = db.Update(func(tx *bolt.Tx) error { - for _, b := range [][]byte{bucketRecommender, bucketSnapshots, bucketPlans, bucketSnapHistory, bucketEvidence} { + for _, b := range [][]byte{ + bucketRecommender, bucketSnapshots, bucketPlans, + bucketSnapHistory, bucketEvidence, bucketProposals, bucketRDSCheckpoints, + } { if _, err := tx.CreateBucketIfNotExists(b); err != nil { return err }