use async_trait::async_trait;
use chrono::Utc;
use everruns_core::error::{AgentLoopError, Result};
use everruns_core::events::{
Event, EventData, EventRequest, INPUT_MESSAGE, OUTPUT_MESSAGE_COMPLETED,
OutputMessageCompletedData, REASON_COMPLETED, REASON_ITEM, TOOL_COMPLETED,
};
use everruns_core::message::{ContentPart, Message};
use everruns_core::tools::ToolResultImage;
use everruns_core::traits::EventEmitter;
use everruns_core::typed_id::EventId;
use everruns_core::typed_id::SessionId;
use everruns_runtime::EventBus;
use std::fs::{File, OpenOptions, TryLockError};
use std::io::{BufRead, BufReader, Seek, SeekFrom, Write};
use std::path::{Path, PathBuf};
use std::sync::Arc;
use tokio::sync::{Mutex, RwLock, broadcast};
const EVENT_BROADCAST_CAPACITY: usize = 1024;
pub fn default_sessions_dir() -> Result<PathBuf> {
dirs::data_dir()
.map(|p| p.join("yolop").join("sessions"))
.ok_or_else(|| {
AgentLoopError::config(
"could not resolve a platform data directory for session logs; \
pass --session-dir <PATH> explicitly",
)
})
}
pub fn session_dir_path(sessions_dir: &Path, session_id: SessionId) -> PathBuf {
sessions_dir.join(session_id.to_string())
}
pub fn session_log_path(session_dir: &Path) -> PathBuf {
session_dir.join("events.jsonl")
}
pub fn legacy_session_log_path(sessions_dir: &Path, session_id: SessionId) -> PathBuf {
sessions_dir.join(format!("{session_id}.jsonl"))
}
pub fn migrate_legacy_session_log(
sessions_dir: &Path,
session_dir: &Path,
session_id: SessionId,
) -> Result<Option<PathBuf>> {
let current = session_log_path(session_dir);
if current.exists() {
return Ok(None);
}
let legacy = legacy_session_log_path(sessions_dir, session_id);
if !legacy.exists() {
return Ok(None);
}
std::fs::create_dir_all(session_dir).map_err(|e| {
AgentLoopError::config(format!("create session dir {}: {e}", session_dir.display()))
})?;
std::fs::copy(&legacy, ¤t).map_err(|e| {
AgentLoopError::config(format!(
"migrate legacy session log {} to {}: {e}",
legacy.display(),
current.display()
))
})?;
Ok(Some(legacy))
}
#[derive(Debug, Default)]
pub struct ReplayedSession {
pub events: Vec<Event>,
pub messages: Vec<Message>,
pub max_sequence: Option<i32>,
}
pub fn replay(path: &Path, expected: SessionId) -> Result<ReplayedSession> {
if !path.exists() {
return Ok(ReplayedSession::default());
}
let file = File::open(path)
.map_err(|e| AgentLoopError::config(format!("open session log {}: {e}", path.display())))?;
let mut out = ReplayedSession::default();
for (i, line) in BufReader::new(file).lines().enumerate() {
let Ok(line) = line else {
tracing::warn!(
line = i + 1,
"session log read error; stopping replay early"
);
break;
};
if line.trim().is_empty() {
continue;
}
match serde_json::from_str::<Event>(&line) {
Ok(event) => {
if event.session_id != expected {
tracing::warn!(
line = i + 1,
found = %event.session_id,
expected = %expected,
"session log line belongs to a different session; skipping"
);
continue;
}
if let Some(seq) = event.sequence {
out.max_sequence = Some(out.max_sequence.map_or(seq, |m| m.max(seq)));
}
if let Some(message) = message_from_event(&event.data) {
out.messages.push(message);
}
out.events.push(event);
}
Err(e) => {
tracing::warn!(line = i + 1, error = %e, "skipping malformed session-log line");
}
}
}
Ok(out)
}
fn is_replay_relevant(event_type: &str) -> bool {
matches!(
event_type,
INPUT_MESSAGE | OUTPUT_MESSAGE_COMPLETED | TOOL_COMPLETED | REASON_COMPLETED | REASON_ITEM
)
}
pub struct JsonlEventEmitter {
events: Arc<RwLock<Vec<Event>>>,
sequence: Arc<RwLock<i32>>,
file: Arc<Mutex<File>>,
live: broadcast::Sender<Event>,
}
impl JsonlEventEmitter {
pub fn open(path: &Path, start_sequence: i32) -> Result<Self> {
if let Some(parent) = path.parent() {
std::fs::create_dir_all(parent).map_err(|e| {
AgentLoopError::config(format!("create session log dir {}: {e}", parent.display()))
})?;
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
std::fs::set_permissions(parent, std::fs::Permissions::from_mode(0o700)).map_err(
|e| {
AgentLoopError::config(format!(
"tighten session dir permissions on {}: {e}",
parent.display()
))
},
)?;
}
}
let mut opts = OpenOptions::new();
opts.create(true).append(true).read(true);
#[cfg(unix)]
{
use std::os::unix::fs::OpenOptionsExt;
opts.mode(0o600);
}
let mut file = opts.open(path).map_err(|e| {
AgentLoopError::config(format!("open session log {}: {e}", path.display()))
})?;
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o600)).map_err(
|e| {
AgentLoopError::config(format!(
"tighten session log permissions on {}: {e}",
path.display()
))
},
)?;
}
match file.try_lock() {
Ok(()) => {}
Err(TryLockError::WouldBlock) => {
return Err(AgentLoopError::config(format!(
"another yolop process is already writing {}; \
refusing to share a session log",
path.display()
)));
}
Err(TryLockError::Error(e)) => {
return Err(AgentLoopError::config(format!(
"lock session log {}: {e}",
path.display()
)));
}
}
let len = file
.seek(SeekFrom::End(0))
.map_err(|e| AgentLoopError::config(format!("stat session log: {e}")))?;
if len > 0 {
let mut last = [0u8; 1];
file.seek(SeekFrom::Start(len - 1))
.and_then(|_| std::io::Read::read_exact(&mut file, &mut last))
.map_err(|e| AgentLoopError::config(format!("read session log tail: {e}")))?;
if last[0] != b'\n' {
tracing::warn!(
"session log {} ends without newline; repairing tail before append",
path.display()
);
writeln!(file)
.map_err(|e| AgentLoopError::config(format!("repair session log tail: {e}")))?;
}
}
let (live, _) = broadcast::channel(EVENT_BROADCAST_CAPACITY);
Ok(Self {
events: Arc::new(RwLock::new(Vec::new())),
sequence: Arc::new(RwLock::new(start_sequence.saturating_sub(1))),
file: Arc::new(Mutex::new(file)),
live,
})
}
pub fn subscribe(&self) -> broadcast::Receiver<Event> {
self.live.subscribe()
}
pub async fn seed_replayed(&self, events: Vec<Event>) {
if events.is_empty() {
return;
}
self.events.write().await.extend(events);
}
}
#[async_trait]
impl EventEmitter for JsonlEventEmitter {
async fn emit(&self, request: EventRequest) -> Result<Event> {
let mut seq = self.sequence.write().await;
*seq = seq.saturating_add(1);
let seq_val = *seq;
drop(seq);
let mut event = request.into_event(EventId::new(), seq_val);
if event.ts.timestamp() == 0 {
event.ts = Utc::now();
}
self.events.write().await.push(event.clone());
if is_replay_relevant(&event.event_type) {
let line = serde_json::to_string(&event).map_err(|e| {
AgentLoopError::config(format!("serialize event for session log: {e}"))
})?;
let mut file = self.file.lock().await;
writeln!(file, "{line}")
.map_err(|e| AgentLoopError::config(format!("write session log line: {e}")))?;
file.flush()
.map_err(|e| AgentLoopError::config(format!("flush session log: {e}")))?;
}
let _ = self.live.send(event.clone());
Ok(event)
}
}
#[async_trait]
impl EventBus for JsonlEventEmitter {
async fn collected_events(&self) -> Vec<Event> {
self.events.read().await.clone()
}
}
fn message_from_event(data: &EventData) -> Option<Message> {
match data {
EventData::InputMessage(d) => Some(d.message.clone()),
EventData::OutputMessageCompleted(OutputMessageCompletedData { message, .. }) => {
Some(message.clone())
}
EventData::ToolCompleted(d) => Some(tool_completed_to_message(d.clone())),
_ => None,
}
}
fn tool_completed_to_message(data: everruns_core::events::ToolCompletedData) -> Message {
let mut images: Vec<ToolResultImage> = Vec::new();
let result = data.result.map(|parts| {
for part in &parts {
if let ContentPart::Image(img) = part
&& let (Some(base64), Some(media_type)) = (&img.base64, &img.media_type)
{
images.push(ToolResultImage {
base64: base64.clone(),
media_type: media_type.clone(),
});
}
}
let text_parts: Vec<&ContentPart> = parts
.iter()
.filter(|part| matches!(part, ContentPart::Text(_)))
.collect();
if text_parts.len() == 1
&& let ContentPart::Text(text) = text_parts[0]
{
return parse_structured_tool_result_text(&text.text);
}
if !text_parts.is_empty() {
serde_json::to_value(&text_parts).unwrap_or_default()
} else {
serde_json::Value::Null
}
});
if images.is_empty() {
Message::tool_result(&data.tool_call_id, result, data.error)
} else {
Message::tool_result_with_images(&data.tool_call_id, result, images)
}
}
fn parse_structured_tool_result_text(text: &str) -> serde_json::Value {
let trimmed = text.trim_start();
if !trimmed.starts_with('{') && !trimmed.starts_with('[') {
return serde_json::Value::String(text.to_string());
}
match serde_json::from_str(text) {
Ok(value @ (serde_json::Value::Object(_) | serde_json::Value::Array(_))) => value,
_ => serde_json::Value::String(text.to_string()),
}
}
#[cfg(test)]
mod tests {
use super::*;
use everruns_core::events::{
EventContext, InputMessageData, OutputMessageCompletedData, ToolCompletedData,
};
use everruns_core::message::Message;
fn input_event(session_id: SessionId, text: &str) -> Event {
Event::new(
session_id,
EventContext::default(),
InputMessageData::new(Message::user(text)),
)
}
#[tokio::test]
async fn seed_replayed_populates_collected_events_without_rewrite() {
let dir = tempfile::tempdir().expect("tempdir");
let session_id = SessionId::from_seed(482);
let session_dir = session_dir_path(dir.path(), session_id);
let path = session_log_path(&session_dir);
let emitter = JsonlEventEmitter::open(&path, 1).expect("open");
let prior = vec![
input_event(session_id, "first prior turn"),
input_event(session_id, "second prior turn"),
];
emitter.seed_replayed(prior.clone()).await;
let collected = emitter.collected_events().await;
assert_eq!(collected.len(), prior.len());
assert_eq!(collected[0].id, prior[0].id);
assert_eq!(collected[1].id, prior[1].id);
let on_disk = std::fs::read_to_string(&path).expect("read");
assert!(
on_disk.is_empty(),
"seed_replayed must not re-persist; found: {on_disk:?}"
);
}
#[tokio::test]
async fn seed_replayed_then_emit_keeps_order() {
let dir = tempfile::tempdir().expect("tempdir");
let session_id = SessionId::from_seed(4820);
let session_dir = session_dir_path(dir.path(), session_id);
let path = session_log_path(&session_dir);
let emitter = JsonlEventEmitter::open(&path, 3).expect("open");
let prior = vec![input_event(session_id, "a"), input_event(session_id, "b")];
emitter.seed_replayed(prior.clone()).await;
let req = EventRequest::new(
session_id,
EventContext::default(),
InputMessageData::new(Message::user("new")),
);
let _new = emitter.emit(req).await.expect("emit");
let collected = emitter.collected_events().await;
assert_eq!(collected.len(), 3, "seeded + 1 new");
assert_eq!(collected[2].sequence, Some(3));
}
#[test]
fn migrate_legacy_session_log_copies_flat_jsonl() {
let dir = tempfile::tempdir().expect("tempdir");
let session_id = SessionId::from_seed(48200);
let legacy_path = legacy_session_log_path(dir.path(), session_id);
std::fs::write(&legacy_path, "legacy event\n").expect("write legacy");
let session_dir = session_dir_path(dir.path(), session_id);
let migrated =
migrate_legacy_session_log(dir.path(), &session_dir, session_id).expect("migrate");
let current_path = session_log_path(&session_dir);
assert_eq!(migrated.as_deref(), Some(legacy_path.as_path()));
assert_eq!(
std::fs::read_to_string(¤t_path).expect("read migrated"),
"legacy event\n"
);
assert!(
legacy_path.exists(),
"migration copies instead of renaming for old yolop compatibility"
);
}
#[test]
fn migrate_legacy_session_log_does_not_overwrite_current_log() {
let dir = tempfile::tempdir().expect("tempdir");
let session_id = SessionId::from_seed(48201);
let legacy_path = legacy_session_log_path(dir.path(), session_id);
std::fs::write(&legacy_path, "legacy event\n").expect("write legacy");
let session_dir = session_dir_path(dir.path(), session_id);
std::fs::create_dir_all(&session_dir).expect("create session dir");
let current_path = session_log_path(&session_dir);
std::fs::write(¤t_path, "current event\n").expect("write current");
let migrated =
migrate_legacy_session_log(dir.path(), &session_dir, session_id).expect("migrate");
assert!(migrated.is_none());
assert_eq!(
std::fs::read_to_string(¤t_path).expect("read current"),
"current event\n"
);
}
#[cfg(unix)]
#[tokio::test]
async fn open_tightens_session_dir_to_owner_only_on_unix() {
use std::os::unix::fs::PermissionsExt;
let dir = tempfile::tempdir().expect("tempdir");
let session_id = SessionId::from_seed(48222);
let session_dir = session_dir_path(dir.path(), session_id);
let path = session_log_path(&session_dir);
let _emitter = JsonlEventEmitter::open(&path, 1).expect("open");
let mode = std::fs::metadata(&session_dir)
.expect("session dir exists")
.permissions()
.mode()
& 0o777;
assert_eq!(
mode, 0o700,
"per-session folder must be owner-only on Unix, got {mode:o}"
);
let file_mode = std::fs::metadata(&path)
.expect("events.jsonl exists")
.permissions()
.mode()
& 0o777;
assert_eq!(
file_mode, 0o600,
"events.jsonl must be owner-only on Unix, got {file_mode:o}"
);
}
#[cfg(unix)]
#[tokio::test]
async fn open_corrects_loose_events_jsonl_permissions_on_resume() {
use std::os::unix::fs::PermissionsExt;
let dir = tempfile::tempdir().expect("tempdir");
let session_id = SessionId::from_seed(48224);
let session_dir = session_dir_path(dir.path(), session_id);
let path = session_log_path(&session_dir);
std::fs::create_dir_all(&session_dir).expect("pre-create");
std::fs::write(&path, "").expect("pre-create file");
std::fs::set_permissions(&path, std::fs::Permissions::from_mode(0o644))
.expect("loosen for test");
let _emitter = JsonlEventEmitter::open(&path, 1).expect("open");
let mode = std::fs::metadata(&path)
.expect("events.jsonl exists")
.permissions()
.mode()
& 0o777;
assert_eq!(
mode, 0o600,
"resume must re-tighten an existing loose events.jsonl, got {mode:o}"
);
}
#[cfg(unix)]
#[tokio::test]
async fn open_corrects_loose_session_dir_permissions_on_resume() {
use std::os::unix::fs::PermissionsExt;
let dir = tempfile::tempdir().expect("tempdir");
let session_id = SessionId::from_seed(48223);
let session_dir = session_dir_path(dir.path(), session_id);
let path = session_log_path(&session_dir);
std::fs::create_dir_all(&session_dir).expect("pre-create");
std::fs::set_permissions(&session_dir, std::fs::Permissions::from_mode(0o755))
.expect("loosen for test");
let _emitter = JsonlEventEmitter::open(&path, 1).expect("open");
let mode = std::fs::metadata(&session_dir)
.expect("session dir exists")
.permissions()
.mode()
& 0o777;
assert_eq!(
mode, 0o700,
"resume must re-tighten an existing loose session folder, got {mode:o}"
);
}
#[tokio::test]
async fn output_message_thinking_is_persisted_for_provider_continuation() {
let dir = tempfile::tempdir().expect("tempdir");
let session_id = SessionId::from_seed(4821);
let session_dir = session_dir_path(dir.path(), session_id);
let path = session_log_path(&session_dir);
let emitter = JsonlEventEmitter::open(&path, 1).expect("open");
let mut message = Message::assistant("I will inspect the files.");
message.thinking = Some("private model reasoning".to_string());
message.thinking_signature = Some("encrypted-thinking-token".to_string());
let req = EventRequest::new(
session_id,
EventContext::default(),
OutputMessageCompletedData::new(message),
);
emitter.emit(req).await.expect("emit");
let on_disk = std::fs::read_to_string(&path).expect("read");
assert!(
on_disk.contains("private model reasoning"),
"session log must persist assistant thinking for restore: {on_disk}"
);
assert!(
on_disk.contains("encrypted-thinking-token"),
"session log must persist thinking_signature for provider continuation: {on_disk}"
);
let replayed = replay(&path, session_id).expect("replay");
let replayed_message = replayed.messages.first().expect("message replayed");
assert_eq!(
replayed_message.thinking.as_deref(),
Some("private model reasoning")
);
assert_eq!(
replayed_message.thinking_signature.as_deref(),
Some("encrypted-thinking-token")
);
}
#[tokio::test]
async fn reason_completed_event_is_persisted_and_replayed() {
use everruns_core::events::ReasonCompletedData;
let dir = tempfile::tempdir().expect("tempdir");
let session_id = SessionId::from_seed(4822);
let session_dir = session_dir_path(dir.path(), session_id);
let path = session_log_path(&session_dir);
let emitter = JsonlEventEmitter::open(&path, 1).expect("open");
let req = EventRequest::new(
session_id,
EventContext::default(),
ReasonCompletedData::success("Will read the lib.rs file", true, 1, Some(120), None),
);
emitter.emit(req).await.expect("emit");
let on_disk = std::fs::read_to_string(&path).expect("read");
assert!(
on_disk.contains("\"reason.completed\""),
"reason.completed should be persisted to JSONL: {on_disk}"
);
assert!(
on_disk.contains("Will read the lib.rs file"),
"text_preview narration should round-trip on disk: {on_disk}"
);
let replayed = replay(&path, session_id).expect("replay");
assert_eq!(replayed.events.len(), 1);
match &replayed.events[0].data {
EventData::ReasonCompleted(data) => {
assert_eq!(
data.text_preview.as_deref(),
Some("Will read the lib.rs file")
);
assert!(data.has_tool_calls);
assert_eq!(data.tool_call_count, 1);
}
other => panic!("expected ReasonCompleted, got {other:?}"),
}
assert!(replayed.messages.is_empty());
}
#[tokio::test]
async fn reason_item_event_is_persisted_and_replayed() {
use everruns_core::events::ReasonItemData;
use everruns_core::typed_id::TurnId;
let dir = tempfile::tempdir().expect("tempdir");
let session_id = SessionId::from_seed(4823);
let session_dir = session_dir_path(dir.path(), session_id);
let path = session_log_path(&session_dir);
let emitter = JsonlEventEmitter::open(&path, 1).expect("open");
let req = EventRequest::new(
session_id,
EventContext::default(),
ReasonItemData {
turn_id: TurnId::new(),
provider: "openai".to_string(),
model: Some("gpt-5".to_string()),
item_id: "rs_abc123".to_string(),
encrypted_content: Some("opaque-encrypted-blob".to_string()),
summary: vec!["Considered file structure.".to_string()],
token_count: Some(42),
},
);
emitter.emit(req).await.expect("emit");
let on_disk = std::fs::read_to_string(&path).expect("read");
assert!(
on_disk.contains("\"reason.item\""),
"reason.item should be persisted: {on_disk}"
);
assert!(
on_disk.contains("opaque-encrypted-blob"),
"encrypted_content must round-trip for provider continuation: {on_disk}"
);
let replayed = replay(&path, session_id).expect("replay");
assert_eq!(replayed.events.len(), 1);
match &replayed.events[0].data {
EventData::ReasonItem(data) => {
assert_eq!(data.provider, "openai");
assert_eq!(data.item_id, "rs_abc123");
assert_eq!(
data.encrypted_content.as_deref(),
Some("opaque-encrypted-blob")
);
assert_eq!(data.summary, vec!["Considered file structure.".to_string()]);
}
other => panic!("expected ReasonItem, got {other:?}"),
}
assert!(replayed.messages.is_empty());
}
#[tokio::test]
async fn subscribe_receives_emitted_events_including_deltas() {
use everruns_core::events::{
OUTPUT_MESSAGE_DELTA, OutputMessageDeltaData, TOOL_OUTPUT_DELTA, ToolOutputDeltaData,
};
use everruns_core::typed_id::TurnId;
let dir = tempfile::tempdir().expect("tempdir");
let session_id = SessionId::from_seed(99);
let session_dir = session_dir_path(dir.path(), session_id);
let path = session_log_path(&session_dir);
let emitter = JsonlEventEmitter::open(&path, 1).expect("open");
let mut rx = emitter.subscribe();
let turn_id = TurnId::new();
let delta_req = EventRequest::new(
session_id,
EventContext::default(),
OutputMessageDeltaData {
turn_id,
delta: "hello".to_string(),
accumulated: "hello".to_string(),
},
);
let _ = emitter.emit(delta_req).await.expect("emit delta");
let tool_req = EventRequest::new(
session_id,
EventContext::default(),
ToolOutputDeltaData {
tool_call_id: "call-1".to_string(),
tool_name: "bash".to_string(),
delta: "running...\n".to_string(),
stream: "stdout".to_string(),
},
);
let _ = emitter.emit(tool_req).await.expect("emit tool delta");
let first = rx.recv().await.expect("first event");
let second = rx.recv().await.expect("second event");
assert_eq!(first.event_type, OUTPUT_MESSAGE_DELTA);
assert_eq!(second.event_type, TOOL_OUTPUT_DELTA);
let on_disk = std::fs::read_to_string(&path).expect("read");
assert!(
on_disk.is_empty(),
"delta events must not be persisted: {on_disk:?}"
);
}
#[tokio::test]
async fn subscribe_does_not_replay_seeded_history() {
let dir = tempfile::tempdir().expect("tempdir");
let session_id = SessionId::from_seed(101);
let session_dir = session_dir_path(dir.path(), session_id);
let path = session_log_path(&session_dir);
let emitter = JsonlEventEmitter::open(&path, 1).expect("open");
emitter
.seed_replayed(vec![input_event(session_id, "old turn")])
.await;
let mut rx = emitter.subscribe();
let drained = tokio::time::timeout(std::time::Duration::from_millis(40), rx.recv()).await;
assert!(
drained.is_err(),
"no event should be delivered, got {drained:?}"
);
}
#[test]
fn tool_completed_replay_preserves_json_result_shape() {
let data = ToolCompletedData::success(
"call_read".to_string(),
"read_file".to_string(),
vec![ContentPart::text(
serde_json::json!({
"path": "/workspace/src/lib.rs",
"content": "1|fn main() {}"
})
.to_string(),
)],
Some(1),
);
let message = tool_completed_to_message(data);
let result = message
.tool_result_content()
.and_then(|content| content.result.as_ref())
.expect("tool result should be present");
assert_eq!(result["path"], "/workspace/src/lib.rs");
assert_eq!(result["content"], "1|fn main() {}");
}
#[test]
fn tool_completed_replay_keeps_scalar_json_as_text() {
let data = ToolCompletedData::success(
"call_scalar".to_string(),
"custom_tool".to_string(),
vec![ContentPart::text("123")],
Some(1),
);
let message = tool_completed_to_message(data);
let result = message
.tool_result_content()
.and_then(|content| content.result.as_ref())
.expect("tool result should be present");
assert_eq!(result, &serde_json::Value::String("123".to_string()));
}
fn serialize_event(event: &Event) -> String {
serde_json::to_string(event).expect("serialize event")
}
#[test]
fn replay_missing_file_returns_empty_session() {
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("nonexistent.jsonl");
let session_id = SessionId::from_seed(70001);
let out = replay(&path, session_id).expect("replay missing file");
assert!(out.events.is_empty());
assert!(out.messages.is_empty());
assert!(out.max_sequence.is_none());
}
#[test]
fn replay_empty_file_returns_empty_session() {
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("events.jsonl");
std::fs::write(&path, b"").expect("write empty file");
let session_id = SessionId::from_seed(70002);
let out = replay(&path, session_id).expect("replay empty file");
assert!(out.events.is_empty());
assert!(out.messages.is_empty());
assert!(out.max_sequence.is_none());
}
#[test]
fn replay_blank_lines_in_middle_are_skipped() {
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("events.jsonl");
let session_id = SessionId::from_seed(70003);
let first = serialize_event(&input_event(session_id, "first"));
let second = serialize_event(&input_event(session_id, "second"));
let content = format!("{first}\n\n \n{second}\n");
std::fs::write(&path, content).expect("write");
let out = replay(&path, session_id).expect("replay");
assert_eq!(
out.events.len(),
2,
"blank/whitespace-only lines must not produce events"
);
}
#[test]
fn replay_malformed_json_line_does_not_stop_replay() {
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("events.jsonl");
let session_id = SessionId::from_seed(70004);
let good_a = serialize_event(&input_event(session_id, "alpha"));
let good_b = serialize_event(&input_event(session_id, "beta"));
let content = format!("{good_a}\n{{not valid json\n{good_b}\n");
std::fs::write(&path, content).expect("write");
let out = replay(&path, session_id).expect("replay");
assert_eq!(out.events.len(), 2, "two good lines must survive");
}
#[test]
fn replay_skips_events_for_different_session_id() {
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("events.jsonl");
let expected = SessionId::from_seed(70005);
let other = SessionId::from_seed(70006);
let mine = serialize_event(&input_event(expected, "mine"));
let theirs = serialize_event(&input_event(other, "theirs"));
let content = format!("{theirs}\n{mine}\n{theirs}\n");
std::fs::write(&path, content).expect("write");
let out = replay(&path, expected).expect("replay");
assert_eq!(
out.events.len(),
1,
"only the line belonging to `expected` should be kept"
);
assert_eq!(out.events[0].session_id, expected);
}
#[test]
fn replay_truncated_last_line_is_skipped_without_dropping_prior() {
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("events.jsonl");
let session_id = SessionId::from_seed(70007);
let good = serialize_event(&input_event(session_id, "kept"));
let partial = "{\"id\":\"evt_\",\"session";
let content = format!("{good}\n{partial}");
std::fs::write(&path, content).expect("write");
let out = replay(&path, session_id).expect("replay");
assert_eq!(
out.events.len(),
1,
"the partial tail must not corrupt the replay"
);
}
#[test]
fn replay_binary_garbage_stops_early_and_drops_suffix() {
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("events.jsonl");
let session_id = SessionId::from_seed(70008);
let prefix = serialize_event(&input_event(session_id, "prefix"));
let suffix = serialize_event(&input_event(session_id, "would-be-suffix"));
let mut bytes: Vec<u8> = Vec::new();
bytes.extend_from_slice(prefix.as_bytes());
bytes.push(b'\n');
bytes.extend_from_slice(&[0xff, 0xfe, 0xfd, b'\n']);
bytes.extend_from_slice(suffix.as_bytes());
bytes.push(b'\n');
std::fs::write(&path, &bytes).expect("write");
let out = replay(&path, session_id).expect("replay must not panic");
assert_eq!(
out.events.len(),
1,
"only the valid prefix line should survive; replay must stop early on read error"
);
let kept_text = match &out.events[0].data {
EventData::InputMessage(d) => d.message.text().unwrap_or_default().to_string(),
other => panic!("expected InputMessage, got {other:?}"),
};
assert!(
kept_text.contains("prefix"),
"expected the prefix event to be the one kept, got: {kept_text}"
);
}
#[test]
fn replay_tracks_highest_sequence_across_all_valid_events() {
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("events.jsonl");
let session_id = SessionId::from_seed(70009);
let mut a = input_event(session_id, "a");
a.sequence = Some(7);
let mut b = input_event(session_id, "b");
b.sequence = Some(3);
let mut c = input_event(session_id, "c");
c.sequence = Some(12);
let content = format!(
"{}\n{}\n{}\n",
serialize_event(&a),
serialize_event(&b),
serialize_event(&c)
);
std::fs::write(&path, content).expect("write");
let out = replay(&path, session_id).expect("replay");
assert_eq!(out.events.len(), 3);
assert_eq!(out.max_sequence, Some(12));
}
#[test]
fn replay_only_blank_lines_returns_empty() {
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("events.jsonl");
std::fs::write(&path, "\n\n \n\t\n").expect("write");
let session_id = SessionId::from_seed(70010);
let out = replay(&path, session_id).expect("replay");
assert!(out.events.is_empty());
assert!(out.max_sequence.is_none());
}
}