Skip to main content

macp_storage/storage/
migration.rs

1use 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    // Already at current version or fresh install (no legacy files, sessions dir exists)
41    let current_version = read_storage_version(base_dir)?;
42    if sessions_dir.exists() && !legacy_sessions.exists() && !legacy_logs.exists() {
43        // v2 → v3: no-op data migration, just bump version.  New LogEntry fields
44        // use #[serde(default)] so existing v2 JSONL lines deserialize fine.
45        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    // Already migrated from v1
57    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    // Load legacy sessions
65    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    // Load legacy logs
78    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    // Migrate each session
93    for (session_id, persisted) in &sessions {
94        let dir = sessions_dir.join(session_id);
95        fs::create_dir_all(&dir)?;
96
97        // Write session.json
98        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        // Write log.jsonl
103        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    // Also migrate logs for sessions that only appear in logs (not in sessions.json)
116    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    // Backup old files instead of deleting
133    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        // A fresh install has nothing to migrate; no sessions dir is required.
250        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        // The malformed legacy file must be left in place for operator
305        // inspection, not renamed away or deleted.
306        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}