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: {message}")]
565 FlowEval {
566 message: String,
569 failed_step: Option<String>,
574 verdict_value: Option<Value>,
578 partial_ctx: Option<Value>,
586 },
587 #[error("pre-dispatch: {0}")]
594 PreDispatch(String),
595}
596
597#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize, schemars::JsonSchema)]
620pub struct TaskInputSpec {
621 #[serde(default)]
623 pub project_root: Option<String>,
624 #[serde(default)]
626 pub work_dir: Option<String>,
627 #[serde(default)]
629 #[schemars(with = "Option<Value>")]
630 pub task_metadata: Option<Value>,
631}
632
633#[derive(Debug, Clone)]
635pub struct TaskLaunchInput {
636 pub blueprint: Blueprint,
638 pub operator_id: String,
640 pub role: Role,
642 pub ttl: Duration,
644 pub operator_kind: Option<OperatorKind>,
653 pub bridge_id: Option<String>,
657 pub hook_id: Option<String>,
660 pub operator_backend_id: Option<String>,
667 pub operator_pin: Option<String>,
690 pub operator_kind_overrides: HashMap<String, OperatorKind>,
695 pub init_ctx: Value,
700 pub task_input: Option<TaskInputSpec>,
707 pub run_ctx: Option<RunContext>,
714 pub check_policy: Option<CheckPolicy>,
737}
738
739impl TaskLaunchInput {
740 pub fn automate(
750 blueprint: Blueprint,
751 operator_id: impl Into<String>,
752 role: Role,
753 ttl: Duration,
754 init_ctx: Value,
755 ) -> Self {
756 Self {
757 blueprint,
758 operator_id: operator_id.into(),
759 role,
760 ttl,
761 operator_kind: None,
762 bridge_id: None,
763 hook_id: None,
764 operator_backend_id: None,
765 operator_pin: None,
766 operator_kind_overrides: HashMap::new(),
767 init_ctx,
768 task_input: None,
769 run_ctx: None,
770 check_policy: None,
771 }
772 }
773}
774
775#[derive(Debug, Clone)]
777pub struct TaskLaunchOutput {
778 pub token: CapToken,
780 pub final_ctx: Value,
784}
785
786pub struct TaskLaunchService {
790 engine: Engine,
791 compiler: Compiler,
792 externs: Arc<dyn Externs + Send + Sync>,
797 binding_provider: Option<Arc<dyn AgentBindingProvider>>,
800 legacy_worker_binding_policy: LegacyWorkerBindingPolicy,
803}
804
805impl TaskLaunchService {
806 pub fn new(engine: Engine, compiler: Compiler) -> Self {
808 Self {
809 engine,
810 compiler,
811 externs: Arc::new(NoExterns),
812 binding_provider: None,
813 legacy_worker_binding_policy: LegacyWorkerBindingPolicy::default(),
814 }
815 }
816
817 pub fn with_externs(mut self, externs: Arc<dyn Externs + Send + Sync>) -> Self {
821 self.externs = externs;
822 self
823 }
824
825 pub fn with_binding_provider(mut self, provider: Arc<dyn AgentBindingProvider>) -> Self {
829 self.binding_provider = Some(provider);
830 self
831 }
832
833 pub fn with_legacy_worker_binding_policy(mut self, policy: LegacyWorkerBindingPolicy) -> Self {
835 self.legacy_worker_binding_policy = policy;
836 self
837 }
838
839 pub fn engine(&self) -> &Engine {
841 &self.engine
842 }
843
844 pub fn compiler(&self) -> &Compiler {
846 &self.compiler
847 }
848
849 pub async fn launch(
861 &self,
862 mut input: TaskLaunchInput,
863 ) -> Result<TaskLaunchOutput, TaskLaunchError> {
864 let pinned_binding_provider: Option<Arc<dyn AgentBindingProvider>> = input
878 .operator_pin
879 .as_deref()
880 .and_then(|pin| self.binding_provider.as_ref()?.pinned_to_session(pin));
881 let binding_provider = pinned_binding_provider
882 .as_deref()
883 .or(self.binding_provider.as_deref());
884 let (bound_agents, snapshot_origin) = load_or_resolve_bound_agents(
885 &input.blueprint,
886 input.run_ctx.as_ref(),
887 binding_provider,
888 self.legacy_worker_binding_policy,
889 )
890 .await?;
891 let binding_digests: HashMap<String, crate::blueprint::BindingDigest> = bound_agents
892 .iter()
893 .map(|bound| (bound.agent.name.clone(), bound.binding_digest.clone()))
894 .collect();
895 if let Some(run_ctx) = input.run_ctx.take() {
896 input.run_ctx = Some(match snapshot_origin {
907 SnapshotOrigin::Launch => run_ctx.with_binding_digests(binding_digests.clone()),
908 SnapshotOrigin::ResumeBackfill => run_ctx,
909 });
910 }
911 input.blueprint = materialize_bound_blueprint(&input.blueprint, &bound_agents);
912 let compiled = self.compiler.compile_bound_pinned(
913 &input.blueprint,
914 &bound_agents,
915 input.operator_pin.as_deref(),
916 )?;
917 self.engine
926 .register_verdict_contracts(compiled.router.verdict_contracts.clone());
927 let spawner = linker::link(
928 compiled.router.clone(),
929 &input.blueprint.spawner_hints.layers,
930 &self.engine,
931 );
932 let (agent_ctx_global, agent_ctx_per_agent) = derive_agent_ctx(&input.blueprint);
948 let (context_policy_default, context_policy_per_agent) =
949 derive_context_policies(&input.blueprint);
950 let spawner = SpawnerStack::new(spawner)
951 .layer(AgentContextMiddleware::new(
952 agent_ctx_global,
953 agent_ctx_per_agent,
954 context_policy_default,
955 context_policy_per_agent,
956 ))
957 .build();
958 let spawner = if let Some(alias) = input.blueprint.metadata.project_name_alias.as_deref() {
965 SpawnerStack::new(spawner)
966 .layer(ProjectNameAliasMiddleware::new(alias))
967 .build()
968 } else {
969 spawner
970 };
971 let worker_bindings = worker_bindings_from_bound_agents(&bound_agents);
975 let spawner = if worker_bindings.is_empty() {
976 spawner
977 } else {
978 SpawnerStack::new(spawner)
979 .layer(WorkerBindingMiddleware::new(worker_bindings))
980 .build()
981 };
982 let audit_defs = derive_audits(&input.blueprint);
992 let spawner = if audit_defs.is_empty() {
993 spawner
994 } else {
995 SpawnerStack::new(spawner)
996 .layer(AfterRunAuditMiddleware::new(
997 audit_defs,
998 compiled.router.clone(),
999 ))
1000 .build()
1001 };
1002
1003 let spawner = match input.task_input.as_ref().and_then(|spec| {
1010 TaskInputMiddleware::new_from_fields(
1011 spec.project_root.clone(),
1012 spec.work_dir.clone(),
1013 spec.task_metadata.clone(),
1014 )
1015 }) {
1016 Some(task_input) => SpawnerStack::new(spawner).layer(task_input).build(),
1017 None => spawner,
1018 };
1019
1020 let bp_agent_kinds = derive_bp_agent_kinds(&input.blueprint);
1025 let bp_global_kind = input
1026 .blueprint
1027 .default_operator_kind
1028 .map(OperatorKind::from);
1029
1030 let token = self
1031 .engine
1032 .attach_with_ids(
1033 input.operator_id,
1034 input.role,
1035 input.ttl,
1036 input.operator_kind,
1037 input.bridge_id,
1038 input.hook_id,
1039 input.operator_backend_id,
1040 input.operator_kind_overrides,
1041 bp_agent_kinds,
1042 bp_global_kind,
1043 )
1044 .await?;
1045 let resolved_check_policy = input.check_policy.or(input.blueprint.check_policy);
1056 let effective_check_policy =
1071 resolved_check_policy.unwrap_or(self.engine.cfg().check_policy);
1072 if effective_check_policy == CheckPolicy::Strict {
1073 let roots_missing = input
1074 .task_input
1075 .as_ref()
1076 .map(|t| t.project_root.is_none() && t.work_dir.is_none())
1077 .unwrap_or(true);
1078 if roots_missing {
1079 return Err(TaskLaunchError::PreDispatch(
1080 "check_policy=strict requires project_root or work_dir, but the launch \
1081 supplied neither"
1082 .to_string(),
1083 ));
1084 }
1085 }
1086 let dispatcher =
1087 EngineDispatcher::with_spawner(self.engine.clone(), token.clone(), spawner);
1088 let dispatcher = dispatcher.with_check_policy(resolved_check_policy);
1089 let map_err_run_ctx = input.run_ctx.clone();
1097 let dispatcher = match input.run_ctx {
1098 Some(run_ctx) => dispatcher.with_run(run_ctx),
1099 None => dispatcher,
1100 };
1101 let dispatcher = dispatcher.with_step_metas(derive_step_metas(&input.blueprint));
1105 let dispatcher = dispatcher.with_binding_digests(binding_digests);
1106 let dispatcher = dispatcher.with_step_naming(compiled.step_naming.clone());
1112 let dispatcher =
1119 dispatcher.with_projection_placement(compiled.projection_placement.clone());
1120 let merged_init_ctx =
1126 merge_init_ctx(input.blueprint.default_init_ctx.as_ref(), &input.init_ctx);
1127 let eval_result = mlua_flow_ir::eval_async_externs(
1128 &input.blueprint.flow,
1129 merged_init_ctx,
1130 &dispatcher,
1131 &*self.externs,
1132 )
1133 .await;
1134 let final_ctx = match eval_result {
1135 Ok(v) => v,
1136 Err(e) => {
1137 let (failed_step, verdict_value) = match &map_err_run_ctx {
1157 Some(rc) => {
1158 let slot = rc.last_failure.lock().ok().and_then(|g| g.clone());
1159 match slot {
1160 Some(lf) => (
1161 lf.step_ref.clone().or_else(|| Some(lf.step_id.to_string())),
1162 Some(lf.verdict_value.clone()),
1163 ),
1164 None => (None, None),
1165 }
1166 }
1167 None => (None, None),
1168 };
1169 let partial_ctx = match &map_err_run_ctx {
1170 Some(rc) => Some(rc.snapshot_partial_ctx().await),
1171 None => None,
1172 };
1173 return Err(TaskLaunchError::FlowEval {
1174 message: e.to_string(),
1175 failed_step,
1176 verdict_value,
1177 partial_ctx,
1178 });
1179 }
1180 };
1181 Ok(TaskLaunchOutput { token, final_ctx })
1182 }
1183}
1184
1185#[cfg(test)]
1190mod tests {
1191 use super::*;
1192 use crate::blueprint::compiler::{RustFnInProcessSpawnerFactory, SpawnerRegistry};
1193 use crate::blueprint::{
1194 current_schema_version, resolve_runner, AgentDef, AgentKind, AgentMeta, AgentProfile,
1195 BlueprintMetadata, CompilerHints, CompilerStrategy, MetaDef, Runner,
1196 };
1197 use crate::core::config::EngineCfg;
1198 use crate::worker::adapter::{WorkerError, WorkerResult};
1199 use mlua_flow_ir::{Expr, JoinMode, Node as FlowNode};
1200 use serde_json::json;
1201 use std::sync::Arc;
1202
1203 fn path(s: &str) -> Expr {
1204 Expr::Path {
1205 at: s.parse().expect("literal test path"),
1206 }
1207 }
1208 fn step(ref_: &str, in_: Expr, out: Expr) -> FlowNode {
1209 FlowNode::Step {
1210 ref_: ref_.to_string(),
1211 in_,
1212 out,
1213 }
1214 }
1215
1216 fn agent(name: &str, fn_id: &str) -> AgentDef {
1217 AgentDef {
1218 name: name.to_string(),
1219 kind: AgentKind::RustFn,
1220 spec: json!({ "fn_id": fn_id }),
1221 profile: None,
1222 meta: Some(AgentMeta::default()),
1223 runner: None,
1224 runner_ref: None,
1225 verdict: None,
1226 }
1227 }
1228
1229 fn build_service(factory: RustFnInProcessSpawnerFactory) -> TaskLaunchService {
1230 let engine = Engine::new(EngineCfg::default());
1231 let mut reg = SpawnerRegistry::new();
1232 reg.register::<RustFnInProcessSpawnerFactory>(Arc::new(factory));
1233 let compiler = Compiler::new(reg);
1234 TaskLaunchService::new(engine, compiler)
1235 }
1236
1237 fn build_service_with_cfg(
1241 factory: RustFnInProcessSpawnerFactory,
1242 cfg: EngineCfg,
1243 ) -> TaskLaunchService {
1244 let engine = Engine::new(cfg);
1245 let mut reg = SpawnerRegistry::new();
1246 reg.register::<RustFnInProcessSpawnerFactory>(Arc::new(factory));
1247 let compiler = Compiler::new(reg);
1248 TaskLaunchService::new(engine, compiler)
1249 }
1250
1251 fn bp(flow: FlowNode, agents: Vec<AgentDef>) -> Blueprint {
1252 Blueprint {
1253 schema_version: current_schema_version(),
1254 id: "ut".into(),
1255 flow,
1256 agents,
1257 operators: vec![],
1258 metas: vec![],
1259 hints: CompilerHints::default(),
1260 strategy: CompilerStrategy::default(),
1261 metadata: BlueprintMetadata::default(),
1262 spawner_hints: Default::default(),
1263 default_agent_kind: AgentKind::Operator,
1264 default_operator_kind: None,
1265 default_init_ctx: None,
1266 default_agent_ctx: None,
1267 default_context_policy: None,
1268 projection_placement: None,
1269 audits: vec![],
1270 degradation_policy: None,
1271 runners: vec![],
1272 default_runner: None,
1273 subprocesses: vec![],
1274 check_policy: None,
1275 blueprint_ref_includes: Vec::new(),
1276 }
1277 }
1278
1279 fn launch_input(blueprint: Blueprint, init_ctx: Value) -> TaskLaunchInput {
1280 TaskLaunchInput::automate(
1281 blueprint,
1282 "ut-op",
1283 Role::Operator,
1284 Duration::from_secs(30),
1285 init_ctx,
1286 )
1287 }
1288
1289 #[test]
1295 fn derive_audits_empty_by_default() {
1296 let blueprint = bp(
1297 step("echo", path("$.input"), path("$.out")),
1298 vec![agent("echo", "echo")],
1299 );
1300 assert!(
1301 derive_audits(&blueprint).is_empty(),
1302 "audits_absent_no_layer: an undeclared audits Vec must stay empty"
1303 );
1304 }
1305
1306 #[test]
1307 fn derive_audits_returns_blueprint_audits_verbatim() {
1308 let mut blueprint = bp(
1309 step("echo", path("$.input"), path("$.out")),
1310 vec![agent("echo", "echo")],
1311 );
1312 blueprint.audits = vec![crate::blueprint::AuditDef {
1313 agent: "auditor".to_string(),
1314 steps: None,
1315 mode: crate::blueprint::AuditMode::Async,
1316 }];
1317 let got = derive_audits(&blueprint);
1318 assert_eq!(got.len(), 1);
1319 assert_eq!(got[0].agent, "auditor");
1320 }
1321
1322 #[tokio::test]
1323 async fn launch_appends_audit_artifact_when_audits_declared() {
1324 use crate::blueprint::{AuditDef, AuditMode};
1325
1326 let factory = RustFnInProcessSpawnerFactory::new()
1327 .register_fn("echo", |inv| async move {
1328 Ok(WorkerResult {
1329 value: json!({ "echoed": inv.prompt }),
1330 ok: true,
1331 stats: None,
1332 })
1333 })
1334 .register_fn("audit-fn", |_inv| async move {
1335 Ok(WorkerResult {
1336 value: json!({ "finding": "clean" }),
1337 ok: true,
1338 stats: None,
1339 })
1340 });
1341 let svc = build_service(factory);
1342 let mut blueprint = bp(
1343 step("echo", path("$.input"), path("$.out")),
1344 vec![agent("echo", "echo"), agent("auditor", "audit-fn")],
1345 );
1346 blueprint.audits = vec![AuditDef {
1347 agent: "auditor".to_string(),
1348 steps: None,
1349 mode: AuditMode::Sync,
1350 }];
1351 let out = svc
1352 .launch(launch_input(blueprint, json!({ "input": "hi" })))
1353 .await
1354 .expect("launch ok — audits must never alter the audited step's outcome");
1355 assert_eq!(out.final_ctx["out"]["echoed"], "hi");
1356
1357 let audited_task_id = svc
1358 .engine()
1359 .with_state("test.find_audited_task", |s| {
1360 s.tasks
1361 .iter()
1362 .find(|(_, t)| t.spec.agent == "echo")
1363 .map(|(id, _)| id.clone())
1364 })
1365 .await
1366 .expect("with_state")
1367 .expect("the echo task must exist");
1368 let tail = svc.engine().output_tail(&audited_task_id, 1).await;
1369 let found = tail.iter().any(|ev| {
1370 matches!(
1371 ev,
1372 crate::worker::output::OutputEvent::Artifact { name, .. } if name == "audit:echo"
1373 )
1374 });
1375 assert!(
1376 found,
1377 "launch() must wire AfterRunAuditMiddleware end-to-end when Blueprint.audits is declared"
1378 );
1379 }
1380
1381 #[tokio::test]
1382 async fn launch_single_step_writes_out_path() {
1383 let factory = RustFnInProcessSpawnerFactory::new().register_fn("echo", |inv| async move {
1384 Ok(WorkerResult {
1385 value: json!({ "echoed": inv.prompt }),
1386 ok: true,
1387 stats: None,
1388 })
1389 });
1390 let svc = build_service(factory);
1391 let blueprint = bp(
1392 step("echo", path("$.input"), path("$.out")),
1393 vec![agent("echo", "echo")],
1394 );
1395 let out = svc
1396 .launch(launch_input(blueprint, json!({ "input": "hi" })))
1397 .await
1398 .expect("launch ok");
1399 assert_eq!(out.final_ctx["out"]["echoed"], "hi");
1400 }
1401
1402 async fn dispatched_check_policy(
1423 launch_policy: Option<CheckPolicy>,
1424 bp_policy: Option<CheckPolicy>,
1425 ) -> Option<CheckPolicy> {
1426 let factory = RustFnInProcessSpawnerFactory::new().register_fn("echo", |inv| async move {
1427 Ok(WorkerResult {
1428 value: json!({ "echoed": inv.prompt }),
1429 ok: true,
1430 stats: None,
1431 })
1432 });
1433 let svc = build_service(factory);
1434 let mut blueprint = bp(
1435 step("echo", path("$.input"), path("$.out")),
1436 vec![agent("echo", "echo")],
1437 );
1438 blueprint.check_policy = bp_policy;
1439 let mut input = launch_input(blueprint, json!({ "input": "hi" }));
1440 input.check_policy = launch_policy;
1441 input.task_input = Some(TaskInputSpec {
1442 project_root: None,
1443 work_dir: Some("/dispatched-check-policy-test-root".to_string()),
1444 task_metadata: None,
1445 });
1446 let _ = svc.launch(input).await;
1447 svc.engine()
1448 .with_state("test.read_dispatched_check_policy", |s| {
1449 s.tasks
1450 .values()
1451 .find(|t| t.spec.agent == "echo")
1452 .and_then(|t| t.spec.check_policy)
1453 })
1454 .await
1455 .expect("with_state")
1456 }
1457
1458 #[tokio::test]
1461 async fn cascade_launch_tier_wins_over_blueprint_tier() {
1462 assert_eq!(
1463 dispatched_check_policy(Some(CheckPolicy::Silent), Some(CheckPolicy::Strict)).await,
1464 Some(CheckPolicy::Silent),
1465 );
1466 }
1467
1468 #[tokio::test]
1471 async fn cascade_blueprint_tier_used_when_launch_absent() {
1472 assert_eq!(
1473 dispatched_check_policy(None, Some(CheckPolicy::Strict)).await,
1474 Some(CheckPolicy::Strict),
1475 );
1476 }
1477
1478 #[tokio::test]
1481 async fn cascade_launch_tier_alone_when_blueprint_absent() {
1482 assert_eq!(
1483 dispatched_check_policy(Some(CheckPolicy::Strict), None).await,
1484 Some(CheckPolicy::Strict),
1485 );
1486 }
1487
1488 #[tokio::test]
1493 async fn cascade_both_none_preserves_server_fallback() {
1494 assert_eq!(dispatched_check_policy(None, None).await, None);
1495 }
1496
1497 #[tokio::test]
1513 async fn strict_blueprint_without_roots_is_rejected_pre_dispatch() {
1514 let factory = RustFnInProcessSpawnerFactory::new().register_fn("echo", |inv| async move {
1515 Ok(WorkerResult {
1516 value: json!({ "echoed": inv.prompt }),
1517 ok: true,
1518 stats: None,
1519 })
1520 });
1521 let svc = build_service(factory);
1522 let mut blueprint = bp(
1523 step("echo", path("$.input"), path("$.out")),
1524 vec![agent("echo", "echo")],
1525 );
1526 blueprint.check_policy = Some(CheckPolicy::Strict);
1527 let err = svc
1529 .launch(launch_input(blueprint, json!({ "input": "hi" })))
1530 .await
1531 .expect_err("strict check_policy + no roots must be rejected before dispatch");
1532 match err {
1533 TaskLaunchError::PreDispatch(message) => {
1534 assert!(
1535 message.contains("strict"),
1536 "message must identify the strict-requires-roots condition: {message}"
1537 );
1538 }
1539 other => panic!("expected TaskLaunchError::PreDispatch, got {other:?}"),
1540 }
1541
1542 let dispatched = svc
1546 .engine()
1547 .with_state("test.no_echo_task_dispatched", |s| {
1548 s.tasks.values().any(|t| t.spec.agent == "echo")
1549 })
1550 .await
1551 .expect("with_state");
1552 assert!(
1553 !dispatched,
1554 "the pre-dispatch guard must reject before any step is dispatched"
1555 );
1556 }
1557
1558 #[tokio::test]
1570 async fn launch_without_any_check_policy_completes_fail_open() {
1571 let factory = RustFnInProcessSpawnerFactory::new().register_fn("echo", |inv| async move {
1572 Ok(WorkerResult {
1573 value: json!({ "echoed": inv.prompt }),
1574 ok: true,
1575 stats: None,
1576 })
1577 });
1578 let svc = build_service(factory);
1579 let blueprint = bp(
1580 step("echo", path("$.input"), path("$.out")),
1581 vec![agent("echo", "echo")],
1582 );
1583 assert_eq!(blueprint.check_policy, None, "BP tier must be unset");
1584 let out = svc
1585 .launch(launch_input(blueprint, json!({ "input": "hi" })))
1586 .await
1587 .expect("warn-mode fail-open must let the launch complete");
1588 assert_eq!(out.final_ctx["out"]["echoed"], "hi");
1589 }
1590
1591 #[tokio::test]
1612 async fn strict_blueprint_with_launch_warn_override_bypasses_pre_dispatch_guard() {
1613 let factory = RustFnInProcessSpawnerFactory::new().register_fn("echo", |inv| async move {
1614 Ok(WorkerResult {
1615 value: json!({ "echoed": inv.prompt }),
1616 ok: true,
1617 stats: None,
1618 })
1619 });
1620 let svc = build_service(factory);
1621 let mut blueprint = bp(
1622 step("echo", path("$.input"), path("$.out")),
1623 vec![agent("echo", "echo")],
1624 );
1625 blueprint.check_policy = Some(CheckPolicy::Strict);
1626 let mut input = launch_input(blueprint, json!({ "input": "hi" }));
1627 input.check_policy = Some(CheckPolicy::Warn);
1628 assert!(input.task_input.is_none(), "no roots supplied at all");
1629 let out = svc
1630 .launch(input)
1631 .await
1632 .expect("launch-tier warn override must bypass the pre-dispatch guard");
1633 assert_eq!(out.final_ctx["out"]["echoed"], "hi");
1634 }
1635
1636 #[tokio::test]
1642 async fn server_tier_strict_alone_triggers_pre_dispatch_guard() {
1643 let factory = RustFnInProcessSpawnerFactory::new().register_fn("echo", |inv| async move {
1644 Ok(WorkerResult {
1645 value: json!({ "echoed": inv.prompt }),
1646 ok: true,
1647 stats: None,
1648 })
1649 });
1650 let svc = build_service_with_cfg(
1651 factory,
1652 EngineCfg {
1653 check_policy: CheckPolicy::Strict,
1654 ..EngineCfg::default()
1655 },
1656 );
1657 let blueprint = bp(
1658 step("echo", path("$.input"), path("$.out")),
1659 vec![agent("echo", "echo")],
1660 );
1661 assert_eq!(blueprint.check_policy, None, "BP tier must be unset");
1662 let input = launch_input(blueprint, json!({ "input": "hi" }));
1663 assert!(input.check_policy.is_none(), "launch tier must be unset");
1664 assert!(input.task_input.is_none(), "no roots supplied");
1665 let err = svc.launch(input).await.expect_err(
1666 "server-tier Strict alone (BP/launch tiers both unset) must trigger the guard",
1667 );
1668 match err {
1669 TaskLaunchError::PreDispatch(message) => {
1670 assert!(
1671 message.contains("strict"),
1672 "expected the strict-requires-roots message, got: {message}"
1673 );
1674 }
1675 other => panic!("expected TaskLaunchError::PreDispatch, got {other:?}"),
1676 }
1677 }
1678
1679 #[tokio::test]
1685 async fn pre_dispatch_guard_rejects_when_task_input_present_but_roots_both_none() {
1686 let factory = RustFnInProcessSpawnerFactory::new().register_fn("echo", |inv| async move {
1687 Ok(WorkerResult {
1688 value: json!({ "echoed": inv.prompt }),
1689 ok: true,
1690 stats: None,
1691 })
1692 });
1693 let svc = build_service(factory);
1694 let mut blueprint = bp(
1695 step("echo", path("$.input"), path("$.out")),
1696 vec![agent("echo", "echo")],
1697 );
1698 blueprint.check_policy = Some(CheckPolicy::Strict);
1699 let mut input = launch_input(blueprint, json!({ "input": "hi" }));
1700 input.task_input = Some(TaskInputSpec {
1701 project_root: None,
1702 work_dir: None,
1703 task_metadata: Some(json!({ "unrelated": true })),
1704 });
1705 let err = svc
1706 .launch(input)
1707 .await
1708 .expect_err("Some(TaskInputSpec) with both roots None must still be roots_missing");
1709 assert!(
1710 matches!(err, TaskLaunchError::PreDispatch(_)),
1711 "expected TaskLaunchError::PreDispatch, got {err:?}"
1712 );
1713 }
1714
1715 #[tokio::test]
1720 async fn pre_dispatch_guard_passes_when_work_dir_present_and_project_root_absent() {
1721 let factory = RustFnInProcessSpawnerFactory::new().register_fn("echo", |inv| async move {
1722 Ok(WorkerResult {
1723 value: json!({ "echoed": inv.prompt }),
1724 ok: true,
1725 stats: None,
1726 })
1727 });
1728 let svc = build_service(factory);
1729 let mut blueprint = bp(
1730 step("echo", path("$.input"), path("$.out")),
1731 vec![agent("echo", "echo")],
1732 );
1733 blueprint.check_policy = Some(CheckPolicy::Strict);
1734 let mut input = launch_input(blueprint, json!({ "input": "hi" }));
1735 input.task_input = Some(TaskInputSpec {
1736 project_root: None,
1737 work_dir: Some("/repo/work".to_string()),
1738 task_metadata: None,
1739 });
1740 let out = svc
1741 .launch(input)
1742 .await
1743 .expect("work_dir alone must satisfy the guard's roots_missing check");
1744 assert_eq!(out.final_ctx["out"]["echoed"], "hi");
1745 }
1746
1747 #[tokio::test]
1748 async fn launch_three_step_seq_threads_ctx_forward() {
1749 let factory = RustFnInProcessSpawnerFactory::new()
1750 .register_fn("upper", |inv| async move {
1751 let s = serde_json::from_str::<String>(&inv.prompt).unwrap_or(inv.prompt);
1752 Ok(WorkerResult {
1753 value: json!(s.to_uppercase()),
1754 ok: true,
1755 stats: None,
1756 })
1757 })
1758 .register_fn("suffix", |inv| async move {
1759 let s = serde_json::from_str::<String>(&inv.prompt).unwrap_or(inv.prompt);
1760 Ok(WorkerResult {
1761 value: json!(format!("{s}!")),
1762 ok: true,
1763 stats: None,
1764 })
1765 })
1766 .register_fn("wrap", |inv| async move {
1767 let s = serde_json::from_str::<String>(&inv.prompt).unwrap_or(inv.prompt);
1768 Ok(WorkerResult {
1769 value: json!(format!("[{s}]")),
1770 ok: true,
1771 stats: None,
1772 })
1773 });
1774 let svc = build_service(factory);
1775 let flow = FlowNode::Seq {
1776 children: vec![
1777 step("upper", path("$.in"), path("$.s1")),
1778 step("suffix", path("$.s1"), path("$.s2")),
1779 step("wrap", path("$.s2"), path("$.s3")),
1780 ],
1781 };
1782 let blueprint = bp(
1783 flow,
1784 vec![
1785 agent("upper", "upper"),
1786 agent("suffix", "suffix"),
1787 agent("wrap", "wrap"),
1788 ],
1789 );
1790 let out = svc
1791 .launch(launch_input(blueprint, json!({ "in": "hello" })))
1792 .await
1793 .expect("launch ok");
1794 assert_eq!(out.final_ctx["s1"], "HELLO");
1795 assert_eq!(out.final_ctx["s2"], "HELLO!");
1796 assert_eq!(out.final_ctx["s3"], "[HELLO!]");
1797 }
1798
1799 #[tokio::test]
1800 async fn launch_fanout_join_all_parallel_completes() {
1801 use std::sync::atomic::{AtomicU32, Ordering};
1802 let counter = Arc::new(AtomicU32::new(0));
1803 let max_seen = Arc::new(AtomicU32::new(0));
1804 let counter_clone = counter.clone();
1805 let max_clone = max_seen.clone();
1806
1807 let factory = RustFnInProcessSpawnerFactory::new().register_fn("para", move |inv| {
1810 let counter = counter_clone.clone();
1811 let max_seen = max_clone.clone();
1812 async move {
1813 let now = counter.fetch_add(1, Ordering::SeqCst) + 1;
1814 let mut prev = max_seen.load(Ordering::SeqCst);
1815 while now > prev {
1816 match max_seen.compare_exchange(prev, now, Ordering::SeqCst, Ordering::SeqCst) {
1817 Ok(_) => break,
1818 Err(p) => prev = p,
1819 }
1820 }
1821 tokio::time::sleep(Duration::from_millis(50)).await;
1822 counter.fetch_sub(1, Ordering::SeqCst);
1823 let s = serde_json::from_str::<String>(&inv.prompt).unwrap_or(inv.prompt);
1824 Ok(WorkerResult {
1825 value: json!(format!("did:{s}")),
1826 ok: true,
1827 stats: None,
1828 })
1829 }
1830 });
1831 let svc = build_service(factory);
1832 let flow = FlowNode::Fanout {
1833 items: path("$.items"),
1834 bind: path("$.item"),
1835 body: Box::new(step("para", path("$.item"), path("$.r"))),
1836 join: JoinMode::All,
1837 out: path("$.results"),
1838 };
1839 let blueprint = bp(flow, vec![agent("para", "para")]);
1840 let out = svc
1841 .launch(launch_input(
1842 blueprint,
1843 json!({ "items": ["a", "b", "c", "d"] }),
1844 ))
1845 .await
1846 .expect("launch ok");
1847 let results = out.final_ctx["results"].as_array().expect("array");
1848 assert_eq!(results.len(), 4);
1849 for (i, expected) in ["a", "b", "c", "d"].iter().enumerate() {
1850 assert_eq!(results[i]["r"], json!(format!("did:{expected}")));
1851 }
1852 let max = max_seen.load(Ordering::SeqCst);
1853 assert!(
1854 max >= 2,
1855 "expected parallel execution (max inflight >= 2), got {max}"
1856 );
1857 }
1858
1859 #[tokio::test]
1860 async fn launch_propagates_worker_error_as_flow_eval_err() {
1861 let factory = RustFnInProcessSpawnerFactory::new()
1862 .register_fn("ok", |inv| async move {
1863 Ok(WorkerResult {
1864 value: json!(inv.prompt),
1865 ok: true,
1866 stats: None,
1867 })
1868 })
1869 .register_fn("boom", |_inv| async move {
1870 Err(WorkerError::Failed("intentional boom".into()))
1871 });
1872 let svc = build_service(factory);
1873 let flow = FlowNode::Seq {
1874 children: vec![
1875 step("ok", path("$.input"), path("$.s1")),
1876 step("boom", path("$.s1"), path("$.s2")),
1877 step("ok", path("$.s2"), path("$.s3")),
1878 ],
1879 };
1880 let blueprint = bp(flow, vec![agent("ok", "ok"), agent("boom", "boom")]);
1881 let err = svc
1882 .launch(launch_input(blueprint, json!({ "input": "x" })))
1883 .await
1884 .expect_err("expected fail");
1885 match err {
1886 TaskLaunchError::FlowEval { message: msg, .. } => {
1887 assert!(
1888 msg.contains("boom") || msg.contains("intentional"),
1889 "expected error to mention worker failure, got: {msg}"
1890 );
1891 }
1892 other => panic!("expected FlowEval error, got {other:?}"),
1893 }
1894 }
1895
1896 #[tokio::test]
1897 async fn launch_resolves_call_extern_via_registered_externs() {
1898 let factory = RustFnInProcessSpawnerFactory::new().register_fn("echo", |inv| async move {
1899 Ok(WorkerResult {
1900 value: json!({ "echoed": inv.prompt }),
1901 ok: true,
1902 stats: None,
1903 })
1904 });
1905 let mut externs = mlua_flow_ir::ExternMap::new();
1906 externs.register("fmt.greet", |args: &[Value]| {
1907 let name = args[0].as_str().unwrap_or("?");
1908 Ok(json!(format!("hello, {name}")))
1909 });
1910 let svc = build_service(factory).with_externs(Arc::new(externs));
1911 let flow = step(
1912 "echo",
1913 Expr::CallExtern {
1914 ref_: "fmt.greet".into(),
1915 args: vec![path("$.who")],
1916 },
1917 path("$.out"),
1918 );
1919 let blueprint = bp(flow, vec![agent("echo", "echo")]);
1920 let out = svc
1921 .launch(launch_input(blueprint, json!({ "who": "swarm" })))
1922 .await
1923 .expect("launch ok");
1924 assert_eq!(out.final_ctx["out"]["echoed"], json!("hello, swarm"));
1925 }
1926
1927 #[tokio::test]
1928 async fn launch_call_extern_without_registry_fails_as_flow_eval() {
1929 let factory = RustFnInProcessSpawnerFactory::new().register_fn("echo", |inv| async move {
1930 Ok(WorkerResult {
1931 value: json!(inv.prompt),
1932 ok: true,
1933 stats: None,
1934 })
1935 });
1936 let svc = build_service(factory); let flow = step(
1938 "echo",
1939 Expr::CallExtern {
1940 ref_: "fmt.greet".into(),
1941 args: vec![],
1942 },
1943 path("$.out"),
1944 );
1945 let blueprint = bp(flow, vec![agent("echo", "echo")]);
1946 let err = svc
1947 .launch(launch_input(blueprint, json!({})))
1948 .await
1949 .expect_err("expected fail");
1950 match err {
1951 TaskLaunchError::FlowEval { message: msg, .. } => {
1952 assert!(msg.contains("extern"), "expected extern error, got: {msg}");
1953 }
1954 other => panic!("expected FlowEval error, got {other:?}"),
1955 }
1956 }
1957
1958 #[tokio::test]
1978 async fn launch_registers_the_blueprints_verdict_contracts_into_the_engine() {
1979 let factory = RustFnInProcessSpawnerFactory::new().register_fn("gate", |inv| async move {
1980 Ok(WorkerResult {
1981 value: json!(inv.prompt),
1982 ok: true,
1983 stats: None,
1984 })
1985 });
1986 let svc = build_service(factory);
1987 let mut gate_agent = agent("gate", "gate");
1988 gate_agent.verdict = Some(mlua_swarm_schema::VerdictContract {
1989 channel: mlua_swarm_schema::VerdictChannel::Body,
1990 values: vec!["PASS".to_string(), "BLOCKED".to_string()],
1991 });
1992 let flow = step("gate", path("$.input"), path("$.out"));
1993 let blueprint = bp(flow, vec![gate_agent]);
1994
1995 let out = svc
1996 .launch(launch_input(blueprint, json!({ "input": "PASS" })))
1997 .await
1998 .expect("launch ok");
1999 assert_eq!(out.final_ctx["out"], json!("PASS"));
2000
2001 let task_id = svc
2006 .engine()
2007 .with_state("test.find_dispatched_task_id", |s| {
2008 s.tasks.keys().next().cloned()
2009 })
2010 .await
2011 .expect("with_state")
2012 .expect("launch must have dispatched exactly one Step (one TaskState)");
2013
2014 let contract = svc
2015 .engine()
2016 .verdict_contract_for_task(&task_id)
2017 .await
2018 .expect(
2019 "TaskLaunchService::launch must have merged this Blueprint's compiled \
2020 verdict_contracts into the engine's runtime registry \
2021 (Engine::register_verdict_contracts, called right after \
2022 compiler.compile succeeds) — verdict_contract_for_task resolving None \
2023 here means that production wiring regressed",
2024 );
2025 assert_eq!(contract.channel, mlua_swarm_schema::VerdictChannel::Body);
2026 assert_eq!(
2027 contract.values,
2028 vec!["PASS".to_string(), "BLOCKED".to_string()]
2029 );
2030 }
2031
2032 #[tokio::test]
2037 async fn launch_with_run_ctx_appends_one_step_entry_per_dispatched_step() {
2038 use crate::store::run::{InMemoryRunStore, RunContext, RunRecord, RunStatus, RunStore};
2039 use crate::types::{RunId, TaskId};
2040
2041 let factory = RustFnInProcessSpawnerFactory::new()
2042 .register_fn("upper", |inv| async move {
2043 Ok(WorkerResult {
2044 value: json!(inv.prompt.to_uppercase()),
2045 ok: true,
2046 stats: None,
2047 })
2048 })
2049 .register_fn("suffix", |inv| async move {
2050 let s = serde_json::from_str::<String>(&inv.prompt).unwrap_or(inv.prompt);
2051 Ok(WorkerResult {
2052 value: json!(format!("{s}!")),
2053 ok: true,
2054 stats: None,
2055 })
2056 });
2057 let svc = build_service(factory);
2058 let flow = FlowNode::Seq {
2059 children: vec![
2060 step("upper", path("$.in"), path("$.s1")),
2061 step("suffix", path("$.s1"), path("$.s2")),
2062 ],
2063 };
2064 let blueprint = bp(
2065 flow,
2066 vec![agent("upper", "upper"), agent("suffix", "suffix")],
2067 );
2068
2069 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
2070 let run_id = RunId::new();
2071 run_store
2072 .create(RunRecord {
2073 id: run_id.clone(),
2074 task_id: TaskId::new(),
2075 status: RunStatus::Running,
2076 step_entries: Vec::new(),
2077 degradations: Vec::new(),
2078 operator_sid: None,
2079 result_ref: None,
2080 input_json: Some("{}".to_string()),
2081 created_at: 0,
2082 updated_at: 0,
2083 })
2084 .await
2085 .expect("seed RunRecord");
2086
2087 let mut input = launch_input(blueprint, json!({ "in": "hi" }));
2088 input.run_ctx = Some(RunContext::new(run_id.clone(), run_store.clone()));
2089
2090 let out = svc.launch(input).await.expect("launch ok");
2091 assert_eq!(out.final_ctx["s2"], "HI!");
2092
2093 let run = run_store.get(&run_id).await.expect("run present");
2094 assert_eq!(
2095 run.step_entries.len(),
2096 2,
2097 "expected one step_entry per dispatched step, got {:?}",
2098 run.step_entries
2099 );
2100 assert_eq!(run.step_entries[0].step_ref, Some("upper".to_string()));
2101 assert_eq!(run.step_entries[0].status, Some("passed".to_string()));
2102 assert!(run.step_entries[0].binding_digest.is_some());
2103 assert_eq!(run.step_entries[1].step_ref, Some("suffix".to_string()));
2104 assert_eq!(run.step_entries[1].status, Some("passed".to_string()));
2105 assert!(run.step_entries[1].binding_digest.is_some());
2106 let snapshot: Value = serde_json::from_str(run.input_json.as_deref().unwrap()).unwrap();
2107 assert_eq!(snapshot["bound_agents"].as_array().unwrap().len(), 2);
2108 }
2109
2110 #[tokio::test]
2111 async fn run_snapshot_reuses_bound_agent_after_blueprint_mutation() {
2112 use crate::store::run::{InMemoryRunStore, RunContext, RunRecord, RunStatus, RunStore};
2113 use crate::types::{RunId, TaskId};
2114
2115 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
2116 let run_id = RunId::new();
2117 run_store
2118 .create(RunRecord {
2119 id: run_id.clone(),
2120 task_id: TaskId::new(),
2121 status: RunStatus::Running,
2122 step_entries: Vec::new(),
2123 degradations: Vec::new(),
2124 operator_sid: None,
2125 result_ref: None,
2126 input_json: Some("{}".to_string()),
2127 created_at: 0,
2128 updated_at: 0,
2129 })
2130 .await
2131 .unwrap();
2132 let run_ctx = RunContext::new(run_id, run_store);
2133 let mut original_agent = agent("worker", "worker");
2134 original_agent.profile = Some(crate::blueprint::AgentProfile {
2135 system_prompt: "original role".to_string(),
2136 ..Default::default()
2137 });
2138 let mut blueprint = bp(
2139 step("worker", path("$.input"), path("$.out")),
2140 vec![original_agent],
2141 );
2142
2143 let (original, _) = load_or_resolve_bound_agents(
2144 &blueprint,
2145 Some(&run_ctx),
2146 None,
2147 LegacyWorkerBindingPolicy::Allow,
2148 )
2149 .await
2150 .unwrap();
2151 blueprint.agents[0].profile.as_mut().unwrap().system_prompt = "mutated role".to_string();
2152 let (restored, _) = load_or_resolve_bound_agents(
2153 &blueprint,
2154 Some(&run_ctx),
2155 None,
2156 LegacyWorkerBindingPolicy::Allow,
2157 )
2158 .await
2159 .unwrap();
2160
2161 assert_eq!(restored[0].binding_digest, original[0].binding_digest);
2162 assert_eq!(
2163 restored[0].agent.profile.as_ref().unwrap().system_prompt,
2164 "original role"
2165 );
2166 }
2167
2168 #[tokio::test]
2169 async fn strict_migration_policy_rejects_fresh_legacy_worker_binding() {
2170 let mut legacy_agent = agent("worker", "worker");
2171 legacy_agent.profile = Some(AgentProfile {
2172 worker_binding: Some("legacy-worker".to_string()),
2173 ..Default::default()
2174 });
2175 let blueprint = bp(
2176 step("worker", path("$.input"), path("$.out")),
2177 vec![legacy_agent],
2178 );
2179
2180 let error =
2181 load_or_resolve_bound_agents(&blueprint, None, None, LegacyWorkerBindingPolicy::Reject)
2182 .await
2183 .expect_err("strict migration policy must reject fallback");
2184 assert!(error
2185 .to_string()
2186 .contains("deprecated profile.worker_binding"));
2187 }
2188
2189 #[tokio::test]
2190 async fn run_snapshot_calls_binding_provider_only_on_first_resolution() {
2191 use crate::binding::{AgentBindingProvider, BindingProviderError};
2192 use crate::blueprint::{BindOutcome, BindReceipt, BindRequest};
2193 use crate::store::run::{InMemoryRunStore, RunContext, RunRecord, RunStatus, RunStore};
2194 use crate::types::{RunId, TaskId};
2195 use std::sync::atomic::{AtomicUsize, Ordering};
2196
2197 struct CountingProvider(AtomicUsize);
2198
2199 #[async_trait::async_trait]
2200 impl AgentBindingProvider for CountingProvider {
2201 async fn bind(
2202 &self,
2203 requests: &[BindRequest],
2204 ) -> Result<Vec<BindOutcome>, BindingProviderError> {
2205 self.0.fetch_add(1, Ordering::SeqCst);
2206 Ok(requests
2207 .iter()
2208 .map(|request| BindOutcome::Bound {
2209 receipt: BindReceipt {
2210 agent: request.agent.clone(),
2211 request_digest: request.request_digest.clone(),
2212 provider_id: "operator-main-ai".to_string(),
2213 provider_revision: Some("test".to_string()),
2214 resolved_model: request.requested_model.clone(),
2215 effective_tools: request.requested_tools.clone(),
2216 launch_variant: request.launch_variant.clone(),
2217 capability_snapshot_digest: None,
2218 },
2219 })
2220 .collect())
2221 }
2222 }
2223
2224 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
2225 let run_id = RunId::new();
2226 run_store
2227 .create(RunRecord {
2228 id: run_id.clone(),
2229 task_id: TaskId::new(),
2230 status: RunStatus::Running,
2231 step_entries: Vec::new(),
2232 degradations: Vec::new(),
2233 operator_sid: None,
2234 result_ref: None,
2235 input_json: Some("{}".to_string()),
2236 created_at: 0,
2237 updated_at: 0,
2238 })
2239 .await
2240 .unwrap();
2241 let run_ctx = RunContext::new(run_id, run_store);
2242 let mut blueprint = bp(
2243 step("worker", path("$.input"), path("$.out")),
2244 vec![agent("worker", "worker")],
2245 );
2246 blueprint.agents[0].runner = Some(Runner::WsClaudeCode {
2247 variant: "mse-worker".to_string(),
2248 tools: vec!["Read".to_string()],
2249 });
2250 let provider = CountingProvider(AtomicUsize::new(0));
2251
2252 let (first, _) = load_or_resolve_bound_agents(
2253 &blueprint,
2254 Some(&run_ctx),
2255 Some(&provider),
2256 LegacyWorkerBindingPolicy::Allow,
2257 )
2258 .await
2259 .unwrap();
2260 let (restored, _) = load_or_resolve_bound_agents(
2261 &blueprint,
2262 Some(&run_ctx),
2263 Some(&provider),
2264 LegacyWorkerBindingPolicy::Allow,
2265 )
2266 .await
2267 .unwrap();
2268
2269 assert_eq!(provider.0.load(Ordering::SeqCst), 1);
2270 assert!(first[0].attestation.is_some());
2271 assert_eq!(restored, first);
2272 }
2273
2274 #[tokio::test]
2283 async fn fresh_resolve_on_launch_persists_launch_origin_no_degradation() {
2284 use crate::store::run::{
2285 InMemoryRunStore, RunContext, RunRecord, RunStatus, RunStore, BOUND_AGENTS_ORIGIN_KEY,
2286 };
2287 use crate::types::{RunId, TaskId};
2288
2289 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
2290 let run_id = RunId::new();
2291 run_store
2292 .create(RunRecord {
2293 id: run_id.clone(),
2294 task_id: TaskId::new(),
2295 status: RunStatus::Running,
2296 step_entries: Vec::new(),
2297 degradations: Vec::new(),
2298 operator_sid: None,
2299 result_ref: None,
2300 input_json: Some("{}".to_string()),
2301 created_at: 0,
2302 updated_at: 0,
2303 })
2304 .await
2305 .unwrap();
2306 let run_ctx = RunContext::new(run_id.clone(), run_store.clone());
2308 let blueprint = bp(
2309 step("worker", path("$.input"), path("$.out")),
2310 vec![agent("worker", "worker")],
2311 );
2312
2313 load_or_resolve_bound_agents(
2314 &blueprint,
2315 Some(&run_ctx),
2316 None,
2317 LegacyWorkerBindingPolicy::Allow,
2318 )
2319 .await
2320 .expect("launch resolve ok");
2321
2322 let run = run_store.get(&run_id).await.expect("run present");
2323 let snapshot: Value = serde_json::from_str(run.input_json.as_deref().unwrap()).unwrap();
2324 assert!(
2325 snapshot["bound_agents"].is_array(),
2326 "bound_agents must be persisted"
2327 );
2328 assert_eq!(snapshot[BOUND_AGENTS_ORIGIN_KEY], json!("launch"));
2329 assert_eq!(
2330 SnapshotOrigin::from_snapshot(&snapshot),
2331 SnapshotOrigin::Launch
2332 );
2333 assert!(
2334 run.degradations.is_empty(),
2335 "an initial-launch resolve is not a degradation"
2336 );
2337 }
2338
2339 #[tokio::test]
2344 async fn backfill_on_resume_persists_resume_origin_and_records_degradation() {
2345 use crate::store::run::{
2346 InMemoryRunStore, RunContext, RunRecord, RunStatus, RunStore, BOUND_AGENTS_ORIGIN_KEY,
2347 };
2348 use crate::types::{RunId, TaskId};
2349
2350 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
2351 let run_id = RunId::new();
2352 run_store
2353 .create(RunRecord {
2354 id: run_id.clone(),
2355 task_id: TaskId::new(),
2356 status: RunStatus::Running,
2357 step_entries: Vec::new(),
2358 degradations: Vec::new(),
2359 operator_sid: None,
2360 result_ref: None,
2361 input_json: Some("{}".to_string()),
2362 created_at: 0,
2363 updated_at: 0,
2364 })
2365 .await
2366 .unwrap();
2367 let run_ctx = RunContext::new(run_id.clone(), run_store.clone()).with_resume();
2368 let blueprint = bp(
2369 step("worker", path("$.input"), path("$.out")),
2370 vec![agent("worker", "worker")],
2371 );
2372
2373 load_or_resolve_bound_agents(
2374 &blueprint,
2375 Some(&run_ctx),
2376 None,
2377 LegacyWorkerBindingPolicy::Allow,
2378 )
2379 .await
2380 .expect("resume backfill ok");
2381
2382 let run = run_store.get(&run_id).await.expect("run present");
2383 let snapshot: Value = serde_json::from_str(run.input_json.as_deref().unwrap()).unwrap();
2384 assert_eq!(snapshot[BOUND_AGENTS_ORIGIN_KEY], json!("resume_backfill"));
2385 assert_eq!(
2386 SnapshotOrigin::from_snapshot(&snapshot),
2387 SnapshotOrigin::ResumeBackfill
2388 );
2389 assert_eq!(
2390 run.degradations.len(),
2391 1,
2392 "a resume backfill must record exactly one degradation"
2393 );
2394 assert_eq!(run.degradations[0].tool, "binding");
2395 assert_eq!(run.degradations[0].fallback, "resume_backfill");
2396 }
2397
2398 fn runner_blueprint(strict_binding: bool) -> Blueprint {
2405 let mut blueprint = bp(
2406 step("worker", path("$.input"), path("$.out")),
2407 vec![agent("worker", "worker")],
2408 );
2409 blueprint.strategy.strict_binding = strict_binding;
2410 blueprint.agents[0].runner = Some(Runner::WsClaudeCode {
2411 variant: "mse-worker".to_string(),
2412 tools: vec!["Read".to_string()],
2413 });
2414 blueprint
2415 }
2416
2417 struct AlwaysUnboundProvider;
2420
2421 #[async_trait::async_trait]
2422 impl AgentBindingProvider for AlwaysUnboundProvider {
2423 async fn bind(
2424 &self,
2425 requests: &[crate::blueprint::BindRequest],
2426 ) -> Result<Vec<crate::blueprint::BindOutcome>, crate::binding::BindingProviderError>
2427 {
2428 Ok(requests
2429 .iter()
2430 .map(|request| crate::blueprint::BindOutcome::Unbound {
2431 agent: request.agent.clone(),
2432 reason: "no capability manifest submitted".to_string(),
2433 })
2434 .collect())
2435 }
2436 }
2437
2438 #[tokio::test]
2442 async fn non_strict_unbound_agent_runs_declaration_only_with_degradation() {
2443 use crate::store::run::{InMemoryRunStore, RunContext, RunRecord, RunStatus, RunStore};
2444 use crate::types::{RunId, TaskId};
2445
2446 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
2447 let run_id = RunId::new();
2448 run_store
2449 .create(RunRecord {
2450 id: run_id.clone(),
2451 task_id: TaskId::new(),
2452 status: RunStatus::Running,
2453 step_entries: Vec::new(),
2454 degradations: Vec::new(),
2455 operator_sid: None,
2456 result_ref: None,
2457 input_json: Some("{}".to_string()),
2458 created_at: 0,
2459 updated_at: 0,
2460 })
2461 .await
2462 .unwrap();
2463 let run_ctx = RunContext::new(run_id.clone(), run_store.clone());
2464
2465 let (bound, _) = load_or_resolve_bound_agents(
2466 &runner_blueprint(false),
2467 Some(&run_ctx),
2468 Some(&AlwaysUnboundProvider),
2469 LegacyWorkerBindingPolicy::Allow,
2470 )
2471 .await
2472 .expect("non-strict launch must succeed even without an attestation");
2473 assert!(
2474 bound[0].attestation.is_none(),
2475 "an unattested agent must stay DeclarationOnly"
2476 );
2477
2478 let run = run_store.get(&run_id).await.expect("run present");
2479 assert_eq!(run.degradations.len(), 1, "expected one degradation entry");
2480 assert_eq!(run.degradations[0].tool, "binding");
2481 assert_eq!(run.degradations[0].fallback, "DeclarationOnly");
2482 assert!(run.degradations[0].error.contains("no capability manifest"));
2483 }
2484
2485 #[tokio::test]
2489 async fn strict_unbound_agent_fails_with_requirements_in_message() {
2490 let error = load_or_resolve_bound_agents(
2491 &runner_blueprint(true),
2492 None,
2493 Some(&AlwaysUnboundProvider),
2494 LegacyWorkerBindingPolicy::Allow,
2495 )
2496 .await
2497 .expect_err("strict + Unbound must reject the launch");
2498 match error {
2499 TaskLaunchError::PreDispatch(message) => {
2500 assert!(message.contains("worker"), "message: {message}");
2501 assert!(message.contains("mse-worker"), "message: {message}");
2502 assert!(message.contains("Read"), "message: {message}");
2503 }
2504 other => panic!("expected PreDispatch, got {other:?}"),
2505 }
2506 }
2507
2508 #[tokio::test]
2511 async fn strict_without_provider_rejects_runner_backed_launch() {
2512 let error = load_or_resolve_bound_agents(
2513 &runner_blueprint(true),
2514 None,
2515 None,
2516 LegacyWorkerBindingPolicy::Allow,
2517 )
2518 .await
2519 .expect_err("strict + no provider must reject a Runner-backed launch");
2520 match error {
2521 TaskLaunchError::PreDispatch(message) => {
2522 assert!(
2523 message.contains("strict_binding requires a binding provider"),
2524 "message: {message}"
2525 );
2526 }
2527 other => panic!("expected PreDispatch, got {other:?}"),
2528 }
2529 }
2530
2531 #[tokio::test]
2534 async fn strict_with_correct_manifest_attests_the_agent() {
2535 use crate::binding::ManifestBindingProvider;
2536 use crate::blueprint::{AgentProviderCapability, AgentProviderManifest};
2537
2538 let provider = ManifestBindingProvider::new(AgentProviderManifest {
2539 provider_id: "operator-main-ai".to_string(),
2540 provider_revision: Some("1".to_string()),
2541 capabilities: vec![AgentProviderCapability {
2542 launch_variant: Some("mse-worker".to_string()),
2543 resolved_model: None,
2544 effective_tools: vec!["Read".to_string()],
2545 capability_snapshot_digest: None,
2546 }],
2547 });
2548 let (bound, _) = load_or_resolve_bound_agents(
2549 &runner_blueprint(true),
2550 None,
2551 Some(&provider),
2552 LegacyWorkerBindingPolicy::Allow,
2553 )
2554 .await
2555 .expect("strict launch with a correct manifest must attest");
2556 assert!(
2557 bound[0].attestation.is_some(),
2558 "a correctly attested agent must carry its attestation"
2559 );
2560 }
2561
2562 #[test]
2571 fn worker_bindings_carry_request_digest_and_model() {
2572 let mut blueprint = runner_blueprint(false);
2573 blueprint.agents[0].profile = Some(AgentProfile {
2574 model: Some("claude-sonnet".to_string()),
2575 ..Default::default()
2576 });
2577 let bound = resolve_bound_agents(&blueprint).expect("resolvable Runner refs");
2578 let bindings = worker_bindings_from_bound_agents(&bound);
2579
2580 let wb = bindings.get("worker").expect("worker binding present");
2581 assert_eq!(
2582 wb.request_digest.as_ref(),
2583 Some(&bound[0].binding_digest),
2584 "the spawn frame must carry the immutable snapshot digest"
2585 );
2586 assert!(wb
2587 .request_digest
2588 .as_ref()
2589 .unwrap()
2590 .as_str()
2591 .starts_with("sha256:"));
2592 assert_eq!(wb.requested_model.as_deref(), Some("claude-sonnet"));
2593 }
2594
2595 #[test]
2598 fn worker_bindings_omit_model_when_profile_has_none() {
2599 let bound = resolve_bound_agents(&runner_blueprint(false)).expect("resolvable Runner refs");
2600 let bindings = worker_bindings_from_bound_agents(&bound);
2601 let wb = bindings.get("worker").expect("worker binding present");
2602 assert!(wb.request_digest.is_some());
2603 assert!(wb.requested_model.is_none());
2604 }
2605
2606 #[tokio::test]
2607 async fn launch_without_run_ctx_appends_no_step_entries() {
2608 let factory = RustFnInProcessSpawnerFactory::new().register_fn("echo", |inv| async move {
2612 Ok(WorkerResult {
2613 value: json!(inv.prompt),
2614 ok: true,
2615 stats: None,
2616 })
2617 });
2618 let svc = build_service(factory);
2619 let blueprint = bp(
2620 step("echo", path("$.input"), path("$.out")),
2621 vec![agent("echo", "echo")],
2622 );
2623 let input = launch_input(blueprint, json!({ "input": "hi" }));
2624 assert!(
2625 input.run_ctx.is_none(),
2626 "automate() defaults run_ctx to None"
2627 );
2628 let out = svc.launch(input).await.expect("launch ok");
2629 assert_eq!(out.final_ctx["out"], "hi");
2630 }
2631
2632 #[tokio::test]
2638 async fn launch_with_task_input_leaves_init_ctx_object_seed_unmutated() {
2639 let factory = RustFnInProcessSpawnerFactory::new().register_fn("echo", |inv| async move {
2646 Ok(WorkerResult {
2647 value: json!({ "echoed": inv.prompt }),
2648 ok: true,
2649 stats: None,
2650 })
2651 });
2652 let svc = build_service(factory);
2653 let blueprint = bp(
2654 step("echo", path("$.input"), path("$.out")),
2655 vec![agent("echo", "echo")],
2656 );
2657 let mut input = launch_input(blueprint, json!({ "input": "hi" }));
2658 input.task_input = Some(TaskInputSpec {
2659 project_root: Some("/repo".to_string()),
2660 work_dir: Some("/repo/work".to_string()),
2661 task_metadata: Some(json!({ "issue": 19 })),
2662 });
2663 let out = svc.launch(input).await.expect("launch ok");
2664 assert_eq!(out.final_ctx["out"]["echoed"], "hi");
2665 assert!(
2666 out.final_ctx.get("project_root").is_none(),
2667 "task_input must not be folded into the flow-ir ctx seed, got {:?}",
2668 out.final_ctx
2669 );
2670 assert!(out.final_ctx.get("work_dir").is_none());
2671 assert!(out.final_ctx.get("task_metadata").is_none());
2672 }
2673
2674 #[tokio::test]
2675 async fn launch_with_task_input_none_is_a_no_op() {
2676 let factory = RustFnInProcessSpawnerFactory::new().register_fn("echo", |inv| async move {
2677 Ok(WorkerResult {
2678 value: json!(inv.prompt),
2679 ok: true,
2680 stats: None,
2681 })
2682 });
2683 let svc = build_service(factory);
2684 let blueprint = bp(
2685 step("echo", path("$.input"), path("$.out")),
2686 vec![agent("echo", "echo")],
2687 );
2688 let mut input = launch_input(blueprint, json!({ "input": "hi" }));
2689 assert!(input.task_input.is_none(), "automate() defaults to None");
2690 input.task_input = None;
2691 let out = svc.launch(input).await.expect("launch ok");
2692 assert_eq!(out.final_ctx["out"], "hi");
2693 }
2694
2695 #[tokio::test]
2696 async fn launch_with_task_input_all_fields_absent_is_a_no_op() {
2697 let factory = RustFnInProcessSpawnerFactory::new().register_fn("echo", |inv| async move {
2701 Ok(WorkerResult {
2702 value: json!(inv.prompt),
2703 ok: true,
2704 stats: None,
2705 })
2706 });
2707 let svc = build_service(factory);
2708 let blueprint = bp(
2709 step("echo", path("$.input"), path("$.out")),
2710 vec![agent("echo", "echo")],
2711 );
2712 let mut input = launch_input(blueprint, json!({ "input": "hi" }));
2713 input.task_input = Some(TaskInputSpec::default());
2714 let out = svc.launch(input).await.expect("launch ok");
2715 assert_eq!(out.final_ctx["out"], "hi");
2716 }
2717
2718 #[test]
2723 fn merge_init_ctx_bp_default_only_passes_through_when_task_is_empty_object() {
2724 let bp_default = json!({ "seeded": "from-bp" });
2725 let task = json!({});
2726 let merged = merge_init_ctx(Some(&bp_default), &task);
2727 assert_eq!(merged, json!({ "seeded": "from-bp" }));
2728 }
2729
2730 #[test]
2731 fn merge_init_ctx_task_only_passes_through_when_bp_default_is_empty_object() {
2732 let bp_default = json!({});
2733 let task = json!({ "seeded": "from-task" });
2734 let merged = merge_init_ctx(Some(&bp_default), &task);
2735 assert_eq!(merged, json!({ "seeded": "from-task" }));
2736 }
2737
2738 #[test]
2739 fn merge_init_ctx_both_objects_task_wins_on_key_collision() {
2740 let bp_default = json!({ "a": "bp", "b": "bp-only" });
2741 let task = json!({ "a": "task", "c": "task-only" });
2742 let merged = merge_init_ctx(Some(&bp_default), &task);
2743 assert_eq!(
2744 merged,
2745 json!({ "a": "task", "b": "bp-only", "c": "task-only" })
2746 );
2747 }
2748
2749 #[test]
2750 fn merge_init_ctx_non_object_task_fully_replaces_bp_default() {
2751 let bp_default = json!({ "seeded": "from-bp" });
2752 let task = json!("plain-string-seed");
2753 let merged = merge_init_ctx(Some(&bp_default), &task);
2754 assert_eq!(merged, json!("plain-string-seed"));
2755 }
2756
2757 #[test]
2758 fn merge_init_ctx_no_bp_default_is_a_no_op() {
2759 let task = json!({ "input": "hi" });
2760 let merged = merge_init_ctx(None, &task);
2761 assert_eq!(merged, task);
2762 }
2763
2764 #[test]
2769 fn merge_init_ctx_3layer_no_run_override_equals_bp_task_merge_only() {
2770 let bp_default = json!({ "a": "bp", "b": "bp-only" });
2774 let task = json!({ "a": "task", "c": "task-only" });
2775 let three_layer = merge_init_ctx_3layer(Some(&bp_default), &task, None);
2776 let two_layer = merge_init_ctx(Some(&bp_default), &task);
2777 assert_eq!(three_layer, two_layer);
2778 assert_eq!(
2779 three_layer,
2780 json!({ "a": "task", "b": "bp-only", "c": "task-only" })
2781 );
2782 }
2783
2784 #[test]
2785 fn merge_init_ctx_3layer_run_object_wins_on_key_collision_over_bp_and_task() {
2786 let bp_default = json!({ "a": "bp", "b": "bp-only" });
2787 let task = json!({ "a": "task", "c": "task-only" });
2788 let run_override = json!({ "a": "run", "d": "run-only" });
2789 let merged = merge_init_ctx_3layer(Some(&bp_default), &task, Some(&run_override));
2790 assert_eq!(
2791 merged,
2792 json!({ "a": "run", "b": "bp-only", "c": "task-only", "d": "run-only" }),
2793 "Run wins on collision (a); BP-only (b) and Task-only (c) keys survive"
2794 );
2795 }
2796
2797 #[test]
2798 fn merge_init_ctx_3layer_run_non_object_fully_replaces_bp_task_merge() {
2799 let bp_default = json!({ "seeded": "from-bp" });
2800 let task = json!({ "seeded": "from-task" });
2801 let run_override = json!("plain-string-run-seed");
2802 let merged = merge_init_ctx_3layer(Some(&bp_default), &task, Some(&run_override));
2803 assert_eq!(merged, json!("plain-string-run-seed"));
2804 }
2805
2806 #[test]
2807 fn merge_init_ctx_3layer_no_bp_default_and_no_run_override_is_task_passthrough() {
2808 let task = json!({ "input": "hi" });
2809 let merged = merge_init_ctx_3layer(None, &task, None);
2810 assert_eq!(merged, task);
2811 }
2812
2813 #[tokio::test]
2814 async fn launch_merges_bp_default_init_ctx_into_task_init_ctx() {
2815 let factory = RustFnInProcessSpawnerFactory::new().register_fn("echo", |inv| async move {
2818 Ok(WorkerResult {
2819 value: json!(inv.prompt),
2820 ok: true,
2821 stats: None,
2822 })
2823 });
2824 let svc = build_service(factory);
2825 let mut blueprint = bp(
2826 step("echo", path("$.greeting"), path("$.out")),
2827 vec![agent("echo", "echo")],
2828 );
2829 blueprint.default_init_ctx = Some(json!({ "greeting": "hello from bp" }));
2830 let out = svc
2832 .launch(launch_input(blueprint, json!({})))
2833 .await
2834 .expect("launch ok");
2835 assert_eq!(out.final_ctx["out"], "hello from bp");
2836 }
2837
2838 fn agent_with_meta(name: &str, fn_id: &str, meta: AgentMeta) -> AgentDef {
2843 AgentDef {
2844 name: name.to_string(),
2845 kind: AgentKind::RustFn,
2846 spec: json!({ "fn_id": fn_id }),
2847 profile: None,
2848 meta: Some(meta),
2849 runner: None,
2850 runner_ref: None,
2851 verdict: None,
2852 }
2853 }
2854
2855 #[test]
2856 fn derive_agent_ctx_empty_blueprint_yields_empty_state() {
2857 let blueprint = bp(step("echo", path("$.in"), path("$.out")), vec![]);
2858 let (global, per_agent) = derive_agent_ctx(&blueprint);
2859 assert_eq!(global, None);
2860 assert!(per_agent.is_empty());
2861 }
2862
2863 #[test]
2864 fn derive_agent_ctx_populated_blueprint_yields_correct_maps() {
2865 let mut blueprint = bp(
2866 step("echo", path("$.in"), path("$.out")),
2867 vec![
2868 agent_with_meta(
2869 "with-ctx",
2870 "echo",
2871 AgentMeta {
2872 ctx: Some(json!({ "org_conventions": "x" })),
2873 ..Default::default()
2874 },
2875 ),
2876 agent("no-ctx", "echo"),
2877 ],
2878 );
2879 blueprint.default_agent_ctx = Some(json!({ "seeded": "from-bp" }));
2880 let (global, per_agent) = derive_agent_ctx(&blueprint);
2881 assert_eq!(global, Some(json!({ "seeded": "from-bp" })));
2882 assert_eq!(
2883 per_agent.len(),
2884 1,
2885 "agents without AgentMeta.ctx are absent, not defaulted to null: {per_agent:?}"
2886 );
2887 assert_eq!(
2888 per_agent.get("with-ctx"),
2889 Some(&json!({ "org_conventions": "x" }))
2890 );
2891 assert!(!per_agent.contains_key("no-ctx"));
2892 }
2893
2894 #[test]
2895 fn derive_context_policies_empty_blueprint_yields_empty_state() {
2896 let blueprint = bp(step("echo", path("$.in"), path("$.out")), vec![]);
2897 let (default_policy, per_agent) = derive_context_policies(&blueprint);
2898 assert_eq!(default_policy, None);
2899 assert!(per_agent.is_empty());
2900 }
2901
2902 #[test]
2903 fn derive_context_policies_populated_blueprint_yields_correct_maps() {
2904 let mut blueprint = bp(
2905 step("echo", path("$.in"), path("$.out")),
2906 vec![
2907 agent_with_meta(
2908 "with-policy",
2909 "echo",
2910 AgentMeta {
2911 context_policy: Some(ContextPolicy {
2912 include: None,
2913 exclude: vec!["work_dir".to_string()],
2914 ..Default::default()
2915 }),
2916 ..Default::default()
2917 },
2918 ),
2919 agent("no-policy", "echo"),
2920 ],
2921 );
2922 blueprint.default_context_policy = Some(ContextPolicy {
2923 include: Some(vec!["project_root".to_string()]),
2924 exclude: vec![],
2925 ..Default::default()
2926 });
2927 let (default_policy, per_agent) = derive_context_policies(&blueprint);
2928 assert_eq!(
2929 default_policy,
2930 Some(ContextPolicy {
2931 include: Some(vec!["project_root".to_string()]),
2932 exclude: vec![],
2933 ..Default::default()
2934 })
2935 );
2936 assert_eq!(per_agent.len(), 1);
2937 assert_eq!(
2938 per_agent.get("with-policy"),
2939 Some(&ContextPolicy {
2940 include: None,
2941 exclude: vec!["work_dir".to_string()],
2942 ..Default::default()
2943 })
2944 );
2945 assert!(!per_agent.contains_key("no-policy"));
2946 }
2947
2948 #[test]
2954 fn derive_step_metas_empty_blueprint_yields_empty_map() {
2955 let blueprint = bp(step("echo", path("$.in"), path("$.out")), vec![]);
2956 assert!(derive_step_metas(&blueprint).is_empty());
2957 }
2958
2959 #[test]
2960 fn derive_step_metas_populated_blueprint_yields_name_to_ctx_map() {
2961 let mut blueprint = bp(step("echo", path("$.in"), path("$.out")), vec![]);
2962 blueprint.metas = vec![
2963 MetaDef {
2964 name: "heavy-scan".to_string(),
2965 ctx: json!({ "work_dir": "/x" }),
2966 },
2967 MetaDef {
2968 name: "light-scan".to_string(),
2969 ctx: json!({ "work_dir": "/y" }),
2970 },
2971 ];
2972 let metas = derive_step_metas(&blueprint);
2973 assert_eq!(metas.len(), 2);
2974 assert_eq!(metas.get("heavy-scan"), Some(&json!({ "work_dir": "/x" })));
2975 assert_eq!(metas.get("light-scan"), Some(&json!({ "work_dir": "/y" })));
2976 }
2977
2978 #[test]
2979 fn derive_agent_ctx_meta_ref_resolves_as_base_under_inline_ctx() {
2980 let mut blueprint = bp(
2981 step("echo", path("$.in"), path("$.out")),
2982 vec![agent_with_meta(
2983 "with-meta-ref",
2984 "echo",
2985 AgentMeta {
2986 ctx: Some(json!({ "work_dir": "/inline-wins" })),
2987 meta_ref: Some("shared".to_string()),
2988 ..Default::default()
2989 },
2990 )],
2991 );
2992 blueprint.metas = vec![MetaDef {
2993 name: "shared".to_string(),
2994 ctx: json!({ "work_dir": "/base", "extra": "from-pool" }),
2995 }];
2996 let (_, per_agent) = derive_agent_ctx(&blueprint);
2997 assert_eq!(
2998 per_agent.get("with-meta-ref"),
2999 Some(&json!({ "work_dir": "/inline-wins", "extra": "from-pool" })),
3000 "inline ctx must win the collided key while pool-only keys survive the merge"
3001 );
3002 }
3003
3004 #[test]
3005 fn derive_agent_ctx_meta_ref_alone_uses_pool_ctx_verbatim() {
3006 let mut blueprint = bp(
3007 step("echo", path("$.in"), path("$.out")),
3008 vec![agent_with_meta(
3009 "with-meta-ref-only",
3010 "echo",
3011 AgentMeta {
3012 meta_ref: Some("shared".to_string()),
3013 ..Default::default()
3014 },
3015 )],
3016 );
3017 blueprint.metas = vec![MetaDef {
3018 name: "shared".to_string(),
3019 ctx: json!({ "work_dir": "/base" }),
3020 }];
3021 let (_, per_agent) = derive_agent_ctx(&blueprint);
3022 assert_eq!(
3023 per_agent.get("with-meta-ref-only"),
3024 Some(&json!({ "work_dir": "/base" }))
3025 );
3026 }
3027
3028 #[test]
3029 fn derive_agent_ctx_unresolved_meta_ref_never_panics_and_falls_back_to_inline() {
3030 let blueprint = bp(
3031 step("echo", path("$.in"), path("$.out")),
3032 vec![agent_with_meta(
3033 "with-unresolved-meta-ref",
3034 "echo",
3035 AgentMeta {
3036 ctx: Some(json!({ "work_dir": "/inline-only" })),
3037 meta_ref: Some("missing".to_string()),
3038 ..Default::default()
3039 },
3040 )],
3041 );
3042 let (_, per_agent) = derive_agent_ctx(&blueprint);
3044 assert_eq!(
3045 per_agent.get("with-unresolved-meta-ref"),
3046 Some(&json!({ "work_dir": "/inline-only" })),
3047 "an unresolved meta_ref must never panic; the agent's own inline ctx still applies"
3048 );
3049 }
3050
3051 #[test]
3067 fn resolve_runner_legacy_fallback_matches_derive_worker_bindings_semantics() {
3068 fn legacy_agent(name: &str, variant: &str, tools: Vec<&str>) -> AgentDef {
3069 AgentDef {
3070 name: name.to_string(),
3071 kind: AgentKind::Operator,
3072 spec: json!({}),
3073 profile: Some(AgentProfile {
3074 worker_binding: Some(variant.to_string()),
3075 tools: tools.into_iter().map(str::to_string).collect(),
3076 ..Default::default()
3077 }),
3078 meta: None,
3079 runner: None,
3080 runner_ref: None,
3081 verdict: None,
3082 }
3083 }
3084
3085 let blueprint = bp(
3086 step("planner", path("$.in"), path("$.out")),
3087 vec![
3088 legacy_agent("planner", "planning-worker", vec!["Read", "Grep"]),
3089 legacy_agent("coder", "code-worker", vec![]),
3090 agent("no-binding", "echo"),
3091 ],
3092 );
3093
3094 let derived = derive_worker_bindings(&blueprint);
3095
3096 for agent_def in &blueprint.agents {
3097 let resolved = resolve_runner(&blueprint, agent_def).expect("no unresolved refs");
3098 match derived.get(&agent_def.name) {
3099 Some(binding) => {
3100 assert_eq!(
3101 resolved,
3102 Some(Runner::WsClaudeCode {
3103 variant: binding.variant.clone(),
3104 tools: binding.tools.clone(),
3105 }),
3106 "resolve_runner must synthesize the same WsClaudeCode Runner \
3107 derive_worker_bindings produces for agent '{}'",
3108 agent_def.name
3109 );
3110 }
3111 None => {
3112 assert_eq!(
3113 resolved, None,
3114 "agent '{}' has no derive_worker_bindings entry, so resolve_runner \
3115 must resolve to None too (no other tier declared)",
3116 agent_def.name
3117 );
3118 }
3119 }
3120 }
3121 }
3122
3123 #[test]
3124 fn ws_operator_runner_projects_into_the_existing_spawn_binding() {
3125 let mut blueprint = bp(
3126 step("reviewer", path("$.in"), path("$.out")),
3127 vec![agent("reviewer", "echo")],
3128 );
3129 blueprint.agents[0].runner = Some(Runner::WsOperator {
3130 variant: "mse-reviewer".to_string(),
3131 tools: vec!["Read".to_string(), "Grep".to_string()],
3132 });
3133
3134 let derived = derive_worker_bindings(&blueprint);
3135 let binding = derived
3136 .get("reviewer")
3137 .expect("ws_operator must feed the canonical spawn binding path");
3138 assert_eq!(binding.variant, "mse-reviewer");
3139 assert_eq!(binding.tools, ["Read", "Grep"]);
3140 }
3141
3142 fn counting_echo_service() -> (TaskLaunchService, Arc<std::sync::atomic::AtomicUsize>) {
3151 use std::sync::atomic::{AtomicUsize, Ordering};
3152 let calls = Arc::new(AtomicUsize::new(0));
3153 let counter = calls.clone();
3154 let factory = RustFnInProcessSpawnerFactory::new().register_fn("echo", move |inv| {
3155 let counter = counter.clone();
3156 async move {
3157 counter.fetch_add(1, Ordering::SeqCst);
3158 Ok(WorkerResult {
3159 value: json!({ "echoed": inv.prompt }),
3160 ok: true,
3161 stats: None,
3162 })
3163 }
3164 });
3165 (build_service(factory), calls)
3166 }
3167
3168 async fn seed_legacy_run(run_store: &Arc<dyn crate::store::run::RunStore>) -> crate::RunId {
3169 use crate::store::run::{RunRecord, RunStatus};
3170 use crate::types::TaskId;
3171 let run_id = crate::RunId::new();
3172 run_store
3173 .create(RunRecord {
3174 id: run_id.clone(),
3175 task_id: TaskId::new(),
3176 status: RunStatus::Running,
3177 step_entries: Vec::new(),
3178 degradations: Vec::new(),
3179 operator_sid: None,
3180 result_ref: None,
3181 input_json: Some("{}".to_string()),
3184 created_at: 0,
3185 updated_at: 0,
3186 })
3187 .await
3188 .expect("seed legacy RunRecord");
3189 run_id
3190 }
3191
3192 #[tokio::test]
3198 async fn backfilled_run_replays_legacy_keys_stably_across_two_resumes() {
3199 use crate::store::replay::{InMemoryReplayStore, ReplayCursor, ReplayStore};
3200 use crate::store::run::{InMemoryRunStore, RunContext, RunStore};
3201 use std::sync::atomic::Ordering;
3202 use std::sync::Mutex;
3203
3204 let (svc, echo_calls) = counting_echo_service();
3205 let blueprint = bp(
3206 step("echo", path("$.input"), path("$.out")),
3207 vec![agent("echo", "echo")],
3208 );
3209 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
3210 let replay_store: Arc<dyn ReplayStore> = Arc::new(InMemoryReplayStore::new());
3211 let run_id = seed_legacy_run(&run_store).await;
3212
3213 let rc1 = RunContext::new(run_id.clone(), run_store.clone())
3217 .with_replay_store(replay_store.clone())
3218 .with_resume();
3219 let mut input1 = launch_input(blueprint.clone(), json!({ "input": "hi" }));
3220 input1.run_ctx = Some(rc1);
3221 let out1 = svc.launch(input1).await.expect("phase-1 resume launch ok");
3222 assert_eq!(out1.final_ctx["out"]["echoed"], "hi");
3223 assert_eq!(
3224 echo_calls.load(Ordering::SeqCst),
3225 1,
3226 "phase 1 dispatches the worker once (nothing to replay yet)"
3227 );
3228 let entries = replay_store
3229 .list_by_run(&run_id)
3230 .await
3231 .expect("list replay rows");
3232 assert_eq!(
3233 entries.len(),
3234 1,
3235 "phase 1 must log exactly one legacy-hashed replay row"
3236 );
3237 let run = run_store.get(&run_id).await.expect("run present");
3238 let snap: Value = serde_json::from_str(run.input_json.as_deref().unwrap()).unwrap();
3239 assert_eq!(
3240 SnapshotOrigin::from_snapshot(&snap),
3241 SnapshotOrigin::ResumeBackfill,
3242 "phase 1 must pin the snapshot as resume_backfill"
3243 );
3244
3245 let cursor = ReplayCursor::from_entries(entries);
3250 let rc2 = RunContext::new(run_id.clone(), run_store.clone())
3251 .with_replay_store(replay_store.clone())
3252 .with_replay_cursor(Arc::new(Mutex::new(cursor)))
3253 .with_resume();
3254 let mut input2 = launch_input(blueprint.clone(), json!({ "input": "hi" }));
3255 input2.run_ctx = Some(rc2);
3256 let out2 = svc.launch(input2).await.expect("phase-2 resume launch ok");
3257 assert_eq!(out2.final_ctx["out"]["echoed"], "hi");
3258 assert_eq!(
3259 echo_calls.load(Ordering::SeqCst),
3260 1,
3261 "phase 2 must REPLAY the legacy-hashed row — the worker must not run again"
3262 );
3263 let run2 = run_store.get(&run_id).await.expect("run present");
3264 let snap2: Value = serde_json::from_str(run2.input_json.as_deref().unwrap()).unwrap();
3265 assert_eq!(
3266 SnapshotOrigin::from_snapshot(&snap2),
3267 SnapshotOrigin::ResumeBackfill,
3268 "origin must stay resume_backfill across resumes (replay key stability)"
3269 );
3270 }
3271
3272 #[tokio::test]
3277 async fn launch_origin_run_uses_digest_keys_and_misses_legacy_replay_row() {
3278 use crate::store::replay::{InMemoryReplayStore, ReplayCursor, ReplayStore};
3279 use crate::store::run::{InMemoryRunStore, RunContext, RunStore};
3280 use std::sync::atomic::Ordering;
3281 use std::sync::Mutex;
3282
3283 let (svc, echo_calls) = counting_echo_service();
3284 let blueprint = bp(
3285 step("echo", path("$.input"), path("$.out")),
3286 vec![agent("echo", "echo")],
3287 );
3288 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
3289 let replay_store: Arc<dyn ReplayStore> = Arc::new(InMemoryReplayStore::new());
3290
3291 let backfill_run = seed_legacy_run(&run_store).await;
3293 let rc_bf = RunContext::new(backfill_run.clone(), run_store.clone())
3294 .with_replay_store(replay_store.clone())
3295 .with_resume();
3296 let mut input_bf = launch_input(blueprint.clone(), json!({ "input": "hi" }));
3297 input_bf.run_ctx = Some(rc_bf);
3298 svc.launch(input_bf).await.expect("backfill launch ok");
3299 assert_eq!(echo_calls.load(Ordering::SeqCst), 1);
3300 let legacy_entries = replay_store
3301 .list_by_run(&backfill_run)
3302 .await
3303 .expect("list legacy rows");
3304 assert_eq!(legacy_entries.len(), 1);
3305
3306 let launch_run = seed_legacy_run(&run_store).await;
3310 let cursor = ReplayCursor::from_entries(legacy_entries);
3311 let rc_launch = RunContext::new(launch_run.clone(), run_store.clone())
3312 .with_replay_store(replay_store.clone())
3313 .with_replay_cursor(Arc::new(Mutex::new(cursor)));
3314 let mut input_launch = launch_input(blueprint.clone(), json!({ "input": "hi" }));
3315 input_launch.run_ctx = Some(rc_launch);
3316 svc.launch(input_launch)
3317 .await
3318 .expect("launch-origin launch ok");
3319 assert_eq!(
3320 echo_calls.load(Ordering::SeqCst),
3321 2,
3322 "a launch-origin Run keys replay by binding digest, so the \
3323 legacy-hashed row must MISS and the worker must run"
3324 );
3325 let run = run_store.get(&launch_run).await.expect("run present");
3326 let snap: Value = serde_json::from_str(run.input_json.as_deref().unwrap()).unwrap();
3327 assert_eq!(SnapshotOrigin::from_snapshot(&snap), SnapshotOrigin::Launch);
3328 }
3329
3330 #[test]
3338 fn task_launch_error_flow_eval_struct_variant_display_preserves_prefix() {
3339 let err = TaskLaunchError::FlowEval {
3340 message: "dispatcher error at ref foo".to_string(),
3341 failed_step: Some("foo".to_string()),
3342 verdict_value: Some(json!({"verdict": "BLOCKED"})),
3343 partial_ctx: Some(json!({"steps": {}})),
3344 };
3345 assert_eq!(err.to_string(), "flow eval: dispatcher error at ref foo");
3346
3347 let err_bare = TaskLaunchError::FlowEval {
3351 message: "unresolved extern".to_string(),
3352 failed_step: None,
3353 verdict_value: None,
3354 partial_ctx: None,
3355 };
3356 assert_eq!(err_bare.to_string(), "flow eval: unresolved extern");
3357 }
3358
3359 #[tokio::test]
3364 async fn task_launch_flow_eval_error_carries_failed_step_and_verdict_value() {
3365 use crate::store::run::{InMemoryRunStore, RunContext, RunRecord, RunStatus, RunStore};
3366 use crate::types::{RunId, TaskId};
3367
3368 let factory = RustFnInProcessSpawnerFactory::new().register_fn("gate", |_inv| async move {
3369 Ok(WorkerResult {
3370 value: json!({ "verdict": "BLOCKED", "reason": "not applicable" }),
3371 ok: false,
3372 stats: None,
3373 })
3374 });
3375 let svc = build_service(factory);
3376 let blueprint = bp(
3377 step("gate", path("$.input"), path("$.out")),
3378 vec![agent("gate", "gate")],
3379 );
3380
3381 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
3382 let run_id = RunId::new();
3383 run_store
3384 .create(RunRecord {
3385 id: run_id.clone(),
3386 task_id: TaskId::new(),
3387 status: RunStatus::Running,
3388 step_entries: Vec::new(),
3389 degradations: Vec::new(),
3390 operator_sid: None,
3391 result_ref: None,
3392 input_json: Some("{}".to_string()),
3393 created_at: 0,
3394 updated_at: 0,
3395 })
3396 .await
3397 .expect("seed RunRecord");
3398
3399 let mut input = launch_input(blueprint, json!({ "input": "hi" }));
3400 input.run_ctx = Some(RunContext::new(run_id, run_store));
3401
3402 let err = svc.launch(input).await.expect_err("expected FlowEval");
3403 match err {
3404 TaskLaunchError::FlowEval {
3405 message,
3406 failed_step,
3407 verdict_value,
3408 partial_ctx,
3409 } => {
3410 assert!(
3411 message.contains("blocked"),
3412 "expected message to mention blocked, got: {message}"
3413 );
3414 assert_eq!(
3415 failed_step,
3416 Some("gate".to_string()),
3417 "failed_step should be the Blueprint step ref, not the opaque StepId"
3418 );
3419 let vv = verdict_value.expect("verdict_value must be Some for Blocked");
3420 assert_eq!(vv["verdict"], "BLOCKED");
3421 assert_eq!(vv["reason"], "not applicable");
3422 assert!(
3426 partial_ctx.is_some(),
3427 "partial_ctx must be Some when a RunContext was supplied"
3428 );
3429 }
3430 other => panic!("expected FlowEval, got {other:?}"),
3431 }
3432 }
3433
3434 #[tokio::test]
3440 async fn task_launch_flow_eval_error_partial_ctx_reconstructs_from_run_store() {
3441 use crate::store::run::{InMemoryRunStore, RunContext, RunRecord, RunStatus, RunStore};
3442 use crate::types::{RunId, TaskId};
3443
3444 let factory = RustFnInProcessSpawnerFactory::new()
3445 .register_fn("upper", |inv| async move {
3446 Ok(WorkerResult {
3447 value: json!(inv.prompt.to_uppercase()),
3448 ok: true,
3449 stats: None,
3450 })
3451 })
3452 .register_fn("gate", |_inv| async move {
3453 Ok(WorkerResult {
3454 value: json!({ "verdict": "BLOCKED" }),
3455 ok: false,
3456 stats: None,
3457 })
3458 })
3459 .register_fn("never", |inv| async move {
3460 Ok(WorkerResult {
3461 value: json!(inv.prompt),
3462 ok: true,
3463 stats: None,
3464 })
3465 });
3466 let svc = build_service(factory);
3467 let flow = FlowNode::Seq {
3468 children: vec![
3469 step("upper", path("$.in"), path("$.s1")),
3470 step("gate", path("$.s1"), path("$.s2")),
3471 step("never", path("$.s2"), path("$.s3")),
3472 ],
3473 };
3474 let blueprint = bp(
3475 flow,
3476 vec![
3477 agent("upper", "upper"),
3478 agent("gate", "gate"),
3479 agent("never", "never"),
3480 ],
3481 );
3482
3483 let run_store: Arc<dyn RunStore> = Arc::new(InMemoryRunStore::new());
3484 let run_id = RunId::new();
3485 run_store
3486 .create(RunRecord {
3487 id: run_id.clone(),
3488 task_id: TaskId::new(),
3489 status: RunStatus::Running,
3490 step_entries: Vec::new(),
3491 degradations: Vec::new(),
3492 operator_sid: None,
3493 result_ref: None,
3494 input_json: Some("{}".to_string()),
3495 created_at: 0,
3496 updated_at: 0,
3497 })
3498 .await
3499 .expect("seed RunRecord");
3500
3501 let mut input = launch_input(blueprint, json!({ "in": "hi" }));
3502 input.run_ctx = Some(RunContext::new(run_id.clone(), run_store.clone()));
3503
3504 let err = svc.launch(input).await.expect_err("expected FlowEval");
3505 let partial_ctx = match err {
3506 TaskLaunchError::FlowEval { partial_ctx, .. } => {
3507 partial_ctx.expect("partial_ctx must be Some")
3508 }
3509 other => panic!("expected FlowEval, got {other:?}"),
3510 };
3511 let steps = partial_ctx
3512 .get("steps")
3513 .and_then(|v| v.as_object())
3514 .expect("partial_ctx.steps object");
3515 assert_eq!(
3519 steps.len(),
3520 2,
3521 "expected 2 step_entries (upper passed + gate blocked), got: {steps:?}"
3522 );
3523 let mut status_by_ref: HashMap<String, String> = HashMap::new();
3524 for (_step_id, entry) in steps {
3525 let step_ref = entry
3526 .get("step_ref")
3527 .and_then(|v| v.as_str())
3528 .expect("step_ref present")
3529 .to_string();
3530 let status = entry
3531 .get("status")
3532 .and_then(|v| v.as_str())
3533 .expect("status present")
3534 .to_string();
3535 status_by_ref.insert(step_ref, status);
3536 }
3537 assert_eq!(
3538 status_by_ref.get("upper").map(String::as_str),
3539 Some("passed")
3540 );
3541 assert_eq!(
3542 status_by_ref.get("gate").map(String::as_str),
3543 Some("blocked")
3544 );
3545 assert!(
3546 !status_by_ref.contains_key("never"),
3547 "step 'never' must not appear — flow-ir stops dispatching after Blocked abort"
3548 );
3549 }
3550
3551 #[tokio::test]
3557 async fn task_launch_flow_eval_error_without_run_ctx_has_none_fields() {
3558 let factory = RustFnInProcessSpawnerFactory::new().register_fn("gate", |_inv| async move {
3559 Ok(WorkerResult {
3560 value: json!({ "verdict": "BLOCKED" }),
3561 ok: false,
3562 stats: None,
3563 })
3564 });
3565 let svc = build_service(factory);
3566 let blueprint = bp(
3567 step("gate", path("$.input"), path("$.out")),
3568 vec![agent("gate", "gate")],
3569 );
3570 let err = svc
3571 .launch(launch_input(blueprint, json!({ "input": "hi" })))
3572 .await
3573 .expect_err("expected FlowEval");
3574 match err {
3575 TaskLaunchError::FlowEval {
3576 failed_step,
3577 verdict_value,
3578 partial_ctx,
3579 ..
3580 } => {
3581 assert_eq!(
3582 failed_step, None,
3583 "failed_step must be None without run_ctx"
3584 );
3585 assert_eq!(
3586 verdict_value, None,
3587 "verdict_value must be None without run_ctx"
3588 );
3589 assert_eq!(
3590 partial_ctx, None,
3591 "partial_ctx must be None without run_ctx"
3592 );
3593 }
3594 other => panic!("expected FlowEval, got {other:?}"),
3595 }
3596 }
3597}