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