use std::collections::HashMap;
use std::path::{Path, PathBuf};
use std::time::Duration;
use tokio::sync::mpsc;
use crate::transcript::{self, SubagentMeta};
use super::bytes::{ReadResult, TailState, read_appended};
use super::{Flow, Source, TailRequest, UiEvent, Update};
const POLL_INTERVAL: Duration = Duration::from_millis(200);
const SWITCH_SCAN_EVERY: u32 = 10;
const SWITCH_IDLE_TICKS: u32 = 150;
pub(crate) struct LiveSession {
pub(crate) project_dir: Option<PathBuf>,
main_path: PathBuf,
subagents_dir: PathBuf,
main_state: TailState,
tracked: HashMap<PathBuf, (Source, TailState)>,
seen_meta: std::collections::HashSet<PathBuf>,
ticks: u32,
idle_ticks: u32,
seed_offsets: HashMap<PathBuf, u64>,
}
#[derive(Debug, Default)]
pub(crate) struct SnapshotSeed {
pub(crate) offsets: HashMap<PathBuf, u64>,
pub(crate) seen_meta: std::collections::HashSet<PathBuf>,
}
impl LiveSession {
pub(crate) fn new(project_dir: PathBuf, main_path: PathBuf) -> Self {
let subagents_dir = transcript::subagents_dir(&main_path).unwrap_or_default();
Self {
project_dir: Some(project_dir),
main_path,
subagents_dir,
main_state: TailState::default(),
tracked: HashMap::new(),
seen_meta: std::collections::HashSet::new(),
ticks: 0,
idle_ticks: 0,
seed_offsets: HashMap::new(),
}
}
fn session_id(&self) -> String {
transcript::session_id_from_path(&self.main_path)
}
pub(crate) fn seed(&mut self, seed: SnapshotSeed) {
if let Some(off) = seed.offsets.get(&self.main_path) {
self.main_state.offset = *off;
}
self.seed_offsets = seed.offsets;
self.seen_meta = seed.seen_meta;
}
fn track(&mut self, path: PathBuf, source: Source) {
if self.tracked.contains_key(&path) {
return;
}
let mut state = TailState::default();
if let Some(off) = self.seed_offsets.remove(&path) {
state.offset = off;
}
self.tracked.insert(path, (source, state));
}
}
pub(crate) fn resolve_live_target(target: &Path) -> Option<(PathBuf, PathBuf)> {
if target.is_dir() {
let main = transcript::latest_session_file(target)?;
Some((target.to_path_buf(), main))
} else {
let project_dir = target.parent().map(Path::to_path_buf).unwrap_or_default();
Some((project_dir, target.to_path_buf()))
}
}
async fn await_first_session(
project_dir: &Path,
req_rx: &mut mpsc::Receiver<TailRequest>,
) -> Result<PathBuf, Flow> {
let mut ticker = tokio::time::interval(POLL_INTERVAL);
ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
loop {
if let Some(main) = transcript::latest_session_file(project_dir) {
return Ok(main);
}
tokio::select! {
req = req_rx.recv() => match req {
Some(TailRequest::Watch(p)) => return Err(Flow::Switch(p)),
None => return Err(Flow::Exit),
},
_ = ticker.tick() => {}
}
}
}
pub(crate) async fn run_live(
target: &Path,
ui_tx: &mpsc::Sender<UiEvent>,
req_rx: &mut mpsc::Receiver<TailRequest>,
) -> Flow {
let pin = !target.is_dir();
let (project_dir, main_path) = match resolve_live_target(target) {
Some(pair) => pair,
None => {
match await_first_session(target, req_rx).await {
Ok(main) => (target.to_path_buf(), main),
Err(flow) => return flow,
}
}
};
let mut session = LiveSession::new(project_dir, main_path);
if pin {
session.project_dir = None;
}
let session_id = session.session_id();
let _ = ui_tx
.send(UiEvent::SessionReset {
session_id: session_id.clone(),
})
.await;
tail_loop(session, session_id, ui_tx, req_rx).await
}
pub(crate) async fn tail_loop(
mut session: LiveSession,
session_id: String,
ui_tx: &mpsc::Sender<UiEvent>,
req_rx: &mut mpsc::Receiver<TailRequest>,
) -> Flow {
let mut ticker = tokio::time::interval(POLL_INTERVAL);
ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
loop {
tokio::select! {
req = req_rx.recv() => {
match req {
Some(TailRequest::Watch(new_path)) => return Flow::Switch(new_path),
None => return Flow::Exit,
}
}
_ = ticker.tick() => {
if let Some(switch) = poll_live(&mut session, &session_id, ui_tx).await {
return Flow::Switch(switch);
}
}
}
}
}
async fn poll_live(
session: &mut LiveSession,
session_id: &str,
ui_tx: &mpsc::Sender<UiEvent>,
) -> Option<PathBuf> {
let mut updates: Vec<Update> = Vec::new();
match read_appended(&session.main_path, &mut session.main_state) {
ReadResult::Reset => {
let _ = ui_tx
.send(UiEvent::SessionReset {
session_id: session_id.to_string(),
})
.await;
return Some(session.main_path.clone());
}
ReadResult::Entries(entries) => {
for entry in entries {
updates.push(Update::Entry {
source: Source::Main,
entry,
});
}
}
ReadResult::NoChange | ReadResult::Missing => {}
}
scan_files(session, &mut updates);
if read_tracked(&mut session.tracked, &mut updates) {
let _ = ui_tx
.send(UiEvent::SessionReset {
session_id: session_id.to_string(),
})
.await;
return Some(session.main_path.clone());
}
let had_activity = !updates.is_empty();
if had_activity {
let _ = ui_tx
.send(UiEvent::Batch {
session_id: session_id.to_string(),
updates,
})
.await;
}
session.idle_ticks = if had_activity {
0
} else {
session.idle_ticks.saturating_add(1)
};
session.ticks = session.ticks.wrapping_add(1);
if session.ticks.is_multiple_of(SWITCH_SCAN_EVERY)
&& session.idle_ticks >= SWITCH_IDLE_TICKS
&& let Some(dir) = &session.project_dir
&& let Some(latest) = transcript::latest_session_file(dir)
&& latest != session.main_path
{
let new_id = transcript::session_id_from_path(&latest);
let _ = ui_tx
.send(UiEvent::SessionReset { session_id: new_id })
.await;
return Some(latest);
}
None
}
fn read_tracked(
tracked: &mut HashMap<PathBuf, (Source, TailState)>,
updates: &mut Vec<Update>,
) -> bool {
let mut reset = false;
for (path, (source, state)) in tracked.iter_mut() {
match read_appended(path, state) {
ReadResult::Entries(entries) => {
for entry in entries {
updates.push(Update::Entry {
source: source.clone(),
entry,
});
}
}
ReadResult::Reset => reset = true,
ReadResult::NoChange | ReadResult::Missing => {}
}
}
reset
}
fn scan_files(session: &mut LiveSession, updates: &mut Vec<Update>) {
let subagents = session.subagents_dir.clone();
for f in transcript::scan_subagent_files(&subagents, None) {
register_subagent(session, f, updates);
}
for wf_id in transcript::scan_workflow_ids(&subagents) {
let journal = transcript::workflow_journal(&subagents, &wf_id);
if journal.is_file() {
session.track(journal, Source::Journal(wf_id.clone()));
}
let wf_dir = transcript::workflow_dir(&subagents, &wf_id);
for f in transcript::scan_subagent_files(&wf_dir, Some(&wf_id)) {
register_subagent(session, f, updates);
}
}
}
fn register_subagent(
session: &mut LiveSession,
f: transcript::SubagentFile,
updates: &mut Vec<Update>,
) {
if f.meta.is_file() {
emit_meta(
&f.meta,
f.agent_id.clone(),
f.workflow,
&mut session.seen_meta,
updates,
);
}
session.track(f.transcript, Source::Sub(f.agent_id));
}
fn emit_meta(
path: &Path,
agent_id: String,
workflow: Option<String>,
seen: &mut std::collections::HashSet<PathBuf>,
updates: &mut Vec<Update>,
) {
if seen.contains(path) {
return;
}
let Ok(bytes) = std::fs::read(path) else {
return;
};
let Ok(meta) = serde_json::from_slice::<SubagentMeta>(&bytes) else {
return;
};
seen.insert(path.to_path_buf());
updates.push(Update::SubagentMeta {
agent_id,
workflow,
meta,
});
}
#[cfg(test)]
mod tests {
use super::*;
use crate::tailer::UiEvent;
#[test]
fn emit_meta_retries_after_failed_parse() {
use std::io::Write;
let mut tmp = std::env::temp_dir();
tmp.push(format!(
"zoetrope_meta_retry_{}.meta.json",
std::process::id()
));
let mut seen = std::collections::HashSet::new();
let mut updates = Vec::new();
std::fs::File::create(&tmp)
.unwrap()
.write_all(b"{\"agentTy")
.unwrap();
emit_meta(&tmp, "a1".into(), None, &mut seen, &mut updates);
assert!(updates.is_empty());
assert!(!seen.contains(&tmp), "failed parse must be retried");
std::fs::File::create(&tmp)
.unwrap()
.write_all(br#"{"agentType":"guide"}"#)
.unwrap();
emit_meta(&tmp, "a1".into(), None, &mut seen, &mut updates);
assert_eq!(updates.len(), 1);
assert!(seen.contains(&tmp));
let _ = std::fs::remove_file(&tmp);
}
#[tokio::test]
async fn auto_switch_follows_a_newer_session_when_idle() {
let mut dir = std::env::temp_dir();
dir.push(format!("zoetrope_switch_test_{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(&dir).unwrap();
let a = dir.join("11111111-1111-1111-1111-111111111111.jsonl");
let b = dir.join("99999999-9999-9999-9999-999999999999.jsonl");
std::fs::write(&a, "").unwrap();
std::fs::write(&b, "").unwrap();
let mut session = LiveSession::new(dir.clone(), a.clone());
session.idle_ticks = SWITCH_IDLE_TICKS;
session.ticks = SWITCH_SCAN_EVERY - 1;
let (tx, mut rx) = mpsc::channel(32);
let switched = poll_live(&mut session, "11111111", &tx).await;
assert_eq!(
switched,
Some(b.clone()),
"an idle session follows the newer file"
);
match rx.try_recv() {
Ok(UiEvent::SessionReset { session_id }) => assert_ne!(
session_id, "11111111-1111-1111-1111-111111111111",
"reset carries the new session id, not the old"
),
other => panic!("expected a SessionReset for the new session, got {other:?}"),
}
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn replay_seed_resumes_where_snapshot_stopped() {
use std::io::Write;
let mut dir = std::env::temp_dir();
dir.push(format!("zoetrope_seed_test_{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(&dir).unwrap();
let main = dir.join("44444444-4444-4444-4444-444444444444.jsonl");
std::fs::write(
&main,
concat!(
r#"{"type":"user","uuid":"u1","parentUuid":null,"timestamp":"2026-06-05T10:00:00.000Z","message":{"role":"user","content":"one"}}"#, "\n",
r#"{"type":"user","uuid":"u2","parentUuid":"u1","timestamp":"2026-06-05T10:01:00.000Z","message":{"role":"user","content":"two"}}"#, "\n",
),
)
.unwrap();
let (items, _info, seed) = crate::tailer::replay::build_replay(&main);
assert_eq!(items.len(), 2);
std::fs::OpenOptions::new()
.append(true)
.open(&main)
.unwrap()
.write_all(
concat!(
r#"{"type":"user","uuid":"u3","parentUuid":"u2","timestamp":"2026-06-05T10:02:00.000Z","message":{"role":"user","content":"three"}}"#, "\n",
)
.as_bytes(),
)
.unwrap();
let mut session = LiveSession::new(dir.clone(), main.clone());
session.project_dir = None;
session.seed(seed);
let (tx, mut rx) = mpsc::channel(32);
poll_live(&mut session, "44444444", &tx).await;
let mut emitted = 0;
while let Ok(ev) = rx.try_recv() {
if let UiEvent::Batch { updates, .. } = ev {
emitted += updates.len();
}
}
assert_eq!(emitted, 1, "only the post-snapshot append is emitted");
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn tracked_file_truncation_reattaches() {
use std::io::Write;
let mut dir = std::env::temp_dir();
dir.push(format!("zoetrope_subtrunc_test_{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
let session_uuid = "55555555-5555-5555-5555-555555555555";
let sub_dir = dir.join(session_uuid).join("subagents");
std::fs::create_dir_all(&sub_dir).unwrap();
let main = dir.join(format!("{session_uuid}.jsonl"));
std::fs::write(&main, concat!(
r#"{"type":"user","uuid":"u1","parentUuid":null,"timestamp":"2026-06-05T10:00:00.000Z","message":{"role":"user","content":"start"}}"#, "\n",
)).unwrap();
let sub = sub_dir.join("agent-bbbbbbbbbbbbbbbbb.jsonl");
std::fs::write(&sub, concat!(
r#"{"type":"user","uuid":"s1","parentUuid":null,"isSidechain":true,"agentId":"bbbbbbbbbbbbbbbbb","timestamp":"2026-06-05T10:01:00.000Z","message":{"role":"user","content":"task"}}"#, "\n",
r#"{"type":"user","uuid":"s2","parentUuid":"s1","isSidechain":true,"agentId":"bbbbbbbbbbbbbbbbb","timestamp":"2026-06-05T10:02:00.000Z","message":{"role":"user","content":"more"}}"#, "\n",
)).unwrap();
let (tx, mut rx) = mpsc::channel(32);
let mut session = LiveSession::new(dir.clone(), main.clone());
assert!(poll_live(&mut session, "55555555", &tx).await.is_none());
std::fs::File::create(&sub)
.unwrap()
.write_all(b"{}\n")
.unwrap();
let switch = poll_live(&mut session, "55555555", &tx).await;
assert_eq!(
switch.as_deref(),
Some(main.as_path()),
"tracked-file truncation must re-attach the session"
);
let mut saw_reset = false;
while let Ok(ev) = rx.try_recv() {
if matches!(&ev, UiEvent::SessionReset { session_id } if session_id == "55555555") {
saw_reset = true;
}
}
assert!(saw_reset);
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn truncation_triggers_full_reattach() {
use std::io::Write;
let mut dir = std::env::temp_dir();
dir.push(format!("zoetrope_trunc_test_{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(&dir).unwrap();
let main = dir.join("33333333-3333-3333-3333-333333333333.jsonl");
std::fs::write(&main, b"{\"type\":\"user\",\"uuid\":\"u1\",\"parentUuid\":null,\"message\":{\"role\":\"user\",\"content\":\"abcdef\"}}\n").unwrap();
let (tx, mut rx) = mpsc::channel(32);
let mut session = LiveSession::new(dir.clone(), main.clone());
poll_live(&mut session, "33333333", &tx).await;
std::fs::File::create(&main)
.unwrap()
.write_all(b"{}\n")
.unwrap();
let switch = poll_live(&mut session, "33333333", &tx).await;
assert_eq!(
switch.as_deref(),
Some(main.as_path()),
"truncation must re-attach, not patch in place"
);
let mut saw_reset = false;
while let Ok(ev) = rx.try_recv() {
if matches!(&ev, UiEvent::SessionReset { session_id } if session_id == "33333333") {
saw_reset = true;
}
}
assert!(saw_reset);
let _ = std::fs::remove_dir_all(&dir);
}
}