use std::path::Path;
use bevy_ecs::entity::Entity;
use leviath_core::run_archive;
use leviath_core::run_meta::{ContextSnapshot, RunMeta, RunStatus};
use leviath_runtime::host::SpawnArgs;
use leviath_runtime::interaction_points::InteractionPointState;
use leviath_runtime::persistence::{RunMetadata, TokenTotals};
use leviath_runtime::restore::restore_agent;
use leviath_runtime::world::PipelineWorld;
use crate::daemon::spawn::{SpawnDeps, build_agent_for_reload};
pub fn reload_persisted_agents(
world: &mut PipelineWorld,
deps: SpawnDeps<'_>,
runs_dir: &Path,
) -> Vec<(String, leviath_runtime::world::AgentId)> {
let mut reloaded: Vec<(RunMeta, Entity)> = Vec::new();
let Ok(dir_entries) = std::fs::read_dir(runs_dir) else {
return Vec::new(); };
let candidates: Vec<(RunMeta, bool)> = dir_entries
.flatten()
.filter_map(|dir_entry| {
let run_dir = dir_entry.path();
let meta = read_meta(&run_dir)?; let parked_on_fanout = run_dir.join("fanout.json").exists();
Some((meta, parked_on_fanout))
})
.collect();
let ordered = leviath_runtime::restore::triage_restores(candidates);
for meta in ordered {
let run_dir = runs_dir.join(&meta.run_id);
match reload_one(world, deps.clone(), &meta, &run_dir) {
Ok(entity) => reloaded.push((meta, entity)),
Err(e) => {
tracing::warn!(run_id = %meta.run_id, error = %e, "skipping un-reloadable run");
mark_crashed(&run_dir, meta, &e.to_string(), deps.now_secs);
}
}
}
relink_tree(world, &reloaded);
restore_fan_outs(world, &reloaded, runs_dir);
reloaded
.into_iter()
.map(|(meta, entity)| (meta.run_id, world.own_agent(entity)))
.collect()
}
pub fn reload_run(
world: &mut PipelineWorld,
deps: SpawnDeps<'_>,
run_id: &str,
runs_dir: &std::path::Path,
) -> Option<leviath_runtime::world::AgentId> {
let run_dir = runs_dir.join(run_id);
let meta = read_meta(&run_dir)?;
if is_terminal(&meta.status) {
return None; }
let entity = reload_one(world, deps, &meta, &run_dir).ok()?;
Some(world.own_agent(entity))
}
fn restore_fan_outs(world: &mut PipelineWorld, reloaded: &[(RunMeta, Entity)], runs_dir: &Path) {
let by_run_id: std::collections::HashMap<&str, Entity> = reloaded
.iter()
.map(|(m, e)| (m.run_id.as_str(), *e))
.collect();
for (meta, entity) in reloaded {
let path = runs_dir.join(&meta.run_id).join("fanout.json");
let Some(state) = std::fs::read_to_string(&path)
.ok()
.and_then(|s| serde_json::from_str::<leviath_runtime::fanout::FanOutState>(&s).ok())
else {
continue;
};
leviath_runtime::fanout::restore_fan_out_waiting(
world.world_mut(),
*entity,
state,
&|rid| by_run_id.get(rid).copied(),
);
}
}
fn relink_tree(world: &mut PipelineWorld, reloaded: &[(RunMeta, Entity)]) {
use leviath_runtime::components::{AgentState, ParentRef, SubAgentChildren};
let by_run_id: std::collections::HashMap<&str, Entity> = reloaded
.iter()
.map(|(m, e)| (m.run_id.as_str(), *e))
.collect();
let w = world.world_mut();
for (meta, entity) in reloaded {
if let Some(parent_id) = &meta.parent_run_id {
match by_run_id.get(parent_id.as_str()) {
Some(&parent_entity) => {
w.entity_mut(*entity).insert(ParentRef {
parent_entity,
parent_agent_id: parent_id.clone(),
depth: meta.depth,
});
}
None => tracing::warn!(
run_id = %meta.run_id, parent = %parent_id,
"parent run did not reload; leaving child unlinked"
),
}
}
if !meta.children.is_empty() {
let children: Vec<Entity> = meta
.children
.iter()
.filter_map(|cid| by_run_id.get(cid.as_str()).copied())
.collect();
if !children.is_empty() {
w.entity_mut(*entity).insert(SubAgentChildren {
children,
max_child_depth: meta.max_child_depth,
});
}
w.get_mut::<AgentState>(*entity)
.expect("a reloaded agent always has AgentState")
.spawned_children_ids = meta.children.clone();
}
}
}
fn read_meta(run_dir: &Path) -> Option<RunMeta> {
let text = std::fs::read_to_string(run_dir.join("meta.json")).ok()?;
serde_json::from_str(&text).ok()
}
fn totals_from(meta: &RunMeta) -> TokenTotals {
TokenTotals {
prompt_tokens: meta.prompt_tokens,
completion_tokens: meta.completion_tokens,
cached_tokens: meta.cached_tokens,
cache_write_tokens: meta.cache_write_tokens,
tool_calls: meta.tool_calls,
}
}
fn mark_crashed(run_dir: &Path, meta: RunMeta, reason: &str, now_secs: i64) {
let crashed = RunMeta {
status: RunStatus::Error,
error: Some(format!(
"the daemon exited while this run was active and it could not be recovered: {reason}"
)),
updated_at: now_secs,
..meta
};
if let Err(e) = crate::runstate::write_meta_to(run_dir, &crashed) {
tracing::warn!(
run_id = %crashed.run_id,
error = %e,
"could not record an un-reloadable run as crashed"
);
}
}
fn is_terminal(status: &RunStatus) -> bool {
matches!(
status,
RunStatus::Complete | RunStatus::Cancelled | RunStatus::Error
)
}
fn reload_one(
world: &mut PipelineWorld,
deps: SpawnDeps<'_>,
meta: &RunMeta,
run_dir: &Path,
) -> Result<Entity, String> {
let args = SpawnArgs {
run_id: meta.run_id.clone(),
blueprint_path: meta.agent_path.clone(),
task: meta.task.clone(),
regions: Default::default(),
model: meta.model.clone(),
workdir: meta.workdir.clone(),
metadata: meta.metadata.clone(),
callback_url: meta.callback_url.clone(),
callback_secret: meta.callback_secret.clone(),
yolo: meta.yolo,
no_seed_commands: true,
allow: Vec::new(),
max_depth: None,
parent_run_id: meta.parent_run_id.clone(),
output: meta.output_request.clone(),
};
let entity = build_agent_for_reload(world.world_mut(), deps, &args)?;
let folded = std::fs::read(run_dir.join("run.lvr"))
.ok()
.and_then(|bytes| run_archive::read_archive_lenient(&mut bytes.as_slice()).ok())
.and_then(|(_version, records)| run_archive::fold(&records));
let (snapshot, stage_index, iteration, totals, pending_batch) = match folded {
Some(folded) => {
let totals = totals_from(&folded.meta);
(
folded.context,
folded.meta.stage_index,
folded.meta.iteration,
totals,
folded.pending_batch,
)
}
None => {
let snapshot = std::fs::read_to_string(run_dir.join("context.json"))
.ok()
.and_then(|s| serde_json::from_str::<ContextSnapshot>(&s).ok())
.unwrap_or_else(|| ContextSnapshot {
stage_name: meta.current_stage.clone(),
total_tokens: 0,
max_tokens: 0,
regions: Vec::new(),
});
(
snapshot,
meta.stage_index,
meta.iteration,
totals_from(meta),
None,
)
}
};
restore_agent(
world.world_mut(),
entity,
&snapshot,
stage_index,
iteration,
totals,
);
leviath_runtime::restore::restore_stage_ledger(
world.world_mut(),
entity,
&crate::runstate::read_stages_index_from(run_dir),
);
if let Some(batch) = pending_batch {
leviath_runtime::restore::restore_pending_batch(
world.world_mut(),
entity,
&batch,
&meta.children,
);
}
{
let mut md = world
.world_mut()
.get_mut::<RunMetadata>(entity)
.expect("build_agent attached run metadata");
md.started_at = meta.started_at;
md.title = meta.title.clone();
md.callback_url = meta.callback_url.clone();
md.callback_secret = meta.callback_secret.clone();
}
{
let mut flags = world
.world_mut()
.get_mut::<leviath_runtime::persistence::RunOutcomeFlags>(entity)
.expect("build_agent attached run outcome flags");
flags.0 = meta.flags.clone();
}
if let Some(output) = read_final_output_from(run_dir, meta) {
world
.world_mut()
.entity_mut(entity)
.insert(leviath_runtime::persistence::FinalOutput(output));
}
if let Some(state) = std::fs::read_to_string(run_dir.join("interactions.json"))
.ok()
.and_then(|s| serde_json::from_str::<InteractionPointState>(&s).ok())
{
let agent = world.own_agent(entity);
leviath_runtime::interaction_points::restore_interaction_point(
world.world_mut(),
agent,
state,
);
}
if meta.status == RunStatus::Paused {
world.pause(world.own_agent(entity));
}
Ok(entity)
}
fn read_final_output_from(dir: &Path, meta: &RunMeta) -> Option<leviath_core::FinalOutput> {
let descriptor = meta.final_output.clone()?;
let content = std::fs::read_to_string(dir.join(leviath_core::FINAL_OUTPUT_FILE)).ok()?;
Some(leviath_core::FinalOutput {
content,
format: descriptor.format,
stage: descriptor.stage,
submitted_at: descriptor.submitted_at,
truncated: descriptor.truncated,
artifacts: descriptor.artifacts,
})
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::Arc;
use leviath_mcp::ToolExecutor;
use leviath_runtime::host::SubAgentOp;
use leviath_runtime::interaction_hub::InteractionHub;
use tokio::sync::Mutex;
use tokio::sync::mpsc::UnboundedSender;
use crate::config::Config;
use crate::daemon::tool_service::CliToolService;
use leviath_runtime::ProviderRegistry;
use leviath_runtime::components::AgentStatus;
use leviath_runtime::inference_pool::InferencePoolConfig;
use tokio::runtime::Handle;
fn sub_tx() -> UnboundedSender<SubAgentOp> {
tokio::sync::mpsc::unbounded_channel().0
}
struct FakeProvider;
#[async_trait::async_trait]
impl leviath_providers::Provider for FakeProvider {
async fn infer(
&self,
_r: &leviath_providers::InferenceRequest,
) -> leviath_providers::Result<leviath_providers::InferenceResponse> {
Err(leviath_providers::ProviderError::Other("t".to_string()))
}
async fn count_tokens(&self, _t: &str, _m: &str) -> usize {
1
}
fn max_context_tokens(&self, _m: &str) -> usize {
1000
}
fn name(&self) -> &str {
"fake"
}
fn capabilities(&self, _m: &str) -> leviath_providers::ModelCapabilities {
leviath_providers::ModelCapabilities::default()
}
}
fn test_world() -> (PipelineWorld, Arc<CliToolService>) {
let cli = Arc::new(CliToolService::new());
let mut registry = ProviderRegistry::new();
for p in ["anthropic", "openai", "ollama"] {
registry.register(p.to_string(), Arc::new(FakeProvider));
}
let world = PipelineWorld::new(
registry,
cli.clone(),
InferencePoolConfig::new(),
1,
None,
Handle::current(),
);
(world, cli)
}
fn coder_manifest() -> String {
crate::test_support::inline_coder_manifest()
}
fn write_run(
runs_dir: &Path,
run_id: &str,
agent_path: &str,
status: RunStatus,
context: Option<&ContextSnapshot>,
) {
write_run_tree(RunFixture {
runs_dir,
run_id,
agent_path,
status,
context,
parent_run_id: None,
children: &[],
depth: 0,
max_child_depth: 0,
});
}
struct RunFixture<'a> {
runs_dir: &'a Path,
run_id: &'a str,
agent_path: &'a str,
status: RunStatus,
context: Option<&'a ContextSnapshot>,
parent_run_id: Option<&'a str>,
children: &'a [&'a str],
depth: usize,
max_child_depth: usize,
}
fn write_run_tree(f: RunFixture<'_>) {
let RunFixture {
runs_dir,
run_id,
agent_path,
status,
context,
parent_run_id,
children,
depth,
max_child_depth,
} = f;
let dir = runs_dir.join(run_id);
std::fs::create_dir_all(&dir).unwrap();
let meta = RunMeta {
run_id: run_id.to_string(),
agent_name: "coder".to_string(),
agent_path: agent_path.to_string(),
task: "resume me".to_string(),
model: None,
pid: 0,
status,
current_stage: "implement".to_string(),
stage_index: 0,
num_stages: 1,
iteration: 5,
prompt_tokens: 42,
completion_tokens: 7,
cached_tokens: 0,
cache_write_tokens: 0,
tool_calls: 3,
workdir: std::env::temp_dir().to_string_lossy().to_string(),
started_at: 111,
updated_at: 222,
last_progress_at: None,
error: None,
title: Some("Resume Me".to_string()),
metadata: std::collections::HashMap::new(),
callback_url: Some("http://cb".to_string()),
callback_secret: None,
parent_run_id: parent_run_id.map(str::to_string),
children: children.iter().map(|s| s.to_string()).collect(),
depth,
max_child_depth,
flags: leviath_core::run_meta::RunFlags {
modified_files: vec!["src/a.rs".to_string()],
modified_file_count: 1,
no_output_tools: true,
..Default::default()
},
yolo: false,
read_paths: None,
final_output: Some(
leviath_core::output::FinalOutput::new(
"already answered",
Some("markdown".to_string()),
"implement".to_string(),
777,
)
.descriptor(),
),
output_request: Some(leviath_core::output::OutputSpec {
format: Some("a2ui".to_string()),
..Default::default()
}),
};
std::fs::write(dir.join("meta.json"), serde_json::to_string(&meta).unwrap()).unwrap();
std::fs::write(
dir.join(leviath_core::FINAL_OUTPUT_FILE),
"already answered",
)
.unwrap();
if let Some(ctx) = context {
std::fs::write(
dir.join("context.json"),
serde_json::to_string(ctx).unwrap(),
)
.unwrap();
}
}
fn agent_dir() -> tempfile::TempDir {
let dir = tempfile::tempdir().unwrap();
std::fs::write(dir.path().join("agent.leviath"), coder_manifest()).unwrap();
dir
}
fn write_run_archive(
runs_dir: &Path,
run_id: &str,
agent_path: &str,
stage_index: usize,
iteration: usize,
prompt_tokens: usize,
context: &ContextSnapshot,
) {
use leviath_core::run_archive::{self, RunIdentity, RunRecord};
let dir = runs_dir.join(run_id);
std::fs::create_dir_all(&dir).unwrap();
let mut meta = RunMeta::new(
run_id.to_string(),
"coder".to_string(),
agent_path.to_string(),
"resume me".to_string(),
None,
std::env::temp_dir().to_string_lossy().to_string(),
1,
);
meta.status = RunStatus::Running;
meta.current_stage = "implement".to_string();
meta.stage_index = stage_index;
meta.iteration = iteration;
meta.prompt_tokens = prompt_tokens;
let mut buf = Vec::new();
run_archive::write_archive_start(&mut buf, run_archive::RUN_ARCHIVE_VERSION).unwrap();
run_archive::write_record(
&mut buf,
&RunRecord::Header {
identity: RunIdentity {
run_id: run_id.to_string(),
machine_id: "m".to_string(),
world_id: "w".to_string(),
created_at: 1,
},
meta: Box::new(meta),
},
)
.unwrap();
run_archive::write_record(
&mut buf,
&RunRecord::ContextCheckpoint {
snapshot: context.clone(),
at: 2,
},
)
.unwrap();
std::fs::write(dir.join("run.lvr"), &buf).unwrap();
}
#[tokio::test]
async fn reload_keeps_a_paused_run_paused() {
let agent = agent_dir();
let manifest = agent.path().join("agent.leviath");
let runs = tempfile::tempdir().unwrap();
write_run(
runs.path(),
"run-paused",
manifest.to_str().unwrap(),
RunStatus::Paused,
None,
);
let (mut world, cli) = test_world();
let hub = InteractionHub::new();
let mcp = Arc::new(Mutex::new(ToolExecutor::new()));
let restored = reload_persisted_agents(
&mut world,
crate::daemon::spawn::SpawnDeps {
tool_service: cli.as_ref(),
config: &Config::default(),
shared_mcp: mcp,
mcp_tool_defs: &[],
hub: &hub,
now_secs: 999,
subagent_tx: sub_tx().clone(),
},
runs.path(),
);
assert_eq!(restored.len(), 1);
let (run_id, entity) = &restored[0];
assert_eq!(run_id, "run-paused");
assert_eq!(world.agent_status(*entity), Some(AgentStatus::Paused));
}
#[test]
fn a_descriptor_without_its_sidecar_restores_nothing() {
let dir = tempfile::tempdir().unwrap();
let mut meta = RunMeta::new(
"run-1".to_string(),
"a".to_string(),
"/p".to_string(),
"t".to_string(),
None,
"/w".to_string(),
1,
);
assert!(read_final_output_from(dir.path(), &meta).is_none());
let answer = leviath_core::output::FinalOutput::new(
"already answered",
Some("markdown".to_string()),
"implement".to_string(),
777,
);
meta.final_output = Some(answer.descriptor());
assert!(read_final_output_from(dir.path(), &meta).is_none());
std::fs::write(
dir.path().join(leviath_core::FINAL_OUTPUT_FILE),
&answer.content,
)
.unwrap();
let restored = read_final_output_from(dir.path(), &meta).expect("both halves");
assert_eq!(restored.content, "already answered");
assert_eq!(restored.stage, "implement");
}
async fn reload_single(runs: &Path, run_id: &str) -> (PipelineWorld, Entity) {
let (mut world, cli) = test_world();
let restored = reload_persisted_agents(
&mut world,
crate::daemon::spawn::SpawnDeps {
tool_service: cli.as_ref(),
config: &Config::default(),
shared_mcp: Arc::new(Mutex::new(ToolExecutor::new())),
mcp_tool_defs: &[],
hub: &InteractionHub::new(),
now_secs: 999,
subagent_tx: sub_tx().clone(),
},
runs,
);
assert_eq!(restored.len(), 1);
assert_eq!(restored[0].0, run_id);
let entity = restored[0].1;
(world, entity.entity())
}
#[tokio::test]
async fn reload_keeps_an_unattended_run_unattended() {
let agent = agent_dir();
let manifest = agent.path().join("agent.leviath");
let runs = tempfile::tempdir().unwrap();
write_run(
runs.path(),
"run-yolo",
manifest.to_str().unwrap(),
RunStatus::Running,
None,
);
let meta_path = runs.path().join("run-yolo").join("meta.json");
let mut meta: RunMeta =
serde_json::from_str(&std::fs::read_to_string(&meta_path).unwrap()).unwrap();
meta.yolo = true;
std::fs::write(&meta_path, serde_json::to_string(&meta).unwrap()).unwrap();
let (world, entity) = reload_single(runs.path(), "run-yolo").await;
assert!(
world
.world()
.get::<RunMetadata>(entity)
.expect("reloaded run has metadata")
.unattended
);
assert!(
world
.world()
.get::<leviath_runtime::components::InteractionAutoApprove>(entity)
.is_some(),
"an unattended reload still auto-approves its checkpoints"
);
}
#[tokio::test]
async fn reload_does_not_invent_unattended() {
let agent = agent_dir();
let manifest = agent.path().join("agent.leviath");
let runs = tempfile::tempdir().unwrap();
write_run(
runs.path(),
"run-plain",
manifest.to_str().unwrap(),
RunStatus::Running,
None,
);
let meta_path = runs.path().join("run-plain").join("meta.json");
let mut raw: serde_json::Value =
serde_json::from_str(&std::fs::read_to_string(&meta_path).unwrap()).unwrap();
raw.as_object_mut().unwrap().remove("yolo");
std::fs::write(&meta_path, serde_json::to_string(&raw).unwrap()).unwrap();
let (world, entity) = reload_single(runs.path(), "run-plain").await;
assert!(
!world
.world()
.get::<RunMetadata>(entity)
.expect("reloaded run has metadata")
.unattended
);
}
#[tokio::test]
async fn reloads_nonterminal_runs_and_restores_state() {
let agent = agent_dir();
let manifest = agent.path().join("agent.leviath");
let runs = tempfile::tempdir().unwrap();
let ctx = ContextSnapshot {
stage_name: "implement".to_string(),
total_tokens: 4,
max_tokens: 100_000,
regions: vec![leviath_core::run_meta::RegionSnapshot {
name: "conversation".to_string(),
kind: "clearable".to_string(),
current_tokens: 4,
max_tokens: 100_000,
entries: vec![leviath_core::run_meta::RegionEntrySnapshot {
content: "earlier turn".to_string(),
tokens: 4,
kind: leviath_core::region::EntryKind::UserMessage,
metadata: None,
key: None,
taint: Default::default(),
}],
}],
};
write_run(
runs.path(),
"run-live",
manifest.to_str().unwrap(),
RunStatus::Running,
Some(&ctx),
);
write_run(
runs.path(),
"run-done",
manifest.to_str().unwrap(),
RunStatus::Complete,
None,
);
let (mut world, cli) = test_world();
let hub = InteractionHub::new();
let mcp = Arc::new(Mutex::new(ToolExecutor::new()));
let restored = reload_persisted_agents(
&mut world,
crate::daemon::spawn::SpawnDeps {
tool_service: cli.as_ref(),
config: &Config::default(),
shared_mcp: mcp,
mcp_tool_defs: &[],
hub: &hub,
now_secs: 999,
subagent_tx: sub_tx().clone(),
},
runs.path(),
);
assert_eq!(restored.len(), 1);
let (run_id, entity) = &restored[0];
assert_eq!(run_id, "run-live");
assert_eq!(world.agent_status(*entity), Some(AgentStatus::Active));
let md = world.world().get::<RunMetadata>(entity.entity()).unwrap();
assert_eq!(md.started_at, 111);
assert_eq!(md.title.as_deref(), Some("Resume Me"));
assert_eq!(md.callback_url.as_deref(), Some("http://cb"));
let totals = world.world().get::<TokenTotals>(entity.entity()).unwrap();
assert_eq!(totals.prompt_tokens, 42);
assert_eq!(totals.tool_calls, 3);
let flags = world
.world()
.get::<leviath_runtime::persistence::RunOutcomeFlags>(entity.entity())
.unwrap();
assert_eq!(flags.0.modified_files, vec!["src/a.rs".to_string()]);
assert_eq!(flags.0.modified_file_count, 1);
assert!(flags.0.no_output_tools);
let output = world
.world()
.get::<leviath_runtime::persistence::FinalOutput>(entity.entity())
.expect("a submitted answer survives the restart");
assert_eq!(output.0.content, "already answered");
assert_eq!(output.0.stage, "implement");
assert_eq!(
md.output_request.as_ref().and_then(|s| s.format.as_deref()),
Some("a2ui")
);
}
fn assert_restored_from_archive(world: &PipelineWorld, entity: Entity) {
use leviath_runtime::components::AgentState;
let state = world.world().get::<AgentState>(entity).unwrap();
assert_eq!(state.current_stage, "fresh-stage");
assert_eq!(state.iteration, 9);
let totals = world.world().get::<TokenTotals>(entity).unwrap();
assert_eq!(totals.prompt_tokens, 99);
}
#[tokio::test]
async fn reload_prefers_the_atomic_journal_over_a_stale_context_json() {
let agent = agent_dir();
let manifest = agent.path().join("agent.leviath");
let mpath = manifest.to_str().unwrap();
let runs = tempfile::tempdir().unwrap();
let stale = ContextSnapshot {
stage_name: "stale-stage".to_string(),
total_tokens: 1,
max_tokens: 100,
regions: vec![],
};
write_run(
runs.path(),
"run-torn",
mpath,
RunStatus::Running,
Some(&stale),
);
let fresh = ContextSnapshot {
stage_name: "fresh-stage".to_string(),
total_tokens: 4,
max_tokens: 100_000,
regions: vec![],
};
write_run_archive(runs.path(), "run-torn", mpath, 0, 9, 99, &fresh);
let (mut world, cli) = test_world();
let hub = InteractionHub::new();
let mcp = Arc::new(Mutex::new(ToolExecutor::new()));
let restored = reload_persisted_agents(
&mut world,
crate::daemon::spawn::SpawnDeps {
tool_service: cli.as_ref(),
config: &Config::default(),
shared_mcp: mcp,
mcp_tool_defs: &[],
hub: &hub,
now_secs: 999,
subagent_tx: sub_tx().clone(),
},
runs.path(),
);
assert_eq!(restored.len(), 1);
assert_restored_from_archive(&world, restored[0].1.entity());
}
#[tokio::test]
async fn reload_tolerates_a_torn_journal_tail() {
let agent = agent_dir();
let manifest = agent.path().join("agent.leviath");
let mpath = manifest.to_str().unwrap();
let runs = tempfile::tempdir().unwrap();
let stale = ContextSnapshot {
stage_name: "stale-stage".to_string(),
total_tokens: 1,
max_tokens: 100,
regions: vec![],
};
write_run(
runs.path(),
"run-torn2",
mpath,
RunStatus::Running,
Some(&stale),
);
let fresh = ContextSnapshot {
stage_name: "fresh-stage".to_string(),
total_tokens: 4,
max_tokens: 100_000,
regions: vec![],
};
write_run_archive(runs.path(), "run-torn2", mpath, 0, 9, 99, &fresh);
{
use std::io::Write;
let mut f = std::fs::OpenOptions::new()
.append(true)
.open(runs.path().join("run-torn2/run.lvr"))
.unwrap();
f.write_all(&[0, 0, 0, 0, 0, 0, 0, 10, 1, 2]).unwrap();
}
let (mut world, cli) = test_world();
let hub = InteractionHub::new();
let mcp = Arc::new(Mutex::new(ToolExecutor::new()));
let restored = reload_persisted_agents(
&mut world,
crate::daemon::spawn::SpawnDeps {
tool_service: cli.as_ref(),
config: &Config::default(),
shared_mcp: mcp,
mcp_tool_defs: &[],
hub: &hub,
now_secs: 999,
subagent_tx: sub_tx().clone(),
},
runs.path(),
);
assert_eq!(restored.len(), 1);
assert_restored_from_archive(&world, restored[0].1.entity());
}
fn append_archive_records(
runs_dir: &Path,
run_id: &str,
records: &[leviath_core::run_archive::RunRecord],
) {
use std::io::Write;
let mut buf = Vec::new();
for r in records {
leviath_core::run_archive::write_record(&mut buf, r).unwrap();
}
let mut f = std::fs::OpenOptions::new()
.append(true)
.open(runs_dir.join(run_id).join("run.lvr"))
.unwrap();
f.write_all(&buf).unwrap();
}
fn batch_call(
id: &str,
name: &str,
result: Option<&str>,
) -> leviath_core::run_archive::ToolCallRecord {
leviath_core::run_archive::ToolCallRecord {
id: id.to_string(),
name: name.to_string(),
arguments: "{}".to_string(),
result: result.map(str::to_string),
thought_signature: None,
}
}
fn conversation_of(world: &PipelineWorld, entity: Entity) -> Vec<leviath_core::RegionEntry> {
world
.world()
.get::<leviath_runtime::components::ContextWindow>(entity)
.unwrap()
.get_region("conversation")
.unwrap()
.content
.clone()
}
#[tokio::test]
async fn reload_replays_a_pending_tool_batch_instead_of_reexecuting() {
use leviath_core::run_archive::RunRecord;
let agent = agent_dir();
let manifest = agent.path().join("agent.leviath");
let mpath = manifest.to_str().unwrap();
let runs = tempfile::tempdir().unwrap();
write_run(runs.path(), "run-batch", mpath, RunStatus::Running, None);
let ctx = ContextSnapshot {
stage_name: "implement".to_string(),
total_tokens: 0,
max_tokens: 100_000,
regions: vec![],
};
write_run_archive(runs.path(), "run-batch", mpath, 0, 9, 99, &ctx);
append_archive_records(
runs.path(),
"run-batch",
&[
RunRecord::ToolBatch {
calls: vec![
batch_call("c_done", "write_file", None),
batch_call("c_lost", "shell", None),
],
at: 3,
stage_index: 0,
iteration: 9,
response: "writing then running".to_string(),
},
RunRecord::ToolCallDone {
iteration: 9,
call_id: "c_done".to_string(),
result: "Wrote 42 bytes to x.txt".to_string(),
at: 4,
},
],
);
let (mut world, cli) = test_world();
let hub = InteractionHub::new();
let mcp = Arc::new(Mutex::new(ToolExecutor::new()));
let restored = reload_persisted_agents(
&mut world,
crate::daemon::spawn::SpawnDeps {
tool_service: cli.as_ref(),
config: &Config::default(),
shared_mcp: mcp,
mcp_tool_defs: &[],
hub: &hub,
now_secs: 999,
subagent_tx: sub_tx().clone(),
},
runs.path(),
);
assert_eq!(restored.len(), 1);
let entity = restored[0].1;
let entries = conversation_of(&world, entity.entity());
assert!(entries.iter().any(|e| matches!(
&e.kind,
leviath_core::region::EntryKind::AssistantTurn { tool_calls } if tool_calls.len() == 2
)));
assert!(
entries
.iter()
.any(|e| e.content == "Wrote 42 bytes to x.txt")
);
assert!(entries.iter().any(|e| e.content.contains("interrupted")
&& e.content.contains("Verify whether it took effect")));
assert!(
world
.world()
.get::<leviath_runtime::pipeline::ReadyToInfer>(entity.entity())
.is_some()
);
}
#[tokio::test]
async fn reload_does_not_replay_a_batch_already_in_the_window() {
use leviath_core::region::EntryKind;
use leviath_core::run_archive::RunRecord;
let agent = agent_dir();
let manifest = agent.path().join("agent.leviath");
let mpath = manifest.to_str().unwrap();
let runs = tempfile::tempdir().unwrap();
write_run(runs.path(), "run-applied", mpath, RunStatus::Running, None);
let ctx = ContextSnapshot {
stage_name: "implement".to_string(),
total_tokens: 2,
max_tokens: 100_000,
regions: vec![leviath_core::run_meta::RegionSnapshot {
name: "conversation".to_string(),
kind: "clearable".to_string(),
current_tokens: 2,
max_tokens: 100_000,
entries: vec![
leviath_core::run_meta::RegionEntrySnapshot {
content: "done".to_string(),
tokens: 1,
kind: EntryKind::AssistantTurn {
tool_calls: vec![leviath_core::region::SerializedToolCall {
id: "c1".to_string(),
name: "write_file".to_string(),
arguments: serde_json::Value::Null,
thought_signature: None,
}],
},
metadata: None,
key: None,
taint: Default::default(),
},
leviath_core::run_meta::RegionEntrySnapshot {
content: "Wrote it".to_string(),
tokens: 1,
kind: EntryKind::ToolResult {
tool_call_id: "c1".to_string(),
tool_name: "write_file".to_string(),
is_error: false,
},
metadata: None,
key: None,
taint: Default::default(),
},
],
}],
};
write_run_archive(runs.path(), "run-applied", mpath, 0, 9, 99, &ctx);
append_archive_records(
runs.path(),
"run-applied",
&[RunRecord::ToolBatch {
calls: vec![batch_call("c1", "write_file", None)],
at: 3,
stage_index: 0,
iteration: 9,
response: "done".to_string(),
}],
);
let (mut world, cli) = test_world();
let hub = InteractionHub::new();
let mcp = Arc::new(Mutex::new(ToolExecutor::new()));
let restored = reload_persisted_agents(
&mut world,
crate::daemon::spawn::SpawnDeps {
tool_service: cli.as_ref(),
config: &Config::default(),
shared_mcp: mcp,
mcp_tool_defs: &[],
hub: &hub,
now_secs: 999,
subagent_tx: sub_tx().clone(),
},
runs.path(),
);
assert_eq!(restored.len(), 1);
let entries = conversation_of(&world, restored[0].1.entity());
assert_eq!(
entries
.iter()
.filter(|e| matches!(&e.kind, EntryKind::AssistantTurn { tool_calls } if !tool_calls.is_empty()))
.count(),
1
);
assert!(!entries.iter().any(|e| e.content.contains("interrupted")));
}
fn interactive_agent_dir() -> tempfile::TempDir {
let dir = tempfile::tempdir().unwrap();
std::fs::write(
dir.path().join("agent.leviath"),
crate::test_support::inline_interactive_manifest(),
)
.unwrap();
dir
}
#[tokio::test]
async fn reload_resumes_a_blocked_interaction_point_in_the_waiting_state() {
let agent = interactive_agent_dir();
let manifest = agent.path().join("agent.leviath");
let runs = tempfile::tempdir().unwrap();
write_run(
runs.path(),
"run-await",
manifest.to_str().unwrap(),
RunStatus::WaitingInput,
None,
);
std::fs::write(
runs.path().join("run-await/interactions.json"),
serde_json::to_string(&InteractionPointState {
cursor: 0,
round: 0,
body: "## Plan\n1. do it".to_string(),
})
.unwrap(),
)
.unwrap();
let (mut world, cli) = test_world();
let hub = InteractionHub::new();
world.insert_interaction_hub(hub.clone()); let mcp = Arc::new(Mutex::new(ToolExecutor::new()));
let restored = reload_persisted_agents(
&mut world,
crate::daemon::spawn::SpawnDeps {
tool_service: cli.as_ref(),
config: &Config::default(),
shared_mcp: mcp,
mcp_tool_defs: &[],
hub: &hub,
now_secs: 999,
subagent_tx: sub_tx().clone(),
},
runs.path(),
);
assert_eq!(restored.len(), 1);
let (run_id, entity) = &restored[0];
assert_eq!(run_id, "run-await");
assert_eq!(world.agent_status(*entity), Some(AgentStatus::Waiting));
assert!(
world
.world()
.get::<leviath_runtime::interaction_points::AwaitingInteractionPoint>(
entity.entity()
)
.is_some()
);
assert!(
world
.world()
.get::<leviath_runtime::pipeline::ReadyToInfer>(entity.entity())
.is_none(),
"the spawn-set ReadyToInfer is cleared so the inference lane won't fire"
);
for _ in 0..8 {
tokio::task::yield_now().await;
}
let pending = hub.pending();
assert_eq!(pending.len(), 1);
assert_eq!(pending[0].0, "run-await");
assert_eq!(pending[0].1.body.as_deref(), Some("## Plan\n1. do it"));
}
#[tokio::test]
async fn reload_restores_actionable_runs_before_blocked_and_skips_terminal() {
let agent = agent_dir();
let mpath = agent.path().join("agent.leviath");
let mpath = mpath.to_str().unwrap();
let runs = tempfile::tempdir().unwrap();
write_run(
runs.path(),
"aaa-blocked",
mpath,
RunStatus::WaitingInput,
None,
);
write_run(runs.path(), "zzz-active", mpath, RunStatus::Running, None);
write_run(runs.path(), "mmm-done", mpath, RunStatus::Complete, None);
let (mut world, cli) = test_world();
let hub = InteractionHub::new();
let mcp = Arc::new(Mutex::new(ToolExecutor::new()));
let restored = reload_persisted_agents(
&mut world,
crate::daemon::spawn::SpawnDeps {
tool_service: cli.as_ref(),
config: &Config::default(),
shared_mcp: mcp,
mcp_tool_defs: &[],
hub: &hub,
now_secs: 999,
subagent_tx: sub_tx().clone(),
},
runs.path(),
);
let order: Vec<&str> = restored.iter().map(|(id, _)| id.as_str()).collect();
assert_eq!(order, vec!["zzz-active", "aaa-blocked"]);
}
#[tokio::test]
async fn reload_run_pages_in_nonterminal_only() {
let agent = agent_dir();
let manifest = agent.path().join("agent.leviath");
let mpath = manifest.to_str().unwrap();
let runs = tempfile::tempdir().unwrap();
write_run(runs.path(), "live", mpath, RunStatus::Running, None);
write_run(runs.path(), "done", mpath, RunStatus::Complete, None);
let (mut world, cli) = test_world();
let hub = InteractionHub::new();
let mcp = Arc::new(Mutex::new(ToolExecutor::new()));
assert!(
reload_run(
&mut world,
crate::daemon::spawn::SpawnDeps {
tool_service: cli.as_ref(),
config: &Config::default(),
shared_mcp: mcp.clone(),
mcp_tool_defs: &[],
hub: &hub,
now_secs: 1,
subagent_tx: sub_tx().clone(),
},
"live",
runs.path(),
)
.is_some()
);
assert!(
reload_run(
&mut world,
crate::daemon::spawn::SpawnDeps {
tool_service: cli.as_ref(),
config: &Config::default(),
shared_mcp: mcp.clone(),
mcp_tool_defs: &[],
hub: &hub,
now_secs: 1,
subagent_tx: sub_tx().clone(),
},
"done",
runs.path(),
)
.is_none()
);
assert!(
reload_run(
&mut world,
crate::daemon::spawn::SpawnDeps {
tool_service: cli.as_ref(),
config: &Config::default(),
shared_mcp: mcp,
mcp_tool_defs: &[],
hub: &hub,
now_secs: 1,
subagent_tx: sub_tx().clone(),
},
"no-such-run",
runs.path(),
)
.is_none()
);
}
#[tokio::test]
async fn resumes_a_parent_parked_mid_fan_out() {
use leviath_core::blueprint::{FanOutConfig, WorkerFailurePolicy};
use leviath_runtime::fanout::{FanOutState, FanOutWaiting};
let agent = agent_dir();
let manifest = agent.path().join("agent.leviath");
let mpath = manifest.to_str().unwrap();
let runs = tempfile::tempdir().unwrap();
write_run(
runs.path(),
"parent-fo",
mpath,
RunStatus::WaitingInput,
None,
);
let state = FanOutState {
config: FanOutConfig {
worker_agent: None,
worker_stage: Some("w".to_string()),
worker_query: None,
merge_stage: None,
max_workers: 1,
on_worker_failure: WorkerFailurePolicy::Continue,
split_prompt: "s".to_string(),
results_region: None,
max_items: None,
},
max_workers: 1,
pending: vec![],
active: vec![("item-1".to_string(), "worker-fo".to_string())],
summaries: vec![],
failures: vec![],
};
std::fs::write(
runs.path().join("parent-fo").join("fanout.json"),
serde_json::to_string(&state).unwrap(),
)
.unwrap();
write_run(runs.path(), "worker-fo", mpath, RunStatus::Running, None);
write_run(runs.path(), "bad-fo", mpath, RunStatus::WaitingInput, None);
std::fs::write(runs.path().join("bad-fo").join("fanout.json"), b"garbage").unwrap();
let (mut world, cli) = test_world();
let hub = InteractionHub::new();
let mcp = Arc::new(Mutex::new(ToolExecutor::new()));
let restored = reload_persisted_agents(
&mut world,
crate::daemon::spawn::SpawnDeps {
tool_service: cli.as_ref(),
config: &Config::default(),
shared_mcp: mcp,
mcp_tool_defs: &[],
hub: &hub,
now_secs: 999,
subagent_tx: sub_tx().clone(),
},
runs.path(),
);
let by_id: std::collections::HashMap<_, _> =
restored.iter().map(|(r, e)| (r.clone(), *e)).collect();
assert!(
world
.world()
.get::<FanOutWaiting>(by_id["parent-fo"].entity())
.is_some()
);
assert!(
world
.world()
.get::<FanOutWaiting>(by_id["bad-fo"].entity())
.is_none()
);
}
#[tokio::test]
async fn rebuilds_parent_child_tree_on_reload() {
use leviath_runtime::components::{ParentRef, SubAgentChildren};
let agent = agent_dir();
let manifest = agent.path().join("agent.leviath");
let mpath = manifest.to_str().unwrap();
let runs = tempfile::tempdir().unwrap();
write_run_tree(RunFixture {
runs_dir: runs.path(),
run_id: "parent",
agent_path: mpath,
status: RunStatus::WaitingInput,
context: None,
parent_run_id: None,
children: &["child-a", "child-b"],
depth: 0,
max_child_depth: 4,
});
write_run_tree(RunFixture {
runs_dir: runs.path(),
run_id: "child-a",
agent_path: mpath,
status: RunStatus::Running,
context: None,
parent_run_id: Some("parent"),
children: &[],
depth: 1,
max_child_depth: 0,
});
write_run_tree(RunFixture {
runs_dir: runs.path(),
run_id: "child-b",
agent_path: mpath,
status: RunStatus::Running,
context: None,
parent_run_id: Some("parent"),
children: &[],
depth: 1,
max_child_depth: 0,
});
let (mut world, cli) = test_world();
let hub = InteractionHub::new();
let mcp = Arc::new(Mutex::new(ToolExecutor::new()));
let restored = reload_persisted_agents(
&mut world,
crate::daemon::spawn::SpawnDeps {
tool_service: cli.as_ref(),
config: &Config::default(),
shared_mcp: mcp,
mcp_tool_defs: &[],
hub: &hub,
now_secs: 999,
subagent_tx: sub_tx().clone(),
},
runs.path(),
);
assert_eq!(restored.len(), 3);
let by_id: std::collections::HashMap<_, _> =
restored.iter().map(|(r, e)| (r.clone(), *e)).collect();
let parent = by_id["parent"];
let child_a = by_id["child-a"];
let child_b = by_id["child-b"];
let kids = world
.world()
.get::<SubAgentChildren>(parent.entity())
.unwrap();
assert_eq!(kids.max_child_depth, 4);
assert_eq!(kids.children.len(), 2);
assert!(
kids.children.contains(&child_a.entity()) && kids.children.contains(&child_b.entity())
);
let pr = world.world().get::<ParentRef>(child_a.entity()).unwrap();
assert_eq!(pr.parent_entity, parent.entity());
assert_eq!(pr.parent_agent_id, "parent");
assert_eq!(pr.depth, 1);
let state = world
.world()
.get::<leviath_runtime::components::AgentState>(parent.entity())
.unwrap();
assert_eq!(state.spawned_children_ids, vec!["child-a", "child-b"]);
}
#[tokio::test]
async fn relink_skips_children_and_parents_that_did_not_reload() {
use leviath_runtime::components::{ParentRef, SubAgentChildren};
let agent = agent_dir();
let manifest = agent.path().join("agent.leviath");
let mpath = manifest.to_str().unwrap();
let runs = tempfile::tempdir().unwrap();
write_run_tree(RunFixture {
runs_dir: runs.path(),
run_id: "lonely-parent",
agent_path: mpath,
status: RunStatus::WaitingInput,
context: None,
parent_run_id: None,
children: &["gone-child"],
depth: 0,
max_child_depth: 2,
});
write_run_tree(RunFixture {
runs_dir: runs.path(),
run_id: "gone-child",
agent_path: mpath,
status: RunStatus::Complete,
context: None,
parent_run_id: Some("lonely-parent"),
children: &[],
depth: 1,
max_child_depth: 0,
});
write_run_tree(RunFixture {
runs_dir: runs.path(),
run_id: "orphan",
agent_path: mpath,
status: RunStatus::Running,
context: None,
parent_run_id: Some("gone-parent"),
children: &[],
depth: 1,
max_child_depth: 0,
});
write_run_tree(RunFixture {
runs_dir: runs.path(),
run_id: "gone-parent",
agent_path: mpath,
status: RunStatus::Error,
context: None,
parent_run_id: None,
children: &["orphan"],
depth: 0,
max_child_depth: 2,
});
let (mut world, cli) = test_world();
let hub = InteractionHub::new();
let mcp = Arc::new(Mutex::new(ToolExecutor::new()));
let restored = reload_persisted_agents(
&mut world,
crate::daemon::spawn::SpawnDeps {
tool_service: cli.as_ref(),
config: &Config::default(),
shared_mcp: mcp,
mcp_tool_defs: &[],
hub: &hub,
now_secs: 999,
subagent_tx: sub_tx().clone(),
},
runs.path(),
);
assert_eq!(restored.len(), 2);
let by_id: std::collections::HashMap<_, _> =
restored.iter().map(|(r, e)| (r.clone(), *e)).collect();
assert!(
world
.world()
.get::<SubAgentChildren>(by_id["lonely-parent"].entity())
.is_none()
);
assert!(
world
.world()
.get::<ParentRef>(by_id["orphan"].entity())
.is_none()
);
}
#[tokio::test]
async fn reload_without_context_json_still_resumes() {
let agent = agent_dir();
let manifest = agent.path().join("agent.leviath");
let runs = tempfile::tempdir().unwrap();
write_run(
runs.path(),
"run-nocontext",
manifest.to_str().unwrap(),
RunStatus::WaitingInput,
None, );
let (mut world, cli) = test_world();
let hub = InteractionHub::new();
let mcp = Arc::new(Mutex::new(ToolExecutor::new()));
let restored = reload_persisted_agents(
&mut world,
crate::daemon::spawn::SpawnDeps {
tool_service: cli.as_ref(),
config: &Config::default(),
shared_mcp: mcp,
mcp_tool_defs: &[],
hub: &hub,
now_secs: 999,
subagent_tx: sub_tx().clone(),
},
runs.path(),
);
assert_eq!(restored.len(), 1);
assert!(
world
.world()
.get::<TokenTotals>(restored[0].1.entity())
.is_some()
);
}
#[tokio::test]
async fn reload_restores_the_persisted_stage_ledger() {
use leviath_core::run_meta::{StageRecord, StageRunStatus};
use leviath_runtime::pipeline::StageLedger;
let agent = agent_dir();
let manifest = agent.path().join("agent.leviath");
let runs = tempfile::tempdir().unwrap();
write_run(
runs.path(),
"run-stages",
manifest.to_str().unwrap(),
RunStatus::Running,
None,
);
let mut analyze = StageRecord::new("analyze".to_string(), 0);
analyze.status = StageRunStatus::Complete;
analyze.entered = true;
analyze.prompt_tokens = 1_234;
analyze.completion_tokens = 56;
analyze.cached_tokens = 7;
analyze.cache_write_tokens = 8;
analyze.first_call_prompt_tokens = Some(400);
analyze.runaway_warned = true;
analyze
.region_tokens
.insert("conversation".to_string(), 900);
analyze.started_at = Some(10);
analyze.ended_at = Some(20);
let mut implement = StageRecord::new("implement".to_string(), 1);
implement.status = StageRunStatus::Active;
implement.entered = true;
implement.prompt_tokens = 77;
implement.started_at = Some(20);
let removed = StageRecord::new("removed_stage".to_string(), 7);
std::fs::write(
runs.path().join("run-stages").join("stages.json"),
serde_json::to_string(&vec![analyze, implement, removed]).unwrap(),
)
.unwrap();
let (mut world, cli) = test_world();
let hub = InteractionHub::new();
let mcp = Arc::new(Mutex::new(ToolExecutor::new()));
let restored = reload_persisted_agents(
&mut world,
crate::daemon::spawn::SpawnDeps {
tool_service: cli.as_ref(),
config: &Config::default(),
shared_mcp: mcp,
mcp_tool_defs: &[],
hub: &hub,
now_secs: 999,
subagent_tx: sub_tx().clone(),
},
runs.path(),
);
assert_eq!(restored.len(), 1);
let ledger = world
.world()
.get::<StageLedger>(restored[0].1.entity())
.expect("a reloaded agent carries a stage ledger");
let names: Vec<&str> = ledger.0.iter().map(|r| r.name.as_str()).collect();
assert_eq!(names, vec!["analyze", "implement", "review"]);
assert_eq!(ledger.0[0].prompt_tokens, 1_234);
assert_eq!(ledger.0[0].completion_tokens, 56);
assert_eq!(ledger.0[0].cached_tokens, 7);
assert_eq!(ledger.0[0].cache_write_tokens, 8);
assert_eq!(ledger.0[0].first_call_prompt_tokens, Some(400));
assert!(ledger.0[0].runaway_warned);
assert_eq!(ledger.0[0].region_tokens.get("conversation"), Some(&900));
assert_eq!(ledger.0[0].started_at, Some(10));
assert_eq!(ledger.0[0].ended_at, Some(20));
assert_eq!(ledger.0[0].status, StageRunStatus::Complete);
assert!(ledger.0[0].entered);
assert_eq!(ledger.0[1].prompt_tokens, 77);
assert!(ledger.0[1].entered);
assert_eq!(ledger.0[2].prompt_tokens, 0);
assert!(!ledger.0[2].entered);
assert_eq!(ledger.0[2].index, 2);
}
#[tokio::test]
async fn skips_missing_dir_junk_and_unreloadable_runs() {
let (mut world, cli) = test_world();
let hub = InteractionHub::new();
let mcp = Arc::new(Mutex::new(ToolExecutor::new()));
assert!(
reload_persisted_agents(
&mut world,
crate::daemon::spawn::SpawnDeps {
tool_service: cli.as_ref(),
config: &Config::default(),
shared_mcp: mcp.clone(),
mcp_tool_defs: &[],
hub: &hub,
now_secs: 1,
subagent_tx: sub_tx().clone(),
},
std::path::Path::new("/no/such/runs/dir"),
)
.is_empty()
);
let runs = tempfile::tempdir().unwrap();
std::fs::create_dir_all(runs.path().join("no-meta")).unwrap();
let corrupt = runs.path().join("corrupt");
std::fs::create_dir_all(&corrupt).unwrap();
std::fs::write(corrupt.join("meta.json"), "not json").unwrap();
write_run(
runs.path(),
"run-badpath",
"/no/such/agent.leviath",
RunStatus::Running,
None,
);
let restored = reload_persisted_agents(
&mut world,
crate::daemon::spawn::SpawnDeps {
tool_service: cli.as_ref(),
config: &Config::default(),
shared_mcp: mcp,
mcp_tool_defs: &[],
hub: &hub,
now_secs: 1,
subagent_tx: sub_tx().clone(),
},
runs.path(),
);
assert!(restored.is_empty());
let meta: RunMeta = serde_json::from_str(
&std::fs::read_to_string(runs.path().join("run-badpath").join("meta.json")).unwrap(),
)
.unwrap();
assert_eq!(meta.status, RunStatus::Error);
let error = meta.error.unwrap_or_default();
assert!(error.contains("could not be recovered"), "got: {error}");
assert_eq!(meta.updated_at, 1);
assert!(!runs.path().join("no-meta").join("meta.json").exists());
assert_eq!(
std::fs::read_to_string(corrupt.join("meta.json")).unwrap(),
"not json"
);
}
#[test]
fn marking_a_crash_is_best_effort() {
let runs = tempfile::tempdir().unwrap();
write_run(
runs.path(),
"run-x",
"/no/such/agent.leviath",
RunStatus::Running,
None,
);
let meta = read_meta(&runs.path().join("run-x")).expect("written above");
mark_crashed(&runs.path().join("gone"), meta, "boom", 7);
assert!(!runs.path().join("gone").exists());
}
#[tokio::test]
async fn fake_provider_methods_are_exercised() {
use leviath_providers::Provider;
let p = FakeProvider;
assert_eq!(p.name(), "fake");
assert_eq!(p.count_tokens("t", "m").await, 1);
assert_eq!(p.max_context_tokens("m"), 1000);
let _ = p.capabilities("m");
assert!(
p.infer(&leviath_providers::InferenceRequest {
system: vec![],
messages: vec![],
model: "m".to_string(),
max_tokens: 1,
temperature: 0.0,
tools: vec![],
extra: serde_json::Value::Null,
request_timeout_secs: None,
})
.await
.is_err()
);
}
#[test]
fn is_terminal_covers_all_statuses() {
assert!(is_terminal(&RunStatus::Complete));
assert!(is_terminal(&RunStatus::Cancelled));
assert!(is_terminal(&RunStatus::Error));
assert!(!is_terminal(&RunStatus::Running));
assert!(!is_terminal(&RunStatus::WaitingInput));
}
}