use crate::study_io::StudyIo;
use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use somatize_core::error::{Result, SomaError};
use somatize_core::event::Event;
use somatize_core::graph::Graph;
use somatize_core::study::{Study, TrialState};
use somatize_core::tracking::{EventEnvelope, RunManifest, RunState, RunStatus};
use somatize_core::viz::{GraphOverlay, NodeStatus};
use std::collections::BTreeMap;
use std::fs;
use std::io::{BufRead, BufReader};
use std::path::{Path, PathBuf};
use super::local_tracker::{load_manifest, load_status};
pub const STALE_HEARTBEAT_SECS: i64 = 300;
pub struct RunReader {
dir: PathBuf,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct RunInfo {
pub run_id: String,
pub kind: String,
pub name: String,
pub state: String,
pub created_at: DateTime<Utc>,
pub finished_at: Option<DateTime<Utc>>,
pub duration_ms: Option<u64>,
pub tags: Vec<String>,
pub dir: String,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub struct NodeSpan {
pub node_id: String,
pub started_ts: Option<DateTime<Utc>>,
pub finished_ts: Option<DateTime<Utc>>,
pub duration_ms: Option<u64>,
pub outcome: String,
pub cache_tier: Option<String>,
pub error: Option<String>,
#[serde(default)]
pub effectful: bool,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct CacheActivity {
pub hits: u64,
pub misses: u64,
pub by_node: BTreeMap<String, NodeCacheCounts>,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct NodeCacheCounts {
pub hits: u64,
pub misses: u64,
pub last_tier: Option<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct MetricPoint {
pub ts: DateTime<Utc>,
pub name: String,
pub value: f64,
pub step: u64,
#[serde(default)]
pub trial_id: Option<String>,
#[serde(default)]
pub node_id: Option<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct HealthFlagRecord {
pub ts: DateTime<Utc>,
pub node_id: String,
pub step: usize,
pub flag: String,
pub detail: String,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct AgenticActivity {
pub turns: u64,
pub input_tokens: u64,
pub output_tokens: u64,
pub effects: u64,
pub replayed: u64,
pub tool_calls: u64,
pub steps_completed: u64,
pub steps_failed: u64,
pub suspensions: u64,
pub by_node: BTreeMap<String, AgentNodeActivity>,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct AgentNodeActivity {
pub turns: u64,
pub input_tokens: u64,
pub output_tokens: u64,
pub duration_ms: u64,
pub effects: u64,
pub effects_by_label: BTreeMap<String, u64>,
pub effect_errors: u64,
pub replayed: u64,
pub tool_calls: u64,
pub tool_errors: u64,
pub handoffs_out: u64,
pub suspensions: u64,
pub spawned: u64,
pub completions: u64,
pub failures: u64,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub struct EffectSpan {
pub node_id: String,
pub turn: usize,
pub effect: String,
pub started_ts: Option<DateTime<Utc>>,
pub finished_ts: Option<DateTime<Utc>>,
pub duration_ms: Option<u64>,
pub replayed: bool,
pub is_error: bool,
pub outcome: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct TrialSpan {
pub trial_id: String,
pub state: String,
pub started_at: Option<DateTime<Utc>>,
pub finished_at: Option<DateTime<Utc>>,
pub duration_ms: Option<u64>,
}
impl RunReader {
pub fn open(run_dir: impl AsRef<Path>) -> Result<Self> {
let dir = run_dir.as_ref().to_path_buf();
load_manifest(&dir)?;
Ok(Self { dir })
}
pub fn dir(&self) -> &Path {
&self.dir
}
pub fn manifest(&self) -> Result<RunManifest> {
load_manifest(&self.dir)
}
pub fn status(&self) -> Result<RunStatus> {
load_status(&self.dir)
}
pub fn info(&self) -> Result<RunInfo> {
let manifest = self.manifest()?;
Ok(run_info(
&self.dir,
manifest,
self.status().ok(),
Utc::now(),
))
}
pub fn events(&self) -> Result<Vec<EventEnvelope>> {
let path = self.dir.join("events.jsonl");
let file = match fs::File::open(&path) {
Ok(f) => f,
Err(_) => return Ok(Vec::new()), };
let mut envelopes = Vec::new();
for line in BufReader::new(file).lines() {
let line = line.map_err(SomaError::Io)?;
if line.trim().is_empty() {
continue;
}
if let Ok(env) = serde_json::from_str::<EventEnvelope>(&line) {
envelopes.push(env);
}
}
Ok(envelopes)
}
pub fn node_timings(&self) -> Result<Vec<NodeSpan>> {
let mut spans: Vec<NodeSpan> = Vec::new();
let mut open: BTreeMap<String, usize> = BTreeMap::new();
for env in self.events()? {
match env.event {
Event::NodeStarted {
node_id, effectful, ..
} => {
open.insert(node_id.clone(), spans.len());
spans.push(NodeSpan {
node_id,
started_ts: Some(env.ts),
finished_ts: None,
duration_ms: None,
outcome: "running".into(),
cache_tier: None,
error: None,
effectful,
});
}
Event::NodeCacheHit {
node_id,
tier,
load_time,
..
} => {
spans.push(NodeSpan {
node_id,
started_ts: Some(env.ts),
finished_ts: Some(env.ts),
duration_ms: Some(load_time.as_millis() as u64),
outcome: "cache_hit".into(),
cache_tier: Some(format!("{tier:?}").to_lowercase()),
error: None,
effectful: false,
});
}
Event::NodeCompleted {
node_id, duration, ..
} => {
let idx = open.remove(&node_id);
let span = match idx {
Some(i) => &mut spans[i],
None => {
spans.push(NodeSpan {
node_id: node_id.clone(),
started_ts: None,
finished_ts: None,
duration_ms: None,
outcome: String::new(),
cache_tier: None,
error: None,
effectful: false,
});
spans.last_mut().expect("just pushed")
}
};
span.finished_ts = Some(env.ts);
span.duration_ms = Some(duration.as_millis() as u64);
span.outcome = "completed".into();
}
Event::NodeFailed { node_id, error, .. } => {
let idx = open.remove(&node_id);
let span = match idx {
Some(i) => &mut spans[i],
None => {
spans.push(NodeSpan {
node_id: node_id.clone(),
started_ts: None,
finished_ts: None,
duration_ms: None,
outcome: String::new(),
cache_tier: None,
error: None,
effectful: false,
});
spans.last_mut().expect("just pushed")
}
};
span.finished_ts = Some(env.ts);
span.outcome = "failed".into();
span.error = Some(error);
}
_ => {}
}
}
Ok(spans)
}
pub fn cache_activity(&self) -> Result<CacheActivity> {
let mut activity = CacheActivity::default();
for env in self.events()? {
match env.event {
Event::NodeCacheHit { node_id, tier, .. } => {
activity.hits += 1;
let counts = activity.by_node.entry(node_id).or_default();
counts.hits += 1;
counts.last_tier = Some(format!("{tier:?}").to_lowercase());
}
Event::NodeCacheMiss { node_id, .. } => {
activity.misses += 1;
activity.by_node.entry(node_id).or_default().misses += 1;
}
_ => {}
}
}
Ok(activity)
}
pub fn metric_series(&self, name: Option<&str>) -> Result<Vec<MetricPoint>> {
let path = self.dir.join("metrics.jsonl");
let mut points: Vec<MetricPoint> = Vec::new();
if let Ok(file) = fs::File::open(&path) {
for line in BufReader::new(file).lines() {
let line = line.map_err(SomaError::Io)?;
if line.trim().is_empty() {
continue;
}
if let Ok(p) = serde_json::from_str::<MetricPoint>(&line) {
points.push(p);
}
}
} else {
for env in self.events()? {
match env.event {
Event::TrialMetric {
trial_id, metric, ..
} => points.push(MetricPoint {
ts: metric.timestamp,
name: metric.name,
value: metric.value,
step: metric.step as u64,
trial_id: Some(trial_id),
node_id: None,
}),
Event::MetricReported {
metric,
node_id,
trial_id,
..
} => points.push(MetricPoint {
ts: metric.timestamp,
name: metric.name,
value: metric.value,
step: metric.step as u64,
trial_id,
node_id,
}),
_ => {}
}
}
}
if let Some(name) = name {
points.retain(|p| p.name == name);
}
Ok(points)
}
pub fn health_flags(&self) -> Result<Vec<HealthFlagRecord>> {
let mut flags = Vec::new();
for env in self.events()? {
if let Event::HealthFlag {
node_id,
step,
flag,
detail,
..
} = env.event
{
flags.push(HealthFlagRecord {
ts: env.ts,
node_id,
step,
flag,
detail,
});
}
}
Ok(flags)
}
pub fn agentic_activity(&self) -> Result<AgenticActivity> {
#[derive(Default)]
struct Pending {
turns: u64,
input_tokens: u64,
output_tokens: u64,
duration_ms: u64,
}
let mut by_node: BTreeMap<String, AgentNodeActivity> = BTreeMap::new();
let mut pending: BTreeMap<String, Pending> = BTreeMap::new();
let mut observed_turns: BTreeMap<String, u64> = BTreeMap::new();
for env in self.events()? {
match env.event {
Event::AgentTurnStarted { node_id, turn, .. } => {
let seen = observed_turns.entry(node_id).or_default();
*seen = (*seen).max(turn as u64 + 1);
}
Event::EffectCompleted {
node_id,
effect,
replayed,
is_error,
..
} => {
let node = by_node.entry(node_id).or_default();
node.effects += 1;
*node.effects_by_label.entry(effect).or_default() += 1;
if replayed {
node.replayed += 1;
}
if is_error {
node.effect_errors += 1;
}
}
Event::ToolCalled {
node_id, is_error, ..
} => {
let node = by_node.entry(node_id).or_default();
node.tool_calls += 1;
if is_error {
node.tool_errors += 1;
}
}
Event::Handoff { from, .. } => {
by_node.entry(from).or_default().handoffs_out += 1;
}
Event::Suspended {
node_id,
turns,
duration,
input_tokens,
output_tokens,
..
} => {
by_node.entry(node_id.clone()).or_default().suspensions += 1;
pending.insert(
node_id,
Pending {
turns: turns as u64,
input_tokens,
output_tokens,
duration_ms: duration.as_millis() as u64,
},
);
}
Event::AgentSpawned {
node_id, children, ..
} => {
by_node.entry(node_id).or_default().spawned += children.len() as u64;
}
Event::AgentStepCompleted {
node_id,
turns,
duration,
input_tokens,
output_tokens,
failed,
..
} => {
let node = by_node.entry(node_id.clone()).or_default();
node.turns += turns as u64;
node.input_tokens += input_tokens;
node.output_tokens += output_tokens;
node.duration_ms += duration.as_millis() as u64;
if failed {
node.failures += 1;
} else {
node.completions += 1;
}
pending.remove(&node_id);
}
_ => {}
}
}
for (node_id, p) in pending {
let node = by_node.entry(node_id).or_default();
node.turns += p.turns;
node.input_tokens += p.input_tokens;
node.output_tokens += p.output_tokens;
node.duration_ms += p.duration_ms;
}
for (node_id, seen) in observed_turns {
let node = by_node.entry(node_id).or_default();
if node.turns == 0 {
node.turns = seen;
}
}
let mut totals = AgenticActivity::default();
for node in by_node.values() {
totals.turns += node.turns;
totals.input_tokens += node.input_tokens;
totals.output_tokens += node.output_tokens;
totals.effects += node.effects;
totals.replayed += node.replayed;
totals.tool_calls += node.tool_calls;
totals.steps_completed += node.completions;
totals.steps_failed += node.failures;
totals.suspensions += node.suspensions;
}
totals.by_node = by_node;
Ok(totals)
}
pub fn agentic_timeline(&self) -> Result<Vec<EffectSpan>> {
let mut spans: Vec<EffectSpan> = Vec::new();
let mut open: BTreeMap<(String, usize, String), Vec<usize>> = BTreeMap::new();
for env in self.events()? {
match env.event {
Event::EffectRequested {
node_id,
turn,
effect,
..
} => {
open.entry((node_id.clone(), turn, effect.clone()))
.or_default()
.push(spans.len());
spans.push(EffectSpan {
node_id,
turn,
effect,
started_ts: Some(env.ts),
finished_ts: None,
duration_ms: None,
replayed: false,
is_error: false,
outcome: "running".into(),
});
}
Event::EffectCompleted {
node_id,
turn,
effect,
duration,
replayed,
is_error,
..
} => {
let key = (node_id.clone(), turn, effect.clone());
let idx = open
.get_mut(&key)
.filter(|v| !v.is_empty())
.map(|v| v.remove(0));
let span = match idx {
Some(i) => &mut spans[i],
None => {
spans.push(EffectSpan {
node_id,
turn,
effect,
started_ts: None,
finished_ts: None,
duration_ms: None,
replayed: false,
is_error: false,
outcome: String::new(),
});
spans.last_mut().expect("just pushed")
}
};
span.finished_ts = Some(env.ts);
span.duration_ms = Some(duration.as_millis() as u64);
span.replayed = replayed;
span.is_error = is_error;
span.outcome = "completed".into();
}
_ => {}
}
}
Ok(spans)
}
pub fn graph(&self) -> Result<Option<Graph>> {
let path = self.dir.join("graph.json");
if !path.exists() {
return Ok(None);
}
let bytes = fs::read(&path)?;
serde_json::from_slice(&bytes)
.map(Some)
.map_err(|e| SomaError::Serialization(e.to_string()))
}
pub fn overlay(&self) -> Result<GraphOverlay> {
let mut overlay = GraphOverlay::default();
let mut counts: BTreeMap<String, u64> = BTreeMap::new();
for span in self.node_timings()? {
let entry = overlay.nodes.entry(span.node_id.clone()).or_default();
*counts.entry(span.node_id).or_default() += 1;
entry.status = Some(match span.outcome.as_str() {
"completed" => NodeStatus::Completed,
"cache_hit" => NodeStatus::Cached,
"failed" => NodeStatus::Failed,
_ => NodeStatus::Running,
});
entry.cache_tier = span.cache_tier;
if let Some(ms) = span.duration_ms {
entry.duration_ms = Some(entry.duration_ms.unwrap_or(0) + ms);
}
}
for (node_id, n) in counts {
if n > 1
&& let Some(entry) = overlay.nodes.get_mut(&node_id)
{
entry.sublabel = Some(format!("×{n}"));
}
}
for flag in self.health_flags()? {
let entry = overlay.nodes.entry(flag.node_id).or_default();
if !entry.flags.contains(&flag.flag) {
entry.flags.push(flag.flag);
}
}
Ok(overlay)
}
pub fn to_mermaid(&self) -> Result<String> {
let graph = self.graph()?.ok_or_else(|| {
SomaError::Other(format!("run dir {} has no graph.json", self.dir.display()))
})?;
Ok(graph.to_mermaid_with(&self.overlay()?))
}
pub fn to_graphviz(&self) -> Result<String> {
let graph = self.graph()?.ok_or_else(|| {
SomaError::Other(format!("run dir {} has no graph.json", self.dir.display()))
})?;
Ok(graph.to_graphviz_with(&self.overlay()?))
}
pub fn to_svg(&self) -> Result<String> {
let graph = self.graph()?.ok_or_else(|| {
SomaError::Other(format!("run dir {} has no graph.json", self.dir.display()))
})?;
Ok(graph.to_svg_with(&self.overlay()?))
}
pub fn study(&self) -> Result<Option<Study>> {
let path = self.dir.join("study.json");
if !path.exists() {
return Ok(None);
}
Study::load(&path).map(Some)
}
pub fn trial_timeline(&self) -> Result<Vec<TrialSpan>> {
let Some(study) = self.study()? else {
return Ok(Vec::new());
};
Ok(study
.trials
.iter()
.map(|t| TrialSpan {
trial_id: t.id.clone(),
state: trial_state_str(&t.state).to_string(),
started_at: t.started_at,
finished_at: t.finished_at,
duration_ms: t.duration_ms,
})
.collect())
}
}
fn trial_state_str(state: &TrialState) -> &'static str {
match state {
TrialState::Pending => "pending",
TrialState::Running => "running",
TrialState::Completed => "completed",
TrialState::Pruned { .. } => "pruned",
TrialState::Failed { .. } => "failed",
}
}
fn run_info(
dir: &Path,
manifest: RunManifest,
status: Option<RunStatus>,
now: DateTime<Utc>,
) -> RunInfo {
let state = match &status {
None => "running".to_string(),
Some(s) => match s.state {
RunState::Completed => "completed".to_string(),
RunState::Failed => "failed".to_string(),
RunState::Running => {
let last_beat = s.heartbeat_at.unwrap_or(s.updated_at);
if (now - last_beat).num_seconds() > STALE_HEARTBEAT_SECS {
"crashed".to_string()
} else {
"running".to_string()
}
}
_ => "running".to_string(),
},
};
let finished_at = status.as_ref().and_then(|s| s.finished_at);
let duration_ms = finished_at
.map(|end| (end - manifest.created_at).num_milliseconds())
.filter(|ms| *ms >= 0)
.map(|ms| ms as u64);
let kind = serde_json::to_value(manifest.kind)
.ok()
.and_then(|v| v.as_str().map(str::to_string))
.unwrap_or_else(|| "other".to_string());
RunInfo {
run_id: manifest.run_id,
kind,
name: manifest.name,
state,
created_at: manifest.created_at,
finished_at,
duration_ms,
tags: manifest.tags,
dir: dir.display().to_string(),
}
}
pub fn list_runs(root: impl AsRef<Path>) -> Result<Vec<RunInfo>> {
let runs_dir = root.as_ref().join("runs");
let entries = match fs::read_dir(&runs_dir) {
Ok(e) => e,
Err(_) => return Ok(Vec::new()), };
let now = Utc::now();
let mut infos: Vec<RunInfo> = entries
.flatten()
.filter(|e| e.path().is_dir())
.filter_map(|e| {
let dir = e.path();
let manifest = load_manifest(&dir).ok()?;
let status = load_status(&dir).ok();
Some(run_info(&dir, manifest, status, now))
})
.collect();
infos.sort_by_key(|info| std::cmp::Reverse(info.created_at));
Ok(infos)
}
#[cfg(test)]
mod tests {
use super::*;
use chrono::Duration as ChronoDuration;
use somatize_core::tracking::RunKind;
fn manifest(run_id: &str) -> RunManifest {
RunManifest::new(run_id, RunKind::Train, "test-run")
}
#[test]
fn run_info_detects_crash_from_stale_heartbeat() {
let now = Utc::now();
let stale = RunStatus {
state: RunState::Running,
updated_at: now - ChronoDuration::seconds(STALE_HEARTBEAT_SECS + 60),
heartbeat_at: Some(now - ChronoDuration::seconds(STALE_HEARTBEAT_SECS + 60)),
finished_at: None,
};
let info = run_info(Path::new("/tmp/r"), manifest("r1"), Some(stale), now);
assert_eq!(info.state, "crashed");
let fresh = RunStatus::running();
let info = run_info(Path::new("/tmp/r"), manifest("r1"), Some(fresh), now);
assert_eq!(info.state, "running");
}
#[test]
fn run_info_duration_and_kind() {
let now = Utc::now();
let mut m = manifest("r2");
m.created_at = now - ChronoDuration::milliseconds(1500);
let status = RunStatus {
state: RunState::Completed,
updated_at: now,
heartbeat_at: Some(now),
finished_at: Some(now),
};
let info = run_info(Path::new("/tmp/r"), m, Some(status), now);
assert_eq!(info.state, "completed");
assert_eq!(info.kind, "train");
assert_eq!(info.duration_ms, Some(1500));
}
}