Skip to main content

kcode_session_log/
lib.rs

1//! `kcode-session-log` provides durable, append-ordered session history.
2//!
3//! A session log stores one three-field header followed by ordered `{role,
4//! text}` events. Event positions are their stable identities and are not
5//! serialized. Pending objects are written to one self-contained file each
6//! before their corresponding events are appended.
7
8use 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}