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 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 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
97pub 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 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 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 if let Some(durable) = self.repository.load(&run_id).await.map_err(storage)? {
550 *final_snapshot = durable;
551 }
552 let step_failure = failure(
553 WorkflowFailureCode::BudgetExceeded,
554 "workflow wall-time budget exceeded",
555 false,
556 );
557 self.fail_timeout_frontier(
558 &mut final_snapshot,
559 &compiled.definition.plan,
560 step_failure.clone(),
561 )
562 .await?;
563 self.fail_active_steps(&mut final_snapshot, step_failure.clone())
564 .await?;
565 self.finish_failed(&mut final_snapshot, step_failure)
566 .await?;
567 }
568 }
569 Ok(final_snapshot.clone())
570 }
571
572 pub async fn progress(
573 &self,
574 run_id: &str,
575 since: u64,
576 ) -> Result<WorkflowProgress, WorkflowRunError> {
577 let snapshot = self
578 .repository
579 .load(run_id)
580 .await
581 .map_err(storage)?
582 .ok_or(WorkflowRunError::NotFound)?;
583 let events = self
584 .repository
585 .events_since(run_id, since)
586 .await
587 .map_err(storage)?;
588 Ok(WorkflowProgress { snapshot, events })
589 }
590
591 pub async fn list_run_ids(&self) -> Result<Vec<String>, WorkflowRunError> {
592 self.repository.list_run_ids().await.map_err(storage)
593 }
594
595 pub async fn cancel(&self, run_id: &str) -> Result<WorkflowRunSnapshot, WorkflowRunError> {
596 let mut snapshot = self
597 .repository
598 .load(run_id)
599 .await
600 .map_err(storage)?
601 .ok_or(WorkflowRunError::NotFound)?;
602 if snapshot.status == WorkflowRunStatus::Cancelled {
603 return Ok(snapshot);
604 }
605 if snapshot.status.is_terminal() {
606 return Err(WorkflowRunError::Terminal);
607 }
608 if let Some(active) = self.active.get(run_id).map(|active| active.clone()) {
609 active.cancellation.cancel();
610 let mut shared = active.snapshot.lock().await;
611 if !shared.status.is_terminal() {
612 self.finish_cancelled(&mut shared).await?;
613 }
614 return Ok(shared.clone());
615 }
616 self.finish_cancelled(&mut snapshot).await?;
617 Ok(snapshot)
618 }
619
620 pub async fn recover(&self) -> Result<Vec<WorkflowRunSnapshot>, WorkflowRunError> {
621 let mut recovered = Vec::new();
622 for run_id in self.repository.list_run_ids().await.map_err(storage)? {
623 let Some(mut snapshot) = self.repository.load(&run_id).await.map_err(storage)? else {
624 continue;
625 };
626 if matches!(
627 snapshot.status,
628 WorkflowRunStatus::Queued | WorkflowRunStatus::Running
629 ) {
630 let reason = "process restarted; explicit safe restart is required".to_string();
631 let active_steps = snapshot
632 .steps
633 .iter()
634 .filter(|(_, step)| {
635 matches!(
636 step.status,
637 WorkflowStepStatus::Queued | WorkflowStepStatus::Running
638 )
639 })
640 .map(|(id, _)| id.clone())
641 .collect::<Vec<_>>();
642 for step_id in active_steps {
643 let state_id = step_id.clone();
644 let step_reason = reason.clone();
645 self.transition(
646 &mut snapshot,
647 Some(step_id),
648 WorkflowRunEventKind::StepSuspended {
649 reason: reason.clone(),
650 },
651 move |snapshot| {
652 if let Some(step) = snapshot.steps.get_mut(&state_id) {
653 step.status = WorkflowStepStatus::Suspended;
654 step.failure = Some(failure(
655 WorkflowFailureCode::RecoverySuspended,
656 step_reason,
657 true,
658 ));
659 }
660 },
661 )
662 .await?;
663 }
664 self.transition(
665 &mut snapshot,
666 None,
667 WorkflowRunEventKind::RunSuspended {
668 reason: reason.clone(),
669 },
670 move |snapshot| {
671 snapshot.status = WorkflowRunStatus::Suspended;
672 snapshot.suspension = Some(WorkflowSuspensionContext::Recovery {
673 reason: reason.clone(),
674 });
675 },
676 )
677 .await?;
678 recovered.push(snapshot);
679 }
680 }
681 Ok(recovered)
682 }
683
684 pub fn subscribe(&self, run_id: &str) -> Option<broadcast::Receiver<WorkflowRunEvent>> {
685 self.events.get(run_id).map(|sender| sender.subscribe())
686 }
687
688 #[cfg(test)]
689 pub(crate) fn runtime_resource_counts(&self) -> (usize, usize) {
690 (self.active.len(), self.events.len())
691 }
692
693 async fn pin_and_validate_bundle(
694 &self,
695 root: &WorkflowRunDefinition,
696 ) -> Result<WorkflowDefinitionBundle, WorkflowRunError> {
697 let bundle =
698 self.definitions.pin_bundle(root).await.map_err(|_| {
699 WorkflowRunError::Preflight("workflow bundle pin failed".to_string())
700 })?;
701 self.validate_bundle(root, &bundle)?;
702 Ok(bundle)
703 }
704
705 fn validate_bundle(
706 &self,
707 root: &WorkflowRunDefinition,
708 bundle: &WorkflowDefinitionBundle,
709 ) -> Result<(), WorkflowRunError> {
710 if bundle.root_id != root.id
711 || bundle.root_revision != root.revision
712 || bundle.root() != Some(root)
713 {
714 return Err(WorkflowRunError::Preflight(
715 "pinned bundle root identity/content mismatch".to_string(),
716 ));
717 }
718 let serialized = serde_json::to_value(bundle).map_err(|_| {
719 WorkflowRunError::Preflight("workflow bundle is not serializable".into())
720 })?;
721 reject_secret_material_in_definition(&serialized).map_err(WorkflowRunError::Preflight)?;
722 let mut stack = vec![(root.id.clone(), root.revision, Vec::<String>::new())];
723 let mut visited = BTreeSet::new();
724 while let Some((id, revision, path)) = stack.pop() {
725 let key = WorkflowDefinitionBundle::key(&id, revision);
726 if path.contains(&key) {
727 return Err(WorkflowRunError::Preflight(format!(
728 "nested workflow cycle includes {key}"
729 )));
730 }
731 if !visited.insert(key.clone()) {
732 continue;
733 }
734 let definition = bundle.get(&id, revision).ok_or_else(|| {
735 WorkflowRunError::Preflight(format!("pinned bundle is missing {key}"))
736 })?;
737 if definition.id != id || definition.revision != revision {
738 return Err(WorkflowRunError::Preflight(
739 "pinned bundle definition identity mismatch".to_string(),
740 ));
741 }
742 let compiled = CompiledWorkflow::compile(definition.clone())?;
743 let mut nested_path = path;
744 nested_path.push(key);
745 for step in compiled.steps.values() {
746 if let WorkflowStepKind::Workflow {
747 workflow_id,
748 revision,
749 args,
750 } = &step.kind
751 {
752 let nested = bundle.get(workflow_id, *revision).ok_or_else(|| {
753 WorkflowRunError::Preflight(format!(
754 "pinned bundle is missing {workflow_id}@{revision}"
755 ))
756 })?;
757 validate_nested_input_contract(args, &nested.input_schema, &compiled)
758 .map_err(WorkflowRunError::Preflight)?;
759 stack.push((workflow_id.clone(), *revision, nested_path.clone()));
760 }
761 }
762 }
763 Ok(())
764 }
765
766 async fn preflight_bundle(
767 &self,
768 bundle: &WorkflowDefinitionBundle,
769 session_id: &str,
770 allowed: &BTreeSet<String>,
771 trusted: bool,
772 root_limits: &WorkflowBudgets,
773 ) -> Result<HashMap<String, NamedAgentSpec>, WorkflowRunError> {
774 let mut pinned_agents = HashMap::<String, NamedAgentSpec>::new();
775 let mut stack = vec![(bundle.root_id.clone(), bundle.root_revision, 0_u32)];
776 let mut visited = BTreeSet::new();
777 while let Some((id, revision, depth)) = stack.pop() {
778 if depth >= root_limits.max_nesting_depth {
779 return Err(WorkflowRunError::Preflight(
780 "nested workflow depth exceeded shared root limit".to_string(),
781 ));
782 }
783 if !visited.insert(WorkflowDefinitionBundle::key(&id, revision)) {
784 continue;
785 }
786 let definition = bundle.get(&id, revision).ok_or_else(|| {
787 WorkflowRunError::Preflight("pinned workflow definition missing".to_string())
788 })?;
789 enforce_budget_within(&definition.budgets, root_limits).map_err(|message| {
790 WorkflowRunError::Preflight(format!(
791 "nested workflow budget expands root: {message}"
792 ))
793 })?;
794 let compiled = CompiledWorkflow::compile(definition.clone())?;
795 for step in compiled.steps.values() {
796 let (target, capabilities) = match &step.kind {
797 WorkflowStepKind::Tool {
798 tool, capabilities, ..
799 } => (WorkflowPolicyTarget::Tool(tool.clone()), capabilities),
800 WorkflowStepKind::Agent {
801 agent,
802 capabilities,
803 ..
804 } => {
805 let spec = if let Some(spec) = pinned_agents.get(agent) {
806 spec.clone()
807 } else {
808 let spec = self
809 .agents
810 .resolve(agent)
811 .await
812 .map_err(|_| {
813 WorkflowRunError::Preflight(
814 "named agent resolution failed".to_string(),
815 )
816 })?
817 .ok_or_else(|| {
818 WorkflowRunError::Preflight(format!(
819 "unknown named agent '{agent}'"
820 ))
821 })?;
822 if spec.name != *agent {
823 return Err(WorkflowRunError::Preflight(
824 "named agent resolver returned mismatched identity".to_string(),
825 ));
826 }
827 pinned_agents.insert(agent.clone(), spec.clone());
828 spec
829 };
830 if !capabilities
831 .iter()
832 .all(|capability| spec.allowed_capabilities.contains(capability))
833 {
834 return Err(WorkflowRunError::Preflight(format!(
835 "agent '{agent}' capability expansion denied"
836 )));
837 }
838 (WorkflowPolicyTarget::Agent(agent.clone()), capabilities)
839 }
840 WorkflowStepKind::Workflow {
841 workflow_id,
842 revision,
843 args,
844 } => {
845 let nested = bundle.get(workflow_id, *revision).ok_or_else(|| {
846 WorkflowRunError::Preflight(format!(
847 "missing pinned workflow {workflow_id}@{revision}"
848 ))
849 })?;
850 validate_nested_input_contract(args, &nested.input_schema, &compiled)
851 .map_err(WorkflowRunError::Preflight)?;
852 stack.push((workflow_id.clone(), *revision, depth + 1));
853 (
854 WorkflowPolicyTarget::Workflow {
855 id: workflow_id.clone(),
856 revision: *revision,
857 },
858 &Vec::new(),
859 )
860 }
861 };
862 let requested = capabilities.iter().cloned().collect::<BTreeSet<_>>();
863 if !requested.is_subset(allowed) {
864 return Err(WorkflowRunError::Preflight(format!(
865 "step '{}' exceeds root capabilities",
866 step.id
867 )));
868 }
869 if let PermissionDecision::Deny(_reason) = self
870 .policy
871 .authorize(session_id, &target, &requested, trusted)
872 .await
873 {
874 return Err(WorkflowRunError::Preflight(
875 "workflow policy denied this step".to_string(),
876 ));
877 }
878 }
879 }
880 Ok(pinned_agents)
881 }
882
883 fn enforce_ceilings(&self, budget: &WorkflowBudgets) -> Result<(), WorkflowRunError> {
884 if budget.max_concurrency > self.ceilings.max_concurrency
885 || budget.max_agents > self.ceilings.max_agents
886 || budget.max_steps > self.ceilings.max_steps
887 || budget.max_retries > self.ceilings.max_retries
888 || budget.max_nesting_depth > self.ceilings.max_nesting_depth
889 || budget.wall_time_ms > self.ceilings.wall_time_ms
890 || exceeds_optional(budget.max_tokens, self.ceilings.max_tokens)
891 || exceeds_optional(budget.max_cost_micros, self.ceilings.max_cost_micros)
892 {
893 return Err(WorkflowRunError::Preflight(
894 "definition exceeds server workflow budget ceilings".to_string(),
895 ));
896 }
897 Ok(())
898 }
899
900 async fn transition(
901 &self,
902 snapshot: &mut WorkflowRunSnapshot,
903 step_id: Option<String>,
904 kind: WorkflowRunEventKind,
905 mutate: impl FnOnce(&mut WorkflowRunSnapshot),
906 ) -> Result<(), WorkflowRunError> {
907 let mut candidate = self
912 .repository
913 .load(&snapshot.run_id)
914 .await
915 .map_err(storage)?
916 .ok_or_else(|| WorkflowRunError::Storage("workflow snapshot missing".to_string()))?;
917 mutate(&mut candidate);
918 candidate.last_sequence += 1;
919 candidate.updated_at = Utc::now();
920 let event = event(&candidate, step_id, kind);
921 let repository = self.repository.clone();
922 let durable_candidate = candidate.clone();
923 let durable_event = event.clone();
924 let commit =
925 tokio::spawn(
926 async move { repository.commit(&durable_candidate, &durable_event).await },
927 );
928 commit
929 .await
930 .map_err(|error| {
931 WorkflowRunError::Storage(format!("workflow commit task failed: {error}"))
932 })?
933 .map_err(storage)?;
934 *snapshot = candidate;
935 self.publish(&event);
936 Ok(())
937 }
938
939 fn publish(&self, event: &WorkflowRunEvent) {
940 if let Some(sender) = self.events.get(&event.run_id) {
941 let _ = sender.send(event.clone());
942 }
943 }
944
945 async fn finish_succeeded(
946 &self,
947 snapshot: &mut WorkflowRunSnapshot,
948 output: Value,
949 ) -> Result<(), WorkflowRunError> {
950 if snapshot.status.is_terminal() {
951 return Ok(());
952 }
953 let copy = output.clone();
954 self.transition(
955 snapshot,
956 None,
957 WorkflowRunEventKind::RunSucceeded { output },
958 move |snapshot| {
959 snapshot.status = WorkflowRunStatus::Succeeded;
960 snapshot.output = Some(copy);
961 },
962 )
963 .await
964 }
965 async fn finish_failed(
966 &self,
967 snapshot: &mut WorkflowRunSnapshot,
968 error: WorkflowFailure,
969 ) -> Result<(), WorkflowRunError> {
970 if snapshot.status.is_terminal() {
971 return Ok(());
972 }
973 let copy = error.clone();
974 self.transition(
975 snapshot,
976 None,
977 WorkflowRunEventKind::RunFailed { failure: error },
978 move |snapshot| {
979 snapshot.status = WorkflowRunStatus::Failed;
980 snapshot.failure = Some(copy);
981 },
982 )
983 .await
984 }
985 async fn finish_cancelled(
986 &self,
987 snapshot: &mut WorkflowRunSnapshot,
988 ) -> Result<(), WorkflowRunError> {
989 if snapshot.status == WorkflowRunStatus::Cancelled {
990 return Ok(());
991 }
992 let active_steps = snapshot
993 .steps
994 .iter()
995 .filter(|(_, step)| {
996 matches!(
997 step.status,
998 WorkflowStepStatus::Queued | WorkflowStepStatus::Running
999 )
1000 })
1001 .map(|(id, _)| id.clone())
1002 .collect::<Vec<_>>();
1003 for step_id in active_steps {
1004 let state_id = step_id.clone();
1005 self.transition(
1006 snapshot,
1007 Some(step_id),
1008 WorkflowRunEventKind::StepCancelled,
1009 move |snapshot| {
1010 if let Some(step) = snapshot.steps.get_mut(&state_id) {
1011 step.status = WorkflowStepStatus::Cancelled;
1012 step.failure = Some(failure(
1013 WorkflowFailureCode::Cancelled,
1014 "workflow cancelled",
1015 false,
1016 ));
1017 }
1018 },
1019 )
1020 .await?;
1021 }
1022 self.transition(
1023 snapshot,
1024 None,
1025 WorkflowRunEventKind::RunCancelled,
1026 |snapshot| {
1027 snapshot.status = WorkflowRunStatus::Cancelled;
1028 snapshot.failure = Some(failure(
1029 WorkflowFailureCode::Cancelled,
1030 "workflow cancelled",
1031 false,
1032 ));
1033 },
1034 )
1035 .await
1036 }
1037
1038 async fn finish_suspended(
1039 &self,
1040 snapshot: &mut WorkflowRunSnapshot,
1041 reason: String,
1042 ) -> Result<(), WorkflowRunError> {
1043 if snapshot.status.is_terminal() {
1044 return Ok(());
1045 }
1046 self.transition(
1047 snapshot,
1048 None,
1049 WorkflowRunEventKind::RunSuspended { reason },
1050 |snapshot| {
1051 snapshot.status = WorkflowRunStatus::Suspended;
1052 },
1053 )
1054 .await
1055 }
1056
1057 async fn fail_active_steps(
1058 &self,
1059 snapshot: &mut WorkflowRunSnapshot,
1060 error: WorkflowFailure,
1061 ) -> Result<(), WorkflowRunError> {
1062 let active_steps = snapshot
1063 .steps
1064 .iter()
1065 .filter(|(_, step)| {
1066 matches!(
1067 step.status,
1068 WorkflowStepStatus::Queued | WorkflowStepStatus::Running
1069 )
1070 })
1071 .map(|(id, _)| id.clone())
1072 .collect::<Vec<_>>();
1073 for step_id in active_steps {
1074 let state_id = step_id.clone();
1075 let copy = error.clone();
1076 self.transition(
1077 snapshot,
1078 Some(step_id),
1079 WorkflowRunEventKind::StepFailed {
1080 failure: error.clone(),
1081 },
1082 move |snapshot| {
1083 if let Some(step) = snapshot.steps.get_mut(&state_id) {
1084 step.status = WorkflowStepStatus::Failed;
1085 step.failure = Some(copy);
1086 }
1087 },
1088 )
1089 .await?;
1090 }
1091 Ok(())
1092 }
1093
1094 async fn fail_timeout_frontier(
1095 &self,
1096 snapshot: &mut WorkflowRunSnapshot,
1097 plan: &WorkflowPlan,
1098 error: WorkflowFailure,
1099 ) -> Result<(), WorkflowRunError> {
1100 if snapshot.steps.values().any(|step| {
1101 matches!(
1102 step.status,
1103 WorkflowStepStatus::Queued | WorkflowStepStatus::Running
1104 )
1105 }) {
1106 return Ok(());
1107 }
1108 for step_id in plan_frontier(plan) {
1109 if snapshot.steps.contains_key(&step_id) {
1110 continue;
1111 }
1112 let state_id = step_id.clone();
1113 let state_error = error.clone();
1114 self.transition(
1115 snapshot,
1116 Some(step_id),
1117 WorkflowRunEventKind::StepFailed {
1118 failure: error.clone(),
1119 },
1120 move |snapshot| {
1121 snapshot.steps.insert(
1122 state_id.clone(),
1123 WorkflowStepSnapshot {
1124 id: state_id,
1125 status: WorkflowStepStatus::Failed,
1126 input_hash: String::new(),
1127 output: None,
1128 failure: Some(state_error),
1129 attempts: 0,
1130 },
1131 );
1132 },
1133 )
1134 .await?;
1135 }
1136 Ok(())
1137 }
1138}
1139
1140impl RunContext {
1141 fn execute_node<'a>(&'a self, plan: &'a WorkflowPlan, path: &'a str) -> NodeFuture<'a> {
1142 Box::pin(async move {
1143 self.check_cancelled()?;
1144 match plan {
1145 WorkflowPlan::Step { step } => self.execute_step(step, path).await,
1146 WorkflowPlan::Sequence { nodes } => {
1147 let mut result = Value::Null;
1148 for (index, node) in nodes.iter().enumerate() {
1149 match self.execute_node(node, &format!("{path}.{index}")).await {
1150 Ok(value) => result = value,
1151 Err(error) => {
1152 if error.code == WorkflowFailureCode::DependencySkipped {
1153 for remaining in &nodes[index + 1..] {
1154 self.skip_plan(
1155 remaining,
1156 "dependency requested skip_dependents",
1157 )
1158 .await?;
1159 }
1160 }
1161 return Err(error);
1162 }
1163 }
1164 }
1165 Ok(result)
1166 }
1167 WorkflowPlan::Parallel { nodes } => {
1168 let parallel_cancellation = self.branch_cancellation.child_token();
1169 let mut futures = FuturesUnordered::new();
1170 for (index, node) in nodes.iter().enumerate() {
1171 let mut child = self.clone();
1172 child.branch_cancellation = parallel_cancellation.clone();
1173 futures.push(async move {
1174 (
1175 index,
1176 child.execute_node(node, &format!("{path}.{index}")).await,
1177 )
1178 });
1179 }
1180 let mut output = vec![Value::Null; nodes.len()];
1181 while let Some((index, result)) = futures.next().await {
1182 match result {
1183 Ok(value) => output[index] = value,
1184 Err(mut error) => {
1185 parallel_cancellation.cancel();
1186 drop(futures);
1192 self.cancel_active_parallel_steps(nodes).await?;
1193 error.message =
1194 format!("parallel branch[{index}] failed: {}", error.message);
1195 return Err(error);
1196 }
1197 }
1198 }
1199 Ok(Value::Array(output))
1200 }
1201 WorkflowPlan::Map { source, item, body } => {
1202 let source = self.resolve_ref(source).await?;
1203 let values = source.as_array().ok_or_else(|| {
1204 failure(
1205 WorkflowFailureCode::InvalidInput,
1206 "map source must be an array",
1207 false,
1208 )
1209 })?;
1210 let used = self.ledger.lock().await.steps as usize;
1211 let remaining = (self.root_limits.max_steps as usize).saturating_sub(used);
1212 let per_item = plan_leaf_count(body).max(1);
1213 if values
1214 .len()
1215 .checked_mul(per_item)
1216 .is_none_or(|required| required > remaining)
1217 {
1218 return Err(failure(
1219 WorkflowFailureCode::BudgetExceeded,
1220 "map cardinality exceeds remaining workflow step budget",
1221 false,
1222 ));
1223 }
1224 let futures = values.iter().cloned().enumerate().map(|(index, value)| {
1225 let mut child = self.clone();
1226 child.items.insert(item.clone(), value);
1227 child.scope = format!("{}[{index}]", self.scope);
1232 async move { child.execute_node(body, &format!("{path}[{index}]")).await }
1233 });
1234 let results = join_all(futures).await;
1235 let mut values = Vec::with_capacity(results.len());
1236 let mut failures = Vec::new();
1237 for (index, result) in results.into_iter().enumerate() {
1238 match result {
1239 Ok(value) => values.push(value),
1240 Err(error) => failures.push((index, error)),
1241 }
1242 }
1243 if failures.is_empty() {
1244 Ok(Value::Array(values))
1245 } else {
1246 let retryable = failures.iter().any(|(_, error)| error.retryable);
1247 let first_code = failures[0].1.code;
1248 let code = if failures
1249 .iter()
1250 .any(|(_, error)| error.code == WorkflowFailureCode::DependencySkipped)
1251 {
1252 WorkflowFailureCode::DependencySkipped
1253 } else if failures.iter().all(|(_, error)| error.code == first_code) {
1254 first_code
1255 } else {
1256 WorkflowFailureCode::ExecutionFailed
1257 };
1258 let diagnostics = failures
1259 .into_iter()
1260 .map(|(index, error)| format!("item[{index}]: {}", error.message))
1261 .collect::<Vec<_>>()
1262 .join("; ");
1263 Err(failure(
1264 code,
1265 format!("map items failed: {diagnostics}"),
1266 retryable,
1267 ))
1268 }
1269 }
1270 WorkflowPlan::Retry {
1271 node,
1272 max_attempts,
1273 delay_ms,
1274 } => {
1275 let limit =
1276 (*max_attempts).min(self.compiled.definition.budgets.max_retries + 1);
1277 let mut last = None;
1278 for attempt in 0..limit {
1279 match self
1280 .execute_node(node, &format!("{path}.retry{attempt}"))
1281 .await
1282 {
1283 Ok(value) => return Ok(value),
1284 Err(error) if error.retryable => {
1285 last = Some(error);
1286 if attempt + 1 < limit {
1287 self.reserve_retry().await?;
1288 self.checkpoint_usage("retry_reserved").await?;
1289 tokio::select! {
1290 _ = self.cancellation.cancelled() => return Err(failure(WorkflowFailureCode::Cancelled, "workflow cancelled", false)),
1291 _ = self.branch_cancellation.cancelled() => return Err(failure(WorkflowFailureCode::Cancelled, "workflow branch cancelled", false)),
1292 _ = tokio::time::sleep(Duration::from_millis(*delay_ms)) => {}
1293 }
1294 } else {
1295 break;
1296 }
1297 }
1298 Err(error) => return Err(error),
1299 }
1300 }
1301 Err(failure(
1302 WorkflowFailureCode::RetryExhausted,
1303 last.map_or_else(|| "retry exhausted".to_string(), |error| error.message),
1304 false,
1305 ))
1306 }
1307 }
1308 })
1309 }
1310
1311 async fn execute_step(&self, step_id: &str, _path: &str) -> Result<Value, WorkflowFailure> {
1312 self.check_cancelled()?;
1313 let step = self.compiled.steps.get(step_id).cloned().ok_or_else(|| {
1314 failure(
1315 WorkflowFailureCode::UnknownReference,
1316 format!("unknown step {step_id}"),
1317 false,
1318 )
1319 })?;
1320 let instance_id = if self.scope == "root" {
1321 step_id.to_string()
1322 } else {
1323 format!("{step_id}@{}", self.scope)
1324 };
1325 let input = match &step.kind {
1326 WorkflowStepKind::Tool { args, .. } | WorkflowStepKind::Workflow { args, .. } => {
1327 self.resolve_template(args).await?
1328 }
1329 WorkflowStepKind::Agent { prompt, .. } => self.resolve_template(prompt).await?,
1330 };
1331 let input_hash = hex::encode(Sha256::digest(
1332 serde_json::to_vec(&input).unwrap_or_default(),
1333 ));
1334 self.reserve_step().await?;
1335 self.checkpoint_usage("step_reserved").await?;
1336 self.step_transition(&instance_id, WorkflowRunEventKind::StepQueued, |snapshot| {
1337 let state = snapshot
1338 .steps
1339 .entry(instance_id.clone())
1340 .or_insert_with(|| WorkflowStepSnapshot {
1341 id: instance_id.clone(),
1342 status: WorkflowStepStatus::Queued,
1343 input_hash: input_hash.clone(),
1344 output: None,
1345 failure: None,
1346 attempts: 0,
1347 });
1348 state.status = WorkflowStepStatus::Queued;
1349 state.input_hash = input_hash;
1350 state.output = None;
1351 state.failure = None;
1352 })
1353 .await?;
1354 let _permit = if matches!(&step.kind, WorkflowStepKind::Workflow { .. }) {
1355 None
1356 } else {
1357 Some(tokio::select! {
1358 _ = self.cancellation.cancelled() => {
1359 let cancelled_id = instance_id.clone();
1360 self.step_transition(&instance_id, WorkflowRunEventKind::StepCancelled, move |snapshot| {
1361 if let Some(state) = snapshot.steps.get_mut(&cancelled_id) { state.status = WorkflowStepStatus::Cancelled; }
1362 }).await?;
1363 return Err(failure(WorkflowFailureCode::Cancelled, "workflow cancelled", false));
1364 }
1365 _ = self.branch_cancellation.cancelled() => {
1366 let cancelled_id = instance_id.clone();
1367 self.step_transition(&instance_id, WorkflowRunEventKind::StepCancelled, move |snapshot| {
1368 if let Some(state) = snapshot.steps.get_mut(&cancelled_id) {
1369 state.status = WorkflowStepStatus::Cancelled;
1370 }
1371 }).await?;
1372 return Err(failure(WorkflowFailureCode::Cancelled, "workflow branch cancelled", false));
1373 }
1374 permit = self.semaphore.acquire() => permit.map_err(|_| failure(WorkflowFailureCode::ExecutionFailed, "workflow semaphore closed", false))?,
1375 })
1376 };
1377 let started_id = instance_id.clone();
1378 self.step_transition(
1379 &instance_id,
1380 WorkflowRunEventKind::StepStarted,
1381 move |snapshot| {
1382 if let Some(state) = snapshot.steps.get_mut(&started_id) {
1383 state.status = WorkflowStepStatus::Running;
1384 state.attempts += 1;
1385 }
1386 },
1387 )
1388 .await?;
1389 let result = self.dispatch(&step, input, &instance_id).await;
1390 let result = match result {
1391 Ok(output) => {
1392 if let Some(schema) = &step.output_schema {
1393 validate_schema(schema, &output)
1394 .map(|()| output)
1395 .map_err(|message| {
1396 failure(WorkflowFailureCode::InvalidOutput, message, false)
1397 })
1398 } else {
1399 Ok(output)
1400 }
1401 }
1402 Err(error) => Err(error),
1403 };
1404 let result = result.and_then(|output| {
1405 reject_secret_material(&output)
1406 .map(|()| output)
1407 .map_err(|message| failure(WorkflowFailureCode::InvalidOutput, message, false))
1408 });
1409 match result {
1410 Ok(output) => {
1411 let copy = output.clone();
1412 let completed_id = instance_id.clone();
1413 self.step_transition(
1414 &instance_id,
1415 WorkflowRunEventKind::StepCompleted {
1416 output: output.clone(),
1417 },
1418 move |snapshot| {
1419 if let Some(state) = snapshot.steps.get_mut(&completed_id) {
1420 state.status = WorkflowStepStatus::Succeeded;
1421 state.output = Some(copy);
1422 }
1423 },
1424 )
1425 .await?;
1426 Ok(output)
1427 }
1428 Err(error) => {
1429 if error.code == WorkflowFailureCode::Cancelled {
1430 let cancelled_id = instance_id.clone();
1431 self.step_transition(
1432 &instance_id,
1433 WorkflowRunEventKind::StepCancelled,
1434 move |snapshot| {
1435 if let Some(state) = snapshot.steps.get_mut(&cancelled_id) {
1436 state.status = WorkflowStepStatus::Cancelled;
1437 state.failure = Some(failure(
1438 WorkflowFailureCode::Cancelled,
1439 "workflow branch cancelled",
1440 false,
1441 ));
1442 }
1443 },
1444 )
1445 .await?;
1446 return Err(error);
1447 }
1448 if error.code == WorkflowFailureCode::Suspended {
1449 let reason = error.message.clone();
1450 let suspended_id = instance_id.clone();
1451 self.step_transition(
1452 &instance_id,
1453 WorkflowRunEventKind::StepSuspended {
1454 reason: reason.clone(),
1455 },
1456 move |snapshot| {
1457 if let Some(state) = snapshot.steps.get_mut(&suspended_id) {
1458 state.status = WorkflowStepStatus::Suspended;
1459 state.failure =
1460 Some(failure(WorkflowFailureCode::Suspended, reason, true));
1461 }
1462 },
1463 )
1464 .await?;
1465 return Err(error);
1466 }
1467 let copy = error.clone();
1468 let failed_id = instance_id.clone();
1469 self.step_transition(
1470 &instance_id,
1471 WorkflowRunEventKind::StepFailed {
1472 failure: error.clone(),
1473 },
1474 move |snapshot| {
1475 if let Some(state) = snapshot.steps.get_mut(&failed_id) {
1476 state.status = WorkflowStepStatus::Failed;
1477 state.failure = Some(copy);
1478 }
1479 },
1480 )
1481 .await?;
1482 match step.failure {
1483 FailurePolicy::ContinueWithError => Ok(serde_json::json!({"error": error})),
1484 FailurePolicy::SkipDependents => Err(failure(
1485 WorkflowFailureCode::DependencySkipped,
1486 format!("{} (dependents skipped)", error.message),
1487 false,
1488 )),
1489 FailurePolicy::FailFast => Err(error),
1490 }
1491 }
1492 }
1493 }
1494
1495 async fn dispatch(
1496 &self,
1497 step: &WorkflowStepDefinition,
1498 input: Value,
1499 instance_id: &str,
1500 ) -> Result<Value, WorkflowFailure> {
1501 let session_id = { self.snapshot.lock().await.session_id.clone() };
1502 match &step.kind {
1503 WorkflowStepKind::Tool {
1504 tool, capabilities, ..
1505 } => {
1506 self.authorize(
1507 &session_id,
1508 WorkflowPolicyTarget::Tool(tool.clone()),
1509 capabilities,
1510 )
1511 .await?;
1512 let (resolved_input, resolved_secrets) =
1513 self.resolve_secret_handles(&input, &session_id).await?;
1514 let arguments = serde_json::to_string(&resolved_input).map_err(|error| {
1515 failure(WorkflowFailureCode::InvalidInput, error.to_string(), false)
1516 })?;
1517 let call = ToolCall {
1518 id: format!("workflow-{}", Uuid::new_v4()),
1519 tool_type: "function".to_string(),
1520 function: FunctionCall {
1521 name: tool.clone(),
1522 arguments,
1523 },
1524 };
1525 let context = ToolExecutionContext {
1526 session_id: Some(&session_id),
1527 tool_call_id: &call.id,
1528 event_tx: None,
1529 available_tool_schemas: None,
1530 bypass_permissions: false,
1531 can_async_resume: false,
1532 bash_completion_sink: None,
1533 pre_parsed_args: Some(&resolved_input),
1534 };
1535 let outcome = self
1536 .engine
1537 .tools
1538 .execute_with_context_outcome(&call, context)
1539 .await
1540 .map_err(|error| {
1541 let (code, message, retryable) = match error {
1542 bamboo_agent_core::tools::ToolError::NotFound(_) => (
1543 WorkflowFailureCode::UnknownReference,
1544 "workflow tool is not available",
1545 false,
1546 ),
1547 bamboo_agent_core::tools::ToolError::InvalidArguments(_) => (
1548 WorkflowFailureCode::InvalidInput,
1549 "workflow tool arguments were rejected",
1550 false,
1551 ),
1552 bamboo_agent_core::tools::ToolError::Execution(_) => (
1553 WorkflowFailureCode::ExecutionFailed,
1554 "workflow tool execution was denied or failed",
1555 true,
1556 ),
1557 };
1558 failure(code, message, retryable)
1559 })?;
1560 match outcome {
1561 ToolOutcome::Completed(result) => {
1562 let output = parse_tool_result(result)?;
1563 if contains_any_secret_material(&output, &resolved_secrets) {
1564 return Err(failure(
1565 WorkflowFailureCode::InvalidOutput,
1566 "workflow tool output contained resolved secret material",
1567 false,
1568 ));
1569 }
1570 Ok(output)
1571 }
1572 ToolOutcome::NeedsHuman { question, .. } => {
1573 self.persist_suspension(WorkflowSuspensionContext::ToolApproval {
1574 step_id: instance_id.to_string(),
1575 tool: tool.clone(),
1576 tool_call_id: question.tool_call_id,
1577 })
1578 .await?;
1579 Err(failure(
1580 WorkflowFailureCode::Suspended,
1581 "workflow tool requires human approval",
1582 true,
1583 ))
1584 }
1585 ToolOutcome::Running(handle) => {
1586 let tool_call_id = handle.tool_call_id.clone();
1587 (handle.kill)();
1588 self.persist_suspension(WorkflowSuspensionContext::ToolRunning {
1589 step_id: instance_id.to_string(),
1590 tool: tool.clone(),
1591 tool_call_id,
1592 killed: true,
1593 })
1594 .await?;
1595 Err(failure(
1596 WorkflowFailureCode::Suspended,
1597 "workflow tool is running without a durable workflow resume handle",
1598 true,
1599 ))
1600 }
1601 }
1602 }
1603 WorkflowStepKind::Agent {
1604 agent,
1605 model,
1606 effort,
1607 capabilities,
1608 structured_output_attempts,
1609 ..
1610 } => {
1611 if contains_secret_handle(&input) {
1612 return Err(failure(
1613 WorkflowFailureCode::PermissionDenied,
1614 "secret capability handles are supported only for tool arguments",
1615 false,
1616 ));
1617 }
1618 self.authorize(
1619 &session_id,
1620 WorkflowPolicyTarget::Agent(agent.clone()),
1621 capabilities,
1622 )
1623 .await?;
1624 let spec = self.pinned_agents.get(agent).cloned().ok_or_else(|| {
1625 failure(
1626 WorkflowFailureCode::PermissionDenied,
1627 "named agent was not pinned during preflight",
1628 false,
1629 )
1630 })?;
1631 let requested = capabilities.iter().cloned().collect::<BTreeSet<_>>();
1632 if !requested.is_subset(&spec.allowed_capabilities) {
1633 return Err(failure(
1634 WorkflowFailureCode::PermissionDenied,
1635 "named agent capability intersection changed",
1636 false,
1637 ));
1638 }
1639 let mut last_error = None;
1640 for _ in 0..*structured_output_attempts {
1641 self.reserve_agent().await?;
1642 self.checkpoint_usage("agent_reserved").await?;
1643 match self
1644 .engine
1645 .agents
1646 .execute(
1647 &spec,
1648 input.clone(),
1649 model.as_deref(),
1650 effort.as_deref(),
1651 &requested,
1652 &session_id,
1653 )
1654 .await
1655 {
1656 Ok(result) => {
1657 let exceeded =
1658 self.record_usage(result.tokens, result.cost_micros).await;
1659 self.checkpoint_usage("agent_usage_recorded").await?;
1660 if let Some(error) = exceeded {
1661 return Err(error);
1662 }
1663 if let Some(schema) = &step.output_schema {
1664 if let Err(error) = validate_schema(schema, &result.output) {
1665 last_error = Some(error);
1666 continue;
1667 }
1668 }
1669 return Ok(result.output);
1670 }
1671 Err(_error) => {
1672 last_error = Some("named agent execution failed".to_string())
1673 }
1674 }
1675 }
1676 Err(failure(
1677 WorkflowFailureCode::InvalidOutput,
1678 last_error.unwrap_or_else(|| "agent structured output exhausted".to_string()),
1679 false,
1680 ))
1681 }
1682 WorkflowStepKind::Workflow {
1683 workflow_id,
1684 revision,
1685 ..
1686 } => {
1687 if self.depth + 1 >= self.root_limits.max_nesting_depth {
1688 return Err(failure(
1689 WorkflowFailureCode::BudgetExceeded,
1690 "nested workflow depth exceeded",
1691 false,
1692 ));
1693 }
1694 let definition = self
1695 .bundle
1696 .get(workflow_id, *revision)
1697 .cloned()
1698 .ok_or_else(|| {
1699 failure(
1700 WorkflowFailureCode::UnknownReference,
1701 format!("persisted bundle missing workflow {workflow_id}@{revision}"),
1702 false,
1703 )
1704 })?;
1705 let nested = StartWorkflowRun {
1706 definition,
1707 args: input,
1708 session_id,
1709 workspace_trusted: self.workspace_trusted,
1710 allowed_capabilities: self.allowed_capabilities.iter().cloned().collect(),
1711 };
1712 let parent_run_id = self.snapshot.lock().await.run_id.clone();
1713 let result = Box::pin(self.engine.run_internal(
1714 nested,
1715 self.bundle.clone(),
1716 self.pinned_agents.clone(),
1717 Some(parent_run_id),
1718 Some(instance_id.to_string()),
1719 self.depth + 1,
1720 self.branch_cancellation.clone(),
1721 self.ledger.clone(),
1722 self.root_limits.clone(),
1723 self.semaphore.clone(),
1724 None,
1725 ))
1726 .await
1727 .map_err(|_error| {
1728 failure(
1729 WorkflowFailureCode::ExecutionFailed,
1730 "nested workflow execution failed",
1731 false,
1732 )
1733 })?;
1734 result.output.ok_or_else(|| {
1735 result.failure.unwrap_or_else(|| {
1736 failure(
1737 WorkflowFailureCode::ExecutionFailed,
1738 "nested workflow returned no output",
1739 false,
1740 )
1741 })
1742 })
1743 }
1744 }
1745 }
1746
1747 async fn authorize(
1748 &self,
1749 session_id: &str,
1750 target: WorkflowPolicyTarget,
1751 capabilities: &[String],
1752 ) -> Result<(), WorkflowFailure> {
1753 let requested = capabilities.iter().cloned().collect::<BTreeSet<_>>();
1754 if !requested.is_subset(&self.allowed_capabilities) {
1755 return Err(failure(
1756 WorkflowFailureCode::PermissionDenied,
1757 "step capability exceeds root policy",
1758 false,
1759 ));
1760 }
1761 match self
1762 .engine
1763 .policy
1764 .authorize(session_id, &target, &requested, self.workspace_trusted)
1765 .await
1766 {
1767 PermissionDecision::Allow => Ok(()),
1768 PermissionDecision::Deny(_reason) => Err(failure(
1769 if self.workspace_trusted {
1770 WorkflowFailureCode::PermissionDenied
1771 } else {
1772 WorkflowFailureCode::UntrustedWorkspace
1773 },
1774 "workflow policy denied this step",
1775 false,
1776 )),
1777 }
1778 }
1779
1780 async fn step_transition(
1781 &self,
1782 step_id: &str,
1783 kind: WorkflowRunEventKind,
1784 mutate: impl FnOnce(&mut WorkflowRunSnapshot),
1785 ) -> Result<(), WorkflowFailure> {
1786 let usage = self.ledger.lock().await.clone();
1787 let mut snapshot = self.snapshot.lock().await;
1788 snapshot.usage = usage;
1789 self.engine
1790 .transition(&mut snapshot, Some(step_id.to_string()), kind, mutate)
1791 .await
1792 .map_err(|error| failure(WorkflowFailureCode::Storage, error.to_string(), false))
1793 }
1794
1795 async fn reserve_step(&self) -> Result<(), WorkflowFailure> {
1796 let mut usage = self.ledger.lock().await;
1797 if usage.steps >= self.root_limits.max_steps {
1798 return Err(failure(
1799 WorkflowFailureCode::BudgetExceeded,
1800 "workflow step budget exceeded",
1801 false,
1802 ));
1803 }
1804 usage.steps += 1;
1805 Ok(())
1806 }
1807
1808 async fn reserve_retry(&self) -> Result<(), WorkflowFailure> {
1809 let mut usage = self.ledger.lock().await;
1810 if usage.retries >= self.root_limits.max_retries {
1811 return Err(failure(
1812 WorkflowFailureCode::BudgetExceeded,
1813 "workflow retry budget exceeded",
1814 false,
1815 ));
1816 }
1817 usage.retries += 1;
1818 Ok(())
1819 }
1820
1821 async fn reserve_agent(&self) -> Result<(), WorkflowFailure> {
1822 let mut usage = self.ledger.lock().await;
1823 if usage.agents >= self.root_limits.max_agents {
1824 return Err(failure(
1825 WorkflowFailureCode::BudgetExceeded,
1826 "workflow agent budget exceeded",
1827 false,
1828 ));
1829 }
1830 usage.agents += 1;
1831 Ok(())
1832 }
1833
1834 async fn record_usage(&self, tokens: u64, cost_micros: u64) -> Option<WorkflowFailure> {
1835 let mut usage = self.ledger.lock().await;
1836 let next_tokens = usage.tokens.saturating_add(tokens);
1837 let next_cost = usage.cost_micros.saturating_add(cost_micros);
1838 usage.tokens = next_tokens;
1839 usage.cost_micros = next_cost;
1840 if self
1841 .root_limits
1842 .max_tokens
1843 .is_some_and(|limit| next_tokens > limit)
1844 || self
1845 .root_limits
1846 .max_cost_micros
1847 .is_some_and(|limit| next_cost > limit)
1848 {
1849 return Some(failure(
1850 WorkflowFailureCode::BudgetExceeded,
1851 "workflow token/cost budget exceeded",
1852 false,
1853 ));
1854 }
1855 None
1856 }
1857
1858 async fn checkpoint_usage(&self, name: &str) -> Result<(), WorkflowFailure> {
1859 let usage = self.ledger.lock().await.clone();
1860 let mut snapshot = self.snapshot.lock().await;
1861 self.engine
1862 .transition(
1863 &mut snapshot,
1864 None,
1865 WorkflowRunEventKind::Phase {
1866 name: name.to_string(),
1867 },
1868 move |snapshot| snapshot.usage = usage,
1869 )
1870 .await
1871 .map_err(|error| failure(WorkflowFailureCode::Storage, error.to_string(), false))
1872 }
1873
1874 async fn persist_suspension(
1875 &self,
1876 context: WorkflowSuspensionContext,
1877 ) -> Result<(), WorkflowFailure> {
1878 let mut snapshot = self.snapshot.lock().await;
1879 self.engine
1880 .transition(
1881 &mut snapshot,
1882 None,
1883 WorkflowRunEventKind::Phase {
1884 name: "suspension_context_persisted".to_string(),
1885 },
1886 move |snapshot| snapshot.suspension = Some(context),
1887 )
1888 .await
1889 .map_err(|error| failure(WorkflowFailureCode::Storage, error.to_string(), false))
1890 }
1891
1892 async fn cancel_active_parallel_steps(
1893 &self,
1894 nodes: &[WorkflowPlan],
1895 ) -> Result<(), WorkflowFailure> {
1896 let sibling_steps = nodes
1897 .iter()
1898 .flat_map(plan_step_ids)
1899 .collect::<BTreeSet<_>>();
1900 let active = {
1901 let snapshot = self.snapshot.lock().await;
1902 snapshot
1903 .steps
1904 .iter()
1905 .filter(|(id, step)| {
1906 matches!(
1907 step.status,
1908 WorkflowStepStatus::Queued | WorkflowStepStatus::Running
1909 ) && sibling_steps
1910 .iter()
1911 .any(|step_id| instance_is_in_scope(id, step_id, &self.scope))
1912 })
1913 .map(|(id, _)| id.clone())
1914 .collect::<Vec<_>>()
1915 };
1916 for step_id in active {
1917 let state_id = step_id.clone();
1918 self.step_transition(
1919 &step_id,
1920 WorkflowRunEventKind::StepCancelled,
1921 move |snapshot| {
1922 if let Some(step) = snapshot.steps.get_mut(&state_id) {
1923 step.status = WorkflowStepStatus::Cancelled;
1924 step.failure = Some(failure(
1925 WorkflowFailureCode::Cancelled,
1926 "parallel sibling cancelled by fail_fast",
1927 false,
1928 ));
1929 }
1930 },
1931 )
1932 .await?;
1933 }
1934 Ok(())
1935 }
1936
1937 async fn resolve_secret_handles(
1938 &self,
1939 value: &Value,
1940 session_id: &str,
1941 ) -> Result<(Value, Vec<String>), WorkflowFailure> {
1942 fn walk<'a>(
1943 context: &'a RunContext,
1944 value: &'a Value,
1945 session_id: &'a str,
1946 ) -> SecretResolutionFuture<'a> {
1947 Box::pin(async move {
1948 match value {
1949 Value::Object(object) if object.contains_key("$secret") => {
1950 let handle: bamboo_domain::WorkflowSecretHandle =
1951 serde_json::from_value(value.clone()).map_err(|_| {
1952 failure(
1953 WorkflowFailureCode::InvalidInput,
1954 "malformed secret capability handle",
1955 false,
1956 )
1957 })?;
1958 let material = context
1959 .engine
1960 .secrets
1961 .resolve(session_id, &handle.capability)
1962 .await
1963 .map_err(|_| {
1964 failure(
1965 WorkflowFailureCode::PermissionDenied,
1966 "secret capability resolution denied",
1967 false,
1968 )
1969 })?;
1970 let material = material.into_exposed();
1971 Ok((Value::String(material.clone()), vec![material]))
1972 }
1973 Value::Object(object) => {
1974 let mut resolved = serde_json::Map::new();
1975 let mut secrets = Vec::new();
1976 for (key, child) in object {
1977 let (child, mut child_secrets) =
1978 walk(context, child, session_id).await?;
1979 resolved.insert(key.clone(), child);
1980 secrets.append(&mut child_secrets);
1981 }
1982 Ok((Value::Object(resolved), secrets))
1983 }
1984 Value::Array(array) => {
1985 let mut resolved = Vec::with_capacity(array.len());
1986 let mut secrets = Vec::new();
1987 for child in array {
1988 let (child, mut child_secrets) =
1989 walk(context, child, session_id).await?;
1990 resolved.push(child);
1991 secrets.append(&mut child_secrets);
1992 }
1993 Ok((Value::Array(resolved), secrets))
1994 }
1995 value => Ok((value.clone(), Vec::new())),
1996 }
1997 })
1998 }
1999 walk(self, value, session_id).await
2000 }
2001
2002 fn skip_plan<'a>(&'a self, plan: &'a WorkflowPlan, reason: &'a str) -> NodeFuture<'a> {
2003 Box::pin(async move {
2004 match plan {
2005 WorkflowPlan::Step { step } => {
2006 let instance_id = if self.scope == "root" {
2007 step.clone()
2008 } else {
2009 format!("{step}@{}", self.scope)
2010 };
2011 let reason_owned = reason.to_string();
2012 let state_id = instance_id.clone();
2013 self.step_transition(
2014 &instance_id,
2015 WorkflowRunEventKind::StepSkipped {
2016 reason: reason.to_string(),
2017 },
2018 move |snapshot| {
2019 let state = snapshot.steps.entry(state_id.clone()).or_insert(
2020 WorkflowStepSnapshot {
2021 id: state_id,
2022 status: WorkflowStepStatus::Skipped,
2023 input_hash: String::new(),
2024 output: None,
2025 failure: Some(failure(
2026 WorkflowFailureCode::DependencySkipped,
2027 reason_owned.clone(),
2028 false,
2029 )),
2030 attempts: 0,
2031 },
2032 );
2033 state.status = WorkflowStepStatus::Skipped;
2034 state.failure = Some(failure(
2035 WorkflowFailureCode::DependencySkipped,
2036 reason_owned,
2037 false,
2038 ));
2039 },
2040 )
2041 .await?;
2042 }
2043 WorkflowPlan::Sequence { nodes } | WorkflowPlan::Parallel { nodes } => {
2044 for node in nodes {
2045 self.skip_plan(node, reason).await?;
2046 }
2047 }
2048 WorkflowPlan::Map { body, .. } | WorkflowPlan::Retry { node: body, .. } => {
2049 self.skip_plan(body, reason).await?;
2050 }
2051 }
2052 Ok(Value::Null)
2053 })
2054 }
2055
2056 async fn resolve_template(&self, value: &Value) -> Result<Value, WorkflowFailure> {
2057 Box::pin(self.resolve_template_inner(value)).await
2058 }
2059
2060 fn resolve_template_inner<'a>(
2061 &'a self,
2062 value: &'a Value,
2063 ) -> Pin<Box<dyn Future<Output = Result<Value, WorkflowFailure>> + Send + 'a>> {
2064 Box::pin(async move {
2065 match value {
2066 Value::Object(object) if object.get("from").is_some() => {
2067 let reference: ValueRef =
2068 serde_json::from_value(value.clone()).map_err(|error| {
2069 failure(
2070 WorkflowFailureCode::InvalidInput,
2071 format!("malformed value reference: {error}"),
2072 false,
2073 )
2074 })?;
2075 self.resolve_ref(&reference).await
2076 }
2077 Value::Object(object) => {
2078 let mut resolved = serde_json::Map::new();
2079 for (key, child) in object {
2080 resolved.insert(key.clone(), self.resolve_template_inner(child).await?);
2081 }
2082 Ok(Value::Object(resolved))
2083 }
2084 Value::Array(array) => {
2085 let mut resolved = Vec::with_capacity(array.len());
2086 for child in array {
2087 resolved.push(self.resolve_template_inner(child).await?);
2088 }
2089 Ok(Value::Array(resolved))
2090 }
2091 value => Ok(value.clone()),
2092 }
2093 })
2094 }
2095
2096 async fn resolve_ref(&self, reference: &ValueRef) -> Result<Value, WorkflowFailure> {
2097 let (root, pointer) = match reference {
2098 ValueRef::Args { pointer } => (
2099 self.snapshot.lock().await.validated_args.clone(),
2100 pointer.as_str(),
2101 ),
2102 ValueRef::Step { step, pointer } => {
2103 let snapshot = self.snapshot.lock().await;
2104 let exact = format!("{step}@{}", self.scope);
2105 let output = snapshot
2106 .steps
2107 .get(&exact)
2108 .or_else(|| snapshot.steps.get(step))
2109 .and_then(|state| state.output.clone())
2110 .ok_or_else(|| {
2111 failure(
2112 WorkflowFailureCode::UnknownReference,
2113 format!("step output '{step}' unavailable in execution scope"),
2114 false,
2115 )
2116 })?;
2117 (output, pointer.as_str())
2118 }
2119 ValueRef::Item { name, pointer } => (
2120 self.items.get(name).cloned().ok_or_else(|| {
2121 failure(
2122 WorkflowFailureCode::UnknownReference,
2123 format!("map item '{name}' unavailable"),
2124 false,
2125 )
2126 })?,
2127 pointer.as_str(),
2128 ),
2129 ValueRef::Literal { value } => return Ok(value.clone()),
2130 };
2131 if pointer.is_empty() {
2132 Ok(root)
2133 } else {
2134 root.pointer(pointer).cloned().ok_or_else(|| {
2135 failure(
2136 WorkflowFailureCode::UnknownReference,
2137 format!("JSON pointer '{pointer}' not found"),
2138 false,
2139 )
2140 })
2141 }
2142 }
2143
2144 fn check_cancelled(&self) -> Result<(), WorkflowFailure> {
2145 if self.cancellation.is_cancelled() || self.branch_cancellation.is_cancelled() {
2146 Err(failure(
2147 WorkflowFailureCode::Cancelled,
2148 "workflow cancelled",
2149 false,
2150 ))
2151 } else {
2152 Ok(())
2153 }
2154 }
2155}
2156
2157fn event(
2158 snapshot: &WorkflowRunSnapshot,
2159 step_id: Option<String>,
2160 kind: WorkflowRunEventKind,
2161) -> WorkflowRunEvent {
2162 WorkflowRunEvent {
2163 run_id: snapshot.run_id.clone(),
2164 sequence: snapshot.last_sequence,
2165 at: Utc::now(),
2166 step_id,
2167 kind,
2168 }
2169}
2170fn failure(
2171 code: WorkflowFailureCode,
2172 message: impl Into<String>,
2173 retryable: bool,
2174) -> WorkflowFailure {
2175 WorkflowFailure {
2176 code,
2177 message: message.into(),
2178 retryable,
2179 }
2180}
2181fn storage(error: std::io::Error) -> WorkflowRunError {
2182 WorkflowRunError::Storage(error.to_string())
2183}
2184fn exceeds_optional(requested: Option<u64>, ceiling: Option<u64>) -> bool {
2185 match (requested, ceiling) {
2186 (Some(requested), Some(ceiling)) => requested > ceiling,
2187 _ => false,
2188 }
2189}
2190
2191fn enforce_budget_within(
2192 requested: &WorkflowBudgets,
2193 ceiling: &WorkflowBudgets,
2194) -> Result<(), &'static str> {
2195 if requested.max_concurrency > ceiling.max_concurrency {
2196 Err("max_concurrency")
2197 } else if requested.max_agents > ceiling.max_agents {
2198 Err("max_agents")
2199 } else if requested.max_steps > ceiling.max_steps {
2200 Err("max_steps")
2201 } else if requested.max_retries > ceiling.max_retries {
2202 Err("max_retries")
2203 } else if requested.max_nesting_depth > ceiling.max_nesting_depth {
2204 Err("max_nesting_depth")
2205 } else if requested.wall_time_ms > ceiling.wall_time_ms {
2206 Err("wall_time_ms")
2207 } else if exceeds_optional(requested.max_tokens, ceiling.max_tokens) {
2208 Err("max_tokens")
2209 } else if exceeds_optional(requested.max_cost_micros, ceiling.max_cost_micros) {
2210 Err("max_cost_micros")
2211 } else {
2212 Ok(())
2213 }
2214}
2215
2216fn definition_bundle_hash(bundle: &WorkflowDefinitionBundle) -> Result<String, WorkflowRunError> {
2217 let bytes = serde_json::to_vec(bundle)
2218 .map_err(|_| WorkflowRunError::Preflight("workflow bundle hashing failed".to_string()))?;
2219 Ok(hex::encode(Sha256::digest(bytes)))
2220}
2221fn parse_tool_result(result: ToolResult) -> Result<Value, WorkflowFailure> {
2222 if !result.success {
2223 return Err(failure(
2224 WorkflowFailureCode::ExecutionFailed,
2225 "workflow tool reported failure",
2226 true,
2227 ));
2228 }
2229 Ok(serde_json::from_str(&result.result).unwrap_or(Value::String(result.result)))
2230}
2231fn reject_secret_material(value: &Value) -> Result<(), String> {
2232 reject_secret_material_inner(value, false)
2233}
2234
2235fn reject_secret_material_in_definition(value: &Value) -> Result<(), String> {
2236 reject_secret_material_inner(value, true)
2237}
2238
2239fn reject_secret_material_inner(value: &Value, allow_bindings: bool) -> Result<(), String> {
2240 fn walk(value: &Value, key: Option<&str>, allow_bindings: bool) -> Result<(), String> {
2241 if value.as_object().is_some_and(|object| {
2242 object.len() == 1
2243 && object
2244 .get("$secret")
2245 .and_then(Value::as_str)
2246 .is_some_and(|handle| !handle.trim().is_empty())
2247 }) {
2248 return Ok(());
2249 }
2250 let safe_binding = allow_bindings
2251 && serde_json::from_value::<ValueRef>(value.clone())
2252 .is_ok_and(|reference| !matches!(reference, ValueRef::Literal { .. }));
2253 if key.is_some_and(|key| {
2254 let normalized = key
2255 .chars()
2256 .filter(|character| character.is_ascii_alphanumeric())
2257 .flat_map(char::to_lowercase)
2258 .collect::<String>();
2259 matches!(
2260 normalized.as_str(),
2261 "secret"
2262 | "token"
2263 | "password"
2264 | "credential"
2265 | "credentials"
2266 | "apikey"
2267 | "accesskey"
2268 | "accesstoken"
2269 | "secretkey"
2270 | "privatekey"
2271 )
2272 }) && !safe_binding
2273 {
2274 return Err("secret-bearing fields are not accepted by workflow runs".to_string());
2275 }
2276 if value.as_str().is_some_and(|value| {
2277 let trimmed = value.trim();
2278 trimmed.starts_with("capability://")
2279 || trimmed.starts_with("Bearer ")
2280 || trimmed.starts_with("sk-")
2281 || trimmed.starts_with("ghp_")
2282 || trimmed.starts_with("github_pat_")
2283 }) {
2284 return Err("opaque credential handles are not enabled for workflows".to_string());
2288 }
2289 match value {
2290 Value::Object(object) => {
2291 for (key, value) in object {
2292 if key == "properties" {
2293 let properties = value.as_object().ok_or_else(|| {
2294 "workflow schema properties must be an object".to_string()
2295 })?;
2296 for schema in properties.values() {
2297 walk(schema, None, allow_bindings)?;
2298 }
2299 } else {
2300 walk(value, Some(key), allow_bindings)?;
2301 }
2302 }
2303 }
2304 Value::Array(array) => {
2305 for value in array {
2306 walk(value, None, allow_bindings)?;
2307 }
2308 }
2309 _ => {}
2310 }
2311 Ok(())
2312 }
2313 walk(value, None, allow_bindings)
2314}
2315
2316fn contains_secret_handle(value: &Value) -> bool {
2317 match value {
2318 Value::Object(object) => {
2319 object.contains_key("$secret") || object.values().any(contains_secret_handle)
2320 }
2321 Value::Array(array) => array.iter().any(contains_secret_handle),
2322 _ => false,
2323 }
2324}
2325
2326fn contains_any_secret_material(value: &Value, secrets: &[String]) -> bool {
2327 let matches = |candidate: &str| {
2328 secrets
2329 .iter()
2330 .any(|secret| !secret.is_empty() && candidate.contains(secret))
2331 };
2332 fn walk(value: &Value, matches: &impl Fn(&str) -> bool) -> bool {
2333 match value {
2334 Value::String(value) => matches(value),
2335 Value::Object(object) => object
2336 .iter()
2337 .any(|(key, value)| matches(key) || walk(value, matches)),
2338 Value::Array(array) => array.iter().any(|value| walk(value, matches)),
2339 _ => false,
2340 }
2341 }
2342 walk(value, &matches)
2343}
2344
2345fn plan_leaf_count(plan: &WorkflowPlan) -> usize {
2346 match plan {
2347 WorkflowPlan::Step { .. } => 1,
2348 WorkflowPlan::Sequence { nodes } | WorkflowPlan::Parallel { nodes } => {
2349 nodes.iter().fold(0usize, |total, node| {
2350 total.saturating_add(plan_leaf_count(node))
2351 })
2352 }
2353 WorkflowPlan::Map { body, .. } | WorkflowPlan::Retry { node: body, .. } => {
2354 plan_leaf_count(body)
2355 }
2356 }
2357}
2358
2359fn plan_step_ids(plan: &WorkflowPlan) -> Vec<String> {
2360 match plan {
2361 WorkflowPlan::Step { step } => vec![step.clone()],
2362 WorkflowPlan::Sequence { nodes } | WorkflowPlan::Parallel { nodes } => {
2363 nodes.iter().flat_map(plan_step_ids).collect()
2364 }
2365 WorkflowPlan::Map { body, .. } | WorkflowPlan::Retry { node: body, .. } => {
2366 plan_step_ids(body)
2367 }
2368 }
2369}
2370
2371fn instance_is_in_scope(instance_id: &str, step_id: &str, scope: &str) -> bool {
2372 if scope == "root" {
2373 instance_id == step_id
2374 || instance_id
2375 .strip_prefix(&format!("{step_id}@root"))
2376 .is_some_and(|suffix| suffix.starts_with('['))
2377 } else {
2378 let exact = format!("{step_id}@{scope}");
2379 instance_id == exact
2380 || instance_id
2381 .strip_prefix(&exact)
2382 .is_some_and(|suffix| suffix.starts_with('['))
2383 }
2384}
2385
2386fn plan_frontier(plan: &WorkflowPlan) -> Vec<String> {
2387 match plan {
2388 WorkflowPlan::Step { step } => vec![step.clone()],
2389 WorkflowPlan::Sequence { nodes } => nodes.first().map_or_else(Vec::new, plan_frontier),
2390 WorkflowPlan::Parallel { nodes } => nodes.iter().flat_map(plan_frontier).collect(),
2391 WorkflowPlan::Map { body, .. } | WorkflowPlan::Retry { node: body, .. } => {
2392 plan_frontier(body)
2393 }
2394 }
2395}
2396
2397fn validate_nested_input_contract(
2398 template: &Value,
2399 target_schema: &Value,
2400 compiled: &CompiledWorkflow,
2401) -> Result<(), String> {
2402 fn contains_ref(value: &Value) -> bool {
2403 match value {
2404 Value::Object(object) => {
2405 object.contains_key("from") || object.values().any(contains_ref)
2406 }
2407 Value::Array(array) => array.iter().any(contains_ref),
2408 _ => false,
2409 }
2410 }
2411 if !contains_ref(template) {
2412 return validate_schema(target_schema, template)
2413 .map_err(|error| format!("nested workflow input is incompatible: {error}"));
2414 }
2415 let reference: ValueRef = serde_json::from_value(template.clone()).map_err(|_| {
2416 "nested dynamic input schema cannot be proven compatible in phase 1".to_string()
2417 })?;
2418 let source_schema = match reference {
2419 ValueRef::Args { pointer } => {
2420 schema_at_pointer(&compiled.definition.input_schema, &pointer)
2421 }
2422 ValueRef::Step { step, pointer } => compiled
2423 .steps
2424 .get(&step)
2425 .and_then(|step| step.output_schema.as_ref())
2426 .and_then(|schema| schema_at_pointer(schema, &pointer)),
2427 ValueRef::Literal { value } => {
2428 return validate_schema(target_schema, &value)
2429 .map_err(|error| format!("nested workflow input is incompatible: {error}"));
2430 }
2431 ValueRef::Item { .. } => None,
2432 }
2433 .ok_or_else(|| {
2434 "nested dynamic input source schema is missing or pointer is invalid".to_string()
2435 })?;
2436 if schema_compatible(source_schema, target_schema) {
2437 Ok(())
2438 } else {
2439 Err("nested workflow input schema is not compatible with its pinned target".to_string())
2440 }
2441}
2442
2443fn schema_at_pointer<'a>(schema: &'a Value, pointer: &str) -> Option<&'a Value> {
2444 if pointer.is_empty() {
2445 return Some(schema);
2446 }
2447 let mut current = schema;
2448 for token in pointer.strip_prefix('/')?.split('/') {
2449 let token = token.replace("~1", "/").replace("~0", "~");
2450 current = if token.parse::<usize>().is_ok() {
2451 current.get("items")?
2452 } else {
2453 current.get("properties")?.get(&token)?
2454 };
2455 }
2456 Some(current)
2457}
2458
2459fn schema_compatible(source: &Value, target: &Value) -> bool {
2460 if source == target {
2461 return true;
2462 }
2463 let source_type = source.get("type").and_then(Value::as_str);
2464 let target_type = target.get("type").and_then(Value::as_str);
2465 source_type.is_some() && source_type == target_type && target_type != Some("object")
2466}
2467
2468fn effective_limits(requested: &WorkflowBudgets, ceilings: &WorkflowBudgets) -> WorkflowBudgets {
2469 WorkflowBudgets {
2470 max_concurrency: requested.max_concurrency.min(ceilings.max_concurrency),
2471 max_agents: requested.max_agents.min(ceilings.max_agents),
2472 max_steps: requested.max_steps.min(ceilings.max_steps),
2473 max_retries: requested.max_retries.min(ceilings.max_retries),
2474 max_nesting_depth: requested.max_nesting_depth.min(ceilings.max_nesting_depth),
2475 wall_time_ms: requested.wall_time_ms.min(ceilings.wall_time_ms),
2476 max_tokens: match (requested.max_tokens, ceilings.max_tokens) {
2477 (Some(requested), Some(ceiling)) => Some(requested.min(ceiling)),
2478 (Some(requested), None) => Some(requested),
2479 (None, ceiling) => ceiling,
2480 },
2481 max_cost_micros: match (requested.max_cost_micros, ceilings.max_cost_micros) {
2482 (Some(requested), Some(ceiling)) => Some(requested.min(ceiling)),
2483 (Some(requested), None) => Some(requested),
2484 (None, ceiling) => ceiling,
2485 },
2486 }
2487}