Skip to main content

codeswarm_core/
persistence.rs

1//! Versioned persistence for the Rust client.
2//!
3//! The first Rust implementation wrote bare [`AgentEvent`] values to JSONL.
4//! The types in this module deliberately accept that format as schema zero and
5//! provide migration helpers for versioned envelopes.  Session
6//! metadata accepts the legacy flattened shape as well as the current envelope.
7//! Metadata remains flattened so older CodeSwarm state files stay readable.
8
9use std::collections::BTreeSet;
10use std::fmt::{Display, Formatter};
11use std::fs::{self, File, OpenOptions};
12use std::io::{BufRead, BufReader, Write};
13use std::ops::Deref;
14use std::path::{Path, PathBuf};
15use std::sync::mpsc::{self, Receiver, Sender};
16
17use serde_json::{Map, Value};
18
19use crate::AgentEvent;
20
21/// The current on-disk schema for Rust persistence.
22pub const CURRENT_SCHEMA_VERSION: u32 = 1;
23const LEGACY_SCHEMA_VERSION: u32 = 0;
24
25/// Persistence operations report malformed input, unsupported versions, and
26/// I/O failures through this type.
27#[derive(Debug)]
28pub enum PersistenceError {
29    Io(std::io::Error),
30    Malformed { kind: &'static str, detail: String },
31    UnsupportedVersion { kind: &'static str, version: u32 },
32}
33
34impl Display for PersistenceError {
35    fn fmt(&self, formatter: &mut Formatter<'_>) -> std::fmt::Result {
36        match self {
37            Self::Io(error) => write!(formatter, "persistence I/O error: {error}"),
38            Self::Malformed { kind, detail } => write!(formatter, "malformed {kind}: {detail}"),
39            Self::UnsupportedVersion { kind, version } => {
40                write!(formatter, "unsupported {kind} schema version {version}")
41            }
42        }
43    }
44}
45
46impl std::error::Error for PersistenceError {}
47
48impl From<std::io::Error> for PersistenceError {
49    fn from(error: std::io::Error) -> Self {
50        Self::Io(error)
51    }
52}
53
54/// Result of converting a file to the current schema.
55#[derive(Clone, Debug, Eq, PartialEq)]
56pub struct MigrationReport {
57    /// `None` means that no source record/version was observed, including when
58    /// the source file is missing or empty.
59    pub source_version: Option<u32>,
60    pub target_version: u32,
61    pub records: usize,
62    pub changed: bool,
63}
64
65/// Events read from a versioned log, including the source versions observed.
66#[derive(Clone, Debug, Eq, PartialEq)]
67pub struct LoadedEvents {
68    pub events: Vec<AgentEvent>,
69    pub source_versions: BTreeSet<u32>,
70}
71
72/// A JSONL event log that accepts legacy bare events and version-zero envelopes.
73#[derive(Clone, Debug)]
74pub struct VersionedEventLog {
75    path: PathBuf,
76}
77
78impl VersionedEventLog {
79    pub fn open(path: impl Into<PathBuf>) -> Self {
80        Self { path: path.into() }
81    }
82
83    pub fn path(&self) -> &Path {
84        &self.path
85    }
86
87    /// Append one current-schema envelope.
88    pub fn append(&self, event: &AgentEvent) -> Result<(), PersistenceError> {
89        let record = serde_json::json!({
90            "schema_version": CURRENT_SCHEMA_VERSION,
91            "event": event,
92        });
93        let mut file = OpenOptions::new()
94            .create(true)
95            .append(true)
96            .open(&self.path)?;
97        serde_json::to_writer(&mut file, &record)
98            .map_err(|error| malformed("event log", error.to_string()))?;
99        file.write_all(b"\n")?;
100        file.sync_data()?;
101        Ok(())
102    }
103
104    /// Read all events. Missing logs are an empty event stream.
105    pub fn read(&self) -> Result<Vec<AgentEvent>, PersistenceError> {
106        Ok(self.read_with_versions()?.events)
107    }
108
109    pub fn read_with_versions(&self) -> Result<LoadedEvents, PersistenceError> {
110        let file = match File::open(&self.path) {
111            Ok(file) => file,
112            Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
113                return Ok(LoadedEvents {
114                    events: Vec::new(),
115                    source_versions: BTreeSet::new(),
116                });
117            }
118            Err(error) => return Err(error.into()),
119        };
120        let mut events = Vec::new();
121        let mut source_versions = BTreeSet::new();
122        for (line_number, result) in BufReader::new(file).lines().enumerate() {
123            let line = result?;
124            if line.trim().is_empty() {
125                continue;
126            }
127            let (version, event) = parse_event_record(&line, line_number + 1)?;
128            source_versions.insert(version);
129            events.push(event);
130        }
131        Ok(LoadedEvents {
132            events,
133            source_versions,
134        })
135    }
136
137    /// Rewrite legacy records atomically. A malformed record leaves the source
138    /// untouched, allowing an operator to repair the bad log manually.
139    pub fn migrate_in_place(&self) -> Result<MigrationReport, PersistenceError> {
140        let loaded = self.read_with_versions()?;
141        if loaded.events.is_empty() && !self.path.exists() {
142            return Ok(MigrationReport {
143                source_version: None,
144                target_version: CURRENT_SCHEMA_VERSION,
145                records: 0,
146                changed: false,
147            });
148        }
149        let source_version = loaded.source_versions.iter().copied().min();
150        let changed = loaded
151            .source_versions
152            .iter()
153            .any(|version| *version != CURRENT_SCHEMA_VERSION);
154        if changed {
155            atomic_write_event_log(&self.path, &loaded.events)?;
156        }
157        Ok(MigrationReport {
158            source_version,
159            target_version: CURRENT_SCHEMA_VERSION,
160            records: loaded.events.len(),
161            changed,
162        })
163    }
164}
165
166fn parse_event_record(
167    line: &str,
168    line_number: usize,
169) -> Result<(u32, AgentEvent), PersistenceError> {
170    let value: Value = serde_json::from_str(line)
171        .map_err(|error| malformed("event log", format!("line {line_number}: {error}")))?;
172    let object = value.as_object().ok_or_else(|| {
173        malformed(
174            "event log",
175            format!("line {line_number} must be a JSON object"),
176        )
177    })?;
178    let has_envelope = object.contains_key("event")
179        || object.contains_key("schema_version")
180        || object.contains_key("version");
181    let version = if has_envelope {
182        read_version(object, "event log", line_number)?
183    } else {
184        LEGACY_SCHEMA_VERSION
185    };
186    ensure_supported(version, "event log")?;
187    let event_value = object.get("event").unwrap_or(&value);
188    let event = serde_json::from_value(event_value.clone())
189        .map_err(|error| malformed("event log", format!("line {line_number} event: {error}")))?;
190    Ok((version, event))
191}
192
193fn atomic_write_event_log(path: &Path, events: &[AgentEvent]) -> Result<(), PersistenceError> {
194    let temporary = temporary_path(path);
195    let result = (|| {
196        let mut file = OpenOptions::new()
197            .create_new(true)
198            .write(true)
199            .open(&temporary)?;
200        for event in events {
201            let record = serde_json::json!({
202                "schema_version": CURRENT_SCHEMA_VERSION,
203                "event": event,
204            });
205            serde_json::to_writer(&mut file, &record)
206                .map_err(|error| malformed("event log", error.to_string()))?;
207            file.write_all(b"\n")?;
208        }
209        file.sync_all()?;
210        fs::rename(&temporary, path)?;
211        Ok(())
212    })();
213    if result.is_err() {
214        let _ = fs::remove_file(&temporary);
215    }
216    result
217}
218
219/// Metadata imported from either a legacy plain object or the Rust envelope.
220/// The map intentionally retains unknown keys for forward compatibility.
221#[derive(Clone, Debug, Eq, PartialEq)]
222pub struct SessionMetadata {
223    data: Map<String, Value>,
224}
225
226impl SessionMetadata {
227    pub fn new(data: Map<String, Value>) -> Self {
228        Self { data }
229    }
230
231    pub fn empty() -> Self {
232        Self::new(Map::new())
233    }
234
235    pub fn schema_version(&self) -> u32 {
236        CURRENT_SCHEMA_VERSION
237    }
238
239    pub fn get(&self, key: &str) -> Option<&Value> {
240        self.data.get(key)
241    }
242
243    /// Replace one metadata value while retaining all other keys.
244    pub fn insert(&mut self, key: impl Into<String>, value: impl Into<Value>) -> Option<Value> {
245        self.data.insert(key.into(), value.into())
246    }
247
248    /// Remove one metadata value, returning the previous value when present.
249    pub fn remove(&mut self, key: &str) -> Option<Value> {
250        self.data.remove(key)
251    }
252
253    pub fn as_object(&self) -> &Map<String, Value> {
254        &self.data
255    }
256
257    /// Flattened JSON keeps existing `roster`/`agent_data` keys addressable.
258    pub fn to_value(&self) -> Value {
259        let mut object = self.data.clone();
260        object.insert("schema_version".into(), Value::from(CURRENT_SCHEMA_VERSION));
261        Value::Object(object)
262    }
263
264    pub fn to_json(&self) -> Result<String, PersistenceError> {
265        serde_json::to_string(&self.to_value())
266            .map_err(|error| malformed("session metadata", error.to_string()))
267    }
268}
269
270/// A loaded metadata value carries the version from which it was imported.
271#[derive(Clone, Debug, Eq, PartialEq)]
272pub struct LoadedSessionMetadata {
273    pub metadata: SessionMetadata,
274    pub source_version: u32,
275}
276
277impl Deref for LoadedSessionMetadata {
278    type Target = SessionMetadata;
279
280    fn deref(&self) -> &Self::Target {
281        &self.metadata
282    }
283}
284
285/// File-backed session metadata with schema migration.
286#[derive(Clone, Debug)]
287pub struct SessionMetadataStore {
288    path: PathBuf,
289}
290
291impl SessionMetadataStore {
292    pub fn open(path: impl Into<PathBuf>) -> Self {
293        Self { path: path.into() }
294    }
295
296    pub fn path(&self) -> &Path {
297        &self.path
298    }
299
300    pub fn read(&self) -> Result<Option<LoadedSessionMetadata>, PersistenceError> {
301        let raw = match fs::read_to_string(&self.path) {
302            Ok(raw) => raw,
303            Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(None),
304            Err(error) => return Err(error.into()),
305        };
306        let value: Value = serde_json::from_str(&raw)
307            .map_err(|error| malformed("session metadata", error.to_string()))?;
308        let object = value
309            .as_object()
310            .ok_or_else(|| malformed("session metadata", "expected a JSON object".into()))?;
311        let source_version = read_version(object, "session metadata", 0)?;
312        ensure_supported(source_version, "session metadata")?;
313        let data = if let Some(metadata) = object.get("metadata") {
314            metadata
315                .as_object()
316                .ok_or_else(|| {
317                    malformed(
318                        "session metadata",
319                        "metadata envelope must be an object".into(),
320                    )
321                })?
322                .clone()
323        } else {
324            object
325                .iter()
326                .filter(|(key, _)| key.as_str() != "schema_version" && key.as_str() != "version")
327                .map(|(key, value)| (key.clone(), value.clone()))
328                .collect()
329        };
330        Ok(Some(LoadedSessionMetadata {
331            metadata: SessionMetadata::new(data),
332            source_version,
333        }))
334    }
335
336    pub fn load(&self) -> Result<Option<LoadedSessionMetadata>, PersistenceError> {
337        self.read()
338    }
339
340    /// Read the current snapshot, apply an in-memory edit, and atomically
341    /// publish the replacement. A missing snapshot is treated as empty
342    /// metadata; malformed or newer snapshots are left untouched and
343    /// returned as errors.
344    pub fn update<F>(&self, edit: F) -> Result<(), PersistenceError>
345    where
346        F: FnOnce(&mut SessionMetadata),
347    {
348        let mut metadata = self
349            .read()?
350            .map(|loaded| loaded.metadata)
351            .unwrap_or_else(SessionMetadata::empty);
352        edit(&mut metadata);
353        self.write(&metadata)
354    }
355
356    pub fn write(&self, metadata: &SessionMetadata) -> Result<(), PersistenceError> {
357        let json = metadata.to_json()?;
358        let temporary = temporary_path(&self.path);
359        let result = (|| {
360            if let Some(parent) = self.path.parent()
361                && !parent.as_os_str().is_empty()
362            {
363                fs::create_dir_all(parent)?;
364            }
365            let mut file = OpenOptions::new()
366                .create_new(true)
367                .write(true)
368                .open(&temporary)?;
369            file.write_all(format!("{json}\n").as_bytes())?;
370            file.sync_all()?;
371            fs::rename(&temporary, &self.path)?;
372            Ok(())
373        })();
374        if result.is_err() {
375            let _ = fs::remove_file(&temporary);
376        }
377        result
378    }
379
380    pub fn migrate_in_place(&self) -> Result<MigrationReport, PersistenceError> {
381        let Some(loaded) = self.read()? else {
382            return Ok(MigrationReport {
383                source_version: None,
384                target_version: CURRENT_SCHEMA_VERSION,
385                records: 0,
386                changed: false,
387            });
388        };
389        let changed = loaded.source_version != CURRENT_SCHEMA_VERSION
390            || serde_json::from_str::<Value>(&fs::read_to_string(&self.path)?)
391                .ok()
392                .and_then(|value| {
393                    value
394                        .as_object()
395                        .map(|object| object.contains_key("metadata"))
396                })
397                .unwrap_or(false);
398        if changed {
399            self.write(&loaded.metadata)?;
400        }
401        Ok(MigrationReport {
402            source_version: Some(loaded.source_version),
403            target_version: CURRENT_SCHEMA_VERSION,
404            records: 1,
405            changed,
406        })
407    }
408
409    /// Start a background metadata writer. Runtime event loops should enqueue
410    /// snapshots through this handle and call [`BufferedSessionMetadataStore::flush`]
411    /// only at lifecycle boundaries; atomic writes and fsync therefore never
412    /// block terminal input or rendering.
413    pub fn buffered(&self) -> std::io::Result<BufferedSessionMetadataStore> {
414        BufferedSessionMetadataStore::open(self.path.clone())
415    }
416}
417
418enum MetadataCommand {
419    Write(SessionMetadata),
420    Flush(Sender<Result<(), PersistenceError>>),
421    Shutdown(Sender<()>),
422}
423
424/// Background writer for runtime session metadata snapshots.
425///
426/// Each queued snapshot is written in order, and every write replaces the
427/// previous file atomically. The handle itself only performs channel sends;
428/// filesystem work happens on its worker thread. `flush` is the explicit
429/// durability boundary used before a session exits.
430pub struct BufferedSessionMetadataStore {
431    sender: Sender<MetadataCommand>,
432    worker: Option<std::thread::JoinHandle<Result<(), PersistenceError>>>,
433}
434
435impl std::fmt::Debug for BufferedSessionMetadataStore {
436    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
437        formatter
438            .debug_struct("BufferedSessionMetadataStore")
439            .field("worker_running", &self.worker.is_some())
440            .finish_non_exhaustive()
441    }
442}
443
444impl BufferedSessionMetadataStore {
445    fn open(path: PathBuf) -> std::io::Result<Self> {
446        let (sender, receiver) = mpsc::channel();
447        let worker = std::thread::Builder::new()
448            .name("codeswarm-session-metadata".into())
449            .spawn(move || metadata_worker(path, receiver))?;
450        Ok(Self {
451            sender,
452            worker: Some(worker),
453        })
454    }
455
456    /// Queue a complete metadata snapshot without doing filesystem I/O in
457    /// the caller.
458    pub fn write(&self, metadata: SessionMetadata) -> Result<(), PersistenceError> {
459        self.sender
460            .send(MetadataCommand::Write(metadata))
461            .map_err(|_| {
462                PersistenceError::Io(std::io::Error::new(
463                    std::io::ErrorKind::BrokenPipe,
464                    "session metadata writer stopped",
465                ))
466            })
467    }
468
469    /// Drain queued snapshots and wait until the latest one is durable.
470    pub fn flush(&self) -> Result<(), PersistenceError> {
471        let (reply, result) = mpsc::channel();
472        self.sender
473            .send(MetadataCommand::Flush(reply))
474            .map_err(|_| {
475                PersistenceError::Io(std::io::Error::new(
476                    std::io::ErrorKind::BrokenPipe,
477                    "session metadata writer stopped",
478                ))
479            })?;
480        result.recv().map_err(|_| {
481            PersistenceError::Io(std::io::Error::new(
482                std::io::ErrorKind::BrokenPipe,
483                "session metadata writer stopped",
484            ))
485        })?
486    }
487}
488
489impl Drop for BufferedSessionMetadataStore {
490    fn drop(&mut self) {
491        let Some(worker) = self.worker.take() else {
492            return;
493        };
494        let (reply, result) = mpsc::channel();
495        if self.sender.send(MetadataCommand::Shutdown(reply)).is_ok() {
496            let _ = result.recv();
497        }
498        let _ = worker.join();
499    }
500}
501
502fn metadata_worker(
503    path: PathBuf,
504    receiver: Receiver<MetadataCommand>,
505) -> Result<(), PersistenceError> {
506    let store = SessionMetadataStore::open(path);
507    let mut first_error = None;
508    while let Ok(command) = receiver.recv() {
509        match command {
510            MetadataCommand::Write(metadata) => {
511                if first_error.is_none()
512                    && let Err(error) = store.write(&metadata)
513                {
514                    first_error = Some(error);
515                }
516            }
517            MetadataCommand::Flush(reply) => {
518                let _ = reply.send(match first_error.take() {
519                    Some(error) => Err(error),
520                    None => Ok(()),
521                });
522            }
523            MetadataCommand::Shutdown(reply) => {
524                let result = first_error.take();
525                let _ = reply.send(());
526                return result.map_or(Ok(()), Err);
527            }
528        }
529    }
530    first_error.map_or(Ok(()), Err)
531}
532
533fn read_version(
534    object: &Map<String, Value>,
535    kind: &'static str,
536    line_number: usize,
537) -> Result<u32, PersistenceError> {
538    let schema = object
539        .get("schema_version")
540        .or_else(|| object.get("version"));
541    let Some(schema) = schema else {
542        return Ok(LEGACY_SCHEMA_VERSION);
543    };
544    let version = schema.as_u64().ok_or_else(|| {
545        malformed(
546            kind,
547            if line_number == 0 {
548                "schema_version must be an unsigned integer".into()
549            } else {
550                format!("line {line_number} schema_version must be an unsigned integer")
551            },
552        )
553    })?;
554    u32::try_from(version).map_err(|_| malformed(kind, "schema_version is too large".into()))
555}
556
557fn ensure_supported(version: u32, kind: &'static str) -> Result<(), PersistenceError> {
558    if version > CURRENT_SCHEMA_VERSION {
559        return Err(PersistenceError::UnsupportedVersion { kind, version });
560    }
561    Ok(())
562}
563
564fn malformed(kind: &'static str, detail: String) -> PersistenceError {
565    PersistenceError::Malformed { kind, detail }
566}
567
568fn temporary_path(path: &Path) -> PathBuf {
569    let file_name = path
570        .file_name()
571        .and_then(|name| name.to_str())
572        .unwrap_or("data");
573    path.with_file_name(format!(".{file_name}.migration-{}.tmp", std::process::id()))
574}
575
576/// Convert a legacy session metadata blob to the current Rust metadata value.
577pub fn import_legacy_session_metadata(
578    value: Option<&str>,
579) -> Result<Option<LoadedSessionMetadata>, PersistenceError> {
580    let Some(value) = value else { return Ok(None) };
581    let parsed: Value = serde_json::from_str(value)
582        .map_err(|error| malformed("session metadata", error.to_string()))?;
583    let object = parsed
584        .as_object()
585        .ok_or_else(|| malformed("session metadata", "expected a JSON object".into()))?;
586    let source_version = read_version(object, "session metadata", 0)?;
587    ensure_supported(source_version, "session metadata")?;
588    let data = object
589        .iter()
590        .filter(|(key, _)| key.as_str() != "schema_version" && key.as_str() != "version")
591        .map(|(key, value)| (key.clone(), value.clone()))
592        .collect();
593    Ok(Some(LoadedSessionMetadata {
594        metadata: SessionMetadata::new(data),
595        source_version,
596    }))
597}
598
599#[cfg(test)]
600mod tests {
601    use std::time::{SystemTime, UNIX_EPOCH};
602
603    use serde_json::{Value, json};
604
605    use super::{
606        CURRENT_SCHEMA_VERSION, PersistenceError, SessionMetadata, SessionMetadataStore,
607        VersionedEventLog, import_legacy_session_metadata,
608    };
609    use crate::AgentEvent;
610
611    fn temp_path(suffix: &str) -> std::path::PathBuf {
612        let unique = SystemTime::now()
613            .duration_since(UNIX_EPOCH)
614            .expect("clock")
615            .as_nanos();
616        std::env::temp_dir().join(format!("codeswarm-persistence-{unique}-{suffix}"))
617    }
618
619    fn event() -> AgentEvent {
620        AgentEvent::Text {
621            slot: 0,
622            text: "hello".into(),
623        }
624    }
625
626    #[test]
627    fn missing_event_log_and_metadata_are_empty() {
628        let event_log = VersionedEventLog::open(temp_path("events.jsonl"));
629        assert!(event_log.read().expect("missing log").is_empty());
630        let metadata = SessionMetadataStore::open(temp_path("metadata.json"));
631        assert!(metadata.read().expect("missing metadata").is_none());
632        assert!(
633            !metadata
634                .migrate_in_place()
635                .expect("missing migration")
636                .changed
637        );
638    }
639
640    #[test]
641    fn malformed_data_is_rejected_without_rewriting_source() {
642        let event_path = temp_path("malformed-events.jsonl");
643        std::fs::write(&event_path, "not-json\n").expect("write");
644        let event_log = VersionedEventLog::open(&event_path);
645        assert!(matches!(
646            event_log.read(),
647            Err(PersistenceError::Malformed { .. })
648        ));
649        assert_eq!(
650            std::fs::read_to_string(&event_path).expect("read"),
651            "not-json\n"
652        );
653        std::fs::remove_file(event_path).expect("cleanup");
654
655        let metadata_path = temp_path("malformed-metadata.json");
656        std::fs::write(&metadata_path, "[]").expect("write");
657        let metadata = SessionMetadataStore::open(&metadata_path);
658        assert!(matches!(
659            metadata.read(),
660            Err(PersistenceError::Malformed { .. })
661        ));
662        assert_eq!(std::fs::read_to_string(&metadata_path).expect("read"), "[]");
663        std::fs::remove_file(metadata_path).expect("cleanup");
664    }
665
666    #[test]
667    fn old_bare_event_log_migrates_to_current_envelope() {
668        let path = temp_path("old-events.jsonl");
669        std::fs::write(
670            &path,
671            serde_json::to_string(&event()).expect("event") + "\n",
672        )
673        .expect("write");
674        let log = VersionedEventLog::open(&path);
675        let report = log.migrate_in_place().expect("migrate");
676        assert_eq!(report.source_version, Some(0));
677        assert!(report.changed);
678        assert_eq!(log.read().expect("read"), vec![event()]);
679        let migrated = std::fs::read_to_string(&path).expect("read raw");
680        assert!(migrated.contains(&format!("\"schema_version\":{CURRENT_SCHEMA_VERSION}")));
681        std::fs::remove_file(path).expect("cleanup");
682    }
683
684    #[test]
685    fn old_metadata_migrates_and_preserves_unknown_keys() {
686        let path = temp_path("old-metadata.json");
687        std::fs::write(
688            &path,
689            r#"{"roster":["openai.com"],"agent_data":{"name":"Codex"}}"#,
690        )
691        .expect("write");
692        let store = SessionMetadataStore::open(&path);
693        let loaded = store.read().expect("read").expect("metadata");
694        assert_eq!(loaded.source_version, 0);
695        assert_eq!(loaded.get("roster"), Some(&json!(["openai.com"])));
696        let report = store.migrate_in_place().expect("migrate");
697        assert!(report.changed);
698        let migrated = std::fs::read_to_string(&path).expect("read");
699        assert!(migrated.contains("\"roster\""));
700        assert!(migrated.contains("\"schema_version\":1"));
701        std::fs::remove_file(path).expect("cleanup");
702    }
703
704    #[test]
705    fn current_event_and_metadata_data_is_not_rewritten() {
706        let event_path = temp_path("current-events.jsonl");
707        let log = VersionedEventLog::open(&event_path);
708        log.append(&event()).expect("append");
709        let before = std::fs::read_to_string(&event_path).expect("read");
710        let report = log.migrate_in_place().expect("migrate");
711        assert_eq!(report.source_version, Some(1));
712        assert!(!report.changed);
713        assert_eq!(std::fs::read_to_string(&event_path).expect("read"), before);
714        std::fs::remove_file(event_path).expect("cleanup");
715
716        let metadata_path = temp_path("current-metadata.json");
717        let store = SessionMetadataStore::open(&metadata_path);
718        let mut data = serde_json::Map::new();
719        data.insert("roster".into(), json!(["agy"]));
720        store.write(&SessionMetadata::new(data)).expect("write");
721        let report = store.migrate_in_place().expect("migrate");
722        assert_eq!(report.source_version, Some(1));
723        assert!(!report.changed);
724        std::fs::remove_file(metadata_path).expect("cleanup");
725    }
726
727    #[test]
728    fn metadata_write_creates_parent_and_replaces_previous_snapshot() {
729        let path = temp_path("nested").join("session.json");
730        let store = SessionMetadataStore::open(&path);
731
732        let mut first = serde_json::Map::new();
733        first.insert("owner".into(), json!("Claude"));
734        store
735            .write(&SessionMetadata::new(first))
736            .expect("first write");
737
738        let mut second = serde_json::Map::new();
739        second.insert("owner".into(), json!("Codex"));
740        second.insert("roster".into(), json!(["openai.com", "claude.ai"]));
741        store
742            .write(&SessionMetadata::new(second))
743            .expect("replacement write");
744
745        let loaded = store.read().expect("read").expect("snapshot");
746        assert_eq!(loaded.source_version, CURRENT_SCHEMA_VERSION);
747        assert_eq!(loaded.get("owner"), Some(&json!("Codex")));
748        assert_eq!(
749            loaded.get("roster"),
750            Some(&json!(["openai.com", "claude.ai"]))
751        );
752        std::fs::remove_dir_all(path.parent().expect("parent")).expect("cleanup");
753    }
754
755    #[test]
756    fn metadata_update_merges_current_values_and_creates_missing_snapshot() {
757        let path = temp_path("update").join("session.json");
758        let store = SessionMetadataStore::open(&path);
759        store
760            .update(|metadata| {
761                metadata.insert("roster", json!(["claude.ai"]));
762            })
763            .expect("create snapshot");
764        store
765            .update(|metadata| {
766                metadata.insert("owner", json!("Claude"));
767                metadata.remove("roster");
768            })
769            .expect("merge snapshot");
770        let loaded = store.read().expect("read").expect("snapshot");
771        assert_eq!(loaded.get("owner"), Some(&json!("Claude")));
772        assert_eq!(loaded.get("roster"), None);
773        std::fs::remove_dir_all(path.parent().expect("parent")).expect("cleanup");
774    }
775
776    #[test]
777    fn buffered_metadata_writes_are_durable_at_flush_and_drop() {
778        let path = temp_path("buffered").join("session.json");
779        let store = SessionMetadataStore::open(&path);
780        let writer = store.buffered().expect("writer");
781        let mut metadata = SessionMetadata::empty();
782        metadata.insert("roster", json!(["claude.ai", "openai.com"]));
783        writer.write(metadata).expect("queue snapshot");
784        writer.flush().expect("flush snapshot");
785        let loaded = store.read().expect("read").expect("snapshot");
786        assert_eq!(
787            loaded.get("roster"),
788            Some(&json!(["claude.ai", "openai.com"]))
789        );
790        drop(writer);
791        std::fs::remove_dir_all(path.parent().expect("parent")).expect("cleanup");
792    }
793
794    #[test]
795    fn legacy_import_accepts_missing_and_current_metadata() {
796        assert!(
797            import_legacy_session_metadata(None)
798                .expect("missing")
799                .is_none()
800        );
801        let loaded = import_legacy_session_metadata(Some(
802            r#"{"schema_version":1,"roster":["agy"],"agent_data":{"name":"Agy"}}"#,
803        ))
804        .expect("current")
805        .expect("metadata");
806        assert_eq!(loaded.source_version, 1);
807        assert_eq!(
808            loaded
809                .get("agent_data")
810                .and_then(Value::as_object)
811                .and_then(|m| m.get("name"))
812                .and_then(Value::as_str),
813            Some("Agy")
814        );
815    }
816
817    #[test]
818    fn future_versions_are_rejected() {
819        let event_path = temp_path("future-events.jsonl");
820        std::fs::write(
821            &event_path,
822            serde_json::to_string(&json!({"schema_version": 99, "event": event()})).expect("json"),
823        )
824        .expect("write");
825        assert!(matches!(
826            VersionedEventLog::open(&event_path).read(),
827            Err(PersistenceError::UnsupportedVersion { version: 99, .. })
828        ));
829        std::fs::remove_file(event_path).expect("cleanup");
830    }
831}