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 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
431pub 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 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 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
625pub 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(); 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(); 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}