use crate::checkpoint::Checkpoint;
use crate::fold::{FoldedRecord, SyncState};
use crate::journal::OplogJournal;
use crate::oplog::{verify_log, ChainError, Hlc, OpRecord};
use serde::{Deserialize, Serialize};
use serde_json::{json, Value};
use std::collections::{BTreeMap, BTreeSet};
use std::fmt;
use std::fs::{self, File};
use std::io::Write;
use std::path::{Path, PathBuf};
pub const RUNS_MAX_PER_AGENT: usize = 50;
pub const RUNS_MAX_AGE_MS: u64 = 30 * 24 * 60 * 60 * 1000;
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum RetentionRule {
KeepAll,
LastN { n: usize },
MaxAgeMs { max_age_ms: u64 },
PerAgentWithAge {
max_per_agent: usize,
max_age_ms: u64,
},
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct RetentionPolicy {
pub rules: BTreeMap<String, RetentionRule>,
pub default_rule: RetentionRule,
}
impl Default for RetentionPolicy {
fn default() -> Self {
Self::keep_all()
}
}
impl RetentionPolicy {
pub fn keep_all() -> Self {
Self {
rules: BTreeMap::new(),
default_rule: RetentionRule::KeepAll,
}
}
pub fn proposal_default(conversation_last_n: usize, trajectory_max_age_ms: u64) -> Self {
let mut rules = BTreeMap::new();
rules.insert(
"conversation".to_string(),
RetentionRule::LastN {
n: conversation_last_n,
},
);
rules.insert(
"run".to_string(),
RetentionRule::PerAgentWithAge {
max_per_agent: RUNS_MAX_PER_AGENT,
max_age_ms: RUNS_MAX_AGE_MS,
},
);
rules.insert(
"trajectory".to_string(),
RetentionRule::MaxAgeMs {
max_age_ms: trajectory_max_age_ms,
},
);
rules.insert("knowledge".to_string(), RetentionRule::KeepAll);
rules.insert("skill".to_string(), RetentionRule::KeepAll);
rules.insert("routing".to_string(), RetentionRule::KeepAll);
Self {
rules,
default_rule: RetentionRule::KeepAll,
}
}
pub fn rule_for(&self, surface_tag: &str) -> &RetentionRule {
self.rules.get(surface_tag).unwrap_or(&self.default_rule)
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct RetentionReport {
pub dropped: BTreeMap<String, usize>,
pub tombstoned: Vec<(String, String)>,
}
fn value_ts(payload: &Value) -> Option<u64> {
payload
.get("timestamp")
.and_then(|v| v.as_u64().or_else(|| v.as_f64().map(|f| f as u64)))
}
fn payload_ts(record: &FoldedRecord) -> Option<u64> {
value_ts(&record.payload)
}
pub fn as_of_from_ops(ops: &[OpRecord]) -> u64 {
ops.iter()
.filter_map(|op| value_ts(&op.payload))
.max()
.unwrap_or(0)
}
fn within_age(record: &FoldedRecord, as_of_ms: u64, max_age_ms: u64) -> bool {
match payload_ts(record) {
Some(ts) => as_of_ms.saturating_sub(ts) <= max_age_ms,
None => true,
}
}
fn recency_sorted(entries: &BTreeMap<String, FoldedRecord>) -> Vec<(&String, &FoldedRecord)> {
let mut sorted: Vec<(&String, &FoldedRecord)> = entries.iter().collect();
sorted.sort_by(|(_, a), (_, b)| {
(payload_ts(a).unwrap_or(u64::MAX), &a.hlc, &a.op_id).cmp(&(
payload_ts(b).unwrap_or(u64::MAX),
&b.hlc,
&b.op_id,
))
});
sorted
}
fn select_retained(
entries: &BTreeMap<String, FoldedRecord>,
rule: &RetentionRule,
as_of_ms: u64,
) -> BTreeSet<String> {
match rule {
RetentionRule::KeepAll => entries.keys().cloned().collect(),
RetentionRule::LastN { n } => recency_sorted(entries)
.into_iter()
.rev()
.take(*n)
.map(|(k, _)| k.clone())
.collect(),
RetentionRule::MaxAgeMs { max_age_ms } => entries
.iter()
.filter(|(_, r)| within_age(r, as_of_ms, *max_age_ms))
.map(|(k, _)| k.clone())
.collect(),
RetentionRule::PerAgentWithAge {
max_per_agent,
max_age_ms,
} => {
let mut per_agent_rank: BTreeMap<&str, usize> = BTreeMap::new();
let mut keep = BTreeSet::new();
for (key, record) in recency_sorted(entries).into_iter().rev() {
let agent = record
.payload
.get("agent_id")
.and_then(Value::as_str)
.unwrap_or("");
let rank = per_agent_rank.entry(agent).or_insert(0);
let over_count = *rank >= *max_per_agent;
*rank += 1;
if !over_count && within_age(record, as_of_ms, *max_age_ms) {
keep.insert(key.clone());
}
}
keep
}
}
}
fn entity_id<'a>(key: &'a str, record: &'a FoldedRecord) -> Option<&'a str> {
key.strip_prefix("id:")
.or_else(|| record.payload.get("id").and_then(Value::as_str))
}
pub fn is_tombstone(record: &FoldedRecord) -> bool {
record
.payload
.get("tombstone")
.and_then(Value::as_bool)
.unwrap_or(false)
}
fn tombstone_of(record: &FoldedRecord, id: &str) -> FoldedRecord {
FoldedRecord {
op_id: record.op_id.clone(),
hlc: record.hlc.clone(),
payload: json!({"id": id, "tombstone": true}),
}
}
pub fn apply_retention(
state: &SyncState,
policy: &RetentionPolicy,
as_of_ms: u64,
) -> Result<(SyncState, RetentionReport), CompactError> {
let routing_tag = crate::oplog::Surface::Routing.tag();
debug_assert!(crate::oplog::Surface::Routing.is_replay_stream());
if state.logs.contains_key(&routing_tag)
&& policy.rule_for(&routing_tag) != &RetentionRule::KeepAll
{
return Err(CompactError::EventStreamRetention {
surface: routing_tag,
});
}
if policy.rule_for(&crate::oplog::Surface::Intent.tag()) != &RetentionRule::KeepAll {
return Err(CompactError::IntentRetention);
}
let mut retained = state.clone();
let mut report = RetentionReport::default();
for (tag, entries) in &state.logs {
let rule = policy.rule_for(tag);
if rule == &RetentionRule::KeepAll {
continue;
}
let live: BTreeMap<String, FoldedRecord> = entries
.iter()
.filter(|(_, record)| !is_tombstone(record))
.map(|(key, record)| (key.clone(), record.clone()))
.collect();
let keep = select_retained(&live, rule, as_of_ms);
let surface = retained.logs.get_mut(tag).expect("cloned from state");
for (key, record) in &live {
if keep.contains(key) {
continue;
}
match entity_id(key, record) {
Some(id) => {
surface.insert(key.clone(), tombstone_of(record, id));
report.tombstoned.push((tag.clone(), key.clone()));
}
None => {
surface.remove(key);
*report.dropped.entry(tag.clone()).or_insert(0) += 1;
}
}
}
}
report.tombstoned.sort();
Ok((retained, report))
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct AckTable {
acked: BTreeMap<String, Hlc>,
}
impl AckTable {
pub fn new() -> Self {
Self::default()
}
pub fn ack(&mut self, device_id: impl Into<String>, frontier: Hlc) -> bool {
let device_id = device_id.into();
match self.acked.get(&device_id) {
Some(current) if frontier <= *current => false,
_ => {
self.acked.insert(device_id, frontier);
true
}
}
}
pub fn get(&self, device_id: &str) -> Option<&Hlc> {
self.acked.get(device_id)
}
pub fn devices(&self) -> impl Iterator<Item = &str> {
self.acked.keys().map(String::as_str)
}
pub fn stable_frontier(&self) -> Option<&Hlc> {
self.acked.values().min()
}
pub fn save(&self, path: &Path) -> std::io::Result<()> {
if let Some(parent) = path.parent() {
if !parent.as_os_str().is_empty() {
fs::create_dir_all(parent)?;
}
}
let tmp_path = {
let mut s = path.as_os_str().to_owned();
s.push(".tmp");
PathBuf::from(s)
};
{
let mut tmp = File::create(&tmp_path)?;
tmp.write_all(
serde_json::to_string(self)
.map_err(std::io::Error::other)?
.as_bytes(),
)?;
tmp.sync_all()?;
}
fs::rename(&tmp_path, path)
}
pub fn load(path: &Path) -> std::io::Result<Self> {
match fs::read_to_string(path) {
Ok(raw) => serde_json::from_str(&raw).map_err(std::io::Error::other),
Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(Self::default()),
Err(e) => Err(e),
}
}
}
#[derive(Debug)]
pub enum CompactError {
Chain(ChainError),
NothingAcked,
UnackedDevice {
device_id: String,
},
EventStreamRetention {
surface: String,
},
IntentRetention,
TruncatedJournal {
checkpoint_hash: String,
},
Io(std::io::Error),
}
impl fmt::Display for CompactError {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
CompactError::Chain(e) => write!(f, "compaction refused: log does not verify: {e}"),
CompactError::NothingAcked => {
write!(
f,
"compaction refused: no acked frontier exists (empty ack table)"
)
}
CompactError::UnackedDevice { device_id } => write!(
f,
"compaction refused: device {device_id} appears in the log but has no acked \
frontier — its ops may include state no other replica has folded"
),
CompactError::EventStreamRetention { surface } => write!(
f,
"retention policy for event-stream surface {surface} must be keep_all — \
observation multisets replay from genesis and cannot be trimmed"
),
CompactError::IntentRetention => write!(
f,
"retention policy for the leased `intent` surface must be keep_all — the \
committed-run idempotency oracle must survive compaction and cannot be trimmed"
),
CompactError::TruncatedJournal { checkpoint_hash } => write!(
f,
"compaction refused: journal already truncated below checkpoint \
{checkpoint_hash} — re-planning from the tail alone would drop that \
checkpoint's state (recompaction over a checkpoint base is a later slice)"
),
CompactError::Io(e) => write!(f, "compaction io error: {e}"),
}
}
}
impl std::error::Error for CompactError {}
#[derive(Debug)]
pub struct CompactionPlan {
pub checkpoint: Checkpoint,
pub retained_ops: Vec<OpRecord>,
pub dropped_ops: usize,
pub frontier: Hlc,
pub as_of_ms: u64,
pub retention: RetentionReport,
}
pub fn plan_compaction(
ops: &[OpRecord],
acks: &AckTable,
policy: &RetentionPolicy,
as_of_ms: Option<u64>,
) -> Result<CompactionPlan, CompactError> {
verify_log(ops).map_err(CompactError::Chain)?;
let frontier = acks
.stable_frontier()
.cloned()
.ok_or(CompactError::NothingAcked)?;
for op in ops {
if acks.get(&op.device_id).is_none() {
return Err(CompactError::UnackedDevice {
device_id: op.device_id.clone(),
});
}
}
let (below, retained_ops): (Vec<OpRecord>, Vec<OpRecord>) =
ops.iter().cloned().partition(|op| op.hlc <= frontier);
let as_of_ms = as_of_ms.unwrap_or_else(|| as_of_from_ops(&below));
let exact = Checkpoint::from_ops(&below).map_err(CompactError::Chain)?;
let (retained_state, retention) = apply_retention(&exact.state, policy, as_of_ms)?;
let checkpoint = Checkpoint::assemble(exact.frontier, exact.scopes, retained_state);
Ok(CompactionPlan {
checkpoint,
retained_ops,
dropped_ops: below.len(),
frontier,
as_of_ms,
retention,
})
}
#[derive(Debug)]
pub struct CompactionOutcome {
pub checkpoint_path: Option<PathBuf>,
pub plan: CompactionPlan,
}
pub fn compact_and_truncate(
journal: &mut OplogJournal,
checkpoint_dir: &Path,
acks: &AckTable,
policy: &RetentionPolicy,
as_of_ms: Option<u64>,
) -> Result<CompactionOutcome, CompactError> {
let (marker, ops) = OplogJournal::load_with_marker(journal.path()).map_err(CompactError::Io)?;
if let Some(marker) = marker {
return Err(CompactError::TruncatedJournal {
checkpoint_hash: marker.checkpoint_hash,
});
}
let plan = plan_compaction(&ops, acks, policy, as_of_ms)?;
if plan.dropped_ops == 0 {
return Ok(CompactionOutcome {
checkpoint_path: None,
plan,
});
}
let checkpoint_path = plan
.checkpoint
.save(checkpoint_dir)
.map_err(CompactError::Io)?;
journal
.truncate_to(&plan.retained_ops, &plan.checkpoint.checkpoint_hash)
.map_err(CompactError::Io)?;
Ok(CompactionOutcome {
checkpoint_path: Some(checkpoint_path),
plan,
})
}
#[cfg(test)]
mod tests {
use super::*;
use crate::fold::fold;
use crate::oplog::{DeviceLog, Scope, Surface};
use serde_json::json;
fn hlc(wall_ms: u64, device: &str) -> Hlc {
Hlc {
wall_ms,
counter: 0,
device_id: device.into(),
}
}
#[test]
fn ack_table_is_monotone_only_and_min_frontier() {
let mut acks = AckTable::new();
assert_eq!(acks.stable_frontier(), None);
assert!(acks.ack("a", hlc(5, "a")));
assert!(acks.ack("b", hlc(9, "b")));
assert_eq!(
acks.stable_frontier(),
Some(&hlc(5, "a")),
"min over devices"
);
assert!(!acks.ack("b", hlc(3, "b")));
assert!(!acks.ack("b", hlc(9, "b")), "equal is not an advance");
assert_eq!(acks.get("b"), Some(&hlc(9, "b")));
assert!(acks.ack("b", hlc(12, "b")));
assert_eq!(acks.get("b"), Some(&hlc(12, "b")));
}
#[test]
fn ack_table_persists_atomically_and_loads_missing_as_empty() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("nested").join("acks.json");
assert_eq!(
AckTable::load(&path).unwrap(),
AckTable::new(),
"missing → empty"
);
let mut acks = AckTable::new();
acks.ack("a", hlc(5, "a"));
acks.ack("b", hlc(9, "b"));
acks.save(&path).unwrap();
assert_eq!(AckTable::load(&path).unwrap(), acks);
assert!(!path.parent().unwrap().join("acks.json.tmp").exists());
}
fn simple_ops() -> Vec<OpRecord> {
let mut a = DeviceLog::new("a");
let mut b = DeviceLog::new("b");
let mut ops = vec![
a.append(Scope::Personal, Surface::Knowledge, json!({"id": "f1"})),
a.append(Scope::Personal, Surface::Knowledge, json!({"id": "f2"})),
a.append(Scope::Personal, Surface::Knowledge, json!({"id": "f3"})),
];
for op in &ops {
b.observe(&op.hlc);
}
ops.push(b.append(Scope::Personal, Surface::Knowledge, json!({"id": "f4"})));
ops
}
#[test]
fn compaction_refuses_without_acks() {
let ops = simple_ops();
let policy = RetentionPolicy::keep_all();
assert!(matches!(
plan_compaction(&ops, &AckTable::new(), &policy, None),
Err(CompactError::NothingAcked)
));
let mut acks = AckTable::new();
acks.ack("a", ops[2].hlc.clone());
assert!(matches!(
plan_compaction(&ops, &acks, &policy, None),
Err(CompactError::UnackedDevice { .. })
));
}
#[test]
fn compaction_never_drops_above_a_lagging_ack() {
let ops = simple_ops();
let mut acks = AckTable::new();
acks.ack("b", ops[3].hlc.clone());
acks.ack("a", ops[1].hlc.clone());
let plan = plan_compaction(&ops, &acks, &RetentionPolicy::keep_all(), None).unwrap();
assert_eq!(
plan.frontier, ops[1].hlc,
"stable frontier = the lagging device's ack"
);
assert_eq!(plan.dropped_ops, 2, "only ops ≤ the lagging frontier drop");
assert_eq!(plan.retained_ops.len(), 2);
assert!(
plan.retained_ops.iter().all(|op| op.hlc > plan.frontier),
"everything a device hasn't seen stays in the journal"
);
}
#[test]
fn retention_conversations_last_n_by_timestamp() {
let mut dev = DeviceLog::new("a");
let ops: Vec<OpRecord> = (0..5)
.map(|i| {
dev.append(
Scope::Personal,
Surface::Conversation,
json!({"speaker": "u", "text": format!("t{i}"), "timestamp": 100 + i}),
)
})
.collect();
let state = fold(&ops);
let policy = RetentionPolicy::proposal_default(2, u64::MAX);
let (retained, report) = apply_retention(&state, &policy, 1_000).unwrap();
let tag = Surface::Conversation.tag();
let texts: Vec<String> = retained
.log_entries(&tag)
.iter()
.map(|r| r.payload["text"].as_str().unwrap().to_string())
.collect();
assert_eq!(texts, vec!["t3", "t4"], "last 2 turns by timestamp survive");
assert_eq!(report.dropped[&tag], 3);
}
#[test]
fn retention_runs_per_agent_and_age_matches_run_store_gc() {
assert_eq!(
RUNS_MAX_PER_AGENT, 50,
"parity with run_store DEFAULT_MAX_RUNS_PER_AGENT"
);
assert_eq!(
RUNS_MAX_AGE_MS,
30 * 24 * 60 * 60 * 1000,
"parity with DEFAULT_MAX_AGE_DAYS"
);
let mut dev = DeviceLog::new("a");
let mut ops = Vec::new();
for (id, agent, ts) in [
("r1", "milo", 10u64), ("r2", "milo", 60), ("r3", "milo", 70), ("r4", "other", 10), ("r5", "other", 80), ] {
ops.push(dev.append(
Scope::Personal,
Surface::Run,
json!({"id": id, "agent_id": agent, "timestamp": ts}),
));
}
let state = fold(&ops);
let mut policy = RetentionPolicy::keep_all();
policy.rules.insert(
"run".to_string(),
RetentionRule::PerAgentWithAge {
max_per_agent: 2,
max_age_ms: 50,
},
);
let (retained, report) = apply_retention(&state, &policy, 100).unwrap();
let entries = retained.log_entries(&Surface::Run.tag());
let kept: Vec<&str> = entries
.iter()
.filter(|r| !is_tombstone(r))
.map(|r| r.payload["id"].as_str().unwrap())
.collect();
assert_eq!(kept, vec!["r2", "r3", "r5"]);
let stubs: Vec<&str> = entries
.iter()
.filter(|r| is_tombstone(r))
.map(|r| r.payload["id"].as_str().unwrap())
.collect();
assert_eq!(stubs, vec!["r1", "r4"]);
assert_eq!(report.tombstoned.len(), 2);
assert!(
report.dropped.is_empty(),
"id-bearing entries are stubbed, never erased"
);
}
#[test]
fn retention_trajectories_by_age_and_undated_never_dropped() {
let mut dev = DeviceLog::new("a");
let ops = vec![
dev.append(
Scope::Personal,
Surface::Trajectory,
json!({"id": "old", "timestamp": 10}),
),
dev.append(
Scope::Personal,
Surface::Trajectory,
json!({"id": "new", "timestamp": 90}),
),
dev.append(
Scope::Personal,
Surface::Trajectory,
json!({"id": "undated"}),
),
];
let state = fold(&ops);
let mut policy = RetentionPolicy::keep_all();
policy.rules.insert(
"trajectory".to_string(),
RetentionRule::MaxAgeMs { max_age_ms: 30 },
);
let (retained, _) = apply_retention(&state, &policy, 100).unwrap();
let entries = retained.log_entries(&Surface::Trajectory.tag());
let kept: Vec<&str> = entries
.iter()
.filter(|r| !is_tombstone(r))
.map(|r| r.payload["id"].as_str().unwrap())
.collect();
assert!(kept.contains(&"new"));
assert!(
kept.contains(&"undated"),
"undated data is never silently age-dropped"
);
assert!(!kept.contains(&"old"));
assert!(entries
.iter()
.any(|r| is_tombstone(r) && r.payload["id"] == json!("old")));
}
#[test]
fn knowledge_and_skills_keep_all_under_the_proposal_default() {
let mut dev = DeviceLog::new("a");
let ops = vec![
dev.append(
Scope::Personal,
Surface::Knowledge,
json!({"id": "f1", "timestamp": 1}),
),
dev.append(
Scope::Personal,
Surface::Skill,
json!({"id": "s1", "timestamp": 1}),
),
];
let state = fold(&ops);
let (retained, report) =
apply_retention(&state, &RetentionPolicy::proposal_default(1, 1), u64::MAX).unwrap();
assert_eq!(retained.logs[&Surface::Knowledge.tag()].len(), 1);
assert_eq!(retained.logs[&Surface::Skill.tag()].len(), 1);
assert!(report.dropped.is_empty());
}
#[test]
fn event_stream_retention_is_rejected() {
let mut dev = DeviceLog::new("a");
let ops = vec![
dev.append(Scope::Personal, Surface::Routing, json!({"sample": 1.0})),
dev.append(Scope::Personal, Surface::Routing, json!({"sample": 0.0})),
];
let state = fold(&ops);
let mut policy = RetentionPolicy::keep_all();
policy
.rules
.insert("routing".to_string(), RetentionRule::LastN { n: 1 });
assert!(matches!(
apply_retention(&state, &policy, 0),
Err(CompactError::EventStreamRetention { .. })
));
let (retained, _) =
apply_retention(&state, &RetentionPolicy::proposal_default(10, 10), u64::MAX).unwrap();
assert_eq!(retained.log_entries(&Surface::Routing.tag()).len(), 2);
}
#[test]
fn every_dropped_id_bearing_entry_leaves_a_tombstone_stub() {
let mut dev = DeviceLog::new("a");
let ops = vec![
dev.append(
Scope::Personal,
Surface::Knowledge,
json!({"id": "f1", "timestamp": 1}),
),
dev.append(
Scope::Personal,
Surface::Knowledge,
json!({"id": "f2", "timestamp": 2, "supersedes": "f1"}),
),
dev.append(
Scope::Personal,
Surface::Knowledge,
json!({"id": "f3", "timestamp": 3, "supersedes": ["f2"]}),
),
dev.append(
Scope::Personal,
Surface::Knowledge,
json!({"id": "f0", "timestamp": 0}),
),
];
let state = fold(&ops);
let mut policy = RetentionPolicy::keep_all();
policy
.rules
.insert("knowledge".to_string(), RetentionRule::LastN { n: 1 });
let (retained, report) = apply_retention(&state, &policy, 10).unwrap();
let tag = Surface::Knowledge.tag();
let surface = &retained.logs[&tag];
assert_eq!(
surface["id:f3"].payload["timestamp"],
json!(3),
"newest stays live"
);
for id in ["f0", "f1", "f2"] {
let stub = &surface[&format!("id:{id}")];
assert!(is_tombstone(stub), "{id} left a stub");
assert_eq!(
stub.payload,
json!({"id": id, "tombstone": true}),
"minimal stub shape"
);
assert_eq!(
stub.op_id,
state.logs[&tag][&format!("id:{id}")].op_id,
"stub keeps the original op identity"
);
}
assert_eq!(report.tombstoned.len(), 3);
assert!(report.dropped.is_empty());
let (again, report2) = apply_retention(&retained, &policy, 10).unwrap();
assert_eq!(again, retained);
assert!(report2.tombstoned.is_empty());
}
#[test]
fn derived_as_of_is_the_max_below_frontier_timestamp() {
let mut dev = DeviceLog::new("a");
let ops = vec![
dev.append(
Scope::Personal,
Surface::Trajectory,
json!({"id": "t1", "timestamp": 40}),
),
dev.append(
Scope::Personal,
Surface::Trajectory,
json!({"id": "t2", "timestamp": 100}),
),
dev.append(Scope::Personal, Surface::Skill, json!({"id": "s1"})), ];
assert_eq!(as_of_from_ops(&ops), 100, "max payload timestamp");
assert_eq!(
as_of_from_ops(&ops[2..]),
0,
"no timestamps → 0 (age rules drop nothing)"
);
let mut acks = AckTable::new();
acks.ack("a", ops[2].hlc.clone());
let mut policy = RetentionPolicy::keep_all();
policy.rules.insert(
"trajectory".to_string(),
RetentionRule::MaxAgeMs { max_age_ms: 30 },
);
let plan = plan_compaction(&ops, &acks, &policy, None).unwrap();
assert_eq!(plan.as_of_ms, 100);
let surface = &plan.checkpoint.state.logs[&Surface::Trajectory.tag()];
assert!(is_tombstone(&surface["id:t1"]));
assert!(!is_tombstone(&surface["id:t2"]));
}
#[test]
fn empty_below_frontier_compaction_is_a_no_op() {
let dir = tempfile::tempdir().unwrap();
let journal_path = dir.path().join("oplog.jsonl");
let ckpt_dir = dir.path().join("checkpoints");
let ops = simple_ops();
let mut journal = OplogJournal::open(&journal_path).unwrap();
for op in &ops {
journal.append(op).unwrap();
}
let mut acks = AckTable::new();
acks.ack(
"a",
Hlc {
wall_ms: 0,
counter: 0,
device_id: "a".into(),
},
);
acks.ack(
"b",
Hlc {
wall_ms: 0,
counter: 0,
device_id: "b".into(),
},
);
let before = fs::read_to_string(&journal_path).unwrap();
let outcome = compact_and_truncate(
&mut journal,
&ckpt_dir,
&acks,
&RetentionPolicy::keep_all(),
None,
)
.unwrap();
assert_eq!(outcome.plan.dropped_ops, 0);
assert!(
outcome.checkpoint_path.is_none(),
"no empty checkpoint file written"
);
assert!(!ckpt_dir.exists(), "checkpoint dir not even created");
assert_eq!(
fs::read_to_string(&journal_path).unwrap(),
before,
"journal untouched (no marker, no rewrite)"
);
assert_eq!(OplogJournal::load(&journal_path).unwrap(), ops);
}
#[test]
fn recompacting_an_already_truncated_journal_is_refused() {
let dir = tempfile::tempdir().unwrap();
let journal_path = dir.path().join("oplog.jsonl");
let ckpt_dir = dir.path().join("checkpoints");
let ops = simple_ops();
let mut journal = OplogJournal::open(&journal_path).unwrap();
for op in &ops {
journal.append(op).unwrap();
}
let mut acks = AckTable::new();
acks.ack("a", ops[1].hlc.clone());
acks.ack("b", ops[1].hlc.clone());
let outcome = compact_and_truncate(
&mut journal,
&ckpt_dir,
&acks,
&RetentionPolicy::keep_all(),
None,
)
.unwrap();
let expected_hash = outcome.plan.checkpoint.checkpoint_hash.clone();
acks.ack("a", ops[3].hlc.clone());
acks.ack("b", ops[3].hlc.clone());
match compact_and_truncate(
&mut journal,
&ckpt_dir,
&acks,
&RetentionPolicy::keep_all(),
None,
) {
Err(CompactError::TruncatedJournal { checkpoint_hash }) => {
assert_eq!(checkpoint_hash, expected_hash)
}
other => panic!("expected TruncatedJournal refusal, got {other:?}"),
}
}
#[test]
fn plan_refuses_an_invalid_log() {
let mut ops = simple_ops();
ops[1].payload = json!({"forged": true});
let mut acks = AckTable::new();
acks.ack("a", ops[2].hlc.clone());
acks.ack("b", ops[3].hlc.clone());
assert!(matches!(
plan_compaction(&ops, &acks, &RetentionPolicy::keep_all(), None),
Err(CompactError::Chain(ChainError::IdMismatch { .. }))
));
}
#[test]
fn intent_retention_rule_is_rejected() {
use crate::lease::{Intent, IntentStatus};
let mut dev = DeviceLog::new("a");
let ops = vec![dev.append(
Scope::Personal,
Surface::Intent,
Intent::new("milo", "R", 1, IntentStatus::Committed).payload(),
)];
let state = fold(&ops);
let mut policy = RetentionPolicy::keep_all();
policy
.rules
.insert("intent".to_string(), RetentionRule::LastN { n: 1 });
assert!(matches!(
apply_retention(&state, &policy, 0),
Err(CompactError::IntentRetention)
));
let (retained, _) = apply_retention(&state, &RetentionPolicy::keep_all(), 0).unwrap();
assert!(retained.committed_run("milo", "R").is_some());
}
#[test]
fn c2_committed_run_survives_a_checkpoint_compaction() {
use crate::fold::fold_onto;
use crate::lease::{Intent, IntentStatus};
let mut a = DeviceLog::new("a");
let mut b = DeviceLog::new("b");
let mut ops = vec![a.append(
Scope::Personal,
Surface::Intent,
Intent::new("milo", "R", 1, IntentStatus::Committed).payload(),
)];
let split = ops.len();
for op in &ops {
b.observe(&op.hlc);
}
ops.push(b.append(
Scope::Personal,
Surface::Intent,
Intent::new("milo", "S", 2, IntentStatus::Pending).payload(),
));
let frontier = ops[..split].iter().map(|o| o.hlc.clone()).max().unwrap();
let mut acks = AckTable::new();
acks.ack("a", frontier.clone());
acks.ack("b", frontier);
let plan = plan_compaction(&ops, &acks, &RetentionPolicy::keep_all(), None).unwrap();
assert_eq!(
plan.dropped_ops, split,
"committed R is below the frontier, dropped from journal"
);
assert!(
plan.checkpoint.state.committed_run("milo", "R").is_some(),
"committed R survives compaction in the checkpoint oracle (C2 fixed)"
);
let reconstructed = fold_onto(&plan.checkpoint.state, &plan.retained_ops);
assert_eq!(
reconstructed,
fold(&ops),
"fold_onto(checkpoint, tail) == fold(full)"
);
assert!(
reconstructed.committed_run("milo", "R").is_some(),
"oracle intact post-compaction"
);
}
}