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