Luigit
repositories / smith

smith

There are many coding harnesses - but this one is fast

owned by admin

smith-core/src/store.rs

Raw
//! Durable session storage: append-only writer, incremental bounded loader,
//! and atomic repair publication (`SMH-SPEC-SPEC0001`, Sessions).

use crate::frame::{Frame, Recovery, encode_entry, read_frames, write_frame};
use crate::session::{Entry, Session};
use smith::config::Config;
use smith::error::{Result, SmithError};
use smith::id::EntryId;
use std::fs::{File, OpenOptions};
use std::io::{BufReader, Seek, SeekFrom, Write};
use std::path::{Path, PathBuf};

/// Exclusive access to one session, held on a sidecar lock file.
///
/// The lock lives beside the session file rather than on it, because repair
/// publishes a replacement inode: a lock taken on the data file would leave
/// a waiter holding a descriptor on the replaced inode, and its appends
/// would vanish. The sidecar is never replaced, so every participant
/// serializes on the same object and opens the data file only after winning.
#[derive(Debug)]
struct SessionLock {
    file: File,
}

impl SessionLock {
    /// Acquire the lock for `path`, creating parent directories and the
    /// sidecar as needed. Blocks while another holder owns the session.
    fn acquire(path: &Path) -> Result<Self> {
        if let Some(parent) = path.parent()
            && !parent.as_os_str().is_empty()
        {
            std::fs::create_dir_all(parent).map_err(io_err("STORE_MKDIR"))?;
        }
        let file = OpenOptions::new()
            .read(true)
            .write(true)
            .create(true)
            .truncate(false)
            .open(lock_path(path))
            .map_err(io_err("STORE_LOCK_OPEN"))?;
        fs2::FileExt::lock_exclusive(&file).map_err(io_err("STORE_LOCK"))?;
        Ok(Self { file })
    }
}

impl Drop for SessionLock {
    fn drop(&mut self) {
        // Best-effort unlock; the OS also releases on close.
        let _ = fs2::FileExt::unlock(&self.file);
    }
}

/// Sidecar lock file path for a session file.
fn lock_path(path: &Path) -> PathBuf {
    let mut name = path.as_os_str().to_os_string();
    name.push(".lock");
    PathBuf::from(name)
}

/// Append-only session writer holding exclusive access to its session.
///
/// Opening repairs a torn tail (truncating to the last complete frame
/// boundary) before appending, so acknowledged appends are always
/// reachable by the loader. The lock serializes writers against each other
/// and against repair.
#[derive(Debug)]
pub struct SessionWriter {
    file: File,
    config: Config,
    path: PathBuf,
    /// Held for the writer's lifetime; released on drop.
    _lock: SessionLock,
}

impl SessionWriter {
    /// Open (creating parent directories) and lock the session file.
    ///
    /// # Errors
    ///
    /// Returns a configuration, locking, or input/output error.
    pub fn open<P: AsRef<Path>>(path: P, config: Config) -> Result<Self> {
        Self::open_with_recovery(path, config).map(|(writer, _)| writer)
    }

    /// Open and lock the session file, also returning what loading recovered,
    /// so callers do not have to walk the file a second time.
    fn open_with_recovery<P: AsRef<Path>>(path: P, config: Config) -> Result<(Self, Recovery)> {
        let path = path.as_ref().to_path_buf();
        // Win the lock before opening the data file, so a repair that is
        // publishing a replacement cannot hand out a stale descriptor.
        let lock = SessionLock::acquire(&path)?;
        Self::attach(path, config, lock)
    }

    /// Take over an already-locked session file.
    fn attach(path: PathBuf, config: Config, lock: SessionLock) -> Result<(Self, Recovery)> {
        // Validation already names the offending field path.
        config.validate()?;
        let mut file = OpenOptions::new()
            .read(true)
            .write(true)
            .create(true)
            .truncate(false)
            .open(&path)
            .map_err(io_err("STORE_OPEN"))?;
        // Drop a torn tail before appending so new frames stay reachable.
        // Complete but damaged frames are left in place and stay reported.
        let recovery = walk_file(&file, config.max_frame_bytes)?;
        let len = file.metadata().map_err(io_err("STORE_STAT"))?.len();
        if recovery.good_end < len {
            file.set_len(recovery.good_end)
                .map_err(io_err("STORE_TRUNCATE"))?;
        }
        file.seek(SeekFrom::Start(recovery.good_end))
            .map_err(io_err("STORE_SEEK"))?;
        Ok((
            Self {
                file,
                config,
                path,
                _lock: lock,
            },
            recovery,
        ))
    }

    /// Append one frame. The encoded size is checked against the configured
    /// limit before any bytes are written; a rejected frame leaves the file
    /// unchanged.
    ///
    /// # Errors
    ///
    /// Returns a framing error for oversized frames, or an input/output error.
    pub fn append_frame(&mut self, frame: &Frame) -> Result<()> {
        let bytes = frame.encode(self.config.max_frame_bytes)?;
        self.write(&bytes)
    }

    /// Append one session record as an entry frame, framing its arena
    /// `content` bytes directly. Size limits apply as in
    /// [`SessionWriter::append_frame`].
    ///
    /// # Errors
    ///
    /// Returns a framing error for oversized frames, or an input/output error.
    pub fn append_entry(&mut self, entry: &Entry, content: &[u8]) -> Result<()> {
        let bytes = encode_entry(entry, content, self.config.max_frame_bytes)?;
        self.write(&bytes)
    }

    fn write(&mut self, bytes: &[u8]) -> Result<()> {
        self.file.write_all(bytes).map_err(io_err("STORE_WRITE"))?;
        self.file.flush().map_err(io_err("STORE_WRITE"))?;
        Ok(())
    }

    /// Path of the underlying session file.
    #[must_use]
    pub fn path(&self) -> &Path {
        &self.path
    }
}

/// A session file opened for use.
#[derive(Debug)]
pub struct OpenSession {
    /// Restored, or newly created, session state.
    pub session: Session,
    /// Locked writer positioned after the last complete frame.
    pub writer: SessionWriter,
    /// Damage and torn-tail information found while loading.
    pub recovery: Recovery,
}

/// Open a session file, writing a header frame when the file is new.
///
/// Reopening restores the persisted session and branch IDs, so state
/// reconstructed from the same file is identical on every load. A file that
/// holds frames but no header cannot be attributed to a session and fails
/// explicitly instead of silently becoming a new session.
///
/// # Errors
///
/// Returns [`SmithError::Session`] when the file holds frames but no header,
/// or a locking, framing, or input/output error.
pub fn open_session(path: &Path, config: Config) -> Result<OpenSession> {
    let (mut writer, recovery) = SessionWriter::open_with_recovery(path, config)?;
    let has_header = recovery
        .frames
        .iter()
        .any(|frame| matches!(frame, Frame::Header { .. }));
    let session = if has_header {
        Session::from_frames(&recovery.frames)?
    } else if recovery.frames.is_empty() && recovery.damage.is_empty() {
        let session = Session::new();
        writer.append_frame(&Frame::Header {
            header: session.header(),
        })?;
        session
    } else {
        return Err(SmithError::Session {
            code: "SESSION_HEADER_MISSING".to_string(),
            message: format!("{} holds frames but no session header", path.display()),
        });
    };
    Ok(OpenSession {
        session,
        writer,
        recovery,
    })
}

/// Publish `session` as a new session file and open it for appending.
///
/// The file is written completely and then published by rename, so an
/// interrupted fork leaves no half-written session behind. An existing
/// non-empty target is refused rather than overwritten.
///
/// # Errors
///
/// Returns [`SmithError::Session`] when the target already holds a session,
/// or a locking, framing, or input/output error.
pub fn create_session(path: &Path, session: &Session, config: Config) -> Result<OpenSession> {
    let lock = SessionLock::acquire(path)?;
    if std::fs::metadata(path).is_ok_and(|meta| meta.len() > 0) {
        return Err(SmithError::Session {
            code: "SESSION_EXISTS".to_string(),
            message: format!("{} already holds a session", path.display()),
        });
    }
    let tmp = path.with_extension("smith-create.tmp");
    {
        let mut out = File::create(&tmp).map_err(io_err("CREATE_TMP"))?;
        write_frame(
            &mut out,
            &Frame::Header {
                header: session.header(),
            },
            config.max_frame_bytes,
        )?;
        for entry in session.active_branch().entries() {
            let bytes = encode_entry(entry, session.bytes(entry), config.max_frame_bytes)?;
            out.write_all(&bytes).map_err(io_err("CREATE_WRITE"))?;
        }
        out.sync_all().map_err(io_err("CREATE_SYNC"))?;
    }
    std::fs::rename(&tmp, path).map_err(io_err("CREATE_PUBLISH"))?;
    let (writer, recovery) = SessionWriter::attach(path.to_path_buf(), config, lock)?;
    Ok(OpenSession {
        session: Session::from_frames(&recovery.frames)?,
        writer,
        recovery,
    })
}

/// Fork the history of `source` ending at `at` into a new session file.
///
/// # Errors
///
/// Returns [`SmithError::Session`] when `at` is unknown or the target already
/// holds a session, or a locking, framing, or input/output error.
pub fn fork_session(
    source: &Session,
    at: EntryId,
    path: &Path,
    config: Config,
) -> Result<OpenSession> {
    let forked = source.fork_at(at)?;
    create_session(path, &forked, config)
}

/// Load a session file, returning recovered frames and explicit damage.
///
/// Complete frames are returned in order with per-frame bounded allocation.
/// A torn tail is located rather than dropped silently, and corrupt or
/// oversized frames are reported with their byte offsets while the frames
/// around them stay readable.
///
/// # Errors
///
/// Returns an input/output error when the file cannot be opened.
pub fn load_session(path: &Path, max_frame_bytes: usize) -> Result<Recovery> {
    let file = File::open(path).map_err(io_err("STORE_OPEN"))?;
    let mut reader = BufReader::with_capacity(64 * 1024, file);
    Ok(read_frames(&mut reader, max_frame_bytes))
}

/// Repair a session file by republishing only its complete frames.
///
/// Damaged frames are dropped and returned in [`Recovery::damage`]. The
/// original content stays intact if any step fails before publication.
/// Repair takes the session lock, so it waits for a live writer to finish
/// instead of replacing the file underneath it.
///
/// # Errors
///
/// Returns a locking, framing, or input/output error.
pub fn repair_session(path: &Path, max_frame_bytes: usize) -> Result<Recovery> {
    let _lock = SessionLock::acquire(path)?;
    let recovery = load_session(path, max_frame_bytes)?;
    let tmp = path.with_extension("smith-repair.tmp");
    {
        let mut out = File::create(&tmp).map_err(io_err("REPAIR_CREATE"))?;
        for frame in &recovery.frames {
            write_frame(&mut out, frame, max_frame_bytes)?;
        }
        out.sync_all().map_err(io_err("REPAIR_SYNC"))?;
    }
    std::fs::rename(&tmp, path).map_err(io_err("REPAIR_PUBLISH"))?;
    Ok(recovery)
}

/// Walk an already-open file from its start using the shared frame reader.
///
/// # Errors
///
/// Returns an input/output error when the file cannot be re-read.
fn walk_file(file: &File, max_frame_bytes: usize) -> Result<Recovery> {
    let mut handle = file.try_clone().map_err(io_err("STORE_OPEN"))?;
    handle
        .seek(SeekFrom::Start(0))
        .map_err(io_err("STORE_SEEK"))?;
    let mut reader = BufReader::with_capacity(64 * 1024, handle);
    Ok(read_frames(&mut reader, max_frame_bytes))
}

fn io_err(code: &'static str) -> impl Fn(std::io::Error) -> SmithError {
    move |e| SmithError::Session {
        code: code.to_string(),
        message: e.to_string(),
    }
}

#[cfg(test)]
#[expect(
    clippy::unwrap_used,
    clippy::panic,
    reason = "tests may panic on invariant violations"
)]
mod tests {
    use super::*;
    use crate::session::COMPACTION_KIND;
    use crate::session::{EntryContent, EntryFrame, EntryKind, Session};
    use smith::message::{ContentBlock, Message, Role};

    fn tmp_path(name: &str) -> PathBuf {
        let dir = std::env::temp_dir().join(format!("smith_store_tests_{}", std::process::id()));
        std::fs::create_dir_all(&dir).unwrap();
        dir.join(name)
    }

    fn message(text: &str) -> EntryContent {
        EntryContent::Message(Message::with_text(Role::User, text))
    }

    fn entry(text: &str) -> EntryFrame {
        EntryFrame::new(None, &message(text)).unwrap()
    }

    fn text(entry: &EntryFrame) -> Option<String> {
        match EntryContent::decode(&entry.content).unwrap() {
            EntryContent::Message(m) => match &m.blocks[0] {
                ContentBlock::Text(t) => Some(t.clone()),
                _ => None,
            },
            _ => None,
        }
    }

    fn frame(text: &str) -> Frame {
        Frame::Known { entry: entry(text) }
    }

    fn text_of(frame: &Frame) -> String {
        match frame {
            Frame::Known { entry } => text(entry).unwrap(),
            other => panic!("unexpected frame {other:?}"),
        }
    }

    #[test]
    fn append_load_round_trip_preserves_entries() {
        let path = tmp_path("round_trip.smh");
        let _ = std::fs::remove_file(&path);
        let cfg = Config::default_valid();
        let first = entry("one");
        let second = EntryFrame::new(Some(first.id), &message("two")).unwrap();
        {
            let mut w = SessionWriter::open(&path, cfg.clone()).unwrap();
            w.append_frame(&Frame::Known {
                entry: first.clone(),
            })
            .unwrap();
            w.append_frame(&Frame::Known {
                entry: second.clone(),
            })
            .unwrap();
        }
        let frames = load_session(&path, cfg.max_frame_bytes).unwrap().frames;
        assert_eq!(frames.len(), 2);
        let session = Session::from_frames(&frames).unwrap();
        assert_eq!(
            session
                .active_branch()
                .entries()
                .iter()
                .map(|e| e.id)
                .collect::<Vec<_>>(),
            vec![first.id, second.id]
        );
        assert_eq!(session.selected_id(), Some(second.id));
        let _ = std::fs::remove_file(&path);
    }

    #[test]
    fn torn_tail_is_truncated_on_open_and_appends_stay_visible() {
        let path = tmp_path("torn_tail.smh");
        let _ = std::fs::remove_file(&path);
        let cfg = Config::default_valid();
        {
            let mut w = SessionWriter::open(&path, cfg.clone()).unwrap();
            w.append_frame(&frame("kept")).unwrap();
        }
        // Simulate a crash mid-append: partial header + payload bytes.
        let good = Frame::Known {
            entry: entry("kept"),
        }
        .encode(cfg.max_frame_bytes)
        .unwrap();
        let mut torn = good;
        torn.extend_from_slice(&[0x10, 0x00, 0x00, 0x00, 0xAA, 0xBB]);
        std::fs::write(&path, &torn).unwrap();
        {
            // Open repairs the tail, then a new append must be loadable.
            let mut w = SessionWriter::open(&path, cfg.clone()).unwrap();
            w.append_frame(&frame("after")).unwrap();
        }
        let frames = load_session(&path, cfg.max_frame_bytes).unwrap().frames;
        let texts: Vec<String> = frames
            .iter()
            .filter_map(|f| match f {
                Frame::Known { entry } => text(entry),
                _ => None,
            })
            .collect();
        assert_eq!(texts, vec!["kept".to_string(), "after".to_string()]);
        let _ = std::fs::remove_file(&path);
    }

    #[test]
    fn loader_stops_at_torn_tail_without_header_fabrication() {
        let path = tmp_path("loader_torn.smh");
        let cfg = Config::default_valid();
        let good = frame("a").encode(cfg.max_frame_bytes).unwrap();
        let mut bytes = good;
        bytes.extend_from_slice(&[0xFF, 0xFF]); // partial header
        std::fs::write(&path, &bytes).unwrap();
        let frames = load_session(&path, cfg.max_frame_bytes).unwrap().frames;
        assert_eq!(frames.len(), 1);
        let _ = std::fs::remove_file(&path);
    }

    #[test]
    fn reopening_a_multi_frame_session_keeps_every_frame() {
        let path = tmp_path("reopen_multi.smh");
        let _ = std::fs::remove_file(&path);
        let cfg = Config::default_valid();
        {
            let mut w = SessionWriter::open(&path, cfg.clone()).unwrap();
            w.append_frame(&frame("one")).unwrap();
            w.append_frame(&frame("two")).unwrap();
        }
        {
            // Reopening must not mistake a later frame boundary for a torn tail.
            let mut w = SessionWriter::open(&path, cfg.clone()).unwrap();
            w.append_frame(&frame("three")).unwrap();
        }
        let recovery = load_session(&path, cfg.max_frame_bytes).unwrap();
        assert_eq!(recovery.frames.len(), 3);
        assert!(recovery.damage.is_empty());
        assert_eq!(recovery.torn_tail_at, None);
        let _ = std::fs::remove_file(&path);
    }

    #[test]
    fn oversized_complete_frame_is_skipped_and_following_frames_reachable() {
        let path = tmp_path("oversized_skip.smh");
        // declared = 5000 > limit 1024, payload fully present, then a valid frame.
        let mut bytes = 5000_u32.to_le_bytes().to_vec();
        bytes.extend_from_slice(&vec![0xAB_u8; 5000]);
        bytes.extend_from_slice(&frame("tail").encode(1 << 20).unwrap());
        std::fs::write(&path, &bytes).unwrap();
        let recovery = load_session(&path, 1024).unwrap();
        assert_eq!(recovery.frames.len(), 1);
        assert!(matches!(recovery.frames[0], Frame::Known { .. }));
        assert_eq!(recovery.damage[0].code(), "FRAME_TOO_LARGE");
        let _ = std::fs::remove_file(&path);
    }

    #[test]
    fn repair_republishes_file_without_torn_tail() {
        let path = tmp_path("repair.smh");
        let cfg = Config::default_valid();
        let good = frame("a").encode(cfg.max_frame_bytes).unwrap();
        let mut bytes = good;
        bytes.extend_from_slice(&frame("b").encode(cfg.max_frame_bytes).unwrap()[..6]);
        std::fs::write(&path, &bytes).unwrap();
        let before = std::fs::metadata(&path).unwrap().len();
        let recovery = repair_session(&path, cfg.max_frame_bytes).unwrap();
        assert_eq!(recovery.frames.len(), 1);
        assert_eq!(recovery.torn_tail_at, Some(before - 6));
        let after = std::fs::metadata(&path).unwrap().len();
        assert!(after < before);
        assert_eq!(
            load_session(&path, cfg.max_frame_bytes)
                .unwrap()
                .frames
                .len(),
            1
        );
        let _ = std::fs::remove_file(&path);
    }

    #[test]
    fn writer_rejects_oversized_frame_without_touching_file() {
        let path = tmp_path("oversize_reject.smh");
        let _ = std::fs::remove_file(&path);
        let mut cfg = Config::default_valid();
        cfg.max_frame_bytes = 4096;
        {
            let mut w = SessionWriter::open(&path, cfg.clone()).unwrap();
            w.append_frame(&frame("small")).unwrap();
            let big = EntryFrame::new(
                None,
                &EntryContent::Message({
                    let mut m = Message::with_text(Role::User, "x");
                    m.add_block(ContentBlock::text("y".repeat(8192)));
                    m
                }),
            )
            .unwrap();
            let err = w.append_frame(&Frame::Known { entry: big }).unwrap_err();
            assert_eq!(err.code(), "FRAME_TOO_LARGE");
        }
        assert_eq!(
            load_session(&path, cfg.max_frame_bytes)
                .unwrap()
                .frames
                .len(),
            1
        );
        let _ = std::fs::remove_file(&path);
    }

    #[test]
    fn lock_serializes_second_writer_handle() {
        let path = tmp_path("locked.smh");
        let _ = std::fs::remove_file(&path);
        let cfg = Config::default_valid();
        let writer = SessionWriter::open(&path, cfg).unwrap();
        // While the first writer holds the session, an independent handle
        // cannot take it: writers cannot interleave frames.
        let probe = OpenOptions::new()
            .read(true)
            .write(true)
            .open(lock_path(&path))
            .unwrap();
        assert!(fs2::FileExt::try_lock_exclusive(&probe).is_err());
        drop(writer);
        assert!(fs2::FileExt::try_lock_exclusive(&probe).is_ok());
        let _ = std::fs::remove_file(&path);
    }

    #[test]
    fn repair_waits_for_a_live_writer_and_keeps_its_entries() {
        let path = tmp_path("repair_race.smh");
        let _ = std::fs::remove_file(&path);
        let _ = std::fs::remove_file(lock_path(&path));
        let cfg = Config::default_valid();
        let mut writer = SessionWriter::open(&path, cfg.clone()).unwrap();
        writer.append_frame(&frame("one")).unwrap();

        let repair_path = path.clone();
        let limit = cfg.max_frame_bytes;
        let (entered, repair_started) = std::sync::mpsc::sync_channel(1);
        let repair = std::thread::spawn(move || {
            entered.send(()).unwrap();
            repair_session(&repair_path, limit).unwrap()
        });

        // synchronize: the repair has started; the writer's lock holds it
        // until the drop, so an append now is still committed to the file
        // the repair will republish.
        repair_started
            .recv_timeout(std::time::Duration::from_secs(10))
            .unwrap();
        writer.append_frame(&frame("two")).unwrap();
        drop(writer);

        let recovery = repair.join().unwrap();
        assert_eq!(recovery.frames.len(), 2);
        let loaded = load_session(&path, cfg.max_frame_bytes).unwrap();
        let texts: Vec<String> = loaded.frames.iter().map(text_of).collect();
        assert_eq!(texts, vec!["one".to_string(), "two".to_string()]);
        let _ = std::fs::remove_file(&path);
        let _ = std::fs::remove_file(lock_path(&path));
    }

    #[test]
    fn unknown_frames_survive_write_and_load() {
        let path = tmp_path("unknown_frames.smh");
        let _ = std::fs::remove_file(&path);
        let cfg = Config::default_valid();
        let unknown = Frame::Unknown {
            kind: 0x4242,
            version: 3,
            payload: vec![1, 2, 3, 4, 5, 6],
        };
        {
            let mut w = SessionWriter::open(&path, cfg.clone()).unwrap();
            w.append_frame(&unknown).unwrap();
            w.append_frame(&frame("known")).unwrap();
        }
        let frames = load_session(&path, cfg.max_frame_bytes).unwrap().frames;
        assert_eq!(frames.len(), 2);
        match &frames[0] {
            Frame::Unknown {
                kind,
                version,
                payload,
            } => {
                assert_eq!(payload, &vec![1, 2, 3, 4, 5, 6]);
                assert_eq!((*kind, *version), (0x4242, 3));
            }
            other => panic!("expected unknown frame, got {other:?}"),
        }
        match &frames[1] {
            Frame::Known { entry } => assert_eq!(text(entry).as_deref(), Some("known")),
            other => panic!("expected known frame, got {other:?}"),
        }
        let _ = std::fs::remove_file(&path);
    }

    #[test]
    fn end_to_end_tool_activity_persists_and_reconstructs() {
        let dir = tmp_path("e2e");
        std::fs::create_dir_all(&dir).unwrap();
        let workdir = dir.join("work");
        std::fs::create_dir_all(&workdir).unwrap();
        let session_path = dir.join("e2e.smh");
        let _ = std::fs::remove_file(&session_path);
        let cfg = Config::default_valid();
        let opened_id;
        {
            let opened = open_session(&session_path, cfg.clone()).unwrap();
            opened_id = opened.session.id();
            let mut ts = crate::tools::ToolSession::with_session(
                workdir.clone(),
                opened.session,
                opened.writer,
            );
            ts.invoke(
                "write",
                &serde_json::json!({"path": "out/hello.txt", "content": "hi"}),
            )
            .unwrap();
            ts.invoke("read", &serde_json::json!({"path": "out/hello.txt"}))
                .unwrap();
        }
        let frames = load_session(&session_path, cfg.max_frame_bytes)
            .unwrap()
            .frames;
        assert_eq!(frames.len(), 5); // header + call/result x2
        let restored = Session::from_frames(&frames).unwrap();
        assert_eq!(restored.id(), opened_id);
        assert_eq!(restored.active_branch().entries().len(), 4);
        // Tool call/result pairs stay adjacent and ordered.
        let kinds: Vec<&str> = restored
            .active_branch()
            .entries()
            .iter()
            .map(|e| match e.kind {
                EntryKind::ToolCall { .. } => "call",
                EntryKind::ToolResult { .. } => "result",
                _ => "other",
            })
            .collect();
        assert_eq!(kinds, vec!["call", "result", "call", "result"]);
        // The file effect actually happened under the invocation workdir.
        assert_eq!(
            std::fs::read_to_string(workdir.join("out/hello.txt")).unwrap(),
            "hi"
        );
        let _ = std::fs::remove_file(&session_path);
        let _ = std::fs::remove_dir_all(&dir);
    }

    #[test]
    fn reopened_session_keeps_its_identity_and_replays_identically() {
        let path = tmp_path("identity.smh");
        let _ = std::fs::remove_file(&path);
        let cfg = Config::default_valid();
        let created = {
            let mut opened = open_session(&path, cfg.clone()).unwrap();
            opened.writer.append_frame(&frame("one")).unwrap();
            opened.session
        };
        // The lock is exclusive, so each load owns the file for its lifetime.
        let first = {
            let opened = open_session(&path, cfg.clone()).unwrap();
            opened.session
        };
        let second = {
            let opened = open_session(&path, cfg).unwrap();
            opened.session
        };
        assert_eq!(first.id(), created.id());
        assert_eq!(first.active_branch().id(), created.active_branch().id());
        assert_eq!(second.id(), first.id());

        // Replay evidence is byte-stable across independent loads.
        let encode = |session: &Session| {
            let mut bytes = Vec::new();
            ciborium::ser::into_writer(&crate::trace::replay(session), &mut bytes).unwrap();
            bytes
        };
        assert_eq!(encode(&first), encode(&second));
        let _ = std::fs::remove_file(&path);
    }

    #[test]
    fn forking_publishes_a_new_session_file_and_leaves_the_source_intact() {
        let dir = tmp_path("fork_dir");
        std::fs::create_dir_all(&dir).unwrap();
        let source_path = dir.join("source.smh");
        let fork_path = dir.join("fork.smh");
        let cfg = Config::default_valid();

        let (source_session, boundary_id, leaf_id) = {
            let mut opened = open_session(&source_path, cfg.clone()).unwrap();
            let root = *opened.session.append(&message("root")).unwrap();
            opened
                .writer
                .append_entry(&root, opened.session.bytes(&root))
                .unwrap();
            let boundary = *opened
                .session
                .append(&EntryContent::Meta {
                    kind: COMPACTION_KIND.to_string(),
                    detail: serde_json::json!({"kept": 1}),
                })
                .unwrap();
            opened
                .writer
                .append_entry(&boundary, opened.session.bytes(&boundary))
                .unwrap();
            let tip = *opened.session.append(&message("after")).unwrap();
            opened
                .writer
                .append_entry(&tip, opened.session.bytes(&tip))
                .unwrap();
            assert_ne!(root.id, tip.id);
            (opened.session, boundary.id, tip.id)
        };

        let fork_id = {
            let forked =
                fork_session(&source_session, boundary_id, &fork_path, cfg.clone()).unwrap();
            assert_ne!(forked.session.id(), source_session.id());
            assert_eq!(forked.session.selected_id(), Some(boundary_id));
            assert_eq!(forked.session.compaction_boundary(), Some(boundary_id));
            assert!(forked.recovery.damage.is_empty());
            forked.session.id()
        };

        // The fork is durable: reopening restores the same identity and leaf.
        {
            let reopened = open_session(&fork_path, cfg.clone()).unwrap();
            assert_eq!(reopened.session.id(), fork_id);
            assert_eq!(reopened.session.selected_id(), Some(boundary_id));
            assert_eq!(reopened.session.active_branch().entries().len(), 2);
            assert_eq!(reopened.session.compaction_boundary(), Some(boundary_id));
            // Published content bytes decode to what the source recorded.
            let content = |session: &Session| {
                session
                    .content(&session.active_branch().entries()[0])
                    .unwrap()
            };
            assert_eq!(content(&reopened.session), content(&source_session));
        }

        // The source keeps its own identity, tip, and entry count.
        {
            let source_again = open_session(&source_path, cfg.clone()).unwrap();
            assert_eq!(source_again.session.id(), source_session.id());
            assert_eq!(source_again.session.selected_id(), Some(leaf_id));
            assert_eq!(source_again.session.active_branch().entries().len(), 3);
        }

        // A populated target is refused instead of overwritten.
        let err = fork_session(&source_session, boundary_id, &fork_path, cfg).unwrap_err();
        assert_eq!(err.code(), "SESSION_EXISTS");
        let _ = std::fs::remove_dir_all(&dir);
    }

    #[test]
    fn session_file_without_header_is_rejected() {
        let path = tmp_path("headerless.smh");
        let cfg = Config::default_valid();
        std::fs::write(&path, frame("orphan").encode(cfg.max_frame_bytes).unwrap()).unwrap();
        let err = open_session(&path, cfg).unwrap_err();
        assert_eq!(err.code(), "SESSION_HEADER_MISSING");
        let _ = std::fs::remove_file(&path);
    }

    #[test]
    fn header_len_matches_protocol() {
        assert_eq!(crate::frame::HEADER_LEN, 4);
        assert_eq!(crate::frame::BODY_PREFIX_LEN, 8);
    }
}