tidaldb/tidal/tests/m9_retroactive_purge.rs
jordan 6f26d03c77 feat(m9): implement purge re-materialization engine
Adds the M9 purge re-materialization feature: a WAL-replay background
engine that rebuilds community cohort aggregates for a (user, community)
pair after retroactive signal purge, restoring ranking correctness without
modifying the immutable WAL.

Key additions:
- cohort::rematerialization module: PurgeJobQueue, RematerializationEngine,
  WAL replay, atomic CohortSignalLedger swap, BLAKE3 audit log, metrics counters
- TidalDb::{submit_purge_job, purge_job_status, rematerialization_metrics}
  public API (db/rematerialization.rs)
- Engine auto-starts in persistent mode; clean shutdown before WAL teardown
- 6 integration tests in tests/m9_purge_remat.rs covering ephemeral and
  persistent modes, job lifecycle, and multi-job independence
- Split oversized files to comply with 600-line limit: db/mod.rs →
  db/from_parts.rs, entities/revocation.rs → revocation/{mod,tests}.rs,
  schema/validation/builders.rs → builders/{mod,tests}.rs,
  signals/warm.rs → warm/{mod,tests,proptests}.rs
- Fix pre-existing bootstrap errors: export AuditKind from session module,
  add overrides_rejected to SessionSnapshot deserialization
2026-03-03 19:18:16 -07:00

430 lines
14 KiB
Rust

//! M9 Retroactive Purge Integration Tests.
//!
//! Exercises the end-to-end retroactive signal purge pipeline:
//! - Contribution log is populated on `signal_with_context` with cohort attribution.
//! - `request_community_purge` retracts contributions from the live cohort ledger.
//! - Purge manifests are persisted to durable storage.
//! - `list_purge_manifests` retrieves persisted manifests.
//! - Second purge on the same user is idempotent (no-op on the ledger).
//! - Scores floor at 0.0 and never go negative.
//! - Purging one user does not affect another user's cohort contributions.
#![allow(clippy::unwrap_used, clippy::cast_precision_loss)]
use std::collections::HashMap;
use std::time::Duration;
use tidaldb::TidalDb;
use tidaldb::cohort::{CohortDef, Predicate};
use tidaldb::schema::{DecaySpec, EntityId, EntityKind, SchemaBuilder, Timestamp, Window};
// ── Shared fixtures ──────────────────────────────────────────────────────────
fn purge_schema() -> tidaldb::schema::Schema {
let mut builder = SchemaBuilder::new();
let _ = builder
.signal(
"view",
EntityKind::Item,
DecaySpec::Exponential {
half_life: Duration::from_secs(7 * 24 * 3600),
},
)
.windows(&[Window::AllTime])
.velocity(false)
.add();
builder.build().expect("purge schema must be valid")
}
fn open_ephemeral() -> TidalDb {
TidalDb::builder()
.ephemeral()
.with_schema(purge_schema())
.open()
.expect("ephemeral open")
}
fn setup_cohort(db: &TidalDb, cohort: &str, field: &str, value: &str) {
db.define_cohort(CohortDef {
name: cohort.to_string(),
predicate: Predicate::Eq {
field: field.into(),
value: value.into(),
},
})
.unwrap();
}
fn user_meta(locale: &str) -> HashMap<String, String> {
let mut m = HashMap::new();
m.insert("locale".to_string(), locale.to_string());
m
}
// ── Test 1: basic_purge_retracts_score ───────────────────────────────────────
/// Signal from user contributes to cohort; purge reduces the cohort score.
#[test]
fn basic_purge_retracts_score() {
let db = open_ephemeral();
setup_cohort(&db, "en_users", "locale", "en");
let user_id = 1u64;
let item_id = EntityId::new(100);
db.write_user(EntityId::new(user_id), &user_meta("en"))
.unwrap();
let ts = Timestamp::now();
db.signal_with_context("view", item_id, 1.0, ts, Some(user_id), None)
.unwrap();
db.signal_with_context("view", item_id, 1.0, ts, Some(user_id), None)
.unwrap();
// Cohort ledger should have score > 0.
let score_before = db
.cohort_ledger()
.read_decay_score("en_users", item_id, "view", 0)
.unwrap()
.unwrap_or(0.0);
assert!(
score_before > 0.0,
"cohort score must be positive before purge"
);
// Purge user 1's contributions from "en_users".
let (purge_id, manifest) = db.request_community_purge(user_id, "en_users").unwrap();
assert_eq!(manifest.user_id, user_id);
assert_eq!(manifest.cohort, "en_users");
assert_eq!(manifest.purge_id, purge_id);
// Manifest must record the contributions that were retracted.
assert!(
!manifest.entries.is_empty(),
"manifest must list retracted entries"
);
let score_after = db
.cohort_ledger()
.read_decay_score("en_users", item_id, "view", 0)
.unwrap()
.unwrap_or(0.0);
assert!(
score_after < score_before,
"cohort score must decrease after purge (before={score_before}, after={score_after})"
);
assert!(score_after >= 0.0, "score must not go negative");
}
// ── Test 2: second_purge_is_idempotent ────────────────────────────────────────
/// A second purge for the same (user, cohort) finds nothing to retract.
#[test]
fn second_purge_is_idempotent() {
let db = open_ephemeral();
setup_cohort(&db, "en_users", "locale", "en");
let user_id = 2u64;
let item_id = EntityId::new(200);
db.write_user(EntityId::new(user_id), &user_meta("en"))
.unwrap();
let ts = Timestamp::now();
db.signal_with_context("view", item_id, 1.0, ts, Some(user_id), None)
.unwrap();
let (_, manifest1) = db.request_community_purge(user_id, "en_users").unwrap();
assert!(
!manifest1.entries.is_empty(),
"first purge must drain entries"
);
let score_after_first = db
.cohort_ledger()
.read_decay_score("en_users", item_id, "view", 0)
.unwrap()
.unwrap_or(0.0);
let (_, manifest2) = db.request_community_purge(user_id, "en_users").unwrap();
assert_eq!(
manifest2.entries.len(),
0,
"second purge must find no entries to drain"
);
// Score unchanged from first purge.
let score_after_second = db
.cohort_ledger()
.read_decay_score("en_users", item_id, "view", 0)
.unwrap()
.unwrap_or(0.0);
let delta = (score_after_second - score_after_first).abs();
assert!(
delta < 1e-9,
"second purge must not change score (delta={delta})"
);
}
// ── Test 3: purge_does_not_affect_other_users ────────────────────────────────
/// Purging user A's contributions must not retract user B's contributions.
#[test]
fn purge_does_not_affect_other_users() {
let db = open_ephemeral();
setup_cohort(&db, "en_users", "locale", "en");
let user_a = 3u64;
let user_b = 4u64;
let item_id = EntityId::new(300);
db.write_user(EntityId::new(user_a), &user_meta("en"))
.unwrap();
db.write_user(EntityId::new(user_b), &user_meta("en"))
.unwrap();
let ts = Timestamp::now();
// Both users signal the same item.
db.signal_with_context("view", item_id, 1.0, ts, Some(user_a), None)
.unwrap();
db.signal_with_context("view", item_id, 1.0, ts, Some(user_b), None)
.unwrap();
let score_before = db
.cohort_ledger()
.read_decay_score("en_users", item_id, "view", 0)
.unwrap()
.unwrap_or(0.0);
assert!(score_before > 0.0);
// Purge only user A.
let (_, manifest) = db.request_community_purge(user_a, "en_users").unwrap();
assert!(
!manifest.entries.is_empty(),
"manifest must list user A entries"
);
// Score should be lower but still positive (user B's contribution remains).
let score_after = db
.cohort_ledger()
.read_decay_score("en_users", item_id, "view", 0)
.unwrap()
.unwrap_or(0.0);
assert!(
score_after < score_before,
"score must decrease after purging user A"
);
assert!(
score_after > 0.0,
"score must remain positive (user B contribution intact)"
);
}
// ── Test 4: manifest_persisted_and_listable ──────────────────────────────────
/// After purge, manifest appears in `list_purge_manifests`.
/// Ephemeral mode uses an in-memory backend so the manifest is persisted.
#[test]
fn manifest_persisted_and_listable() {
let db = open_ephemeral();
setup_cohort(&db, "en_users", "locale", "en");
let user_id = 5u64;
let item_id = EntityId::new(500);
db.write_user(EntityId::new(user_id), &user_meta("en"))
.unwrap();
let ts = Timestamp::now();
db.signal_with_context("view", item_id, 1.0, ts, Some(user_id), None)
.unwrap();
let (purge_id, _) = db.request_community_purge(user_id, "en_users").unwrap();
// Manifest must appear in the list.
let manifests = db.list_purge_manifests(user_id).unwrap();
assert_eq!(manifests.len(), 1, "one manifest must be persisted");
assert_eq!(
manifests[0].purge_id, purge_id,
"manifest purge_id must match"
);
assert_eq!(manifests[0].user_id, user_id);
assert_eq!(manifests[0].cohort, "en_users");
// A second purge creates a second manifest.
db.request_community_purge(user_id, "en_users").unwrap();
let manifests2 = db.list_purge_manifests(user_id).unwrap();
assert_eq!(manifests2.len(), 2, "two manifests must be persisted");
}
// ── Test 5: purge_unknown_cohort_is_noop ─────────────────────────────────────
/// Purging a user from a cohort they never belonged to returns an empty manifest.
#[test]
fn purge_unknown_cohort_is_noop() {
let db = open_ephemeral();
let user_id = 6u64;
db.write_user(EntityId::new(user_id), &user_meta("en"))
.unwrap();
// No signals were recorded, no cohort was defined.
let (_, manifest) = db
.request_community_purge(user_id, "nonexistent_cohort")
.unwrap();
assert_eq!(
manifest.entries.len(),
0,
"purge of unknown cohort must return empty manifest"
);
}
// ── Test 6: score_never_goes_negative ────────────────────────────────────────
/// Even if the contribution log is replayed twice (bug scenario), the ledger
/// score must be floored at 0.0.
#[test]
fn score_never_goes_negative() {
let db = open_ephemeral();
setup_cohort(&db, "en_users", "locale", "en");
let user_id = 7u64;
let item_id = EntityId::new(700);
db.write_user(EntityId::new(user_id), &user_meta("en"))
.unwrap();
let ts = Timestamp::now();
db.signal_with_context("view", item_id, 1.0, ts, Some(user_id), None)
.unwrap();
// Normal purge.
db.request_community_purge(user_id, "en_users").unwrap();
let score = db
.cohort_ledger()
.read_decay_score("en_users", item_id, "view", 0)
.unwrap()
.unwrap_or(0.0);
assert!(
score >= 0.0,
"cohort score must be non-negative after purge (score={score})"
);
}
// ── Test 7: multiple_items_purged_together ────────────────────────────────────
/// Signals to multiple items by the same user are all retracted in one purge.
#[test]
fn multiple_items_purged_together() {
let db = open_ephemeral();
setup_cohort(&db, "en_users", "locale", "en");
let user_id = 8u64;
let item_a = EntityId::new(801);
let item_b = EntityId::new(802);
db.write_user(EntityId::new(user_id), &user_meta("en"))
.unwrap();
let ts = Timestamp::now();
db.signal_with_context("view", item_a, 1.0, ts, Some(user_id), None)
.unwrap();
db.signal_with_context("view", item_b, 1.0, ts, Some(user_id), None)
.unwrap();
let score_a_before = db
.cohort_ledger()
.read_decay_score("en_users", item_a, "view", 0)
.unwrap()
.unwrap_or(0.0);
let score_b_before = db
.cohort_ledger()
.read_decay_score("en_users", item_b, "view", 0)
.unwrap()
.unwrap_or(0.0);
assert!(score_a_before > 0.0);
assert!(score_b_before > 0.0);
let (_, manifest) = db.request_community_purge(user_id, "en_users").unwrap();
// Both items should appear in the manifest.
assert_eq!(
manifest.entries.len(),
2,
"manifest must cover all items signalled by the user"
);
let score_a_after = db
.cohort_ledger()
.read_decay_score("en_users", item_a, "view", 0)
.unwrap()
.unwrap_or(0.0);
let score_b_after = db
.cohort_ledger()
.read_decay_score("en_users", item_b, "view", 0)
.unwrap()
.unwrap_or(0.0);
assert!(
score_a_after < score_a_before,
"item A score must decrease after purge"
);
assert!(
score_b_after < score_b_before,
"item B score must decrease after purge"
);
}
// ── Test 8: only_signaling_cohort_member_contributes ─────────────────────────
/// A user who is not in a cohort does not appear in the contribution log for
/// that cohort, so purging them is a no-op on the ledger.
#[test]
fn non_cohort_member_purge_is_noop() {
let db = open_ephemeral();
setup_cohort(&db, "en_users", "locale", "en");
let en_user = 9u64;
let fr_user = 10u64;
let item_id = EntityId::new(900);
db.write_user(EntityId::new(en_user), &user_meta("en"))
.unwrap();
db.write_user(EntityId::new(fr_user), &user_meta("fr"))
.unwrap();
let ts = Timestamp::now();
// Both users signal; only en_user lands in the cohort.
db.signal_with_context("view", item_id, 1.0, ts, Some(en_user), None)
.unwrap();
db.signal_with_context("view", item_id, 1.0, ts, Some(fr_user), None)
.unwrap();
let score_before = db
.cohort_ledger()
.read_decay_score("en_users", item_id, "view", 0)
.unwrap()
.unwrap_or(0.0);
assert!(score_before > 0.0, "en_users must have score from en_user");
// Purge fr_user from en_users (they were never in the log for this cohort).
let (_, manifest) = db.request_community_purge(fr_user, "en_users").unwrap();
assert_eq!(
manifest.entries.len(),
0,
"fr_user has no contributions to en_users"
);
let score_after = db
.cohort_ledger()
.read_decay_score("en_users", item_id, "view", 0)
.unwrap()
.unwrap_or(0.0);
let delta = (score_after - score_before).abs();
assert!(
delta < 1e-9,
"score must be unchanged when purging non-member (delta={delta})"
);
}