1use 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#[derive(Debug, Clone, Serialize, Deserialize)]
29#[serde(tag = "kind", rename_all = "snake_case")]
30pub enum BlueprintRef {
31 Inline {
34 value: Box<Blueprint>,
36 },
37 Id {
39 id: BlueprintId,
41 #[serde(default)]
43 version: VersionSelector,
44 },
45}
46
47#[derive(Debug, Clone, Default, Serialize, Deserialize)]
49#[serde(tag = "kind", rename_all = "snake_case")]
50pub enum VersionSelector {
51 #[default]
53 Latest,
54 Fixed {
56 value: BlueprintVersion,
58 },
59 SemverReq {
62 req: semver::VersionReq,
65 },
66}
67
68#[derive(Debug, Clone)]
71pub struct TaskApplicationInput {
72 pub blueprint: BlueprintRef,
75 pub operator_id: String,
77 pub role: Role,
79 pub ttl: Duration,
81 pub init_ctx: Value,
83 pub operator_kind: Option<crate::core::ctx::OperatorKind>,
92 pub bridge_id: Option<String>,
96 pub hook_id: Option<String>,
99 pub operator_backend_id: Option<String>,
102 pub operator_kind_overrides: HashMap<String, OperatorKind>,
107 pub task_input: Option<TaskInputSpec>,
115 pub check_policy: Option<CheckPolicy>,
122}
123
124impl TaskApplicationInput {
125 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#[derive(Debug, Clone)]
158pub struct TaskApplicationOutput {
159 pub token: CapToken,
161 pub final_ctx: Value,
163 pub bound_version: Option<BlueprintVersion>,
166}
167
168#[derive(Debug, Error)]
171pub enum TaskApplicationError {
172 #[error("store not configured (BlueprintRef::Id requires store)")]
175 NoStore,
176 #[error("store: {0}")]
178 Store(#[from] BlueprintStoreError),
179 #[error("launch: {0}")]
181 Launch(#[from] TaskLaunchError),
182 #[error("invalid semver version_label {label:?}: {source}")]
184 InvalidSemver {
185 label: String,
187 #[source]
189 source: semver::Error,
190 },
191 #[error("no version matches semver req: {req}")]
193 NoMatchingVersion {
194 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
213pub struct TaskApplication {
216 launch: Arc<TaskLaunchService>,
217 store: Option<Arc<dyn BlueprintStore>>,
220}
221
222impl TaskApplication {
223 pub fn new(launch: Arc<TaskLaunchService>, store: Arc<dyn BlueprintStore>) -> Self {
226 Self {
227 launch,
228 store: Some(store),
229 }
230 }
231
232 pub fn new_inline_only(launch: Arc<TaskLaunchService>) -> Self {
236 Self {
237 launch,
238 store: None,
239 }
240 }
241
242 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 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 async fn handle(&self, input: Self::Input) -> Result<Self::Output, Self::Error> {
327 self.handle_with_run(input, None).await
328 }
329}
330
331#[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 #[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}