macp_storage/storage/
recovery.rs1use 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 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 cleanup_temp_files(dir.path());
169 }
170}