use std::{
collections::HashMap,
path::{Path, PathBuf},
sync::{Arc, Mutex, OnceLock, Weak},
};
use anyhow::{Context as _, bail, ensure};
use kcode_checksummed_frame_log::{
AppendDurability, IncompleteTailPolicy, append as append_frame, create as create_frame_log,
visit_frames,
};
use kcode_pending_object_store::{PendingObjectStore, StoredPendingObject};
use serde::{Deserialize, Serialize};
pub const FORMAT_VERSION: &str = "0.2.1";
const SESSION_MAGIC: &[u8] = b"KSESSIONLOG\n";
const HEADER_FRAME: u8 = 1;
const EVENT_FRAME: u8 = 2;
const SEALED_FRAME: u8 = 3;
#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct SessionHeader {
pub format_version: String,
pub session_id: String,
pub created_at: String,
}
#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq, Serialize, Deserialize)]
#[serde(rename_all = "kebab-case")]
pub enum Role {
SystemMessage,
SystemError,
UserMessage,
KennedyMessage,
KennedyToolCall,
ToolResult,
ToolError,
Object,
PendingObject,
}
#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
pub struct SessionEvent {
pub role: Role,
pub text: String,
}
#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
pub struct SessionLog {
pub header: SessionHeader,
pub events: Vec<SessionEvent>,
}
#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
pub struct EventPosition(pub u64);
impl EventPosition {
pub fn index(self) -> u64 {
self.0
}
}
impl std::fmt::Display for EventPosition {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
self.0.fmt(formatter)
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct PendingObject {
pub event_position: EventPosition,
pub text: String,
pub file_name: String,
pub media_type: String,
pub bytes: Vec<u8>,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct SealedSession {
log: SessionLog,
directory: PathBuf,
}
impl SealedSession {
pub fn list(&self) -> &SessionLog {
&self.log
}
pub fn pending_objects(&self) -> anyhow::Result<Vec<PendingObject>> {
let store = pending_object_store(&self.directory, &self.log.header.session_id)?;
self.log
.events
.iter()
.enumerate()
.filter(|(_, event)| event.role == Role::PendingObject)
.map(|(position, event)| {
read_pending_object(&store, EventPosition(position as u64), event.text.clone())
})
.collect()
}
}
#[derive(Clone, Debug)]
pub struct SessionStore {
directory: PathBuf,
}
impl SessionStore {
pub fn new(directory: impl Into<PathBuf>) -> Self {
Self {
directory: directory.into(),
}
}
pub fn directory(&self) -> &Path {
&self.directory
}
pub fn create_session(
&self,
session_id: impl Into<String>,
created_at: impl Into<String>,
) -> anyhow::Result<Session> {
let session_id = session_id.into();
let created_at = created_at.into();
validate_session_id(&session_id)?;
ensure!(
!created_at.trim().is_empty(),
"session creation time cannot be empty"
);
std::fs::create_dir_all(&self.directory)
.with_context(|| format!("creating {}", self.directory.display()))?;
let path = session_path(&self.directory, &session_id);
let append_lock = session_lock(&path);
let _guard = append_lock
.lock()
.map_err(|_| anyhow::anyhow!("session-log append lock is poisoned"))?;
let header = SessionHeader {
format_version: FORMAT_VERSION.into(),
session_id,
created_at,
};
create_frame_log(
&path,
SESSION_MAGIC,
HEADER_FRAME,
&serde_json::to_vec(&header)?,
)?;
drop(_guard);
Ok(Session {
directory: self.directory.clone(),
path,
log: SessionLog {
header,
events: Vec::new(),
},
sealed: false,
append_lock,
})
}
pub fn open_session(&self, session_id: &str) -> anyhow::Result<Session> {
validate_session_id(session_id)?;
let path = session_path(&self.directory, session_id);
let append_lock = session_lock(&path);
let _guard = append_lock
.lock()
.map_err(|_| anyhow::anyhow!("session-log append lock is poisoned"))?;
let loaded = load_session_file(&path, true)?;
ensure!(
loaded.log.header.session_id == session_id,
"session filename and header identity differ"
);
let object_store = pending_object_store(&self.directory, &loaded.log.header.session_id)?;
object_store.reconcile(&referenced_pending_positions(&loaded.log))?;
drop(_guard);
Ok(Session {
directory: self.directory.clone(),
path,
log: loaded.log,
sealed: loaded.sealed,
append_lock,
})
}
pub fn session_ids(&self) -> anyhow::Result<Vec<String>> {
if !self.directory.exists() {
return Ok(Vec::new());
}
let mut ids = std::fs::read_dir(&self.directory)?
.filter_map(Result::ok)
.filter_map(|entry| {
entry
.file_name()
.to_str()
.and_then(|name| name.strip_suffix(".session-log"))
.map(str::to_owned)
})
.filter(|id| validate_session_id(id).is_ok())
.collect::<Vec<_>>();
ids.sort();
Ok(ids)
}
}
pub struct Session {
directory: PathBuf,
path: PathBuf,
log: SessionLog,
sealed: bool,
append_lock: Arc<Mutex<()>>,
}
impl Session {
pub fn path(&self) -> &Path {
&self.path
}
pub fn list(&self) -> SessionLog {
self.log.clone()
}
pub fn is_sealed(&self) -> bool {
self.sealed
}
pub fn add_event(
&mut self,
role: Role,
text: impl Into<String>,
) -> anyhow::Result<EventPosition> {
ensure!(
role != Role::PendingObject,
"pending objects must be added through add_pending_object"
);
let event = SessionEvent {
role,
text: text.into(),
};
let append_lock = self.append_lock.clone();
let _guard = append_lock
.lock()
.map_err(|_| anyhow::anyhow!("session-log append lock is poisoned"))?;
self.refresh_locked()?;
ensure!(!self.sealed, "session is sealed");
let position = EventPosition(self.log.events.len() as u64);
append_event_file(&self.path, &event)?;
self.log.events.push(event);
Ok(position)
}
pub fn add_pending_object(
&mut self,
text: impl Into<String>,
file_name: impl Into<String>,
media_type: impl Into<String>,
bytes: &[u8],
) -> anyhow::Result<EventPosition> {
let text = text.into();
let file_name = file_name.into();
let media_type = media_type.into();
ensure!(
!file_name.trim().is_empty(),
"object filename cannot be empty"
);
ensure!(
!media_type.trim().is_empty(),
"object media type cannot be empty"
);
let append_lock = self.append_lock.clone();
let _guard = append_lock
.lock()
.map_err(|_| anyhow::anyhow!("session-log append lock is poisoned"))?;
self.refresh_locked()?;
ensure!(!self.sealed, "session is sealed");
let position = EventPosition(self.log.events.len() as u64);
let object_store = pending_object_store(&self.directory, &self.log.header.session_id)?;
object_store.install(position.0, &file_name, &media_type, bytes)?;
let event = SessionEvent {
role: Role::PendingObject,
text,
};
append_event_file(&self.path, &event)?;
self.log.events.push(event);
Ok(position)
}
pub fn read_pending_object(&self, position: EventPosition) -> anyhow::Result<PendingObject> {
let index = usize::try_from(position.0).context("event position does not fit memory")?;
let event = self
.log
.events
.get(index)
.with_context(|| format!("event {position} does not exist"))?;
ensure!(
event.role == Role::PendingObject,
"event {position} is not a pending object"
);
let store = pending_object_store(&self.directory, &self.log.header.session_id)?;
read_pending_object(&store, position, event.text.clone())
}
pub fn seal(&mut self) -> anyhow::Result<SealedSession> {
let append_lock = self.append_lock.clone();
let _guard = append_lock
.lock()
.map_err(|_| anyhow::anyhow!("session-log append lock is poisoned"))?;
self.refresh_locked()?;
if !self.sealed {
verify_referenced_objects(&self.directory, &self.log)?;
append_frame(
&self.path,
SEALED_FRAME,
&[],
AppendDurability::FileAndParentDirectory,
)?;
self.sealed = true;
}
Ok(SealedSession {
log: self.log.clone(),
directory: self.directory.clone(),
})
}
pub fn delete_committed(self) -> anyhow::Result<()> {
self.delete_files()
}
pub fn delete_abandoned(self) -> anyhow::Result<()> {
self.delete_files()
}
fn refresh_locked(&mut self) -> anyhow::Result<()> {
let loaded = load_session_file(&self.path, true)?;
ensure!(
loaded.log.header.session_id == self.log.header.session_id,
"session identity changed on disk"
);
self.log = loaded.log;
self.sealed = loaded.sealed;
Ok(())
}
fn delete_files(self) -> anyhow::Result<()> {
let append_lock = self.append_lock.clone();
let _guard = append_lock
.lock()
.map_err(|_| anyhow::anyhow!("session-log append lock is poisoned"))?;
let object_store = pending_object_store(&self.directory, &self.log.header.session_id)?;
if self.path.exists() {
std::fs::remove_file(&self.path)
.with_context(|| format!("removing {}", self.path.display()))?;
}
object_store.delete_all()
}
}
struct LoadedSession {
log: SessionLog,
sealed: bool,
}
struct SessionFrameVisitor {
header: Option<SessionHeader>,
events: Vec<SessionEvent>,
sealed: bool,
frame_count: u64,
}
impl SessionFrameVisitor {
fn new() -> Self {
Self {
header: None,
events: Vec::new(),
sealed: false,
frame_count: 0,
}
}
fn visit(&mut self, kind: u8, payload: &[u8]) -> anyhow::Result<()> {
match kind {
HEADER_FRAME => {
ensure!(self.frame_count == 0, "duplicate session header");
let value: SessionHeader = serde_json::from_slice(payload)?;
ensure!(
value.format_version == FORMAT_VERSION,
"unsupported session-log format {}",
value.format_version
);
validate_session_id(&value.session_id)?;
ensure!(
!value.created_at.trim().is_empty(),
"session creation time cannot be empty"
);
self.header = Some(value);
}
EVENT_FRAME => {
ensure!(self.header.is_some(), "session event precedes header");
ensure!(!self.sealed, "session event follows sealed footer");
self.events.push(serde_json::from_slice(payload)?);
}
SEALED_FRAME => {
ensure!(self.header.is_some(), "sealed footer precedes header");
ensure!(payload.is_empty(), "sealed footer payload must be empty");
ensure!(!self.sealed, "duplicate sealed footer");
self.sealed = true;
}
other => bail!("unknown complete session-log frame kind {other}"),
}
self.frame_count += 1;
Ok(())
}
fn finish(self) -> anyhow::Result<LoadedSession> {
let header = self.header.context("session log has no header")?;
Ok(LoadedSession {
log: SessionLog {
header,
events: self.events,
},
sealed: self.sealed,
})
}
}
fn validate_session_id(session_id: &str) -> anyhow::Result<()> {
ensure!(!session_id.is_empty(), "session ID cannot be empty");
ensure!(session_id.len() <= 255, "session ID exceeds 255 characters");
ensure!(
session_id
.bytes()
.all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'-' | b'_')),
"session ID contains characters that are unsafe in filenames"
);
Ok(())
}
fn session_path(directory: &Path, session_id: &str) -> PathBuf {
directory.join(format!("{session_id}.session-log"))
}
fn append_event_file(path: &Path, event: &SessionEvent) -> anyhow::Result<()> {
append_frame(
path,
EVENT_FRAME,
&serde_json::to_vec(event)?,
AppendDurability::FileOnly,
)
}
fn load_session_file(path: &Path, repair_tail: bool) -> anyhow::Result<LoadedSession> {
let tail_policy = if repair_tail {
IncompleteTailPolicy::TruncateAndSync
} else {
IncompleteTailPolicy::Reject
};
let mut visitor = SessionFrameVisitor::new();
visit_frames(path, SESSION_MAGIC, tail_policy, |kind, payload| {
visitor.visit(kind, payload)
})?;
visitor.finish()
}
fn pending_object_store(directory: &Path, session_id: &str) -> anyhow::Result<PendingObjectStore> {
PendingObjectStore::new(directory, session_id, FORMAT_VERSION)
}
fn read_pending_object(
store: &PendingObjectStore,
position: EventPosition,
text: String,
) -> anyhow::Result<PendingObject> {
let StoredPendingObject {
file_name,
media_type,
bytes,
} = store.read(position.0)?;
Ok(PendingObject {
event_position: position,
text,
file_name,
media_type,
bytes,
})
}
fn referenced_pending_positions(log: &SessionLog) -> Vec<u64> {
log.events
.iter()
.enumerate()
.filter_map(|(position, event)| {
(event.role == Role::PendingObject).then_some(position as u64)
})
.collect()
}
fn verify_referenced_objects(directory: &Path, log: &SessionLog) -> anyhow::Result<()> {
pending_object_store(directory, &log.header.session_id)?
.verify_all(&referenced_pending_positions(log))
}
fn session_lock(path: &Path) -> Arc<Mutex<()>> {
static LOCKS: OnceLock<Mutex<HashMap<PathBuf, Weak<Mutex<()>>>>> = OnceLock::new();
let locks = LOCKS.get_or_init(|| Mutex::new(HashMap::new()));
let mut locks = locks.lock().expect("session-log lock registry is poisoned");
locks.retain(|_, lock| lock.strong_count() > 0);
if let Some(lock) = locks.get(path).and_then(Weak::upgrade) {
return lock;
}
let lock = Arc::new(Mutex::new(()));
locks.insert(path.to_path_buf(), Arc::downgrade(&lock));
lock
}
#[cfg(test)]
mod tests {
use std::{
fs::OpenOptions,
io::Write,
time::{SystemTime, UNIX_EPOCH},
};
use super::*;
const FRAME_HEADER_LEN: usize = 1 + 8 + 32;
fn directory(label: &str) -> PathBuf {
std::env::temp_dir().join(format!(
"session-log-{label}-{}-{}",
std::process::id(),
SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap()
.as_nanos()
))
}
#[test]
fn events_are_an_ordered_role_and_text_array_without_serialized_ids() {
let directory = directory("ordered");
let store = SessionStore::new(&directory);
let mut session = store
.create_session("session-1", "2026-07-24T00:00:00Z")
.unwrap();
assert_eq!(
session.add_event(Role::SystemMessage, "system").unwrap(),
EventPosition(0)
);
assert_eq!(
session.add_event(Role::UserMessage, "hello").unwrap(),
EventPosition(1)
);
drop(session);
let reopened = store.open_session("session-1").unwrap();
assert_eq!(
reopened.list(),
SessionLog {
header: SessionHeader {
format_version: "0.2.1".into(),
session_id: "session-1".into(),
created_at: "2026-07-24T00:00:00Z".into(),
},
events: vec![
SessionEvent {
role: Role::SystemMessage,
text: "system".into(),
},
SessionEvent {
role: Role::UserMessage,
text: "hello".into(),
},
],
}
);
let serialized = serde_json::to_string(&reopened.list().events).unwrap();
assert!(!serialized.contains("\"id\""));
std::fs::remove_dir_all(directory).unwrap();
}
#[test]
fn pending_object_is_durable_before_its_event_and_uses_event_position() {
let directory = directory("object");
let store = SessionStore::new(&directory);
let mut session = store
.create_session("object-session", "2026-07-24T00:00:00Z")
.unwrap();
session.add_event(Role::UserMessage, "upload").unwrap();
let position = session
.add_pending_object("notes.txt", "notes.txt", "text/plain", b"durable bytes")
.unwrap();
assert_eq!(position, EventPosition(1));
assert!(directory.join("object-session-1.pending-object").exists());
let object = session.read_pending_object(position).unwrap();
assert_eq!(object.file_name, "notes.txt");
assert_eq!(object.media_type, "text/plain");
assert_eq!(object.bytes, b"durable bytes");
std::fs::remove_dir_all(directory).unwrap();
}
#[test]
fn open_removes_unreferenced_final_and_temporary_objects() {
let directory = directory("orphans");
let store = SessionStore::new(&directory);
let session = store
.create_session("orphan-session", "2026-07-24T00:00:00Z")
.unwrap();
drop(session);
std::fs::write(directory.join("orphan-session-0.pending-object"), b"orphan").unwrap();
std::fs::write(
directory.join("orphan-session-1.pending-object.tmp"),
b"temporary",
)
.unwrap();
let leading_zero = directory.join("orphan-session-00.pending-object");
let similar = directory.join("orphan-session-other-2.pending-object");
std::fs::write(&leading_zero, b"keep").unwrap();
std::fs::write(&similar, b"keep").unwrap();
store.open_session("orphan-session").unwrap();
assert!(!directory.join("orphan-session-0.pending-object").exists());
assert!(
!directory
.join("orphan-session-1.pending-object.tmp")
.exists()
);
assert!(leading_zero.exists());
assert!(similar.exists());
std::fs::remove_dir_all(directory).unwrap();
}
#[test]
fn incomplete_event_tail_is_discarded() {
let directory = directory("tail");
let store = SessionStore::new(&directory);
let mut session = store
.create_session("tail-session", "2026-07-24T00:00:00Z")
.unwrap();
session.add_event(Role::UserMessage, "complete").unwrap();
let path = session.path().to_path_buf();
drop(session);
let valid_len = std::fs::metadata(&path).unwrap().len();
OpenOptions::new()
.append(true)
.open(&path)
.unwrap()
.write_all(&[EVENT_FRAME, 20, 0, 0])
.unwrap();
let reopened = store.open_session("tail-session").unwrap();
assert_eq!(reopened.list().events.len(), 1);
assert_eq!(std::fs::metadata(&path).unwrap().len(), valid_len);
std::fs::remove_dir_all(directory).unwrap();
}
#[test]
fn checksum_invalid_complete_frame_is_corruption_not_a_recoverable_tail() {
let directory = directory("checksum");
let store = SessionStore::new(&directory);
let mut session = store
.create_session("checksum-session", "2026-07-24T00:00:00Z")
.unwrap();
session.add_event(Role::UserMessage, "complete").unwrap();
let path = session.path().to_path_buf();
drop(session);
let mut bytes = std::fs::read(&path).unwrap();
let payload_byte = SESSION_MAGIC.len() + FRAME_HEADER_LEN;
bytes[payload_byte] ^= 0xff;
std::fs::write(&path, &bytes).unwrap();
assert!(store.open_session("checksum-session").is_err());
assert_eq!(std::fs::read(&path).unwrap(), bytes);
std::fs::remove_dir_all(directory).unwrap();
}
#[test]
fn seal_is_durable_idempotent_and_rejects_later_events() {
let directory = directory("seal");
let store = SessionStore::new(&directory);
let mut session = store
.create_session("sealed-session", "2026-07-24T00:00:00Z")
.unwrap();
session.add_event(Role::UserMessage, "hello").unwrap();
session.seal().unwrap();
session.seal().unwrap();
assert!(session.add_event(Role::KennedyMessage, "too late").is_err());
drop(session);
assert!(store.open_session("sealed-session").unwrap().is_sealed());
std::fs::remove_dir_all(directory).unwrap();
}
#[test]
fn deletion_matches_only_exact_session_object_names() {
let directory = directory("delete");
let store = SessionStore::new(&directory);
let session = store.create_session("abc", "2026-07-24T00:00:00Z").unwrap();
let final_object = directory.join("abc-0.pending-object");
let temporary_object = directory.join("abc-1.pending-object.tmp");
let similar = directory.join("abc-other-0.pending-object");
let malformed = directory.join("abc-not-a-number.pending-object");
let leading_zero = directory.join("abc-01.pending-object");
std::fs::write(&final_object, b"remove").unwrap();
std::fs::write(&temporary_object, b"remove").unwrap();
std::fs::write(&similar, b"keep").unwrap();
std::fs::write(&malformed, b"keep").unwrap();
std::fs::write(&leading_zero, b"keep").unwrap();
session.delete_abandoned().unwrap();
assert!(!final_object.exists());
assert!(!temporary_object.exists());
assert!(similar.exists());
assert!(malformed.exists());
assert!(leading_zero.exists());
std::fs::remove_dir_all(directory).unwrap();
}
#[test]
fn frame_kinds_and_json_payloads_remain_legacy_compatible() {
let directory = directory("frame-mapping");
let store = SessionStore::new(&directory);
let mut session = store
.create_session("mapping-session", "2026-07-24T00:00:00Z")
.unwrap();
session.add_event(Role::ToolResult, "ok").unwrap();
session
.add_pending_object(
"pending text",
"blob.bin",
"application/octet-stream",
b"\0\xff",
)
.unwrap();
let path = session.path().to_path_buf();
session.seal().unwrap();
let mut frames = Vec::new();
visit_frames(
&path,
SESSION_MAGIC,
IncompleteTailPolicy::Reject,
|kind, payload| {
frames.push((kind, payload.to_vec()));
Ok(())
},
)
.unwrap();
assert_eq!(
frames,
vec![
(
HEADER_FRAME,
br#"{"formatVersion":"0.2.1","sessionId":"mapping-session","createdAt":"2026-07-24T00:00:00Z"}"#
.to_vec(),
),
(
EVENT_FRAME,
br#"{"role":"tool-result","text":"ok"}"#.to_vec(),
),
(
EVENT_FRAME,
br#"{"role":"pending-object","text":"pending text"}"#.to_vec(),
),
(SEALED_FRAME, Vec::new()),
]
);
std::fs::remove_dir_all(directory).unwrap();
}
#[test]
fn pending_object_sidecar_mapping_matches_extracted_store_exactly() {
let directory = directory("sidecar-mapping");
let store = SessionStore::new(&directory);
let mut session = store
.create_session("sidecar-session", "2026-07-24T00:00:00Z")
.unwrap();
let position = session
.add_pending_object("display text", "notes.txt", "text/plain", b"durable bytes")
.unwrap();
let expected_directory = directory.join("expected");
std::fs::create_dir(&expected_directory).unwrap();
PendingObjectStore::new(&expected_directory, "sidecar-session", "0.2.1")
.unwrap()
.install(0, "notes.txt", "text/plain", b"durable bytes")
.unwrap();
let actual_path = directory.join("sidecar-session-0.pending-object");
let expected_path = expected_directory.join("sidecar-session-0.pending-object");
assert_eq!(
std::fs::read(&actual_path).unwrap(),
std::fs::read(&expected_path).unwrap()
);
assert!(
!directory
.join("sidecar-session-0.pending-object.tmp")
.exists()
);
assert_eq!(
session.read_pending_object(position).unwrap(),
PendingObject {
event_position: EventPosition(0),
text: "display text".into(),
file_name: "notes.txt".into(),
media_type: "text/plain".into(),
bytes: b"durable bytes".to_vec(),
}
);
std::fs::remove_dir_all(directory).unwrap();
}
#[test]
fn opens_legacy_0_2_1_log_and_sidecar_bytes() {
let directory = directory("legacy");
std::fs::create_dir_all(&directory).unwrap();
let path = directory.join("legacy-session.session-log");
create_frame_log(
&path,
SESSION_MAGIC,
HEADER_FRAME,
br#"{"formatVersion":"0.2.1","sessionId":"legacy-session","createdAt":"2026-07-24T00:00:00Z"}"#,
)
.unwrap();
append_frame(
&path,
EVENT_FRAME,
br#"{"role":"user-message","text":"legacy event"}"#,
AppendDurability::FileOnly,
)
.unwrap();
let object_store = PendingObjectStore::new(&directory, "legacy-session", "0.2.1").unwrap();
object_store
.install(1, "legacy.bin", "application/octet-stream", b"legacy bytes")
.unwrap();
append_frame(
&path,
EVENT_FRAME,
br#"{"role":"pending-object","text":"legacy object"}"#,
AppendDurability::FileOnly,
)
.unwrap();
append_frame(
&path,
SEALED_FRAME,
&[],
AppendDurability::FileAndParentDirectory,
)
.unwrap();
let store = SessionStore::new(&directory);
let reopened = store.open_session("legacy-session").unwrap();
assert!(reopened.is_sealed());
assert_eq!(reopened.list().header.format_version, "0.2.1");
assert_eq!(
reopened.list().events,
vec![
SessionEvent {
role: Role::UserMessage,
text: "legacy event".into(),
},
SessionEvent {
role: Role::PendingObject,
text: "legacy object".into(),
},
]
);
assert_eq!(
reopened
.read_pending_object(EventPosition(1))
.unwrap()
.bytes,
b"legacy bytes"
);
assert!(directory.join("legacy-session.session-log").exists());
assert!(directory.join("legacy-session-1.pending-object").exists());
std::fs::remove_dir_all(directory).unwrap();
}
}