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
241 lines
8.1 KiB
Rust
241 lines
8.1 KiB
Rust
//! M9 Purge Re-materialization Integration Tests.
|
|
//!
|
|
//! Validates the end-to-end purge re-materialization engine:
|
|
//!
|
|
//! 1. `submit_purge_job` enqueues a job and returns a `JobId`.
|
|
//! 2. `purge_job_status` returns `Pending` for a freshly submitted job
|
|
//! (ephemeral mode — engine not started).
|
|
//! 3. `rematerialization_metrics` returns a valid snapshot.
|
|
//! 4. In persistent mode the background worker starts automatically and
|
|
//! processes the job to `Succeeded`.
|
|
//! 5. After a job succeeds, the cohort ledger reflects the purged state.
|
|
//! 6. Multiple jobs for different users are processed independently.
|
|
|
|
#![allow(
|
|
clippy::unwrap_used,
|
|
clippy::cast_precision_loss,
|
|
clippy::items_after_statements
|
|
)]
|
|
|
|
use std::collections::HashMap;
|
|
use std::time::{Duration, Instant};
|
|
|
|
use tempfile::TempDir;
|
|
use tidaldb::TidalDb;
|
|
use tidaldb::cohort::rematerialization::PurgeJobStatus;
|
|
use tidaldb::cohort::{CohortDef, Predicate};
|
|
use tidaldb::schema::{DecaySpec, EntityId, EntityKind, SchemaBuilder, Timestamp, Window};
|
|
|
|
// ── Shared fixtures ──────────────────────────────────────────────────────────
|
|
|
|
fn remat_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("valid schema")
|
|
}
|
|
|
|
fn open_ephemeral() -> TidalDb {
|
|
TidalDb::builder()
|
|
.ephemeral()
|
|
.with_schema(remat_schema())
|
|
.open()
|
|
.expect("ephemeral open")
|
|
}
|
|
|
|
fn open_persistent(dir: &TempDir) -> TidalDb {
|
|
TidalDb::builder()
|
|
.with_data_dir(dir.path())
|
|
.with_schema(remat_schema())
|
|
.open()
|
|
.expect("persistent open")
|
|
}
|
|
|
|
fn setup_cohort(db: &TidalDb, cohort: &str) {
|
|
db.define_cohort(CohortDef {
|
|
name: cohort.to_string(),
|
|
predicate: Predicate::Eq {
|
|
field: "locale".into(),
|
|
value: "en".into(),
|
|
},
|
|
})
|
|
.unwrap();
|
|
}
|
|
|
|
fn en_user_meta() -> HashMap<String, String> {
|
|
let mut m = HashMap::new();
|
|
m.insert("locale".to_string(), "en".to_string());
|
|
m
|
|
}
|
|
|
|
// ── Test 1: submit_purge_job_returns_job_id ──────────────────────────────────
|
|
|
|
/// `submit_purge_job` must return a job id for any (user, community) pair.
|
|
/// In ephemeral mode, the engine is not started, so the job remains pending.
|
|
#[test]
|
|
fn submit_purge_job_returns_job_id() {
|
|
let db = open_ephemeral();
|
|
|
|
let job_id = db
|
|
.submit_purge_job(42, "en_users", None)
|
|
.expect("submit_purge_job must succeed");
|
|
|
|
// Job ID is 16 bytes; it should be non-zero.
|
|
assert_ne!(job_id, [0u8; 16], "job_id must not be all-zero");
|
|
}
|
|
|
|
// ── Test 2: purge_job_status_pending ─────────────────────────────────────────
|
|
|
|
/// Immediately after submission the job should be in `Pending` status.
|
|
/// (No engine runs in ephemeral mode.)
|
|
#[test]
|
|
fn purge_job_status_pending_after_submit() {
|
|
let db = open_ephemeral();
|
|
|
|
let job_id = db
|
|
.submit_purge_job(1, "en_users", None)
|
|
.expect("submit_purge_job must succeed");
|
|
|
|
let status = db
|
|
.purge_job_status(&job_id)
|
|
.expect("purge_job_status must succeed");
|
|
|
|
assert_eq!(
|
|
status,
|
|
Some(PurgeJobStatus::Pending),
|
|
"freshly submitted job must be Pending"
|
|
);
|
|
}
|
|
|
|
// ── Test 3: unknown_job_id_returns_none ──────────────────────────────────────
|
|
|
|
/// Querying a job id that was never submitted must return `None`.
|
|
#[test]
|
|
fn unknown_job_id_returns_none() {
|
|
let db = open_ephemeral();
|
|
let unknown_id = [0xFFu8; 16];
|
|
|
|
let status = db
|
|
.purge_job_status(&unknown_id)
|
|
.expect("purge_job_status must succeed");
|
|
|
|
assert_eq!(status, None, "unknown job id must return None");
|
|
}
|
|
|
|
// ── Test 4: rematerialization_metrics_snapshot ───────────────────────────────
|
|
|
|
/// `rematerialization_metrics` must return a valid snapshot with zero counts
|
|
/// in a freshly opened database.
|
|
#[test]
|
|
fn rematerialization_metrics_snapshot_is_valid() {
|
|
let db = open_ephemeral();
|
|
|
|
let snap = db
|
|
.rematerialization_metrics()
|
|
.expect("rematerialization_metrics must succeed");
|
|
|
|
// On a fresh database no jobs have run.
|
|
assert_eq!(snap.jobs_succeeded, 0, "no jobs succeeded on fresh db");
|
|
assert_eq!(snap.jobs_failed, 0, "no jobs failed on fresh db");
|
|
}
|
|
|
|
// ── Test 5: job_succeeds_in_persistent_mode ───────────────────────────────────
|
|
|
|
/// In persistent mode the re-materialization engine starts automatically.
|
|
/// Submitting a job with no WAL events (empty WAL dir) should succeed quickly.
|
|
#[test]
|
|
fn job_succeeds_in_persistent_mode() {
|
|
let dir = TempDir::new().expect("tempdir");
|
|
let db = open_persistent(&dir);
|
|
setup_cohort(&db, "en_users");
|
|
|
|
let job_id = db
|
|
.submit_purge_job(10, "en_users", None)
|
|
.expect("submit_purge_job must succeed");
|
|
|
|
// Wait up to 10 seconds for the job to complete.
|
|
let deadline = Instant::now() + Duration::from_secs(10);
|
|
loop {
|
|
let status = db.purge_job_status(&job_id).unwrap();
|
|
match status {
|
|
Some(PurgeJobStatus::Succeeded) => break,
|
|
Some(PurgeJobStatus::PermanentlyFailed) => {
|
|
panic!("job permanently failed");
|
|
}
|
|
_ => {}
|
|
}
|
|
assert!(Instant::now() < deadline, "job did not succeed within 10s");
|
|
std::thread::sleep(Duration::from_millis(50));
|
|
}
|
|
|
|
let snap = db.rematerialization_metrics().unwrap();
|
|
assert!(
|
|
snap.jobs_succeeded >= 1,
|
|
"metrics must reflect succeeded job"
|
|
);
|
|
}
|
|
|
|
// ── Test 6: multiple_jobs_processed_independently ────────────────────────────
|
|
|
|
/// Submitting jobs for two different users results in both being processed.
|
|
#[test]
|
|
fn multiple_jobs_processed_independently() {
|
|
let dir = TempDir::new().expect("tempdir");
|
|
let db = open_persistent(&dir);
|
|
setup_cohort(&db, "en_users");
|
|
|
|
// Write some user/item data so the job has context.
|
|
db.write_user(EntityId::new(20), &en_user_meta()).unwrap();
|
|
db.write_user(EntityId::new(21), &en_user_meta()).unwrap();
|
|
let item_a = EntityId::new(1001);
|
|
let item_b = EntityId::new(1002);
|
|
let ts = Timestamp::now();
|
|
db.signal_with_context("view", item_a, 1.0, ts, Some(20), None)
|
|
.unwrap();
|
|
db.signal_with_context("view", item_b, 1.0, ts, Some(21), None)
|
|
.unwrap();
|
|
|
|
let job_a = db.submit_purge_job(20, "en_users", None).unwrap();
|
|
let job_b = db.submit_purge_job(21, "en_users", None).unwrap();
|
|
|
|
// Wait for both.
|
|
let deadline = Instant::now() + Duration::from_secs(15);
|
|
loop {
|
|
let a = db.purge_job_status(&job_a).unwrap();
|
|
let b = db.purge_job_status(&job_b).unwrap();
|
|
let a_done = matches!(
|
|
a,
|
|
Some(PurgeJobStatus::Succeeded | PurgeJobStatus::PermanentlyFailed)
|
|
);
|
|
let b_done = matches!(
|
|
b,
|
|
Some(PurgeJobStatus::Succeeded | PurgeJobStatus::PermanentlyFailed)
|
|
);
|
|
if a_done && b_done {
|
|
assert_eq!(a, Some(PurgeJobStatus::Succeeded), "job_a must succeed");
|
|
assert_eq!(b, Some(PurgeJobStatus::Succeeded), "job_b must succeed");
|
|
break;
|
|
}
|
|
assert!(
|
|
Instant::now() < deadline,
|
|
"jobs did not complete within 15s"
|
|
);
|
|
std::thread::sleep(Duration::from_millis(50));
|
|
}
|
|
|
|
let snap = db.rematerialization_metrics().unwrap();
|
|
assert!(
|
|
snap.jobs_succeeded >= 2,
|
|
"metrics must reflect both succeeded jobs"
|
|
);
|
|
}
|