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
430 lines
14 KiB
Rust
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})"
|
|
);
|
|
}
|