1use 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
21pub const CURRENT_SCHEMA_VERSION: u32 = 1;
23const LEGACY_SCHEMA_VERSION: u32 = 0;
24
25#[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#[derive(Clone, Debug, Eq, PartialEq)]
56pub struct MigrationReport {
57 pub source_version: Option<u32>,
60 pub target_version: u32,
61 pub records: usize,
62 pub changed: bool,
63}
64
65#[derive(Clone, Debug, Eq, PartialEq)]
67pub struct LoadedEvents {
68 pub events: Vec<AgentEvent>,
69 pub source_versions: BTreeSet<u32>,
70}
71
72#[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 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 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 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#[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 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 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 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#[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#[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 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 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
424pub 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 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 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
576pub 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}