use async_trait::async_trait;
use chrono::{DateTime, Utc};
use everruns_core::error::{AgentLoopError, Result};
use everruns_core::events::{
Event, EventData, EventRequest, INPUT_MESSAGE, OUTPUT_MESSAGE_COMPLETED,
OutputMessageCompletedData, REASON_COMPLETED, REASON_ITEM, SESSION_TITLE_UPDATED,
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 serde::{Deserialize, Serialize};
use std::fs::{File, OpenOptions, TryLockError};
use std::io::{BufRead, BufReader, Seek, SeekFrom, Write};
use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::sync::atomic::{AtomicI32, Ordering};
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")
}
fn session_workspace_path(session_dir: &Path) -> PathBuf {
session_dir.join("workspace.json")
}
pub fn legacy_session_log_path(sessions_dir: &Path, session_id: SessionId) -> PathBuf {
sessions_dir.join(format!("{session_id}.jsonl"))
}
#[derive(Debug, Clone, Deserialize, Serialize)]
pub struct WorktreeMetadata {
pub path: PathBuf,
pub branch: String,
pub base_ref: String,
#[serde(default, skip_serializing_if = "String::is_empty")]
pub slug: String,
}
#[derive(Debug, Clone, Deserialize, Serialize)]
pub struct SessionWorkspaceMetadata {
#[serde(alias = "workspace_root")]
pub active_root: PathBuf,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub repo_root: Option<PathBuf>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub canonical_repo_root: Option<PathBuf>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub project_id: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub title: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub summary: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub created_at: Option<DateTime<Utc>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub updated_at: Option<DateTime<Utc>>,
#[serde(default)]
pub session_kind: SessionKind,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub parent_session_id: Option<SessionId>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub worktree: Option<WorktreeMetadata>,
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq, Default)]
#[serde(rename_all = "snake_case")]
pub enum SessionKind {
#[default]
Interactive,
Print,
Nested,
Acp,
Test,
Automated,
}
impl SessionWorkspaceMetadata {
pub fn new(active_root: PathBuf, repo_root: Option<PathBuf>) -> Self {
let now = Utc::now();
let canonical_repo_root = repo_root.as_ref().map(|path| canonicalize_lossy(path));
let project_id = canonical_repo_root
.as_ref()
.map(|path| format!("file:{}", path.display()));
Self {
active_root,
repo_root,
canonical_repo_root,
project_id,
title: None,
summary: None,
created_at: Some(now),
updated_at: Some(now),
session_kind: SessionKind::Interactive,
parent_session_id: None,
worktree: None,
}
}
pub fn apply_initial_prompt(&mut self, prompt: &str) {
let title = prompt_title(prompt);
if title.is_empty() {
return;
}
if self.title.is_none() {
self.title = Some(title.clone());
}
if self.summary.is_none() {
self.summary = Some(title);
}
}
}
fn canonicalize_lossy(path: &Path) -> PathBuf {
std::fs::canonicalize(path).unwrap_or_else(|_| path.to_path_buf())
}
fn prompt_title(prompt: &str) -> String {
const MAX_TITLE_CHARS: usize = 80;
let collapsed = prompt.split_whitespace().collect::<Vec<_>>().join(" ");
let mut chars = collapsed.chars();
let title: String = chars.by_ref().take(MAX_TITLE_CHARS).collect();
if chars.next().is_some() {
format!("{title}…")
} else {
title
}
}
pub fn read_session_workspace_metadata(
session_dir: &Path,
) -> Result<Option<SessionWorkspaceMetadata>> {
let path = session_workspace_path(session_dir);
if !path.exists() {
return Ok(None);
}
let bytes = match std::fs::read(&path) {
Ok(bytes) => bytes,
Err(e) => {
tracing::warn!(
path = %path.display(),
error = %e,
"ignoring unreadable session workspace metadata"
);
return Ok(None);
}
};
let metadata: SessionWorkspaceMetadata = match serde_json::from_slice(&bytes) {
Ok(metadata) => metadata,
Err(e) => {
tracing::warn!(
path = %path.display(),
error = %e,
"ignoring malformed session workspace metadata"
);
return Ok(None);
}
};
Ok(Some(metadata))
}
pub fn read_session_workspace(session_dir: &Path) -> Result<Option<PathBuf>> {
Ok(read_session_workspace_metadata(session_dir)?.map(|m| m.active_root))
}
pub fn write_session_workspace(
session_dir: &Path,
metadata: &SessionWorkspaceMetadata,
) -> Result<()> {
std::fs::create_dir_all(session_dir).map_err(|e| {
AgentLoopError::config(format!("create session dir {}: {e}", session_dir.display()))
})?;
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
std::fs::set_permissions(session_dir, std::fs::Permissions::from_mode(0o700)).map_err(
|e| {
AgentLoopError::config(format!(
"tighten session dir permissions on {}: {e}",
session_dir.display()
))
},
)?;
}
let path = session_workspace_path(session_dir);
let mut metadata = metadata.clone();
metadata.updated_at = Some(Utc::now());
if metadata.created_at.is_none() {
metadata.created_at = metadata.updated_at;
}
let bytes = serde_json::to_vec_pretty(&metadata)
.map_err(|e| AgentLoopError::config(format!("serialize session workspace: {e}")))?;
write_private_file(&path, &bytes)
}
pub fn update_session_workspace_title(session_dir: &Path, title: &str) -> Result<bool> {
let Some(mut metadata) = read_session_workspace_metadata(session_dir)? else {
return Err(AgentLoopError::config(format!(
"session workspace metadata not found: {}",
session_workspace_path(session_dir).display()
)));
};
if metadata.title.as_deref() == Some(title) {
return Ok(false);
}
metadata.title = Some(title.to_string());
write_session_workspace(session_dir, &metadata)?;
Ok(true)
}
fn write_private_file(path: &Path, bytes: &[u8]) -> Result<()> {
let mut opts = OpenOptions::new();
opts.create(true).write(true).truncate(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 private file {}: {e}", path.display()))
})?;
file.write_all(bytes).map_err(|e| {
AgentLoopError::config(format!("write private file {}: {e}", path.display()))
})?;
file.write_all(b"\n").map_err(|e| {
AgentLoopError::config(format!("write private file {}: {e}", path.display()))
})?;
file.flush().map_err(|e| {
AgentLoopError::config(format!("flush private file {}: {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 private file permissions on {}: {e}",
path.display()
))
})?;
}
Ok(())
}
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
| SESSION_TITLE_UPDATED
)
}
pub(crate) fn latest_session_title(events: &[Event]) -> Option<String> {
events.iter().rev().find_map(|event| match &event.data {
EventData::SessionTitleUpdated(data) => Some(data.title.clone()),
_ => None,
})
}
pub struct JsonlEventEmitter {
events: Arc<RwLock<Vec<Event>>>,
sequence: Arc<AtomicI32>,
persisted_sequence: Arc<AtomicI32>,
file: Arc<Mutex<File>>,
session_dir: PathBuf,
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(AtomicI32::new(start_sequence.saturating_sub(1))),
persisted_sequence: Arc::new(AtomicI32::new(start_sequence.saturating_sub(1))),
file: Arc::new(Mutex::new(file)),
session_dir: path
.parent()
.unwrap_or_else(|| Path::new("."))
.to_path_buf(),
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);
}
pub async fn replace_collected_events(&self, events: Vec<Event>) {
*self.events.write().await = events;
}
pub fn last_sequence(&self) -> i32 {
self.persisted_sequence.load(Ordering::Acquire)
}
}
#[async_trait]
impl EventEmitter for JsonlEventEmitter {
async fn emit(&self, request: EventRequest) -> Result<Event> {
let previous = self
.sequence
.fetch_update(Ordering::AcqRel, Ordering::Acquire, |sequence| {
Some(sequence.saturating_add(1))
})
.expect("sequence update closure always succeeds");
let seq_val = previous.saturating_add(1);
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}")))?;
self.persisted_sequence.fetch_max(seq_val, Ordering::AcqRel);
}
if let EventData::SessionTitleUpdated(data) = &event.data
&& let Err(error) = update_session_workspace_title(&self.session_dir, &data.title)
{
tracing::warn!(
session_id = %event.session_id,
error = %error,
"failed to project session title into workspace metadata"
);
}
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,
}
}
pub(crate) fn messages_from_events(events: &[Event]) -> Vec<Message> {
events
.iter()
.filter_map(|event| message_from_event(&event.data))
.collect()
}
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, SessionTitleUpdatedData,
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)),
)
}
fn title_event(session_id: SessionId, previous_title: Option<&str>, title: &str) -> Event {
Event::new(
session_id,
EventContext::default(),
SessionTitleUpdatedData {
previous_title: previous_title.map(str::to_string),
title: title.to_string(),
},
)
}
#[tokio::test]
async fn title_event_is_persisted_and_projected() {
let dir = tempfile::tempdir().expect("tempdir");
let session_id = SessionId::from_seed(483);
let session_dir = session_dir_path(dir.path(), session_id);
write_session_workspace(
&session_dir,
&SessionWorkspaceMetadata::new(dir.path().to_path_buf(), None),
)
.expect("seed workspace metadata");
let emitter = JsonlEventEmitter::open(&session_log_path(&session_dir), 1).expect("open");
emitter
.emit(EventRequest::new(
session_id,
EventContext::default(),
SessionTitleUpdatedData {
previous_title: None,
title: "Automatic session titles".to_string(),
},
))
.await
.expect("emit title update");
let replayed = replay(&session_log_path(&session_dir), session_id).expect("replay");
assert_eq!(
latest_session_title(&replayed.events).as_deref(),
Some("Automatic session titles")
);
let metadata = read_session_workspace_metadata(&session_dir)
.expect("read metadata")
.expect("metadata present");
assert_eq!(metadata.title.as_deref(), Some("Automatic session titles"));
}
#[test]
fn latest_title_uses_last_semantic_update() {
let session_id = SessionId::from_seed(484);
let events = vec![
title_event(session_id, None, "First theme"),
input_event(session_id, "switch topics"),
title_event(session_id, Some("First theme"), "Second theme"),
];
assert_eq!(
latest_session_title(&events).as_deref(),
Some("Second theme")
);
}
#[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"
);
}
#[test]
fn read_session_workspace_ignores_malformed_metadata() {
let dir = tempfile::tempdir().expect("tempdir");
let session_id = SessionId::from_seed(48202);
let session_dir = session_dir_path(dir.path(), session_id);
std::fs::create_dir_all(&session_dir).expect("create session dir");
std::fs::write(session_workspace_path(&session_dir), "{not json")
.expect("write malformed metadata");
let workspace = read_session_workspace(&session_dir).expect("read workspace metadata");
assert!(workspace.is_none());
}
#[test]
fn read_legacy_workspace_metadata_defaults_new_fields() {
let json = r#"{
"active_root": "/tmp/yolop-workspace",
"repo_root": "/tmp/yolop-repo"
}"#;
let metadata: SessionWorkspaceMetadata =
serde_json::from_str(json).expect("legacy workspace metadata deserializes");
assert_eq!(metadata.active_root, PathBuf::from("/tmp/yolop-workspace"));
assert_eq!(metadata.repo_root, Some(PathBuf::from("/tmp/yolop-repo")));
assert_eq!(metadata.canonical_repo_root, None);
assert_eq!(metadata.project_id, None);
assert_eq!(metadata.title, None);
assert_eq!(metadata.summary, None);
assert_eq!(metadata.created_at, None);
assert_eq!(metadata.updated_at, None);
assert_eq!(metadata.session_kind, SessionKind::Interactive);
assert_eq!(metadata.parent_session_id, None);
assert!(metadata.worktree.is_none());
}
#[test]
fn write_workspace_metadata_persists_agf_listing_fields() {
let dir = tempfile::tempdir().expect("tempdir");
let session_id = SessionId::from_seed(48203);
let session_dir = session_dir_path(dir.path(), session_id);
let canonical_dir = dir.path().canonicalize().expect("canonical tempdir");
let mut metadata = SessionWorkspaceMetadata::new(
dir.path().join("worktree"),
Some(dir.path().to_path_buf()),
);
metadata.title = Some("Implement workspace metadata".to_string());
metadata.summary = Some("Add list-friendly Yolop session metadata.".to_string());
metadata.session_kind = SessionKind::Nested;
metadata.parent_session_id = Some(SessionId::from_seed(48204));
metadata.worktree = Some(WorktreeMetadata {
path: dir.path().join("worktree"),
branch: "feature/agf-metadata".to_string(),
base_ref: "origin/main".to_string(),
slug: "agf-metadata".to_string(),
});
write_session_workspace(&session_dir, &metadata).expect("write workspace metadata");
let written = read_session_workspace_metadata(&session_dir)
.expect("read workspace metadata")
.expect("metadata present");
assert_eq!(
written.title.as_deref(),
Some("Implement workspace metadata")
);
assert_eq!(
written.summary.as_deref(),
Some("Add list-friendly Yolop session metadata.")
);
assert_eq!(written.session_kind, SessionKind::Nested);
assert_eq!(written.parent_session_id, Some(SessionId::from_seed(48204)));
assert_eq!(written.canonical_repo_root, Some(canonical_dir.clone()));
assert_eq!(
written.project_id.as_deref(),
Some(format!("file:{}", canonical_dir.display()).as_str())
);
assert!(written.created_at.is_some(), "created_at is populated");
assert!(written.updated_at.is_some(), "updated_at is populated");
assert!(
written.updated_at >= written.created_at,
"updated_at should not precede created_at"
);
assert_eq!(
written.worktree.map(|worktree| worktree.branch),
Some("feature/agf-metadata".to_string())
);
}
#[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::{MessageId, 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,
message_id: MessageId::new(),
delta: "hello".to_string(),
accumulated: "hello".to_string(),
phase: None,
},
);
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": "/repo/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"], "/repo/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());
}
}