Skip to main content

macp_storage/storage/
recovery.rs

1use crate::log_store::LogEntry;
2use macp_core::session::Session;
3use std::fs;
4use std::path::Path;
5
6pub fn recover_session(session: &mut Session, log_entries: &[LogEntry]) {
7    // Ensure all log entry message IDs are in the session's dedup set.
8    // If the runtime crashed after writing a log entry but before persisting
9    // the session snapshot, there will be entries in the log not reflected
10    // in seen_message_ids.
11    let mut recovered = 0usize;
12    for entry in log_entries {
13        if !entry.message_id.is_empty() && session.seen_message_ids.insert(entry.message_id.clone())
14        {
15            recovered += 1;
16        }
17    }
18    if recovered > 0 {
19        eprintln!(
20            "recovery: session '{}' reconciled {} log entries into dedup state",
21            session.session_id, recovered
22        );
23    }
24}
25
26pub fn cleanup_temp_files(base_dir: &Path) {
27    let sessions_dir = base_dir.join("sessions");
28    if !sessions_dir.exists() {
29        return;
30    }
31    if let Ok(entries) = fs::read_dir(&sessions_dir) {
32        for entry in entries.flatten() {
33            if !entry.file_type().map(|ft| ft.is_dir()).unwrap_or(false) {
34                continue;
35            }
36            let dir = entry.path();
37            if let Ok(files) = fs::read_dir(&dir) {
38                for file in files.flatten() {
39                    let path = file.path();
40                    if path.extension().and_then(|e| e.to_str()) == Some("tmp") {
41                        eprintln!("recovery: removing orphaned temp file {}", path.display());
42                        let _ = fs::remove_file(&path);
43                    }
44                }
45            }
46        }
47    }
48}
49
50#[cfg(test)]
51mod tests {
52    use super::*;
53    use crate::log_store::EntryKind;
54
55    use std::collections::HashSet;
56
57    fn sample_session() -> Session {
58        Session::builder("s1", "macp.mode.decision.v1", "alice")
59            .ttl_expiry(61_000)
60            .ttl_ms(60_000)
61            .started_at_unix_ms(1_000)
62            .participants(vec!["alice".into()])
63            .seen_message_ids(HashSet::from(["m1".into()]))
64            .mode_version("1.0.0")
65            .configuration_version("cfg-1")
66            .policy_version("pol-1")
67            .build()
68    }
69
70    fn sample_entry(id: &str) -> LogEntry {
71        LogEntry {
72            message_id: id.into(),
73            received_at_ms: 1_700_000_000_000,
74            sender: "alice".into(),
75            message_type: "Message".into(),
76            raw_payload: vec![],
77            entry_kind: EntryKind::Incoming,
78            session_id: String::new(),
79            mode: String::new(),
80            macp_version: String::new(),
81            timestamp_unix_ms: 1_700_000_000_000,
82            bound_mode_version: None,
83            semantics_rev: 0,
84            bound_max_suspend_ms: None,
85            compacted_incoming_ordinals: 0,
86        }
87    }
88
89    #[test]
90    fn crash_recovery_reconciles_dedup_state() {
91        let mut session = sample_session();
92        assert!(session.seen_message_ids.contains("m1"));
93        assert!(!session.seen_message_ids.contains("m2"));
94        assert!(!session.seen_message_ids.contains("m3"));
95
96        let entries = vec![sample_entry("m1"), sample_entry("m2"), sample_entry("m3")];
97
98        recover_session(&mut session, &entries);
99
100        assert!(session.seen_message_ids.contains("m1"));
101        assert!(session.seen_message_ids.contains("m2"));
102        assert!(session.seen_message_ids.contains("m3"));
103    }
104
105    #[test]
106    fn recovery_ignores_entries_with_empty_message_id() {
107        let mut session = sample_session();
108        let before = session.seen_message_ids.len();
109
110        let entries = vec![sample_entry(""), sample_entry("m2")];
111        recover_session(&mut session, &entries);
112
113        assert!(!session.seen_message_ids.contains(""));
114        assert!(session.seen_message_ids.contains("m2"));
115        assert_eq!(session.seen_message_ids.len(), before + 1);
116    }
117
118    #[test]
119    fn recovery_is_idempotent() {
120        let mut session = sample_session();
121        let entries = vec![sample_entry("m1"), sample_entry("m2"), sample_entry("m3")];
122
123        recover_session(&mut session, &entries);
124        let after_first = session.seen_message_ids.clone();
125
126        recover_session(&mut session, &entries);
127        assert_eq!(session.seen_message_ids, after_first);
128        assert_eq!(session.seen_message_ids.len(), after_first.len());
129    }
130
131    #[test]
132    fn cleanup_temp_files_removes_orphans() {
133        let dir = tempfile::tempdir().unwrap();
134        let base = dir.path();
135        let sessions_dir = base.join("sessions").join("s1");
136        fs::create_dir_all(&sessions_dir).unwrap();
137
138        fs::write(sessions_dir.join("session.json.tmp"), b"partial").unwrap();
139        assert!(sessions_dir.join("session.json.tmp").exists());
140
141        cleanup_temp_files(base);
142
143        assert!(!sessions_dir.join("session.json.tmp").exists());
144    }
145
146    #[test]
147    fn cleanup_temp_files_preserves_non_tmp_files() {
148        let dir = tempfile::tempdir().unwrap();
149        let base = dir.path();
150        let session_dir = base.join("sessions").join("s1");
151        fs::create_dir_all(&session_dir).unwrap();
152
153        fs::write(session_dir.join("session.json"), b"{}").unwrap();
154        fs::write(session_dir.join("log.jsonl"), b"").unwrap();
155        fs::write(session_dir.join("session.json.tmp"), b"partial").unwrap();
156
157        cleanup_temp_files(base);
158
159        assert!(session_dir.join("session.json").exists());
160        assert!(session_dir.join("log.jsonl").exists());
161        assert!(!session_dir.join("session.json.tmp").exists());
162    }
163
164    #[test]
165    fn cleanup_temp_files_handles_missing_sessions_dir() {
166        let dir = tempfile::tempdir().unwrap();
167        // No sessions/ dir at all: must be a silent no-op, not a panic.
168        cleanup_temp_files(dir.path());
169    }
170}