macp_storage/storage/
migration.rs1use crate::log_store::LogEntry;
2use crate::registry::PersistedSession;
3use std::collections::HashMap;
4use std::fs;
5use std::io;
6use std::path::Path;
7
8const STORAGE_VERSION: u32 = 3;
9
10#[derive(serde::Serialize, serde::Deserialize)]
11struct StorageVersion {
12 version: u32,
13}
14
15pub fn write_storage_version(base_dir: &Path) -> io::Result<()> {
16 let sv = StorageVersion {
17 version: STORAGE_VERSION,
18 };
19 let bytes = serde_json::to_vec_pretty(&sv)
20 .map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?;
21 fs::write(base_dir.join("storage_version.json"), bytes)
22}
23
24pub fn read_storage_version(base_dir: &Path) -> io::Result<Option<u32>> {
25 let path = base_dir.join("storage_version.json");
26 if !path.exists() {
27 return Ok(None);
28 }
29 let bytes = fs::read(&path)?;
30 let sv: StorageVersion = serde_json::from_slice(&bytes)
31 .map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?;
32 Ok(Some(sv.version))
33}
34
35pub fn migrate_if_needed(base_dir: &Path) -> io::Result<()> {
36 let sessions_dir = base_dir.join("sessions");
37 let legacy_sessions = base_dir.join("sessions.json");
38 let legacy_logs = base_dir.join("logs.json");
39
40 let current_version = read_storage_version(base_dir)?;
42 if sessions_dir.exists() && !legacy_sessions.exists() && !legacy_logs.exists() {
43 if current_version.unwrap_or(0) < STORAGE_VERSION {
46 write_storage_version(base_dir)?;
47 }
48 return Ok(());
49 }
50
51 if !legacy_sessions.exists() && !legacy_logs.exists() && !sessions_dir.exists() {
52 write_storage_version(base_dir)?;
53 return Ok(());
54 }
55
56 if sessions_dir.exists() && !legacy_sessions.exists() && !legacy_logs.exists() {
58 write_storage_version(base_dir)?;
59 return Ok(());
60 }
61
62 println!("Migrating legacy storage format to per-session directories...");
63
64 let sessions: HashMap<String, PersistedSession> = if legacy_sessions.exists() {
66 let bytes = fs::read(&legacy_sessions)?;
67 serde_json::from_slice(&bytes).map_err(|e| {
68 io::Error::new(
69 io::ErrorKind::InvalidData,
70 format!("failed to parse legacy sessions.json: {e}"),
71 )
72 })?
73 } else {
74 HashMap::new()
75 };
76
77 let logs: HashMap<String, Vec<LogEntry>> = if legacy_logs.exists() {
79 let bytes = fs::read(&legacy_logs)?;
80 serde_json::from_slice(&bytes).map_err(|e| {
81 io::Error::new(
82 io::ErrorKind::InvalidData,
83 format!("failed to parse legacy logs.json: {e}"),
84 )
85 })?
86 } else {
87 HashMap::new()
88 };
89
90 fs::create_dir_all(&sessions_dir)?;
91
92 for (session_id, persisted) in &sessions {
94 let dir = sessions_dir.join(session_id);
95 fs::create_dir_all(&dir)?;
96
97 let session_bytes = serde_json::to_vec_pretty(persisted)
99 .map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?;
100 fs::write(dir.join("session.json"), session_bytes)?;
101
102 if let Some(entries) = logs.get(session_id) {
104 let mut log_data = String::new();
105 for entry in entries {
106 let line = serde_json::to_string(entry)
107 .map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?;
108 log_data.push_str(&line);
109 log_data.push('\n');
110 }
111 fs::write(dir.join("log.jsonl"), log_data)?;
112 }
113 }
114
115 for (session_id, entries) in &logs {
117 if sessions.contains_key(session_id) {
118 continue;
119 }
120 let dir = sessions_dir.join(session_id);
121 fs::create_dir_all(&dir)?;
122 let mut log_data = String::new();
123 for entry in entries {
124 let line = serde_json::to_string(entry)
125 .map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?;
126 log_data.push_str(&line);
127 log_data.push('\n');
128 }
129 fs::write(dir.join("log.jsonl"), log_data)?;
130 }
131
132 if legacy_sessions.exists() {
134 fs::rename(&legacy_sessions, base_dir.join("sessions.json.migrated"))?;
135 }
136 if legacy_logs.exists() {
137 fs::rename(&legacy_logs, base_dir.join("logs.json.migrated"))?;
138 }
139
140 write_storage_version(base_dir)?;
141 println!(
142 "Migration complete: {} sessions, {} log sets migrated.",
143 sessions.len(),
144 logs.len()
145 );
146 Ok(())
147}
148
149#[cfg(test)]
150mod tests {
151 use super::*;
152 use crate::log_store::EntryKind;
153 use crate::storage::StorageBackend;
154 use macp_core::session::Session;
155 use std::collections::HashSet;
156
157 fn sample_session(id: &str) -> Session {
158 Session::builder(id, "macp.mode.decision.v1", "alice")
159 .ttl_expiry(61_000)
160 .ttl_ms(60_000)
161 .started_at_unix_ms(1_000)
162 .mode_state(vec![1, 2, 3])
163 .participants(vec!["alice".into(), "bob".into()])
164 .seen_message_ids(HashSet::from(["m1".into()]))
165 .intent("test intent")
166 .mode_version("1.0.0")
167 .configuration_version("cfg-1")
168 .policy_version("pol-1")
169 .context_id("test-ctx")
170 .roots(vec![macp_pb::pb::Root {
171 uri: "root://1".into(),
172 name: "r1".into(),
173 }])
174 .build()
175 }
176
177 fn sample_entry(id: &str) -> LogEntry {
178 LogEntry {
179 message_id: id.into(),
180 received_at_ms: 1_700_000_000_000,
181 sender: "alice".into(),
182 message_type: "Message".into(),
183 raw_payload: vec![],
184 entry_kind: EntryKind::Incoming,
185 session_id: String::new(),
186 mode: String::new(),
187 macp_version: String::new(),
188 timestamp_unix_ms: 1_700_000_000_000,
189 bound_mode_version: None,
190 semantics_rev: 0,
191 bound_max_suspend_ms: None,
192 compacted_incoming_ordinals: 0,
193 }
194 }
195
196 #[tokio::test]
197 async fn migration_from_legacy_format() {
198 let dir = tempfile::tempdir().unwrap();
199 let base = dir.path();
200
201 let session = sample_session("s1");
202 let persisted = PersistedSession::from(&session);
203 let sessions_map: HashMap<String, PersistedSession> =
204 [("s1".into(), persisted)].into_iter().collect();
205 fs::write(
206 base.join("sessions.json"),
207 serde_json::to_vec_pretty(&sessions_map).unwrap(),
208 )
209 .unwrap();
210
211 let entries = vec![sample_entry("m1"), sample_entry("m2")];
212 let logs_map: HashMap<String, Vec<LogEntry>> =
213 [("s1".into(), entries)].into_iter().collect();
214 fs::write(
215 base.join("logs.json"),
216 serde_json::to_vec_pretty(&logs_map).unwrap(),
217 )
218 .unwrap();
219
220 migrate_if_needed(base).unwrap();
221
222 assert!(base.join("sessions/s1/session.json").exists());
223 assert!(base.join("sessions/s1/log.jsonl").exists());
224
225 assert!(base.join("sessions.json.migrated").exists());
226 assert!(base.join("logs.json.migrated").exists());
227 assert!(!base.join("sessions.json").exists());
228 assert!(!base.join("logs.json").exists());
229
230 assert_eq!(read_storage_version(base).unwrap(), Some(STORAGE_VERSION));
231
232 let backend = crate::storage::FileBackend::new(base.to_path_buf()).unwrap();
233 let loaded = backend.load_session("s1").await.unwrap().unwrap();
234 assert_eq!(loaded.session_id, "s1");
235 assert_eq!(loaded.ttl_ms, 60_000);
236
237 let log = backend.load_log("s1").await.unwrap();
238 assert_eq!(log.len(), 2);
239 }
240
241 #[test]
242 fn fresh_install_writes_current_version() {
243 let dir = tempfile::tempdir().unwrap();
244 let base = dir.path();
245
246 migrate_if_needed(base).unwrap();
247
248 assert_eq!(read_storage_version(base).unwrap(), Some(STORAGE_VERSION));
249 assert!(!base.join("sessions").exists());
251 }
252
253 #[test]
254 fn already_current_dir_is_noop() {
255 let dir = tempfile::tempdir().unwrap();
256 let base = dir.path();
257 fs::create_dir_all(base.join("sessions")).unwrap();
258 write_storage_version(base).unwrap();
259 let before = fs::read(base.join("storage_version.json")).unwrap();
260
261 migrate_if_needed(base).unwrap();
262
263 assert_eq!(read_storage_version(base).unwrap(), Some(STORAGE_VERSION));
264 let after = fs::read(base.join("storage_version.json")).unwrap();
265 assert_eq!(before, after);
266 assert!(!base.join("sessions.json.migrated").exists());
267 assert!(!base.join("logs.json.migrated").exists());
268 }
269
270 #[test]
271 fn v2_dir_gets_version_bumped_to_current() {
272 let dir = tempfile::tempdir().unwrap();
273 let base = dir.path();
274 fs::create_dir_all(base.join("sessions")).unwrap();
275 fs::write(base.join("storage_version.json"), br#"{"version":2}"#).unwrap();
276 assert_eq!(read_storage_version(base).unwrap(), Some(2));
277
278 migrate_if_needed(base).unwrap();
279
280 assert_eq!(read_storage_version(base).unwrap(), Some(STORAGE_VERSION));
281 }
282
283 #[test]
284 fn sessions_dir_without_version_file_gets_version_written() {
285 let dir = tempfile::tempdir().unwrap();
286 let base = dir.path();
287 fs::create_dir_all(base.join("sessions")).unwrap();
288 assert_eq!(read_storage_version(base).unwrap(), None);
289
290 migrate_if_needed(base).unwrap();
291
292 assert_eq!(read_storage_version(base).unwrap(), Some(STORAGE_VERSION));
293 }
294
295 #[test]
296 fn malformed_legacy_sessions_json_errors_and_preserves_data() {
297 let dir = tempfile::tempdir().unwrap();
298 let base = dir.path();
299 fs::write(base.join("sessions.json"), b"{not json").unwrap();
300
301 let err = migrate_if_needed(base).unwrap_err();
302 assert_eq!(err.kind(), io::ErrorKind::InvalidData);
303
304 assert!(base.join("sessions.json").exists());
307 assert_eq!(fs::read(base.join("sessions.json")).unwrap(), b"{not json");
308 assert!(!base.join("sessions.json.migrated").exists());
309 }
310
311 #[test]
312 fn legacy_logs_without_sessions_creates_log_only_dir() {
313 let dir = tempfile::tempdir().unwrap();
314 let base = dir.path();
315
316 let entries = vec![sample_entry("m1"), sample_entry("m2")];
317 let logs_map: HashMap<String, Vec<LogEntry>> =
318 [("orphan".into(), entries)].into_iter().collect();
319 fs::write(
320 base.join("logs.json"),
321 serde_json::to_vec_pretty(&logs_map).unwrap(),
322 )
323 .unwrap();
324
325 migrate_if_needed(base).unwrap();
326
327 assert!(base.join("sessions/orphan/log.jsonl").exists());
328 assert!(!base.join("sessions/orphan/session.json").exists());
329 assert!(base.join("logs.json.migrated").exists());
330 assert!(!base.join("logs.json").exists());
331 assert_eq!(read_storage_version(base).unwrap(), Some(STORAGE_VERSION));
332
333 let content = fs::read_to_string(base.join("sessions/orphan/log.jsonl")).unwrap();
334 assert_eq!(content.lines().count(), 2);
335 }
336}