1use std::path::{Path, PathBuf};
12use std::sync::Arc;
13
14use bevy_ecs::entity::Entity;
15use bevy_ecs::world::World;
16use leviath_core::blueprint::FanOutConfig;
17use leviath_providers::Tool;
18use leviath_runtime::fanout::FanOutSpawner;
19use leviath_runtime::host::SubAgentOp;
20use leviath_runtime::interaction_hub::InteractionHub;
21use leviath_runtime::persistence::RunMetadata;
22use leviath_runtime::pipeline::{AgentBlueprint, force_transition};
23use tokio::sync::Mutex;
24use tokio::sync::mpsc::UnboundedSender;
25
26use crate::config::Config;
27use crate::daemon::client::resolve_spawn_args;
28use crate::daemon::spawn::build_agent;
29use crate::daemon::tool_service::CliToolService;
30
31#[derive(Clone)]
34pub struct DaemonFanOutSpawner {
35 pub config: Config,
36 pub shared_mcp: Arc<Mutex<leviath_mcp::ToolExecutor>>,
37 pub mcp_tool_defs: Vec<Tool>,
38 pub mcp_pool: Arc<crate::daemon::mcp_pool::McpPool>,
42 pub hub: InteractionHub,
43 pub subagent_tx: UnboundedSender<SubAgentOp>,
44 pub tool_service: Arc<CliToolService>,
45 pub agents_dir: Option<PathBuf>,
48 pub now_secs: fn() -> i64,
49}
50
51impl DaemonFanOutSpawner {
52 fn worker_mcp_defs(&self, blueprint_path: &str) -> Vec<Tool> {
59 let mut defs = self.mcp_tool_defs.clone();
60 let Ok(toml) = std::fs::read_to_string(blueprint_path) else {
61 return defs;
62 };
63 let servers = crate::daemon::mcp_pool::parse_blueprint_mcp_servers(&toml);
64 if servers.is_empty() {
65 return defs;
66 }
67 defs.extend(self.mcp_pool.cached_defs_for(&servers));
68 tokio::runtime::Handle::current().spawn(self.mcp_pool.clone().ensure_all(servers));
71 defs
72 }
73}
74
75impl FanOutSpawner for DaemonFanOutSpawner {
76 fn spawn_worker(
77 &self,
78 world: &mut World,
79 parent: Entity,
80 config: &FanOutConfig,
81 item_id: &str,
82 item_context: &serde_json::Value,
83 ) -> Result<Entity, String> {
84 let (parent_path, workdir, parent_run_id) = world
87 .get::<RunMetadata>(parent)
88 .map(|md| (md.agent_path.clone(), md.workdir.clone(), md.run_id.clone()))
89 .ok_or_else(|| "fan-out parent has no run metadata".to_string())?;
90
91 let (resolve_path, entry_stage) =
92 resolve_worker_source(config, &parent_path, self.agents_dir.as_deref())?;
93
94 let task = format_worker_task(item_id, item_context);
95 let mut args = resolve_spawn_args(
96 &resolve_path,
97 &task,
98 None,
99 &workdir,
100 false,
101 Vec::new(),
102 None,
103 std::collections::HashMap::new(),
105 true,
110 )
111 .map_err(|e| format!("resolve worker blueprint: {e}"))?;
112 args.parent_run_id = Some(parent_run_id);
114
115 let mcp_defs = self.worker_mcp_defs(&args.blueprint_path);
120
121 let child = build_agent(
122 world,
123 self.tool_service.as_ref(),
124 &self.config,
125 self.shared_mcp.clone(),
126 &mcp_defs,
127 &self.hub,
128 &args,
129 (self.now_secs)(),
130 self.subagent_tx.clone(),
131 )?;
132
133 if let Some(stage) = entry_stage {
136 match world
137 .get::<AgentBlueprint>(child)
138 .and_then(|bp| bp.0.stages.iter().position(|s| s.name == stage))
139 {
140 Some(idx) => force_transition(world, child, idx),
141 None => return Err(format!("worker stage '{stage}' not found in blueprint")),
142 }
143 }
144 Ok(child)
145 }
146}
147
148pub(crate) fn resolve_worker_source(
153 config: &FanOutConfig,
154 parent_path: &str,
155 agents_dir: Option<&Path>,
156) -> Result<(String, Option<String>), String> {
157 if let Some(stage) = &config.worker_stage {
158 return Ok((parent_path.to_string(), Some(stage.clone())));
159 }
160 if let Some(agent) = &config.worker_agent {
161 return Ok((agent.clone(), None));
162 }
163 if let Some(query) = &config.worker_query {
164 let path = discover_worker(agents_dir, query)?;
165 return Ok((path.to_string_lossy().to_string(), None));
166 }
167 Err("fan-out config has no worker source".to_string())
168}
169
170fn format_worker_task(item_id: &str, item_context: &serde_json::Value) -> String {
172 format!("Work item id: {item_id}\nContext: {item_context}")
173}
174
175fn discover_worker(agents_dir: Option<&Path>, query: &str) -> Result<PathBuf, String> {
178 let dir = agents_dir.ok_or_else(|| "no agents directory to search for a worker".to_string())?;
179 let needle = query.to_lowercase();
180 let entries =
181 std::fs::read_dir(dir).map_err(|e| format!("read agents dir '{}': {e}", dir.display()))?;
182 for entry in entries.flatten() {
183 let path = entry.path();
184 let manifest = path.join("agent.leviath");
185 if !manifest.is_file() {
186 continue;
187 }
188 let name_matches = path
189 .file_name()
190 .and_then(|n| n.to_str())
191 .is_some_and(|n| n.to_lowercase().contains(&needle));
192 let desc_matches = std::fs::read_to_string(&manifest)
193 .ok()
194 .and_then(|c| leviath_core::manifest::parse_manifest(&c).ok())
195 .is_some_and(|bp| bp.description.to_lowercase().contains(&needle));
196 if name_matches || desc_matches {
197 return Ok(path);
198 }
199 }
200 Err(format!("no installed agent matches worker query '{query}'"))
201}
202
203#[cfg(test)]
204mod tests {
205 use super::*;
206 use leviath_core::blueprint::WorkerFailurePolicy;
207
208 fn cfg(stage: Option<&str>, agent: Option<&str>, query: Option<&str>) -> FanOutConfig {
209 FanOutConfig {
210 worker_agent: agent.map(String::from),
211 worker_stage: stage.map(String::from),
212 worker_query: query.map(String::from),
213 merge_stage: None,
214 max_workers: 4,
215 on_worker_failure: WorkerFailurePolicy::Continue,
216 split_prompt: "split".to_string(),
217 }
218 }
219
220 #[test]
221 fn format_worker_task_includes_id_and_context() {
222 let t = format_worker_task("t1", &serde_json::json!({"file": "a.rs"}));
223 assert!(t.contains("Work item id: t1"));
224 assert!(t.contains("a.rs"));
225 }
226
227 #[test]
228 fn resolve_worker_source_picks_the_configured_source() {
229 assert_eq!(
231 resolve_worker_source(&cfg(Some("w"), None, None), "/p/agent.leviath", None).unwrap(),
232 ("/p/agent.leviath".to_string(), Some("w".to_string()))
233 );
234 assert_eq!(
236 resolve_worker_source(&cfg(None, Some("fixer"), None), "/p", None).unwrap(),
237 ("fixer".to_string(), None)
238 );
239 }
240
241 #[test]
242 fn resolve_worker_source_worker_query_discovers_an_agent() {
243 let dir = tempfile::tempdir().unwrap();
244 let agent = dir.path().join("test-fixer");
245 std::fs::create_dir_all(&agent).unwrap();
246 std::fs::write(
247 agent.join("agent.leviath"),
248 "[agent]\nname = \"test-fixer\"\nversion = \"0.1.0\"\ndescription = \"fixes tests\"\n\n[stages.main]\nmodel = { provider = \"anthropic\", model = \"claude-sonnet-4-6\" }\n",
249 )
250 .unwrap();
251 let (path, entry) =
252 resolve_worker_source(&cfg(None, None, Some("fixer")), "/p", Some(dir.path())).unwrap();
253 assert!(path.contains("test-fixer"));
254 assert_eq!(entry, None);
255
256 let empty = tempfile::tempdir().unwrap();
258 assert!(
259 resolve_worker_source(&cfg(None, None, Some("zzz")), "/p", Some(empty.path())).is_err()
260 );
261 }
262
263 #[test]
264 fn resolve_worker_source_errors_without_a_source() {
265 let empty = cfg(None, None, None);
266 assert!(resolve_worker_source(&empty, "/p", None).is_err());
267 }
268
269 #[test]
270 fn discover_worker_matches_name_or_description_and_reports_misses() {
271 let dir = tempfile::tempdir().unwrap();
272 let a = dir.path().join("alpha");
274 std::fs::create_dir_all(&a).unwrap();
275 std::fs::write(
276 a.join("agent.leviath"),
277 "[agent]\nname = \"alpha\"\nversion = \"0.1.0\"\ndescription = \"a widget wrangler\"\n\n[stages.main]\nmodel = { provider = \"anthropic\", model = \"claude-sonnet-4-6\" }\n",
278 )
279 .unwrap();
280 std::fs::create_dir_all(dir.path().join("not-an-agent")).unwrap();
282
283 assert!(
285 discover_worker(Some(dir.path()), "widget")
286 .unwrap()
287 .ends_with("alpha")
288 );
289 assert!(
291 discover_worker(Some(dir.path()), "ALPHA")
292 .unwrap()
293 .ends_with("alpha")
294 );
295 assert!(discover_worker(Some(dir.path()), "nonexistent").is_err());
297 assert!(discover_worker(None, "x").is_err());
299 let file = dir.path().join("alpha").join("agent.leviath");
301 assert!(discover_worker(Some(&file), "x").is_err());
302 }
303
304 use leviath_runtime::components::AgentStatus;
307 use leviath_runtime::host::SpawnArgs;
308 use leviath_runtime::inference_pool::InferencePoolConfig;
309 use leviath_runtime::pipeline::StageCursor;
310 use leviath_runtime::world::PipelineWorld;
311 use std::collections::HashMap;
312 use tokio::runtime::Handle;
313
314 struct FakeProvider;
315 #[async_trait::async_trait]
316 impl leviath_providers::Provider for FakeProvider {
317 async fn infer(
318 &self,
319 _r: leviath_providers::InferenceRequest,
320 ) -> leviath_providers::Result<leviath_providers::InferenceResponse> {
321 Err(leviath_providers::ProviderError::Other("test".to_string()))
322 }
323 async fn count_tokens(&self, _t: &str, _m: &str) -> usize {
324 1
325 }
326 fn max_context_tokens(&self, _m: &str) -> usize {
327 100_000
328 }
329 fn name(&self) -> &str {
330 "fake"
331 }
332 fn capabilities(&self, _m: &str) -> leviath_providers::ModelCapabilities {
333 leviath_providers::ModelCapabilities::default()
334 }
335 }
336
337 fn two_stage_manifest() -> String {
339 "[agent]\nname = \"host\"\nversion = \"0.1.0\"\ndescription = \"d\"\nentry_stage = \"first\"\n\n\
340 [stages.first]\nmodel = { provider = \"anthropic\", model = \"m\" }\nsystem_prompt = \"first\"\n\n\
341 [stages.second]\nmodel = { provider = \"anthropic\", model = \"m\" }\nallow_as_worker = true\nsystem_prompt = \"second\"\n"
342 .to_string()
343 }
344
345 fn spawner_with(tool_service: Arc<CliToolService>) -> DaemonFanOutSpawner {
346 let shared_mcp = Arc::new(Mutex::new(leviath_mcp::ToolExecutor::new()));
347 DaemonFanOutSpawner {
348 config: Config::default(),
349 shared_mcp: shared_mcp.clone(),
350 mcp_tool_defs: vec![],
351 mcp_pool: crate::daemon::mcp_pool::McpPool::for_daemon(shared_mcp, &[]),
352 hub: InteractionHub::new(),
353 subagent_tx: tokio::sync::mpsc::unbounded_channel().0,
354 tool_service,
355 agents_dir: None,
356 now_secs: || 100,
357 }
358 }
359
360 #[tokio::test]
361 async fn worker_mcp_defs_advertises_cached_servers_and_falls_back_to_global() {
362 let mut spawner = spawner_with(Arc::new(CliToolService::new()));
363 spawner.mcp_tool_defs = vec![Tool {
364 name: "global_tool".to_string(),
365 description: String::new(),
366 parameters: serde_json::json!({}),
367 }];
368 let dir = tempfile::tempdir().unwrap();
370 let manifest = dir.path().join("agent.leviath");
371 std::fs::write(
372 &manifest,
373 "[agent]\nname = \"w\"\n\n[[mcp_servers]]\nname = \"srv\"\ncommand = \"python3\"\nargs = [\"pass\"]\n",
374 )
375 .unwrap();
376 let servers = crate::daemon::mcp_pool::parse_blueprint_mcp_servers(
379 &std::fs::read_to_string(&manifest).unwrap(),
380 );
381 spawner.mcp_pool.seed(
382 &servers[0],
383 vec![Tool {
384 name: "srv_tool".to_string(),
385 description: String::new(),
386 parameters: serde_json::json!({}),
387 }],
388 );
389 let defs = spawner.worker_mcp_defs(&manifest.to_string_lossy());
390 let names: Vec<&str> = defs.iter().map(|t| t.name.as_str()).collect();
391 assert_eq!(names, vec!["global_tool", "srv_tool"]);
392 let only_global = spawner.worker_mcp_defs("/no/such/agent.leviath");
394 assert_eq!(only_global.len(), 1);
395 assert_eq!(only_global[0].name, "global_tool");
396 let empty = dir.path().join("empty.leviath");
398 std::fs::write(&empty, "[agent]\nname = \"e\"\n").unwrap();
399 assert_eq!(spawner.worker_mcp_defs(&empty.to_string_lossy()).len(), 1);
400 }
401
402 fn world_with_parent(manifest_path: &str) -> (PipelineWorld, DaemonFanOutSpawner, Entity) {
405 let cli = Arc::new(CliToolService::new());
406 let mut registry = leviath_runtime::ProviderRegistry::new();
407 registry.register("anthropic".to_string(), Arc::new(FakeProvider));
408 let mut world = PipelineWorld::new(
409 registry,
410 cli.clone(),
411 InferencePoolConfig::new(),
412 1,
413 std::env::temp_dir(),
414 Handle::current(),
415 );
416 let spawner = spawner_with(cli.clone());
417 let args = SpawnArgs {
418 run_id: "parent".to_string(),
419 blueprint_path: manifest_path.to_string(),
420 task: "parent task".to_string(),
421 regions: HashMap::new(),
422 model: None,
423 workdir: std::env::temp_dir().to_string_lossy().to_string(),
424 metadata: HashMap::new(),
425 callback_url: None,
426 callback_secret: None,
427 yolo: false,
428 no_seed_commands: false,
429 allow: Vec::new(),
430 max_depth: None,
431 parent_run_id: None,
432 };
433 let parent = build_agent(
434 world.world_mut(),
435 cli.as_ref(),
436 &spawner.config,
437 spawner.shared_mcp.clone(),
438 &spawner.mcp_tool_defs,
439 &spawner.hub,
440 &args,
441 100,
442 spawner.subagent_tx.clone(),
443 )
444 .expect("parent spawns");
445 (world, spawner, parent)
446 }
447
448 #[tokio::test]
449 async fn spawn_worker_worker_stage_enters_the_worker_stage() {
450 let dir = tempfile::tempdir().unwrap();
451 let manifest = dir.path().join("agent.leviath");
452 std::fs::write(&manifest, two_stage_manifest()).unwrap();
453 let (mut world, spawner, parent) = world_with_parent(&manifest.to_string_lossy());
454
455 let child = spawner
456 .spawn_worker(
457 world.world_mut(),
458 parent,
459 &cfg(Some("second"), None, None),
460 "item-1",
461 &serde_json::json!({"k": "v"}),
462 )
463 .expect("worker spawns");
464 assert_eq!(world.world().get::<StageCursor>(child).unwrap().index, 1);
466 assert_eq!(world.agent_status(child), Some(AgentStatus::Active));
467 }
468
469 #[tokio::test]
470 async fn spawn_worker_worker_agent_uses_a_separate_blueprint() {
471 let dir = tempfile::tempdir().unwrap();
472 let parent_m = dir.path().join("agent.leviath");
473 std::fs::write(&parent_m, two_stage_manifest()).unwrap();
474 let (mut world, spawner, parent) = world_with_parent(&parent_m.to_string_lossy());
475
476 let worker_dir = dir.path().join("worker");
478 std::fs::create_dir_all(&worker_dir).unwrap();
479 std::fs::write(worker_dir.join("agent.leviath"), two_stage_manifest()).unwrap();
480 let child = spawner
481 .spawn_worker(
482 world.world_mut(),
483 parent,
484 &cfg(None, Some(&worker_dir.to_string_lossy()), None),
485 "item-1",
486 &serde_json::json!({}),
487 )
488 .expect("worker spawns");
489 assert_eq!(world.world().get::<StageCursor>(child).unwrap().index, 0);
491 }
492
493 #[tokio::test]
494 async fn spawn_worker_errors_without_parent_metadata() {
495 let dir = tempfile::tempdir().unwrap();
496 let manifest = dir.path().join("agent.leviath");
497 std::fs::write(&manifest, two_stage_manifest()).unwrap();
498 let (mut world, spawner, _parent) = world_with_parent(&manifest.to_string_lossy());
499 let bare = world.world_mut().spawn_empty().id();
501 let err = spawner
502 .spawn_worker(
503 world.world_mut(),
504 bare,
505 &cfg(Some("second"), None, None),
506 "i",
507 &serde_json::json!({}),
508 )
509 .unwrap_err();
510 assert!(err.contains("no run metadata"));
511 }
512
513 #[tokio::test]
514 async fn spawn_worker_errors_when_worker_stage_is_missing() {
515 let dir = tempfile::tempdir().unwrap();
516 let manifest = dir.path().join("agent.leviath");
517 std::fs::write(&manifest, two_stage_manifest()).unwrap();
518 let (mut world, spawner, parent) = world_with_parent(&manifest.to_string_lossy());
519 let err = spawner
520 .spawn_worker(
521 world.world_mut(),
522 parent,
523 &cfg(Some("ghost"), None, None),
524 "i",
525 &serde_json::json!({}),
526 )
527 .unwrap_err();
528 assert!(err.contains("ghost"));
529 }
530
531 #[tokio::test]
532 async fn spawn_worker_propagates_a_build_error() {
533 let dir = tempfile::tempdir().unwrap();
534 let manifest = dir.path().join("agent.leviath");
535 std::fs::write(&manifest, two_stage_manifest()).unwrap();
536 let (mut world, spawner, parent) = world_with_parent(&manifest.to_string_lossy());
537
538 let bad_dir = dir.path().join("bad");
541 std::fs::create_dir_all(&bad_dir).unwrap();
542 std::fs::write(
543 bad_dir.join("agent.leviath"),
544 "[agent]\nname = \"bad\"\nversion = \"0.1.0\"\ndescription = \"d\"\n\n\
545 [stages.only]\nmodel = { provider = \"anthropic\", model = \"m\" }\n\n\
546 [stages.only.transitions.nowhere]\n",
547 )
548 .unwrap();
549 let err = spawner
550 .spawn_worker(
551 world.world_mut(),
552 parent,
553 &cfg(None, Some(&bad_dir.to_string_lossy()), None),
554 "i",
555 &serde_json::json!({}),
556 )
557 .unwrap_err();
558 assert!(err.contains("invalid blueprint"));
559 }
560
561 #[tokio::test]
562 async fn spawn_worker_propagates_a_worker_source_error() {
563 let dir = tempfile::tempdir().unwrap();
564 let manifest = dir.path().join("agent.leviath");
565 std::fs::write(&manifest, two_stage_manifest()).unwrap();
566 let (mut world, spawner, parent) = world_with_parent(&manifest.to_string_lossy());
567 let err = spawner
569 .spawn_worker(
570 world.world_mut(),
571 parent,
572 &cfg(None, None, None),
573 "i",
574 &serde_json::json!({}),
575 )
576 .unwrap_err();
577 assert!(err.contains("no worker source"));
578 }
579
580 #[tokio::test]
581 async fn fake_provider_metadata_is_exercised() {
582 use leviath_providers::Provider;
583 let p = FakeProvider;
584 assert_eq!(p.name(), "fake");
585 assert_eq!(p.count_tokens("t", "m").await, 1);
586 assert_eq!(p.max_context_tokens("m"), 100_000);
587 let _ = p.capabilities("m");
588 }
589
590 #[tokio::test]
591 async fn fake_provider_infer_errors() {
592 use leviath_providers::Provider;
593 let p = FakeProvider;
594 assert!(
595 p.infer(leviath_providers::InferenceRequest {
596 system: vec![],
597 messages: vec![],
598 model: "m".to_string(),
599 max_tokens: 1,
600 temperature: 0.0,
601 tools: vec![],
602 extra: serde_json::Value::Null,
603 request_timeout_secs: None,
604 })
605 .await
606 .is_err()
607 );
608 }
609
610 #[tokio::test]
611 async fn spawn_worker_propagates_a_resolve_error() {
612 let dir = tempfile::tempdir().unwrap();
613 let manifest = dir.path().join("agent.leviath");
614 std::fs::write(&manifest, two_stage_manifest()).unwrap();
615 let (mut world, spawner, parent) = world_with_parent(&manifest.to_string_lossy());
616 let err = spawner
618 .spawn_worker(
619 world.world_mut(),
620 parent,
621 &cfg(None, Some("/no/such/agent/xyz"), None),
622 "i",
623 &serde_json::json!({}),
624 )
625 .unwrap_err();
626 assert!(err.contains("resolve worker blueprint"));
627 }
628
629 #[test]
630 fn discover_worker_skips_agents_with_unparsable_manifests() {
631 let dir = tempfile::tempdir().unwrap();
632 let bad = dir.path().join("broken");
633 std::fs::create_dir_all(&bad).unwrap();
634 std::fs::write(bad.join("agent.leviath"), "this is not valid toml : : :").unwrap();
635 assert!(discover_worker(Some(dir.path()), "zzz").is_err());
637 assert!(
639 discover_worker(Some(dir.path()), "broken")
640 .unwrap()
641 .ends_with("broken")
642 );
643 }
644}