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);
}
}