Skip to main content

leviath_cli/daemon/
fanout_spawner.rs

1//! The daemon-side [`FanOutSpawner`]: resolves a fan-out worker's blueprint
2//! (self-at-worker-stage / a named agent / a capability query) and starts it in
3//! the shared world, seeded with its work item.
4//!
5//! The runtime's fan-out systems only *start and track* workers; *finding* the
6//! blueprint is CLI policy, so it lives here - mirroring how [`build_agent`]
7//! resolves any spawn. For `worker_stage` the worker runs the parent's own
8//! blueprint entered at that stage (via [`leviath_runtime::pipeline::force_transition`]);
9//! for `worker_agent` / `worker_query` it runs a separate installed blueprint.
10
11use 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/// Everything [`build_agent`] needs, captured so a fan-out worker can be spawned
31/// from inside a world-system (which has no access to the daemon's context).
32#[derive(Clone)]
33pub struct DaemonFanOutSpawner {
34    /// Spawn-time config, read fresh per worker so a `config.toml` edit reaches
35    /// fan-out workers too (shared with the daemon's main spawner).
36    pub config: Arc<crate::daemon::config_reload::ConfigReloader>,
37    /// The daemon-wide MCP executor, for servers configured globally.
38    pub shared_mcp: Arc<Mutex<leviath_mcp::ToolExecutor>>,
39    /// Tool definitions from `shared_mcp`, resolved once rather than per worker.
40    pub mcp_tool_defs: Vec<Tool>,
41    /// Shared MCP pool for per-agent `[[mcp_servers]]` - a fan-out worker
42    /// advertises its blueprint's already-connected servers and lazily warms any
43    /// uncached ones for subsequent workers of the same type.
44    pub mcp_pool: Arc<crate::daemon::mcp_pool::McpPool>,
45    /// Where a worker's prompts go. Shared with the daemon, so a fan-out
46    /// worker's question reaches the same place as any other run's.
47    pub hub: InteractionHub,
48    /// Where a worker's own sub-agent spawns are sent.
49    pub subagent_tx: UnboundedSender<SubAgentOp>,
50    /// The tool dispatcher every worker's calls go through.
51    pub tool_service: Arc<CliToolService>,
52    /// `~/.leviath/agents`, for resolving `worker_query`. `None` when there is no
53    /// home directory.
54    pub agents_dir: Option<PathBuf>,
55    /// The clock, injected so a test can stamp a worker deterministically.
56    pub now_secs: fn() -> i64,
57}
58
59impl DaemonFanOutSpawner {
60    /// A fan-out worker's advertised MCP defs: the global servers' defs plus its
61    /// blueprint's already-cached servers. Any uncached server is warmed on a
62    /// detached task (via the current runtime handle - `spawn_worker` runs inside
63    /// the daemon's tick, on the runtime) so a subsequent worker of the same type
64    /// advertises it. Returns just the global defs when the manifest is unreadable
65    /// or declares no servers.
66    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        // Warm any not-yet-connected servers for the next worker of this type, on
77        // a detached task (we run inside the tick, on the runtime).
78        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        // The parent supplies the worker's workdir, run id (parentage), and (for
93        // `worker_stage`) its blueprint path.
94        // `unattended` rides along: a worker of an unattended parent is a worker
95        // nobody is watching either, and one that stops on an approval prompt
96        // parks the parent behind it for good.
97        // The requested output shape rides along too: a caller who asked the
98        // parent for a2ui wants its workers' contributions in the same shape,
99        // and the worker's answer is what the merge stage reads.
100        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            // Fan-out workers get their split of the parent task via `task`.
127            regions: std::collections::HashMap::new(),
128            // Workers share the parent's workdir and are splits of a task the
129            // parent already scoped, so re-running a repo-scan command seed once
130            // per worker would be pure waste (and up to `max_workers` copies of
131            // the same output).
132            no_seed_commands: true,
133            output_request,
134        })
135        .map_err(|e| format!("resolve worker blueprint: {e}"))?;
136        // Nest the worker under its fan-out parent in the run tree.
137        args.parent_run_id = Some(parent_run_id);
138
139        // Per-agent MCP (issue #97): advertise the worker blueprint's servers that
140        // are already connected in the shared pool (a `worker_stage` worker shares
141        // the parent's - already warmed by the parent's preprocessor; the first
142        // `worker_agent`/`worker_query` worker warms them here for its siblings).
143        let mcp_defs = self.worker_mcp_defs(&args.blueprint_path);
144        // The worker holds its blueprint's per-agent servers open like any
145        // other run; the reap hook releases the lease when the worker ends.
146        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        // `worker_stage` runs the same blueprint entered at that stage rather than
165        // its default entry stage.
166        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
183/// Resolve a fan-out config's worker source to `(path_to_resolve, entry_stage)`.
184/// `path_to_resolve` is fed to [`resolve_spawn_args`] (which resolves a file,
185/// directory, or installed-agent name); `entry_stage` is `Some` only for
186/// `worker_stage` (self-as-worker).
187pub(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
205/// The task text seeded into a worker: its id plus the compact JSON context.
206fn format_worker_task(item_id: &str, item_context: &serde_json::Value) -> String {
207    format!("Work item id: {item_id}\nContext: {item_context}")
208}
209
210/// Find an installed agent whose directory name or manifest description contains
211/// `query` (case-insensitive). Returns the agent's directory.
212fn 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        // worker_stage → parent path + entry stage.
268        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        // worker_agent → the agent name, no entry stage.
273        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        // A worker_query with no match propagates discover_worker's error.
295        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        // An agent whose description (not name) matches.
311        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        // A directory without a manifest is skipped.
319        std::fs::create_dir_all(dir.path().join("not-an-agent")).unwrap();
320
321        // Matches by description.
322        assert!(
323            discover_worker(Some(dir.path()), "widget")
324                .unwrap()
325                .ends_with("alpha")
326        );
327        // Matches by directory name.
328        assert!(
329            discover_worker(Some(dir.path()), "ALPHA")
330                .unwrap()
331                .ends_with("alpha")
332        );
333        // No match.
334        assert!(discover_worker(Some(dir.path()), "nonexistent").is_err());
335        // No agents dir.
336        assert!(discover_worker(None, "x").is_err());
337        // Unreadable dir (path is a file).
338        let file = dir.path().join("alpha").join("agent.leviath");
339        assert!(discover_worker(Some(&file), "x").is_err());
340    }
341
342    // ── spawn_worker (integration over build_agent) ───────────────────────────
343
344    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    /// A two-stage blueprint whose second stage opts in as a fan-out worker.
376    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        // A worker blueprint declaring an MCP server.
410        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        // Seed the pool so the server's tool is already advertised (and the
418        // detached warm task hits the cache - no real connection).
419        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        // An unreadable manifest → just the global defs (read-error arm).
434        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        // A manifest with no [[mcp_servers]] → just the global defs.
438        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    /// Build a world with a live parent agent (from `manifest`) and return the
444    /// world, the spawner, and the parent entity.
445    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        // The worker entered the `second` stage (index 1).
516        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    /// A worker of an unattended parent is unattended. Spawning workers attended
524    /// under a `--yolo` parent left them stopping on approval prompts nobody was
525    /// watching for, with the parent parked behind them (issue #184).
526    #[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        // worker_agent given as a directory containing agent.leviath.
564        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        // A separate blueprint enters at its own entry stage (index 0).
577        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        // A bare entity with no RunMetadata.
587        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        // A worker blueprint that parses but fails validation (transition to a
626        // stage that doesn't exist) - resolve succeeds, build_agent errors.
627        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            // The `task` region is not incidental: a fan-out hands each worker
632            // its item as the task, and a worker with nowhere to put one is
633            // refused before validation ever runs - which would make this test
634            // pass on the wrong error.
635            "[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        // A config with no worker source ⇒ resolve_worker_source errors.
660        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        // worker_agent pointing at a nonexistent blueprint.
709        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        // Name doesn't match and the manifest won't parse → skipped → miss.
728        assert!(discover_worker(Some(dir.path()), "zzz").is_err());
729        // But the directory name still matches even when the manifest is broken.
730        assert!(
731            discover_worker(Some(dir.path()), "broken")
732                .unwrap()
733                .ends_with("broken")
734        );
735    }
736}