1use std::{
9 collections::{HashMap, HashSet},
10 fs::{File, OpenOptions},
11 io::{Read, Seek, SeekFrom, Write},
12 path::{Path, PathBuf},
13 sync::{Arc, Mutex, OnceLock, Weak},
14};
15
16use anyhow::{Context as _, bail, ensure};
17use serde::{Deserialize, Serialize};
18use sha2::{Digest, Sha256};
19
20pub const FORMAT_VERSION: &str = "0.2.1";
21
22const SESSION_MAGIC: &[u8] = b"KSESSIONLOG\n";
23const OBJECT_MAGIC: &[u8] = b"KSPENDING01\n";
24const FRAME_HEADER_BYTES: u64 = 1 + 8 + 32;
25const HEADER_FRAME: u8 = 1;
26const EVENT_FRAME: u8 = 2;
27const SEALED_FRAME: u8 = 3;
28const OBJECT_FIXED_HEADER_BYTES: usize = 2 + 4 + 4 + 8 + 32;
29
30#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
31#[serde(rename_all = "camelCase")]
32pub struct SessionHeader {
33 pub format_version: String,
34 pub session_id: String,
35 pub created_at: String,
36}
37
38#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq, Serialize, Deserialize)]
39#[serde(rename_all = "kebab-case")]
40pub enum Role {
41 SystemMessage,
42 SystemError,
43 UserMessage,
44 KennedyMessage,
45 KennedyToolCall,
46 ToolResult,
47 ToolError,
48 Object,
49 PendingObject,
50}
51
52#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
53pub struct SessionEvent {
54 pub role: Role,
55 pub text: String,
56}
57
58#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
59pub struct SessionLog {
60 pub header: SessionHeader,
61 pub events: Vec<SessionEvent>,
62}
63
64#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
65pub struct EventPosition(pub u64);
66
67impl EventPosition {
68 pub fn index(self) -> u64 {
69 self.0
70 }
71}
72
73impl std::fmt::Display for EventPosition {
74 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
75 self.0.fmt(formatter)
76 }
77}
78
79#[derive(Clone, Debug, Eq, PartialEq)]
80pub struct PendingObject {
81 pub event_position: EventPosition,
82 pub text: String,
83 pub file_name: String,
84 pub media_type: String,
85 pub bytes: Vec<u8>,
86}
87
88#[derive(Clone, Debug, Eq, PartialEq)]
89pub struct SealedSession {
90 log: SessionLog,
91 directory: PathBuf,
92}
93
94impl SealedSession {
95 pub fn list(&self) -> &SessionLog {
96 &self.log
97 }
98
99 pub fn pending_objects(&self) -> anyhow::Result<Vec<PendingObject>> {
100 self.log
101 .events
102 .iter()
103 .enumerate()
104 .filter(|(_, event)| event.role == Role::PendingObject)
105 .map(|(position, event)| {
106 read_pending_object_file(
107 &self.directory,
108 &self.log.header.session_id,
109 EventPosition(position as u64),
110 event.text.clone(),
111 )
112 })
113 .collect()
114 }
115}
116
117#[derive(Clone, Debug)]
118pub struct SessionStore {
119 directory: PathBuf,
120}
121
122impl SessionStore {
123 pub fn new(directory: impl Into<PathBuf>) -> Self {
124 Self {
125 directory: directory.into(),
126 }
127 }
128
129 pub fn directory(&self) -> &Path {
130 &self.directory
131 }
132
133 pub fn create_session(
134 &self,
135 session_id: impl Into<String>,
136 created_at: impl Into<String>,
137 ) -> anyhow::Result<Session> {
138 let session_id = session_id.into();
139 let created_at = created_at.into();
140 validate_session_id(&session_id)?;
141 ensure!(
142 !created_at.trim().is_empty(),
143 "session creation time cannot be empty"
144 );
145 std::fs::create_dir_all(&self.directory)
146 .with_context(|| format!("creating {}", self.directory.display()))?;
147 let path = session_path(&self.directory, &session_id);
148 let append_lock = session_lock(&path);
149 let _guard = append_lock
150 .lock()
151 .map_err(|_| anyhow::anyhow!("session-log append lock is poisoned"))?;
152 let header = SessionHeader {
153 format_version: FORMAT_VERSION.into(),
154 session_id,
155 created_at,
156 };
157 let mut file = OpenOptions::new()
158 .create_new(true)
159 .write(true)
160 .open(&path)
161 .with_context(|| format!("creating {}", path.display()))?;
162 file.write_all(SESSION_MAGIC)?;
163 append_frame(&mut file, HEADER_FRAME, &serde_json::to_vec(&header)?)?;
164 file.sync_all()?;
165 sync_directory(&self.directory)?;
166 drop(_guard);
167 Ok(Session {
168 directory: self.directory.clone(),
169 path,
170 log: SessionLog {
171 header,
172 events: Vec::new(),
173 },
174 sealed: false,
175 append_lock,
176 })
177 }
178
179 pub fn open_session(&self, session_id: &str) -> anyhow::Result<Session> {
180 validate_session_id(session_id)?;
181 let path = session_path(&self.directory, session_id);
182 let append_lock = session_lock(&path);
183 let _guard = append_lock
184 .lock()
185 .map_err(|_| anyhow::anyhow!("session-log append lock is poisoned"))?;
186 let loaded = load_session_file(&path, true)?;
187 ensure!(
188 loaded.log.header.session_id == session_id,
189 "session filename and header identity differ"
190 );
191 cleanup_orphan_objects(&self.directory, &loaded.log)?;
192 verify_referenced_objects(&self.directory, &loaded.log)?;
193 drop(_guard);
194 Ok(Session {
195 directory: self.directory.clone(),
196 path,
197 log: loaded.log,
198 sealed: loaded.sealed,
199 append_lock,
200 })
201 }
202
203 pub fn session_ids(&self) -> anyhow::Result<Vec<String>> {
204 if !self.directory.exists() {
205 return Ok(Vec::new());
206 }
207 let mut ids = std::fs::read_dir(&self.directory)?
208 .filter_map(Result::ok)
209 .filter_map(|entry| {
210 entry
211 .file_name()
212 .to_str()
213 .and_then(|name| name.strip_suffix(".session-log"))
214 .map(str::to_owned)
215 })
216 .filter(|id| validate_session_id(id).is_ok())
217 .collect::<Vec<_>>();
218 ids.sort();
219 Ok(ids)
220 }
221}
222
223pub struct Session {
224 directory: PathBuf,
225 path: PathBuf,
226 log: SessionLog,
227 sealed: bool,
228 append_lock: Arc<Mutex<()>>,
229}
230
231impl Session {
232 pub fn path(&self) -> &Path {
233 &self.path
234 }
235
236 pub fn list(&self) -> SessionLog {
237 self.log.clone()
238 }
239
240 pub fn is_sealed(&self) -> bool {
241 self.sealed
242 }
243
244 pub fn add_event(
245 &mut self,
246 role: Role,
247 text: impl Into<String>,
248 ) -> anyhow::Result<EventPosition> {
249 ensure!(
250 role != Role::PendingObject,
251 "pending objects must be added through add_pending_object"
252 );
253 let event = SessionEvent {
254 role,
255 text: text.into(),
256 };
257 let append_lock = self.append_lock.clone();
258 let _guard = append_lock
259 .lock()
260 .map_err(|_| anyhow::anyhow!("session-log append lock is poisoned"))?;
261 self.refresh_locked()?;
262 ensure!(!self.sealed, "session is sealed");
263 let position = EventPosition(self.log.events.len() as u64);
264 append_event_file(&self.path, &event)?;
265 self.log.events.push(event);
266 Ok(position)
267 }
268
269 pub fn add_pending_object(
270 &mut self,
271 text: impl Into<String>,
272 file_name: impl Into<String>,
273 media_type: impl Into<String>,
274 bytes: &[u8],
275 ) -> anyhow::Result<EventPosition> {
276 let text = text.into();
277 let file_name = file_name.into();
278 let media_type = media_type.into();
279 ensure!(
280 !file_name.trim().is_empty(),
281 "object filename cannot be empty"
282 );
283 ensure!(
284 !media_type.trim().is_empty(),
285 "object media type cannot be empty"
286 );
287 let append_lock = self.append_lock.clone();
288 let _guard = append_lock
289 .lock()
290 .map_err(|_| anyhow::anyhow!("session-log append lock is poisoned"))?;
291 self.refresh_locked()?;
292 ensure!(!self.sealed, "session is sealed");
293 let position = EventPosition(self.log.events.len() as u64);
294 write_pending_object_file(
295 &self.directory,
296 &self.log.header.session_id,
297 position,
298 &file_name,
299 &media_type,
300 bytes,
301 )?;
302 let event = SessionEvent {
303 role: Role::PendingObject,
304 text,
305 };
306 append_event_file(&self.path, &event)?;
307 self.log.events.push(event);
308 Ok(position)
309 }
310
311 pub fn read_pending_object(&self, position: EventPosition) -> anyhow::Result<PendingObject> {
312 let index = usize::try_from(position.0).context("event position does not fit memory")?;
313 let event = self
314 .log
315 .events
316 .get(index)
317 .with_context(|| format!("event {position} does not exist"))?;
318 ensure!(
319 event.role == Role::PendingObject,
320 "event {position} is not a pending object"
321 );
322 read_pending_object_file(
323 &self.directory,
324 &self.log.header.session_id,
325 position,
326 event.text.clone(),
327 )
328 }
329
330 pub fn seal(&mut self) -> anyhow::Result<SealedSession> {
331 let append_lock = self.append_lock.clone();
332 let _guard = append_lock
333 .lock()
334 .map_err(|_| anyhow::anyhow!("session-log append lock is poisoned"))?;
335 self.refresh_locked()?;
336 if !self.sealed {
337 verify_referenced_objects(&self.directory, &self.log)?;
338 let mut file = OpenOptions::new().append(true).open(&self.path)?;
339 append_frame(&mut file, SEALED_FRAME, &[])?;
340 file.sync_all()?;
341 sync_directory(&self.directory)?;
342 self.sealed = true;
343 }
344 Ok(SealedSession {
345 log: self.log.clone(),
346 directory: self.directory.clone(),
347 })
348 }
349
350 pub fn delete_committed(self) -> anyhow::Result<()> {
351 self.delete_files()
352 }
353
354 pub fn delete_abandoned(self) -> anyhow::Result<()> {
355 self.delete_files()
356 }
357
358 fn refresh_locked(&mut self) -> anyhow::Result<()> {
359 let loaded = load_session_file(&self.path, true)?;
360 ensure!(
361 loaded.log.header.session_id == self.log.header.session_id,
362 "session identity changed on disk"
363 );
364 self.log = loaded.log;
365 self.sealed = loaded.sealed;
366 Ok(())
367 }
368
369 fn delete_files(self) -> anyhow::Result<()> {
370 let append_lock = self.append_lock.clone();
371 let _guard = append_lock
372 .lock()
373 .map_err(|_| anyhow::anyhow!("session-log append lock is poisoned"))?;
374 let session_id = self.log.header.session_id;
375 if self.path.exists() {
376 std::fs::remove_file(&self.path)
377 .with_context(|| format!("removing {}", self.path.display()))?;
378 }
379 if self.directory.exists() {
380 for entry in std::fs::read_dir(&self.directory)? {
381 let entry = entry?;
382 let Some(name) = entry.file_name().to_str().map(str::to_owned) else {
383 continue;
384 };
385 if parse_object_filename(&session_id, &name).is_some() {
386 std::fs::remove_file(entry.path())
387 .with_context(|| format!("removing {}", entry.path().display()))?;
388 }
389 }
390 }
391 sync_directory(&self.directory)?;
392 Ok(())
393 }
394}
395
396struct LoadedSession {
397 log: SessionLog,
398 sealed: bool,
399}
400
401fn validate_session_id(session_id: &str) -> anyhow::Result<()> {
402 ensure!(!session_id.is_empty(), "session ID cannot be empty");
403 ensure!(session_id.len() <= 255, "session ID exceeds 255 characters");
404 ensure!(
405 session_id
406 .bytes()
407 .all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'-' | b'_')),
408 "session ID contains characters that are unsafe in filenames"
409 );
410 Ok(())
411}
412
413fn session_path(directory: &Path, session_id: &str) -> PathBuf {
414 directory.join(format!("{session_id}.session-log"))
415}
416
417fn pending_object_path(directory: &Path, session_id: &str, position: EventPosition) -> PathBuf {
418 directory.join(format!("{session_id}-{}.pending-object", position.0))
419}
420
421fn pending_object_temp_path(
422 directory: &Path,
423 session_id: &str,
424 position: EventPosition,
425) -> PathBuf {
426 directory.join(format!("{session_id}-{}.pending-object.tmp", position.0))
427}
428
429fn parse_object_filename(session_id: &str, name: &str) -> Option<(EventPosition, bool)> {
430 let tail = name.strip_prefix(&format!("{session_id}-"))?;
431 let (number, temporary) = if let Some(number) = tail.strip_suffix(".pending-object.tmp") {
432 (number, true)
433 } else {
434 (tail.strip_suffix(".pending-object")?, false)
435 };
436 if number.is_empty() || number.starts_with('0') && number != "0" {
437 return None;
438 }
439 let position = number.parse::<u64>().ok()?;
440 (position.to_string() == number).then_some((EventPosition(position), temporary))
441}
442
443fn append_event_file(path: &Path, event: &SessionEvent) -> anyhow::Result<()> {
444 let mut file = OpenOptions::new().append(true).open(path)?;
445 append_frame(&mut file, EVENT_FRAME, &serde_json::to_vec(event)?)?;
446 file.sync_all()?;
447 Ok(())
448}
449
450fn load_session_file(path: &Path, repair_tail: bool) -> anyhow::Result<LoadedSession> {
451 let mut file = OpenOptions::new()
452 .read(true)
453 .write(repair_tail)
454 .open(path)
455 .with_context(|| format!("opening {}", path.display()))?;
456 let mut magic = vec![0_u8; SESSION_MAGIC.len()];
457 file.read_exact(&mut magic)?;
458 ensure!(
459 magic == SESSION_MAGIC,
460 "{} is not a session log",
461 path.display()
462 );
463 let file_len = file.metadata()?.len();
464 let mut cursor = SESSION_MAGIC.len() as u64;
465 let mut header = None;
466 let mut events = Vec::new();
467 let mut sealed = false;
468 while cursor < file_len {
469 let remaining = file_len - cursor;
470 if remaining < FRAME_HEADER_BYTES {
471 if repair_tail {
472 file.set_len(cursor)?;
473 file.sync_all()?;
474 break;
475 }
476 bail!("session log has an incomplete trailing frame header");
477 }
478 file.seek(SeekFrom::Start(cursor))?;
479 let mut frame_header = [0_u8; FRAME_HEADER_BYTES as usize];
480 file.read_exact(&mut frame_header)?;
481 let kind = frame_header[0];
482 let payload_len = u64::from_le_bytes(frame_header[1..9].try_into().unwrap());
483 let frame_end = cursor
484 .checked_add(FRAME_HEADER_BYTES)
485 .and_then(|value| value.checked_add(payload_len))
486 .context("session-log frame length overflow")?;
487 if frame_end > file_len {
488 if repair_tail {
489 file.set_len(cursor)?;
490 file.sync_all()?;
491 break;
492 }
493 bail!("session log has an incomplete trailing frame");
494 }
495 let payload_len_usize =
496 usize::try_from(payload_len).context("session-log frame does not fit memory")?;
497 let mut payload = vec![0_u8; payload_len_usize];
498 file.read_exact(&mut payload)?;
499 if Sha256::digest(&payload).as_slice() != &frame_header[9..] {
500 bail!("session log has a checksum-invalid complete frame");
501 }
502 match kind {
503 HEADER_FRAME => {
504 ensure!(
505 cursor == SESSION_MAGIC.len() as u64,
506 "duplicate session header"
507 );
508 let value: SessionHeader = serde_json::from_slice(&payload)?;
509 ensure!(
510 value.format_version == FORMAT_VERSION,
511 "unsupported session-log format {}",
512 value.format_version
513 );
514 validate_session_id(&value.session_id)?;
515 ensure!(
516 !value.created_at.trim().is_empty(),
517 "session creation time cannot be empty"
518 );
519 header = Some(value);
520 }
521 EVENT_FRAME => {
522 ensure!(header.is_some(), "session event precedes header");
523 ensure!(!sealed, "session event follows sealed footer");
524 events.push(serde_json::from_slice(&payload)?);
525 }
526 SEALED_FRAME => {
527 ensure!(header.is_some(), "sealed footer precedes header");
528 ensure!(payload.is_empty(), "sealed footer payload must be empty");
529 ensure!(!sealed, "duplicate sealed footer");
530 sealed = true;
531 }
532 other => bail!("unknown complete session-log frame kind {other}"),
533 }
534 cursor = frame_end;
535 }
536 let header = header.context("session log has no header")?;
537 Ok(LoadedSession {
538 log: SessionLog { header, events },
539 sealed,
540 })
541}
542
543fn append_frame(file: &mut File, kind: u8, payload: &[u8]) -> anyhow::Result<()> {
544 file.write_all(&[kind])?;
545 file.write_all(&(payload.len() as u64).to_le_bytes())?;
546 file.write_all(&Sha256::digest(payload))?;
547 file.write_all(payload)?;
548 Ok(())
549}
550
551fn write_pending_object_file(
552 directory: &Path,
553 session_id: &str,
554 position: EventPosition,
555 file_name: &str,
556 media_type: &str,
557 bytes: &[u8],
558) -> anyhow::Result<()> {
559 let version = FORMAT_VERSION.as_bytes();
560 let version_len = u16::try_from(version.len()).context("format version is too long")?;
561 let file_name_len = u32::try_from(file_name.len()).context("object filename exceeds 4 GiB")?;
562 let media_type_len =
563 u32::try_from(media_type.len()).context("object media type exceeds 4 GiB")?;
564 let object_len = u64::try_from(bytes.len()).context("object exceeds addressable size")?;
565 let final_path = pending_object_path(directory, session_id, position);
566 let temp_path = pending_object_temp_path(directory, session_id, position);
567 ensure!(
568 !final_path.exists(),
569 "pending object file {} already exists",
570 final_path.display()
571 );
572 if temp_path.exists() {
573 std::fs::remove_file(&temp_path)?;
574 }
575 let mut file = OpenOptions::new()
576 .create_new(true)
577 .write(true)
578 .open(&temp_path)?;
579 file.write_all(OBJECT_MAGIC)?;
580 file.write_all(&version_len.to_le_bytes())?;
581 file.write_all(&file_name_len.to_le_bytes())?;
582 file.write_all(&media_type_len.to_le_bytes())?;
583 file.write_all(&object_len.to_le_bytes())?;
584 file.write_all(&Sha256::digest(bytes))?;
585 file.write_all(version)?;
586 file.write_all(file_name.as_bytes())?;
587 file.write_all(media_type.as_bytes())?;
588 file.write_all(bytes)?;
589 file.sync_all()?;
590 std::fs::rename(&temp_path, &final_path)?;
591 sync_directory(directory)?;
592 Ok(())
593}
594
595fn read_pending_object_file(
596 directory: &Path,
597 session_id: &str,
598 position: EventPosition,
599 text: String,
600) -> anyhow::Result<PendingObject> {
601 let path = pending_object_path(directory, session_id, position);
602 let mut file =
603 File::open(&path).with_context(|| format!("opening pending object {}", path.display()))?;
604 let mut magic = vec![0_u8; OBJECT_MAGIC.len()];
605 file.read_exact(&mut magic)?;
606 ensure!(
607 magic == OBJECT_MAGIC,
608 "{} is not a pending-object file",
609 path.display()
610 );
611 let mut fixed = [0_u8; OBJECT_FIXED_HEADER_BYTES];
612 file.read_exact(&mut fixed)?;
613 let version_len = u16::from_le_bytes(fixed[0..2].try_into().unwrap()) as usize;
614 let file_name_len = u32::from_le_bytes(fixed[2..6].try_into().unwrap()) as usize;
615 let media_type_len = u32::from_le_bytes(fixed[6..10].try_into().unwrap()) as usize;
616 let object_len = u64::from_le_bytes(fixed[10..18].try_into().unwrap());
617 let checksum = &fixed[18..50];
618 let variable_len = version_len
619 .checked_add(file_name_len)
620 .and_then(|value| value.checked_add(media_type_len))
621 .context("pending-object header length overflow")?;
622 let expected_file_len = u64::try_from(OBJECT_MAGIC.len() + OBJECT_FIXED_HEADER_BYTES)
623 .context("pending-object fixed header does not fit u64")?
624 .checked_add(
625 u64::try_from(variable_len)
626 .context("pending-object variable header does not fit u64")?,
627 )
628 .and_then(|value| value.checked_add(object_len))
629 .context("pending-object declared length overflow")?;
630 ensure!(
631 file.metadata()?.len() == expected_file_len,
632 "pending-object declared length differs from file length"
633 );
634 let mut variable = vec![0_u8; variable_len];
635 file.read_exact(&mut variable)?;
636 let version = std::str::from_utf8(&variable[..version_len])?;
637 ensure!(
638 version == FORMAT_VERSION,
639 "unsupported pending-object format {version}"
640 );
641 let file_name_end = version_len + file_name_len;
642 let file_name = std::str::from_utf8(&variable[version_len..file_name_end])?.to_owned();
643 let media_type =
644 std::str::from_utf8(&variable[file_name_end..file_name_end + media_type_len])?.to_owned();
645 ensure!(!file_name.trim().is_empty(), "object filename is empty");
646 ensure!(!media_type.trim().is_empty(), "object media type is empty");
647 let object_len_usize =
648 usize::try_from(object_len).context("pending object does not fit memory")?;
649 let mut bytes = vec![0_u8; object_len_usize];
650 file.read_exact(&mut bytes)?;
651 ensure!(
652 file.stream_position()? == file.metadata()?.len(),
653 "pending-object file has trailing bytes"
654 );
655 ensure!(
656 Sha256::digest(&bytes).as_slice() == checksum,
657 "pending-object checksum mismatch"
658 );
659 Ok(PendingObject {
660 event_position: position,
661 text,
662 file_name,
663 media_type,
664 bytes,
665 })
666}
667
668fn referenced_pending_positions(log: &SessionLog) -> HashSet<EventPosition> {
669 log.events
670 .iter()
671 .enumerate()
672 .filter_map(|(position, event)| {
673 (event.role == Role::PendingObject).then_some(EventPosition(position as u64))
674 })
675 .collect()
676}
677
678fn verify_referenced_objects(directory: &Path, log: &SessionLog) -> anyhow::Result<()> {
679 for position in referenced_pending_positions(log) {
680 let event = &log.events[position.0 as usize];
681 read_pending_object_file(
682 directory,
683 &log.header.session_id,
684 position,
685 event.text.clone(),
686 )?;
687 }
688 Ok(())
689}
690
691fn cleanup_orphan_objects(directory: &Path, log: &SessionLog) -> anyhow::Result<()> {
692 let referenced = referenced_pending_positions(log);
693 if !directory.exists() {
694 return Ok(());
695 }
696 let mut removed = false;
697 for entry in std::fs::read_dir(directory)? {
698 let entry = entry?;
699 let Some(name) = entry.file_name().to_str().map(str::to_owned) else {
700 continue;
701 };
702 let Some((position, temporary)) = parse_object_filename(&log.header.session_id, &name)
703 else {
704 continue;
705 };
706 if temporary || !referenced.contains(&position) {
707 std::fs::remove_file(entry.path())?;
708 removed = true;
709 }
710 }
711 if removed {
712 sync_directory(directory)?;
713 }
714 Ok(())
715}
716
717fn session_lock(path: &Path) -> Arc<Mutex<()>> {
718 static LOCKS: OnceLock<Mutex<HashMap<PathBuf, Weak<Mutex<()>>>>> = OnceLock::new();
719 let locks = LOCKS.get_or_init(|| Mutex::new(HashMap::new()));
720 let mut locks = locks.lock().expect("session-log lock registry is poisoned");
721 locks.retain(|_, lock| lock.strong_count() > 0);
722 if let Some(lock) = locks.get(path).and_then(Weak::upgrade) {
723 return lock;
724 }
725 let lock = Arc::new(Mutex::new(()));
726 locks.insert(path.to_path_buf(), Arc::downgrade(&lock));
727 lock
728}
729
730fn sync_directory(directory: &Path) -> anyhow::Result<()> {
731 File::open(directory)
732 .with_context(|| format!("opening directory {}", directory.display()))?
733 .sync_all()
734 .with_context(|| format!("synchronizing directory {}", directory.display()))
735}
736
737#[cfg(test)]
738mod tests {
739 use std::{
740 fs::OpenOptions,
741 io::Write,
742 time::{SystemTime, UNIX_EPOCH},
743 };
744
745 use super::*;
746
747 fn directory(label: &str) -> PathBuf {
748 std::env::temp_dir().join(format!(
749 "session-log-{label}-{}-{}",
750 std::process::id(),
751 SystemTime::now()
752 .duration_since(UNIX_EPOCH)
753 .unwrap()
754 .as_nanos()
755 ))
756 }
757
758 #[test]
759 fn events_are_an_ordered_role_and_text_array_without_serialized_ids() {
760 let directory = directory("ordered");
761 let store = SessionStore::new(&directory);
762 let mut session = store
763 .create_session("session-1", "2026-07-24T00:00:00Z")
764 .unwrap();
765 assert_eq!(
766 session.add_event(Role::SystemMessage, "system").unwrap(),
767 EventPosition(0)
768 );
769 assert_eq!(
770 session.add_event(Role::UserMessage, "hello").unwrap(),
771 EventPosition(1)
772 );
773 drop(session);
774
775 let reopened = store.open_session("session-1").unwrap();
776 assert_eq!(
777 reopened.list(),
778 SessionLog {
779 header: SessionHeader {
780 format_version: "0.2.1".into(),
781 session_id: "session-1".into(),
782 created_at: "2026-07-24T00:00:00Z".into(),
783 },
784 events: vec![
785 SessionEvent {
786 role: Role::SystemMessage,
787 text: "system".into(),
788 },
789 SessionEvent {
790 role: Role::UserMessage,
791 text: "hello".into(),
792 },
793 ],
794 }
795 );
796 let serialized = serde_json::to_string(&reopened.list().events).unwrap();
797 assert!(!serialized.contains("\"id\""));
798 std::fs::remove_dir_all(directory).unwrap();
799 }
800
801 #[test]
802 fn pending_object_is_durable_before_its_event_and_uses_event_position() {
803 let directory = directory("object");
804 let store = SessionStore::new(&directory);
805 let mut session = store
806 .create_session("object-session", "2026-07-24T00:00:00Z")
807 .unwrap();
808 session.add_event(Role::UserMessage, "upload").unwrap();
809 let position = session
810 .add_pending_object("notes.txt", "notes.txt", "text/plain", b"durable bytes")
811 .unwrap();
812 assert_eq!(position, EventPosition(1));
813 assert!(directory.join("object-session-1.pending-object").exists());
814 let object = session.read_pending_object(position).unwrap();
815 assert_eq!(object.file_name, "notes.txt");
816 assert_eq!(object.media_type, "text/plain");
817 assert_eq!(object.bytes, b"durable bytes");
818 std::fs::remove_dir_all(directory).unwrap();
819 }
820
821 #[test]
822 fn open_removes_unreferenced_final_and_temporary_objects() {
823 let directory = directory("orphans");
824 let store = SessionStore::new(&directory);
825 let session = store
826 .create_session("orphan-session", "2026-07-24T00:00:00Z")
827 .unwrap();
828 drop(session);
829 std::fs::write(directory.join("orphan-session-0.pending-object"), b"orphan").unwrap();
830 std::fs::write(
831 directory.join("orphan-session-1.pending-object.tmp"),
832 b"temporary",
833 )
834 .unwrap();
835 store.open_session("orphan-session").unwrap();
836 assert!(!directory.join("orphan-session-0.pending-object").exists());
837 assert!(
838 !directory
839 .join("orphan-session-1.pending-object.tmp")
840 .exists()
841 );
842 std::fs::remove_dir_all(directory).unwrap();
843 }
844
845 #[test]
846 fn incomplete_event_tail_is_discarded() {
847 let directory = directory("tail");
848 let store = SessionStore::new(&directory);
849 let mut session = store
850 .create_session("tail-session", "2026-07-24T00:00:00Z")
851 .unwrap();
852 session.add_event(Role::UserMessage, "complete").unwrap();
853 let path = session.path().to_path_buf();
854 drop(session);
855 let valid_len = std::fs::metadata(&path).unwrap().len();
856 OpenOptions::new()
857 .append(true)
858 .open(&path)
859 .unwrap()
860 .write_all(&[EVENT_FRAME, 20, 0, 0])
861 .unwrap();
862 let reopened = store.open_session("tail-session").unwrap();
863 assert_eq!(reopened.list().events.len(), 1);
864 assert_eq!(std::fs::metadata(&path).unwrap().len(), valid_len);
865 std::fs::remove_dir_all(directory).unwrap();
866 }
867
868 #[test]
869 fn checksum_invalid_complete_frame_is_corruption_not_a_recoverable_tail() {
870 let directory = directory("checksum");
871 let store = SessionStore::new(&directory);
872 let mut session = store
873 .create_session("checksum-session", "2026-07-24T00:00:00Z")
874 .unwrap();
875 session.add_event(Role::UserMessage, "complete").unwrap();
876 let path = session.path().to_path_buf();
877 drop(session);
878 let mut bytes = std::fs::read(&path).unwrap();
879 let payload_byte = SESSION_MAGIC.len() + FRAME_HEADER_BYTES as usize;
880 bytes[payload_byte] ^= 0xff;
881 std::fs::write(&path, &bytes).unwrap();
882 assert!(store.open_session("checksum-session").is_err());
883 assert_eq!(std::fs::read(&path).unwrap(), bytes);
884 std::fs::remove_dir_all(directory).unwrap();
885 }
886
887 #[test]
888 fn seal_is_durable_idempotent_and_rejects_later_events() {
889 let directory = directory("seal");
890 let store = SessionStore::new(&directory);
891 let mut session = store
892 .create_session("sealed-session", "2026-07-24T00:00:00Z")
893 .unwrap();
894 session.add_event(Role::UserMessage, "hello").unwrap();
895 session.seal().unwrap();
896 session.seal().unwrap();
897 assert!(session.add_event(Role::KennedyMessage, "too late").is_err());
898 drop(session);
899 assert!(store.open_session("sealed-session").unwrap().is_sealed());
900 std::fs::remove_dir_all(directory).unwrap();
901 }
902
903 #[test]
904 fn deletion_matches_only_exact_session_object_names() {
905 let directory = directory("delete");
906 let store = SessionStore::new(&directory);
907 let session = store.create_session("abc", "2026-07-24T00:00:00Z").unwrap();
908 std::fs::write(directory.join("abc-other-0.pending-object"), b"keep").unwrap();
909 std::fs::write(directory.join("abc-not-a-number.pending-object"), b"keep").unwrap();
910 session.delete_abandoned().unwrap();
911 assert!(directory.join("abc-other-0.pending-object").exists());
912 assert!(directory.join("abc-not-a-number.pending-object").exists());
913 std::fs::remove_dir_all(directory).unwrap();
914 }
915}