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                    session_id: Some(&session_id),
1571                    tool_call_id: &call.id,
1572                    event_tx: None,
1573                    available_tool_schemas: None,
1574                    bypass_permissions: session_flags.bypass_permissions,
1575                    auto_approve_permissions: session_flags.auto_approve_permissions,
1576                    plan_read_only: session_flags.plan_read_only,
1577                    can_async_resume: false,
1578                    bash_completion_sink: None,
1579                    pre_parsed_args: Some(&resolved_input),
1580                };
1581                if session_flags.plan_read_only
1582                    && !bamboo_tools::orchestrator::plan_mode_allows_tool(&call.function.name)
1583                {
1584                    return Err(failure(
1585                        WorkflowFailureCode::ExecutionFailed,
1586                        format!("Plan mode: {} operation blocked", call.function.name),
1587                        false,
1588                    ));
1589                }
1590                let outcome = self
1591                    .engine
1592                    .tools
1593                    .execute_with_context_outcome(&call, context)
1594                    .await
1595                    .map_err(|error| {
1596                        let (code, message, retryable) = match error {
1597                            bamboo_agent_core::tools::ToolError::NotFound(_) => (
1598                                WorkflowFailureCode::UnknownReference,
1599                                "workflow tool is not available",
1600                                false,
1601                            ),
1602                            bamboo_agent_core::tools::ToolError::InvalidArguments(_) => (
1603                                WorkflowFailureCode::InvalidInput,
1604                                "workflow tool arguments were rejected",
1605                                false,
1606                            ),
1607                            bamboo_agent_core::tools::ToolError::Execution(_) => (
1608                                WorkflowFailureCode::ExecutionFailed,
1609                                "workflow tool execution was denied or failed",
1610                                true,
1611                            ),
1612                        };
1613                        failure(code, message, retryable)
1614                    })?;
1615                match outcome {
1616                    ToolOutcome::Completed(result) => {
1617                        let output = parse_tool_result(result)?;
1618                        if contains_any_secret_material(&output, &resolved_secrets) {
1619                            return Err(failure(
1620                                WorkflowFailureCode::InvalidOutput,
1621                                "workflow tool output contained resolved secret material",
1622                                false,
1623                            ));
1624                        }
1625                        Ok(output)
1626                    }
1627                    ToolOutcome::NeedsHuman { question, .. } => {
1628                        self.persist_suspension(WorkflowSuspensionContext::ToolApproval {
1629                            step_id: instance_id.to_string(),
1630                            tool: tool.clone(),
1631                            tool_call_id: question.tool_call_id,
1632                        })
1633                        .await?;
1634                        Err(failure(
1635                            WorkflowFailureCode::Suspended,
1636                            "workflow tool requires human approval",
1637                            true,
1638                        ))
1639                    }
1640                    ToolOutcome::Running(handle) => {
1641                        let tool_call_id = handle.tool_call_id.clone();
1642                        (handle.kill)();
1643                        self.persist_suspension(WorkflowSuspensionContext::ToolRunning {
1644                            step_id: instance_id.to_string(),
1645                            tool: tool.clone(),
1646                            tool_call_id,
1647                            killed: true,
1648                        })
1649                        .await?;
1650                        Err(failure(
1651                            WorkflowFailureCode::Suspended,
1652                            "workflow tool is running without a durable workflow resume handle",
1653                            true,
1654                        ))
1655                    }
1656                }
1657            }
1658            WorkflowStepKind::Agent {
1659                agent,
1660                model,
1661                effort,
1662                capabilities,
1663                structured_output_attempts,
1664                ..
1665            } => {
1666                if contains_secret_handle(&input) {
1667                    return Err(failure(
1668                        WorkflowFailureCode::PermissionDenied,
1669                        "secret capability handles are supported only for tool arguments",
1670                        false,
1671                    ));
1672                }
1673                self.authorize(
1674                    &session_id,
1675                    WorkflowPolicyTarget::Agent(agent.clone()),
1676                    capabilities,
1677                )
1678                .await?;
1679                let spec = self.pinned_agents.get(agent).cloned().ok_or_else(|| {
1680                    failure(
1681                        WorkflowFailureCode::PermissionDenied,
1682                        "named agent was not pinned during preflight",
1683                        false,
1684                    )
1685                })?;
1686                let requested = capabilities.iter().cloned().collect::<BTreeSet<_>>();
1687                if !requested.is_subset(&spec.allowed_capabilities) {
1688                    return Err(failure(
1689                        WorkflowFailureCode::PermissionDenied,
1690                        "named agent capability intersection changed",
1691                        false,
1692                    ));
1693                }
1694                let mut last_error = None;
1695                for _ in 0..*structured_output_attempts {
1696                    self.ensure_agent_usage_budget_available().await?;
1697                    self.reserve_agent().await?;
1698                    self.checkpoint_usage("agent_reserved").await?;
1699                    match self
1700                        .engine
1701                        .agents
1702                        .execute(
1703                            &spec,
1704                            input.clone(),
1705                            model.as_deref(),
1706                            effort.as_deref(),
1707                            &requested,
1708                            &session_id,
1709                        )
1710                        .await
1711                    {
1712                        Ok(result) => {
1713                            let exceeded =
1714                                self.record_usage(result.tokens, result.cost_micros).await;
1715                            self.checkpoint_usage("agent_usage_recorded").await?;
1716                            if let Some(error) = exceeded {
1717                                return Err(error);
1718                            }
1719                            if let Some(schema) = &step.output_schema {
1720                                if let Err(error) = validate_schema(schema, &result.output) {
1721                                    last_error = Some(error);
1722                                    continue;
1723                                }
1724                            }
1725                            return Ok(result.output);
1726                        }
1727                        Err(_error) => {
1728                            last_error = Some("named agent execution failed".to_string())
1729                        }
1730                    }
1731                }
1732                Err(failure(
1733                    WorkflowFailureCode::InvalidOutput,
1734                    last_error.unwrap_or_else(|| "agent structured output exhausted".to_string()),
1735                    false,
1736                ))
1737            }
1738            WorkflowStepKind::Workflow {
1739                workflow_id,
1740                revision,
1741                ..
1742            } => {
1743                if self.depth + 1 >= self.root_limits.max_nesting_depth {
1744                    return Err(failure(
1745                        WorkflowFailureCode::BudgetExceeded,
1746                        "nested workflow depth exceeded",
1747                        false,
1748                    ));
1749                }
1750                let definition = self
1751                    .bundle
1752                    .get(workflow_id, *revision)
1753                    .cloned()
1754                    .ok_or_else(|| {
1755                        failure(
1756                            WorkflowFailureCode::UnknownReference,
1757                            format!("persisted bundle missing workflow {workflow_id}@{revision}"),
1758                            false,
1759                        )
1760                    })?;
1761                let nested = StartWorkflowRun {
1762                    definition,
1763                    args: input,
1764                    session_id,
1765                    workspace_trusted: self.workspace_trusted,
1766                    allowed_capabilities: self.allowed_capabilities.iter().cloned().collect(),
1767                };
1768                let parent_run_id = self.snapshot.lock().await.run_id.clone();
1769                let result = Box::pin(self.engine.run_internal(
1770                    nested,
1771                    self.bundle.clone(),
1772                    self.pinned_agents.clone(),
1773                    Some(parent_run_id),
1774                    Some(instance_id.to_string()),
1775                    self.depth + 1,
1776                    self.branch_cancellation.clone(),
1777                    self.ledger.clone(),
1778                    self.root_limits.clone(),
1779                    self.semaphore.clone(),
1780                    None,
1781                ))
1782                .await
1783                .map_err(|_error| {
1784                    failure(
1785                        WorkflowFailureCode::ExecutionFailed,
1786                        "nested workflow execution failed",
1787                        false,
1788                    )
1789                })?;
1790                result.output.ok_or_else(|| {
1791                    result.failure.unwrap_or_else(|| {
1792                        failure(
1793                            WorkflowFailureCode::ExecutionFailed,
1794                            "nested workflow returned no output",
1795                            false,
1796                        )
1797                    })
1798                })
1799            }
1800        }
1801    }
1802
1803    async fn authorize(
1804        &self,
1805        session_id: &str,
1806        target: WorkflowPolicyTarget,
1807        capabilities: &[String],
1808    ) -> Result<(), WorkflowFailure> {
1809        let requested = capabilities.iter().cloned().collect::<BTreeSet<_>>();
1810        if !requested.is_subset(&self.allowed_capabilities) {
1811            return Err(failure(
1812                WorkflowFailureCode::PermissionDenied,
1813                "step capability exceeds root policy",
1814                false,
1815            ));
1816        }
1817        match self
1818            .engine
1819            .policy
1820            .authorize(session_id, &target, &requested, self.workspace_trusted)
1821            .await
1822        {
1823            PermissionDecision::Allow => Ok(()),
1824            PermissionDecision::Deny(_reason) => Err(failure(
1825                if self.workspace_trusted {
1826                    WorkflowFailureCode::PermissionDenied
1827                } else {
1828                    WorkflowFailureCode::UntrustedWorkspace
1829                },
1830                "workflow policy denied this step",
1831                false,
1832            )),
1833        }
1834    }
1835
1836    async fn step_transition(
1837        &self,
1838        step_id: &str,
1839        kind: WorkflowRunEventKind,
1840        mutate: impl FnOnce(&mut WorkflowRunSnapshot),
1841    ) -> Result<(), WorkflowFailure> {
1842        let usage = self.ledger.lock().await.clone();
1843        let mut snapshot = self.snapshot.lock().await;
1844        snapshot.usage = usage;
1845        self.engine
1846            .transition(&mut snapshot, Some(step_id.to_string()), kind, mutate)
1847            .await
1848            .map_err(|error| failure(WorkflowFailureCode::Storage, error.to_string(), false))
1849    }
1850
1851    async fn reserve_step(&self) -> Result<(), WorkflowFailure> {
1852        let mut usage = self.ledger.lock().await;
1853        if usage.steps >= self.root_limits.max_steps {
1854            return Err(failure(
1855                WorkflowFailureCode::BudgetExceeded,
1856                "workflow step budget exceeded",
1857                false,
1858            ));
1859        }
1860        usage.steps += 1;
1861        Ok(())
1862    }
1863
1864    async fn reserve_retry(&self) -> Result<(), WorkflowFailure> {
1865        let mut usage = self.ledger.lock().await;
1866        if usage.retries >= self.root_limits.max_retries {
1867            return Err(failure(
1868                WorkflowFailureCode::BudgetExceeded,
1869                "workflow retry budget exceeded",
1870                false,
1871            ));
1872        }
1873        usage.retries += 1;
1874        Ok(())
1875    }
1876
1877    async fn reserve_agent(&self) -> Result<(), WorkflowFailure> {
1878        let mut usage = self.ledger.lock().await;
1879        if usage.agents >= self.root_limits.max_agents {
1880            return Err(failure(
1881                WorkflowFailureCode::BudgetExceeded,
1882                "workflow agent budget exceeded",
1883                false,
1884            ));
1885        }
1886        usage.agents += 1;
1887        Ok(())
1888    }
1889
1890    async fn record_usage(&self, tokens: u64, cost_micros: u64) -> Option<WorkflowFailure> {
1891        let mut usage = self.ledger.lock().await;
1892        let next_tokens = usage.tokens.saturating_add(tokens);
1893        let next_cost = usage.cost_micros.saturating_add(cost_micros);
1894        usage.tokens = next_tokens;
1895        usage.cost_micros = next_cost;
1896        if self
1897            .root_limits
1898            .max_tokens
1899            .is_some_and(|limit| next_tokens > limit)
1900            || self
1901                .root_limits
1902                .max_cost_micros
1903                .is_some_and(|limit| next_cost > limit)
1904        {
1905            return Some(failure(
1906                WorkflowFailureCode::BudgetExceeded,
1907                "workflow token/cost budget exceeded",
1908                false,
1909            ));
1910        }
1911        None
1912    }
1913
1914    async fn ensure_agent_usage_budget_available(&self) -> Result<(), WorkflowFailure> {
1915        let usage = self.ledger.lock().await;
1916        if self
1917            .root_limits
1918            .max_tokens
1919            .is_some_and(|limit| usage.tokens >= limit)
1920            || self
1921                .root_limits
1922                .max_cost_micros
1923                .is_some_and(|limit| usage.cost_micros >= limit)
1924        {
1925            return Err(failure(
1926                WorkflowFailureCode::BudgetExceeded,
1927                "workflow token/cost budget exhausted before agent dispatch",
1928                false,
1929            ));
1930        }
1931        Ok(())
1932    }
1933
1934    async fn checkpoint_usage(&self, name: &str) -> Result<(), WorkflowFailure> {
1935        let usage = self.ledger.lock().await.clone();
1936        let mut snapshot = self.snapshot.lock().await;
1937        self.engine
1938            .transition(
1939                &mut snapshot,
1940                None,
1941                WorkflowRunEventKind::Phase {
1942                    name: name.to_string(),
1943                },
1944                move |snapshot| snapshot.usage = usage,
1945            )
1946            .await
1947            .map_err(|error| failure(WorkflowFailureCode::Storage, error.to_string(), false))
1948    }
1949
1950    async fn persist_suspension(
1951        &self,
1952        context: WorkflowSuspensionContext,
1953    ) -> Result<(), WorkflowFailure> {
1954        let mut snapshot = self.snapshot.lock().await;
1955        self.engine
1956            .transition(
1957                &mut snapshot,
1958                None,
1959                WorkflowRunEventKind::Phase {
1960                    name: "suspension_context_persisted".to_string(),
1961                },
1962                move |snapshot| snapshot.suspension = Some(context),
1963            )
1964            .await
1965            .map_err(|error| failure(WorkflowFailureCode::Storage, error.to_string(), false))
1966    }
1967
1968    async fn cancel_active_parallel_steps(
1969        &self,
1970        nodes: &[WorkflowPlan],
1971    ) -> Result<(), WorkflowFailure> {
1972        let sibling_steps = nodes
1973            .iter()
1974            .flat_map(plan_step_ids)
1975            .collect::<BTreeSet<_>>();
1976        let active = {
1977            let snapshot = self.snapshot.lock().await;
1978            snapshot
1979                .steps
1980                .iter()
1981                .filter(|(id, step)| {
1982                    matches!(
1983                        step.status,
1984                        WorkflowStepStatus::Queued | WorkflowStepStatus::Running
1985                    ) && sibling_steps
1986                        .iter()
1987                        .any(|step_id| instance_is_in_scope(id, step_id, &self.scope))
1988                })
1989                .map(|(id, _)| id.clone())
1990                .collect::<Vec<_>>()
1991        };
1992        for step_id in active {
1993            let state_id = step_id.clone();
1994            self.step_transition(
1995                &step_id,
1996                WorkflowRunEventKind::StepCancelled,
1997                move |snapshot| {
1998                    if let Some(step) = snapshot.steps.get_mut(&state_id) {
1999                        step.status = WorkflowStepStatus::Cancelled;
2000                        step.failure = Some(failure(
2001                            WorkflowFailureCode::Cancelled,
2002                            "parallel sibling cancelled by fail_fast",
2003                            false,
2004                        ));
2005                    }
2006                },
2007            )
2008            .await?;
2009        }
2010        Ok(())
2011    }
2012
2013    async fn resolve_secret_handles(
2014        &self,
2015        value: &Value,
2016        session_id: &str,
2017    ) -> Result<(Value, Vec<String>), WorkflowFailure> {
2018        fn walk<'a>(
2019            context: &'a RunContext,
2020            value: &'a Value,
2021            session_id: &'a str,
2022        ) -> SecretResolutionFuture<'a> {
2023            Box::pin(async move {
2024                match value {
2025                    Value::Object(object) if object.contains_key("$secret") => {
2026                        let handle: bamboo_domain::WorkflowSecretHandle =
2027                            serde_json::from_value(value.clone()).map_err(|_| {
2028                                failure(
2029                                    WorkflowFailureCode::InvalidInput,
2030                                    "malformed secret capability handle",
2031                                    false,
2032                                )
2033                            })?;
2034                        let material = context
2035                            .engine
2036                            .secrets
2037                            .resolve(session_id, &handle.capability)
2038                            .await
2039                            .map_err(|_| {
2040                                failure(
2041                                    WorkflowFailureCode::PermissionDenied,
2042                                    "secret capability resolution denied",
2043                                    false,
2044                                )
2045                            })?;
2046                        let material = material.into_exposed();
2047                        Ok((Value::String(material.clone()), vec![material]))
2048                    }
2049                    Value::Object(object) => {
2050                        let mut resolved = serde_json::Map::new();
2051                        let mut secrets = Vec::new();
2052                        for (key, child) in object {
2053                            let (child, mut child_secrets) =
2054                                walk(context, child, session_id).await?;
2055                            resolved.insert(key.clone(), child);
2056                            secrets.append(&mut child_secrets);
2057                        }
2058                        Ok((Value::Object(resolved), secrets))
2059                    }
2060                    Value::Array(array) => {
2061                        let mut resolved = Vec::with_capacity(array.len());
2062                        let mut secrets = Vec::new();
2063                        for child in array {
2064                            let (child, mut child_secrets) =
2065                                walk(context, child, session_id).await?;
2066                            resolved.push(child);
2067                            secrets.append(&mut child_secrets);
2068                        }
2069                        Ok((Value::Array(resolved), secrets))
2070                    }
2071                    value => Ok((value.clone(), Vec::new())),
2072                }
2073            })
2074        }
2075        walk(self, value, session_id).await
2076    }
2077
2078    fn skip_plan<'a>(&'a self, plan: &'a WorkflowPlan, reason: &'a str) -> NodeFuture<'a> {
2079        Box::pin(async move {
2080            match plan {
2081                WorkflowPlan::Step { step } => {
2082                    let instance_id = if self.scope == "root" {
2083                        step.clone()
2084                    } else {
2085                        format!("{step}@{}", self.scope)
2086                    };
2087                    let reason_owned = reason.to_string();
2088                    let state_id = instance_id.clone();
2089                    self.step_transition(
2090                        &instance_id,
2091                        WorkflowRunEventKind::StepSkipped {
2092                            reason: reason.to_string(),
2093                        },
2094                        move |snapshot| {
2095                            let state = snapshot.steps.entry(state_id.clone()).or_insert(
2096                                WorkflowStepSnapshot {
2097                                    id: state_id,
2098                                    status: WorkflowStepStatus::Skipped,
2099                                    input_hash: String::new(),
2100                                    output: None,
2101                                    failure: Some(failure(
2102                                        WorkflowFailureCode::DependencySkipped,
2103                                        reason_owned.clone(),
2104                                        false,
2105                                    )),
2106                                    attempts: 0,
2107                                },
2108                            );
2109                            state.status = WorkflowStepStatus::Skipped;
2110                            state.failure = Some(failure(
2111                                WorkflowFailureCode::DependencySkipped,
2112                                reason_owned,
2113                                false,
2114                            ));
2115                        },
2116                    )
2117                    .await?;
2118                }
2119                WorkflowPlan::Sequence { nodes } | WorkflowPlan::Parallel { nodes } => {
2120                    for node in nodes {
2121                        self.skip_plan(node, reason).await?;
2122                    }
2123                }
2124                WorkflowPlan::Map { body, .. } | WorkflowPlan::Retry { node: body, .. } => {
2125                    self.skip_plan(body, reason).await?;
2126                }
2127            }
2128            Ok(Value::Null)
2129        })
2130    }
2131
2132    async fn resolve_template(&self, value: &Value) -> Result<Value, WorkflowFailure> {
2133        Box::pin(self.resolve_template_inner(value)).await
2134    }
2135
2136    fn resolve_template_inner<'a>(
2137        &'a self,
2138        value: &'a Value,
2139    ) -> Pin<Box<dyn Future<Output = Result<Value, WorkflowFailure>> + Send + 'a>> {
2140        Box::pin(async move {
2141            match value {
2142                Value::Object(object) if object.get("from").is_some() => {
2143                    let reference: ValueRef =
2144                        serde_json::from_value(value.clone()).map_err(|error| {
2145                            failure(
2146                                WorkflowFailureCode::InvalidInput,
2147                                format!("malformed value reference: {error}"),
2148                                false,
2149                            )
2150                        })?;
2151                    self.resolve_ref(&reference).await
2152                }
2153                Value::Object(object) => {
2154                    let mut resolved = serde_json::Map::new();
2155                    for (key, child) in object {
2156                        resolved.insert(key.clone(), self.resolve_template_inner(child).await?);
2157                    }
2158                    Ok(Value::Object(resolved))
2159                }
2160                Value::Array(array) => {
2161                    let mut resolved = Vec::with_capacity(array.len());
2162                    for child in array {
2163                        resolved.push(self.resolve_template_inner(child).await?);
2164                    }
2165                    Ok(Value::Array(resolved))
2166                }
2167                value => Ok(value.clone()),
2168            }
2169        })
2170    }
2171
2172    async fn resolve_ref(&self, reference: &ValueRef) -> Result<Value, WorkflowFailure> {
2173        let (root, pointer) = match reference {
2174            ValueRef::Args { pointer } => (
2175                self.snapshot.lock().await.validated_args.clone(),
2176                pointer.as_str(),
2177            ),
2178            ValueRef::Step { step, pointer } => {
2179                let snapshot = self.snapshot.lock().await;
2180                let exact = format!("{step}@{}", self.scope);
2181                let output = snapshot
2182                    .steps
2183                    .get(&exact)
2184                    .or_else(|| snapshot.steps.get(step))
2185                    .and_then(|state| state.output.clone())
2186                    .ok_or_else(|| {
2187                        failure(
2188                            WorkflowFailureCode::UnknownReference,
2189                            format!("step output '{step}' unavailable in execution scope"),
2190                            false,
2191                        )
2192                    })?;
2193                (output, pointer.as_str())
2194            }
2195            ValueRef::Item { name, pointer } => (
2196                self.items.get(name).cloned().ok_or_else(|| {
2197                    failure(
2198                        WorkflowFailureCode::UnknownReference,
2199                        format!("map item '{name}' unavailable"),
2200                        false,
2201                    )
2202                })?,
2203                pointer.as_str(),
2204            ),
2205            ValueRef::Literal { value } => return Ok(value.clone()),
2206        };
2207        if pointer.is_empty() {
2208            Ok(root)
2209        } else {
2210            root.pointer(pointer).cloned().ok_or_else(|| {
2211                failure(
2212                    WorkflowFailureCode::UnknownReference,
2213                    format!("JSON pointer '{pointer}' not found"),
2214                    false,
2215                )
2216            })
2217        }
2218    }
2219
2220    fn check_cancelled(&self) -> Result<(), WorkflowFailure> {
2221        if self.cancellation.is_cancelled() || self.branch_cancellation.is_cancelled() {
2222            Err(failure(
2223                WorkflowFailureCode::Cancelled,
2224                "workflow cancelled",
2225                false,
2226            ))
2227        } else {
2228            Ok(())
2229        }
2230    }
2231}
2232
2233fn event(
2234    snapshot: &WorkflowRunSnapshot,
2235    step_id: Option<String>,
2236    kind: WorkflowRunEventKind,
2237) -> WorkflowRunEvent {
2238    WorkflowRunEvent {
2239        run_id: snapshot.run_id.clone(),
2240        sequence: snapshot.last_sequence,
2241        at: Utc::now(),
2242        step_id,
2243        kind,
2244    }
2245}
2246fn failure(
2247    code: WorkflowFailureCode,
2248    message: impl Into<String>,
2249    retryable: bool,
2250) -> WorkflowFailure {
2251    WorkflowFailure {
2252        code,
2253        message: message.into(),
2254        retryable,
2255    }
2256}
2257fn storage(error: std::io::Error) -> WorkflowRunError {
2258    WorkflowRunError::Storage(error.to_string())
2259}
2260fn exceeds_optional(requested: Option<u64>, ceiling: Option<u64>) -> bool {
2261    match (requested, ceiling) {
2262        (Some(requested), Some(ceiling)) => requested > ceiling,
2263        _ => false,
2264    }
2265}
2266
2267fn enforce_budget_within(
2268    requested: &WorkflowBudgets,
2269    ceiling: &WorkflowBudgets,
2270) -> Result<(), &'static str> {
2271    if requested.max_concurrency > ceiling.max_concurrency {
2272        Err("max_concurrency")
2273    } else if requested.max_agents > ceiling.max_agents {
2274        Err("max_agents")
2275    } else if requested.max_steps > ceiling.max_steps {
2276        Err("max_steps")
2277    } else if requested.max_retries > ceiling.max_retries {
2278        Err("max_retries")
2279    } else if requested.max_nesting_depth > ceiling.max_nesting_depth {
2280        Err("max_nesting_depth")
2281    } else if requested.wall_time_ms > ceiling.wall_time_ms {
2282        Err("wall_time_ms")
2283    } else if exceeds_optional(requested.max_tokens, ceiling.max_tokens) {
2284        Err("max_tokens")
2285    } else if exceeds_optional(requested.max_cost_micros, ceiling.max_cost_micros) {
2286        Err("max_cost_micros")
2287    } else {
2288        Ok(())
2289    }
2290}
2291
2292fn definition_bundle_hash(bundle: &WorkflowDefinitionBundle) -> Result<String, WorkflowRunError> {
2293    let bytes = serde_json::to_vec(bundle)
2294        .map_err(|_| WorkflowRunError::Preflight("workflow bundle hashing failed".to_string()))?;
2295    Ok(hex::encode(Sha256::digest(bytes)))
2296}
2297fn parse_tool_result(result: ToolResult) -> Result<Value, WorkflowFailure> {
2298    if !result.success {
2299        return Err(failure(
2300            WorkflowFailureCode::ExecutionFailed,
2301            "workflow tool reported failure",
2302            true,
2303        ));
2304    }
2305    Ok(serde_json::from_str(&result.result).unwrap_or(Value::String(result.result)))
2306}
2307fn reject_secret_material(value: &Value) -> Result<(), String> {
2308    reject_secret_material_inner(value, false)
2309}
2310
2311fn reject_secret_material_in_definition(value: &Value) -> Result<(), String> {
2312    reject_secret_material_inner(value, true)
2313}
2314
2315fn reject_secret_material_inner(value: &Value, allow_bindings: bool) -> Result<(), String> {
2316    fn walk(value: &Value, key: Option<&str>, allow_bindings: bool) -> Result<(), String> {
2317        if value.as_object().is_some_and(|object| {
2318            object.len() == 1
2319                && object
2320                    .get("$secret")
2321                    .and_then(Value::as_str)
2322                    .is_some_and(|handle| !handle.trim().is_empty())
2323        }) {
2324            return Ok(());
2325        }
2326        let safe_binding = allow_bindings
2327            && serde_json::from_value::<ValueRef>(value.clone())
2328                .is_ok_and(|reference| !matches!(reference, ValueRef::Literal { .. }));
2329        if key.is_some_and(|key| {
2330            let normalized = key
2331                .chars()
2332                .filter(|character| character.is_ascii_alphanumeric())
2333                .flat_map(char::to_lowercase)
2334                .collect::<String>();
2335            matches!(
2336                normalized.as_str(),
2337                "secret"
2338                    | "token"
2339                    | "password"
2340                    | "credential"
2341                    | "credentials"
2342                    | "apikey"
2343                    | "accesskey"
2344                    | "accesstoken"
2345                    | "secretkey"
2346                    | "privatekey"
2347            )
2348        }) && !safe_binding
2349        {
2350            return Err("secret-bearing fields are not accepted by workflow runs".to_string());
2351        }
2352        if value.as_str().is_some_and(|value| {
2353            let trimmed = value.trim();
2354            trimmed.starts_with("capability://")
2355                || trimmed.starts_with("Bearer ")
2356                || trimmed.starts_with("sk-")
2357                || trimmed.starts_with("ghp_")
2358                || trimmed.starts_with("github_pat_")
2359        }) {
2360            // No production capability resolver is part of #578. Treating an
2361            // arbitrary caller string as an opaque handle would be an injection
2362            // channel, so handles and common raw credential forms fail closed.
2363            return Err("opaque credential handles are not enabled for workflows".to_string());
2364        }
2365        match value {
2366            Value::Object(object) => {
2367                for (key, value) in object {
2368                    if key == "properties" {
2369                        let properties = value.as_object().ok_or_else(|| {
2370                            "workflow schema properties must be an object".to_string()
2371                        })?;
2372                        for schema in properties.values() {
2373                            walk(schema, None, allow_bindings)?;
2374                        }
2375                    } else {
2376                        walk(value, Some(key), allow_bindings)?;
2377                    }
2378                }
2379            }
2380            Value::Array(array) => {
2381                for value in array {
2382                    walk(value, None, allow_bindings)?;
2383                }
2384            }
2385            _ => {}
2386        }
2387        Ok(())
2388    }
2389    walk(value, None, allow_bindings)
2390}
2391
2392fn contains_secret_handle(value: &Value) -> bool {
2393    match value {
2394        Value::Object(object) => {
2395            object.contains_key("$secret") || object.values().any(contains_secret_handle)
2396        }
2397        Value::Array(array) => array.iter().any(contains_secret_handle),
2398        _ => false,
2399    }
2400}
2401
2402fn contains_any_secret_material(value: &Value, secrets: &[String]) -> bool {
2403    let matches = |candidate: &str| {
2404        secrets
2405            .iter()
2406            .any(|secret| !secret.is_empty() && candidate.contains(secret))
2407    };
2408    fn walk(value: &Value, matches: &impl Fn(&str) -> bool) -> bool {
2409        match value {
2410            Value::String(value) => matches(value),
2411            Value::Object(object) => object
2412                .iter()
2413                .any(|(key, value)| matches(key) || walk(value, matches)),
2414            Value::Array(array) => array.iter().any(|value| walk(value, matches)),
2415            _ => false,
2416        }
2417    }
2418    walk(value, &matches)
2419}
2420
2421fn plan_leaf_count(plan: &WorkflowPlan) -> usize {
2422    match plan {
2423        WorkflowPlan::Step { .. } => 1,
2424        WorkflowPlan::Sequence { nodes } | WorkflowPlan::Parallel { nodes } => {
2425            nodes.iter().fold(0usize, |total, node| {
2426                total.saturating_add(plan_leaf_count(node))
2427            })
2428        }
2429        WorkflowPlan::Map { body, .. } | WorkflowPlan::Retry { node: body, .. } => {
2430            plan_leaf_count(body)
2431        }
2432    }
2433}
2434
2435fn plan_step_ids(plan: &WorkflowPlan) -> Vec<String> {
2436    match plan {
2437        WorkflowPlan::Step { step } => vec![step.clone()],
2438        WorkflowPlan::Sequence { nodes } | WorkflowPlan::Parallel { nodes } => {
2439            nodes.iter().flat_map(plan_step_ids).collect()
2440        }
2441        WorkflowPlan::Map { body, .. } | WorkflowPlan::Retry { node: body, .. } => {
2442            plan_step_ids(body)
2443        }
2444    }
2445}
2446
2447fn instance_is_in_scope(instance_id: &str, step_id: &str, scope: &str) -> bool {
2448    if scope == "root" {
2449        instance_id == step_id
2450            || instance_id
2451                .strip_prefix(&format!("{step_id}@root"))
2452                .is_some_and(|suffix| suffix.starts_with('['))
2453    } else {
2454        let exact = format!("{step_id}@{scope}");
2455        instance_id == exact
2456            || instance_id
2457                .strip_prefix(&exact)
2458                .is_some_and(|suffix| suffix.starts_with('['))
2459    }
2460}
2461
2462fn plan_frontier(plan: &WorkflowPlan) -> Vec<String> {
2463    match plan {
2464        WorkflowPlan::Step { step } => vec![step.clone()],
2465        WorkflowPlan::Sequence { nodes } => nodes.first().map_or_else(Vec::new, plan_frontier),
2466        WorkflowPlan::Parallel { nodes } => nodes.iter().flat_map(plan_frontier).collect(),
2467        WorkflowPlan::Map { body, .. } | WorkflowPlan::Retry { node: body, .. } => {
2468            plan_frontier(body)
2469        }
2470    }
2471}
2472
2473fn validate_nested_input_contract(
2474    template: &Value,
2475    target_schema: &Value,
2476    compiled: &CompiledWorkflow,
2477) -> Result<(), String> {
2478    fn contains_ref(value: &Value) -> bool {
2479        match value {
2480            Value::Object(object) => {
2481                object.contains_key("from") || object.values().any(contains_ref)
2482            }
2483            Value::Array(array) => array.iter().any(contains_ref),
2484            _ => false,
2485        }
2486    }
2487    if !contains_ref(template) {
2488        return validate_schema(target_schema, template)
2489            .map_err(|error| format!("nested workflow input is incompatible: {error}"));
2490    }
2491    let reference: ValueRef = serde_json::from_value(template.clone()).map_err(|_| {
2492        "nested dynamic input schema cannot be proven compatible in phase 1".to_string()
2493    })?;
2494    let source_schema = match reference {
2495        ValueRef::Args { pointer } => {
2496            schema_at_pointer(&compiled.definition.input_schema, &pointer)
2497        }
2498        ValueRef::Step { step, pointer } => compiled
2499            .steps
2500            .get(&step)
2501            .and_then(|step| step.output_schema.as_ref())
2502            .and_then(|schema| schema_at_pointer(schema, &pointer)),
2503        ValueRef::Literal { value } => {
2504            return validate_schema(target_schema, &value)
2505                .map_err(|error| format!("nested workflow input is incompatible: {error}"));
2506        }
2507        ValueRef::Item { .. } => None,
2508    }
2509    .ok_or_else(|| {
2510        "nested dynamic input source schema is missing or pointer is invalid".to_string()
2511    })?;
2512    if schema_compatible(source_schema, target_schema) {
2513        Ok(())
2514    } else {
2515        Err("nested workflow input schema is not compatible with its pinned target".to_string())
2516    }
2517}
2518
2519fn schema_at_pointer<'a>(schema: &'a Value, pointer: &str) -> Option<&'a Value> {
2520    if pointer.is_empty() {
2521        return Some(schema);
2522    }
2523    let mut current = schema;
2524    for token in pointer.strip_prefix('/')?.split('/') {
2525        let token = token.replace("~1", "/").replace("~0", "~");
2526        current = if token.parse::<usize>().is_ok() {
2527            current.get("items")?
2528        } else {
2529            current.get("properties")?.get(&token)?
2530        };
2531    }
2532    Some(current)
2533}
2534
2535fn schema_compatible(source: &Value, target: &Value) -> bool {
2536    if source == target {
2537        return true;
2538    }
2539    let source_type = source.get("type").and_then(Value::as_str);
2540    let target_type = target.get("type").and_then(Value::as_str);
2541    source_type.is_some() && source_type == target_type && target_type != Some("object")
2542}
2543
2544fn effective_limits(requested: &WorkflowBudgets, ceilings: &WorkflowBudgets) -> WorkflowBudgets {
2545    WorkflowBudgets {
2546        max_concurrency: requested.max_concurrency.min(ceilings.max_concurrency),
2547        max_agents: requested.max_agents.min(ceilings.max_agents),
2548        max_steps: requested.max_steps.min(ceilings.max_steps),
2549        max_retries: requested.max_retries.min(ceilings.max_retries),
2550        max_nesting_depth: requested.max_nesting_depth.min(ceilings.max_nesting_depth),
2551        wall_time_ms: requested.wall_time_ms.min(ceilings.wall_time_ms),
2552        max_tokens: match (requested.max_tokens, ceilings.max_tokens) {
2553            (Some(requested), Some(ceiling)) => Some(requested.min(ceiling)),
2554            (Some(requested), None) => Some(requested),
2555            (None, ceiling) => ceiling,
2556        },
2557        max_cost_micros: match (requested.max_cost_micros, ceilings.max_cost_micros) {
2558            (Some(requested), Some(ceiling)) => Some(requested.min(ceiling)),
2559            (Some(requested), None) => Some(requested),
2560            (None, ceiling) => ceiling,
2561        },
2562    }
2563}