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