Skip to main content

mlua_swarm/application/
task.rs

1//! `TaskApplication` — the `POST /v1/tasks` entry point.
2//!
3//! Input: `BlueprintRef` (Inline / Id) plus a `TaskSpec`. Output:
4//! `(CapToken, StepId, version)`. Once the Blueprint is resolved, the
5//! engine-side operations (`bind` + `attach` + `start_task`) are
6//! delegated to [`TaskLaunchService`].
7
8use super::semver_resolve::SemverResolveError;
9use super::Application;
10use crate::blueprint::store::{BlueprintId, BlueprintStore, BlueprintStoreError, BlueprintVersion};
11use crate::blueprint::Blueprint;
12use crate::core::config::CheckPolicy;
13use crate::core::ctx::OperatorKind;
14use crate::service::{
15    TaskInputSpec, TaskLaunchError, TaskLaunchInput, TaskLaunchOutput, TaskLaunchService,
16};
17use crate::store::run::RunContext;
18use crate::types::{CapToken, Role};
19use async_trait::async_trait;
20use serde::{Deserialize, Serialize};
21use serde_json::Value;
22use std::collections::HashMap;
23use std::sync::Arc;
24use std::time::Duration;
25use thiserror::Error;
26
27/// How a task entry says the Blueprint should be resolved.
28#[derive(Debug, Clone, Serialize, Deserialize)]
29#[serde(tag = "kind", rename_all = "snake_case")]
30pub enum BlueprintRef {
31    /// The Blueprint value is embedded directly in the request; no
32    /// store lookup happens.
33    Inline {
34        /// The Blueprint to run as-is.
35        value: Box<Blueprint>,
36    },
37    /// Resolve the Blueprint from the `BlueprintStore` by id.
38    Id {
39        /// The `BlueprintId` to look up in the store.
40        id: BlueprintId,
41        /// Which generation to pick; defaults to `Latest`.
42        #[serde(default)]
43        version: VersionSelector,
44    },
45}
46
47/// How to pick a generation — a `version` inside the store.
48#[derive(Debug, Clone, Default, Serialize, Deserialize)]
49#[serde(tag = "kind", rename_all = "snake_case")]
50pub enum VersionSelector {
51    /// Use the store's current head version.
52    #[default]
53    Latest,
54    /// Use one exact, previously-committed version.
55    Fixed {
56        /// The exact version to read.
57        value: BlueprintVersion,
58    },
59    /// Scan the store's history and pick the highest version whose
60    /// `BlueprintMetadata.version_label` satisfies `req`.
61    SemverReq {
62        /// The semver requirement every candidate label is matched
63        /// against.
64        req: semver::VersionReq,
65    },
66}
67
68/// Input to [`TaskApplication::handle`] — the `POST /v1/tasks` request
69/// body once decoded.
70#[derive(Debug, Clone)]
71pub struct TaskApplicationInput {
72    /// Accepts both Inline (a Blueprint value directly) and Id
73    /// (store fetch + a `VersionSelector`).
74    pub blueprint: BlueprintRef,
75    /// Caller-supplied id for the Operator that owns this run.
76    pub operator_id: String,
77    /// The Operator's role for this run.
78    pub role: Role,
79    /// How long the attached session is allowed to live.
80    pub ttl: Duration,
81    /// Initial `ctx` for flow.ir `eval`. Read by every `Step.in`.
82    pub init_ctx: Value,
83    /// "Runtime Global" tier of the `OperatorKind` cascade. `Some(_)` is
84    /// always an explicit request — including `Some(OperatorKind::Automate)`
85    /// — that outranks the BP-level tiers (`OperatorDef.kind` /
86    /// `Blueprint.default_operator_kind`); `None` leaves it unspecified so
87    /// those tiers / the final default decide. Under `MainAi` /
88    /// `Composite`, `MainAIMiddleware`'s `spawn_hook` before/after
89    /// callbacks become effective. See
90    /// `crate::core::ctx::collapse_operator_kind`.
91    pub operator_kind: Option<crate::core::ctx::OperatorKind>,
92    /// `SeniorBridge` registry ID. `None` — none in use;
93    /// `Some(id)` — attach a bridge previously registered on the
94    /// engine.
95    pub bridge_id: Option<String>,
96    /// `SpawnHook` registry ID. Same shape as above — attach a hook
97    /// previously registered on the engine.
98    pub hook_id: Option<String>,
99    /// Operator registry ID — used on the path that hands the whole
100    /// spawn off to an external Operator.
101    pub operator_backend_id: Option<String>,
102    /// "Runtime Agent-level" tier (highest priority) of the `OperatorKind`
103    /// cascade — per-agent override, keyed by `AgentDef.name`. Empty by
104    /// default. See `crate::core::ctx::collapse_operator_kind` for the full tier
105    /// list.
106    pub operator_kind_overrides: HashMap<String, OperatorKind>,
107    /// Task-level canonical execution context (issue #19 ST2). When
108    /// `Some`, the resolved sibling fields (`project_root` / `work_dir`
109    /// / `task_metadata`) are threaded down to [`TaskLaunchInput`] and
110    /// consumed by
111    /// [`crate::middleware::task_input::TaskInputMiddleware::new_from_fields`].
112    /// `None` — no Task-level context is layered on the spawner stack
113    /// (default; keeps the wire body opt-in).
114    pub task_input: Option<TaskInputSpec>,
115    /// The "launch request" tier (tier 1) of the
116    /// `check_policy` cascade, threaded straight down to
117    /// [`TaskLaunchInput::check_policy`]. `None` (the default via
118    /// [`Self::automate`]) leaves this tier unspecified — the Blueprint tier
119    /// / server-wide default decide. Wired from the `POST /v1/tasks`
120    /// request body's top-level `check_policy` field.
121    pub check_policy: Option<CheckPolicy>,
122}
123
124impl TaskApplicationInput {
125    /// Helper for existing callers on the default path — no hooks and no
126    /// per-agent `OperatorKind` overrides. Leaves the "Runtime Global" tier
127    /// unspecified (`None`), so the BP-level tiers / final default
128    /// (`OperatorKind::Automate`) decide — this preserves today's
129    /// behaviour for every existing caller without silently forcing
130    /// `Automate` as an explicit override that would outrank a BP-declared
131    /// `MainAi`/`Composite` kind.
132    pub fn automate(
133        blueprint: BlueprintRef,
134        operator_id: impl Into<String>,
135        role: Role,
136        ttl: Duration,
137        init_ctx: Value,
138    ) -> Self {
139        Self {
140            blueprint,
141            operator_id: operator_id.into(),
142            role,
143            ttl,
144            init_ctx,
145            operator_kind: None,
146            bridge_id: None,
147            hook_id: None,
148            operator_backend_id: None,
149            operator_kind_overrides: HashMap::new(),
150            task_input: None,
151            check_policy: None,
152        }
153    }
154}
155
156/// Result of a successful [`TaskApplication::handle`] call.
157#[derive(Debug, Clone)]
158pub struct TaskApplicationOutput {
159    /// The capability token for the attached session.
160    pub token: CapToken,
161    /// The final `ctx` after the flow ran to completion.
162    pub final_ctx: Value,
163    /// Only `Some` when resolution went through the store
164    /// (`BlueprintRef::Id`); `None` on the Inline path.
165    pub bound_version: Option<BlueprintVersion>,
166}
167
168/// Failure modes of [`TaskApplication::handle`] and
169/// [`TaskApplication::resolve`].
170#[derive(Debug, Error)]
171pub enum TaskApplicationError {
172    /// `BlueprintRef::Id` was used but this `TaskApplication` was
173    /// built via [`TaskApplication::new_inline_only`] (no store).
174    #[error("store not configured (BlueprintRef::Id requires store)")]
175    NoStore,
176    /// The `BlueprintStore` returned an error while resolving the ref.
177    #[error("store: {0}")]
178    Store(#[from] BlueprintStoreError),
179    /// `TaskLaunchService::launch` failed after resolution succeeded.
180    #[error("launch: {0}")]
181    Launch(#[from] TaskLaunchError),
182    /// A stored version's `version_label` is not valid semver.
183    #[error("invalid semver version_label {label:?}: {source}")]
184    InvalidSemver {
185        /// The offending label string.
186        label: String,
187        /// The underlying semver parse error.
188        #[source]
189        source: semver::Error,
190    },
191    /// No stored version's label satisfies the `SemverReq`.
192    #[error("no version matches semver req: {req}")]
193    NoMatchingVersion {
194        /// The requirement string that matched nothing.
195        req: String,
196    },
197}
198
199impl From<SemverResolveError> for TaskApplicationError {
200    fn from(e: SemverResolveError) -> Self {
201        match e {
202            SemverResolveError::Store(e) => TaskApplicationError::Store(e),
203            SemverResolveError::InvalidSemver { label, source } => {
204                TaskApplicationError::InvalidSemver { label, source }
205            }
206            SemverResolveError::NoMatchingVersion { req } => {
207                TaskApplicationError::NoMatchingVersion { req }
208            }
209        }
210    }
211}
212
213/// The `POST /v1/tasks` [`Application`] — resolves a `BlueprintRef` and
214/// runs it to completion through [`TaskLaunchService`].
215pub struct TaskApplication {
216    launch: Arc<TaskLaunchService>,
217    /// Only needed when resolving `BlueprintRef::Id`; `None` in
218    /// Inline-only mode.
219    store: Option<Arc<dyn BlueprintStore>>,
220}
221
222impl TaskApplication {
223    /// Build a `TaskApplication` that can resolve both `Inline` and
224    /// `Id` `BlueprintRef`s (the `Id` path reads through `store`).
225    pub fn new(launch: Arc<TaskLaunchService>, store: Arc<dyn BlueprintStore>) -> Self {
226        Self {
227            launch,
228            store: Some(store),
229        }
230    }
231
232    /// Build a `TaskApplication` restricted to `Inline` `BlueprintRef`s
233    /// — no store is configured, so `Id` resolution always fails with
234    /// `TaskApplicationError::NoStore`.
235    pub fn new_inline_only(launch: Arc<TaskLaunchService>) -> Self {
236        Self {
237            launch,
238            store: None,
239        }
240    }
241
242    /// Resolve a `BlueprintRef` and return the real Blueprint plus,
243    /// when it went through the store, the resolved version.
244    pub async fn resolve(
245        &self,
246        bp_ref: &BlueprintRef,
247    ) -> Result<(Blueprint, Option<BlueprintVersion>), TaskApplicationError> {
248        match bp_ref {
249            BlueprintRef::Inline { value } => Ok((value.as_ref().clone(), None)),
250            BlueprintRef::Id { id, version } => {
251                let store = self.store.as_ref().ok_or(TaskApplicationError::NoStore)?;
252                let bp_id = id.clone();
253                let traced = match version {
254                    VersionSelector::Latest => store.read_head(&bp_id).await?,
255                    VersionSelector::Fixed { value } => store.read_version(&bp_id, *value).await?,
256                    VersionSelector::SemverReq { req } => {
257                        let v = super::semver_resolve::resolve_semver(store.as_ref(), &bp_id, req)
258                            .await?;
259                        store.read_version(&bp_id, v).await?
260                    }
261                };
262                let ver = traced.trace.version;
263                Ok((traced.value, Some(ver)))
264            }
265        }
266    }
267
268    /// Resolve the `BlueprintRef` (Inline / Id) and run the flow to
269    /// completion through `TaskLaunchService::launch`, threading `run_ctx`
270    /// (issue #13 run_id propagation) into the launch input.
271    ///
272    /// [`Application::handle`] delegates here with `run_ctx: None` — a
273    /// separate method rather than a new field on [`TaskApplicationInput`]
274    /// so the pre-existing exhaustive `TaskApplicationInput { .. }` struct
275    /// literal in `mlua-swarm-cli`'s MCP adapter (which has no `run_ctx`)
276    /// keeps compiling unchanged. Server entry points that mint a `RunId`
277    /// up front (`POST /v1/tasks`, `POST /v1/tasks/:id/runs`) call this
278    /// directly with `Some(run_ctx)`.
279    pub async fn handle_with_run(
280        &self,
281        input: TaskApplicationInput,
282        run_ctx: Option<RunContext>,
283    ) -> Result<TaskApplicationOutput, TaskApplicationError> {
284        let (blueprint, bound_version) = self.resolve(&input.blueprint).await?;
285        let TaskLaunchOutput { token, final_ctx } = self
286            .launch
287            .launch(TaskLaunchInput {
288                blueprint,
289                operator_id: input.operator_id,
290                role: input.role,
291                ttl: input.ttl,
292                operator_kind: input.operator_kind,
293                bridge_id: input.bridge_id,
294                hook_id: input.hook_id,
295                operator_backend_id: input.operator_backend_id,
296                operator_kind_overrides: input.operator_kind_overrides,
297                init_ctx: input.init_ctx,
298                run_ctx,
299                task_input: input.task_input,
300                check_policy: input.check_policy,
301            })
302            .await?;
303        Ok(TaskApplicationOutput {
304            token,
305            final_ctx,
306            bound_version,
307        })
308    }
309}
310
311#[async_trait]
312impl Application for TaskApplication {
313    type Input = TaskApplicationInput;
314    type Output = TaskApplicationOutput;
315    type Error = TaskApplicationError;
316
317    fn name(&self) -> &str {
318        "task"
319    }
320
321    /// Resolve the `BlueprintRef` (Inline / Id) and run the flow to
322    /// completion through `TaskLaunchService::launch`. Delegates to
323    /// [`TaskApplication::handle_with_run`] with `run_ctx: None` (no run
324    /// tracing) — callers that need `RunRecord.step_entries` tracing call
325    /// `handle_with_run` directly instead.
326    async fn handle(&self, input: Self::Input) -> Result<Self::Output, Self::Error> {
327        self.handle_with_run(input, None).await
328    }
329}
330
331// ──────────────────────────────────────────────────────────────────────────
332// UT
333// ──────────────────────────────────────────────────────────────────────────
334
335#[cfg(test)]
336mod tests {
337    use super::*;
338    use crate::blueprint::compiler::{Compiler, SpawnerRegistry};
339    use crate::blueprint::store::{
340        blueprint_version, BlueprintId, BlueprintStore, BlueprintStoreError, CommitMetadata,
341        InMemoryBlueprintStore,
342    };
343    use crate::blueprint::{
344        current_schema_version, AgentKind, Blueprint, BlueprintMetadata, CompilerHints,
345        CompilerStrategy,
346    };
347    use crate::core::config::EngineCfg;
348    use crate::core::ctx::OperatorKind;
349    use crate::core::engine::Engine;
350    use mlua_flow_ir::Node as FlowNode;
351
352    fn empty_bp() -> Blueprint {
353        Blueprint {
354            schema_version: current_schema_version(),
355            id: "ut-bp".into(),
356            flow: FlowNode::Seq { children: vec![] },
357            agents: vec![],
358            operators: vec![],
359            metas: vec![],
360            hints: CompilerHints::default(),
361            strategy: CompilerStrategy::default(),
362            metadata: BlueprintMetadata::default(),
363            spawner_hints: Default::default(),
364            default_agent_kind: AgentKind::Operator,
365            default_operator_kind: None,
366            default_init_ctx: None,
367            default_agent_ctx: None,
368            default_context_policy: None,
369            projection_placement: None,
370            audits: vec![],
371            degradation_policy: None,
372            runners: vec![],
373            default_runner: None,
374            check_policy: None,
375        }
376    }
377
378    fn bp_with_label(id: &str, version_label: Option<&str>) -> Blueprint {
379        Blueprint {
380            schema_version: current_schema_version(),
381            id: id.into(),
382            flow: FlowNode::Seq { children: vec![] },
383            agents: vec![],
384            operators: vec![],
385            metas: vec![],
386            hints: CompilerHints::default(),
387            strategy: CompilerStrategy::default(),
388            metadata: BlueprintMetadata {
389                description: None,
390                origin: Default::default(),
391                tags: vec![],
392                version_label: version_label.map(|s| s.to_string()),
393                project_name_alias: None,
394                default_run_ttl_secs: None,
395                strict_verdict_handling: None,
396            },
397            spawner_hints: Default::default(),
398            default_agent_kind: AgentKind::Operator,
399            default_operator_kind: None,
400            default_init_ctx: None,
401            default_agent_ctx: None,
402            default_context_policy: None,
403            projection_placement: None,
404            audits: vec![],
405            degradation_policy: None,
406            runners: vec![],
407            default_runner: None,
408            check_policy: None,
409        }
410    }
411
412    fn build_app_with_store() -> (TaskApplication, Arc<dyn BlueprintStore>) {
413        let reg = SpawnerRegistry::new();
414        let compiler = Compiler::new(reg);
415        let engine = Engine::new(EngineCfg::default());
416        let launch = Arc::new(TaskLaunchService::new(engine, compiler));
417        let store: Arc<dyn BlueprintStore> = Arc::new(InMemoryBlueprintStore::new());
418        (TaskApplication::new(launch, store.clone()), store)
419    }
420
421    fn build_app_inline_only() -> TaskApplication {
422        let reg = SpawnerRegistry::new();
423        let compiler = Compiler::new(reg);
424        let engine = Engine::new(EngineCfg::default());
425        let launch = Arc::new(TaskLaunchService::new(engine, compiler));
426        TaskApplication::new_inline_only(launch)
427    }
428
429    async fn seed(store: &Arc<dyn BlueprintStore>, bp: &Blueprint) -> BlueprintVersion {
430        let id = bp.id.clone();
431        let v = blueprint_version(bp).expect("hash");
432        store
433            .write_new(&id, bp, &[], CommitMetadata::seed(id.clone(), v, 0))
434            .await
435            .expect("seed");
436        v
437    }
438
439    #[test]
440    fn automate_helper_sets_defaults() {
441        let input = TaskApplicationInput::automate(
442            BlueprintRef::Inline {
443                value: Box::new(empty_bp()),
444            },
445            "op-1",
446            Role::Operator,
447            Duration::from_secs(10),
448            serde_json::json!({}),
449        );
450        assert!(
451            input.operator_kind.is_none(),
452            "automate() leaves the Runtime Global tier unspecified (None), \
453             not an explicit Some(Automate) override"
454        );
455        assert!(input.bridge_id.is_none());
456        assert!(input.hook_id.is_none());
457        assert_eq!(input.operator_id, "op-1");
458    }
459
460    #[test]
461    fn struct_literal_allows_callback_ids() {
462        let input = TaskApplicationInput {
463            blueprint: BlueprintRef::Inline {
464                value: Box::new(empty_bp()),
465            },
466            operator_id: "op-2".into(),
467            role: Role::Operator,
468            ttl: Duration::from_secs(5),
469            init_ctx: serde_json::json!({}),
470            operator_kind: Some(OperatorKind::MainAi),
471            bridge_id: Some("br-x".into()),
472            hook_id: Some("hk-y".into()),
473            operator_backend_id: None,
474            operator_kind_overrides: HashMap::new(),
475            task_input: None,
476            check_policy: None,
477        };
478        assert!(matches!(input.operator_kind, Some(OperatorKind::MainAi)));
479        assert_eq!(input.bridge_id.as_deref(), Some("br-x"));
480        assert_eq!(input.hook_id.as_deref(), Some("hk-y"));
481    }
482
483    // ──────────────────────────────────────────────────────────────────
484    // resolve / resolve_semver carve
485    // ──────────────────────────────────────────────────────────────────
486
487    #[tokio::test]
488    async fn resolve_inline_returns_bp_and_no_version() {
489        let app = build_app_inline_only();
490        let bp = empty_bp();
491        let (got, ver) = app
492            .resolve(&BlueprintRef::Inline {
493                value: Box::new(bp.clone()),
494            })
495            .await
496            .expect("resolve inline ok");
497        assert_eq!(got.id, bp.id);
498        assert!(ver.is_none(), "the Inline path yields bound_version=None");
499    }
500
501    #[tokio::test]
502    async fn resolve_id_latest_returns_bp_and_version() {
503        let (app, store) = build_app_with_store();
504        let bp = bp_with_label("rid-latest", Some("0.1.0"));
505        let v = seed(&store, &bp).await;
506        let (got, ver) = app
507            .resolve(&BlueprintRef::Id {
508                id: bp.id.clone(),
509                version: VersionSelector::Latest,
510            })
511            .await
512            .expect("resolve id latest ok");
513        assert_eq!(got.id, bp.id);
514        assert_eq!(ver, Some(v), "Latest = seed version");
515    }
516
517    #[tokio::test]
518    async fn resolve_id_fixed_picks_exact_version() {
519        let (app, store) = build_app_with_store();
520        let id = "rid-fixed";
521        let bp1 = bp_with_label(id, Some("1.0.0"));
522        let bp2 = bp_with_label(id, Some("2.0.0"));
523        let v1 = seed(&store, &bp1).await;
524        let _v2 = seed(&store, &bp2).await;
525        let (got, ver) = app
526            .resolve(&BlueprintRef::Id {
527                id: BlueprintId::new(id),
528                version: VersionSelector::Fixed { value: v1 },
529            })
530            .await
531            .expect("resolve id fixed ok");
532        assert_eq!(ver, Some(v1));
533        assert_eq!(
534            got.metadata.version_label.as_deref(),
535            Some("1.0.0"),
536            "Fixed{{v1}} resolves to v1 = 1.0.0"
537        );
538    }
539
540    #[tokio::test]
541    async fn resolve_id_semver_picks_highest_matching() {
542        let (app, store) = build_app_with_store();
543        let id = "rid-semver";
544        let _ = seed(&store, &bp_with_label(id, Some("1.0.0"))).await;
545        let _ = seed(&store, &bp_with_label(id, Some("1.2.0"))).await;
546        let _ = seed(&store, &bp_with_label(id, Some("2.0.0"))).await;
547        let req = semver::VersionReq::parse("^1").expect("req");
548        let (got, ver) = app
549            .resolve(&BlueprintRef::Id {
550                id: BlueprintId::new(id),
551                version: VersionSelector::SemverReq { req },
552            })
553            .await
554            .expect("resolve semver ok");
555        assert!(ver.is_some());
556        assert_eq!(
557            got.metadata.version_label.as_deref(),
558            Some("1.2.0"),
559            "^1 max = 1.2.0 (2.0.0 is out of range; 1.0.0 is lower)"
560        );
561    }
562
563    #[tokio::test]
564    async fn resolve_id_semver_no_match_errs() {
565        let (app, store) = build_app_with_store();
566        let id = "rid-semver-nomatch";
567        let _ = seed(&store, &bp_with_label(id, Some("1.0.0"))).await;
568        let req = semver::VersionReq::parse("^3").expect("req");
569        let err = app
570            .resolve(&BlueprintRef::Id {
571                id: BlueprintId::new(id),
572                version: VersionSelector::SemverReq { req },
573            })
574            .await
575            .expect_err("expected NoMatchingVersion");
576        match err {
577            TaskApplicationError::NoMatchingVersion { req } => {
578                assert!(req.contains("^3"), "req string carry: {req}");
579            }
580            other => panic!("expected NoMatchingVersion, got {other:?}"),
581        }
582    }
583
584    #[tokio::test]
585    async fn resolve_id_semver_invalid_label_errs() {
586        let (app, store) = build_app_with_store();
587        let id = "rid-semver-bad";
588        let _ = seed(&store, &bp_with_label(id, Some("not-semver"))).await;
589        let req = semver::VersionReq::parse("^1").expect("req");
590        let err = app
591            .resolve(&BlueprintRef::Id {
592                id: BlueprintId::new(id),
593                version: VersionSelector::SemverReq { req },
594            })
595            .await
596            .expect_err("expected InvalidSemver");
597        match err {
598            TaskApplicationError::InvalidSemver { label, .. } => {
599                assert_eq!(label, "not-semver");
600            }
601            other => panic!("expected InvalidSemver, got {other:?}"),
602        }
603    }
604
605    #[tokio::test]
606    async fn resolve_id_without_store_errs_no_store() {
607        let app = build_app_inline_only();
608        let err = app
609            .resolve(&BlueprintRef::Id {
610                id: BlueprintId::new("anything"),
611                version: VersionSelector::Latest,
612            })
613            .await
614            .expect_err("expected NoStore");
615        assert!(matches!(err, TaskApplicationError::NoStore), "got {err:?}");
616    }
617
618    #[tokio::test]
619    async fn resolve_id_not_found_errs_store() {
620        let (app, _store) = build_app_with_store();
621        let err = app
622            .resolve(&BlueprintRef::Id {
623                id: BlueprintId::new("never-seeded"),
624                version: VersionSelector::Latest,
625            })
626            .await
627            .expect_err("expected Store(IdNotFound|HeadEmpty)");
628        match err {
629            TaskApplicationError::Store(
630                BlueprintStoreError::IdNotFound(_) | BlueprintStoreError::HeadEmpty(_),
631            ) => {}
632            other => panic!("expected Store(IdNotFound|HeadEmpty), got {other:?}"),
633        }
634    }
635}