1use crate::checkpoint::Checkpoint;
63use crate::fold::{FoldedRecord, SyncState};
64use crate::journal::OplogJournal;
65use crate::oplog::{verify_log, ChainError, Hlc, OpRecord};
66use serde::{Deserialize, Serialize};
67use serde_json::{json, Value};
68use std::collections::{BTreeMap, BTreeSet};
69use std::fmt;
70use std::fs::{self, File};
71use std::io::Write;
72use std::path::{Path, PathBuf};
73
74pub const RUNS_MAX_PER_AGENT: usize = 50;
77pub const RUNS_MAX_AGE_MS: u64 = 30 * 24 * 60 * 60 * 1000;
80
81#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
83#[serde(rename_all = "snake_case")]
84pub enum RetentionRule {
85 KeepAll,
88 LastN { n: usize },
91 MaxAgeMs { max_age_ms: u64 },
94 PerAgentWithAge {
98 max_per_agent: usize,
99 max_age_ms: u64,
100 },
101}
102
103#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
108pub struct RetentionPolicy {
109 pub rules: BTreeMap<String, RetentionRule>,
110 pub default_rule: RetentionRule,
111}
112
113impl Default for RetentionPolicy {
114 fn default() -> Self {
115 Self::keep_all()
116 }
117}
118
119impl RetentionPolicy {
120 pub fn keep_all() -> Self {
123 Self {
124 rules: BTreeMap::new(),
125 default_rule: RetentionRule::KeepAll,
126 }
127 }
128
129 pub fn proposal_default(conversation_last_n: usize, trajectory_max_age_ms: u64) -> Self {
138 let mut rules = BTreeMap::new();
139 rules.insert(
140 "conversation".to_string(),
141 RetentionRule::LastN {
142 n: conversation_last_n,
143 },
144 );
145 rules.insert(
146 "run".to_string(),
147 RetentionRule::PerAgentWithAge {
148 max_per_agent: RUNS_MAX_PER_AGENT,
149 max_age_ms: RUNS_MAX_AGE_MS,
150 },
151 );
152 rules.insert(
153 "trajectory".to_string(),
154 RetentionRule::MaxAgeMs {
155 max_age_ms: trajectory_max_age_ms,
156 },
157 );
158 rules.insert("knowledge".to_string(), RetentionRule::KeepAll);
159 rules.insert("skill".to_string(), RetentionRule::KeepAll);
160 rules.insert("routing".to_string(), RetentionRule::KeepAll);
161 Self {
162 rules,
163 default_rule: RetentionRule::KeepAll,
164 }
165 }
166
167 pub fn rule_for(&self, surface_tag: &str) -> &RetentionRule {
168 self.rules.get(surface_tag).unwrap_or(&self.default_rule)
169 }
170}
171
172#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
176pub struct RetentionReport {
177 pub dropped: BTreeMap<String, usize>,
178 pub tombstoned: Vec<(String, String)>,
179}
180
181fn value_ts(payload: &Value) -> Option<u64> {
182 payload
183 .get("timestamp")
184 .and_then(|v| v.as_u64().or_else(|| v.as_f64().map(|f| f as u64)))
185}
186
187fn payload_ts(record: &FoldedRecord) -> Option<u64> {
188 value_ts(&record.payload)
189}
190
191pub fn as_of_from_ops(ops: &[OpRecord]) -> u64 {
199 ops.iter()
200 .filter_map(|op| value_ts(&op.payload))
201 .max()
202 .unwrap_or(0)
203}
204
205fn within_age(record: &FoldedRecord, as_of_ms: u64, max_age_ms: u64) -> bool {
206 match payload_ts(record) {
207 Some(ts) => as_of_ms.saturating_sub(ts) <= max_age_ms,
209 None => true,
211 }
212}
213
214fn recency_sorted(entries: &BTreeMap<String, FoldedRecord>) -> Vec<(&String, &FoldedRecord)> {
217 let mut sorted: Vec<(&String, &FoldedRecord)> = entries.iter().collect();
218 sorted.sort_by(|(_, a), (_, b)| {
219 (payload_ts(a).unwrap_or(u64::MAX), &a.hlc, &a.op_id).cmp(&(
220 payload_ts(b).unwrap_or(u64::MAX),
221 &b.hlc,
222 &b.op_id,
223 ))
224 });
225 sorted
226}
227
228fn select_retained(
229 entries: &BTreeMap<String, FoldedRecord>,
230 rule: &RetentionRule,
231 as_of_ms: u64,
232) -> BTreeSet<String> {
233 match rule {
234 RetentionRule::KeepAll => entries.keys().cloned().collect(),
235 RetentionRule::LastN { n } => recency_sorted(entries)
236 .into_iter()
237 .rev()
238 .take(*n)
239 .map(|(k, _)| k.clone())
240 .collect(),
241 RetentionRule::MaxAgeMs { max_age_ms } => entries
242 .iter()
243 .filter(|(_, r)| within_age(r, as_of_ms, *max_age_ms))
244 .map(|(k, _)| k.clone())
245 .collect(),
246 RetentionRule::PerAgentWithAge {
247 max_per_agent,
248 max_age_ms,
249 } => {
250 let mut per_agent_rank: BTreeMap<&str, usize> = BTreeMap::new();
253 let mut keep = BTreeSet::new();
254 for (key, record) in recency_sorted(entries).into_iter().rev() {
255 let agent = record
256 .payload
257 .get("agent_id")
258 .and_then(Value::as_str)
259 .unwrap_or("");
260 let rank = per_agent_rank.entry(agent).or_insert(0);
261 let over_count = *rank >= *max_per_agent;
262 *rank += 1;
263 if !over_count && within_age(record, as_of_ms, *max_age_ms) {
264 keep.insert(key.clone());
265 }
266 }
267 keep
268 }
269 }
270}
271
272fn entity_id<'a>(key: &'a str, record: &'a FoldedRecord) -> Option<&'a str> {
275 key.strip_prefix("id:")
276 .or_else(|| record.payload.get("id").and_then(Value::as_str))
277}
278
279pub fn is_tombstone(record: &FoldedRecord) -> bool {
282 record
283 .payload
284 .get("tombstone")
285 .and_then(Value::as_bool)
286 .unwrap_or(false)
287}
288
289fn tombstone_of(record: &FoldedRecord, id: &str) -> FoldedRecord {
294 FoldedRecord {
295 op_id: record.op_id.clone(),
296 hlc: record.hlc.clone(),
297 payload: json!({"id": id, "tombstone": true}),
298 }
299}
300
301pub fn apply_retention(
312 state: &SyncState,
313 policy: &RetentionPolicy,
314 as_of_ms: u64,
315) -> Result<(SyncState, RetentionReport), CompactError> {
316 let routing_tag = crate::oplog::Surface::Routing.tag();
327 debug_assert!(crate::oplog::Surface::Routing.is_replay_stream());
328 if state.logs.contains_key(&routing_tag)
329 && policy.rule_for(&routing_tag) != &RetentionRule::KeepAll
330 {
331 return Err(CompactError::EventStreamRetention {
332 surface: routing_tag,
333 });
334 }
335
336 if policy.rule_for(&crate::oplog::Surface::Intent.tag()) != &RetentionRule::KeepAll {
346 return Err(CompactError::IntentRetention);
347 }
348
349 let mut retained = state.clone();
350 let mut report = RetentionReport::default();
351 for (tag, entries) in &state.logs {
352 let rule = policy.rule_for(tag);
353 if rule == &RetentionRule::KeepAll {
354 continue;
355 }
356 let live: BTreeMap<String, FoldedRecord> = entries
359 .iter()
360 .filter(|(_, record)| !is_tombstone(record))
361 .map(|(key, record)| (key.clone(), record.clone()))
362 .collect();
363 let keep = select_retained(&live, rule, as_of_ms);
364 let surface = retained.logs.get_mut(tag).expect("cloned from state");
365 for (key, record) in &live {
366 if keep.contains(key) {
367 continue;
368 }
369 match entity_id(key, record) {
370 Some(id) => {
371 surface.insert(key.clone(), tombstone_of(record, id));
372 report.tombstoned.push((tag.clone(), key.clone()));
373 }
374 None => {
375 surface.remove(key);
376 *report.dropped.entry(tag.clone()).or_insert(0) += 1;
377 }
378 }
379 }
380 }
381 report.tombstoned.sort();
382 Ok((retained, report))
383}
384
385#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
398pub struct AckTable {
399 acked: BTreeMap<String, Hlc>,
400}
401
402impl AckTable {
403 pub fn new() -> Self {
404 Self::default()
405 }
406
407 pub fn ack(&mut self, device_id: impl Into<String>, frontier: Hlc) -> bool {
411 let device_id = device_id.into();
412 match self.acked.get(&device_id) {
413 Some(current) if frontier <= *current => false,
414 _ => {
415 self.acked.insert(device_id, frontier);
416 true
417 }
418 }
419 }
420
421 pub fn get(&self, device_id: &str) -> Option<&Hlc> {
422 self.acked.get(device_id)
423 }
424
425 pub fn devices(&self) -> impl Iterator<Item = &str> {
426 self.acked.keys().map(String::as_str)
427 }
428
429 pub fn stable_frontier(&self) -> Option<&Hlc> {
434 self.acked.values().min()
435 }
436
437 pub fn save(&self, path: &Path) -> std::io::Result<()> {
441 if let Some(parent) = path.parent() {
442 if !parent.as_os_str().is_empty() {
443 fs::create_dir_all(parent)?;
444 }
445 }
446 let tmp_path = {
447 let mut s = path.as_os_str().to_owned();
448 s.push(".tmp");
449 PathBuf::from(s)
450 };
451 {
452 let mut tmp = File::create(&tmp_path)?;
453 tmp.write_all(
454 serde_json::to_string(self)
455 .map_err(std::io::Error::other)?
456 .as_bytes(),
457 )?;
458 tmp.sync_all()?;
459 }
460 fs::rename(&tmp_path, path)
461 }
462
463 pub fn load(path: &Path) -> std::io::Result<Self> {
467 match fs::read_to_string(path) {
468 Ok(raw) => serde_json::from_str(&raw).map_err(std::io::Error::other),
469 Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(Self::default()),
470 Err(e) => Err(e),
471 }
472 }
473}
474
475#[derive(Debug)]
478pub enum CompactError {
479 Chain(ChainError),
482 NothingAcked,
485 UnackedDevice {
488 device_id: String,
489 },
490 EventStreamRetention {
494 surface: String,
495 },
496 IntentRetention,
500 TruncatedJournal {
505 checkpoint_hash: String,
506 },
507 Io(std::io::Error),
508}
509
510impl fmt::Display for CompactError {
511 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
512 match self {
513 CompactError::Chain(e) => write!(f, "compaction refused: log does not verify: {e}"),
514 CompactError::NothingAcked => {
515 write!(
516 f,
517 "compaction refused: no acked frontier exists (empty ack table)"
518 )
519 }
520 CompactError::UnackedDevice { device_id } => write!(
521 f,
522 "compaction refused: device {device_id} appears in the log but has no acked \
523 frontier — its ops may include state no other replica has folded"
524 ),
525 CompactError::EventStreamRetention { surface } => write!(
526 f,
527 "retention policy for event-stream surface {surface} must be keep_all — \
528 observation multisets replay from genesis and cannot be trimmed"
529 ),
530 CompactError::IntentRetention => write!(
531 f,
532 "retention policy for the leased `intent` surface must be keep_all — the \
533 committed-run idempotency oracle must survive compaction and cannot be trimmed"
534 ),
535 CompactError::TruncatedJournal { checkpoint_hash } => write!(
536 f,
537 "compaction refused: journal already truncated below checkpoint \
538 {checkpoint_hash} — re-planning from the tail alone would drop that \
539 checkpoint's state (recompaction over a checkpoint base is a later slice)"
540 ),
541 CompactError::Io(e) => write!(f, "compaction io error: {e}"),
542 }
543 }
544}
545
546impl std::error::Error for CompactError {}
547
548#[derive(Debug)]
551pub struct CompactionPlan {
552 pub checkpoint: Checkpoint,
555 pub retained_ops: Vec<OpRecord>,
557 pub dropped_ops: usize,
560 pub frontier: Hlc,
562 pub as_of_ms: u64,
565 pub retention: RetentionReport,
567}
568
569pub fn plan_compaction(
580 ops: &[OpRecord],
581 acks: &AckTable,
582 policy: &RetentionPolicy,
583 as_of_ms: Option<u64>,
584) -> Result<CompactionPlan, CompactError> {
585 verify_log(ops).map_err(CompactError::Chain)?;
586 let frontier = acks
587 .stable_frontier()
588 .cloned()
589 .ok_or(CompactError::NothingAcked)?;
590 for op in ops {
591 if acks.get(&op.device_id).is_none() {
592 return Err(CompactError::UnackedDevice {
593 device_id: op.device_id.clone(),
594 });
595 }
596 }
597
598 let (below, retained_ops): (Vec<OpRecord>, Vec<OpRecord>) =
599 ops.iter().cloned().partition(|op| op.hlc <= frontier);
600 let as_of_ms = as_of_ms.unwrap_or_else(|| as_of_from_ops(&below));
601
602 let exact = Checkpoint::from_ops(&below).map_err(CompactError::Chain)?;
605 let (retained_state, retention) = apply_retention(&exact.state, policy, as_of_ms)?;
606 let checkpoint = Checkpoint::assemble(exact.frontier, exact.scopes, retained_state);
607
608 Ok(CompactionPlan {
609 checkpoint,
610 retained_ops,
611 dropped_ops: below.len(),
612 frontier,
613 as_of_ms,
614 retention,
615 })
616}
617
618#[derive(Debug)]
620pub struct CompactionOutcome {
621 pub checkpoint_path: Option<PathBuf>,
627 pub plan: CompactionPlan,
628}
629
630pub fn compact_and_truncate(
653 journal: &mut OplogJournal,
654 checkpoint_dir: &Path,
655 acks: &AckTable,
656 policy: &RetentionPolicy,
657 as_of_ms: Option<u64>,
658) -> Result<CompactionOutcome, CompactError> {
659 let (marker, ops) = OplogJournal::load_with_marker(journal.path()).map_err(CompactError::Io)?;
660 if let Some(marker) = marker {
661 return Err(CompactError::TruncatedJournal {
662 checkpoint_hash: marker.checkpoint_hash,
663 });
664 }
665 let plan = plan_compaction(&ops, acks, policy, as_of_ms)?;
666 if plan.dropped_ops == 0 {
667 return Ok(CompactionOutcome {
670 checkpoint_path: None,
671 plan,
672 });
673 }
674 let checkpoint_path = plan
676 .checkpoint
677 .save(checkpoint_dir)
678 .map_err(CompactError::Io)?;
679 journal
680 .truncate_to(&plan.retained_ops, &plan.checkpoint.checkpoint_hash)
681 .map_err(CompactError::Io)?;
682 Ok(CompactionOutcome {
683 checkpoint_path: Some(checkpoint_path),
684 plan,
685 })
686}
687
688#[cfg(test)]
689mod tests {
690 use super::*;
691 use crate::fold::fold;
692 use crate::oplog::{DeviceLog, Scope, Surface};
693 use serde_json::json;
694
695 fn hlc(wall_ms: u64, device: &str) -> Hlc {
696 Hlc {
697 wall_ms,
698 counter: 0,
699 device_id: device.into(),
700 }
701 }
702
703 #[test]
704 fn ack_table_is_monotone_only_and_min_frontier() {
705 let mut acks = AckTable::new();
706 assert_eq!(acks.stable_frontier(), None);
707 assert!(acks.ack("a", hlc(5, "a")));
708 assert!(acks.ack("b", hlc(9, "b")));
709 assert_eq!(
710 acks.stable_frontier(),
711 Some(&hlc(5, "a")),
712 "min over devices"
713 );
714
715 assert!(!acks.ack("b", hlc(3, "b")));
717 assert!(!acks.ack("b", hlc(9, "b")), "equal is not an advance");
718 assert_eq!(acks.get("b"), Some(&hlc(9, "b")));
719 assert!(acks.ack("b", hlc(12, "b")));
720 assert_eq!(acks.get("b"), Some(&hlc(12, "b")));
721 }
722
723 #[test]
724 fn ack_table_persists_atomically_and_loads_missing_as_empty() {
725 let dir = tempfile::tempdir().unwrap();
726 let path = dir.path().join("nested").join("acks.json");
727 assert_eq!(
728 AckTable::load(&path).unwrap(),
729 AckTable::new(),
730 "missing → empty"
731 );
732
733 let mut acks = AckTable::new();
734 acks.ack("a", hlc(5, "a"));
735 acks.ack("b", hlc(9, "b"));
736 acks.save(&path).unwrap();
737 assert_eq!(AckTable::load(&path).unwrap(), acks);
738 assert!(!path.parent().unwrap().join("acks.json.tmp").exists());
740 }
741
742 fn simple_ops() -> Vec<OpRecord> {
744 let mut a = DeviceLog::new("a");
745 let mut b = DeviceLog::new("b");
746 let mut ops = vec![
747 a.append(Scope::Personal, Surface::Knowledge, json!({"id": "f1"})),
748 a.append(Scope::Personal, Surface::Knowledge, json!({"id": "f2"})),
749 a.append(Scope::Personal, Surface::Knowledge, json!({"id": "f3"})),
750 ];
751 for op in &ops {
752 b.observe(&op.hlc);
753 }
754 ops.push(b.append(Scope::Personal, Surface::Knowledge, json!({"id": "f4"})));
755 ops
756 }
757
758 #[test]
759 fn compaction_refuses_without_acks() {
760 let ops = simple_ops();
761 let policy = RetentionPolicy::keep_all();
762 assert!(matches!(
763 plan_compaction(&ops, &AckTable::new(), &policy, None),
764 Err(CompactError::NothingAcked)
765 ));
766
767 let mut acks = AckTable::new();
770 acks.ack("a", ops[2].hlc.clone());
771 assert!(matches!(
772 plan_compaction(&ops, &acks, &policy, None),
773 Err(CompactError::UnackedDevice { .. })
774 ));
775 }
776
777 #[test]
778 fn compaction_never_drops_above_a_lagging_ack() {
779 let ops = simple_ops();
780 let mut acks = AckTable::new();
781 acks.ack("b", ops[3].hlc.clone());
783 acks.ack("a", ops[1].hlc.clone());
784
785 let plan = plan_compaction(&ops, &acks, &RetentionPolicy::keep_all(), None).unwrap();
786 assert_eq!(
787 plan.frontier, ops[1].hlc,
788 "stable frontier = the lagging device's ack"
789 );
790 assert_eq!(plan.dropped_ops, 2, "only ops ≤ the lagging frontier drop");
791 assert_eq!(plan.retained_ops.len(), 2);
792 assert!(
793 plan.retained_ops.iter().all(|op| op.hlc > plan.frontier),
794 "everything a device hasn't seen stays in the journal"
795 );
796 }
797
798 #[test]
799 fn retention_conversations_last_n_by_timestamp() {
800 let mut dev = DeviceLog::new("a");
801 let ops: Vec<OpRecord> = (0..5)
802 .map(|i| {
803 dev.append(
804 Scope::Personal,
805 Surface::Conversation,
806 json!({"speaker": "u", "text": format!("t{i}"), "timestamp": 100 + i}),
807 )
808 })
809 .collect();
810 let state = fold(&ops);
811 let policy = RetentionPolicy::proposal_default(2, u64::MAX);
812 let (retained, report) = apply_retention(&state, &policy, 1_000).unwrap();
813 let tag = Surface::Conversation.tag();
814 let texts: Vec<String> = retained
815 .log_entries(&tag)
816 .iter()
817 .map(|r| r.payload["text"].as_str().unwrap().to_string())
818 .collect();
819 assert_eq!(texts, vec!["t3", "t4"], "last 2 turns by timestamp survive");
820 assert_eq!(report.dropped[&tag], 3);
821 }
822
823 #[test]
824 fn retention_runs_per_agent_and_age_matches_run_store_gc() {
825 assert_eq!(
826 RUNS_MAX_PER_AGENT, 50,
827 "parity with run_store DEFAULT_MAX_RUNS_PER_AGENT"
828 );
829 assert_eq!(
830 RUNS_MAX_AGE_MS,
831 30 * 24 * 60 * 60 * 1000,
832 "parity with DEFAULT_MAX_AGE_DAYS"
833 );
834
835 let mut dev = DeviceLog::new("a");
838 let mut ops = Vec::new();
839 for (id, agent, ts) in [
840 ("r1", "milo", 10u64), ("r2", "milo", 60), ("r3", "milo", 70), ("r4", "other", 10), ("r5", "other", 80), ] {
846 ops.push(dev.append(
847 Scope::Personal,
848 Surface::Run,
849 json!({"id": id, "agent_id": agent, "timestamp": ts}),
850 ));
851 }
852 let state = fold(&ops);
853 let mut policy = RetentionPolicy::keep_all();
854 policy.rules.insert(
855 "run".to_string(),
856 RetentionRule::PerAgentWithAge {
857 max_per_agent: 2,
858 max_age_ms: 50,
859 },
860 );
861 let (retained, report) = apply_retention(&state, &policy, 100).unwrap();
862 let entries = retained.log_entries(&Surface::Run.tag());
863 let kept: Vec<&str> = entries
864 .iter()
865 .filter(|r| !is_tombstone(r))
866 .map(|r| r.payload["id"].as_str().unwrap())
867 .collect();
868 assert_eq!(kept, vec!["r2", "r3", "r5"]);
869 let stubs: Vec<&str> = entries
871 .iter()
872 .filter(|r| is_tombstone(r))
873 .map(|r| r.payload["id"].as_str().unwrap())
874 .collect();
875 assert_eq!(stubs, vec!["r1", "r4"]);
876 assert_eq!(report.tombstoned.len(), 2);
877 assert!(
878 report.dropped.is_empty(),
879 "id-bearing entries are stubbed, never erased"
880 );
881 }
882
883 #[test]
884 fn retention_trajectories_by_age_and_undated_never_dropped() {
885 let mut dev = DeviceLog::new("a");
886 let ops = vec![
887 dev.append(
888 Scope::Personal,
889 Surface::Trajectory,
890 json!({"id": "old", "timestamp": 10}),
891 ),
892 dev.append(
893 Scope::Personal,
894 Surface::Trajectory,
895 json!({"id": "new", "timestamp": 90}),
896 ),
897 dev.append(
898 Scope::Personal,
899 Surface::Trajectory,
900 json!({"id": "undated"}),
901 ),
902 ];
903 let state = fold(&ops);
904 let mut policy = RetentionPolicy::keep_all();
905 policy.rules.insert(
906 "trajectory".to_string(),
907 RetentionRule::MaxAgeMs { max_age_ms: 30 },
908 );
909 let (retained, _) = apply_retention(&state, &policy, 100).unwrap();
910 let entries = retained.log_entries(&Surface::Trajectory.tag());
911 let kept: Vec<&str> = entries
912 .iter()
913 .filter(|r| !is_tombstone(r))
914 .map(|r| r.payload["id"].as_str().unwrap())
915 .collect();
916 assert!(kept.contains(&"new"));
917 assert!(
918 kept.contains(&"undated"),
919 "undated data is never silently age-dropped"
920 );
921 assert!(!kept.contains(&"old"));
922 assert!(entries
924 .iter()
925 .any(|r| is_tombstone(r) && r.payload["id"] == json!("old")));
926 }
927
928 #[test]
929 fn knowledge_and_skills_keep_all_under_the_proposal_default() {
930 let mut dev = DeviceLog::new("a");
931 let ops = vec![
932 dev.append(
933 Scope::Personal,
934 Surface::Knowledge,
935 json!({"id": "f1", "timestamp": 1}),
936 ),
937 dev.append(
938 Scope::Personal,
939 Surface::Skill,
940 json!({"id": "s1", "timestamp": 1}),
941 ),
942 ];
943 let state = fold(&ops);
944 let (retained, report) =
946 apply_retention(&state, &RetentionPolicy::proposal_default(1, 1), u64::MAX).unwrap();
947 assert_eq!(retained.logs[&Surface::Knowledge.tag()].len(), 1);
948 assert_eq!(retained.logs[&Surface::Skill.tag()].len(), 1);
949 assert!(report.dropped.is_empty());
950 }
951
952 #[test]
953 fn event_stream_retention_is_rejected() {
954 let mut dev = DeviceLog::new("a");
955 let ops = vec![
956 dev.append(Scope::Personal, Surface::Routing, json!({"sample": 1.0})),
957 dev.append(Scope::Personal, Surface::Routing, json!({"sample": 0.0})),
958 ];
959 let state = fold(&ops);
960 let mut policy = RetentionPolicy::keep_all();
961 policy
962 .rules
963 .insert("routing".to_string(), RetentionRule::LastN { n: 1 });
964 assert!(matches!(
965 apply_retention(&state, &policy, 0),
966 Err(CompactError::EventStreamRetention { .. })
967 ));
968 let (retained, _) =
970 apply_retention(&state, &RetentionPolicy::proposal_default(10, 10), u64::MAX).unwrap();
971 assert_eq!(retained.log_entries(&Surface::Routing.tag()).len(), 2);
972 }
973
974 #[test]
975 fn every_dropped_id_bearing_entry_leaves_a_tombstone_stub() {
976 let mut dev = DeviceLog::new("a");
981 let ops = vec![
982 dev.append(
983 Scope::Personal,
984 Surface::Knowledge,
985 json!({"id": "f1", "timestamp": 1}),
986 ),
987 dev.append(
988 Scope::Personal,
989 Surface::Knowledge,
990 json!({"id": "f2", "timestamp": 2, "supersedes": "f1"}),
991 ),
992 dev.append(
993 Scope::Personal,
994 Surface::Knowledge,
995 json!({"id": "f3", "timestamp": 3, "supersedes": ["f2"]}),
996 ),
997 dev.append(
998 Scope::Personal,
999 Surface::Knowledge,
1000 json!({"id": "f0", "timestamp": 0}),
1001 ),
1002 ];
1003 let state = fold(&ops);
1004 let mut policy = RetentionPolicy::keep_all();
1005 policy
1006 .rules
1007 .insert("knowledge".to_string(), RetentionRule::LastN { n: 1 });
1008 let (retained, report) = apply_retention(&state, &policy, 10).unwrap();
1009 let tag = Surface::Knowledge.tag();
1010 let surface = &retained.logs[&tag];
1011 assert_eq!(
1012 surface["id:f3"].payload["timestamp"],
1013 json!(3),
1014 "newest stays live"
1015 );
1016 for id in ["f0", "f1", "f2"] {
1017 let stub = &surface[&format!("id:{id}")];
1018 assert!(is_tombstone(stub), "{id} left a stub");
1019 assert_eq!(
1020 stub.payload,
1021 json!({"id": id, "tombstone": true}),
1022 "minimal stub shape"
1023 );
1024 assert_eq!(
1025 stub.op_id,
1026 state.logs[&tag][&format!("id:{id}")].op_id,
1027 "stub keeps the original op identity"
1028 );
1029 }
1030 assert_eq!(report.tombstoned.len(), 3);
1031 assert!(report.dropped.is_empty());
1032
1033 let (again, report2) = apply_retention(&retained, &policy, 10).unwrap();
1036 assert_eq!(again, retained);
1037 assert!(report2.tombstoned.is_empty());
1038 }
1039
1040 #[test]
1041 fn derived_as_of_is_the_max_below_frontier_timestamp() {
1042 let mut dev = DeviceLog::new("a");
1043 let ops = vec![
1044 dev.append(
1045 Scope::Personal,
1046 Surface::Trajectory,
1047 json!({"id": "t1", "timestamp": 40}),
1048 ),
1049 dev.append(
1050 Scope::Personal,
1051 Surface::Trajectory,
1052 json!({"id": "t2", "timestamp": 100}),
1053 ),
1054 dev.append(Scope::Personal, Surface::Skill, json!({"id": "s1"})), ];
1056 assert_eq!(as_of_from_ops(&ops), 100, "max payload timestamp");
1057 assert_eq!(
1058 as_of_from_ops(&ops[2..]),
1059 0,
1060 "no timestamps → 0 (age rules drop nothing)"
1061 );
1062
1063 let mut acks = AckTable::new();
1067 acks.ack("a", ops[2].hlc.clone());
1068 let mut policy = RetentionPolicy::keep_all();
1069 policy.rules.insert(
1070 "trajectory".to_string(),
1071 RetentionRule::MaxAgeMs { max_age_ms: 30 },
1072 );
1073 let plan = plan_compaction(&ops, &acks, &policy, None).unwrap();
1074 assert_eq!(plan.as_of_ms, 100);
1075 let surface = &plan.checkpoint.state.logs[&Surface::Trajectory.tag()];
1077 assert!(is_tombstone(&surface["id:t1"]));
1078 assert!(!is_tombstone(&surface["id:t2"]));
1079 }
1080
1081 #[test]
1082 fn empty_below_frontier_compaction_is_a_no_op() {
1083 let dir = tempfile::tempdir().unwrap();
1084 let journal_path = dir.path().join("oplog.jsonl");
1085 let ckpt_dir = dir.path().join("checkpoints");
1086
1087 let ops = simple_ops();
1090 let mut journal = OplogJournal::open(&journal_path).unwrap();
1091 for op in &ops {
1092 journal.append(op).unwrap();
1093 }
1094 let mut acks = AckTable::new();
1095 acks.ack(
1096 "a",
1097 Hlc {
1098 wall_ms: 0,
1099 counter: 0,
1100 device_id: "a".into(),
1101 },
1102 );
1103 acks.ack(
1104 "b",
1105 Hlc {
1106 wall_ms: 0,
1107 counter: 0,
1108 device_id: "b".into(),
1109 },
1110 );
1111
1112 let before = fs::read_to_string(&journal_path).unwrap();
1113 let outcome = compact_and_truncate(
1114 &mut journal,
1115 &ckpt_dir,
1116 &acks,
1117 &RetentionPolicy::keep_all(),
1118 None,
1119 )
1120 .unwrap();
1121 assert_eq!(outcome.plan.dropped_ops, 0);
1122 assert!(
1123 outcome.checkpoint_path.is_none(),
1124 "no empty checkpoint file written"
1125 );
1126 assert!(!ckpt_dir.exists(), "checkpoint dir not even created");
1127 assert_eq!(
1128 fs::read_to_string(&journal_path).unwrap(),
1129 before,
1130 "journal untouched (no marker, no rewrite)"
1131 );
1132 assert_eq!(OplogJournal::load(&journal_path).unwrap(), ops);
1134 }
1135
1136 #[test]
1137 fn recompacting_an_already_truncated_journal_is_refused() {
1138 let dir = tempfile::tempdir().unwrap();
1139 let journal_path = dir.path().join("oplog.jsonl");
1140 let ckpt_dir = dir.path().join("checkpoints");
1141 let ops = simple_ops();
1142 let mut journal = OplogJournal::open(&journal_path).unwrap();
1143 for op in &ops {
1144 journal.append(op).unwrap();
1145 }
1146 let mut acks = AckTable::new();
1147 acks.ack("a", ops[1].hlc.clone());
1148 acks.ack("b", ops[1].hlc.clone());
1149 let outcome = compact_and_truncate(
1150 &mut journal,
1151 &ckpt_dir,
1152 &acks,
1153 &RetentionPolicy::keep_all(),
1154 None,
1155 )
1156 .unwrap();
1157 let expected_hash = outcome.plan.checkpoint.checkpoint_hash.clone();
1158
1159 acks.ack("a", ops[3].hlc.clone());
1163 acks.ack("b", ops[3].hlc.clone());
1164 match compact_and_truncate(
1165 &mut journal,
1166 &ckpt_dir,
1167 &acks,
1168 &RetentionPolicy::keep_all(),
1169 None,
1170 ) {
1171 Err(CompactError::TruncatedJournal { checkpoint_hash }) => {
1172 assert_eq!(checkpoint_hash, expected_hash)
1173 }
1174 other => panic!("expected TruncatedJournal refusal, got {other:?}"),
1175 }
1176 }
1177
1178 #[test]
1179 fn plan_refuses_an_invalid_log() {
1180 let mut ops = simple_ops();
1181 ops[1].payload = json!({"forged": true});
1182 let mut acks = AckTable::new();
1183 acks.ack("a", ops[2].hlc.clone());
1184 acks.ack("b", ops[3].hlc.clone());
1185 assert!(matches!(
1186 plan_compaction(&ops, &acks, &RetentionPolicy::keep_all(), None),
1187 Err(CompactError::Chain(ChainError::IdMismatch { .. }))
1188 ));
1189 }
1190
1191 #[test]
1192 fn intent_retention_rule_is_rejected() {
1193 use crate::lease::{Intent, IntentStatus};
1197 let mut dev = DeviceLog::new("a");
1198 let ops = vec![dev.append(
1199 Scope::Personal,
1200 Surface::Intent,
1201 Intent::new("milo", "R", 1, IntentStatus::Committed).payload(),
1202 )];
1203 let state = fold(&ops);
1204 let mut policy = RetentionPolicy::keep_all();
1205 policy
1206 .rules
1207 .insert("intent".to_string(), RetentionRule::LastN { n: 1 });
1208 assert!(matches!(
1209 apply_retention(&state, &policy, 0),
1210 Err(CompactError::IntentRetention)
1211 ));
1212 let (retained, _) = apply_retention(&state, &RetentionPolicy::keep_all(), 0).unwrap();
1214 assert!(retained.committed_run("milo", "R").is_some());
1215 }
1216
1217 #[test]
1218 fn c2_committed_run_survives_a_checkpoint_compaction() {
1219 use crate::fold::fold_onto;
1224 use crate::lease::{Intent, IntentStatus};
1225
1226 let mut a = DeviceLog::new("a");
1227 let mut b = DeviceLog::new("b");
1228 let mut ops = vec![a.append(
1229 Scope::Personal,
1230 Surface::Intent,
1231 Intent::new("milo", "R", 1, IntentStatus::Committed).payload(),
1232 )];
1233 let split = ops.len();
1234 for op in &ops {
1235 b.observe(&op.hlc);
1236 }
1237 ops.push(b.append(
1239 Scope::Personal,
1240 Surface::Intent,
1241 Intent::new("milo", "S", 2, IntentStatus::Pending).payload(),
1242 ));
1243
1244 let frontier = ops[..split].iter().map(|o| o.hlc.clone()).max().unwrap();
1246 let mut acks = AckTable::new();
1247 acks.ack("a", frontier.clone());
1248 acks.ack("b", frontier);
1249 let plan = plan_compaction(&ops, &acks, &RetentionPolicy::keep_all(), None).unwrap();
1250 assert_eq!(
1251 plan.dropped_ops, split,
1252 "committed R is below the frontier, dropped from journal"
1253 );
1254
1255 assert!(
1258 plan.checkpoint.state.committed_run("milo", "R").is_some(),
1259 "committed R survives compaction in the checkpoint oracle (C2 fixed)"
1260 );
1261 let reconstructed = fold_onto(&plan.checkpoint.state, &plan.retained_ops);
1262 assert_eq!(
1263 reconstructed,
1264 fold(&ops),
1265 "fold_onto(checkpoint, tail) == fold(full)"
1266 );
1267 assert!(
1268 reconstructed.committed_run("milo", "R").is_some(),
1269 "oracle intact post-compaction"
1270 );
1271 }
1272}