Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -145,7 +145,7 @@ What works end-to-end:
- Period lifecycle (Open ↔ Closed) with **frozen-snapshot semantics**: closing a period captures `(frozen_quantity, frozen_event_count, watermark_at_close_ms)` on the manifest entry; closed-period queries return the frozen snapshot + pending Correction/Retraction adjustments + `net_total`, so an invoice line stays stable even as late corrections trickle in. New `Usage` events in a closed period are rejected; reopening discards the frozen snapshot.
- Shutdown drain — `ctrl+c` flushes the memtable before exiting
- proptest property tests for spec §19 invariants (`tests/properties.rs`): raw=rollup totals, dedupe idempotence under retry, compaction-preserves-sum, recovery-preserves-sum (with and without prior flush), rollup-tick idempotence, payload-conflict detection — 32 randomized cases per property, regenerated on every CI run
- Deterministic Simulation Testing (`tests/dst.rs`): state-machine driver that runs random sequences of `Ingest / Flush / RollupTick / CompactTick / Restart / ClosePeriod` ops against a parallel reference model, asserting `raw == model` after every step and `raw == rollup` at the end. Failures shrink to the minimum-length op trace, so any §19 violation surfaces as a reproducible recipe.
- Deterministic Simulation Testing (`tests/dst.rs`): state-machine driver that runs random sequences of `Ingest / Flush / RollupTick / CompactTick / Restart / ClosePeriod / CorruptLatestManifestAndRestart` ops against a parallel reference model, asserting `raw == model` after every step and `raw == rollup` at the end. Manifest-corruption sequences exercise the generation-rollback path; the model resyncs to whatever the SUT believes after recovery, so any post-rollback divergence (panic, raw-vs-rollup drift, ghost events) shows up as a failure. A dedicated fail-closed test verifies that a corrupt-only-generation manifest (no older generation to fall back to) refuses to start. Failures shrink to the minimum-length op trace, so any §19 violation surfaces as a reproducible recipe.

Manifest layout (Phase A — generations):

Expand Down
132 changes: 132 additions & 0 deletions tests/dst.rs
Original file line number Diff line number Diff line change
Expand Up @@ -100,6 +100,13 @@ enum Op {
Restart,
/// Close the (account, 2026-01) period.
ClosePeriod { account: String },
/// Corrupt the latest manifest generation file on disk, then
/// restart. Recovery should fall back to the previous generation;
/// the model resyncs to whatever the SUT believes is true after
/// the rollback. No-op if there's nothing to roll back to (i.e.
/// `current_generation < 2`) — the corruption step is skipped and
/// the restart runs as normal.
CorruptLatestManifestAndRestart,
}

fn arb_event() -> impl Strategy<Value = UsageEvent> {
Expand Down Expand Up @@ -139,6 +146,7 @@ fn arb_op() -> impl Strategy<Value = Op> {
1 => Just(Op::Restart),
1 => prop::sample::select(ACCOUNTS)
.prop_map(|a| Op::ClosePeriod { account: a.to_string() }),
1 => Just(Op::CorruptLatestManifestAndRestart),
]
}

Expand Down Expand Up @@ -191,6 +199,30 @@ impl Model {
fn sum_for(&self, account: &str) -> i128 {
self.sum.get(account).copied().unwrap_or(0)
}

/// Replace this model's state with the SUT's. Used after a
/// corruption-and-restart op: the manifest rollback may have
/// invalidated some of the prior model state (segments only
/// referenced by the corrupt generation are no longer visible to
/// queries), so the reference model resyncs to whatever the SUT
/// believes is true after recovery. Subsequent ops can then be
/// modeled normally on top of the new baseline.
async fn resync_from_sut(&mut self, harness: &DstHarness) {
self.acked.clear();
self.sum.clear();
self.closed.clear();
for acc in ACCOUNTS {
for ev in harness.scan_account_events(acc).await {
let (id_hash, payload_hash) = compute_event_hashes(&ev);
self.acked.insert(id_hash, payload_hash);
*self.sum.entry(acc.to_string()).or_insert(0) += ev.quantity;
}
}
let manifest = harness.state.manifest.read().await;
for cp in &manifest.closed_periods {
self.closed.insert((cp.account_id.clone(), cp.year, cp.month));
}
}
}

/// Single-process harness that drives ingest / flush / rollup /
Expand Down Expand Up @@ -472,6 +504,50 @@ impl DstHarness {
.unwrap_or(0)
}

/// Scan all raw events for an account (segments + memtable). Used
/// by the model's `resync_from_sut` to re-derive `acked` after a
/// corruption-rollback restart.
async fn scan_account_events(&self, account: &str) -> Vec<UsageEvent> {
let plan = QueryPlan {
source: QuerySource::RawEvents,
account_id: Some(account.into()),
from_ms: 0,
to_ms: i64::MAX,
filters: vec![],
group_by: vec![],
metrics: HashMap::new(),
limit: None,
};
execute_plan(&self.state, &plan)
.await
.into_iter()
.filter_map(|v| serde_json::from_value::<UsageEvent>(v).ok())
.collect()
}

/// Overwrite the latest manifest generation file on disk with
/// garbage, then restart. Recovery should fall back to generation
/// N-1 (per `Manifest::load_from_generations`). If only one
/// generation exists (or none), the corruption step is skipped
/// and the restart runs as a vanilla restart — the alternative is
/// a fail-closed `Recovery::run_startup_recovery` Err that would
/// crash the test, and the fail-closed path is covered by a
/// dedicated unit test below.
fn corrupt_latest_manifest_and_restart(self) -> Self {
let manifest_dir = self.state.config.db_root.join("manifest");
let current_path = manifest_dir.join("CURRENT");
let current_gen: Option<u64> = std::fs::read_to_string(&current_path)
.ok()
.and_then(|s| s.trim().parse().ok());
if let Some(g) = current_gen {
if g >= 2 {
let gen_path = manifest_dir.join(format!("manifest-{:06}.json", g));
std::fs::write(&gen_path, b"{ corrupt").expect("corrupt manifest write");
}
}
self.restart()
}

async fn raw_sums(&self) -> HashMap<String, i128> {
let mut out = HashMap::new();
for acc in ACCOUNTS {
Expand Down Expand Up @@ -534,6 +610,10 @@ async fn step(mut harness: DstHarness, model: &mut Model, op: Op) -> DstHarness
model.close(&account, TARGET_YEAR, TARGET_MONTH);
harness.close_period(&account, TARGET_YEAR, TARGET_MONTH).await;
}
Op::CorruptLatestManifestAndRestart => {
harness = harness.corrupt_latest_manifest_and_restart();
model.resync_from_sut(&harness).await;
}
}
harness
}
Expand Down Expand Up @@ -571,9 +651,61 @@ fn op_label(op: &Op) -> &'static str {
Op::CompactTick => "CompactTick",
Op::Restart => "Restart",
Op::ClosePeriod { .. } => "ClosePeriod",
Op::CorruptLatestManifestAndRestart => "CorruptLatestManifestAndRestart",
}
}

/// Fail-closed coverage: if `manifest/CURRENT` references a generation
/// number but that generation file is corrupt *and there's no older
/// generation to fall back to*, recovery must return an error rather
/// than silently start with an empty DB. (The state-machine driver
/// above intentionally skips this case — it would crash the test —
/// so we cover it here.)
#[test]
fn corrupt_only_manifest_generation_is_fatal() {
let rt = rt();
rt.block_on(async {
let harness = DstHarness::new();
let ev = UsageEvent {
event_id: EventId("evt".into()),
kind: EventKind::Usage,
correction_ref: None,
account_id: AccountId("acc_a".into()),
subscription_id: Some(SubscriptionId("sub_1".into())),
product_id: ProductId("ai_gateway".into()),
meter_id: MeterId("tokens.input".into()),
timestamp_ms: base_ts() + 1000,
quantity: 1,
unit: Unit("token".into()),
source: SourceId("test".into()),
model_id: Some(ModelId("m1".into())),
dimensions: SmallDimensions::default(),
ingested_at_ms: 0,
};
harness.ingest(vec![ev]).await;
harness.flush().await; // → generation 1

let db_root = harness.state.config.db_root.clone();
let manifest_dir = db_root.join("manifest");
std::fs::write(manifest_dir.join("manifest-000001.json"), b"{ corrupt")
.expect("corrupt write");

// Take ownership of the TempDir guard so it survives past the
// recovery call. Dropping `harness` (which holds the guard)
// would otherwise delete db_root before recovery runs.
let DstHarness { state, _tmp } = harness;
drop(state);

let recovery = Recovery::new(db_root);
let result = recovery.run_startup_recovery(100_000);
assert!(
result.is_err(),
"recovery must fail closed when no manifest generation parses"
);
drop(_tmp);
});
}

proptest! {
// 64 sequences ≈ 1–2 seconds locally; each step does real disk I/O
// (WAL fsync, segment write, manifest rename). The shrinker
Expand Down
Loading