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 blueprint_ref_includes: Vec::new(),
376 }
377 }
378
379 fn bp_with_label(id: &str, version_label: Option<&str>) -> Blueprint {
380 Blueprint {
381 schema_version: current_schema_version(),
382 id: id.into(),
383 flow: FlowNode::Seq { children: vec![] },
384 agents: vec![],
385 operators: vec![],
386 metas: vec![],
387 hints: CompilerHints::default(),
388 strategy: CompilerStrategy::default(),
389 metadata: BlueprintMetadata {
390 description: None,
391 origin: Default::default(),
392 tags: vec![],
393 version_label: version_label.map(|s| s.to_string()),
394 project_name_alias: None,
395 default_run_ttl_secs: None,
396 strict_verdict_handling: None,
397 },
398 spawner_hints: Default::default(),
399 default_agent_kind: AgentKind::Operator,
400 default_operator_kind: None,
401 default_init_ctx: None,
402 default_agent_ctx: None,
403 default_context_policy: None,
404 projection_placement: None,
405 audits: vec![],
406 degradation_policy: None,
407 runners: vec![],
408 default_runner: None,
409 check_policy: None,
410 blueprint_ref_includes: Vec::new(),
411 }
412 }
413
414 fn build_app_with_store() -> (TaskApplication, Arc<dyn BlueprintStore>) {
415 let reg = SpawnerRegistry::new();
416 let compiler = Compiler::new(reg);
417 let engine = Engine::new(EngineCfg::default());
418 let launch = Arc::new(TaskLaunchService::new(engine, compiler));
419 let store: Arc<dyn BlueprintStore> = Arc::new(InMemoryBlueprintStore::new());
420 (TaskApplication::new(launch, store.clone()), store)
421 }
422
423 fn build_app_inline_only() -> TaskApplication {
424 let reg = SpawnerRegistry::new();
425 let compiler = Compiler::new(reg);
426 let engine = Engine::new(EngineCfg::default());
427 let launch = Arc::new(TaskLaunchService::new(engine, compiler));
428 TaskApplication::new_inline_only(launch)
429 }
430
431 async fn seed(store: &Arc<dyn BlueprintStore>, bp: &Blueprint) -> BlueprintVersion {
432 let id = bp.id.clone();
433 let v = blueprint_version(bp).expect("hash");
434 store
435 .write_new(&id, bp, &[], CommitMetadata::seed(id.clone(), v, 0))
436 .await
437 .expect("seed");
438 v
439 }
440
441 #[test]
442 fn automate_helper_sets_defaults() {
443 let input = TaskApplicationInput::automate(
444 BlueprintRef::Inline {
445 value: Box::new(empty_bp()),
446 },
447 "op-1",
448 Role::Operator,
449 Duration::from_secs(10),
450 serde_json::json!({}),
451 );
452 assert!(
453 input.operator_kind.is_none(),
454 "automate() leaves the Runtime Global tier unspecified (None), \
455 not an explicit Some(Automate) override"
456 );
457 assert!(input.bridge_id.is_none());
458 assert!(input.hook_id.is_none());
459 assert_eq!(input.operator_id, "op-1");
460 }
461
462 #[test]
463 fn struct_literal_allows_callback_ids() {
464 let input = TaskApplicationInput {
465 blueprint: BlueprintRef::Inline {
466 value: Box::new(empty_bp()),
467 },
468 operator_id: "op-2".into(),
469 role: Role::Operator,
470 ttl: Duration::from_secs(5),
471 init_ctx: serde_json::json!({}),
472 operator_kind: Some(OperatorKind::MainAi),
473 bridge_id: Some("br-x".into()),
474 hook_id: Some("hk-y".into()),
475 operator_backend_id: None,
476 operator_kind_overrides: HashMap::new(),
477 task_input: None,
478 check_policy: None,
479 };
480 assert!(matches!(input.operator_kind, Some(OperatorKind::MainAi)));
481 assert_eq!(input.bridge_id.as_deref(), Some("br-x"));
482 assert_eq!(input.hook_id.as_deref(), Some("hk-y"));
483 }
484
485 #[tokio::test]
490 async fn resolve_inline_returns_bp_and_no_version() {
491 let app = build_app_inline_only();
492 let bp = empty_bp();
493 let (got, ver) = app
494 .resolve(&BlueprintRef::Inline {
495 value: Box::new(bp.clone()),
496 })
497 .await
498 .expect("resolve inline ok");
499 assert_eq!(got.id, bp.id);
500 assert!(ver.is_none(), "the Inline path yields bound_version=None");
501 }
502
503 #[tokio::test]
504 async fn resolve_id_latest_returns_bp_and_version() {
505 let (app, store) = build_app_with_store();
506 let bp = bp_with_label("rid-latest", Some("0.1.0"));
507 let v = seed(&store, &bp).await;
508 let (got, ver) = app
509 .resolve(&BlueprintRef::Id {
510 id: bp.id.clone(),
511 version: VersionSelector::Latest,
512 })
513 .await
514 .expect("resolve id latest ok");
515 assert_eq!(got.id, bp.id);
516 assert_eq!(ver, Some(v), "Latest = seed version");
517 }
518
519 #[tokio::test]
520 async fn resolve_id_fixed_picks_exact_version() {
521 let (app, store) = build_app_with_store();
522 let id = "rid-fixed";
523 let bp1 = bp_with_label(id, Some("1.0.0"));
524 let bp2 = bp_with_label(id, Some("2.0.0"));
525 let v1 = seed(&store, &bp1).await;
526 let _v2 = seed(&store, &bp2).await;
527 let (got, ver) = app
528 .resolve(&BlueprintRef::Id {
529 id: BlueprintId::new(id),
530 version: VersionSelector::Fixed { value: v1 },
531 })
532 .await
533 .expect("resolve id fixed ok");
534 assert_eq!(ver, Some(v1));
535 assert_eq!(
536 got.metadata.version_label.as_deref(),
537 Some("1.0.0"),
538 "Fixed{{v1}} resolves to v1 = 1.0.0"
539 );
540 }
541
542 #[tokio::test]
543 async fn resolve_id_semver_picks_highest_matching() {
544 let (app, store) = build_app_with_store();
545 let id = "rid-semver";
546 let _ = seed(&store, &bp_with_label(id, Some("1.0.0"))).await;
547 let _ = seed(&store, &bp_with_label(id, Some("1.2.0"))).await;
548 let _ = seed(&store, &bp_with_label(id, Some("2.0.0"))).await;
549 let req = semver::VersionReq::parse("^1").expect("req");
550 let (got, ver) = app
551 .resolve(&BlueprintRef::Id {
552 id: BlueprintId::new(id),
553 version: VersionSelector::SemverReq { req },
554 })
555 .await
556 .expect("resolve semver ok");
557 assert!(ver.is_some());
558 assert_eq!(
559 got.metadata.version_label.as_deref(),
560 Some("1.2.0"),
561 "^1 max = 1.2.0 (2.0.0 is out of range; 1.0.0 is lower)"
562 );
563 }
564
565 #[tokio::test]
566 async fn resolve_id_semver_no_match_errs() {
567 let (app, store) = build_app_with_store();
568 let id = "rid-semver-nomatch";
569 let _ = seed(&store, &bp_with_label(id, Some("1.0.0"))).await;
570 let req = semver::VersionReq::parse("^3").expect("req");
571 let err = app
572 .resolve(&BlueprintRef::Id {
573 id: BlueprintId::new(id),
574 version: VersionSelector::SemverReq { req },
575 })
576 .await
577 .expect_err("expected NoMatchingVersion");
578 match err {
579 TaskApplicationError::NoMatchingVersion { req } => {
580 assert!(req.contains("^3"), "req string carry: {req}");
581 }
582 other => panic!("expected NoMatchingVersion, got {other:?}"),
583 }
584 }
585
586 #[tokio::test]
587 async fn resolve_id_semver_invalid_label_errs() {
588 let (app, store) = build_app_with_store();
589 let id = "rid-semver-bad";
590 let _ = seed(&store, &bp_with_label(id, Some("not-semver"))).await;
591 let req = semver::VersionReq::parse("^1").expect("req");
592 let err = app
593 .resolve(&BlueprintRef::Id {
594 id: BlueprintId::new(id),
595 version: VersionSelector::SemverReq { req },
596 })
597 .await
598 .expect_err("expected InvalidSemver");
599 match err {
600 TaskApplicationError::InvalidSemver { label, .. } => {
601 assert_eq!(label, "not-semver");
602 }
603 other => panic!("expected InvalidSemver, got {other:?}"),
604 }
605 }
606
607 #[tokio::test]
608 async fn resolve_id_without_store_errs_no_store() {
609 let app = build_app_inline_only();
610 let err = app
611 .resolve(&BlueprintRef::Id {
612 id: BlueprintId::new("anything"),
613 version: VersionSelector::Latest,
614 })
615 .await
616 .expect_err("expected NoStore");
617 assert!(matches!(err, TaskApplicationError::NoStore), "got {err:?}");
618 }
619
620 #[tokio::test]
621 async fn resolve_id_not_found_errs_store() {
622 let (app, _store) = build_app_with_store();
623 let err = app
624 .resolve(&BlueprintRef::Id {
625 id: BlueprintId::new("never-seeded"),
626 version: VersionSelector::Latest,
627 })
628 .await
629 .expect_err("expected Store(IdNotFound|HeadEmpty)");
630 match err {
631 TaskApplicationError::Store(
632 BlueprintStoreError::IdNotFound(_) | BlueprintStoreError::HeadEmpty(_),
633 ) => {}
634 other => panic!("expected Store(IdNotFound|HeadEmpty), got {other:?}"),
635 }
636 }
637}