Skip to main content

bamboo_engine/workflow_run/
executor.rs

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