Skip to main content

bamboo_engine/workflow_run/
executor.rs

1use std::collections::{BTreeMap, BTreeSet, HashMap};
2use std::future::Future;
3use std::pin::Pin;
4use std::sync::{Arc, Weak};
5use std::time::Duration;
6
7use async_trait::async_trait;
8use bamboo_agent_core::tools::{
9    FunctionCall, ToolCall, ToolExecutionContext, ToolExecutor, ToolOutcome, ToolResult,
10};
11use bamboo_domain::{
12    validate_schema, CompiledWorkflow, FailurePolicy, StartWorkflowRun, ValueRef,
13    WorkflowBudgetUsage, WorkflowBudgets, WorkflowCompileError, WorkflowDefinitionBundle,
14    WorkflowFailure, WorkflowFailureCode, WorkflowPlan, WorkflowProgress, WorkflowRunDefinition,
15    WorkflowRunEvent, WorkflowRunEventKind, WorkflowRunSnapshot, WorkflowRunStatus,
16    WorkflowStepDefinition, WorkflowStepKind, WorkflowStepSnapshot, WorkflowStepStatus,
17    WorkflowSuspensionContext,
18};
19use chrono::Utc;
20use dashmap::DashMap;
21use futures::{future::join_all, stream::FuturesUnordered, StreamExt};
22use serde_json::Value;
23use sha2::{Digest, Sha256};
24use thiserror::Error;
25use tokio::sync::{broadcast, Mutex, Semaphore};
26use tokio_util::sync::CancellationToken;
27use uuid::Uuid;
28
29use super::repository::WorkflowRunRepository;
30
31type SecretResolutionFuture<'a> =
32    Pin<Box<dyn Future<Output = Result<(Value, Vec<String>), WorkflowFailure>> + Send + 'a>>;
33
34#[derive(Debug, Clone, PartialEq, Eq)]
35pub struct NamedAgentSpec {
36    pub name: String,
37    pub allowed_capabilities: BTreeSet<String>,
38}
39
40#[derive(Debug, Clone)]
41pub struct AgentStepResult {
42    pub output: Value,
43    pub tokens: u64,
44    pub cost_micros: u64,
45}
46
47#[async_trait]
48pub trait AgentStepPort: Send + Sync {
49    /// #563 seam. Unknown names must return `Ok(None)` and fail preflight.
50    async fn resolve(&self, name: &str) -> Result<Option<NamedAgentSpec>, String>;
51    async fn execute(
52        &self,
53        spec: &NamedAgentSpec,
54        prompt: Value,
55        model: Option<&str>,
56        effort: Option<&str>,
57        capabilities: &BTreeSet<String>,
58        session_id: &str,
59    ) -> Result<AgentStepResult, String>;
60}
61
62#[async_trait]
63pub trait WorkflowDefinitionPort: Send + Sync {
64    /// Pin the root and every transitively referenced nested definition from one
65    /// immutable catalog publication. Implementations must never re-read a live
66    /// source while constructing the returned bundle.
67    async fn pin_bundle(
68        &self,
69        root: &WorkflowRunDefinition,
70    ) -> Result<WorkflowDefinitionBundle, String>;
71}
72
73#[derive(Debug, Clone, PartialEq, Eq)]
74pub enum PermissionDecision {
75    Allow,
76    Deny(String),
77}
78
79#[derive(Debug, Clone, PartialEq, Eq)]
80pub enum WorkflowPolicyTarget {
81    Tool(String),
82    Agent(String),
83    Workflow { id: String, revision: u64 },
84}
85
86#[async_trait]
87pub trait WorkflowPolicyPort: Send + Sync {
88    async fn authorize(
89        &self,
90        session_id: &str,
91        target: &WorkflowPolicyTarget,
92        requested: &BTreeSet<String>,
93        workspace_trusted: bool,
94    ) -> PermissionDecision;
95}
96
97/// Resolves a typed, persisted-safe capability handle to ephemeral secret
98/// material. Implementations own access control; raw values are never returned
99/// in snapshots/events/errors.
100pub struct WorkflowSecretMaterial(String);
101
102impl WorkflowSecretMaterial {
103    pub fn new(value: String) -> Self {
104        Self(value)
105    }
106
107    fn into_exposed(self) -> String {
108        self.0
109    }
110}
111
112#[async_trait]
113pub trait WorkflowSecretResolverPort: Send + Sync {
114    async fn resolve(
115        &self,
116        session_id: &str,
117        capability: &str,
118    ) -> Result<WorkflowSecretMaterial, String>;
119}
120
121#[derive(Debug, Error)]
122pub enum WorkflowRunError {
123    #[error(transparent)]
124    Compile(#[from] WorkflowCompileError),
125    #[error("invalid workflow input: {0}")]
126    InvalidInput(String),
127    #[error("workflow preflight failed: {0}")]
128    Preflight(String),
129    #[error("workflow storage failed: {0}")]
130    Storage(String),
131    #[error("workflow run not found")]
132    NotFound,
133    #[error("workflow run is already terminal")]
134    Terminal,
135}
136
137pub struct WorkflowRunEngine {
138    repository: Arc<dyn WorkflowRunRepository>,
139    tools: Arc<dyn ToolExecutor>,
140    agents: Arc<dyn AgentStepPort>,
141    definitions: Arc<dyn WorkflowDefinitionPort>,
142    policy: Arc<dyn WorkflowPolicyPort>,
143    secrets: Arc<dyn WorkflowSecretResolverPort>,
144    ceilings: WorkflowBudgets,
145    active: DashMap<String, Arc<ActiveRun>>,
146    events: DashMap<String, broadcast::Sender<WorkflowRunEvent>>,
147}
148
149struct ActiveRun {
150    cancellation: CancellationToken,
151    snapshot: Arc<Mutex<WorkflowRunSnapshot>>,
152}
153
154struct RuntimeRegistration {
155    engine: Weak<WorkflowRunEngine>,
156    run_id: String,
157}
158
159impl Drop for RuntimeRegistration {
160    fn drop(&mut self) {
161        if let Some(engine) = self.engine.upgrade() {
162            engine.active.remove(&self.run_id);
163            engine.events.remove(&self.run_id);
164        }
165    }
166}
167
168struct RunContext {
169    engine: Arc<WorkflowRunEngine>,
170    compiled: Arc<CompiledWorkflow>,
171    bundle: Arc<WorkflowDefinitionBundle>,
172    pinned_agents: Arc<HashMap<String, NamedAgentSpec>>,
173    snapshot: Arc<Mutex<WorkflowRunSnapshot>>,
174    cancellation: CancellationToken,
175    branch_cancellation: CancellationToken,
176    allowed_capabilities: BTreeSet<String>,
177    workspace_trusted: bool,
178    semaphore: Arc<Semaphore>,
179    items: HashMap<String, Value>,
180    scope: String,
181    depth: u32,
182    ledger: Arc<Mutex<WorkflowBudgetUsage>>,
183    root_limits: WorkflowBudgets,
184}
185
186impl Clone for RunContext {
187    fn clone(&self) -> Self {
188        Self {
189            engine: self.engine.clone(),
190            compiled: self.compiled.clone(),
191            bundle: self.bundle.clone(),
192            pinned_agents: self.pinned_agents.clone(),
193            snapshot: self.snapshot.clone(),
194            cancellation: self.cancellation.clone(),
195            branch_cancellation: self.branch_cancellation.clone(),
196            allowed_capabilities: self.allowed_capabilities.clone(),
197            workspace_trusted: self.workspace_trusted,
198            semaphore: self.semaphore.clone(),
199            items: self.items.clone(),
200            scope: self.scope.clone(),
201            depth: self.depth,
202            ledger: self.ledger.clone(),
203            root_limits: self.root_limits.clone(),
204        }
205    }
206}
207
208type NodeFuture<'a> = Pin<Box<dyn Future<Output = Result<Value, WorkflowFailure>> + Send + 'a>>;
209type StartSignal =
210    Arc<Mutex<Option<tokio::sync::oneshot::Sender<Result<WorkflowRunSnapshot, WorkflowRunError>>>>>;
211
212impl WorkflowRunEngine {
213    pub fn new(
214        repository: Arc<dyn WorkflowRunRepository>,
215        tools: Arc<dyn ToolExecutor>,
216        agents: Arc<dyn AgentStepPort>,
217        definitions: Arc<dyn WorkflowDefinitionPort>,
218        policy: Arc<dyn WorkflowPolicyPort>,
219        secrets: Arc<dyn WorkflowSecretResolverPort>,
220        ceilings: WorkflowBudgets,
221    ) -> Arc<Self> {
222        Arc::new(Self {
223            repository,
224            tools,
225            agents,
226            definitions,
227            policy,
228            secrets,
229            ceilings,
230            active: DashMap::new(),
231            events: DashMap::new(),
232        })
233    }
234
235    pub async fn run(
236        self: &Arc<Self>,
237        request: StartWorkflowRun,
238    ) -> Result<WorkflowRunSnapshot, WorkflowRunError> {
239        let bundle = self.pin_and_validate_bundle(&request.definition).await?;
240        self.run_pinned(request, bundle).await
241    }
242
243    pub async fn run_pinned(
244        self: &Arc<Self>,
245        request: StartWorkflowRun,
246        bundle: WorkflowDefinitionBundle,
247    ) -> Result<WorkflowRunSnapshot, WorkflowRunError> {
248        self.validate_bundle(&request.definition, &bundle)?;
249        let cancellation = CancellationToken::new();
250        let ledger = Arc::new(Mutex::new(WorkflowBudgetUsage::default()));
251        let limits = effective_limits(&request.definition.budgets, &self.ceilings);
252        let semaphore = Arc::new(Semaphore::new(limits.max_concurrency));
253        let pinned_agents = Arc::new(
254            self.preflight_bundle(
255                &bundle,
256                &request.session_id,
257                &request.allowed_capabilities.iter().cloned().collect(),
258                request.workspace_trusted,
259                &limits,
260            )
261            .await?,
262        );
263        self.run_internal(
264            request,
265            Arc::new(bundle),
266            pinned_agents,
267            None,
268            None,
269            0,
270            cancellation,
271            ledger,
272            limits,
273            semaphore,
274            None,
275        )
276        .await
277    }
278
279    /// Start in the background and return only after the running snapshot is
280    /// durable. HTTP and tool adapters use this non-blocking entrypoint.
281    pub async fn start(
282        self: &Arc<Self>,
283        request: StartWorkflowRun,
284    ) -> Result<WorkflowRunSnapshot, WorkflowRunError> {
285        let bundle = self.pin_and_validate_bundle(&request.definition).await?;
286        self.start_pinned(request, bundle).await
287    }
288
289    pub async fn start_pinned(
290        self: &Arc<Self>,
291        request: StartWorkflowRun,
292        bundle: WorkflowDefinitionBundle,
293    ) -> Result<WorkflowRunSnapshot, WorkflowRunError> {
294        self.validate_bundle(&request.definition, &bundle)?;
295        let cancellation = CancellationToken::new();
296        let ledger = Arc::new(Mutex::new(WorkflowBudgetUsage::default()));
297        let limits = effective_limits(&request.definition.budgets, &self.ceilings);
298        let semaphore = Arc::new(Semaphore::new(limits.max_concurrency));
299        let pinned_agents = Arc::new(
300            self.preflight_bundle(
301                &bundle,
302                &request.session_id,
303                &request.allowed_capabilities.iter().cloned().collect(),
304                request.workspace_trusted,
305                &limits,
306            )
307            .await?,
308        );
309        let (tx, rx) = tokio::sync::oneshot::channel();
310        let signal = Arc::new(Mutex::new(Some(tx)));
311        let engine = self.clone();
312        tokio::spawn(async move {
313            let result = engine
314                .run_internal(
315                    request,
316                    Arc::new(bundle),
317                    pinned_agents,
318                    None,
319                    None,
320                    0,
321                    cancellation,
322                    ledger,
323                    limits,
324                    semaphore,
325                    Some(signal.clone()),
326                )
327                .await;
328            if let Err(error) = result {
329                if let Some(sender) = signal.lock().await.take() {
330                    let _ = sender.send(Err(error));
331                } else {
332                    tracing::error!("background workflow run failed after start");
333                }
334            } else if signal.lock().await.is_some() {
335                tracing::error!("workflow task completed without publishing a start snapshot");
336            }
337        });
338        rx.await.map_err(|_| {
339            WorkflowRunError::Storage("workflow task exited before durable start".to_string())
340        })?
341    }
342
343    /// Phase-1 safe restart starts a fresh run from the suspended run's pinned
344    /// definition snapshot. Prefix/script resume remains explicitly out of scope (#581).
345    pub async fn restart(
346        self: &Arc<Self>,
347        run_id: &str,
348        workspace_trusted: bool,
349        allowed_capabilities: Vec<String>,
350    ) -> Result<WorkflowRunSnapshot, WorkflowRunError> {
351        let previous = self
352            .repository
353            .load(run_id)
354            .await
355            .map_err(storage)?
356            .ok_or(WorkflowRunError::NotFound)?;
357        if previous.status != WorkflowRunStatus::Suspended {
358            return Err(if previous.status.is_terminal() {
359                WorkflowRunError::Terminal
360            } else {
361                WorkflowRunError::Preflight("only suspended workflows can restart".to_string())
362            });
363        }
364        if matches!(
365            previous.suspension,
366            Some(
367                WorkflowSuspensionContext::ToolApproval { .. }
368                    | WorkflowSuspensionContext::ToolRunning { .. }
369            )
370        ) {
371            return Err(WorkflowRunError::Preflight(
372                "workflow has durable suspension context and requires explicit resume handling"
373                    .to_string(),
374            ));
375        }
376        let bundle = previous.definition_bundle;
377        self.start_pinned(
378            StartWorkflowRun {
379                definition: previous.definition,
380                args: previous.validated_args,
381                session_id: previous.session_id,
382                workspace_trusted,
383                allowed_capabilities,
384            },
385            bundle,
386        )
387        .await
388    }
389
390    #[allow(clippy::too_many_arguments)]
391    async fn run_internal(
392        self: &Arc<Self>,
393        request: StartWorkflowRun,
394        bundle: Arc<WorkflowDefinitionBundle>,
395        pinned_agents: Arc<HashMap<String, NamedAgentSpec>>,
396        parent_run_id: Option<String>,
397        parent_step_id: Option<String>,
398        depth: u32,
399        cancellation: CancellationToken,
400        ledger: Arc<Mutex<WorkflowBudgetUsage>>,
401        root_limits: WorkflowBudgets,
402        semaphore: Arc<Semaphore>,
403        started: Option<StartSignal>,
404    ) -> Result<WorkflowRunSnapshot, WorkflowRunError> {
405        if request.definition.steps.len() > self.ceilings.max_steps as usize {
406            return Err(WorkflowRunError::Preflight(
407                "workflow definition exceeds server step-count ceiling".to_string(),
408            ));
409        }
410        let compiled = Arc::new(CompiledWorkflow::compile(request.definition)?);
411        self.enforce_ceilings(&compiled.definition.budgets)?;
412        let definition_value = serde_json::to_value(&compiled.definition)
413            .map_err(|error| WorkflowRunError::Preflight(error.to_string()))?;
414        reject_secret_material_in_definition(&definition_value)
415            .map_err(WorkflowRunError::Preflight)?;
416        compiled
417            .validate_input(&request.args)
418            .map_err(WorkflowRunError::InvalidInput)?;
419        reject_secret_material(&request.args).map_err(WorkflowRunError::InvalidInput)?;
420        let allowed_capabilities = request
421            .allowed_capabilities
422            .into_iter()
423            .collect::<BTreeSet<_>>();
424        enforce_budget_within(&compiled.definition.budgets, &root_limits).map_err(|message| {
425            WorkflowRunError::Preflight(format!("nested workflow budget expands root: {message}"))
426        })?;
427
428        let run_id = Uuid::new_v4().to_string();
429        let now = Utc::now();
430        let snapshot = WorkflowRunSnapshot {
431            run_id: run_id.clone(),
432            parent_run_id,
433            parent_step_id,
434            session_id: request.session_id,
435            definition: compiled.definition.clone(),
436            definition_bundle: bundle.as_ref().clone(),
437            definition_bundle_hash: definition_bundle_hash(&bundle)?,
438            validated_args: request.args,
439            status: WorkflowRunStatus::Queued,
440            steps: BTreeMap::new(),
441            usage: WorkflowBudgetUsage::default(),
442            last_sequence: 1,
443            output: None,
444            failure: None,
445            suspension: None,
446            created_at: now,
447            updated_at: now,
448        };
449        let queued = event(&snapshot, None, WorkflowRunEventKind::RunQueued);
450        self.repository
451            .create(&snapshot, &queued)
452            .await
453            .map_err(storage)?;
454        let (sender, _) = broadcast::channel(256);
455        self.events.insert(run_id.clone(), sender);
456        self.publish(&queued);
457        let snapshot = Arc::new(Mutex::new(snapshot));
458        self.active.insert(
459            run_id.clone(),
460            Arc::new(ActiveRun {
461                cancellation: cancellation.clone(),
462                snapshot: snapshot.clone(),
463            }),
464        );
465        let _registration = RuntimeRegistration {
466            engine: Arc::downgrade(self),
467            run_id: run_id.clone(),
468        };
469        {
470            let mut shared = snapshot.lock().await;
471            let start_result = if cancellation.is_cancelled() {
472                self.finish_cancelled(&mut shared).await
473            } else {
474                self.transition(
475                    &mut shared,
476                    None,
477                    WorkflowRunEventKind::RunStarted,
478                    |snapshot| {
479                        snapshot.status = WorkflowRunStatus::Running;
480                    },
481                )
482                .await
483            };
484            start_result?;
485            if let Some(started) = started {
486                if let Some(sender) = started.lock().await.take() {
487                    let _ = sender.send(Ok(shared.clone()));
488                }
489            }
490        }
491        let context = RunContext {
492            engine: self.clone(),
493            compiled: compiled.clone(),
494            bundle,
495            pinned_agents,
496            snapshot: snapshot.clone(),
497            cancellation: cancellation.clone(),
498            branch_cancellation: cancellation.child_token(),
499            allowed_capabilities,
500            workspace_trusted: request.workspace_trusted,
501            semaphore,
502            items: HashMap::new(),
503            scope: "root".to_string(),
504            depth,
505            ledger,
506            root_limits,
507        };
508        let result = tokio::time::timeout(
509            Duration::from_millis(compiled.definition.budgets.wall_time_ms),
510            context.execute_node(&compiled.definition.plan, "root"),
511        )
512        .await;
513        let mut final_snapshot = snapshot.lock().await;
514        if final_snapshot.status.is_terminal() {
515            return Ok(final_snapshot.clone());
516        }
517        match result {
518            Ok(Ok(output)) if cancellation.is_cancelled() => {
519                let _ = output;
520                self.finish_cancelled(&mut final_snapshot).await?;
521            }
522            Ok(Ok(output)) => {
523                if let Some(schema) = &compiled.definition.output_schema {
524                    if let Err(message) = validate_schema(schema, &output) {
525                        let failure = failure(WorkflowFailureCode::InvalidOutput, message, false);
526                        self.finish_failed(&mut final_snapshot, failure).await?;
527                    } else {
528                        self.finish_succeeded(&mut final_snapshot, output).await?;
529                    }
530                } else {
531                    self.finish_succeeded(&mut final_snapshot, output).await?;
532                }
533            }
534            Ok(Err(error)) if error.code == WorkflowFailureCode::Cancelled => {
535                self.finish_cancelled(&mut final_snapshot).await?;
536            }
537            Ok(Err(error)) if error.code == WorkflowFailureCode::Suspended => {
538                self.finish_suspended(&mut final_snapshot, error.message)
539                    .await?;
540            }
541            Ok(Err(error)) => self.finish_failed(&mut final_snapshot, error).await?,
542            Err(_) => {
543                cancellation.cancel();
544                // `timeout` cancels the node future at an arbitrary await,
545                // including a repository commit. Reconcile the in-memory copy
546                // from the journal/temp recovery protocol before allocating the
547                // next sequence, so a partially committed StepStarted cannot
548                // make the timeout terminal transition skip a sequence.
549                if let Some(durable) = self.repository.load(&run_id).await.map_err(storage)? {
550                    *final_snapshot = durable;
551                }
552                let step_failure = failure(
553                    WorkflowFailureCode::BudgetExceeded,
554                    "workflow wall-time budget exceeded",
555                    false,
556                );
557                self.fail_timeout_frontier(
558                    &mut final_snapshot,
559                    &compiled.definition.plan,
560                    step_failure.clone(),
561                )
562                .await?;
563                self.fail_active_steps(&mut final_snapshot, step_failure.clone())
564                    .await?;
565                self.finish_failed(&mut final_snapshot, step_failure)
566                    .await?;
567            }
568        }
569        Ok(final_snapshot.clone())
570    }
571
572    pub async fn progress(
573        &self,
574        run_id: &str,
575        since: u64,
576    ) -> Result<WorkflowProgress, WorkflowRunError> {
577        let snapshot = self
578            .repository
579            .load(run_id)
580            .await
581            .map_err(storage)?
582            .ok_or(WorkflowRunError::NotFound)?;
583        let events = self
584            .repository
585            .events_since(run_id, since)
586            .await
587            .map_err(storage)?;
588        Ok(WorkflowProgress { snapshot, events })
589    }
590
591    pub async fn list_run_ids(&self) -> Result<Vec<String>, WorkflowRunError> {
592        self.repository.list_run_ids().await.map_err(storage)
593    }
594
595    pub async fn cancel(&self, run_id: &str) -> Result<WorkflowRunSnapshot, WorkflowRunError> {
596        let mut snapshot = self
597            .repository
598            .load(run_id)
599            .await
600            .map_err(storage)?
601            .ok_or(WorkflowRunError::NotFound)?;
602        if snapshot.status == WorkflowRunStatus::Cancelled {
603            return Ok(snapshot);
604        }
605        if snapshot.status.is_terminal() {
606            return Err(WorkflowRunError::Terminal);
607        }
608        if let Some(active) = self.active.get(run_id).map(|active| active.clone()) {
609            active.cancellation.cancel();
610            let mut shared = active.snapshot.lock().await;
611            if !shared.status.is_terminal() {
612                self.finish_cancelled(&mut shared).await?;
613            }
614            return Ok(shared.clone());
615        }
616        self.finish_cancelled(&mut snapshot).await?;
617        Ok(snapshot)
618    }
619
620    pub async fn recover(&self) -> Result<Vec<WorkflowRunSnapshot>, WorkflowRunError> {
621        let mut recovered = Vec::new();
622        for run_id in self.repository.list_run_ids().await.map_err(storage)? {
623            let Some(mut snapshot) = self.repository.load(&run_id).await.map_err(storage)? else {
624                continue;
625            };
626            if matches!(
627                snapshot.status,
628                WorkflowRunStatus::Queued | WorkflowRunStatus::Running
629            ) {
630                let reason = "process restarted; explicit safe restart is required".to_string();
631                let active_steps = snapshot
632                    .steps
633                    .iter()
634                    .filter(|(_, step)| {
635                        matches!(
636                            step.status,
637                            WorkflowStepStatus::Queued | WorkflowStepStatus::Running
638                        )
639                    })
640                    .map(|(id, _)| id.clone())
641                    .collect::<Vec<_>>();
642                for step_id in active_steps {
643                    let state_id = step_id.clone();
644                    let step_reason = reason.clone();
645                    self.transition(
646                        &mut snapshot,
647                        Some(step_id),
648                        WorkflowRunEventKind::StepSuspended {
649                            reason: reason.clone(),
650                        },
651                        move |snapshot| {
652                            if let Some(step) = snapshot.steps.get_mut(&state_id) {
653                                step.status = WorkflowStepStatus::Suspended;
654                                step.failure = Some(failure(
655                                    WorkflowFailureCode::RecoverySuspended,
656                                    step_reason,
657                                    true,
658                                ));
659                            }
660                        },
661                    )
662                    .await?;
663                }
664                self.transition(
665                    &mut snapshot,
666                    None,
667                    WorkflowRunEventKind::RunSuspended {
668                        reason: reason.clone(),
669                    },
670                    move |snapshot| {
671                        snapshot.status = WorkflowRunStatus::Suspended;
672                        snapshot.suspension = Some(WorkflowSuspensionContext::Recovery {
673                            reason: reason.clone(),
674                        });
675                    },
676                )
677                .await?;
678                recovered.push(snapshot);
679            }
680        }
681        Ok(recovered)
682    }
683
684    pub fn subscribe(&self, run_id: &str) -> Option<broadcast::Receiver<WorkflowRunEvent>> {
685        self.events.get(run_id).map(|sender| sender.subscribe())
686    }
687
688    #[cfg(test)]
689    pub(crate) fn runtime_resource_counts(&self) -> (usize, usize) {
690        (self.active.len(), self.events.len())
691    }
692
693    async fn pin_and_validate_bundle(
694        &self,
695        root: &WorkflowRunDefinition,
696    ) -> Result<WorkflowDefinitionBundle, WorkflowRunError> {
697        let bundle =
698            self.definitions.pin_bundle(root).await.map_err(|_| {
699                WorkflowRunError::Preflight("workflow bundle pin failed".to_string())
700            })?;
701        self.validate_bundle(root, &bundle)?;
702        Ok(bundle)
703    }
704
705    fn validate_bundle(
706        &self,
707        root: &WorkflowRunDefinition,
708        bundle: &WorkflowDefinitionBundle,
709    ) -> Result<(), WorkflowRunError> {
710        if bundle.root_id != root.id
711            || bundle.root_revision != root.revision
712            || bundle.root() != Some(root)
713        {
714            return Err(WorkflowRunError::Preflight(
715                "pinned bundle root identity/content mismatch".to_string(),
716            ));
717        }
718        let serialized = serde_json::to_value(bundle).map_err(|_| {
719            WorkflowRunError::Preflight("workflow bundle is not serializable".into())
720        })?;
721        reject_secret_material_in_definition(&serialized).map_err(WorkflowRunError::Preflight)?;
722        let mut stack = vec![(root.id.clone(), root.revision, Vec::<String>::new())];
723        let mut visited = BTreeSet::new();
724        while let Some((id, revision, path)) = stack.pop() {
725            let key = WorkflowDefinitionBundle::key(&id, revision);
726            if path.contains(&key) {
727                return Err(WorkflowRunError::Preflight(format!(
728                    "nested workflow cycle includes {key}"
729                )));
730            }
731            if !visited.insert(key.clone()) {
732                continue;
733            }
734            let definition = bundle.get(&id, revision).ok_or_else(|| {
735                WorkflowRunError::Preflight(format!("pinned bundle is missing {key}"))
736            })?;
737            if definition.id != id || definition.revision != revision {
738                return Err(WorkflowRunError::Preflight(
739                    "pinned bundle definition identity mismatch".to_string(),
740                ));
741            }
742            let compiled = CompiledWorkflow::compile(definition.clone())?;
743            let mut nested_path = path;
744            nested_path.push(key);
745            for step in compiled.steps.values() {
746                if let WorkflowStepKind::Workflow {
747                    workflow_id,
748                    revision,
749                    args,
750                } = &step.kind
751                {
752                    let nested = bundle.get(workflow_id, *revision).ok_or_else(|| {
753                        WorkflowRunError::Preflight(format!(
754                            "pinned bundle is missing {workflow_id}@{revision}"
755                        ))
756                    })?;
757                    validate_nested_input_contract(args, &nested.input_schema, &compiled)
758                        .map_err(WorkflowRunError::Preflight)?;
759                    stack.push((workflow_id.clone(), *revision, nested_path.clone()));
760                }
761            }
762        }
763        Ok(())
764    }
765
766    async fn preflight_bundle(
767        &self,
768        bundle: &WorkflowDefinitionBundle,
769        session_id: &str,
770        allowed: &BTreeSet<String>,
771        trusted: bool,
772        root_limits: &WorkflowBudgets,
773    ) -> Result<HashMap<String, NamedAgentSpec>, WorkflowRunError> {
774        let mut pinned_agents = HashMap::<String, NamedAgentSpec>::new();
775        let mut stack = vec![(bundle.root_id.clone(), bundle.root_revision, 0_u32)];
776        let mut visited = BTreeSet::new();
777        while let Some((id, revision, depth)) = stack.pop() {
778            if depth >= root_limits.max_nesting_depth {
779                return Err(WorkflowRunError::Preflight(
780                    "nested workflow depth exceeded shared root limit".to_string(),
781                ));
782            }
783            if !visited.insert(WorkflowDefinitionBundle::key(&id, revision)) {
784                continue;
785            }
786            let definition = bundle.get(&id, revision).ok_or_else(|| {
787                WorkflowRunError::Preflight("pinned workflow definition missing".to_string())
788            })?;
789            enforce_budget_within(&definition.budgets, root_limits).map_err(|message| {
790                WorkflowRunError::Preflight(format!(
791                    "nested workflow budget expands root: {message}"
792                ))
793            })?;
794            let compiled = CompiledWorkflow::compile(definition.clone())?;
795            for step in compiled.steps.values() {
796                let (target, capabilities) = match &step.kind {
797                    WorkflowStepKind::Tool {
798                        tool, capabilities, ..
799                    } => (WorkflowPolicyTarget::Tool(tool.clone()), capabilities),
800                    WorkflowStepKind::Agent {
801                        agent,
802                        capabilities,
803                        ..
804                    } => {
805                        let spec = if let Some(spec) = pinned_agents.get(agent) {
806                            spec.clone()
807                        } else {
808                            let spec = self
809                                .agents
810                                .resolve(agent)
811                                .await
812                                .map_err(|_| {
813                                    WorkflowRunError::Preflight(
814                                        "named agent resolution failed".to_string(),
815                                    )
816                                })?
817                                .ok_or_else(|| {
818                                    WorkflowRunError::Preflight(format!(
819                                        "unknown named agent '{agent}'"
820                                    ))
821                                })?;
822                            if spec.name != *agent {
823                                return Err(WorkflowRunError::Preflight(
824                                    "named agent resolver returned mismatched identity".to_string(),
825                                ));
826                            }
827                            pinned_agents.insert(agent.clone(), spec.clone());
828                            spec
829                        };
830                        if !capabilities
831                            .iter()
832                            .all(|capability| spec.allowed_capabilities.contains(capability))
833                        {
834                            return Err(WorkflowRunError::Preflight(format!(
835                                "agent '{agent}' capability expansion denied"
836                            )));
837                        }
838                        (WorkflowPolicyTarget::Agent(agent.clone()), capabilities)
839                    }
840                    WorkflowStepKind::Workflow {
841                        workflow_id,
842                        revision,
843                        args,
844                    } => {
845                        let nested = bundle.get(workflow_id, *revision).ok_or_else(|| {
846                            WorkflowRunError::Preflight(format!(
847                                "missing pinned workflow {workflow_id}@{revision}"
848                            ))
849                        })?;
850                        validate_nested_input_contract(args, &nested.input_schema, &compiled)
851                            .map_err(WorkflowRunError::Preflight)?;
852                        stack.push((workflow_id.clone(), *revision, depth + 1));
853                        (
854                            WorkflowPolicyTarget::Workflow {
855                                id: workflow_id.clone(),
856                                revision: *revision,
857                            },
858                            &Vec::new(),
859                        )
860                    }
861                };
862                let requested = capabilities.iter().cloned().collect::<BTreeSet<_>>();
863                if !requested.is_subset(allowed) {
864                    return Err(WorkflowRunError::Preflight(format!(
865                        "step '{}' exceeds root capabilities",
866                        step.id
867                    )));
868                }
869                if let PermissionDecision::Deny(_reason) = self
870                    .policy
871                    .authorize(session_id, &target, &requested, trusted)
872                    .await
873                {
874                    return Err(WorkflowRunError::Preflight(
875                        "workflow policy denied this step".to_string(),
876                    ));
877                }
878            }
879        }
880        Ok(pinned_agents)
881    }
882
883    fn enforce_ceilings(&self, budget: &WorkflowBudgets) -> Result<(), WorkflowRunError> {
884        if budget.max_concurrency > self.ceilings.max_concurrency
885            || budget.max_agents > self.ceilings.max_agents
886            || budget.max_steps > self.ceilings.max_steps
887            || budget.max_retries > self.ceilings.max_retries
888            || budget.max_nesting_depth > self.ceilings.max_nesting_depth
889            || budget.wall_time_ms > self.ceilings.wall_time_ms
890            || exceeds_optional(budget.max_tokens, self.ceilings.max_tokens)
891            || exceeds_optional(budget.max_cost_micros, self.ceilings.max_cost_micros)
892        {
893            return Err(WorkflowRunError::Preflight(
894                "definition exceeds server workflow budget ceilings".to_string(),
895            ));
896        }
897        Ok(())
898    }
899
900    async fn transition(
901        &self,
902        snapshot: &mut WorkflowRunSnapshot,
903        step_id: Option<String>,
904        kind: WorkflowRunEventKind,
905        mutate: impl FnOnce(&mut WorkflowRunSnapshot),
906    ) -> Result<(), WorkflowRunError> {
907        // Always derive a candidate from durable state. The commit itself runs
908        // in an owned task, so dropping this caller at a timeout/cancel boundary
909        // cannot interrupt rename/fsync halfway through. In-memory state advances
910        // only after that task confirms the durable commit.
911        let mut candidate = self
912            .repository
913            .load(&snapshot.run_id)
914            .await
915            .map_err(storage)?
916            .ok_or_else(|| WorkflowRunError::Storage("workflow snapshot missing".to_string()))?;
917        mutate(&mut candidate);
918        candidate.last_sequence += 1;
919        candidate.updated_at = Utc::now();
920        let event = event(&candidate, step_id, kind);
921        let repository = self.repository.clone();
922        let durable_candidate = candidate.clone();
923        let durable_event = event.clone();
924        let commit =
925            tokio::spawn(
926                async move { repository.commit(&durable_candidate, &durable_event).await },
927            );
928        commit
929            .await
930            .map_err(|error| {
931                WorkflowRunError::Storage(format!("workflow commit task failed: {error}"))
932            })?
933            .map_err(storage)?;
934        *snapshot = candidate;
935        self.publish(&event);
936        Ok(())
937    }
938
939    fn publish(&self, event: &WorkflowRunEvent) {
940        if let Some(sender) = self.events.get(&event.run_id) {
941            let _ = sender.send(event.clone());
942        }
943    }
944
945    async fn finish_succeeded(
946        &self,
947        snapshot: &mut WorkflowRunSnapshot,
948        output: Value,
949    ) -> Result<(), WorkflowRunError> {
950        if snapshot.status.is_terminal() {
951            return Ok(());
952        }
953        let copy = output.clone();
954        self.transition(
955            snapshot,
956            None,
957            WorkflowRunEventKind::RunSucceeded { output },
958            move |snapshot| {
959                snapshot.status = WorkflowRunStatus::Succeeded;
960                snapshot.output = Some(copy);
961            },
962        )
963        .await
964    }
965    async fn finish_failed(
966        &self,
967        snapshot: &mut WorkflowRunSnapshot,
968        error: WorkflowFailure,
969    ) -> Result<(), WorkflowRunError> {
970        if snapshot.status.is_terminal() {
971            return Ok(());
972        }
973        let copy = error.clone();
974        self.transition(
975            snapshot,
976            None,
977            WorkflowRunEventKind::RunFailed { failure: error },
978            move |snapshot| {
979                snapshot.status = WorkflowRunStatus::Failed;
980                snapshot.failure = Some(copy);
981            },
982        )
983        .await
984    }
985    async fn finish_cancelled(
986        &self,
987        snapshot: &mut WorkflowRunSnapshot,
988    ) -> Result<(), WorkflowRunError> {
989        if snapshot.status == WorkflowRunStatus::Cancelled {
990            return Ok(());
991        }
992        let active_steps = snapshot
993            .steps
994            .iter()
995            .filter(|(_, step)| {
996                matches!(
997                    step.status,
998                    WorkflowStepStatus::Queued | WorkflowStepStatus::Running
999                )
1000            })
1001            .map(|(id, _)| id.clone())
1002            .collect::<Vec<_>>();
1003        for step_id in active_steps {
1004            let state_id = step_id.clone();
1005            self.transition(
1006                snapshot,
1007                Some(step_id),
1008                WorkflowRunEventKind::StepCancelled,
1009                move |snapshot| {
1010                    if let Some(step) = snapshot.steps.get_mut(&state_id) {
1011                        step.status = WorkflowStepStatus::Cancelled;
1012                        step.failure = Some(failure(
1013                            WorkflowFailureCode::Cancelled,
1014                            "workflow cancelled",
1015                            false,
1016                        ));
1017                    }
1018                },
1019            )
1020            .await?;
1021        }
1022        self.transition(
1023            snapshot,
1024            None,
1025            WorkflowRunEventKind::RunCancelled,
1026            |snapshot| {
1027                snapshot.status = WorkflowRunStatus::Cancelled;
1028                snapshot.failure = Some(failure(
1029                    WorkflowFailureCode::Cancelled,
1030                    "workflow cancelled",
1031                    false,
1032                ));
1033            },
1034        )
1035        .await
1036    }
1037
1038    async fn finish_suspended(
1039        &self,
1040        snapshot: &mut WorkflowRunSnapshot,
1041        reason: String,
1042    ) -> Result<(), WorkflowRunError> {
1043        if snapshot.status.is_terminal() {
1044            return Ok(());
1045        }
1046        self.transition(
1047            snapshot,
1048            None,
1049            WorkflowRunEventKind::RunSuspended { reason },
1050            |snapshot| {
1051                snapshot.status = WorkflowRunStatus::Suspended;
1052            },
1053        )
1054        .await
1055    }
1056
1057    async fn fail_active_steps(
1058        &self,
1059        snapshot: &mut WorkflowRunSnapshot,
1060        error: WorkflowFailure,
1061    ) -> Result<(), WorkflowRunError> {
1062        let active_steps = snapshot
1063            .steps
1064            .iter()
1065            .filter(|(_, step)| {
1066                matches!(
1067                    step.status,
1068                    WorkflowStepStatus::Queued | WorkflowStepStatus::Running
1069                )
1070            })
1071            .map(|(id, _)| id.clone())
1072            .collect::<Vec<_>>();
1073        for step_id in active_steps {
1074            let state_id = step_id.clone();
1075            let copy = error.clone();
1076            self.transition(
1077                snapshot,
1078                Some(step_id),
1079                WorkflowRunEventKind::StepFailed {
1080                    failure: error.clone(),
1081                },
1082                move |snapshot| {
1083                    if let Some(step) = snapshot.steps.get_mut(&state_id) {
1084                        step.status = WorkflowStepStatus::Failed;
1085                        step.failure = Some(copy);
1086                    }
1087                },
1088            )
1089            .await?;
1090        }
1091        Ok(())
1092    }
1093
1094    async fn fail_timeout_frontier(
1095        &self,
1096        snapshot: &mut WorkflowRunSnapshot,
1097        plan: &WorkflowPlan,
1098        error: WorkflowFailure,
1099    ) -> Result<(), WorkflowRunError> {
1100        if snapshot.steps.values().any(|step| {
1101            matches!(
1102                step.status,
1103                WorkflowStepStatus::Queued | WorkflowStepStatus::Running
1104            )
1105        }) {
1106            return Ok(());
1107        }
1108        for step_id in plan_frontier(plan) {
1109            if snapshot.steps.contains_key(&step_id) {
1110                continue;
1111            }
1112            let state_id = step_id.clone();
1113            let state_error = error.clone();
1114            self.transition(
1115                snapshot,
1116                Some(step_id),
1117                WorkflowRunEventKind::StepFailed {
1118                    failure: error.clone(),
1119                },
1120                move |snapshot| {
1121                    snapshot.steps.insert(
1122                        state_id.clone(),
1123                        WorkflowStepSnapshot {
1124                            id: state_id,
1125                            status: WorkflowStepStatus::Failed,
1126                            input_hash: String::new(),
1127                            output: None,
1128                            failure: Some(state_error),
1129                            attempts: 0,
1130                        },
1131                    );
1132                },
1133            )
1134            .await?;
1135        }
1136        Ok(())
1137    }
1138}
1139
1140impl RunContext {
1141    fn execute_node<'a>(&'a self, plan: &'a WorkflowPlan, path: &'a str) -> NodeFuture<'a> {
1142        Box::pin(async move {
1143            self.check_cancelled()?;
1144            match plan {
1145                WorkflowPlan::Step { step } => self.execute_step(step, path).await,
1146                WorkflowPlan::Sequence { nodes } => {
1147                    let mut result = Value::Null;
1148                    for (index, node) in nodes.iter().enumerate() {
1149                        match self.execute_node(node, &format!("{path}.{index}")).await {
1150                            Ok(value) => result = value,
1151                            Err(error) => {
1152                                if error.code == WorkflowFailureCode::DependencySkipped {
1153                                    for remaining in &nodes[index + 1..] {
1154                                        self.skip_plan(
1155                                            remaining,
1156                                            "dependency requested skip_dependents",
1157                                        )
1158                                        .await?;
1159                                    }
1160                                }
1161                                return Err(error);
1162                            }
1163                        }
1164                    }
1165                    Ok(result)
1166                }
1167                WorkflowPlan::Parallel { nodes } => {
1168                    let parallel_cancellation = self.branch_cancellation.child_token();
1169                    let mut futures = FuturesUnordered::new();
1170                    for (index, node) in nodes.iter().enumerate() {
1171                        let mut child = self.clone();
1172                        child.branch_cancellation = parallel_cancellation.clone();
1173                        futures.push(async move {
1174                            (
1175                                index,
1176                                child.execute_node(node, &format!("{path}.{index}")).await,
1177                            )
1178                        });
1179                    }
1180                    let mut output = vec![Value::Null; nodes.len()];
1181                    while let Some((index, result)) = futures.next().await {
1182                        match result {
1183                            Ok(value) => output[index] = value,
1184                            Err(mut error) => {
1185                                parallel_cancellation.cancel();
1186                                // Drop sibling futures before awaiting durable
1187                                // cancellation transitions. A sibling may hold
1188                                // the snapshot mutex across its shielded commit;
1189                                // leaving it parked inside FuturesUnordered would
1190                                // deadlock this reconciliation.
1191                                drop(futures);
1192                                self.cancel_active_parallel_steps(nodes).await?;
1193                                error.message =
1194                                    format!("parallel branch[{index}] failed: {}", error.message);
1195                                return Err(error);
1196                            }
1197                        }
1198                    }
1199                    Ok(Value::Array(output))
1200                }
1201                WorkflowPlan::Map { source, item, body } => {
1202                    let source = self.resolve_ref(source).await?;
1203                    let values = source.as_array().ok_or_else(|| {
1204                        failure(
1205                            WorkflowFailureCode::InvalidInput,
1206                            "map source must be an array",
1207                            false,
1208                        )
1209                    })?;
1210                    let used = self.ledger.lock().await.steps as usize;
1211                    let remaining = (self.root_limits.max_steps as usize).saturating_sub(used);
1212                    let per_item = plan_leaf_count(body).max(1);
1213                    if values
1214                        .len()
1215                        .checked_mul(per_item)
1216                        .is_none_or(|required| required > remaining)
1217                    {
1218                        return Err(failure(
1219                            WorkflowFailureCode::BudgetExceeded,
1220                            "map cardinality exceeds remaining workflow step budget",
1221                            false,
1222                        ));
1223                    }
1224                    let futures = values.iter().cloned().enumerate().map(|(index, value)| {
1225                        let mut child = self.clone();
1226                        child.items.insert(item.clone(), value);
1227                        // Scope identifies the logical map item, not a retry
1228                        // attempt's diagnostic path. This keeps durable attempts
1229                        // cumulative when Retry wraps Map and gives nested
1230                        // Parallel invocations an item-local cancellation domain.
1231                        child.scope = format!("{}[{index}]", self.scope);
1232                        async move { child.execute_node(body, &format!("{path}[{index}]")).await }
1233                    });
1234                    let results = join_all(futures).await;
1235                    let mut values = Vec::with_capacity(results.len());
1236                    let mut failures = Vec::new();
1237                    for (index, result) in results.into_iter().enumerate() {
1238                        match result {
1239                            Ok(value) => values.push(value),
1240                            Err(error) => failures.push((index, error)),
1241                        }
1242                    }
1243                    if failures.is_empty() {
1244                        Ok(Value::Array(values))
1245                    } else {
1246                        let retryable = failures.iter().any(|(_, error)| error.retryable);
1247                        let first_code = failures[0].1.code;
1248                        let code = if failures
1249                            .iter()
1250                            .any(|(_, error)| error.code == WorkflowFailureCode::DependencySkipped)
1251                        {
1252                            WorkflowFailureCode::DependencySkipped
1253                        } else if failures.iter().all(|(_, error)| error.code == first_code) {
1254                            first_code
1255                        } else {
1256                            WorkflowFailureCode::ExecutionFailed
1257                        };
1258                        let diagnostics = failures
1259                            .into_iter()
1260                            .map(|(index, error)| format!("item[{index}]: {}", error.message))
1261                            .collect::<Vec<_>>()
1262                            .join("; ");
1263                        Err(failure(
1264                            code,
1265                            format!("map items failed: {diagnostics}"),
1266                            retryable,
1267                        ))
1268                    }
1269                }
1270                WorkflowPlan::Retry {
1271                    node,
1272                    max_attempts,
1273                    delay_ms,
1274                } => {
1275                    let limit =
1276                        (*max_attempts).min(self.compiled.definition.budgets.max_retries + 1);
1277                    let mut last = None;
1278                    for attempt in 0..limit {
1279                        match self
1280                            .execute_node(node, &format!("{path}.retry{attempt}"))
1281                            .await
1282                        {
1283                            Ok(value) => return Ok(value),
1284                            Err(error) if error.retryable => {
1285                                last = Some(error);
1286                                if attempt + 1 < limit {
1287                                    self.reserve_retry().await?;
1288                                    self.checkpoint_usage("retry_reserved").await?;
1289                                    tokio::select! {
1290                                        _ = self.cancellation.cancelled() => return Err(failure(WorkflowFailureCode::Cancelled, "workflow cancelled", false)),
1291                                        _ = self.branch_cancellation.cancelled() => return Err(failure(WorkflowFailureCode::Cancelled, "workflow branch cancelled", false)),
1292                                        _ = tokio::time::sleep(Duration::from_millis(*delay_ms)) => {}
1293                                    }
1294                                } else {
1295                                    break;
1296                                }
1297                            }
1298                            Err(error) => return Err(error),
1299                        }
1300                    }
1301                    Err(failure(
1302                        WorkflowFailureCode::RetryExhausted,
1303                        last.map_or_else(|| "retry exhausted".to_string(), |error| error.message),
1304                        false,
1305                    ))
1306                }
1307            }
1308        })
1309    }
1310
1311    async fn execute_step(&self, step_id: &str, _path: &str) -> Result<Value, WorkflowFailure> {
1312        self.check_cancelled()?;
1313        let step = self.compiled.steps.get(step_id).cloned().ok_or_else(|| {
1314            failure(
1315                WorkflowFailureCode::UnknownReference,
1316                format!("unknown step {step_id}"),
1317                false,
1318            )
1319        })?;
1320        let instance_id = if self.scope == "root" {
1321            step_id.to_string()
1322        } else {
1323            format!("{step_id}@{}", self.scope)
1324        };
1325        let input = match &step.kind {
1326            WorkflowStepKind::Tool { args, .. } | WorkflowStepKind::Workflow { args, .. } => {
1327                self.resolve_template(args).await?
1328            }
1329            WorkflowStepKind::Agent { prompt, .. } => self.resolve_template(prompt).await?,
1330        };
1331        let input_hash = hex::encode(Sha256::digest(
1332            serde_json::to_vec(&input).unwrap_or_default(),
1333        ));
1334        self.reserve_step().await?;
1335        self.checkpoint_usage("step_reserved").await?;
1336        self.step_transition(&instance_id, WorkflowRunEventKind::StepQueued, |snapshot| {
1337            let state = snapshot
1338                .steps
1339                .entry(instance_id.clone())
1340                .or_insert_with(|| WorkflowStepSnapshot {
1341                    id: instance_id.clone(),
1342                    status: WorkflowStepStatus::Queued,
1343                    input_hash: input_hash.clone(),
1344                    output: None,
1345                    failure: None,
1346                    attempts: 0,
1347                });
1348            state.status = WorkflowStepStatus::Queued;
1349            state.input_hash = input_hash;
1350            state.output = None;
1351            state.failure = None;
1352        })
1353        .await?;
1354        let _permit = if matches!(&step.kind, WorkflowStepKind::Workflow { .. }) {
1355            None
1356        } else {
1357            Some(tokio::select! {
1358                _ = self.cancellation.cancelled() => {
1359                    let cancelled_id = instance_id.clone();
1360                    self.step_transition(&instance_id, WorkflowRunEventKind::StepCancelled, move |snapshot| {
1361                        if let Some(state) = snapshot.steps.get_mut(&cancelled_id) { state.status = WorkflowStepStatus::Cancelled; }
1362                    }).await?;
1363                    return Err(failure(WorkflowFailureCode::Cancelled, "workflow cancelled", false));
1364                }
1365                _ = self.branch_cancellation.cancelled() => {
1366                    let cancelled_id = instance_id.clone();
1367                    self.step_transition(&instance_id, WorkflowRunEventKind::StepCancelled, move |snapshot| {
1368                        if let Some(state) = snapshot.steps.get_mut(&cancelled_id) {
1369                            state.status = WorkflowStepStatus::Cancelled;
1370                        }
1371                    }).await?;
1372                    return Err(failure(WorkflowFailureCode::Cancelled, "workflow branch cancelled", false));
1373                }
1374                permit = self.semaphore.acquire() => permit.map_err(|_| failure(WorkflowFailureCode::ExecutionFailed, "workflow semaphore closed", false))?,
1375            })
1376        };
1377        let started_id = instance_id.clone();
1378        self.step_transition(
1379            &instance_id,
1380            WorkflowRunEventKind::StepStarted,
1381            move |snapshot| {
1382                if let Some(state) = snapshot.steps.get_mut(&started_id) {
1383                    state.status = WorkflowStepStatus::Running;
1384                    state.attempts += 1;
1385                }
1386            },
1387        )
1388        .await?;
1389        let result = self.dispatch(&step, input, &instance_id).await;
1390        let result = match result {
1391            Ok(output) => {
1392                if let Some(schema) = &step.output_schema {
1393                    validate_schema(schema, &output)
1394                        .map(|()| output)
1395                        .map_err(|message| {
1396                            failure(WorkflowFailureCode::InvalidOutput, message, false)
1397                        })
1398                } else {
1399                    Ok(output)
1400                }
1401            }
1402            Err(error) => Err(error),
1403        };
1404        let result = result.and_then(|output| {
1405            reject_secret_material(&output)
1406                .map(|()| output)
1407                .map_err(|message| failure(WorkflowFailureCode::InvalidOutput, message, false))
1408        });
1409        match result {
1410            Ok(output) => {
1411                let copy = output.clone();
1412                let completed_id = instance_id.clone();
1413                self.step_transition(
1414                    &instance_id,
1415                    WorkflowRunEventKind::StepCompleted {
1416                        output: output.clone(),
1417                    },
1418                    move |snapshot| {
1419                        if let Some(state) = snapshot.steps.get_mut(&completed_id) {
1420                            state.status = WorkflowStepStatus::Succeeded;
1421                            state.output = Some(copy);
1422                        }
1423                    },
1424                )
1425                .await?;
1426                Ok(output)
1427            }
1428            Err(error) => {
1429                if error.code == WorkflowFailureCode::Cancelled {
1430                    let cancelled_id = instance_id.clone();
1431                    self.step_transition(
1432                        &instance_id,
1433                        WorkflowRunEventKind::StepCancelled,
1434                        move |snapshot| {
1435                            if let Some(state) = snapshot.steps.get_mut(&cancelled_id) {
1436                                state.status = WorkflowStepStatus::Cancelled;
1437                                state.failure = Some(failure(
1438                                    WorkflowFailureCode::Cancelled,
1439                                    "workflow branch cancelled",
1440                                    false,
1441                                ));
1442                            }
1443                        },
1444                    )
1445                    .await?;
1446                    return Err(error);
1447                }
1448                if error.code == WorkflowFailureCode::Suspended {
1449                    let reason = error.message.clone();
1450                    let suspended_id = instance_id.clone();
1451                    self.step_transition(
1452                        &instance_id,
1453                        WorkflowRunEventKind::StepSuspended {
1454                            reason: reason.clone(),
1455                        },
1456                        move |snapshot| {
1457                            if let Some(state) = snapshot.steps.get_mut(&suspended_id) {
1458                                state.status = WorkflowStepStatus::Suspended;
1459                                state.failure =
1460                                    Some(failure(WorkflowFailureCode::Suspended, reason, true));
1461                            }
1462                        },
1463                    )
1464                    .await?;
1465                    return Err(error);
1466                }
1467                let copy = error.clone();
1468                let failed_id = instance_id.clone();
1469                self.step_transition(
1470                    &instance_id,
1471                    WorkflowRunEventKind::StepFailed {
1472                        failure: error.clone(),
1473                    },
1474                    move |snapshot| {
1475                        if let Some(state) = snapshot.steps.get_mut(&failed_id) {
1476                            state.status = WorkflowStepStatus::Failed;
1477                            state.failure = Some(copy);
1478                        }
1479                    },
1480                )
1481                .await?;
1482                match step.failure {
1483                    FailurePolicy::ContinueWithError => Ok(serde_json::json!({"error": error})),
1484                    FailurePolicy::SkipDependents => Err(failure(
1485                        WorkflowFailureCode::DependencySkipped,
1486                        format!("{} (dependents skipped)", error.message),
1487                        false,
1488                    )),
1489                    FailurePolicy::FailFast => Err(error),
1490                }
1491            }
1492        }
1493    }
1494
1495    async fn dispatch(
1496        &self,
1497        step: &WorkflowStepDefinition,
1498        input: Value,
1499        instance_id: &str,
1500    ) -> Result<Value, WorkflowFailure> {
1501        let session_id = { self.snapshot.lock().await.session_id.clone() };
1502        match &step.kind {
1503            WorkflowStepKind::Tool {
1504                tool, capabilities, ..
1505            } => {
1506                self.authorize(
1507                    &session_id,
1508                    WorkflowPolicyTarget::Tool(tool.clone()),
1509                    capabilities,
1510                )
1511                .await?;
1512                let (resolved_input, resolved_secrets) =
1513                    self.resolve_secret_handles(&input, &session_id).await?;
1514                let arguments = serde_json::to_string(&resolved_input).map_err(|error| {
1515                    failure(WorkflowFailureCode::InvalidInput, error.to_string(), false)
1516                })?;
1517                let call = ToolCall {
1518                    id: format!("workflow-{}", Uuid::new_v4()),
1519                    tool_type: "function".to_string(),
1520                    function: FunctionCall {
1521                        name: tool.clone(),
1522                        arguments,
1523                    },
1524                };
1525                let context = ToolExecutionContext {
1526                    session_id: Some(&session_id),
1527                    tool_call_id: &call.id,
1528                    event_tx: None,
1529                    available_tool_schemas: None,
1530                    bypass_permissions: false,
1531                    can_async_resume: false,
1532                    bash_completion_sink: None,
1533                    pre_parsed_args: Some(&resolved_input),
1534                };
1535                let outcome = self
1536                    .engine
1537                    .tools
1538                    .execute_with_context_outcome(&call, context)
1539                    .await
1540                    .map_err(|error| {
1541                        let (code, message, retryable) = match error {
1542                            bamboo_agent_core::tools::ToolError::NotFound(_) => (
1543                                WorkflowFailureCode::UnknownReference,
1544                                "workflow tool is not available",
1545                                false,
1546                            ),
1547                            bamboo_agent_core::tools::ToolError::InvalidArguments(_) => (
1548                                WorkflowFailureCode::InvalidInput,
1549                                "workflow tool arguments were rejected",
1550                                false,
1551                            ),
1552                            bamboo_agent_core::tools::ToolError::Execution(_) => (
1553                                WorkflowFailureCode::ExecutionFailed,
1554                                "workflow tool execution was denied or failed",
1555                                true,
1556                            ),
1557                        };
1558                        failure(code, message, retryable)
1559                    })?;
1560                match outcome {
1561                    ToolOutcome::Completed(result) => {
1562                        let output = parse_tool_result(result)?;
1563                        if contains_any_secret_material(&output, &resolved_secrets) {
1564                            return Err(failure(
1565                                WorkflowFailureCode::InvalidOutput,
1566                                "workflow tool output contained resolved secret material",
1567                                false,
1568                            ));
1569                        }
1570                        Ok(output)
1571                    }
1572                    ToolOutcome::NeedsHuman { question, .. } => {
1573                        self.persist_suspension(WorkflowSuspensionContext::ToolApproval {
1574                            step_id: instance_id.to_string(),
1575                            tool: tool.clone(),
1576                            tool_call_id: question.tool_call_id,
1577                        })
1578                        .await?;
1579                        Err(failure(
1580                            WorkflowFailureCode::Suspended,
1581                            "workflow tool requires human approval",
1582                            true,
1583                        ))
1584                    }
1585                    ToolOutcome::Running(handle) => {
1586                        let tool_call_id = handle.tool_call_id.clone();
1587                        (handle.kill)();
1588                        self.persist_suspension(WorkflowSuspensionContext::ToolRunning {
1589                            step_id: instance_id.to_string(),
1590                            tool: tool.clone(),
1591                            tool_call_id,
1592                            killed: true,
1593                        })
1594                        .await?;
1595                        Err(failure(
1596                            WorkflowFailureCode::Suspended,
1597                            "workflow tool is running without a durable workflow resume handle",
1598                            true,
1599                        ))
1600                    }
1601                }
1602            }
1603            WorkflowStepKind::Agent {
1604                agent,
1605                model,
1606                effort,
1607                capabilities,
1608                structured_output_attempts,
1609                ..
1610            } => {
1611                if contains_secret_handle(&input) {
1612                    return Err(failure(
1613                        WorkflowFailureCode::PermissionDenied,
1614                        "secret capability handles are supported only for tool arguments",
1615                        false,
1616                    ));
1617                }
1618                self.authorize(
1619                    &session_id,
1620                    WorkflowPolicyTarget::Agent(agent.clone()),
1621                    capabilities,
1622                )
1623                .await?;
1624                let spec = self.pinned_agents.get(agent).cloned().ok_or_else(|| {
1625                    failure(
1626                        WorkflowFailureCode::PermissionDenied,
1627                        "named agent was not pinned during preflight",
1628                        false,
1629                    )
1630                })?;
1631                let requested = capabilities.iter().cloned().collect::<BTreeSet<_>>();
1632                if !requested.is_subset(&spec.allowed_capabilities) {
1633                    return Err(failure(
1634                        WorkflowFailureCode::PermissionDenied,
1635                        "named agent capability intersection changed",
1636                        false,
1637                    ));
1638                }
1639                let mut last_error = None;
1640                for _ in 0..*structured_output_attempts {
1641                    self.reserve_agent().await?;
1642                    self.checkpoint_usage("agent_reserved").await?;
1643                    match self
1644                        .engine
1645                        .agents
1646                        .execute(
1647                            &spec,
1648                            input.clone(),
1649                            model.as_deref(),
1650                            effort.as_deref(),
1651                            &requested,
1652                            &session_id,
1653                        )
1654                        .await
1655                    {
1656                        Ok(result) => {
1657                            let exceeded =
1658                                self.record_usage(result.tokens, result.cost_micros).await;
1659                            self.checkpoint_usage("agent_usage_recorded").await?;
1660                            if let Some(error) = exceeded {
1661                                return Err(error);
1662                            }
1663                            if let Some(schema) = &step.output_schema {
1664                                if let Err(error) = validate_schema(schema, &result.output) {
1665                                    last_error = Some(error);
1666                                    continue;
1667                                }
1668                            }
1669                            return Ok(result.output);
1670                        }
1671                        Err(_error) => {
1672                            last_error = Some("named agent execution failed".to_string())
1673                        }
1674                    }
1675                }
1676                Err(failure(
1677                    WorkflowFailureCode::InvalidOutput,
1678                    last_error.unwrap_or_else(|| "agent structured output exhausted".to_string()),
1679                    false,
1680                ))
1681            }
1682            WorkflowStepKind::Workflow {
1683                workflow_id,
1684                revision,
1685                ..
1686            } => {
1687                if self.depth + 1 >= self.root_limits.max_nesting_depth {
1688                    return Err(failure(
1689                        WorkflowFailureCode::BudgetExceeded,
1690                        "nested workflow depth exceeded",
1691                        false,
1692                    ));
1693                }
1694                let definition = self
1695                    .bundle
1696                    .get(workflow_id, *revision)
1697                    .cloned()
1698                    .ok_or_else(|| {
1699                        failure(
1700                            WorkflowFailureCode::UnknownReference,
1701                            format!("persisted bundle missing workflow {workflow_id}@{revision}"),
1702                            false,
1703                        )
1704                    })?;
1705                let nested = StartWorkflowRun {
1706                    definition,
1707                    args: input,
1708                    session_id,
1709                    workspace_trusted: self.workspace_trusted,
1710                    allowed_capabilities: self.allowed_capabilities.iter().cloned().collect(),
1711                };
1712                let parent_run_id = self.snapshot.lock().await.run_id.clone();
1713                let result = Box::pin(self.engine.run_internal(
1714                    nested,
1715                    self.bundle.clone(),
1716                    self.pinned_agents.clone(),
1717                    Some(parent_run_id),
1718                    Some(instance_id.to_string()),
1719                    self.depth + 1,
1720                    self.branch_cancellation.clone(),
1721                    self.ledger.clone(),
1722                    self.root_limits.clone(),
1723                    self.semaphore.clone(),
1724                    None,
1725                ))
1726                .await
1727                .map_err(|_error| {
1728                    failure(
1729                        WorkflowFailureCode::ExecutionFailed,
1730                        "nested workflow execution failed",
1731                        false,
1732                    )
1733                })?;
1734                result.output.ok_or_else(|| {
1735                    result.failure.unwrap_or_else(|| {
1736                        failure(
1737                            WorkflowFailureCode::ExecutionFailed,
1738                            "nested workflow returned no output",
1739                            false,
1740                        )
1741                    })
1742                })
1743            }
1744        }
1745    }
1746
1747    async fn authorize(
1748        &self,
1749        session_id: &str,
1750        target: WorkflowPolicyTarget,
1751        capabilities: &[String],
1752    ) -> Result<(), WorkflowFailure> {
1753        let requested = capabilities.iter().cloned().collect::<BTreeSet<_>>();
1754        if !requested.is_subset(&self.allowed_capabilities) {
1755            return Err(failure(
1756                WorkflowFailureCode::PermissionDenied,
1757                "step capability exceeds root policy",
1758                false,
1759            ));
1760        }
1761        match self
1762            .engine
1763            .policy
1764            .authorize(session_id, &target, &requested, self.workspace_trusted)
1765            .await
1766        {
1767            PermissionDecision::Allow => Ok(()),
1768            PermissionDecision::Deny(_reason) => Err(failure(
1769                if self.workspace_trusted {
1770                    WorkflowFailureCode::PermissionDenied
1771                } else {
1772                    WorkflowFailureCode::UntrustedWorkspace
1773                },
1774                "workflow policy denied this step",
1775                false,
1776            )),
1777        }
1778    }
1779
1780    async fn step_transition(
1781        &self,
1782        step_id: &str,
1783        kind: WorkflowRunEventKind,
1784        mutate: impl FnOnce(&mut WorkflowRunSnapshot),
1785    ) -> Result<(), WorkflowFailure> {
1786        let usage = self.ledger.lock().await.clone();
1787        let mut snapshot = self.snapshot.lock().await;
1788        snapshot.usage = usage;
1789        self.engine
1790            .transition(&mut snapshot, Some(step_id.to_string()), kind, mutate)
1791            .await
1792            .map_err(|error| failure(WorkflowFailureCode::Storage, error.to_string(), false))
1793    }
1794
1795    async fn reserve_step(&self) -> Result<(), WorkflowFailure> {
1796        let mut usage = self.ledger.lock().await;
1797        if usage.steps >= self.root_limits.max_steps {
1798            return Err(failure(
1799                WorkflowFailureCode::BudgetExceeded,
1800                "workflow step budget exceeded",
1801                false,
1802            ));
1803        }
1804        usage.steps += 1;
1805        Ok(())
1806    }
1807
1808    async fn reserve_retry(&self) -> Result<(), WorkflowFailure> {
1809        let mut usage = self.ledger.lock().await;
1810        if usage.retries >= self.root_limits.max_retries {
1811            return Err(failure(
1812                WorkflowFailureCode::BudgetExceeded,
1813                "workflow retry budget exceeded",
1814                false,
1815            ));
1816        }
1817        usage.retries += 1;
1818        Ok(())
1819    }
1820
1821    async fn reserve_agent(&self) -> Result<(), WorkflowFailure> {
1822        let mut usage = self.ledger.lock().await;
1823        if usage.agents >= self.root_limits.max_agents {
1824            return Err(failure(
1825                WorkflowFailureCode::BudgetExceeded,
1826                "workflow agent budget exceeded",
1827                false,
1828            ));
1829        }
1830        usage.agents += 1;
1831        Ok(())
1832    }
1833
1834    async fn record_usage(&self, tokens: u64, cost_micros: u64) -> Option<WorkflowFailure> {
1835        let mut usage = self.ledger.lock().await;
1836        let next_tokens = usage.tokens.saturating_add(tokens);
1837        let next_cost = usage.cost_micros.saturating_add(cost_micros);
1838        usage.tokens = next_tokens;
1839        usage.cost_micros = next_cost;
1840        if self
1841            .root_limits
1842            .max_tokens
1843            .is_some_and(|limit| next_tokens > limit)
1844            || self
1845                .root_limits
1846                .max_cost_micros
1847                .is_some_and(|limit| next_cost > limit)
1848        {
1849            return Some(failure(
1850                WorkflowFailureCode::BudgetExceeded,
1851                "workflow token/cost budget exceeded",
1852                false,
1853            ));
1854        }
1855        None
1856    }
1857
1858    async fn checkpoint_usage(&self, name: &str) -> Result<(), WorkflowFailure> {
1859        let usage = self.ledger.lock().await.clone();
1860        let mut snapshot = self.snapshot.lock().await;
1861        self.engine
1862            .transition(
1863                &mut snapshot,
1864                None,
1865                WorkflowRunEventKind::Phase {
1866                    name: name.to_string(),
1867                },
1868                move |snapshot| snapshot.usage = usage,
1869            )
1870            .await
1871            .map_err(|error| failure(WorkflowFailureCode::Storage, error.to_string(), false))
1872    }
1873
1874    async fn persist_suspension(
1875        &self,
1876        context: WorkflowSuspensionContext,
1877    ) -> Result<(), WorkflowFailure> {
1878        let mut snapshot = self.snapshot.lock().await;
1879        self.engine
1880            .transition(
1881                &mut snapshot,
1882                None,
1883                WorkflowRunEventKind::Phase {
1884                    name: "suspension_context_persisted".to_string(),
1885                },
1886                move |snapshot| snapshot.suspension = Some(context),
1887            )
1888            .await
1889            .map_err(|error| failure(WorkflowFailureCode::Storage, error.to_string(), false))
1890    }
1891
1892    async fn cancel_active_parallel_steps(
1893        &self,
1894        nodes: &[WorkflowPlan],
1895    ) -> Result<(), WorkflowFailure> {
1896        let sibling_steps = nodes
1897            .iter()
1898            .flat_map(plan_step_ids)
1899            .collect::<BTreeSet<_>>();
1900        let active = {
1901            let snapshot = self.snapshot.lock().await;
1902            snapshot
1903                .steps
1904                .iter()
1905                .filter(|(id, step)| {
1906                    matches!(
1907                        step.status,
1908                        WorkflowStepStatus::Queued | WorkflowStepStatus::Running
1909                    ) && sibling_steps
1910                        .iter()
1911                        .any(|step_id| instance_is_in_scope(id, step_id, &self.scope))
1912                })
1913                .map(|(id, _)| id.clone())
1914                .collect::<Vec<_>>()
1915        };
1916        for step_id in active {
1917            let state_id = step_id.clone();
1918            self.step_transition(
1919                &step_id,
1920                WorkflowRunEventKind::StepCancelled,
1921                move |snapshot| {
1922                    if let Some(step) = snapshot.steps.get_mut(&state_id) {
1923                        step.status = WorkflowStepStatus::Cancelled;
1924                        step.failure = Some(failure(
1925                            WorkflowFailureCode::Cancelled,
1926                            "parallel sibling cancelled by fail_fast",
1927                            false,
1928                        ));
1929                    }
1930                },
1931            )
1932            .await?;
1933        }
1934        Ok(())
1935    }
1936
1937    async fn resolve_secret_handles(
1938        &self,
1939        value: &Value,
1940        session_id: &str,
1941    ) -> Result<(Value, Vec<String>), WorkflowFailure> {
1942        fn walk<'a>(
1943            context: &'a RunContext,
1944            value: &'a Value,
1945            session_id: &'a str,
1946        ) -> SecretResolutionFuture<'a> {
1947            Box::pin(async move {
1948                match value {
1949                    Value::Object(object) if object.contains_key("$secret") => {
1950                        let handle: bamboo_domain::WorkflowSecretHandle =
1951                            serde_json::from_value(value.clone()).map_err(|_| {
1952                                failure(
1953                                    WorkflowFailureCode::InvalidInput,
1954                                    "malformed secret capability handle",
1955                                    false,
1956                                )
1957                            })?;
1958                        let material = context
1959                            .engine
1960                            .secrets
1961                            .resolve(session_id, &handle.capability)
1962                            .await
1963                            .map_err(|_| {
1964                                failure(
1965                                    WorkflowFailureCode::PermissionDenied,
1966                                    "secret capability resolution denied",
1967                                    false,
1968                                )
1969                            })?;
1970                        let material = material.into_exposed();
1971                        Ok((Value::String(material.clone()), vec![material]))
1972                    }
1973                    Value::Object(object) => {
1974                        let mut resolved = serde_json::Map::new();
1975                        let mut secrets = Vec::new();
1976                        for (key, child) in object {
1977                            let (child, mut child_secrets) =
1978                                walk(context, child, session_id).await?;
1979                            resolved.insert(key.clone(), child);
1980                            secrets.append(&mut child_secrets);
1981                        }
1982                        Ok((Value::Object(resolved), secrets))
1983                    }
1984                    Value::Array(array) => {
1985                        let mut resolved = Vec::with_capacity(array.len());
1986                        let mut secrets = Vec::new();
1987                        for child in array {
1988                            let (child, mut child_secrets) =
1989                                walk(context, child, session_id).await?;
1990                            resolved.push(child);
1991                            secrets.append(&mut child_secrets);
1992                        }
1993                        Ok((Value::Array(resolved), secrets))
1994                    }
1995                    value => Ok((value.clone(), Vec::new())),
1996                }
1997            })
1998        }
1999        walk(self, value, session_id).await
2000    }
2001
2002    fn skip_plan<'a>(&'a self, plan: &'a WorkflowPlan, reason: &'a str) -> NodeFuture<'a> {
2003        Box::pin(async move {
2004            match plan {
2005                WorkflowPlan::Step { step } => {
2006                    let instance_id = if self.scope == "root" {
2007                        step.clone()
2008                    } else {
2009                        format!("{step}@{}", self.scope)
2010                    };
2011                    let reason_owned = reason.to_string();
2012                    let state_id = instance_id.clone();
2013                    self.step_transition(
2014                        &instance_id,
2015                        WorkflowRunEventKind::StepSkipped {
2016                            reason: reason.to_string(),
2017                        },
2018                        move |snapshot| {
2019                            let state = snapshot.steps.entry(state_id.clone()).or_insert(
2020                                WorkflowStepSnapshot {
2021                                    id: state_id,
2022                                    status: WorkflowStepStatus::Skipped,
2023                                    input_hash: String::new(),
2024                                    output: None,
2025                                    failure: Some(failure(
2026                                        WorkflowFailureCode::DependencySkipped,
2027                                        reason_owned.clone(),
2028                                        false,
2029                                    )),
2030                                    attempts: 0,
2031                                },
2032                            );
2033                            state.status = WorkflowStepStatus::Skipped;
2034                            state.failure = Some(failure(
2035                                WorkflowFailureCode::DependencySkipped,
2036                                reason_owned,
2037                                false,
2038                            ));
2039                        },
2040                    )
2041                    .await?;
2042                }
2043                WorkflowPlan::Sequence { nodes } | WorkflowPlan::Parallel { nodes } => {
2044                    for node in nodes {
2045                        self.skip_plan(node, reason).await?;
2046                    }
2047                }
2048                WorkflowPlan::Map { body, .. } | WorkflowPlan::Retry { node: body, .. } => {
2049                    self.skip_plan(body, reason).await?;
2050                }
2051            }
2052            Ok(Value::Null)
2053        })
2054    }
2055
2056    async fn resolve_template(&self, value: &Value) -> Result<Value, WorkflowFailure> {
2057        Box::pin(self.resolve_template_inner(value)).await
2058    }
2059
2060    fn resolve_template_inner<'a>(
2061        &'a self,
2062        value: &'a Value,
2063    ) -> Pin<Box<dyn Future<Output = Result<Value, WorkflowFailure>> + Send + 'a>> {
2064        Box::pin(async move {
2065            match value {
2066                Value::Object(object) if object.get("from").is_some() => {
2067                    let reference: ValueRef =
2068                        serde_json::from_value(value.clone()).map_err(|error| {
2069                            failure(
2070                                WorkflowFailureCode::InvalidInput,
2071                                format!("malformed value reference: {error}"),
2072                                false,
2073                            )
2074                        })?;
2075                    self.resolve_ref(&reference).await
2076                }
2077                Value::Object(object) => {
2078                    let mut resolved = serde_json::Map::new();
2079                    for (key, child) in object {
2080                        resolved.insert(key.clone(), self.resolve_template_inner(child).await?);
2081                    }
2082                    Ok(Value::Object(resolved))
2083                }
2084                Value::Array(array) => {
2085                    let mut resolved = Vec::with_capacity(array.len());
2086                    for child in array {
2087                        resolved.push(self.resolve_template_inner(child).await?);
2088                    }
2089                    Ok(Value::Array(resolved))
2090                }
2091                value => Ok(value.clone()),
2092            }
2093        })
2094    }
2095
2096    async fn resolve_ref(&self, reference: &ValueRef) -> Result<Value, WorkflowFailure> {
2097        let (root, pointer) = match reference {
2098            ValueRef::Args { pointer } => (
2099                self.snapshot.lock().await.validated_args.clone(),
2100                pointer.as_str(),
2101            ),
2102            ValueRef::Step { step, pointer } => {
2103                let snapshot = self.snapshot.lock().await;
2104                let exact = format!("{step}@{}", self.scope);
2105                let output = snapshot
2106                    .steps
2107                    .get(&exact)
2108                    .or_else(|| snapshot.steps.get(step))
2109                    .and_then(|state| state.output.clone())
2110                    .ok_or_else(|| {
2111                        failure(
2112                            WorkflowFailureCode::UnknownReference,
2113                            format!("step output '{step}' unavailable in execution scope"),
2114                            false,
2115                        )
2116                    })?;
2117                (output, pointer.as_str())
2118            }
2119            ValueRef::Item { name, pointer } => (
2120                self.items.get(name).cloned().ok_or_else(|| {
2121                    failure(
2122                        WorkflowFailureCode::UnknownReference,
2123                        format!("map item '{name}' unavailable"),
2124                        false,
2125                    )
2126                })?,
2127                pointer.as_str(),
2128            ),
2129            ValueRef::Literal { value } => return Ok(value.clone()),
2130        };
2131        if pointer.is_empty() {
2132            Ok(root)
2133        } else {
2134            root.pointer(pointer).cloned().ok_or_else(|| {
2135                failure(
2136                    WorkflowFailureCode::UnknownReference,
2137                    format!("JSON pointer '{pointer}' not found"),
2138                    false,
2139                )
2140            })
2141        }
2142    }
2143
2144    fn check_cancelled(&self) -> Result<(), WorkflowFailure> {
2145        if self.cancellation.is_cancelled() || self.branch_cancellation.is_cancelled() {
2146            Err(failure(
2147                WorkflowFailureCode::Cancelled,
2148                "workflow cancelled",
2149                false,
2150            ))
2151        } else {
2152            Ok(())
2153        }
2154    }
2155}
2156
2157fn event(
2158    snapshot: &WorkflowRunSnapshot,
2159    step_id: Option<String>,
2160    kind: WorkflowRunEventKind,
2161) -> WorkflowRunEvent {
2162    WorkflowRunEvent {
2163        run_id: snapshot.run_id.clone(),
2164        sequence: snapshot.last_sequence,
2165        at: Utc::now(),
2166        step_id,
2167        kind,
2168    }
2169}
2170fn failure(
2171    code: WorkflowFailureCode,
2172    message: impl Into<String>,
2173    retryable: bool,
2174) -> WorkflowFailure {
2175    WorkflowFailure {
2176        code,
2177        message: message.into(),
2178        retryable,
2179    }
2180}
2181fn storage(error: std::io::Error) -> WorkflowRunError {
2182    WorkflowRunError::Storage(error.to_string())
2183}
2184fn exceeds_optional(requested: Option<u64>, ceiling: Option<u64>) -> bool {
2185    match (requested, ceiling) {
2186        (Some(requested), Some(ceiling)) => requested > ceiling,
2187        _ => false,
2188    }
2189}
2190
2191fn enforce_budget_within(
2192    requested: &WorkflowBudgets,
2193    ceiling: &WorkflowBudgets,
2194) -> Result<(), &'static str> {
2195    if requested.max_concurrency > ceiling.max_concurrency {
2196        Err("max_concurrency")
2197    } else if requested.max_agents > ceiling.max_agents {
2198        Err("max_agents")
2199    } else if requested.max_steps > ceiling.max_steps {
2200        Err("max_steps")
2201    } else if requested.max_retries > ceiling.max_retries {
2202        Err("max_retries")
2203    } else if requested.max_nesting_depth > ceiling.max_nesting_depth {
2204        Err("max_nesting_depth")
2205    } else if requested.wall_time_ms > ceiling.wall_time_ms {
2206        Err("wall_time_ms")
2207    } else if exceeds_optional(requested.max_tokens, ceiling.max_tokens) {
2208        Err("max_tokens")
2209    } else if exceeds_optional(requested.max_cost_micros, ceiling.max_cost_micros) {
2210        Err("max_cost_micros")
2211    } else {
2212        Ok(())
2213    }
2214}
2215
2216fn definition_bundle_hash(bundle: &WorkflowDefinitionBundle) -> Result<String, WorkflowRunError> {
2217    let bytes = serde_json::to_vec(bundle)
2218        .map_err(|_| WorkflowRunError::Preflight("workflow bundle hashing failed".to_string()))?;
2219    Ok(hex::encode(Sha256::digest(bytes)))
2220}
2221fn parse_tool_result(result: ToolResult) -> Result<Value, WorkflowFailure> {
2222    if !result.success {
2223        return Err(failure(
2224            WorkflowFailureCode::ExecutionFailed,
2225            "workflow tool reported failure",
2226            true,
2227        ));
2228    }
2229    Ok(serde_json::from_str(&result.result).unwrap_or(Value::String(result.result)))
2230}
2231fn reject_secret_material(value: &Value) -> Result<(), String> {
2232    reject_secret_material_inner(value, false)
2233}
2234
2235fn reject_secret_material_in_definition(value: &Value) -> Result<(), String> {
2236    reject_secret_material_inner(value, true)
2237}
2238
2239fn reject_secret_material_inner(value: &Value, allow_bindings: bool) -> Result<(), String> {
2240    fn walk(value: &Value, key: Option<&str>, allow_bindings: bool) -> Result<(), String> {
2241        if value.as_object().is_some_and(|object| {
2242            object.len() == 1
2243                && object
2244                    .get("$secret")
2245                    .and_then(Value::as_str)
2246                    .is_some_and(|handle| !handle.trim().is_empty())
2247        }) {
2248            return Ok(());
2249        }
2250        let safe_binding = allow_bindings
2251            && serde_json::from_value::<ValueRef>(value.clone())
2252                .is_ok_and(|reference| !matches!(reference, ValueRef::Literal { .. }));
2253        if key.is_some_and(|key| {
2254            let normalized = key
2255                .chars()
2256                .filter(|character| character.is_ascii_alphanumeric())
2257                .flat_map(char::to_lowercase)
2258                .collect::<String>();
2259            matches!(
2260                normalized.as_str(),
2261                "secret"
2262                    | "token"
2263                    | "password"
2264                    | "credential"
2265                    | "credentials"
2266                    | "apikey"
2267                    | "accesskey"
2268                    | "accesstoken"
2269                    | "secretkey"
2270                    | "privatekey"
2271            )
2272        }) && !safe_binding
2273        {
2274            return Err("secret-bearing fields are not accepted by workflow runs".to_string());
2275        }
2276        if value.as_str().is_some_and(|value| {
2277            let trimmed = value.trim();
2278            trimmed.starts_with("capability://")
2279                || trimmed.starts_with("Bearer ")
2280                || trimmed.starts_with("sk-")
2281                || trimmed.starts_with("ghp_")
2282                || trimmed.starts_with("github_pat_")
2283        }) {
2284            // No production capability resolver is part of #578. Treating an
2285            // arbitrary caller string as an opaque handle would be an injection
2286            // channel, so handles and common raw credential forms fail closed.
2287            return Err("opaque credential handles are not enabled for workflows".to_string());
2288        }
2289        match value {
2290            Value::Object(object) => {
2291                for (key, value) in object {
2292                    if key == "properties" {
2293                        let properties = value.as_object().ok_or_else(|| {
2294                            "workflow schema properties must be an object".to_string()
2295                        })?;
2296                        for schema in properties.values() {
2297                            walk(schema, None, allow_bindings)?;
2298                        }
2299                    } else {
2300                        walk(value, Some(key), allow_bindings)?;
2301                    }
2302                }
2303            }
2304            Value::Array(array) => {
2305                for value in array {
2306                    walk(value, None, allow_bindings)?;
2307                }
2308            }
2309            _ => {}
2310        }
2311        Ok(())
2312    }
2313    walk(value, None, allow_bindings)
2314}
2315
2316fn contains_secret_handle(value: &Value) -> bool {
2317    match value {
2318        Value::Object(object) => {
2319            object.contains_key("$secret") || object.values().any(contains_secret_handle)
2320        }
2321        Value::Array(array) => array.iter().any(contains_secret_handle),
2322        _ => false,
2323    }
2324}
2325
2326fn contains_any_secret_material(value: &Value, secrets: &[String]) -> bool {
2327    let matches = |candidate: &str| {
2328        secrets
2329            .iter()
2330            .any(|secret| !secret.is_empty() && candidate.contains(secret))
2331    };
2332    fn walk(value: &Value, matches: &impl Fn(&str) -> bool) -> bool {
2333        match value {
2334            Value::String(value) => matches(value),
2335            Value::Object(object) => object
2336                .iter()
2337                .any(|(key, value)| matches(key) || walk(value, matches)),
2338            Value::Array(array) => array.iter().any(|value| walk(value, matches)),
2339            _ => false,
2340        }
2341    }
2342    walk(value, &matches)
2343}
2344
2345fn plan_leaf_count(plan: &WorkflowPlan) -> usize {
2346    match plan {
2347        WorkflowPlan::Step { .. } => 1,
2348        WorkflowPlan::Sequence { nodes } | WorkflowPlan::Parallel { nodes } => {
2349            nodes.iter().fold(0usize, |total, node| {
2350                total.saturating_add(plan_leaf_count(node))
2351            })
2352        }
2353        WorkflowPlan::Map { body, .. } | WorkflowPlan::Retry { node: body, .. } => {
2354            plan_leaf_count(body)
2355        }
2356    }
2357}
2358
2359fn plan_step_ids(plan: &WorkflowPlan) -> Vec<String> {
2360    match plan {
2361        WorkflowPlan::Step { step } => vec![step.clone()],
2362        WorkflowPlan::Sequence { nodes } | WorkflowPlan::Parallel { nodes } => {
2363            nodes.iter().flat_map(plan_step_ids).collect()
2364        }
2365        WorkflowPlan::Map { body, .. } | WorkflowPlan::Retry { node: body, .. } => {
2366            plan_step_ids(body)
2367        }
2368    }
2369}
2370
2371fn instance_is_in_scope(instance_id: &str, step_id: &str, scope: &str) -> bool {
2372    if scope == "root" {
2373        instance_id == step_id
2374            || instance_id
2375                .strip_prefix(&format!("{step_id}@root"))
2376                .is_some_and(|suffix| suffix.starts_with('['))
2377    } else {
2378        let exact = format!("{step_id}@{scope}");
2379        instance_id == exact
2380            || instance_id
2381                .strip_prefix(&exact)
2382                .is_some_and(|suffix| suffix.starts_with('['))
2383    }
2384}
2385
2386fn plan_frontier(plan: &WorkflowPlan) -> Vec<String> {
2387    match plan {
2388        WorkflowPlan::Step { step } => vec![step.clone()],
2389        WorkflowPlan::Sequence { nodes } => nodes.first().map_or_else(Vec::new, plan_frontier),
2390        WorkflowPlan::Parallel { nodes } => nodes.iter().flat_map(plan_frontier).collect(),
2391        WorkflowPlan::Map { body, .. } | WorkflowPlan::Retry { node: body, .. } => {
2392            plan_frontier(body)
2393        }
2394    }
2395}
2396
2397fn validate_nested_input_contract(
2398    template: &Value,
2399    target_schema: &Value,
2400    compiled: &CompiledWorkflow,
2401) -> Result<(), String> {
2402    fn contains_ref(value: &Value) -> bool {
2403        match value {
2404            Value::Object(object) => {
2405                object.contains_key("from") || object.values().any(contains_ref)
2406            }
2407            Value::Array(array) => array.iter().any(contains_ref),
2408            _ => false,
2409        }
2410    }
2411    if !contains_ref(template) {
2412        return validate_schema(target_schema, template)
2413            .map_err(|error| format!("nested workflow input is incompatible: {error}"));
2414    }
2415    let reference: ValueRef = serde_json::from_value(template.clone()).map_err(|_| {
2416        "nested dynamic input schema cannot be proven compatible in phase 1".to_string()
2417    })?;
2418    let source_schema = match reference {
2419        ValueRef::Args { pointer } => {
2420            schema_at_pointer(&compiled.definition.input_schema, &pointer)
2421        }
2422        ValueRef::Step { step, pointer } => compiled
2423            .steps
2424            .get(&step)
2425            .and_then(|step| step.output_schema.as_ref())
2426            .and_then(|schema| schema_at_pointer(schema, &pointer)),
2427        ValueRef::Literal { value } => {
2428            return validate_schema(target_schema, &value)
2429                .map_err(|error| format!("nested workflow input is incompatible: {error}"));
2430        }
2431        ValueRef::Item { .. } => None,
2432    }
2433    .ok_or_else(|| {
2434        "nested dynamic input source schema is missing or pointer is invalid".to_string()
2435    })?;
2436    if schema_compatible(source_schema, target_schema) {
2437        Ok(())
2438    } else {
2439        Err("nested workflow input schema is not compatible with its pinned target".to_string())
2440    }
2441}
2442
2443fn schema_at_pointer<'a>(schema: &'a Value, pointer: &str) -> Option<&'a Value> {
2444    if pointer.is_empty() {
2445        return Some(schema);
2446    }
2447    let mut current = schema;
2448    for token in pointer.strip_prefix('/')?.split('/') {
2449        let token = token.replace("~1", "/").replace("~0", "~");
2450        current = if token.parse::<usize>().is_ok() {
2451            current.get("items")?
2452        } else {
2453            current.get("properties")?.get(&token)?
2454        };
2455    }
2456    Some(current)
2457}
2458
2459fn schema_compatible(source: &Value, target: &Value) -> bool {
2460    if source == target {
2461        return true;
2462    }
2463    let source_type = source.get("type").and_then(Value::as_str);
2464    let target_type = target.get("type").and_then(Value::as_str);
2465    source_type.is_some() && source_type == target_type && target_type != Some("object")
2466}
2467
2468fn effective_limits(requested: &WorkflowBudgets, ceilings: &WorkflowBudgets) -> WorkflowBudgets {
2469    WorkflowBudgets {
2470        max_concurrency: requested.max_concurrency.min(ceilings.max_concurrency),
2471        max_agents: requested.max_agents.min(ceilings.max_agents),
2472        max_steps: requested.max_steps.min(ceilings.max_steps),
2473        max_retries: requested.max_retries.min(ceilings.max_retries),
2474        max_nesting_depth: requested.max_nesting_depth.min(ceilings.max_nesting_depth),
2475        wall_time_ms: requested.wall_time_ms.min(ceilings.wall_time_ms),
2476        max_tokens: match (requested.max_tokens, ceilings.max_tokens) {
2477            (Some(requested), Some(ceiling)) => Some(requested.min(ceiling)),
2478            (Some(requested), None) => Some(requested),
2479            (None, ceiling) => ceiling,
2480        },
2481        max_cost_micros: match (requested.max_cost_micros, ceilings.max_cost_micros) {
2482            (Some(requested), Some(ceiling)) => Some(requested.min(ceiling)),
2483            (Some(requested), None) => Some(requested),
2484            (None, ceiling) => ceiling,
2485        },
2486    }
2487}