1use crate::binding::{
23 attest_bound_agents, binding_requests, validate_bound_agent_snapshots, AgentBindingProvider,
24 LegacyWorkerBindingPolicy, UnboundAgent,
25};
26use crate::blueprint::compiler::{materialize_bound_blueprint, CompileError, Compiler};
27use crate::blueprint::{
28 resolve_bound_agents, AuditDef, Blueprint, BoundAgent, EngineDispatcher, Runner,
29};
30use crate::core::agent_context::ContextPolicy;
31use crate::core::config::CheckPolicy;
32use crate::core::ctx::OperatorKind;
33use crate::core::engine::Engine;
34use crate::core::errors::EngineError;
35use crate::middleware::agent_context::AgentContextMiddleware;
36use crate::middleware::project_name_alias::ProjectNameAliasMiddleware;
37use crate::middleware::task_input::TaskInputMiddleware;
38use crate::middleware::worker_binding::WorkerBindingMiddleware;
39use crate::middleware::{AfterRunAuditMiddleware, SpawnerStack};
40use crate::operator::WorkerBinding;
41use crate::service::linker;
42use crate::store::run::{RunContext, SnapshotOrigin};
43use crate::types::{CapToken, Role};
44use mlua_flow_ir::{Externs, NoExterns};
45use serde::{Deserialize, Serialize};
46use serde_json::Value;
47use std::collections::HashMap;
48use std::sync::Arc;
49use std::time::Duration;
50use thiserror::Error;
51
52#[cfg(test)]
79pub(crate) fn derive_worker_bindings(blueprint: &Blueprint) -> HashMap<String, WorkerBinding> {
80 let bound_agents = resolve_bound_agents(blueprint)
83 .expect("derive_worker_bindings requires a Blueprint with resolvable Runner refs");
84 worker_bindings_from_bound_agents(&bound_agents)
85}
86
87fn worker_bindings_from_bound_agents(
88 bound_agents: &[BoundAgent],
89) -> HashMap<String, WorkerBinding> {
90 bound_agents
91 .iter()
92 .filter_map(|bound| match &bound.runner {
93 Some(Runner::WsOperator { variant, tools })
94 | Some(Runner::WsClaudeCode { variant, tools }) => Some((
95 bound.agent.name.clone(),
96 WorkerBinding {
97 variant: variant.clone(),
98 tools: tools.clone(),
99 request_digest: Some(bound.binding_digest.clone()),
100 requested_model: bound.agent.profile.as_ref().and_then(|p| p.model.clone()),
101 },
102 )),
103 _ => None,
104 })
105 .collect()
106}
107
108async fn attest_or_gate_fresh(
123 bound_agents: &mut [BoundAgent],
124 binding_provider: Option<&dyn AgentBindingProvider>,
125 strict: bool,
126 run_ctx: Option<&RunContext>,
127) -> Result<(), TaskLaunchError> {
128 match binding_provider {
129 Some(provider) => {
130 let unbound = attest_bound_agents(provider, bound_agents, strict)
131 .await
132 .map_err(|error| TaskLaunchError::PreDispatch(error.to_string()))?;
133 for agent in &unbound {
134 record_unbound_degradation(agent, run_ctx).await;
135 }
136 Ok(())
137 }
138 None => {
139 if strict && !binding_requests(bound_agents).is_empty() {
140 return Err(TaskLaunchError::PreDispatch(format!(
141 "strict_binding requires a binding provider but none is injected; \
142 {} Runner-backed agent(s) cannot be attested",
143 binding_requests(bound_agents).len()
144 )));
145 }
146 Ok(())
147 }
148 }
149}
150
151async fn record_unbound_degradation(agent: &UnboundAgent, run_ctx: Option<&RunContext>) {
157 tracing::warn!(
158 agent = %agent.agent,
159 reason = %agent.reason,
160 "binding_unattested: agent runs DeclarationOnly (strict_binding is off)"
161 );
162 let Some(run_ctx) = run_ctx else {
163 return;
164 };
165 let entry = crate::store::run::DegradationEntry {
166 tool: "binding".to_string(),
167 error: agent.reason.clone(),
168 fallback: "DeclarationOnly".to_string(),
169 note: Some(format!(
170 "agent '{}' launched without a binding attestation (strict_binding off)",
171 agent.agent
172 )),
173 step_ref: None,
174 attempt: None,
175 at: crate::types::now_unix(),
176 };
177 if let Err(error) = run_ctx
178 .run_store
179 .append_degradation(&run_ctx.run_id, entry)
180 .await
181 {
182 tracing::warn!(
183 agent = %agent.agent,
184 %error,
185 "binding_unattested: failed to record degradation entry"
186 );
187 }
188}
189
190async fn load_or_resolve_bound_agents(
191 blueprint: &Blueprint,
192 run_ctx: Option<&RunContext>,
193 binding_provider: Option<&dyn AgentBindingProvider>,
194 legacy_worker_binding_policy: LegacyWorkerBindingPolicy,
195) -> Result<(Vec<BoundAgent>, SnapshotOrigin), TaskLaunchError> {
196 let strict = blueprint.strategy.strict_binding;
199 let resolve_fresh = || match legacy_worker_binding_policy {
200 LegacyWorkerBindingPolicy::Allow => resolve_bound_agents(blueprint),
201 LegacyWorkerBindingPolicy::Reject => {
202 crate::blueprint::resolve_bound_agents_strict(blueprint)
203 }
204 };
205 let Some(run_ctx) = run_ctx else {
206 let mut bound_agents = resolve_fresh().map_err(CompileError::from)?;
209 attest_or_gate_fresh(&mut bound_agents, binding_provider, strict, None).await?;
210 return Ok((bound_agents, SnapshotOrigin::Launch));
211 };
212
213 let record = run_ctx
214 .run_store
215 .get(&run_ctx.run_id)
216 .await
217 .map_err(|e| TaskLaunchError::PreDispatch(format!("load Run binding snapshot: {e}")))?;
218 if let Some(input_json) = record.input_json.as_deref() {
219 let snapshot: Value = serde_json::from_str(input_json).map_err(|e| {
220 TaskLaunchError::PreDispatch(format!("decode Run launch snapshot: {e}"))
221 })?;
222 if let Some(value) = snapshot.get("bound_agents") {
223 let bound_agents: Vec<BoundAgent> =
224 serde_json::from_value(value.clone()).map_err(|e| {
225 TaskLaunchError::PreDispatch(format!("decode Run BoundAgent snapshot: {e}"))
226 })?;
227 validate_bound_agent_snapshots(&bound_agents).map_err(|error| {
228 TaskLaunchError::PreDispatch(format!("validate Run BoundAgent snapshot: {error}"))
229 })?;
230 let origin = SnapshotOrigin::from_snapshot(&snapshot);
238 return Ok((bound_agents, origin));
239 }
240 }
241
242 let mut bound_agents = resolve_fresh().map_err(CompileError::from)?;
243 attest_or_gate_fresh(&mut bound_agents, binding_provider, strict, Some(run_ctx)).await?;
244 let origin = if run_ctx.resume {
250 SnapshotOrigin::ResumeBackfill
251 } else {
252 SnapshotOrigin::Launch
253 };
254 if let Some(input_json) = record.input_json {
255 let mut snapshot: Value = serde_json::from_str(&input_json).map_err(|e| {
256 TaskLaunchError::PreDispatch(format!("decode Run launch snapshot: {e}"))
257 })?;
258 let object = snapshot.as_object_mut().ok_or_else(|| {
259 TaskLaunchError::PreDispatch("Run launch snapshot must be a JSON object".to_string())
260 })?;
261 object.insert(
266 "bound_agents".to_string(),
267 serde_json::to_value(&bound_agents).map_err(|e| {
268 TaskLaunchError::PreDispatch(format!("encode Run BoundAgent snapshot: {e}"))
269 })?,
270 );
271 object.insert(
272 crate::store::run::BOUND_AGENTS_ORIGIN_KEY.to_string(),
273 serde_json::to_value(origin).map_err(|e| {
274 TaskLaunchError::PreDispatch(format!("encode Run BoundAgent origin: {e}"))
275 })?,
276 );
277 run_ctx
278 .run_store
279 .set_input_json(
280 &run_ctx.run_id,
281 serde_json::to_string(&snapshot).map_err(|e| {
282 TaskLaunchError::PreDispatch(format!("encode Run launch snapshot: {e}"))
283 })?,
284 )
285 .await
286 .map_err(|e| {
287 TaskLaunchError::PreDispatch(format!("persist Run BoundAgent snapshot: {e}"))
288 })?;
289 }
290 if origin == SnapshotOrigin::ResumeBackfill {
294 record_backfill_degradation(run_ctx, blueprint.id.as_str()).await;
295 }
296 Ok((bound_agents, origin))
297}
298
299async fn record_backfill_degradation(run_ctx: &RunContext, blueprint_id: &str) {
306 tracing::warn!(
307 run_id = %run_ctx.run_id,
308 blueprint = %blueprint_id,
309 "binding_backfill: resumed Run had no binding snapshot; bound_agents \
310 re-derived from the current Blueprint (not a launch-time pin)"
311 );
312 let entry = crate::store::run::DegradationEntry {
313 tool: "binding".to_string(),
314 error: "run resumed without a launch-pinned binding snapshot".to_string(),
315 fallback: "resume_backfill".to_string(),
316 note: Some(format!(
317 "run '{}' backfilled bound_agents from Blueprint '{}' at resume time",
318 run_ctx.run_id, blueprint_id
319 )),
320 step_ref: None,
321 attempt: None,
322 at: crate::types::now_unix(),
323 };
324 if let Err(error) = run_ctx
325 .run_store
326 .append_degradation(&run_ctx.run_id, entry)
327 .await
328 {
329 tracing::warn!(
330 run_id = %run_ctx.run_id,
331 %error,
332 "binding_backfill: failed to record degradation entry"
333 );
334 }
335}
336
337fn derive_audits(blueprint: &Blueprint) -> Vec<AuditDef> {
346 blueprint.audits.clone()
347}
348
349pub(crate) fn derive_agent_ctx(blueprint: &Blueprint) -> (Option<Value>, HashMap<String, Value>) {
378 let global = blueprint.default_agent_ctx.clone();
379 let meta_pool = derive_step_metas(blueprint);
380 let per_agent = blueprint
381 .agents
382 .iter()
383 .filter_map(|ad| {
384 let meta = ad.meta.as_ref()?;
385 let inline = meta.ctx.clone();
386 let base = meta.meta_ref.as_ref().and_then(|name| {
387 let resolved = meta_pool.get(name).cloned();
388 if resolved.is_none() {
389 tracing::warn!(
390 agent = %ad.name,
391 meta_ref = %name,
392 "derive_agent_ctx: AgentMeta.meta_ref names an undefined Blueprint.metas entry; skipping the base layer"
393 );
394 }
395 resolved
396 });
397 let merged = match (base, inline) {
398 (None, None) => None,
399 (Some(base), None) => Some(base),
400 (None, Some(inline)) => Some(inline),
401 (Some(base), Some(inline)) => Some(shallow_merge_inline_wins(base, inline)),
402 };
403 merged.map(|ctx| (ad.name.clone(), ctx))
404 })
405 .collect();
406 (global, per_agent)
407}
408
409pub(crate) fn shallow_merge_inline_wins(base: Value, inline: Value) -> Value {
416 match (base, inline) {
417 (Value::Object(mut base), Value::Object(inline)) => {
418 for (k, v) in inline {
419 base.insert(k, v);
420 }
421 Value::Object(base)
422 }
423 (_, inline) => inline,
424 }
425}
426
427fn derive_step_metas(blueprint: &Blueprint) -> HashMap<String, Value> {
435 blueprint
436 .metas
437 .iter()
438 .map(|m| (m.name.clone(), m.ctx.clone()))
439 .collect()
440}
441
442fn derive_context_policies(
450 blueprint: &Blueprint,
451) -> (Option<ContextPolicy>, HashMap<String, ContextPolicy>) {
452 let default_policy = blueprint.default_context_policy.clone();
453 let per_agent = blueprint
454 .agents
455 .iter()
456 .filter_map(|ad| {
457 let meta = ad.meta.as_ref()?;
458 let policy = meta.context_policy.clone()?;
459 Some((ad.name.clone(), policy))
460 })
461 .collect();
462 (default_policy, per_agent)
463}
464
465fn merge_init_ctx(bp_default: Option<&Value>, task_init_ctx: &Value) -> Value {
481 match (bp_default, task_init_ctx) {
482 (Some(Value::Object(bp_map)), Value::Object(task_map)) => {
483 let mut merged = bp_map.clone();
484 for (k, v) in task_map {
485 merged.insert(k.clone(), v.clone());
486 }
487 Value::Object(merged)
488 }
489 (None, _) => task_init_ctx.clone(),
490 (_, task) => task.clone(),
491 }
492}
493
494pub fn merge_init_ctx_3layer(
510 bp_default: Option<&Value>,
511 task_init_ctx: &Value,
512 run_override: Option<&Value>,
513) -> Value {
514 let bp_task = merge_init_ctx(bp_default, task_init_ctx);
515 match run_override {
516 Some(run) => merge_init_ctx(Some(&bp_task), run),
517 None => bp_task,
518 }
519}
520
521fn derive_bp_agent_kinds(blueprint: &Blueprint) -> HashMap<String, OperatorKind> {
522 let mut out = HashMap::new();
523 if blueprint.operators.is_empty() {
524 return out;
525 }
526 for agent in &blueprint.agents {
527 let Some(op_ref) = agent.spec.get("operator_ref").and_then(|v| v.as_str()) else {
528 continue;
529 };
530 let Some(op_def) = blueprint.operators.iter().find(|o| o.name == op_ref) else {
531 continue;
532 };
533 if let Some(kind) = op_def.kind {
534 out.insert(agent.name.clone(), OperatorKind::from(kind));
535 }
536 }
537 out
538}
539
540#[derive(Debug, Error)]
542pub enum TaskLaunchError {
543 #[error("compile: {0}")]
545 Compile(#[from] CompileError),
546 #[error("engine: {0}")]
548 Engine(#[from] EngineError),
549 #[error("flow eval: {0}")]
552 FlowEval(String),
553 #[error("pre-dispatch: {0}")]
560 PreDispatch(String),
561}
562
563#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize, schemars::JsonSchema)]
586pub struct TaskInputSpec {
587 #[serde(default)]
589 pub project_root: Option<String>,
590 #[serde(default)]
592 pub work_dir: Option<String>,
593 #[serde(default)]
595 #[schemars(with = "Option<Value>")]
596 pub task_metadata: Option<Value>,
597}
598
599#[derive(Debug, Clone)]
601pub struct TaskLaunchInput {
602 pub blueprint: Blueprint,
604 pub operator_id: String,
606 pub role: Role,
608 pub ttl: Duration,
610 pub operator_kind: Option<OperatorKind>,
619 pub bridge_id: Option<String>,
623 pub hook_id: Option<String>,
626 pub operator_backend_id: Option<String>,
633 pub operator_kind_overrides: HashMap<String, OperatorKind>,
638 pub init_ctx: Value,
643 pub task_input: Option<TaskInputSpec>,
650 pub run_ctx: Option<RunContext>,
657 pub check_policy: Option<CheckPolicy>,
680}
681
682impl TaskLaunchInput {
683 pub fn automate(
693 blueprint: Blueprint,
694 operator_id: impl Into<String>,
695 role: Role,
696 ttl: Duration,
697 init_ctx: Value,
698 ) -> Self {
699 Self {
700 blueprint,
701 operator_id: operator_id.into(),
702 role,
703 ttl,
704 operator_kind: None,
705 bridge_id: None,
706 hook_id: None,
707 operator_backend_id: None,
708 operator_kind_overrides: HashMap::new(),
709 init_ctx,
710 task_input: None,
711 run_ctx: None,
712 check_policy: None,
713 }
714 }
715}
716
717#[derive(Debug, Clone)]
719pub struct TaskLaunchOutput {
720 pub token: CapToken,
722 pub final_ctx: Value,
726}
727
728pub struct TaskLaunchService {
732 engine: Engine,
733 compiler: Compiler,
734 externs: Arc<dyn Externs + Send + Sync>,
739 binding_provider: Option<Arc<dyn AgentBindingProvider>>,
742 legacy_worker_binding_policy: LegacyWorkerBindingPolicy,
745}
746
747impl TaskLaunchService {
748 pub fn new(engine: Engine, compiler: Compiler) -> Self {
750 Self {
751 engine,
752 compiler,
753 externs: Arc::new(NoExterns),
754 binding_provider: None,
755 legacy_worker_binding_policy: LegacyWorkerBindingPolicy::default(),
756 }
757 }
758
759 pub fn with_externs(mut self, externs: Arc<dyn Externs + Send + Sync>) -> Self {
763 self.externs = externs;
764 self
765 }
766
767 pub fn with_binding_provider(mut self, provider: Arc<dyn AgentBindingProvider>) -> Self {
771 self.binding_provider = Some(provider);
772 self
773 }
774
775 pub fn with_legacy_worker_binding_policy(mut self, policy: LegacyWorkerBindingPolicy) -> Self {
777 self.legacy_worker_binding_policy = policy;
778 self
779 }
780
781 pub fn engine(&self) -> &Engine {
783 &self.engine
784 }
785
786 pub fn compiler(&self) -> &Compiler {
788 &self.compiler
789 }
790
791 pub async fn launch(
803 &self,
804 mut input: TaskLaunchInput,
805 ) -> Result<TaskLaunchOutput, TaskLaunchError> {
806 let (bound_agents, snapshot_origin) = load_or_resolve_bound_agents(
815 &input.blueprint,
816 input.run_ctx.as_ref(),
817 self.binding_provider.as_deref(),
818 self.legacy_worker_binding_policy,
819 )
820 .await?;
821 let binding_digests: HashMap<String, crate::blueprint::BindingDigest> = bound_agents
822 .iter()
823 .map(|bound| (bound.agent.name.clone(), bound.binding_digest.clone()))
824 .collect();
825 if let Some(run_ctx) = input.run_ctx.take() {
826 input.run_ctx = Some(match snapshot_origin {
837 SnapshotOrigin::Launch => run_ctx.with_binding_digests(binding_digests.clone()),
838 SnapshotOrigin::ResumeBackfill => run_ctx,
839 });
840 }
841 input.blueprint = materialize_bound_blueprint(&input.blueprint, &bound_agents);
842 let compiled = self
843 .compiler
844 .compile_bound(&input.blueprint, &bound_agents)?;
845 self.engine
854 .register_verdict_contracts(compiled.router.verdict_contracts.clone());
855 let spawner = linker::link(
856 compiled.router.clone(),
857 &input.blueprint.spawner_hints.layers,
858 &self.engine,
859 );
860 let (agent_ctx_global, agent_ctx_per_agent) = derive_agent_ctx(&input.blueprint);
876 let (context_policy_default, context_policy_per_agent) =
877 derive_context_policies(&input.blueprint);
878 let spawner = SpawnerStack::new(spawner)
879 .layer(AgentContextMiddleware::new(
880 agent_ctx_global,
881 agent_ctx_per_agent,
882 context_policy_default,
883 context_policy_per_agent,
884 ))
885 .build();
886 let spawner = if let Some(alias) = input.blueprint.metadata.project_name_alias.as_deref() {
893 SpawnerStack::new(spawner)
894 .layer(ProjectNameAliasMiddleware::new(alias))
895 .build()
896 } else {
897 spawner
898 };
899 let worker_bindings = worker_bindings_from_bound_agents(&bound_agents);
903 let spawner = if worker_bindings.is_empty() {
904 spawner
905 } else {
906 SpawnerStack::new(spawner)
907 .layer(WorkerBindingMiddleware::new(worker_bindings))
908 .build()
909 };
910 let audit_defs = derive_audits(&input.blueprint);
920 let spawner = if audit_defs.is_empty() {
921 spawner
922 } else {
923 SpawnerStack::new(spawner)
924 .layer(AfterRunAuditMiddleware::new(
925 audit_defs,
926 compiled.router.clone(),
927 ))
928 .build()
929 };
930
931 let spawner = match input.task_input.as_ref().and_then(|spec| {
938 TaskInputMiddleware::new_from_fields(
939 spec.project_root.clone(),
940 spec.work_dir.clone(),
941 spec.task_metadata.clone(),
942 )
943 }) {
944 Some(task_input) => SpawnerStack::new(spawner).layer(task_input).build(),
945 None => spawner,
946 };
947
948 let bp_agent_kinds = derive_bp_agent_kinds(&input.blueprint);
953 let bp_global_kind = input
954 .blueprint
955 .default_operator_kind
956 .map(OperatorKind::from);
957
958 let token = self
959 .engine
960 .attach_with_ids(
961 input.operator_id,
962 input.role,
963 input.ttl,
964 input.operator_kind,
965 input.bridge_id,
966 input.hook_id,
967 input.operator_backend_id,
968 input.operator_kind_overrides,
969 bp_agent_kinds,
970 bp_global_kind,
971 )
972 .await?;
973 let resolved_check_policy = input.check_policy.or(input.blueprint.check_policy);
984 let effective_check_policy =
999 resolved_check_policy.unwrap_or(self.engine.cfg().check_policy);
1000 if effective_check_policy == CheckPolicy::Strict {
1001 let roots_missing = input
1002 .task_input
1003 .as_ref()
1004 .map(|t| t.project_root.is_none() && t.work_dir.is_none())
1005 .unwrap_or(true);
1006 if roots_missing {
1007 return Err(TaskLaunchError::PreDispatch(
1008 "check_policy=strict requires project_root or work_dir, but the launch \
1009 supplied neither"
1010 .to_string(),
1011 ));
1012 }
1013 }
1014 let dispatcher =
1015 EngineDispatcher::with_spawner(self.engine.clone(), token.clone(), spawner);
1016 let dispatcher = dispatcher.with_check_policy(resolved_check_policy);
1017 let dispatcher = match input.run_ctx {
1018 Some(run_ctx) => dispatcher.with_run(run_ctx),
1019 None => dispatcher,
1020 };
1021 let dispatcher = dispatcher.with_step_metas(derive_step_metas(&input.blueprint));
1025 let dispatcher = dispatcher.with_binding_digests(binding_digests);
1026 let dispatcher = dispatcher.with_step_naming(compiled.step_naming.clone());
1032 let dispatcher =
1039 dispatcher.with_projection_placement(compiled.projection_placement.clone());
1040 let merged_init_ctx =
1046 merge_init_ctx(input.blueprint.default_init_ctx.as_ref(), &input.init_ctx);
1047 let final_ctx = mlua_flow_ir::eval_async_externs(
1048 &input.blueprint.flow,
1049 merged_init_ctx,
1050 &dispatcher,
1051 &*self.externs,
1052 )
1053 .await
1054 .map_err(|e| TaskLaunchError::FlowEval(e.to_string()))?;
1055 Ok(TaskLaunchOutput { token, final_ctx })
1056 }
1057}
1058
1059#[cfg(test)]
1064mod tests {
1065 use super::*;
1066 use crate::blueprint::compiler::{RustFnInProcessSpawnerFactory, SpawnerRegistry};
1067 use crate::blueprint::{
1068 current_schema_version, resolve_runner, AgentDef, AgentKind, AgentMeta, AgentProfile,
1069 BlueprintMetadata, CompilerHints, CompilerStrategy, MetaDef, Runner,
1070 };
1071 use crate::core::config::EngineCfg;
1072 use crate::worker::adapter::{WorkerError, WorkerResult};
1073 use mlua_flow_ir::{Expr, JoinMode, Node as FlowNode};
1074 use serde_json::json;
1075 use std::sync::Arc;
1076
1077 fn path(s: &str) -> Expr {
1078 Expr::Path {
1079 at: s.parse().expect("literal test path"),
1080 }
1081 }
1082 fn step(ref_: &str, in_: Expr, out: Expr) -> FlowNode {
1083 FlowNode::Step {
1084 ref_: ref_.to_string(),
1085 in_,
1086 out,
1087 }
1088 }
1089
1090 fn agent(name: &str, fn_id: &str) -> AgentDef {
1091 AgentDef {
1092 name: name.to_string(),
1093 kind: AgentKind::RustFn,
1094 spec: json!({ "fn_id": fn_id }),
1095 profile: None,
1096 meta: Some(AgentMeta::default()),
1097 runner: None,
1098 runner_ref: None,
1099 verdict: None,
1100 }
1101 }
1102
1103 fn build_service(factory: RustFnInProcessSpawnerFactory) -> TaskLaunchService {
1104 let engine = Engine::new(EngineCfg::default());
1105 let mut reg = SpawnerRegistry::new();
1106 reg.register::<RustFnInProcessSpawnerFactory>(Arc::new(factory));
1107 let compiler = Compiler::new(reg);
1108 TaskLaunchService::new(engine, compiler)
1109 }
1110
1111 fn build_service_with_cfg(
1115 factory: RustFnInProcessSpawnerFactory,
1116 cfg: EngineCfg,
1117 ) -> TaskLaunchService {
1118 let engine = Engine::new(cfg);
1119 let mut reg = SpawnerRegistry::new();
1120 reg.register::<RustFnInProcessSpawnerFactory>(Arc::new(factory));
1121 let compiler = Compiler::new(reg);
1122 TaskLaunchService::new(engine, compiler)
1123 }
1124
1125 fn bp(flow: FlowNode, agents: Vec<AgentDef>) -> Blueprint {
1126 Blueprint {
1127 schema_version: current_schema_version(),
1128 id: "ut".into(),
1129 flow,
1130 agents,
1131 operators: vec![],
1132 metas: vec![],
1133 hints: CompilerHints::default(),
1134 strategy: CompilerStrategy::default(),
1135 metadata: BlueprintMetadata::default(),
1136 spawner_hints: Default::default(),
1137 default_agent_kind: AgentKind::Operator,
1138 default_operator_kind: None,
1139 default_init_ctx: None,
1140 default_agent_ctx: None,
1141 default_context_policy: None,
1142 projection_placement: None,
1143 audits: vec![],
1144 degradation_policy: None,
1145 runners: vec![],
1146 default_runner: None,
1147 check_policy: None,
1148 blueprint_ref_includes: Vec::new(),
1149 }
1150 }
1151
1152 fn launch_input(blueprint: Blueprint, init_ctx: Value) -> TaskLaunchInput {
1153 TaskLaunchInput::automate(
1154 blueprint,
1155 "ut-op",
1156 Role::Operator,
1157 Duration::from_secs(30),
1158 init_ctx,
1159 )
1160 }
1161
1162 #[test]
1168 fn derive_audits_empty_by_default() {
1169 let blueprint = bp(
1170 step("echo", path("$.input"), path("$.out")),
1171 vec![agent("echo", "echo")],
1172 );
1173 assert!(
1174 derive_audits(&blueprint).is_empty(),
1175 "audits_absent_no_layer: an undeclared audits Vec must stay empty"
1176 );
1177 }
1178
1179 #[test]
1180 fn derive_audits_returns_blueprint_audits_verbatim() {
1181 let mut blueprint = bp(
1182 step("echo", path("$.input"), path("$.out")),
1183 vec![agent("echo", "echo")],
1184 );
1185 blueprint.audits = vec![crate::blueprint::AuditDef {
1186 agent: "auditor".to_string(),
1187 steps: None,
1188 mode: crate::blueprint::AuditMode::Async,
1189 }];
1190 let got = derive_audits(&blueprint);
1191 assert_eq!(got.len(), 1);
1192 assert_eq!(got[0].agent, "auditor");
1193 }
1194
1195 #[tokio::test]
1196 async fn launch_appends_audit_artifact_when_audits_declared() {
1197 use crate::blueprint::{AuditDef, AuditMode};
1198
1199 let factory = RustFnInProcessSpawnerFactory::new()
1200 .register_fn("echo", |inv| async move {
1201 Ok(WorkerResult {
1202 value: json!({ "echoed": inv.prompt }),
1203 ok: true,
1204 })
1205 })
1206 .register_fn("audit-fn", |_inv| async move {
1207 Ok(WorkerResult {
1208 value: json!({ "finding": "clean" }),
1209 ok: true,
1210 })
1211 });
1212 let svc = build_service(factory);
1213 let mut blueprint = bp(
1214 step("echo", path("$.input"), path("$.out")),
1215 vec![agent("echo", "echo"), agent("auditor", "audit-fn")],
1216 );
1217 blueprint.audits = vec![AuditDef {
1218 agent: "auditor".to_string(),
1219 steps: None,
1220 mode: AuditMode::Sync,
1221 }];
1222 let out = svc
1223 .launch(launch_input(blueprint, json!({ "input": "hi" })))
1224 .await
1225 .expect("launch ok — audits must never alter the audited step's outcome");
1226 assert_eq!(out.final_ctx["out"]["echoed"], "hi");
1227
1228 let audited_task_id = svc
1229 .engine()
1230 .with_state("test.find_audited_task", |s| {
1231 s.tasks
1232 .iter()
1233 .find(|(_, t)| t.spec.agent == "echo")
1234 .map(|(id, _)| id.clone())
1235 })
1236 .await
1237 .expect("with_state")
1238 .expect("the echo task must exist");
1239 let tail = svc.engine().output_tail(&audited_task_id, 1).await;
1240 let found = tail.iter().any(|ev| {
1241 matches!(
1242 ev,
1243 crate::worker::output::OutputEvent::Artifact { name, .. } if name == "audit:echo"
1244 )
1245 });
1246 assert!(
1247 found,
1248 "launch() must wire AfterRunAuditMiddleware end-to-end when Blueprint.audits is declared"
1249 );
1250 }
1251
1252 #[tokio::test]
1253 async fn launch_single_step_writes_out_path() {
1254 let factory = RustFnInProcessSpawnerFactory::new().register_fn("echo", |inv| async move {
1255 Ok(WorkerResult {
1256 value: json!({ "echoed": inv.prompt }),
1257 ok: true,
1258 })
1259 });
1260 let svc = build_service(factory);
1261 let blueprint = bp(
1262 step("echo", path("$.input"), path("$.out")),
1263 vec![agent("echo", "echo")],
1264 );
1265 let out = svc
1266 .launch(launch_input(blueprint, json!({ "input": "hi" })))
1267 .await
1268 .expect("launch ok");
1269 assert_eq!(out.final_ctx["out"]["echoed"], "hi");
1270 }
1271
1272 async fn dispatched_check_policy(
1293 launch_policy: Option<CheckPolicy>,
1294 bp_policy: Option<CheckPolicy>,
1295 ) -> Option<CheckPolicy> {
1296 let factory = RustFnInProcessSpawnerFactory::new().register_fn("echo", |inv| async move {
1297 Ok(WorkerResult {
1298 value: json!({ "echoed": inv.prompt }),
1299 ok: true,
1300 })
1301 });
1302 let svc = build_service(factory);
1303 let mut blueprint = bp(
1304 step("echo", path("$.input"), path("$.out")),
1305 vec![agent("echo", "echo")],
1306 );
1307 blueprint.check_policy = bp_policy;
1308 let mut input = launch_input(blueprint, json!({ "input": "hi" }));
1309 input.check_policy = launch_policy;
1310 input.task_input = Some(TaskInputSpec {
1311 project_root: None,
1312 work_dir: Some("/dispatched-check-policy-test-root".to_string()),
1313 task_metadata: None,
1314 });
1315 let _ = svc.launch(input).await;
1316 svc.engine()
1317 .with_state("test.read_dispatched_check_policy", |s| {
1318 s.tasks
1319 .values()
1320 .find(|t| t.spec.agent == "echo")
1321 .and_then(|t| t.spec.check_policy)
1322 })
1323 .await
1324 .expect("with_state")
1325 }
1326
1327 #[tokio::test]
1330 async fn cascade_launch_tier_wins_over_blueprint_tier() {
1331 assert_eq!(
1332 dispatched_check_policy(Some(CheckPolicy::Silent), Some(CheckPolicy::Strict)).await,
1333 Some(CheckPolicy::Silent),
1334 );
1335 }
1336
1337 #[tokio::test]
1340 async fn cascade_blueprint_tier_used_when_launch_absent() {
1341 assert_eq!(
1342 dispatched_check_policy(None, Some(CheckPolicy::Strict)).await,
1343 Some(CheckPolicy::Strict),
1344 );
1345 }
1346
1347 #[tokio::test]
1350 async fn cascade_launch_tier_alone_when_blueprint_absent() {
1351 assert_eq!(
1352 dispatched_check_policy(Some(CheckPolicy::Strict), None).await,
1353 Some(CheckPolicy::Strict),
1354 );
1355 }
1356
1357 #[tokio::test]
1362 async fn cascade_both_none_preserves_server_fallback() {
1363 assert_eq!(dispatched_check_policy(None, None).await, None);
1364 }
1365
1366 #[tokio::test]
1382 async fn strict_blueprint_without_roots_is_rejected_pre_dispatch() {
1383 let factory = RustFnInProcessSpawnerFactory::new().register_fn("echo", |inv| async move {
1384 Ok(WorkerResult {
1385 value: json!({ "echoed": inv.prompt }),
1386 ok: true,
1387 })
1388 });
1389 let svc = build_service(factory);
1390 let mut blueprint = bp(
1391 step("echo", path("$.input"), path("$.out")),
1392 vec![agent("echo", "echo")],
1393 );
1394 blueprint.check_policy = Some(CheckPolicy::Strict);
1395 let err = svc
1397 .launch(launch_input(blueprint, json!({ "input": "hi" })))
1398 .await
1399 .expect_err("strict check_policy + no roots must be rejected before dispatch");
1400 match err {
1401 TaskLaunchError::PreDispatch(message) => {
1402 assert!(
1403 message.contains("strict"),
1404 "message must identify the strict-requires-roots condition: {message}"
1405 );
1406 }
1407 other => panic!("expected TaskLaunchError::PreDispatch, got {other:?}"),
1408 }
1409
1410 let dispatched = svc
1414 .engine()
1415 .with_state("test.no_echo_task_dispatched", |s| {
1416 s.tasks.values().any(|t| t.spec.agent == "echo")
1417 })
1418 .await
1419 .expect("with_state");
1420 assert!(
1421 !dispatched,
1422 "the pre-dispatch guard must reject before any step is dispatched"
1423 );
1424 }
1425
1426 #[tokio::test]
1438 async fn launch_without_any_check_policy_completes_fail_open() {
1439 let factory = RustFnInProcessSpawnerFactory::new().register_fn("echo", |inv| async move {
1440 Ok(WorkerResult {
1441 value: json!({ "echoed": inv.prompt }),
1442 ok: true,
1443 })
1444 });
1445 let svc = build_service(factory);
1446 let blueprint = bp(
1447 step("echo", path("$.input"), path("$.out")),
1448 vec![agent("echo", "echo")],
1449 );
1450 assert_eq!(blueprint.check_policy, None, "BP tier must be unset");
1451 let out = svc
1452 .launch(launch_input(blueprint, json!({ "input": "hi" })))
1453 .await
1454 .expect("warn-mode fail-open must let the launch complete");
1455 assert_eq!(out.final_ctx["out"]["echoed"], "hi");
1456 }
1457
1458 #[tokio::test]
1479 async fn strict_blueprint_with_launch_warn_override_bypasses_pre_dispatch_guard() {
1480 let factory = RustFnInProcessSpawnerFactory::new().register_fn("echo", |inv| async move {
1481 Ok(WorkerResult {
1482 value: json!({ "echoed": inv.prompt }),
1483 ok: true,
1484 })
1485 });
1486 let svc = build_service(factory);
1487 let mut blueprint = bp(
1488 step("echo", path("$.input"), path("$.out")),
1489 vec![agent("echo", "echo")],
1490 );
1491 blueprint.check_policy = Some(CheckPolicy::Strict);
1492 let mut input = launch_input(blueprint, json!({ "input": "hi" }));
1493 input.check_policy = Some(CheckPolicy::Warn);
1494 assert!(input.task_input.is_none(), "no roots supplied at all");
1495 let out = svc
1496 .launch(input)
1497 .await
1498 .expect("launch-tier warn override must bypass the pre-dispatch guard");
1499 assert_eq!(out.final_ctx["out"]["echoed"], "hi");
1500 }
1501
1502 #[tokio::test]
1508 async fn server_tier_strict_alone_triggers_pre_dispatch_guard() {
1509 let factory = RustFnInProcessSpawnerFactory::new().register_fn("echo", |inv| async move {
1510 Ok(WorkerResult {
1511 value: json!({ "echoed": inv.prompt }),
1512 ok: true,
1513 })
1514 });
1515 let svc = build_service_with_cfg(
1516 factory,
1517 EngineCfg {
1518 check_policy: CheckPolicy::Strict,
1519 ..EngineCfg::default()
1520 },
1521 );
1522 let blueprint = bp(
1523 step("echo", path("$.input"), path("$.out")),
1524 vec![agent("echo", "echo")],
1525 );
1526 assert_eq!(blueprint.check_policy, None, "BP tier must be unset");
1527 let input = launch_input(blueprint, json!({ "input": "hi" }));
1528 assert!(input.check_policy.is_none(), "launch tier must be unset");
1529 assert!(input.task_input.is_none(), "no roots supplied");
1530 let err = svc.launch(input).await.expect_err(
1531 "server-tier Strict alone (BP/launch tiers both unset) must trigger the guard",
1532 );
1533 match err {
1534 TaskLaunchError::PreDispatch(message) => {
1535 assert!(
1536 message.contains("strict"),
1537 "expected the strict-requires-roots message, got: {message}"
1538 );
1539 }
1540 other => panic!("expected TaskLaunchError::PreDispatch, got {other:?}"),
1541 }
1542 }
1543
1544 #[tokio::test]
1550 async fn pre_dispatch_guard_rejects_when_task_input_present_but_roots_both_none() {
1551 let factory = RustFnInProcessSpawnerFactory::new().register_fn("echo", |inv| async move {
1552 Ok(WorkerResult {
1553 value: json!({ "echoed": inv.prompt }),
1554 ok: true,
1555 })
1556 });
1557 let svc = build_service(factory);
1558 let mut blueprint = bp(
1559 step("echo", path("$.input"), path("$.out")),
1560 vec![agent("echo", "echo")],
1561 );
1562 blueprint.check_policy = Some(CheckPolicy::Strict);
1563 let mut input = launch_input(blueprint, json!({ "input": "hi" }));
1564 input.task_input = Some(TaskInputSpec {
1565 project_root: None,
1566 work_dir: None,
1567 task_metadata: Some(json!({ "unrelated": true })),
1568 });
1569 let err = svc
1570 .launch(input)
1571 .await
1572 .expect_err("Some(TaskInputSpec) with both roots None must still be roots_missing");
1573 assert!(
1574 matches!(err, TaskLaunchError::PreDispatch(_)),
1575 "expected TaskLaunchError::PreDispatch, got {err:?}"
1576 );
1577 }
1578
1579 #[tokio::test]
1584 async fn pre_dispatch_guard_passes_when_work_dir_present_and_project_root_absent() {
1585 let factory = RustFnInProcessSpawnerFactory::new().register_fn("echo", |inv| async move {
1586 Ok(WorkerResult {
1587 value: json!({ "echoed": inv.prompt }),
1588 ok: true,
1589 })
1590 });
1591 let svc = build_service(factory);
1592 let mut blueprint = bp(
1593 step("echo", path("$.input"), path("$.out")),
1594 vec![agent("echo", "echo")],
1595 );
1596 blueprint.check_policy = Some(CheckPolicy::Strict);
1597 let mut input = launch_input(blueprint, json!({ "input": "hi" }));
1598 input.task_input = Some(TaskInputSpec {
1599 project_root: None,
1600 work_dir: Some("/repo/work".to_string()),
1601 task_metadata: None,
1602 });
1603 let out = svc
1604 .launch(input)
1605 .await
1606 .expect("work_dir alone must satisfy the guard's roots_missing check");
1607 assert_eq!(out.final_ctx["out"]["echoed"], "hi");
1608 }
1609
1610 #[tokio::test]
1611 async fn launch_three_step_seq_threads_ctx_forward() {
1612 let factory = RustFnInProcessSpawnerFactory::new()
1613 .register_fn("upper", |inv| async move {
1614 let s = serde_json::from_str::<String>(&inv.prompt).unwrap_or(inv.prompt);
1615 Ok(WorkerResult {
1616 value: json!(s.to_uppercase()),
1617 ok: true,
1618 })
1619 })
1620 .register_fn("suffix", |inv| async move {
1621 let s = serde_json::from_str::<String>(&inv.prompt).unwrap_or(inv.prompt);
1622 Ok(WorkerResult {
1623 value: json!(format!("{s}!")),
1624 ok: true,
1625 })
1626 })
1627 .register_fn("wrap", |inv| async move {
1628 let s = serde_json::from_str::<String>(&inv.prompt).unwrap_or(inv.prompt);
1629 Ok(WorkerResult {
1630 value: json!(format!("[{s}]")),
1631 ok: true,
1632 })
1633 });
1634 let svc = build_service(factory);
1635 let flow = FlowNode::Seq {
1636 children: vec![
1637 step("upper", path("$.in"), path("$.s1")),
1638 step("suffix", path("$.s1"), path("$.s2")),
1639 step("wrap", path("$.s2"), path("$.s3")),
1640 ],
1641 };
1642 let blueprint = bp(
1643 flow,
1644 vec![
1645 agent("upper", "upper"),
1646 agent("suffix", "suffix"),
1647 agent("wrap", "wrap"),
1648 ],
1649 );
1650 let out = svc
1651 .launch(launch_input(blueprint, json!({ "in": "hello" })))
1652 .await
1653 .expect("launch ok");
1654 assert_eq!(out.final_ctx["s1"], "HELLO");
1655 assert_eq!(out.final_ctx["s2"], "HELLO!");
1656 assert_eq!(out.final_ctx["s3"], "[HELLO!]");
1657 }
1658
1659 #[tokio::test]
1660 async fn launch_fanout_join_all_parallel_completes() {
1661 use std::sync::atomic::{AtomicU32, Ordering};
1662 let counter = Arc::new(AtomicU32::new(0));
1663 let max_seen = Arc::new(AtomicU32::new(0));
1664 let counter_clone = counter.clone();
1665 let max_clone = max_seen.clone();
1666
1667 let factory = RustFnInProcessSpawnerFactory::new().register_fn("para", move |inv| {
1670 let counter = counter_clone.clone();
1671 let max_seen = max_clone.clone();
1672 async move {
1673 let now = counter.fetch_add(1, Ordering::SeqCst) + 1;
1674 let mut prev = max_seen.load(Ordering::SeqCst);
1675 while now > prev {
1676 match max_seen.compare_exchange(prev, now, Ordering::SeqCst, Ordering::SeqCst) {
1677 Ok(_) => break,
1678 Err(p) => prev = p,
1679 }
1680 }
1681 tokio::time::sleep(Duration::from_millis(50)).await;
1682 counter.fetch_sub(1, Ordering::SeqCst);
1683 let s = serde_json::from_str::<String>(&inv.prompt).unwrap_or(inv.prompt);
1684 Ok(WorkerResult {
1685 value: json!(format!("did:{s}")),
1686 ok: true,
1687 })
1688 }
1689 });
1690 let svc = build_service(factory);
1691 let flow = FlowNode::Fanout {
1692 items: path("$.items"),
1693 bind: path("$.item"),
1694 body: Box::new(step("para", path("$.item"), path("$.r"))),
1695 join: JoinMode::All,
1696 out: path("$.results"),
1697 };
1698 let blueprint = bp(flow, vec![agent("para", "para")]);
1699 let out = svc
1700 .launch(launch_input(
1701 blueprint,
1702 json!({ "items": ["a", "b", "c", "d"] }),
1703 ))
1704 .await
1705 .expect("launch ok");
1706 let results = out.final_ctx["results"].as_array().expect("array");
1707 assert_eq!(results.len(), 4);
1708 for (i, expected) in ["a", "b", "c", "d"].iter().enumerate() {
1709 assert_eq!(results[i]["r"], json!(format!("did:{expected}")));
1710 }
1711 let max = max_seen.load(Ordering::SeqCst);
1712 assert!(
1713 max >= 2,
1714 "expected parallel execution (max inflight >= 2), got {max}"
1715 );
1716 }
1717
1718 #[tokio::test]
1719 async fn launch_propagates_worker_error_as_flow_eval_err() {
1720 let factory = RustFnInProcessSpawnerFactory::new()
1721 .register_fn("ok", |inv| async move {
1722 Ok(WorkerResult {
1723 value: json!(inv.prompt),
1724 ok: true,
1725 })
1726 })
1727 .register_fn("boom", |_inv| async move {
1728 Err(WorkerError::Failed("intentional boom".into()))
1729 });
1730 let svc = build_service(factory);
1731 let flow = FlowNode::Seq {
1732 children: vec![
1733 step("ok", path("$.input"), path("$.s1")),
1734 step("boom", path("$.s1"), path("$.s2")),
1735 step("ok", path("$.s2"), path("$.s3")),
1736 ],
1737 };
1738 let blueprint = bp(flow, vec![agent("ok", "ok"), agent("boom", "boom")]);
1739 let err = svc
1740 .launch(launch_input(blueprint, json!({ "input": "x" })))
1741 .await
1742 .expect_err("expected fail");
1743 match err {
1744 TaskLaunchError::FlowEval(msg) => {
1745 assert!(
1746 msg.contains("boom") || msg.contains("intentional"),
1747 "expected error to mention worker failure, got: {msg}"
1748 );
1749 }
1750 other => panic!("expected FlowEval error, got {other:?}"),
1751 }
1752 }
1753
1754 #[tokio::test]
1755 async fn launch_resolves_call_extern_via_registered_externs() {
1756 let factory = RustFnInProcessSpawnerFactory::new().register_fn("echo", |inv| async move {
1757 Ok(WorkerResult {
1758 value: json!({ "echoed": inv.prompt }),
1759 ok: true,
1760 })
1761 });
1762 let mut externs = mlua_flow_ir::ExternMap::new();
1763 externs.register("fmt.greet", |args: &[Value]| {
1764 let name = args[0].as_str().unwrap_or("?");
1765 Ok(json!(format!("hello, {name}")))
1766 });
1767 let svc = build_service(factory).with_externs(Arc::new(externs));
1768 let flow = step(
1769 "echo",
1770 Expr::CallExtern {
1771 ref_: "fmt.greet".into(),
1772 args: vec![path("$.who")],
1773 },
1774 path("$.out"),
1775 );
1776 let blueprint = bp(flow, vec![agent("echo", "echo")]);
1777 let out = svc
1778 .launch(launch_input(blueprint, json!({ "who": "swarm" })))
1779 .await
1780 .expect("launch ok");
1781 assert_eq!(out.final_ctx["out"]["echoed"], json!("hello, swarm"));
1782 }
1783
1784 #[tokio::test]
1785 async fn launch_call_extern_without_registry_fails_as_flow_eval() {
1786 let factory = RustFnInProcessSpawnerFactory::new().register_fn("echo", |inv| async move {
1787 Ok(WorkerResult {
1788 value: json!(inv.prompt),
1789 ok: true,
1790 })
1791 });
1792 let svc = build_service(factory); let flow = step(
1794 "echo",
1795 Expr::CallExtern {
1796 ref_: "fmt.greet".into(),
1797 args: vec![],
1798 },
1799 path("$.out"),
1800 );
1801 let blueprint = bp(flow, vec![agent("echo", "echo")]);
1802 let err = svc
1803 .launch(launch_input(blueprint, json!({})))
1804 .await
1805 .expect_err("expected fail");
1806 match err {
1807 TaskLaunchError::FlowEval(msg) => {
1808 assert!(msg.contains("extern"), "expected extern error, got: {msg}");
1809 }
1810 other => panic!("expected FlowEval error, got {other:?}"),
1811 }
1812 }
1813
1814 #[tokio::test]
1834 async fn launch_registers_the_blueprints_verdict_contracts_into_the_engine() {
1835 let factory = RustFnInProcessSpawnerFactory::new().register_fn("gate", |inv| async move {
1836 Ok(WorkerResult {
1837 value: json!(inv.prompt),
1838 ok: true,
1839 })
1840 });
1841 let svc = build_service(factory);
1842 let mut gate_agent = agent("gate", "gate");
1843 gate_agent.verdict = Some(mlua_swarm_schema::VerdictContract {
1844 channel: mlua_swarm_schema::VerdictChannel::Body,
1845 values: vec!["PASS".to_string(), "BLOCKED".to_string()],
1846 });
1847 let flow = step("gate", path("$.input"), path("$.out"));
1848 let blueprint = bp(flow, vec![gate_agent]);
1849
1850 let out = svc
1851 .launch(launch_input(blueprint, json!({ "input": "PASS" })))
1852 .await
1853 .expect("launch ok");
1854 assert_eq!(out.final_ctx["out"], json!("PASS"));
1855
1856 let task_id = svc
1861 .engine()
1862 .with_state("test.find_dispatched_task_id", |s| {
1863 s.tasks.keys().next().cloned()
1864 })
1865 .await
1866 .expect("with_state")
1867 .expect("launch must have dispatched exactly one Step (one TaskState)");
1868
1869 let contract = svc
1870 .engine()
1871 .verdict_contract_for_task(&task_id)
1872 .await
1873 .expect(
1874 "TaskLaunchService::launch must have merged this Blueprint's compiled \
1875 verdict_contracts into the engine's runtime registry \
1876 (Engine::register_verdict_contracts, called right after \
1877 compiler.compile succeeds) — verdict_contract_for_task resolving None \
1878 here means that production wiring regressed",
1879 );
1880 assert_eq!(contract.channel, mlua_swarm_schema::VerdictChannel::Body);
1881 assert_eq!(
1882 contract.values,
1883 vec!["PASS".to_string(), "BLOCKED".to_string()]
1884 );
1885 }
1886
1887 #[tokio::test]
1892 async fn launch_with_run_ctx_appends_one_step_entry_per_dispatched_step() {
1893 use crate::store::run::{InMemoryRunStore, RunContext, RunRecord, RunStatus, RunStore};
1894 use crate::types::{RunId, TaskId};
1895
1896 let factory = RustFnInProcessSpawnerFactory::new()
1897 .register_fn("upper", |inv| async move {
1898 Ok(WorkerResult {
1899 value: json!(inv.prompt.to_uppercase()),
1900 ok: true,
1901 })
1902 })
1903 .register_fn("suffix", |inv| async move {
1904 let s = serde_json::from_str::<String>(&inv.prompt).unwrap_or(inv.prompt);
1905 Ok(WorkerResult {
1906 value: json!(format!("{s}!")),
1907 ok: true,
1908 })
1909 });
1910 let svc = build_service(factory);
1911 let flow = FlowNode::Seq {
1912 children: vec![
1913 step("upper", path("$.in"), path("$.s1")),
1914 step("suffix", path("$.s1"), path("$.s2")),
1915 ],
1916 };
1917 let blueprint = bp(
1918 flow,
1919 vec![agent("upper", "upper"), agent("suffix", "suffix")],
1920 );
1921
1922 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
1923 let run_id = RunId::new();
1924 run_store
1925 .create(RunRecord {
1926 id: run_id.clone(),
1927 task_id: TaskId::new(),
1928 status: RunStatus::Running,
1929 step_entries: Vec::new(),
1930 degradations: Vec::new(),
1931 operator_sid: None,
1932 result_ref: None,
1933 input_json: Some("{}".to_string()),
1934 created_at: 0,
1935 updated_at: 0,
1936 })
1937 .await
1938 .expect("seed RunRecord");
1939
1940 let mut input = launch_input(blueprint, json!({ "in": "hi" }));
1941 input.run_ctx = Some(RunContext::new(run_id.clone(), run_store.clone()));
1942
1943 let out = svc.launch(input).await.expect("launch ok");
1944 assert_eq!(out.final_ctx["s2"], "HI!");
1945
1946 let run = run_store.get(&run_id).await.expect("run present");
1947 assert_eq!(
1948 run.step_entries.len(),
1949 2,
1950 "expected one step_entry per dispatched step, got {:?}",
1951 run.step_entries
1952 );
1953 assert_eq!(run.step_entries[0].step_ref, Some("upper".to_string()));
1954 assert_eq!(run.step_entries[0].status, Some("passed".to_string()));
1955 assert!(run.step_entries[0].binding_digest.is_some());
1956 assert_eq!(run.step_entries[1].step_ref, Some("suffix".to_string()));
1957 assert_eq!(run.step_entries[1].status, Some("passed".to_string()));
1958 assert!(run.step_entries[1].binding_digest.is_some());
1959 let snapshot: Value = serde_json::from_str(run.input_json.as_deref().unwrap()).unwrap();
1960 assert_eq!(snapshot["bound_agents"].as_array().unwrap().len(), 2);
1961 }
1962
1963 #[tokio::test]
1964 async fn run_snapshot_reuses_bound_agent_after_blueprint_mutation() {
1965 use crate::store::run::{InMemoryRunStore, RunContext, RunRecord, RunStatus, RunStore};
1966 use crate::types::{RunId, TaskId};
1967
1968 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
1969 let run_id = RunId::new();
1970 run_store
1971 .create(RunRecord {
1972 id: run_id.clone(),
1973 task_id: TaskId::new(),
1974 status: RunStatus::Running,
1975 step_entries: Vec::new(),
1976 degradations: Vec::new(),
1977 operator_sid: None,
1978 result_ref: None,
1979 input_json: Some("{}".to_string()),
1980 created_at: 0,
1981 updated_at: 0,
1982 })
1983 .await
1984 .unwrap();
1985 let run_ctx = RunContext::new(run_id, run_store);
1986 let mut original_agent = agent("worker", "worker");
1987 original_agent.profile = Some(crate::blueprint::AgentProfile {
1988 system_prompt: "original role".to_string(),
1989 ..Default::default()
1990 });
1991 let mut blueprint = bp(
1992 step("worker", path("$.input"), path("$.out")),
1993 vec![original_agent],
1994 );
1995
1996 let (original, _) = load_or_resolve_bound_agents(
1997 &blueprint,
1998 Some(&run_ctx),
1999 None,
2000 LegacyWorkerBindingPolicy::Allow,
2001 )
2002 .await
2003 .unwrap();
2004 blueprint.agents[0].profile.as_mut().unwrap().system_prompt = "mutated role".to_string();
2005 let (restored, _) = load_or_resolve_bound_agents(
2006 &blueprint,
2007 Some(&run_ctx),
2008 None,
2009 LegacyWorkerBindingPolicy::Allow,
2010 )
2011 .await
2012 .unwrap();
2013
2014 assert_eq!(restored[0].binding_digest, original[0].binding_digest);
2015 assert_eq!(
2016 restored[0].agent.profile.as_ref().unwrap().system_prompt,
2017 "original role"
2018 );
2019 }
2020
2021 #[tokio::test]
2022 async fn strict_migration_policy_rejects_fresh_legacy_worker_binding() {
2023 let mut legacy_agent = agent("worker", "worker");
2024 legacy_agent.profile = Some(AgentProfile {
2025 worker_binding: Some("legacy-worker".to_string()),
2026 ..Default::default()
2027 });
2028 let blueprint = bp(
2029 step("worker", path("$.input"), path("$.out")),
2030 vec![legacy_agent],
2031 );
2032
2033 let error =
2034 load_or_resolve_bound_agents(&blueprint, None, None, LegacyWorkerBindingPolicy::Reject)
2035 .await
2036 .expect_err("strict migration policy must reject fallback");
2037 assert!(error
2038 .to_string()
2039 .contains("deprecated profile.worker_binding"));
2040 }
2041
2042 #[tokio::test]
2043 async fn run_snapshot_calls_binding_provider_only_on_first_resolution() {
2044 use crate::binding::{AgentBindingProvider, BindingProviderError};
2045 use crate::blueprint::{BindOutcome, BindReceipt, BindRequest};
2046 use crate::store::run::{InMemoryRunStore, RunContext, RunRecord, RunStatus, RunStore};
2047 use crate::types::{RunId, TaskId};
2048 use std::sync::atomic::{AtomicUsize, Ordering};
2049
2050 struct CountingProvider(AtomicUsize);
2051
2052 #[async_trait::async_trait]
2053 impl AgentBindingProvider for CountingProvider {
2054 async fn bind(
2055 &self,
2056 requests: &[BindRequest],
2057 ) -> Result<Vec<BindOutcome>, BindingProviderError> {
2058 self.0.fetch_add(1, Ordering::SeqCst);
2059 Ok(requests
2060 .iter()
2061 .map(|request| BindOutcome::Bound {
2062 receipt: BindReceipt {
2063 agent: request.agent.clone(),
2064 request_digest: request.request_digest.clone(),
2065 provider_id: "operator-main-ai".to_string(),
2066 provider_revision: Some("test".to_string()),
2067 resolved_model: request.requested_model.clone(),
2068 effective_tools: request.requested_tools.clone(),
2069 launch_variant: request.launch_variant.clone(),
2070 capability_snapshot_digest: None,
2071 },
2072 })
2073 .collect())
2074 }
2075 }
2076
2077 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
2078 let run_id = RunId::new();
2079 run_store
2080 .create(RunRecord {
2081 id: run_id.clone(),
2082 task_id: TaskId::new(),
2083 status: RunStatus::Running,
2084 step_entries: Vec::new(),
2085 degradations: Vec::new(),
2086 operator_sid: None,
2087 result_ref: None,
2088 input_json: Some("{}".to_string()),
2089 created_at: 0,
2090 updated_at: 0,
2091 })
2092 .await
2093 .unwrap();
2094 let run_ctx = RunContext::new(run_id, run_store);
2095 let mut blueprint = bp(
2096 step("worker", path("$.input"), path("$.out")),
2097 vec![agent("worker", "worker")],
2098 );
2099 blueprint.agents[0].runner = Some(Runner::WsClaudeCode {
2100 variant: "mse-worker".to_string(),
2101 tools: vec!["Read".to_string()],
2102 });
2103 let provider = CountingProvider(AtomicUsize::new(0));
2104
2105 let (first, _) = load_or_resolve_bound_agents(
2106 &blueprint,
2107 Some(&run_ctx),
2108 Some(&provider),
2109 LegacyWorkerBindingPolicy::Allow,
2110 )
2111 .await
2112 .unwrap();
2113 let (restored, _) = load_or_resolve_bound_agents(
2114 &blueprint,
2115 Some(&run_ctx),
2116 Some(&provider),
2117 LegacyWorkerBindingPolicy::Allow,
2118 )
2119 .await
2120 .unwrap();
2121
2122 assert_eq!(provider.0.load(Ordering::SeqCst), 1);
2123 assert!(first[0].attestation.is_some());
2124 assert_eq!(restored, first);
2125 }
2126
2127 #[tokio::test]
2136 async fn fresh_resolve_on_launch_persists_launch_origin_no_degradation() {
2137 use crate::store::run::{
2138 InMemoryRunStore, RunContext, RunRecord, RunStatus, RunStore, BOUND_AGENTS_ORIGIN_KEY,
2139 };
2140 use crate::types::{RunId, TaskId};
2141
2142 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
2143 let run_id = RunId::new();
2144 run_store
2145 .create(RunRecord {
2146 id: run_id.clone(),
2147 task_id: TaskId::new(),
2148 status: RunStatus::Running,
2149 step_entries: Vec::new(),
2150 degradations: Vec::new(),
2151 operator_sid: None,
2152 result_ref: None,
2153 input_json: Some("{}".to_string()),
2154 created_at: 0,
2155 updated_at: 0,
2156 })
2157 .await
2158 .unwrap();
2159 let run_ctx = RunContext::new(run_id.clone(), run_store.clone());
2161 let blueprint = bp(
2162 step("worker", path("$.input"), path("$.out")),
2163 vec![agent("worker", "worker")],
2164 );
2165
2166 load_or_resolve_bound_agents(
2167 &blueprint,
2168 Some(&run_ctx),
2169 None,
2170 LegacyWorkerBindingPolicy::Allow,
2171 )
2172 .await
2173 .expect("launch resolve ok");
2174
2175 let run = run_store.get(&run_id).await.expect("run present");
2176 let snapshot: Value = serde_json::from_str(run.input_json.as_deref().unwrap()).unwrap();
2177 assert!(
2178 snapshot["bound_agents"].is_array(),
2179 "bound_agents must be persisted"
2180 );
2181 assert_eq!(snapshot[BOUND_AGENTS_ORIGIN_KEY], json!("launch"));
2182 assert_eq!(
2183 SnapshotOrigin::from_snapshot(&snapshot),
2184 SnapshotOrigin::Launch
2185 );
2186 assert!(
2187 run.degradations.is_empty(),
2188 "an initial-launch resolve is not a degradation"
2189 );
2190 }
2191
2192 #[tokio::test]
2197 async fn backfill_on_resume_persists_resume_origin_and_records_degradation() {
2198 use crate::store::run::{
2199 InMemoryRunStore, RunContext, RunRecord, RunStatus, RunStore, BOUND_AGENTS_ORIGIN_KEY,
2200 };
2201 use crate::types::{RunId, TaskId};
2202
2203 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
2204 let run_id = RunId::new();
2205 run_store
2206 .create(RunRecord {
2207 id: run_id.clone(),
2208 task_id: TaskId::new(),
2209 status: RunStatus::Running,
2210 step_entries: Vec::new(),
2211 degradations: Vec::new(),
2212 operator_sid: None,
2213 result_ref: None,
2214 input_json: Some("{}".to_string()),
2215 created_at: 0,
2216 updated_at: 0,
2217 })
2218 .await
2219 .unwrap();
2220 let run_ctx = RunContext::new(run_id.clone(), run_store.clone()).with_resume();
2221 let blueprint = bp(
2222 step("worker", path("$.input"), path("$.out")),
2223 vec![agent("worker", "worker")],
2224 );
2225
2226 load_or_resolve_bound_agents(
2227 &blueprint,
2228 Some(&run_ctx),
2229 None,
2230 LegacyWorkerBindingPolicy::Allow,
2231 )
2232 .await
2233 .expect("resume backfill ok");
2234
2235 let run = run_store.get(&run_id).await.expect("run present");
2236 let snapshot: Value = serde_json::from_str(run.input_json.as_deref().unwrap()).unwrap();
2237 assert_eq!(snapshot[BOUND_AGENTS_ORIGIN_KEY], json!("resume_backfill"));
2238 assert_eq!(
2239 SnapshotOrigin::from_snapshot(&snapshot),
2240 SnapshotOrigin::ResumeBackfill
2241 );
2242 assert_eq!(
2243 run.degradations.len(),
2244 1,
2245 "a resume backfill must record exactly one degradation"
2246 );
2247 assert_eq!(run.degradations[0].tool, "binding");
2248 assert_eq!(run.degradations[0].fallback, "resume_backfill");
2249 }
2250
2251 fn runner_blueprint(strict_binding: bool) -> Blueprint {
2258 let mut blueprint = bp(
2259 step("worker", path("$.input"), path("$.out")),
2260 vec![agent("worker", "worker")],
2261 );
2262 blueprint.strategy.strict_binding = strict_binding;
2263 blueprint.agents[0].runner = Some(Runner::WsClaudeCode {
2264 variant: "mse-worker".to_string(),
2265 tools: vec!["Read".to_string()],
2266 });
2267 blueprint
2268 }
2269
2270 struct AlwaysUnboundProvider;
2273
2274 #[async_trait::async_trait]
2275 impl AgentBindingProvider for AlwaysUnboundProvider {
2276 async fn bind(
2277 &self,
2278 requests: &[crate::blueprint::BindRequest],
2279 ) -> Result<Vec<crate::blueprint::BindOutcome>, crate::binding::BindingProviderError>
2280 {
2281 Ok(requests
2282 .iter()
2283 .map(|request| crate::blueprint::BindOutcome::Unbound {
2284 agent: request.agent.clone(),
2285 reason: "no capability manifest submitted".to_string(),
2286 })
2287 .collect())
2288 }
2289 }
2290
2291 #[tokio::test]
2295 async fn non_strict_unbound_agent_runs_declaration_only_with_degradation() {
2296 use crate::store::run::{InMemoryRunStore, RunContext, RunRecord, RunStatus, RunStore};
2297 use crate::types::{RunId, TaskId};
2298
2299 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
2300 let run_id = RunId::new();
2301 run_store
2302 .create(RunRecord {
2303 id: run_id.clone(),
2304 task_id: TaskId::new(),
2305 status: RunStatus::Running,
2306 step_entries: Vec::new(),
2307 degradations: Vec::new(),
2308 operator_sid: None,
2309 result_ref: None,
2310 input_json: Some("{}".to_string()),
2311 created_at: 0,
2312 updated_at: 0,
2313 })
2314 .await
2315 .unwrap();
2316 let run_ctx = RunContext::new(run_id.clone(), run_store.clone());
2317
2318 let (bound, _) = load_or_resolve_bound_agents(
2319 &runner_blueprint(false),
2320 Some(&run_ctx),
2321 Some(&AlwaysUnboundProvider),
2322 LegacyWorkerBindingPolicy::Allow,
2323 )
2324 .await
2325 .expect("non-strict launch must succeed even without an attestation");
2326 assert!(
2327 bound[0].attestation.is_none(),
2328 "an unattested agent must stay DeclarationOnly"
2329 );
2330
2331 let run = run_store.get(&run_id).await.expect("run present");
2332 assert_eq!(run.degradations.len(), 1, "expected one degradation entry");
2333 assert_eq!(run.degradations[0].tool, "binding");
2334 assert_eq!(run.degradations[0].fallback, "DeclarationOnly");
2335 assert!(run.degradations[0].error.contains("no capability manifest"));
2336 }
2337
2338 #[tokio::test]
2342 async fn strict_unbound_agent_fails_with_requirements_in_message() {
2343 let error = load_or_resolve_bound_agents(
2344 &runner_blueprint(true),
2345 None,
2346 Some(&AlwaysUnboundProvider),
2347 LegacyWorkerBindingPolicy::Allow,
2348 )
2349 .await
2350 .expect_err("strict + Unbound must reject the launch");
2351 match error {
2352 TaskLaunchError::PreDispatch(message) => {
2353 assert!(message.contains("worker"), "message: {message}");
2354 assert!(message.contains("mse-worker"), "message: {message}");
2355 assert!(message.contains("Read"), "message: {message}");
2356 }
2357 other => panic!("expected PreDispatch, got {other:?}"),
2358 }
2359 }
2360
2361 #[tokio::test]
2364 async fn strict_without_provider_rejects_runner_backed_launch() {
2365 let error = load_or_resolve_bound_agents(
2366 &runner_blueprint(true),
2367 None,
2368 None,
2369 LegacyWorkerBindingPolicy::Allow,
2370 )
2371 .await
2372 .expect_err("strict + no provider must reject a Runner-backed launch");
2373 match error {
2374 TaskLaunchError::PreDispatch(message) => {
2375 assert!(
2376 message.contains("strict_binding requires a binding provider"),
2377 "message: {message}"
2378 );
2379 }
2380 other => panic!("expected PreDispatch, got {other:?}"),
2381 }
2382 }
2383
2384 #[tokio::test]
2387 async fn strict_with_correct_manifest_attests_the_agent() {
2388 use crate::binding::ManifestBindingProvider;
2389 use crate::blueprint::{AgentProviderCapability, AgentProviderManifest};
2390
2391 let provider = ManifestBindingProvider::new(AgentProviderManifest {
2392 provider_id: "operator-main-ai".to_string(),
2393 provider_revision: Some("1".to_string()),
2394 capabilities: vec![AgentProviderCapability {
2395 launch_variant: Some("mse-worker".to_string()),
2396 resolved_model: None,
2397 effective_tools: vec!["Read".to_string()],
2398 capability_snapshot_digest: None,
2399 }],
2400 });
2401 let (bound, _) = load_or_resolve_bound_agents(
2402 &runner_blueprint(true),
2403 None,
2404 Some(&provider),
2405 LegacyWorkerBindingPolicy::Allow,
2406 )
2407 .await
2408 .expect("strict launch with a correct manifest must attest");
2409 assert!(
2410 bound[0].attestation.is_some(),
2411 "a correctly attested agent must carry its attestation"
2412 );
2413 }
2414
2415 #[test]
2424 fn worker_bindings_carry_request_digest_and_model() {
2425 let mut blueprint = runner_blueprint(false);
2426 blueprint.agents[0].profile = Some(AgentProfile {
2427 model: Some("claude-sonnet".to_string()),
2428 ..Default::default()
2429 });
2430 let bound = resolve_bound_agents(&blueprint).expect("resolvable Runner refs");
2431 let bindings = worker_bindings_from_bound_agents(&bound);
2432
2433 let wb = bindings.get("worker").expect("worker binding present");
2434 assert_eq!(
2435 wb.request_digest.as_ref(),
2436 Some(&bound[0].binding_digest),
2437 "the spawn frame must carry the immutable snapshot digest"
2438 );
2439 assert!(wb
2440 .request_digest
2441 .as_ref()
2442 .unwrap()
2443 .as_str()
2444 .starts_with("sha256:"));
2445 assert_eq!(wb.requested_model.as_deref(), Some("claude-sonnet"));
2446 }
2447
2448 #[test]
2451 fn worker_bindings_omit_model_when_profile_has_none() {
2452 let bound = resolve_bound_agents(&runner_blueprint(false)).expect("resolvable Runner refs");
2453 let bindings = worker_bindings_from_bound_agents(&bound);
2454 let wb = bindings.get("worker").expect("worker binding present");
2455 assert!(wb.request_digest.is_some());
2456 assert!(wb.requested_model.is_none());
2457 }
2458
2459 #[tokio::test]
2460 async fn launch_without_run_ctx_appends_no_step_entries() {
2461 let factory = RustFnInProcessSpawnerFactory::new().register_fn("echo", |inv| async move {
2465 Ok(WorkerResult {
2466 value: json!(inv.prompt),
2467 ok: true,
2468 })
2469 });
2470 let svc = build_service(factory);
2471 let blueprint = bp(
2472 step("echo", path("$.input"), path("$.out")),
2473 vec![agent("echo", "echo")],
2474 );
2475 let input = launch_input(blueprint, json!({ "input": "hi" }));
2476 assert!(
2477 input.run_ctx.is_none(),
2478 "automate() defaults run_ctx to None"
2479 );
2480 let out = svc.launch(input).await.expect("launch ok");
2481 assert_eq!(out.final_ctx["out"], "hi");
2482 }
2483
2484 #[tokio::test]
2490 async fn launch_with_task_input_leaves_init_ctx_object_seed_unmutated() {
2491 let factory = RustFnInProcessSpawnerFactory::new().register_fn("echo", |inv| async move {
2498 Ok(WorkerResult {
2499 value: json!({ "echoed": inv.prompt }),
2500 ok: true,
2501 })
2502 });
2503 let svc = build_service(factory);
2504 let blueprint = bp(
2505 step("echo", path("$.input"), path("$.out")),
2506 vec![agent("echo", "echo")],
2507 );
2508 let mut input = launch_input(blueprint, json!({ "input": "hi" }));
2509 input.task_input = Some(TaskInputSpec {
2510 project_root: Some("/repo".to_string()),
2511 work_dir: Some("/repo/work".to_string()),
2512 task_metadata: Some(json!({ "issue": 19 })),
2513 });
2514 let out = svc.launch(input).await.expect("launch ok");
2515 assert_eq!(out.final_ctx["out"]["echoed"], "hi");
2516 assert!(
2517 out.final_ctx.get("project_root").is_none(),
2518 "task_input must not be folded into the flow-ir ctx seed, got {:?}",
2519 out.final_ctx
2520 );
2521 assert!(out.final_ctx.get("work_dir").is_none());
2522 assert!(out.final_ctx.get("task_metadata").is_none());
2523 }
2524
2525 #[tokio::test]
2526 async fn launch_with_task_input_none_is_a_no_op() {
2527 let factory = RustFnInProcessSpawnerFactory::new().register_fn("echo", |inv| async move {
2528 Ok(WorkerResult {
2529 value: json!(inv.prompt),
2530 ok: true,
2531 })
2532 });
2533 let svc = build_service(factory);
2534 let blueprint = bp(
2535 step("echo", path("$.input"), path("$.out")),
2536 vec![agent("echo", "echo")],
2537 );
2538 let mut input = launch_input(blueprint, json!({ "input": "hi" }));
2539 assert!(input.task_input.is_none(), "automate() defaults to None");
2540 input.task_input = None;
2541 let out = svc.launch(input).await.expect("launch ok");
2542 assert_eq!(out.final_ctx["out"], "hi");
2543 }
2544
2545 #[tokio::test]
2546 async fn launch_with_task_input_all_fields_absent_is_a_no_op() {
2547 let factory = RustFnInProcessSpawnerFactory::new().register_fn("echo", |inv| async move {
2551 Ok(WorkerResult {
2552 value: json!(inv.prompt),
2553 ok: true,
2554 })
2555 });
2556 let svc = build_service(factory);
2557 let blueprint = bp(
2558 step("echo", path("$.input"), path("$.out")),
2559 vec![agent("echo", "echo")],
2560 );
2561 let mut input = launch_input(blueprint, json!({ "input": "hi" }));
2562 input.task_input = Some(TaskInputSpec::default());
2563 let out = svc.launch(input).await.expect("launch ok");
2564 assert_eq!(out.final_ctx["out"], "hi");
2565 }
2566
2567 #[test]
2572 fn merge_init_ctx_bp_default_only_passes_through_when_task_is_empty_object() {
2573 let bp_default = json!({ "seeded": "from-bp" });
2574 let task = json!({});
2575 let merged = merge_init_ctx(Some(&bp_default), &task);
2576 assert_eq!(merged, json!({ "seeded": "from-bp" }));
2577 }
2578
2579 #[test]
2580 fn merge_init_ctx_task_only_passes_through_when_bp_default_is_empty_object() {
2581 let bp_default = json!({});
2582 let task = json!({ "seeded": "from-task" });
2583 let merged = merge_init_ctx(Some(&bp_default), &task);
2584 assert_eq!(merged, json!({ "seeded": "from-task" }));
2585 }
2586
2587 #[test]
2588 fn merge_init_ctx_both_objects_task_wins_on_key_collision() {
2589 let bp_default = json!({ "a": "bp", "b": "bp-only" });
2590 let task = json!({ "a": "task", "c": "task-only" });
2591 let merged = merge_init_ctx(Some(&bp_default), &task);
2592 assert_eq!(
2593 merged,
2594 json!({ "a": "task", "b": "bp-only", "c": "task-only" })
2595 );
2596 }
2597
2598 #[test]
2599 fn merge_init_ctx_non_object_task_fully_replaces_bp_default() {
2600 let bp_default = json!({ "seeded": "from-bp" });
2601 let task = json!("plain-string-seed");
2602 let merged = merge_init_ctx(Some(&bp_default), &task);
2603 assert_eq!(merged, json!("plain-string-seed"));
2604 }
2605
2606 #[test]
2607 fn merge_init_ctx_no_bp_default_is_a_no_op() {
2608 let task = json!({ "input": "hi" });
2609 let merged = merge_init_ctx(None, &task);
2610 assert_eq!(merged, task);
2611 }
2612
2613 #[test]
2618 fn merge_init_ctx_3layer_no_run_override_equals_bp_task_merge_only() {
2619 let bp_default = json!({ "a": "bp", "b": "bp-only" });
2623 let task = json!({ "a": "task", "c": "task-only" });
2624 let three_layer = merge_init_ctx_3layer(Some(&bp_default), &task, None);
2625 let two_layer = merge_init_ctx(Some(&bp_default), &task);
2626 assert_eq!(three_layer, two_layer);
2627 assert_eq!(
2628 three_layer,
2629 json!({ "a": "task", "b": "bp-only", "c": "task-only" })
2630 );
2631 }
2632
2633 #[test]
2634 fn merge_init_ctx_3layer_run_object_wins_on_key_collision_over_bp_and_task() {
2635 let bp_default = json!({ "a": "bp", "b": "bp-only" });
2636 let task = json!({ "a": "task", "c": "task-only" });
2637 let run_override = json!({ "a": "run", "d": "run-only" });
2638 let merged = merge_init_ctx_3layer(Some(&bp_default), &task, Some(&run_override));
2639 assert_eq!(
2640 merged,
2641 json!({ "a": "run", "b": "bp-only", "c": "task-only", "d": "run-only" }),
2642 "Run wins on collision (a); BP-only (b) and Task-only (c) keys survive"
2643 );
2644 }
2645
2646 #[test]
2647 fn merge_init_ctx_3layer_run_non_object_fully_replaces_bp_task_merge() {
2648 let bp_default = json!({ "seeded": "from-bp" });
2649 let task = json!({ "seeded": "from-task" });
2650 let run_override = json!("plain-string-run-seed");
2651 let merged = merge_init_ctx_3layer(Some(&bp_default), &task, Some(&run_override));
2652 assert_eq!(merged, json!("plain-string-run-seed"));
2653 }
2654
2655 #[test]
2656 fn merge_init_ctx_3layer_no_bp_default_and_no_run_override_is_task_passthrough() {
2657 let task = json!({ "input": "hi" });
2658 let merged = merge_init_ctx_3layer(None, &task, None);
2659 assert_eq!(merged, task);
2660 }
2661
2662 #[tokio::test]
2663 async fn launch_merges_bp_default_init_ctx_into_task_init_ctx() {
2664 let factory = RustFnInProcessSpawnerFactory::new().register_fn("echo", |inv| async move {
2667 Ok(WorkerResult {
2668 value: json!(inv.prompt),
2669 ok: true,
2670 })
2671 });
2672 let svc = build_service(factory);
2673 let mut blueprint = bp(
2674 step("echo", path("$.greeting"), path("$.out")),
2675 vec![agent("echo", "echo")],
2676 );
2677 blueprint.default_init_ctx = Some(json!({ "greeting": "hello from bp" }));
2678 let out = svc
2680 .launch(launch_input(blueprint, json!({})))
2681 .await
2682 .expect("launch ok");
2683 assert_eq!(out.final_ctx["out"], "hello from bp");
2684 }
2685
2686 fn agent_with_meta(name: &str, fn_id: &str, meta: AgentMeta) -> AgentDef {
2691 AgentDef {
2692 name: name.to_string(),
2693 kind: AgentKind::RustFn,
2694 spec: json!({ "fn_id": fn_id }),
2695 profile: None,
2696 meta: Some(meta),
2697 runner: None,
2698 runner_ref: None,
2699 verdict: None,
2700 }
2701 }
2702
2703 #[test]
2704 fn derive_agent_ctx_empty_blueprint_yields_empty_state() {
2705 let blueprint = bp(step("echo", path("$.in"), path("$.out")), vec![]);
2706 let (global, per_agent) = derive_agent_ctx(&blueprint);
2707 assert_eq!(global, None);
2708 assert!(per_agent.is_empty());
2709 }
2710
2711 #[test]
2712 fn derive_agent_ctx_populated_blueprint_yields_correct_maps() {
2713 let mut blueprint = bp(
2714 step("echo", path("$.in"), path("$.out")),
2715 vec![
2716 agent_with_meta(
2717 "with-ctx",
2718 "echo",
2719 AgentMeta {
2720 ctx: Some(json!({ "org_conventions": "x" })),
2721 ..Default::default()
2722 },
2723 ),
2724 agent("no-ctx", "echo"),
2725 ],
2726 );
2727 blueprint.default_agent_ctx = Some(json!({ "seeded": "from-bp" }));
2728 let (global, per_agent) = derive_agent_ctx(&blueprint);
2729 assert_eq!(global, Some(json!({ "seeded": "from-bp" })));
2730 assert_eq!(
2731 per_agent.len(),
2732 1,
2733 "agents without AgentMeta.ctx are absent, not defaulted to null: {per_agent:?}"
2734 );
2735 assert_eq!(
2736 per_agent.get("with-ctx"),
2737 Some(&json!({ "org_conventions": "x" }))
2738 );
2739 assert!(!per_agent.contains_key("no-ctx"));
2740 }
2741
2742 #[test]
2743 fn derive_context_policies_empty_blueprint_yields_empty_state() {
2744 let blueprint = bp(step("echo", path("$.in"), path("$.out")), vec![]);
2745 let (default_policy, per_agent) = derive_context_policies(&blueprint);
2746 assert_eq!(default_policy, None);
2747 assert!(per_agent.is_empty());
2748 }
2749
2750 #[test]
2751 fn derive_context_policies_populated_blueprint_yields_correct_maps() {
2752 let mut blueprint = bp(
2753 step("echo", path("$.in"), path("$.out")),
2754 vec![
2755 agent_with_meta(
2756 "with-policy",
2757 "echo",
2758 AgentMeta {
2759 context_policy: Some(ContextPolicy {
2760 include: None,
2761 exclude: vec!["work_dir".to_string()],
2762 ..Default::default()
2763 }),
2764 ..Default::default()
2765 },
2766 ),
2767 agent("no-policy", "echo"),
2768 ],
2769 );
2770 blueprint.default_context_policy = Some(ContextPolicy {
2771 include: Some(vec!["project_root".to_string()]),
2772 exclude: vec![],
2773 ..Default::default()
2774 });
2775 let (default_policy, per_agent) = derive_context_policies(&blueprint);
2776 assert_eq!(
2777 default_policy,
2778 Some(ContextPolicy {
2779 include: Some(vec!["project_root".to_string()]),
2780 exclude: vec![],
2781 ..Default::default()
2782 })
2783 );
2784 assert_eq!(per_agent.len(), 1);
2785 assert_eq!(
2786 per_agent.get("with-policy"),
2787 Some(&ContextPolicy {
2788 include: None,
2789 exclude: vec!["work_dir".to_string()],
2790 ..Default::default()
2791 })
2792 );
2793 assert!(!per_agent.contains_key("no-policy"));
2794 }
2795
2796 #[test]
2802 fn derive_step_metas_empty_blueprint_yields_empty_map() {
2803 let blueprint = bp(step("echo", path("$.in"), path("$.out")), vec![]);
2804 assert!(derive_step_metas(&blueprint).is_empty());
2805 }
2806
2807 #[test]
2808 fn derive_step_metas_populated_blueprint_yields_name_to_ctx_map() {
2809 let mut blueprint = bp(step("echo", path("$.in"), path("$.out")), vec![]);
2810 blueprint.metas = vec![
2811 MetaDef {
2812 name: "heavy-scan".to_string(),
2813 ctx: json!({ "work_dir": "/x" }),
2814 },
2815 MetaDef {
2816 name: "light-scan".to_string(),
2817 ctx: json!({ "work_dir": "/y" }),
2818 },
2819 ];
2820 let metas = derive_step_metas(&blueprint);
2821 assert_eq!(metas.len(), 2);
2822 assert_eq!(metas.get("heavy-scan"), Some(&json!({ "work_dir": "/x" })));
2823 assert_eq!(metas.get("light-scan"), Some(&json!({ "work_dir": "/y" })));
2824 }
2825
2826 #[test]
2827 fn derive_agent_ctx_meta_ref_resolves_as_base_under_inline_ctx() {
2828 let mut blueprint = bp(
2829 step("echo", path("$.in"), path("$.out")),
2830 vec![agent_with_meta(
2831 "with-meta-ref",
2832 "echo",
2833 AgentMeta {
2834 ctx: Some(json!({ "work_dir": "/inline-wins" })),
2835 meta_ref: Some("shared".to_string()),
2836 ..Default::default()
2837 },
2838 )],
2839 );
2840 blueprint.metas = vec![MetaDef {
2841 name: "shared".to_string(),
2842 ctx: json!({ "work_dir": "/base", "extra": "from-pool" }),
2843 }];
2844 let (_, per_agent) = derive_agent_ctx(&blueprint);
2845 assert_eq!(
2846 per_agent.get("with-meta-ref"),
2847 Some(&json!({ "work_dir": "/inline-wins", "extra": "from-pool" })),
2848 "inline ctx must win the collided key while pool-only keys survive the merge"
2849 );
2850 }
2851
2852 #[test]
2853 fn derive_agent_ctx_meta_ref_alone_uses_pool_ctx_verbatim() {
2854 let mut blueprint = bp(
2855 step("echo", path("$.in"), path("$.out")),
2856 vec![agent_with_meta(
2857 "with-meta-ref-only",
2858 "echo",
2859 AgentMeta {
2860 meta_ref: Some("shared".to_string()),
2861 ..Default::default()
2862 },
2863 )],
2864 );
2865 blueprint.metas = vec![MetaDef {
2866 name: "shared".to_string(),
2867 ctx: json!({ "work_dir": "/base" }),
2868 }];
2869 let (_, per_agent) = derive_agent_ctx(&blueprint);
2870 assert_eq!(
2871 per_agent.get("with-meta-ref-only"),
2872 Some(&json!({ "work_dir": "/base" }))
2873 );
2874 }
2875
2876 #[test]
2877 fn derive_agent_ctx_unresolved_meta_ref_never_panics_and_falls_back_to_inline() {
2878 let blueprint = bp(
2879 step("echo", path("$.in"), path("$.out")),
2880 vec![agent_with_meta(
2881 "with-unresolved-meta-ref",
2882 "echo",
2883 AgentMeta {
2884 ctx: Some(json!({ "work_dir": "/inline-only" })),
2885 meta_ref: Some("missing".to_string()),
2886 ..Default::default()
2887 },
2888 )],
2889 );
2890 let (_, per_agent) = derive_agent_ctx(&blueprint);
2892 assert_eq!(
2893 per_agent.get("with-unresolved-meta-ref"),
2894 Some(&json!({ "work_dir": "/inline-only" })),
2895 "an unresolved meta_ref must never panic; the agent's own inline ctx still applies"
2896 );
2897 }
2898
2899 #[test]
2915 fn resolve_runner_legacy_fallback_matches_derive_worker_bindings_semantics() {
2916 fn legacy_agent(name: &str, variant: &str, tools: Vec<&str>) -> AgentDef {
2917 AgentDef {
2918 name: name.to_string(),
2919 kind: AgentKind::Operator,
2920 spec: json!({}),
2921 profile: Some(AgentProfile {
2922 worker_binding: Some(variant.to_string()),
2923 tools: tools.into_iter().map(str::to_string).collect(),
2924 ..Default::default()
2925 }),
2926 meta: None,
2927 runner: None,
2928 runner_ref: None,
2929 verdict: None,
2930 }
2931 }
2932
2933 let blueprint = bp(
2934 step("planner", path("$.in"), path("$.out")),
2935 vec![
2936 legacy_agent("planner", "mse-worker-planner", vec!["Read", "Grep"]),
2937 legacy_agent("coder", "mse-worker-coder", vec![]),
2938 agent("no-binding", "echo"),
2939 ],
2940 );
2941
2942 let derived = derive_worker_bindings(&blueprint);
2943
2944 for agent_def in &blueprint.agents {
2945 let resolved = resolve_runner(&blueprint, agent_def).expect("no unresolved refs");
2946 match derived.get(&agent_def.name) {
2947 Some(binding) => {
2948 assert_eq!(
2949 resolved,
2950 Some(Runner::WsClaudeCode {
2951 variant: binding.variant.clone(),
2952 tools: binding.tools.clone(),
2953 }),
2954 "resolve_runner must synthesize the same WsClaudeCode Runner \
2955 derive_worker_bindings produces for agent '{}'",
2956 agent_def.name
2957 );
2958 }
2959 None => {
2960 assert_eq!(
2961 resolved, None,
2962 "agent '{}' has no derive_worker_bindings entry, so resolve_runner \
2963 must resolve to None too (no other tier declared)",
2964 agent_def.name
2965 );
2966 }
2967 }
2968 }
2969 }
2970
2971 #[test]
2972 fn ws_operator_runner_projects_into_the_existing_spawn_binding() {
2973 let mut blueprint = bp(
2974 step("reviewer", path("$.in"), path("$.out")),
2975 vec![agent("reviewer", "echo")],
2976 );
2977 blueprint.agents[0].runner = Some(Runner::WsOperator {
2978 variant: "mse-reviewer".to_string(),
2979 tools: vec!["Read".to_string(), "Grep".to_string()],
2980 });
2981
2982 let derived = derive_worker_bindings(&blueprint);
2983 let binding = derived
2984 .get("reviewer")
2985 .expect("ws_operator must feed the canonical spawn binding path");
2986 assert_eq!(binding.variant, "mse-reviewer");
2987 assert_eq!(binding.tools, ["Read", "Grep"]);
2988 }
2989
2990 fn counting_echo_service() -> (TaskLaunchService, Arc<std::sync::atomic::AtomicUsize>) {
2999 use std::sync::atomic::{AtomicUsize, Ordering};
3000 let calls = Arc::new(AtomicUsize::new(0));
3001 let counter = calls.clone();
3002 let factory = RustFnInProcessSpawnerFactory::new().register_fn("echo", move |inv| {
3003 let counter = counter.clone();
3004 async move {
3005 counter.fetch_add(1, Ordering::SeqCst);
3006 Ok(WorkerResult {
3007 value: json!({ "echoed": inv.prompt }),
3008 ok: true,
3009 })
3010 }
3011 });
3012 (build_service(factory), calls)
3013 }
3014
3015 async fn seed_legacy_run(run_store: &Arc<dyn crate::store::run::RunStore>) -> crate::RunId {
3016 use crate::store::run::{RunRecord, RunStatus};
3017 use crate::types::TaskId;
3018 let run_id = crate::RunId::new();
3019 run_store
3020 .create(RunRecord {
3021 id: run_id.clone(),
3022 task_id: TaskId::new(),
3023 status: RunStatus::Running,
3024 step_entries: Vec::new(),
3025 degradations: Vec::new(),
3026 operator_sid: None,
3027 result_ref: None,
3028 input_json: Some("{}".to_string()),
3031 created_at: 0,
3032 updated_at: 0,
3033 })
3034 .await
3035 .expect("seed legacy RunRecord");
3036 run_id
3037 }
3038
3039 #[tokio::test]
3045 async fn backfilled_run_replays_legacy_keys_stably_across_two_resumes() {
3046 use crate::store::replay::{InMemoryReplayStore, ReplayCursor, ReplayStore};
3047 use crate::store::run::{InMemoryRunStore, RunContext, RunStore};
3048 use std::sync::atomic::Ordering;
3049 use std::sync::Mutex;
3050
3051 let (svc, echo_calls) = counting_echo_service();
3052 let blueprint = bp(
3053 step("echo", path("$.input"), path("$.out")),
3054 vec![agent("echo", "echo")],
3055 );
3056 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
3057 let replay_store: Arc<dyn ReplayStore> = Arc::new(InMemoryReplayStore::new());
3058 let run_id = seed_legacy_run(&run_store).await;
3059
3060 let rc1 = RunContext::new(run_id.clone(), run_store.clone())
3064 .with_replay_store(replay_store.clone())
3065 .with_resume();
3066 let mut input1 = launch_input(blueprint.clone(), json!({ "input": "hi" }));
3067 input1.run_ctx = Some(rc1);
3068 let out1 = svc.launch(input1).await.expect("phase-1 resume launch ok");
3069 assert_eq!(out1.final_ctx["out"]["echoed"], "hi");
3070 assert_eq!(
3071 echo_calls.load(Ordering::SeqCst),
3072 1,
3073 "phase 1 dispatches the worker once (nothing to replay yet)"
3074 );
3075 let entries = replay_store
3076 .list_by_run(&run_id)
3077 .await
3078 .expect("list replay rows");
3079 assert_eq!(
3080 entries.len(),
3081 1,
3082 "phase 1 must log exactly one legacy-hashed replay row"
3083 );
3084 let run = run_store.get(&run_id).await.expect("run present");
3085 let snap: Value = serde_json::from_str(run.input_json.as_deref().unwrap()).unwrap();
3086 assert_eq!(
3087 SnapshotOrigin::from_snapshot(&snap),
3088 SnapshotOrigin::ResumeBackfill,
3089 "phase 1 must pin the snapshot as resume_backfill"
3090 );
3091
3092 let cursor = ReplayCursor::from_entries(entries);
3097 let rc2 = RunContext::new(run_id.clone(), run_store.clone())
3098 .with_replay_store(replay_store.clone())
3099 .with_replay_cursor(Arc::new(Mutex::new(cursor)))
3100 .with_resume();
3101 let mut input2 = launch_input(blueprint.clone(), json!({ "input": "hi" }));
3102 input2.run_ctx = Some(rc2);
3103 let out2 = svc.launch(input2).await.expect("phase-2 resume launch ok");
3104 assert_eq!(out2.final_ctx["out"]["echoed"], "hi");
3105 assert_eq!(
3106 echo_calls.load(Ordering::SeqCst),
3107 1,
3108 "phase 2 must REPLAY the legacy-hashed row — the worker must not run again"
3109 );
3110 let run2 = run_store.get(&run_id).await.expect("run present");
3111 let snap2: Value = serde_json::from_str(run2.input_json.as_deref().unwrap()).unwrap();
3112 assert_eq!(
3113 SnapshotOrigin::from_snapshot(&snap2),
3114 SnapshotOrigin::ResumeBackfill,
3115 "origin must stay resume_backfill across resumes (replay key stability)"
3116 );
3117 }
3118
3119 #[tokio::test]
3124 async fn launch_origin_run_uses_digest_keys_and_misses_legacy_replay_row() {
3125 use crate::store::replay::{InMemoryReplayStore, ReplayCursor, ReplayStore};
3126 use crate::store::run::{InMemoryRunStore, RunContext, RunStore};
3127 use std::sync::atomic::Ordering;
3128 use std::sync::Mutex;
3129
3130 let (svc, echo_calls) = counting_echo_service();
3131 let blueprint = bp(
3132 step("echo", path("$.input"), path("$.out")),
3133 vec![agent("echo", "echo")],
3134 );
3135 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
3136 let replay_store: Arc<dyn ReplayStore> = Arc::new(InMemoryReplayStore::new());
3137
3138 let backfill_run = seed_legacy_run(&run_store).await;
3140 let rc_bf = RunContext::new(backfill_run.clone(), run_store.clone())
3141 .with_replay_store(replay_store.clone())
3142 .with_resume();
3143 let mut input_bf = launch_input(blueprint.clone(), json!({ "input": "hi" }));
3144 input_bf.run_ctx = Some(rc_bf);
3145 svc.launch(input_bf).await.expect("backfill launch ok");
3146 assert_eq!(echo_calls.load(Ordering::SeqCst), 1);
3147 let legacy_entries = replay_store
3148 .list_by_run(&backfill_run)
3149 .await
3150 .expect("list legacy rows");
3151 assert_eq!(legacy_entries.len(), 1);
3152
3153 let launch_run = seed_legacy_run(&run_store).await;
3157 let cursor = ReplayCursor::from_entries(legacy_entries);
3158 let rc_launch = RunContext::new(launch_run.clone(), run_store.clone())
3159 .with_replay_store(replay_store.clone())
3160 .with_replay_cursor(Arc::new(Mutex::new(cursor)));
3161 let mut input_launch = launch_input(blueprint.clone(), json!({ "input": "hi" }));
3162 input_launch.run_ctx = Some(rc_launch);
3163 svc.launch(input_launch)
3164 .await
3165 .expect("launch-origin launch ok");
3166 assert_eq!(
3167 echo_calls.load(Ordering::SeqCst),
3168 2,
3169 "a launch-origin Run keys replay by binding digest, so the \
3170 legacy-hashed row must MISS and the worker must run"
3171 );
3172 let run = run_store.get(&launch_run).await.expect("run present");
3173 let snap: Value = serde_json::from_str(run.input_json.as_deref().unwrap()).unwrap();
3174 assert_eq!(SnapshotOrigin::from_snapshot(&snap), SnapshotOrigin::Launch);
3175 }
3176}