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