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