Skip to main content

codeswarm_adapters/
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        self.buffered_with_errors(|_| {})
415    }
416
417    pub fn buffered_with_errors(
418        &self,
419        on_error: impl Fn(String) + Send + 'static,
420    ) -> std::io::Result<BufferedSessionMetadataStore> {
421        BufferedSessionMetadataStore::open(self.path.clone(), Box::new(on_error))
422    }
423}
424
425enum MetadataCommand {
426    Write(SessionMetadata),
427    Flush(Sender<Result<(), PersistenceError>>),
428    Shutdown(Sender<()>),
429}
430
431/// Background writer for runtime session metadata snapshots.
432///
433/// Each queued snapshot is written in order, and every write replaces the
434/// previous file atomically. The handle itself only performs channel sends;
435/// filesystem work happens on its worker thread. `flush` is the explicit
436/// durability boundary used before a session exits.
437pub struct BufferedSessionMetadataStore {
438    sender: Sender<MetadataCommand>,
439    worker: Option<std::thread::JoinHandle<Result<(), PersistenceError>>>,
440}
441
442impl std::fmt::Debug for BufferedSessionMetadataStore {
443    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
444        formatter
445            .debug_struct("BufferedSessionMetadataStore")
446            .field("worker_running", &self.worker.is_some())
447            .finish_non_exhaustive()
448    }
449}
450
451impl BufferedSessionMetadataStore {
452    fn open(path: PathBuf, on_error: Box<dyn Fn(String) + Send>) -> std::io::Result<Self> {
453        let (sender, receiver) = mpsc::channel();
454        let worker = std::thread::Builder::new()
455            .name("codeswarm-session-metadata".into())
456            .spawn(move || metadata_worker(path, receiver, on_error))?;
457        Ok(Self {
458            sender,
459            worker: Some(worker),
460        })
461    }
462
463    /// Queue a complete metadata snapshot without doing filesystem I/O in
464    /// the caller.
465    pub fn write(&self, metadata: SessionMetadata) -> Result<(), PersistenceError> {
466        self.sender
467            .send(MetadataCommand::Write(metadata))
468            .map_err(|_| {
469                PersistenceError::Io(std::io::Error::new(
470                    std::io::ErrorKind::BrokenPipe,
471                    "session metadata writer stopped",
472                ))
473            })
474    }
475
476    /// Drain queued snapshots and wait until the latest one is durable.
477    pub fn flush(&self) -> Result<(), PersistenceError> {
478        let (reply, result) = mpsc::channel();
479        self.sender
480            .send(MetadataCommand::Flush(reply))
481            .map_err(|_| {
482                PersistenceError::Io(std::io::Error::new(
483                    std::io::ErrorKind::BrokenPipe,
484                    "session metadata writer stopped",
485                ))
486            })?;
487        result.recv().map_err(|_| {
488            PersistenceError::Io(std::io::Error::new(
489                std::io::ErrorKind::BrokenPipe,
490                "session metadata writer stopped",
491            ))
492        })?
493    }
494}
495
496impl Drop for BufferedSessionMetadataStore {
497    fn drop(&mut self) {
498        let Some(worker) = self.worker.take() else {
499            return;
500        };
501        let (reply, result) = mpsc::channel();
502        if self.sender.send(MetadataCommand::Shutdown(reply)).is_ok() {
503            let _ = result.recv();
504        }
505        let _ = worker.join();
506    }
507}
508
509fn metadata_worker(
510    path: PathBuf,
511    receiver: Receiver<MetadataCommand>,
512    on_error: Box<dyn Fn(String) + Send>,
513) -> Result<(), PersistenceError> {
514    let store = SessionMetadataStore::open(path);
515    let mut first_error = None;
516    let mut pending = None;
517    loop {
518        let command = if pending.is_some() {
519            match receiver.recv_timeout(std::time::Duration::from_secs(1)) {
520                Ok(command) => command,
521                Err(mpsc::RecvTimeoutError::Timeout) => {
522                    if let Some(metadata) = &pending
523                        && store.write(metadata).is_ok()
524                    {
525                        pending = None;
526                        first_error = None;
527                    }
528                    continue;
529                }
530                Err(mpsc::RecvTimeoutError::Disconnected) => break,
531            }
532        } else {
533            match receiver.recv() {
534                Ok(command) => command,
535                Err(_) => break,
536            }
537        };
538        match command {
539            MetadataCommand::Write(metadata) => match store.write(&metadata) {
540                Ok(()) => {
541                    pending = None;
542                    first_error = None;
543                }
544                Err(error) => {
545                    if first_error.is_none() {
546                        on_error(error.to_string());
547                    }
548                    pending = Some(metadata);
549                    first_error = Some(error);
550                }
551            },
552            MetadataCommand::Flush(reply) => {
553                if let Some(metadata) = &pending {
554                    match store.write(metadata) {
555                        Ok(()) => {
556                            pending = None;
557                            first_error = None;
558                        }
559                        Err(error) => {
560                            first_error = Some(error);
561                        }
562                    }
563                }
564                let _ = reply.send(match first_error.take() {
565                    Some(error) => Err(error),
566                    None => Ok(()),
567                });
568            }
569            MetadataCommand::Shutdown(reply) => {
570                if let Some(metadata) = &pending {
571                    first_error = store.write(metadata).err();
572                }
573                let result = first_error.take();
574                let _ = reply.send(());
575                return result.map_or(Ok(()), Err);
576            }
577        }
578    }
579    first_error.map_or(Ok(()), Err)
580}
581
582fn read_version(
583    object: &Map<String, Value>,
584    kind: &'static str,
585    line_number: usize,
586) -> Result<u32, PersistenceError> {
587    let schema = object
588        .get("schema_version")
589        .or_else(|| object.get("version"));
590    let Some(schema) = schema else {
591        return Ok(LEGACY_SCHEMA_VERSION);
592    };
593    let version = schema.as_u64().ok_or_else(|| {
594        malformed(
595            kind,
596            if line_number == 0 {
597                "schema_version must be an unsigned integer".into()
598            } else {
599                format!("line {line_number} schema_version must be an unsigned integer")
600            },
601        )
602    })?;
603    u32::try_from(version).map_err(|_| malformed(kind, "schema_version is too large".into()))
604}
605
606fn ensure_supported(version: u32, kind: &'static str) -> Result<(), PersistenceError> {
607    if version > CURRENT_SCHEMA_VERSION {
608        return Err(PersistenceError::UnsupportedVersion { kind, version });
609    }
610    Ok(())
611}
612
613fn malformed(kind: &'static str, detail: String) -> PersistenceError {
614    PersistenceError::Malformed { kind, detail }
615}
616
617fn temporary_path(path: &Path) -> PathBuf {
618    let file_name = path
619        .file_name()
620        .and_then(|name| name.to_str())
621        .unwrap_or("data");
622    path.with_file_name(format!(".{file_name}.migration-{}.tmp", std::process::id()))
623}
624
625/// Convert a legacy session metadata blob to the current Rust metadata value.
626pub fn import_legacy_session_metadata(
627    value: Option<&str>,
628) -> Result<Option<LoadedSessionMetadata>, PersistenceError> {
629    let Some(value) = value else { return Ok(None) };
630    let parsed: Value = serde_json::from_str(value)
631        .map_err(|error| malformed("session metadata", error.to_string()))?;
632    let object = parsed
633        .as_object()
634        .ok_or_else(|| malformed("session metadata", "expected a JSON object".into()))?;
635    let source_version = read_version(object, "session metadata", 0)?;
636    ensure_supported(source_version, "session metadata")?;
637    let data = object
638        .iter()
639        .filter(|(key, _)| key.as_str() != "schema_version" && key.as_str() != "version")
640        .map(|(key, value)| (key.clone(), value.clone()))
641        .collect();
642    Ok(Some(LoadedSessionMetadata {
643        metadata: SessionMetadata::new(data),
644        source_version,
645    }))
646}
647
648#[cfg(test)]
649mod tests {
650    use std::time::{SystemTime, UNIX_EPOCH};
651
652    use serde_json::{Value, json};
653
654    use super::{
655        CURRENT_SCHEMA_VERSION, PersistenceError, SessionMetadata, SessionMetadataStore,
656        VersionedEventLog, import_legacy_session_metadata,
657    };
658    use crate::AgentEvent;
659
660    fn temp_path(suffix: &str) -> std::path::PathBuf {
661        let unique = SystemTime::now()
662            .duration_since(UNIX_EPOCH)
663            .expect("clock")
664            .as_nanos();
665        std::env::temp_dir().join(format!("codeswarm-persistence-{unique}-{suffix}"))
666    }
667
668    fn event() -> AgentEvent {
669        AgentEvent::Text {
670            slot: 0,
671            text: "hello".into(),
672        }
673    }
674
675    #[test]
676    fn missing_event_log_and_metadata_are_empty() {
677        let event_log = VersionedEventLog::open(temp_path("events.jsonl"));
678        assert!(event_log.read().expect("missing log").is_empty());
679        let metadata = SessionMetadataStore::open(temp_path("metadata.json"));
680        assert!(metadata.read().expect("missing metadata").is_none());
681        assert!(
682            !metadata
683                .migrate_in_place()
684                .expect("missing migration")
685                .changed
686        );
687    }
688
689    #[test]
690    fn malformed_data_is_rejected_without_rewriting_source() {
691        let event_path = temp_path("malformed-events.jsonl");
692        std::fs::write(&event_path, "not-json\n").expect("write");
693        let event_log = VersionedEventLog::open(&event_path);
694        assert!(matches!(
695            event_log.read(),
696            Err(PersistenceError::Malformed { .. })
697        ));
698        assert_eq!(
699            std::fs::read_to_string(&event_path).expect("read"),
700            "not-json\n"
701        );
702        std::fs::remove_file(event_path).expect("cleanup");
703
704        let metadata_path = temp_path("malformed-metadata.json");
705        std::fs::write(&metadata_path, "[]").expect("write");
706        let metadata = SessionMetadataStore::open(&metadata_path);
707        assert!(matches!(
708            metadata.read(),
709            Err(PersistenceError::Malformed { .. })
710        ));
711        assert_eq!(std::fs::read_to_string(&metadata_path).expect("read"), "[]");
712        std::fs::remove_file(metadata_path).expect("cleanup");
713    }
714
715    #[test]
716    fn old_bare_event_log_migrates_to_current_envelope() {
717        let path = temp_path("old-events.jsonl");
718        std::fs::write(
719            &path,
720            serde_json::to_string(&event()).expect("event") + "\n",
721        )
722        .expect("write");
723        let log = VersionedEventLog::open(&path);
724        let report = log.migrate_in_place().expect("migrate");
725        assert_eq!(report.source_version, Some(0));
726        assert!(report.changed);
727        assert_eq!(log.read().expect("read"), vec![event()]);
728        let migrated = std::fs::read_to_string(&path).expect("read raw");
729        assert!(migrated.contains(&format!("\"schema_version\":{CURRENT_SCHEMA_VERSION}")));
730        std::fs::remove_file(path).expect("cleanup");
731    }
732
733    #[test]
734    fn old_metadata_migrates_and_preserves_unknown_keys() {
735        let path = temp_path("old-metadata.json");
736        std::fs::write(
737            &path,
738            r#"{"roster":["openai.com"],"agent_data":{"name":"Codex"}}"#,
739        )
740        .expect("write");
741        let store = SessionMetadataStore::open(&path);
742        let loaded = store.read().expect("read").expect("metadata");
743        assert_eq!(loaded.source_version, 0);
744        assert_eq!(loaded.get("roster"), Some(&json!(["openai.com"])));
745        let report = store.migrate_in_place().expect("migrate");
746        assert!(report.changed);
747        let migrated = std::fs::read_to_string(&path).expect("read");
748        assert!(migrated.contains("\"roster\""));
749        assert!(migrated.contains("\"schema_version\":1"));
750        std::fs::remove_file(path).expect("cleanup");
751    }
752
753    #[test]
754    fn current_event_and_metadata_data_is_not_rewritten() {
755        let event_path = temp_path("current-events.jsonl");
756        let log = VersionedEventLog::open(&event_path);
757        log.append(&event()).expect("append");
758        let before = std::fs::read_to_string(&event_path).expect("read");
759        let report = log.migrate_in_place().expect("migrate");
760        assert_eq!(report.source_version, Some(1));
761        assert!(!report.changed);
762        assert_eq!(std::fs::read_to_string(&event_path).expect("read"), before);
763        std::fs::remove_file(event_path).expect("cleanup");
764
765        let metadata_path = temp_path("current-metadata.json");
766        let store = SessionMetadataStore::open(&metadata_path);
767        let mut data = serde_json::Map::new();
768        data.insert("roster".into(), json!(["agy"]));
769        store.write(&SessionMetadata::new(data)).expect("write");
770        let report = store.migrate_in_place().expect("migrate");
771        assert_eq!(report.source_version, Some(1));
772        assert!(!report.changed);
773        std::fs::remove_file(metadata_path).expect("cleanup");
774    }
775
776    #[test]
777    fn metadata_write_creates_parent_and_replaces_previous_snapshot() {
778        let path = temp_path("nested").join("session.json");
779        let store = SessionMetadataStore::open(&path);
780
781        let mut first = serde_json::Map::new();
782        first.insert("owner".into(), json!("Claude"));
783        store
784            .write(&SessionMetadata::new(first))
785            .expect("first write");
786
787        let mut second = serde_json::Map::new();
788        second.insert("owner".into(), json!("Codex"));
789        second.insert("roster".into(), json!(["openai.com", "claude.ai"]));
790        store
791            .write(&SessionMetadata::new(second))
792            .expect("replacement write");
793
794        let loaded = store.read().expect("read").expect("snapshot");
795        assert_eq!(loaded.source_version, CURRENT_SCHEMA_VERSION);
796        assert_eq!(loaded.get("owner"), Some(&json!("Codex")));
797        assert_eq!(
798            loaded.get("roster"),
799            Some(&json!(["openai.com", "claude.ai"]))
800        );
801        std::fs::remove_dir_all(path.parent().expect("parent")).expect("cleanup");
802    }
803
804    #[test]
805    fn metadata_update_merges_current_values_and_creates_missing_snapshot() {
806        let path = temp_path("update").join("session.json");
807        let store = SessionMetadataStore::open(&path);
808        store
809            .update(|metadata| {
810                metadata.insert("roster", json!(["claude.ai"]));
811            })
812            .expect("create snapshot");
813        store
814            .update(|metadata| {
815                metadata.insert("owner", json!("Claude"));
816                metadata.remove("roster");
817            })
818            .expect("merge snapshot");
819        let loaded = store.read().expect("read").expect("snapshot");
820        assert_eq!(loaded.get("owner"), Some(&json!("Claude")));
821        assert_eq!(loaded.get("roster"), None);
822        std::fs::remove_dir_all(path.parent().expect("parent")).expect("cleanup");
823    }
824
825    #[test]
826    fn failed_metadata_snapshot_is_retained_and_newer_writes_recover() {
827        let root = temp_path("recovery");
828        std::fs::create_dir_all(&root).unwrap();
829        let path = root.join("metadata");
830        std::fs::create_dir(&path).unwrap(); // A directory cannot be replaced by a JSON file.
831        let (sender, errors) = std::sync::mpsc::channel();
832        let store = SessionMetadataStore::open(&path);
833        let writer = store
834            .buffered_with_errors(move |error| {
835                let _ = sender.send(error);
836            })
837            .unwrap();
838        let snapshot = |value| {
839            SessionMetadata::new(serde_json::Map::from_iter([("value".into(), json!(value))]))
840        };
841        writer.write(snapshot(1)).unwrap();
842        assert!(
843            errors
844                .recv_timeout(std::time::Duration::from_secs(2))
845                .is_ok()
846        );
847        std::fs::remove_dir(&path).unwrap();
848        writer.write(snapshot(2)).unwrap();
849        writer.flush().unwrap();
850        assert_eq!(
851            store.read().unwrap().unwrap().metadata.get("value"),
852            Some(&json!(2))
853        );
854        std::fs::remove_file(&path).unwrap();
855        std::fs::create_dir(&path).unwrap();
856        writer.write(snapshot(3)).unwrap();
857        assert!(
858            errors
859                .recv_timeout(std::time::Duration::from_secs(2))
860                .is_ok()
861        );
862        std::fs::remove_dir(&path).unwrap();
863        writer.flush().unwrap(); // Retries the retained snapshot without another write.
864        assert_eq!(
865            store.read().unwrap().unwrap().metadata.get("value"),
866            Some(&json!(3))
867        );
868        std::fs::remove_file(&path).unwrap();
869        std::fs::create_dir(&path).unwrap();
870        writer.write(snapshot(4)).unwrap();
871        assert!(
872            errors
873                .recv_timeout(std::time::Duration::from_secs(2))
874                .is_ok()
875        );
876        std::fs::remove_dir(&path).unwrap();
877        let deadline = std::time::Instant::now() + std::time::Duration::from_secs(3);
878        while !path.is_file() && std::time::Instant::now() < deadline {
879            std::thread::sleep(std::time::Duration::from_millis(10));
880        }
881        assert_eq!(
882            store.read().unwrap().unwrap().metadata.get("value"),
883            Some(&json!(4))
884        );
885        drop(writer);
886        std::fs::remove_dir_all(root).unwrap();
887    }
888
889    #[test]
890    fn buffered_metadata_writes_are_durable_at_flush_and_drop() {
891        let path = temp_path("buffered").join("session.json");
892        let store = SessionMetadataStore::open(&path);
893        let writer = store.buffered().expect("writer");
894        let mut metadata = SessionMetadata::empty();
895        metadata.insert("roster", json!(["claude.ai", "openai.com"]));
896        writer.write(metadata).expect("queue snapshot");
897        writer.flush().expect("flush snapshot");
898        let loaded = store.read().expect("read").expect("snapshot");
899        assert_eq!(
900            loaded.get("roster"),
901            Some(&json!(["claude.ai", "openai.com"]))
902        );
903        drop(writer);
904        std::fs::remove_dir_all(path.parent().expect("parent")).expect("cleanup");
905    }
906
907    #[test]
908    fn legacy_import_accepts_missing_and_current_metadata() {
909        assert!(
910            import_legacy_session_metadata(None)
911                .expect("missing")
912                .is_none()
913        );
914        let loaded = import_legacy_session_metadata(Some(
915            r#"{"schema_version":1,"roster":["agy"],"agent_data":{"name":"Agy"}}"#,
916        ))
917        .expect("current")
918        .expect("metadata");
919        assert_eq!(loaded.source_version, 1);
920        assert_eq!(
921            loaded
922                .get("agent_data")
923                .and_then(Value::as_object)
924                .and_then(|m| m.get("name"))
925                .and_then(Value::as_str),
926            Some("Agy")
927        );
928    }
929
930    #[test]
931    fn future_versions_are_rejected() {
932        let event_path = temp_path("future-events.jsonl");
933        std::fs::write(
934            &event_path,
935            serde_json::to_string(&json!({"schema_version": 99, "event": event()})).expect("json"),
936        )
937        .expect("write");
938        assert!(matches!(
939            VersionedEventLog::open(&event_path).read(),
940            Err(PersistenceError::UnsupportedVersion { version: 99, .. })
941        ));
942        std::fs::remove_file(event_path).expect("cleanup");
943    }
944}