1use std::collections::{BTreeMap, BTreeSet};
2use std::path::{Path, PathBuf};
3use std::sync::atomic::{AtomicBool, Ordering};
4use std::sync::Arc;
5
6use async_trait::async_trait;
7use bamboo_agent_core::tools::{
8 Tool, ToolClass, ToolCtx, ToolError, ToolExecutionSessionFlags, ToolOutcome, ToolResult,
9};
10use bamboo_domain::{
11 StartWorkflowRun, WorkflowBudgetUsage, WorkflowBudgets, WorkflowDefinitionBundle,
12 WorkflowFailure, WorkflowFailureCode, WorkflowPlan, WorkflowProgress, WorkflowRunDefinition,
13 WorkflowRunEvent, WorkflowRunEventKind, WorkflowRunSnapshot, WorkflowRunStatus,
14 WorkflowStepKind, WorkflowStepSnapshot, WorkflowStepStatus, WorkflowSuspensionContext,
15};
16use bamboo_engine::{
17 AgentStepPort, AgentStepResult, FileWorkflowRunRepository, NamedAgentSpec, PermissionDecision,
18 WorkflowDefinitionPort, WorkflowPolicyPort, WorkflowPolicyTarget, WorkflowRunEngine,
19 WorkflowRunError, WorkflowSecretMaterial, WorkflowSecretResolverPort,
20 WorkflowSessionPermissionPort,
21};
22use bamboo_skills::SkillManager;
23use serde::{Deserialize, Serialize};
24use serde_json::{json, Value};
25
26const MAX_CONCURRENCY: usize = 8;
27const MAX_AGENTS: u32 = 16;
28const MAX_STEPS: u32 = 512;
29const MAX_RETRIES: u32 = 16;
30const MAX_NESTING_DEPTH: u32 = 8;
31const MAX_WALL_TIME_MS: u64 = 60 * 60 * 1000;
32const MAX_TOKENS: u64 = 2_000_000;
33const MAX_COST_MICROS: u64 = 100_000_000;
34const MAX_PINNED_DEFINITIONS_PER_RUN: usize = 32;
35const MAX_PINNED_BUNDLE_BYTES_PER_RUN: usize = 512 * 1024;
36const MAX_WORKFLOW_RUN_IDS_PER_SESSION: usize = 256;
37const SAFE_UNTRUSTED_WORKFLOW_TOOLS: &[&str] = &[
38 "Read",
39 "read_file",
40 "GetFileInfo",
41 "Glob",
42 "list_directory",
43 "Grep",
44];
45
46#[derive(Clone)]
49pub struct WorkflowRunAccess {
50 engine: Arc<WorkflowRunEngine>,
51 skills: Arc<SkillManager>,
52 sessions: bamboo_engine::SessionRepository,
53}
54
55struct ServerWorkflowSessionPermissions {
56 sessions: bamboo_engine::SessionRepository,
57 permission_config: Arc<bamboo_tools::permission::PermissionConfig>,
58}
59
60#[async_trait]
61impl WorkflowSessionPermissionPort for ServerWorkflowSessionPermissions {
62 async fn flags_for_session(
63 &self,
64 session_id: &str,
65 ) -> Result<ToolExecutionSessionFlags, String> {
66 let session = self
67 .sessions
68 .try_load(session_id)
69 .await
70 .map_err(|error| error.to_string())?
71 .ok_or_else(|| format!("workflow session '{session_id}' does not exist"))?;
72 let configured = if session
73 .agent_runtime_state
74 .as_ref()
75 .is_some_and(|state| state.plan_mode.is_some())
76 {
77 bamboo_domain::PermissionMode::Plan
78 } else {
79 self.permission_config.mode()
80 };
81 Ok(ToolExecutionSessionFlags::from_session_and_configured_mode(
82 &session, configured,
83 ))
84 }
85}
86
87impl WorkflowRunAccess {
88 pub async fn new(
89 data_dir: &Path,
90 tools: Arc<dyn bamboo_agent_core::tools::ToolExecutor>,
91 skills: Arc<SkillManager>,
92 sessions: bamboo_engine::SessionRepository,
93 ) -> Result<Self, String> {
94 Self::new_with_permission_config(data_dir, tools, skills, sessions, None).await
95 }
96
97 pub async fn new_with_permission_config(
98 data_dir: &Path,
99 tools: Arc<dyn bamboo_agent_core::tools::ToolExecutor>,
100 skills: Arc<SkillManager>,
101 sessions: bamboo_engine::SessionRepository,
102 permission_config: Option<Arc<bamboo_tools::permission::PermissionConfig>>,
103 ) -> Result<Self, String> {
104 let repository = Arc::new(
105 FileWorkflowRunRepository::new(data_dir.join("workflow-runs"))
106 .map_err(|error| format!("failed to initialize workflow journal: {error}"))?,
107 );
108 let engine = WorkflowRunEngine::new(
109 repository,
110 tools,
111 Arc::new(UnavailableAgentPort),
112 Arc::new(ExternallyPinnedDefinitions),
113 Arc::new(ServerWorkflowPolicy),
114 Arc::new(UnavailableSecretResolver),
115 WorkflowBudgets {
116 max_concurrency: MAX_CONCURRENCY,
117 max_agents: MAX_AGENTS,
118 max_steps: MAX_STEPS,
119 max_retries: MAX_RETRIES,
120 max_nesting_depth: MAX_NESTING_DEPTH,
121 wall_time_ms: MAX_WALL_TIME_MS,
122 max_tokens: Some(MAX_TOKENS),
123 max_cost_micros: Some(MAX_COST_MICROS),
124 },
125 );
126 if let Some(permission_config) = permission_config {
127 engine.set_session_permission_port(Arc::new(ServerWorkflowSessionPermissions {
128 sessions: sessions.clone(),
129 permission_config,
130 }));
131 }
132 engine
133 .recover()
134 .await
135 .map_err(|error| format!("failed to recover workflow journal: {error}"))?;
136 Ok(Self {
137 engine,
138 skills,
139 sessions,
140 })
141 }
142
143 async fn session_context(
144 &self,
145 session_id: &str,
146 ) -> Result<(Option<PathBuf>, bool), WorkflowRunError> {
147 let session =
148 self.sessions.try_load(session_id).await.map_err(|_| {
149 WorkflowRunError::Preflight("session state is unavailable".to_string())
150 })?;
151 let session = session.ok_or_else(|| {
152 WorkflowRunError::Preflight("workflow session does not exist".to_string())
153 })?;
154 let preferred = session.workspace.map(PathBuf::from);
155 let workspace =
156 bamboo_agent_core::workspace_state::ensure_session_workspace(session_id, preferred)
157 .or_else(|| {
158 Some(
162 bamboo_agent_core::workspace_state::workspace_or_process_cwd(Some(
163 session_id,
164 )),
165 )
166 });
167 Ok((workspace, false))
173 }
174
175 pub async fn start(
176 &self,
177 session_id: &str,
178 workflow_id: &str,
179 revision: u64,
180 args: Value,
181 budget: Option<WorkflowBudgets>,
182 ) -> Result<WorkflowRunSnapshot, WorkflowRunError> {
183 self.start_for_invoker(session_id, workflow_id, revision, args, budget, false)
184 .await
185 }
186
187 pub async fn start_from_tool(
188 &self,
189 session_id: &str,
190 workflow_id: &str,
191 revision: u64,
192 args: Value,
193 budget: Option<WorkflowBudgets>,
194 ) -> Result<WorkflowRunSnapshot, WorkflowRunError> {
195 self.start_for_invoker(session_id, workflow_id, revision, args, budget, true)
196 .await
197 }
198
199 async fn start_for_invoker(
200 &self,
201 session_id: &str,
202 workflow_id: &str,
203 revision: u64,
204 args: Value,
205 budget: Option<WorkflowBudgets>,
206 model_started: bool,
207 ) -> Result<WorkflowRunSnapshot, WorkflowRunError> {
208 self.ensure_run_index_capacity(session_id).await?;
209 let (workspace, workspace_trusted) = self.session_context(session_id).await?;
210 let store = self
211 .skills
212 .store_for_workspace(workspace.as_deref())
213 .await
214 .map_err(|_| {
215 WorkflowRunError::Preflight("workflow catalog is unavailable".to_string())
216 })?;
217 let catalog = store.workflow_catalog_snapshot().await;
218 let entry = catalog
219 .entries
220 .iter()
221 .find(|entry| {
222 entry.winner
223 && entry.id == workflow_id
224 && entry.revision == revision
225 && entry.status == bamboo_skills::WorkflowStatus::Valid
226 })
227 .ok_or_else(|| {
228 WorkflowRunError::Preflight(
229 "requested workflow revision is unavailable".to_string(),
230 )
231 })?;
232 if entry.kind != bamboo_skills::WorkflowKind::Orchestration {
233 return Err(WorkflowRunError::Preflight(
234 "instruction workflows must be activated with load_skill".to_string(),
235 ));
236 }
237 if model_started {
238 let session = self
239 .sessions
240 .try_load(session_id)
241 .await
242 .map_err(|_| {
243 WorkflowRunError::Preflight("session state is unavailable".to_string())
244 })?
245 .ok_or(WorkflowRunError::NotFound)?;
246 let opted_in = session
247 .metadata
248 .get(bamboo_skills::WORKFLOW_ORCHESTRATION_OPT_IN_METADATA_KEY)
249 .is_some_and(|value| value.eq_ignore_ascii_case("true"));
250 if !opted_in {
251 return Err(WorkflowRunError::Preflight(
252 "model-started orchestration requires explicit session opt-in".to_string(),
253 ));
254 }
255 }
256 let mut bundle = self
257 .skills
258 .pin_workflow_definition_bundle(workspace.as_deref(), workflow_id, revision)
259 .await
260 .map_err(|_| WorkflowRunError::Preflight("workflow catalog pin failed".to_string()))?;
261 let policy = if model_started {
262 "automatic"
263 } else {
264 "explicit"
265 };
266 if bundle.root_invocation_policy[policy].as_bool() != Some(true) {
267 return Err(WorkflowRunError::Preflight(format!(
268 "pinned workflow invocation policy denies {policy} start"
269 )));
270 }
271 let bundle_bytes = serde_json::to_vec(&bundle)
272 .map_err(|_| WorkflowRunError::Preflight("workflow bundle is invalid".to_string()))?
273 .len();
274 enforce_pinned_bundle_limits(bundle.definitions.len(), bundle_bytes)?;
275 let mut definition = bundle.root().cloned().ok_or_else(|| {
276 WorkflowRunError::Preflight("pinned workflow root is missing".to_string())
277 })?;
278 if let Some(requested) = budget {
279 validate_requested_budget(&requested)?;
280 definition.budgets = tighten_workflow_budget(&definition.budgets, &requested);
281 let root_key = WorkflowDefinitionBundle::key(&definition.id, definition.revision);
282 bundle.definitions.insert(root_key, definition.clone());
283 }
284 let snapshot = self
285 .engine
286 .start_pinned(
287 StartWorkflowRun {
288 definition,
289 args,
290 session_id: session_id.to_string(),
291 workspace_trusted,
292 allowed_capabilities: vec!["read".to_string()],
297 },
298 bundle,
299 )
300 .await?;
301 self.index_started_run_or_compensate(session_id, snapshot)
302 .await
303 }
304
305 async fn index_started_run_or_compensate(
306 &self,
307 session_id: &str,
308 snapshot: WorkflowRunSnapshot,
309 ) -> Result<WorkflowRunSnapshot, WorkflowRunError> {
310 if let Err(error) = self.remember_run_id(session_id, &snapshot.run_id).await {
311 return match self.engine.cancel(&snapshot.run_id).await {
312 Ok(cancelled) if cancelled.status.is_terminal() => Err(WorkflowRunError::Storage(
313 format!(
314 "run index persistence failed; run {} reached terminal {:?}: {error}",
315 snapshot.run_id, cancelled.status
316 ),
317 )),
318 Ok(cancelled) => Err(WorkflowRunError::Storage(format!(
319 "run index persistence failed; orphan run {} remains {:?}; repair with this run id: {error}",
320 snapshot.run_id, cancelled.status
321 ))),
322 Err(cancel_error) => Err(WorkflowRunError::Storage(format!(
323 "run index persistence failed; orphan run {} could not be cancelled ({cancel_error}); repair with this run id: {error}",
324 snapshot.run_id
325 ))),
326 };
327 }
328 Ok(snapshot)
329 }
330
331 async fn ensure_run_index_capacity(&self, session_id: &str) -> Result<(), WorkflowRunError> {
332 let session = self
333 .sessions
334 .try_load(session_id)
335 .await
336 .map_err(|_| WorkflowRunError::Storage("session state unavailable".to_string()))?
337 .ok_or(WorkflowRunError::NotFound)?;
338 let ids = session
339 .metadata
340 .get(bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY)
341 .and_then(|raw| serde_json::from_str::<Vec<String>>(raw).ok())
342 .unwrap_or_default();
343 if ids.len() < MAX_WORKFLOW_RUN_IDS_PER_SESSION {
344 return Ok(());
345 }
346 let mut evictable = BTreeSet::new();
347 for run_id in &ids {
348 match self.engine.progress(run_id, u64::MAX).await {
349 Ok(progress) if progress.snapshot.status.is_terminal() => {
350 evictable.insert(run_id.clone());
351 }
352 Err(WorkflowRunError::NotFound) => {
353 evictable.insert(run_id.clone());
354 }
355 Ok(_) => {}
356 Err(error) => return Err(error),
357 }
358 }
359 if !evictable.is_empty() {
360 self.sessions
361 .update_runtime_session(
362 session_id,
363 &[bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY],
364 move |session| {
365 let mut ids = session
366 .metadata
367 .get(bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY)
368 .and_then(|raw| serde_json::from_str::<Vec<String>>(raw).ok())
369 .unwrap_or_default();
370 ids.retain(|id| !evictable.contains(id));
371 session.metadata.insert(
372 bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY.to_string(),
373 serde_json::to_string(&ids)
374 .expect("string vector serialization cannot fail"),
375 );
376 },
377 )
378 .await
379 .map_err(|_| WorkflowRunError::Storage("run index pruning failed".to_string()))?
380 .ok_or(WorkflowRunError::NotFound)?;
381 }
382 let remaining = self
383 .sessions
384 .try_load(session_id)
385 .await
386 .map_err(|_| WorkflowRunError::Storage("session state unavailable".to_string()))?
387 .and_then(|session| {
388 session
389 .metadata
390 .get(bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY)
391 .and_then(|raw| serde_json::from_str::<Vec<String>>(raw).ok())
392 })
393 .unwrap_or_default()
394 .len();
395 if remaining >= MAX_WORKFLOW_RUN_IDS_PER_SESSION {
396 return Err(WorkflowRunError::Preflight(
397 "workflow run index is full of active runs".to_string(),
398 ));
399 }
400 Ok(())
401 }
402
403 async fn remember_run_id(
404 &self,
405 session_id: &str,
406 run_id: &str,
407 ) -> Result<(), WorkflowRunError> {
408 let run_id = run_id.to_string();
409 let index_full = Arc::new(AtomicBool::new(false));
410 let index_full_in_transaction = index_full.clone();
411 self.sessions
412 .update_runtime_session(
413 session_id,
414 &[bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY],
415 move |session| {
416 let mut ids = session
417 .metadata
418 .get(bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY)
419 .and_then(|raw| serde_json::from_str::<Vec<String>>(raw).ok())
420 .unwrap_or_default();
421 ids.retain(|existing| existing != &run_id);
422 if ids.len() >= MAX_WORKFLOW_RUN_IDS_PER_SESSION {
423 index_full_in_transaction.store(true, Ordering::SeqCst);
424 return;
425 }
426 ids.push(run_id);
427 session.metadata.insert(
428 bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY.to_string(),
429 serde_json::to_string(&ids)
430 .expect("string vector serialization cannot fail"),
431 );
432 },
433 )
434 .await
435 .map_err(|_| WorkflowRunError::Storage("run index persistence failed".to_string()))?
436 .ok_or(WorkflowRunError::NotFound)
437 .and_then(|_| {
438 if index_full.load(Ordering::SeqCst) {
439 Err(WorkflowRunError::Storage(
440 "workflow run index reached its active-run capacity".to_string(),
441 ))
442 } else {
443 Ok(())
444 }
445 })
446 }
447
448 pub async fn list_for_session(
449 &self,
450 session_id: &str,
451 ) -> Result<Vec<WorkflowRunSnapshot>, WorkflowRunError> {
452 self.session_context(session_id).await?;
453 let session = self
454 .sessions
455 .try_load(session_id)
456 .await
457 .map_err(|_| WorkflowRunError::Storage("session state unavailable".to_string()))?
458 .ok_or(WorkflowRunError::NotFound)?;
459 let run_ids = session
460 .metadata
461 .get(bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY)
462 .and_then(|raw| serde_json::from_str::<Vec<String>>(raw).ok())
463 .unwrap_or_default();
464 let mut snapshots = Vec::new();
465 let mut stale = BTreeSet::new();
466 for run_id in run_ids {
467 match self.engine.progress(&run_id, u64::MAX).await {
468 Ok(progress) if progress.snapshot.session_id == session_id => {
469 snapshots.push(progress.snapshot);
470 }
471 Ok(_) | Err(WorkflowRunError::NotFound) => {
472 stale.insert(run_id);
473 }
474 Err(error) => return Err(error),
475 }
476 }
477 if !stale.is_empty() {
478 let _ = self
479 .sessions
480 .update_runtime_session(
481 session_id,
482 &[bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY],
483 move |session| {
484 let mut ids = session
485 .metadata
486 .get(bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY)
487 .and_then(|raw| serde_json::from_str::<Vec<String>>(raw).ok())
488 .unwrap_or_default();
489 ids.retain(|id| !stale.contains(id));
490 session.metadata.insert(
491 bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY.to_string(),
492 serde_json::to_string(&ids)
493 .expect("string vector serialization cannot fail"),
494 );
495 },
496 )
497 .await;
498 }
499 snapshots.sort_by_key(|snapshot| std::cmp::Reverse(snapshot.created_at));
500 Ok(snapshots)
501 }
502
503 pub async fn progress_for_session(
504 &self,
505 session_id: &str,
506 run_id: &str,
507 since: u64,
508 ) -> Result<WorkflowProgress, WorkflowRunError> {
509 let progress = self.engine.progress(run_id, since).await?;
510 if progress.snapshot.session_id != session_id {
511 return Err(WorkflowRunError::NotFound);
512 }
513 Ok(progress)
514 }
515
516 pub async fn cancel_for_session(
517 &self,
518 session_id: &str,
519 run_id: &str,
520 ) -> Result<WorkflowRunSnapshot, WorkflowRunError> {
521 let progress = self
522 .progress_for_session(session_id, run_id, u64::MAX)
523 .await?;
524 ensure_workflow_cancel_allowed(&progress.snapshot)?;
525 self.engine.cancel(run_id).await
526 }
527
528 pub async fn restart_for_session(
529 &self,
530 session_id: &str,
531 run_id: &str,
532 ) -> Result<WorkflowRunSnapshot, WorkflowRunError> {
533 let progress = self
534 .progress_for_session(session_id, run_id, u64::MAX)
535 .await?;
536 ensure_workflow_restart_as_new_run_allowed(&progress.snapshot)?;
537 self.ensure_run_index_capacity(session_id).await?;
538 let (_, workspace_trusted) = self.session_context(session_id).await?;
539 let snapshot = self
540 .engine
541 .restart(run_id, workspace_trusted, vec!["read".to_string()])
542 .await?;
543 self.index_started_run_or_compensate(session_id, snapshot)
544 .await
545 }
546
547 pub async fn restart_from_tool(
548 &self,
549 session_id: &str,
550 run_id: &str,
551 ) -> Result<WorkflowRunSnapshot, WorkflowRunError> {
552 let progress = self
553 .progress_for_session(session_id, run_id, u64::MAX)
554 .await?;
555 if progress.snapshot.definition_bundle.root_invocation_policy["automatic"].as_bool()
556 != Some(true)
557 {
558 return Err(WorkflowRunError::Preflight(
559 "pinned workflow invocation policy denies automatic restart".to_string(),
560 ));
561 }
562 let session = self
563 .sessions
564 .try_load(session_id)
565 .await
566 .map_err(|_| WorkflowRunError::Preflight("session state is unavailable".to_string()))?
567 .ok_or(WorkflowRunError::NotFound)?;
568 let opted_in = session
569 .metadata
570 .get(bamboo_skills::WORKFLOW_ORCHESTRATION_OPT_IN_METADATA_KEY)
571 .is_some_and(|value| value.eq_ignore_ascii_case("true"));
572 if !opted_in {
573 return Err(WorkflowRunError::Preflight(
574 "model-started orchestration restart requires explicit session opt-in".to_string(),
575 ));
576 }
577 self.restart_for_session(session_id, run_id).await
578 }
579}
580
581fn ensure_workflow_cancel_allowed(snapshot: &WorkflowRunSnapshot) -> Result<(), WorkflowRunError> {
582 match snapshot.status {
583 WorkflowRunStatus::Succeeded | WorkflowRunStatus::Failed => Err(WorkflowRunError::Terminal),
584 WorkflowRunStatus::Queued
585 | WorkflowRunStatus::Running
586 | WorkflowRunStatus::Suspended
587 | WorkflowRunStatus::Cancelled => Ok(()),
588 }
589}
590
591fn ensure_workflow_restart_as_new_run_allowed(
592 snapshot: &WorkflowRunSnapshot,
593) -> Result<(), WorkflowRunError> {
594 match (snapshot.status, snapshot.suspension.as_ref()) {
595 (WorkflowRunStatus::Suspended, Some(WorkflowSuspensionContext::Recovery { .. })) => Ok(()),
596 (status, _) if status.is_terminal() => Err(WorkflowRunError::Terminal),
597 _ => Err(WorkflowRunError::Preflight(
598 "only recovery-suspended workflows can restart as a new run".to_string(),
599 )),
600 }
601}
602
603fn tighten_workflow_budget(
604 definition: &WorkflowBudgets,
605 requested: &WorkflowBudgets,
606) -> WorkflowBudgets {
607 WorkflowBudgets {
608 max_concurrency: definition.max_concurrency.min(requested.max_concurrency),
609 max_agents: definition.max_agents.min(requested.max_agents),
610 max_steps: definition.max_steps.min(requested.max_steps),
611 max_retries: definition.max_retries.min(requested.max_retries),
612 max_nesting_depth: definition
613 .max_nesting_depth
614 .min(requested.max_nesting_depth),
615 wall_time_ms: definition.wall_time_ms.min(requested.wall_time_ms),
616 max_tokens: match (definition.max_tokens, requested.max_tokens) {
617 (Some(left), Some(right)) => Some(left.min(right)),
618 (Some(value), None) | (None, Some(value)) => Some(value),
619 (None, None) => None,
620 },
621 max_cost_micros: match (definition.max_cost_micros, requested.max_cost_micros) {
622 (Some(left), Some(right)) => Some(left.min(right)),
623 (Some(value), None) | (None, Some(value)) => Some(value),
624 (None, None) => None,
625 },
626 }
627}
628
629fn validate_requested_budget(requested: &WorkflowBudgets) -> Result<(), WorkflowRunError> {
630 if requested.max_concurrency == 0
631 || requested.max_steps == 0
632 || requested.max_nesting_depth == 0
633 || requested.wall_time_ms == 0
634 {
635 return Err(WorkflowRunError::InvalidInput(
636 "workflow execution limits must be positive".to_string(),
637 ));
638 }
639 Ok(())
640}
641
642fn enforce_pinned_bundle_limits(
643 definition_count: usize,
644 serialized_bytes: usize,
645) -> Result<(), WorkflowRunError> {
646 if definition_count > MAX_PINNED_DEFINITIONS_PER_RUN {
647 return Err(WorkflowRunError::Preflight(
648 "workflow dependency count exceeds the server limit".to_string(),
649 ));
650 }
651 if serialized_bytes > MAX_PINNED_BUNDLE_BYTES_PER_RUN {
652 return Err(WorkflowRunError::Preflight(
653 "workflow definition bundle exceeds the server size limit".to_string(),
654 ));
655 }
656 Ok(())
657}
658
659#[derive(Debug, Clone, Serialize, PartialEq)]
667pub(crate) struct PublicWorkflowRunSnapshot {
668 pub run_id: String,
669 #[serde(skip_serializing_if = "Option::is_none")]
670 pub parent_run_id: Option<String>,
671 #[serde(skip_serializing_if = "Option::is_none")]
672 pub parent_step_id: Option<String>,
673 pub session_id: String,
674 pub workflow_id: String,
675 pub workflow_revision: u64,
676 pub definition_bundle_hash: String,
677 pub status: WorkflowRunStatus,
678 pub can_cancel: bool,
679 pub can_restart_as_new_run: bool,
680 pub planned_steps: BTreeMap<String, PublicWorkflowPlannedStep>,
681 pub plan: PublicWorkflowPlan,
682 pub steps: BTreeMap<String, PublicWorkflowStepSnapshot>,
683 pub budget: WorkflowBudgets,
684 pub usage: WorkflowBudgetUsage,
685 pub child_agent_count: u32,
686 pub last_sequence: u64,
687 #[serde(skip_serializing_if = "Option::is_none")]
688 pub failure: Option<PublicWorkflowFailure>,
689 #[serde(skip_serializing_if = "Option::is_none")]
690 pub suspension: Option<PublicWorkflowSuspension>,
691 pub created_at: chrono::DateTime<chrono::Utc>,
692 pub updated_at: chrono::DateTime<chrono::Utc>,
693}
694
695#[derive(Debug, Clone, Serialize, PartialEq)]
696pub(crate) struct PublicWorkflowStepSnapshot {
697 pub id: String,
698 pub status: WorkflowStepStatus,
699 #[serde(skip_serializing_if = "Option::is_none")]
700 pub failure: Option<PublicWorkflowFailure>,
701 pub attempts: u32,
702}
703
704#[derive(Debug, Clone, Serialize, PartialEq, Eq)]
705pub(crate) struct PublicWorkflowPlannedStep {
706 pub id: String,
707 pub kind: PublicWorkflowPlannedStepKind,
708}
709
710#[derive(Debug, Clone, Copy, Serialize, PartialEq, Eq)]
711#[serde(rename_all = "snake_case")]
712pub(crate) enum PublicWorkflowPlannedStepKind {
713 Tool,
714 Agent,
715 Workflow,
716}
717
718#[derive(Debug, Clone, Serialize, PartialEq, Eq)]
721#[serde(tag = "type", rename_all = "snake_case")]
722pub(crate) enum PublicWorkflowPlan {
723 Step {
724 step: String,
725 },
726 Sequence {
727 nodes: Vec<PublicWorkflowPlan>,
728 },
729 Parallel {
730 nodes: Vec<PublicWorkflowPlan>,
731 },
732 Map {
733 body: Box<PublicWorkflowPlan>,
734 },
735 Retry {
736 node: Box<PublicWorkflowPlan>,
737 max_attempts: u32,
738 },
739}
740
741#[derive(Debug, Clone, Serialize, PartialEq, Eq)]
742pub(crate) struct PublicWorkflowFailure {
743 pub code: WorkflowFailureCode,
744 pub message: &'static str,
745 pub retryable: bool,
746}
747
748#[derive(Debug, Clone, Serialize, PartialEq, Eq)]
749#[serde(tag = "type", rename_all = "snake_case")]
750pub(crate) enum PublicWorkflowSuspension {
751 ToolApproval { step_id: String },
752 ToolRunning { step_id: String, killed: bool },
753 Recovery,
754}
755
756#[derive(Debug, Clone, Serialize, PartialEq)]
760pub(crate) struct PublicWorkflowRunEvent {
761 pub run_id: String,
762 pub sequence: u64,
763 pub at: chrono::DateTime<chrono::Utc>,
764 #[serde(skip_serializing_if = "Option::is_none")]
765 pub step_id: Option<String>,
766 #[serde(flatten)]
767 pub kind: PublicWorkflowRunEventKind,
768}
769
770#[derive(Debug, Clone, Serialize, PartialEq)]
771#[serde(tag = "type", rename_all = "snake_case")]
772pub(crate) enum PublicWorkflowRunEventKind {
773 RunQueued,
774 RunStarted,
775 Phase { name: &'static str },
776 StepQueued,
777 StepStarted,
778 StepSuspended,
779 StepCompleted,
780 StepFailed { failure: PublicWorkflowFailure },
781 StepCancelled,
782 StepSkipped,
783 RunSuspended,
784 RunSucceeded,
785 RunFailed { failure: PublicWorkflowFailure },
786 RunCancelled,
787}
788
789fn public_workflow_plan(plan: &WorkflowPlan) -> PublicWorkflowPlan {
790 match plan {
791 WorkflowPlan::Step { step } => PublicWorkflowPlan::Step { step: step.clone() },
792 WorkflowPlan::Sequence { nodes } => PublicWorkflowPlan::Sequence {
793 nodes: nodes.iter().map(public_workflow_plan).collect(),
794 },
795 WorkflowPlan::Parallel { nodes } => PublicWorkflowPlan::Parallel {
796 nodes: nodes.iter().map(public_workflow_plan).collect(),
797 },
798 WorkflowPlan::Map { body, .. } => PublicWorkflowPlan::Map {
799 body: Box::new(public_workflow_plan(body)),
800 },
801 WorkflowPlan::Retry {
802 node, max_attempts, ..
803 } => PublicWorkflowPlan::Retry {
804 node: Box::new(public_workflow_plan(node)),
805 max_attempts: *max_attempts,
806 },
807 }
808}
809
810fn public_planned_steps(
811 definition: &WorkflowRunDefinition,
812) -> BTreeMap<String, PublicWorkflowPlannedStep> {
813 definition
814 .steps
815 .iter()
816 .map(|step| {
817 let kind = match &step.kind {
818 WorkflowStepKind::Tool { .. } => PublicWorkflowPlannedStepKind::Tool,
819 WorkflowStepKind::Agent { .. } => PublicWorkflowPlannedStepKind::Agent,
820 WorkflowStepKind::Workflow { .. } => PublicWorkflowPlannedStepKind::Workflow,
821 };
822 (
823 step.id.clone(),
824 PublicWorkflowPlannedStep {
825 id: step.id.clone(),
826 kind,
827 },
828 )
829 })
830 .collect()
831}
832
833fn public_workflow_failure(failure: WorkflowFailure) -> PublicWorkflowFailure {
834 let message = match failure.code {
835 WorkflowFailureCode::InvalidDefinition => "Workflow definition is invalid",
836 WorkflowFailureCode::InvalidInput => "Workflow input is invalid",
837 WorkflowFailureCode::InvalidOutput => "Workflow output is invalid",
838 WorkflowFailureCode::UnknownReference => "Workflow reference is unavailable",
839 WorkflowFailureCode::PermissionDenied => "Workflow permission was denied",
840 WorkflowFailureCode::UntrustedWorkspace => "Workflow workspace is not trusted",
841 WorkflowFailureCode::BudgetExceeded => "Workflow execution budget was exceeded",
842 WorkflowFailureCode::RetryExhausted => "Workflow retry budget was exhausted",
843 WorkflowFailureCode::ExecutionFailed => "Workflow execution failed",
844 WorkflowFailureCode::Cancelled => "Workflow execution was cancelled",
845 WorkflowFailureCode::RecoverySuspended => "Workflow recovery requires attention",
846 WorkflowFailureCode::Suspended => "Workflow execution is suspended",
847 WorkflowFailureCode::DependencySkipped => "Workflow dependency was skipped",
848 WorkflowFailureCode::Storage => "Workflow storage is unavailable",
849 };
850 PublicWorkflowFailure {
851 code: failure.code,
852 message,
853 retryable: failure.retryable,
854 }
855}
856
857fn public_workflow_step(step: WorkflowStepSnapshot) -> PublicWorkflowStepSnapshot {
858 PublicWorkflowStepSnapshot {
859 id: step.id,
860 status: step.status,
861 failure: step.failure.map(public_workflow_failure),
862 attempts: step.attempts,
863 }
864}
865
866fn public_workflow_suspension(suspension: WorkflowSuspensionContext) -> PublicWorkflowSuspension {
867 match suspension {
868 WorkflowSuspensionContext::ToolApproval { step_id, .. } => {
869 PublicWorkflowSuspension::ToolApproval { step_id }
870 }
871 WorkflowSuspensionContext::ToolRunning {
872 step_id, killed, ..
873 } => PublicWorkflowSuspension::ToolRunning { step_id, killed },
874 WorkflowSuspensionContext::Recovery { .. } => PublicWorkflowSuspension::Recovery,
875 }
876}
877
878fn public_workflow_phase(name: &str) -> &'static str {
879 match name {
880 "retry_reserved" => "retry_reserved",
881 "step_reserved" => "step_reserved",
882 "agent_reserved" => "agent_reserved",
883 "agent_usage_recorded" => "agent_usage_recorded",
884 "suspension_context_persisted" => "suspension_context_persisted",
885 _ => "workflow_progressed",
886 }
887}
888
889pub(crate) fn public_workflow_snapshot(snapshot: WorkflowRunSnapshot) -> PublicWorkflowRunSnapshot {
890 let can_cancel = ensure_workflow_cancel_allowed(&snapshot).is_ok();
891 let can_restart_as_new_run = ensure_workflow_restart_as_new_run_allowed(&snapshot).is_ok();
892 let planned_steps = public_planned_steps(&snapshot.definition);
893 let plan = public_workflow_plan(&snapshot.definition.plan);
894 let child_agent_count = snapshot.usage.agents;
895 PublicWorkflowRunSnapshot {
896 run_id: snapshot.run_id,
897 parent_run_id: snapshot.parent_run_id,
898 parent_step_id: snapshot.parent_step_id,
899 session_id: snapshot.session_id,
900 workflow_id: snapshot.definition.id,
901 workflow_revision: snapshot.definition.revision,
902 definition_bundle_hash: snapshot.definition_bundle_hash,
903 status: snapshot.status,
904 can_cancel,
905 can_restart_as_new_run,
906 planned_steps,
907 plan,
908 steps: snapshot
909 .steps
910 .into_iter()
911 .map(|(id, step)| (id, public_workflow_step(step)))
912 .collect(),
913 budget: snapshot.definition.budgets,
914 usage: snapshot.usage,
915 child_agent_count,
916 last_sequence: snapshot.last_sequence,
917 failure: snapshot.failure.map(public_workflow_failure),
918 suspension: snapshot.suspension.map(public_workflow_suspension),
919 created_at: snapshot.created_at,
920 updated_at: snapshot.updated_at,
921 }
922}
923
924pub(crate) fn public_workflow_event(event: WorkflowRunEvent) -> PublicWorkflowRunEvent {
925 let kind = match event.kind {
926 WorkflowRunEventKind::RunQueued => PublicWorkflowRunEventKind::RunQueued,
927 WorkflowRunEventKind::RunStarted => PublicWorkflowRunEventKind::RunStarted,
928 WorkflowRunEventKind::Phase { name } => PublicWorkflowRunEventKind::Phase {
929 name: public_workflow_phase(&name),
930 },
931 WorkflowRunEventKind::StepQueued => PublicWorkflowRunEventKind::StepQueued,
932 WorkflowRunEventKind::StepStarted => PublicWorkflowRunEventKind::StepStarted,
933 WorkflowRunEventKind::StepSuspended { .. } => PublicWorkflowRunEventKind::StepSuspended,
934 WorkflowRunEventKind::StepCompleted { .. } => PublicWorkflowRunEventKind::StepCompleted,
935 WorkflowRunEventKind::StepFailed { failure } => PublicWorkflowRunEventKind::StepFailed {
936 failure: public_workflow_failure(failure),
937 },
938 WorkflowRunEventKind::StepCancelled => PublicWorkflowRunEventKind::StepCancelled,
939 WorkflowRunEventKind::StepSkipped { .. } => PublicWorkflowRunEventKind::StepSkipped,
940 WorkflowRunEventKind::RunSuspended { .. } => PublicWorkflowRunEventKind::RunSuspended,
941 WorkflowRunEventKind::RunSucceeded { .. } => PublicWorkflowRunEventKind::RunSucceeded,
942 WorkflowRunEventKind::RunFailed { failure } => PublicWorkflowRunEventKind::RunFailed {
943 failure: public_workflow_failure(failure),
944 },
945 WorkflowRunEventKind::RunCancelled => PublicWorkflowRunEventKind::RunCancelled,
946 };
947 PublicWorkflowRunEvent {
948 run_id: event.run_id,
949 sequence: event.sequence,
950 at: event.at,
951 step_id: event.step_id,
952 kind,
953 }
954}
955
956struct ExternallyPinnedDefinitions;
960
961#[async_trait]
962impl WorkflowDefinitionPort for ExternallyPinnedDefinitions {
963 async fn pin_bundle(
964 &self,
965 _root: &WorkflowRunDefinition,
966 ) -> Result<WorkflowDefinitionBundle, String> {
967 Err("server workflows must be pinned through SkillManager".to_string())
968 }
969}
970
971struct UnavailableAgentPort;
974
975#[async_trait]
976impl AgentStepPort for UnavailableAgentPort {
977 async fn resolve(&self, _name: &str) -> Result<Option<NamedAgentSpec>, String> {
978 Ok(None)
979 }
980
981 async fn execute(
982 &self,
983 _spec: &NamedAgentSpec,
984 _prompt: Value,
985 _model: Option<&str>,
986 _effort: Option<&str>,
987 _capabilities: &BTreeSet<String>,
988 _session_id: &str,
989 ) -> Result<AgentStepResult, String> {
990 Err("named-agent execution is not available".to_string())
991 }
992}
993
994struct ServerWorkflowPolicy;
995
996#[async_trait]
997impl WorkflowPolicyPort for ServerWorkflowPolicy {
998 async fn authorize(
999 &self,
1000 _session_id: &str,
1001 target: &WorkflowPolicyTarget,
1002 requested: &BTreeSet<String>,
1003 _workspace_trusted: bool,
1004 ) -> PermissionDecision {
1005 match target {
1011 WorkflowPolicyTarget::Workflow { .. } if requested.is_empty() => {
1015 PermissionDecision::Allow
1016 }
1017 WorkflowPolicyTarget::Tool(name) => {
1018 let safe_target = SAFE_UNTRUSTED_WORKFLOW_TOOLS
1019 .iter()
1020 .any(|candidate| candidate.eq_ignore_ascii_case(name));
1021 if safe_target && requested.iter().all(|capability| capability == "read") {
1022 PermissionDecision::Allow
1023 } else {
1024 PermissionDecision::Deny(
1025 "workflow capability authority is not available".to_string(),
1026 )
1027 }
1028 }
1029 WorkflowPolicyTarget::Agent(_) | WorkflowPolicyTarget::Workflow { .. } => {
1030 PermissionDecision::Deny(
1031 "workflow capability authority is not available".to_string(),
1032 )
1033 }
1034 }
1035 }
1036}
1037
1038struct UnavailableSecretResolver;
1039
1040#[async_trait]
1041impl WorkflowSecretResolverPort for UnavailableSecretResolver {
1042 async fn resolve(
1043 &self,
1044 _session_id: &str,
1045 _capability: &str,
1046 ) -> Result<WorkflowSecretMaterial, String> {
1047 Err("workflow secret capability resolver is not available".to_string())
1048 }
1049}
1050
1051#[derive(Debug, Deserialize)]
1052#[serde(tag = "action", rename_all = "snake_case", deny_unknown_fields)]
1053enum WorkflowToolInput {
1054 Start {
1055 workflow_id: String,
1056 revision: u64,
1057 #[serde(default = "empty_object")]
1058 args: Value,
1059 #[serde(default)]
1060 budget: Option<WorkflowBudgets>,
1061 },
1062 List {},
1063 Get {
1064 run_id: String,
1065 },
1066 Events {
1067 run_id: String,
1068 #[serde(default)]
1069 since: u64,
1070 },
1071 Cancel {
1072 run_id: String,
1073 },
1074 Restart {
1075 run_id: String,
1076 },
1077}
1078
1079fn empty_object() -> Value {
1080 json!({})
1081}
1082
1083pub struct WorkflowRunTool {
1084 access: WorkflowRunAccess,
1085}
1086
1087impl WorkflowRunTool {
1088 pub fn new(access: WorkflowRunAccess) -> Self {
1089 Self { access }
1090 }
1091}
1092
1093#[async_trait]
1094impl Tool for WorkflowRunTool {
1095 fn name(&self) -> &str {
1096 "workflow_run"
1097 }
1098
1099 fn description(&self) -> &str {
1100 "Start, inspect, cancel, or safely restart a catalog-pinned workflow run"
1101 }
1102
1103 fn parameters_schema(&self) -> Value {
1104 json!({
1105 "type": "object",
1106 "properties": {
1107 "action": {
1108 "type": "string",
1109 "enum": ["start", "list", "get", "events", "cancel", "restart"]
1110 },
1111 "workflow_id": {"type": "string", "minLength": 1},
1112 "revision": {"type": "integer", "minimum": 1},
1113 "args": {"type": "object", "default": {}},
1114 "budget": {
1115 "type": "object",
1116 "properties": {
1117 "max_concurrency": {"type": "integer", "minimum": 1},
1118 "max_agents": {"type": "integer", "minimum": 0},
1119 "max_steps": {"type": "integer", "minimum": 1},
1120 "max_retries": {"type": "integer", "minimum": 0},
1121 "max_nesting_depth": {"type": "integer", "minimum": 1},
1122 "wall_time_ms": {"type": "integer", "minimum": 1},
1123 "max_tokens": {"type": "integer", "minimum": 0},
1124 "max_cost_micros": {"type": "integer", "minimum": 0}
1125 },
1126 "required": [
1127 "max_concurrency",
1128 "max_agents",
1129 "max_steps",
1130 "max_retries",
1131 "max_nesting_depth",
1132 "wall_time_ms"
1133 ],
1134 "additionalProperties": false
1135 },
1136 "run_id": {"type": "string"},
1137 "since": {"type": "integer", "minimum": 0}
1138 },
1139 "required": ["action"],
1140 "additionalProperties": false
1141 })
1142 }
1143
1144 fn classify(&self, args: &Value) -> ToolClass {
1145 match args.get("action").and_then(Value::as_str) {
1146 Some("get" | "list" | "events") => ToolClass::READONLY_PARALLEL,
1147 _ => ToolClass::MUTATING_SERIAL,
1148 }
1149 }
1150
1151 async fn invoke(&self, args: Value, ctx: ToolCtx) -> Result<ToolOutcome, ToolError> {
1152 let input: WorkflowToolInput = serde_json::from_value(args)
1153 .map_err(|error| ToolError::InvalidArguments(error.to_string()))?;
1154 let session_id = ctx.session_id().ok_or_else(|| {
1155 ToolError::InvalidArguments("workflow_run requires a session".to_string())
1156 })?;
1157 let result = match input {
1158 WorkflowToolInput::Start {
1159 workflow_id,
1160 revision,
1161 args,
1162 budget,
1163 } => serde_json::to_value(public_workflow_snapshot(
1164 self.access
1165 .start_from_tool(session_id, &workflow_id, revision, args, budget)
1166 .await
1167 .map_err(workflow_tool_error)?,
1168 )),
1169 WorkflowToolInput::List {} => serde_json::to_value(
1170 self.access
1171 .list_for_session(session_id)
1172 .await
1173 .map_err(workflow_tool_error)?
1174 .into_iter()
1175 .map(public_workflow_snapshot)
1176 .collect::<Vec<_>>(),
1177 ),
1178 WorkflowToolInput::Get { run_id } => {
1179 let progress = self
1180 .access
1181 .progress_for_session(session_id, &run_id, u64::MAX)
1182 .await
1183 .map_err(workflow_tool_error)?;
1184 serde_json::to_value(public_workflow_snapshot(progress.snapshot))
1185 }
1186 WorkflowToolInput::Events { run_id, since } => {
1187 let progress = self
1188 .access
1189 .progress_for_session(session_id, &run_id, since)
1190 .await
1191 .map_err(workflow_tool_error)?;
1192 serde_json::to_value(
1193 progress
1194 .events
1195 .into_iter()
1196 .map(public_workflow_event)
1197 .collect::<Vec<_>>(),
1198 )
1199 }
1200 WorkflowToolInput::Cancel { run_id } => serde_json::to_value(public_workflow_snapshot(
1201 self.access
1202 .cancel_for_session(session_id, &run_id)
1203 .await
1204 .map_err(workflow_tool_error)?,
1205 )),
1206 WorkflowToolInput::Restart { run_id } => {
1207 serde_json::to_value(public_workflow_snapshot(
1208 self.access
1209 .restart_from_tool(session_id, &run_id)
1210 .await
1211 .map_err(workflow_tool_error)?,
1212 ))
1213 }
1214 }
1215 .map_err(|error| ToolError::Execution(error.to_string()))?;
1216 Ok(ToolOutcome::Completed(ToolResult::text(
1217 true,
1218 serde_json::to_string(&result)
1219 .map_err(|error| ToolError::Execution(error.to_string()))?,
1220 )))
1221 }
1222}
1223
1224fn workflow_tool_error(error: WorkflowRunError) -> ToolError {
1225 match error {
1226 WorkflowRunError::InvalidInput(_) => {
1227 ToolError::InvalidArguments("workflow input is invalid".to_string())
1228 }
1229 WorkflowRunError::Compile(_) => {
1230 ToolError::InvalidArguments("workflow definition is invalid".to_string())
1231 }
1232 WorkflowRunError::Preflight(_) => {
1233 ToolError::InvalidArguments("workflow preflight failed".to_string())
1234 }
1235 WorkflowRunError::Storage(details) => {
1236 tracing::error!(%details, "workflow tool storage unavailable");
1237 ToolError::Execution("workflow storage unavailable".to_string())
1238 }
1239 WorkflowRunError::NotFound => ToolError::Execution("workflow run not found".to_string()),
1240 WorkflowRunError::Terminal => {
1241 ToolError::Execution("workflow run is already terminal".to_string())
1242 }
1243 }
1244}
1245
1246#[cfg(test)]
1247mod tests {
1248 use super::*;
1249 use bamboo_agent_core::storage::Storage;
1250 use bamboo_agent_core::tools::{
1251 FunctionSchema, ToolCall, ToolExecutor, ToolResult, ToolSchema,
1252 };
1253 use bamboo_agent_core::Session;
1254 use bamboo_engine::WorkflowRunRepository;
1255 use bamboo_llm::protocol::{gemini::GeminiTool, ToProvider};
1256 use std::collections::HashMap;
1257 use tokio::sync::{RwLock, Semaphore};
1258
1259 #[derive(Default)]
1260 struct WorkflowTestStorage {
1261 sessions: RwLock<HashMap<String, Session>>,
1262 fail_saves: AtomicBool,
1263 }
1264
1265 impl WorkflowTestStorage {
1266 fn set_fail_saves(&self, fail: bool) {
1267 self.fail_saves.store(fail, Ordering::SeqCst);
1268 }
1269 }
1270
1271 #[async_trait]
1272 impl Storage for WorkflowTestStorage {
1273 async fn save_session(&self, session: &Session) -> std::io::Result<()> {
1274 if self.fail_saves.load(Ordering::SeqCst) {
1275 return Err(std::io::Error::other(
1276 "injected workflow session persistence failure",
1277 ));
1278 }
1279 self.sessions
1280 .write()
1281 .await
1282 .insert(session.id.clone(), session.clone());
1283 Ok(())
1284 }
1285
1286 async fn load_session(&self, session_id: &str) -> std::io::Result<Option<Session>> {
1287 Ok(self.sessions.read().await.get(session_id).cloned())
1288 }
1289
1290 async fn delete_session(&self, session_id: &str) -> std::io::Result<bool> {
1291 Ok(self.sessions.write().await.remove(session_id).is_some())
1292 }
1293 }
1294
1295 struct WorkflowReadTool;
1296
1297 #[async_trait]
1298 impl ToolExecutor for WorkflowReadTool {
1299 async fn execute(
1300 &self,
1301 _call: &ToolCall,
1302 ) -> bamboo_agent_core::tools::executor::Result<ToolResult> {
1303 Ok(ToolResult::text(true, r#"{"ok":true}"#))
1304 }
1305
1306 fn list_tools(&self) -> Vec<ToolSchema> {
1307 vec![ToolSchema {
1308 schema_type: "function".to_string(),
1309 function: FunctionSchema {
1310 name: "Read".to_string(),
1311 description: "read".to_string(),
1312 parameters: serde_json::json!({"type":"object"}),
1313 },
1314 }]
1315 }
1316 }
1317
1318 struct BlockingWorkflowReadTool {
1319 entered: Arc<Semaphore>,
1320 }
1321
1322 #[async_trait]
1323 impl ToolExecutor for BlockingWorkflowReadTool {
1324 async fn execute(
1325 &self,
1326 _call: &ToolCall,
1327 ) -> bamboo_agent_core::tools::executor::Result<ToolResult> {
1328 self.entered.add_permits(1);
1329 std::future::pending().await
1330 }
1331
1332 fn list_tools(&self) -> Vec<ToolSchema> {
1333 WorkflowReadTool.list_tools()
1334 }
1335 }
1336
1337 async fn workflow_test_access_with_tools(
1338 tools: Arc<dyn ToolExecutor>,
1339 ) -> (
1340 WorkflowRunAccess,
1341 bamboo_engine::SessionRepository,
1342 tempfile::TempDir,
1343 ) {
1344 let (access, repo, directory, _) = workflow_test_access_with_tools_and_storage(tools).await;
1345 (access, repo, directory)
1346 }
1347
1348 async fn workflow_test_access_with_tools_and_storage(
1349 tools: Arc<dyn ToolExecutor>,
1350 ) -> (
1351 WorkflowRunAccess,
1352 bamboo_engine::SessionRepository,
1353 tempfile::TempDir,
1354 Arc<WorkflowTestStorage>,
1355 ) {
1356 let directory = tempfile::tempdir().expect("tempdir");
1357 let skills_dir = directory.path().join("skills");
1358 let root = skills_dir.join("review-flow");
1359 std::fs::create_dir_all(&root).expect("workflow dir");
1360 std::fs::write(
1361 root.join("SKILL.md"),
1362 "---\nname: review-flow\ndescription: Review flow\n---\nRun review flow.\n",
1363 )
1364 .expect("skill");
1365 std::fs::write(
1366 root.join("workflow.yaml"),
1367 "workflow_schema: 1\nid: review-flow\nrevision: 42\ninvocation_policy: {explicit: true, automatic: true}\ninput_schema:\n type: object\n additionalProperties: true\nsteps:\n - id: inspect\n type: tool\n tool: Read\n args: {}\n capabilities: [read]\n output_schema:\n type: object\n additionalProperties: true\nplan:\n type: step\n step: inspect\nbudgets:\n max_concurrency: 2\n max_agents: 1\n max_steps: 4\n max_retries: 2\n max_nesting_depth: 2\n wall_time_ms: 10000\n max_tokens: 1000\n max_cost_micros: 1000\n",
1368 )
1369 .expect("workflow");
1370 let skills = Arc::new(SkillManager::with_config(bamboo_skills::SkillStoreConfig {
1371 skills_dir,
1372 ..Default::default()
1373 }));
1374 skills.initialize().await.expect("skills initialize");
1375 let storage = Arc::new(WorkflowTestStorage::default());
1376 let storage_port: Arc<dyn Storage> = storage.clone();
1377 let persistence = Arc::new(bamboo_storage::LockedSessionStore::new(
1378 storage_port.clone(),
1379 ));
1380 let cache = Arc::default();
1381 let repo = bamboo_engine::SessionRepository::new(cache, storage_port, persistence);
1382 let access = WorkflowRunAccess::new(directory.path(), tools, skills, repo.clone())
1383 .await
1384 .expect("workflow access");
1385 (access, repo, directory, storage)
1386 }
1387
1388 async fn workflow_test_access() -> (
1389 WorkflowRunAccess,
1390 bamboo_engine::SessionRepository,
1391 tempfile::TempDir,
1392 ) {
1393 workflow_test_access_with_tools(Arc::new(WorkflowReadTool)).await
1394 }
1395
1396 async fn seed_durable_workflow_run(
1397 access: &WorkflowRunAccess,
1398 directory: &Path,
1399 workspace: &Path,
1400 session_id: &str,
1401 run_id: &str,
1402 status: WorkflowRunStatus,
1403 suspension: Option<WorkflowSuspensionContext>,
1404 ) -> WorkflowRunSnapshot {
1405 let bundle = access
1406 .skills
1407 .pin_workflow_definition_bundle(Some(workspace), "review-flow", 42)
1408 .await
1409 .expect("pin workflow test bundle");
1410 let definition = bundle.root().cloned().expect("pinned root definition");
1411 let step_status = match status {
1412 WorkflowRunStatus::Queued => WorkflowStepStatus::Queued,
1413 WorkflowRunStatus::Running => WorkflowStepStatus::Running,
1414 WorkflowRunStatus::Suspended => WorkflowStepStatus::Suspended,
1415 WorkflowRunStatus::Succeeded => WorkflowStepStatus::Succeeded,
1416 WorkflowRunStatus::Failed => WorkflowStepStatus::Failed,
1417 WorkflowRunStatus::Cancelled => WorkflowStepStatus::Cancelled,
1418 };
1419 let now = chrono::Utc::now();
1420 let failure = (status == WorkflowRunStatus::Failed).then(|| WorkflowFailure {
1421 code: WorkflowFailureCode::ExecutionFailed,
1422 message: "seeded workflow failure".to_string(),
1423 retryable: false,
1424 });
1425 let snapshot = WorkflowRunSnapshot {
1426 run_id: run_id.to_string(),
1427 parent_run_id: None,
1428 parent_step_id: None,
1429 session_id: session_id.to_string(),
1430 definition: definition.clone(),
1431 definition_bundle: bundle,
1432 definition_bundle_hash: "seeded-public-bundle-hash".to_string(),
1433 validated_args: json!({}),
1434 status,
1435 steps: definition
1436 .steps
1437 .iter()
1438 .map(|step| {
1439 (
1440 step.id.clone(),
1441 WorkflowStepSnapshot {
1442 id: step.id.clone(),
1443 status: step_status,
1444 input_hash: "seeded-input-hash".to_string(),
1445 output: None,
1446 failure: failure.clone(),
1447 attempts: 0,
1448 },
1449 )
1450 })
1451 .collect(),
1452 usage: WorkflowBudgetUsage::default(),
1453 last_sequence: 1,
1454 output: None,
1455 failure: failure.clone(),
1456 suspension,
1457 created_at: now,
1458 updated_at: now,
1459 };
1460 persist_durable_workflow_run(directory, &snapshot).await;
1461 snapshot
1462 }
1463
1464 async fn persist_durable_workflow_run(directory: &Path, snapshot: &WorkflowRunSnapshot) {
1465 let kind = match snapshot.status {
1466 WorkflowRunStatus::Queued => WorkflowRunEventKind::RunQueued,
1467 WorkflowRunStatus::Running => WorkflowRunEventKind::RunStarted,
1468 WorkflowRunStatus::Suspended => WorkflowRunEventKind::RunSuspended {
1469 reason: "seeded suspension".to_string(),
1470 },
1471 WorkflowRunStatus::Succeeded => WorkflowRunEventKind::RunSucceeded {
1472 output: Value::Null,
1473 },
1474 WorkflowRunStatus::Failed => WorkflowRunEventKind::RunFailed {
1475 failure: snapshot.failure.clone().expect("failed run has failure"),
1476 },
1477 WorkflowRunStatus::Cancelled => WorkflowRunEventKind::RunCancelled,
1478 };
1479 let repository = FileWorkflowRunRepository::new(directory.join("workflow-runs"))
1480 .expect("open workflow test repository");
1481 repository
1482 .create(
1483 snapshot,
1484 &WorkflowRunEvent {
1485 run_id: snapshot.run_id.clone(),
1486 sequence: snapshot.last_sequence,
1487 at: snapshot.updated_at,
1488 step_id: None,
1489 kind,
1490 },
1491 )
1492 .await
1493 .expect("seed durable workflow run");
1494 }
1495
1496 async fn replace_workflow_run_index(
1497 repo: &bamboo_engine::SessionRepository,
1498 session_id: &str,
1499 run_ids: Vec<String>,
1500 ) {
1501 repo.update_runtime_session(
1502 session_id,
1503 &[bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY],
1504 move |session| {
1505 session.metadata.insert(
1506 bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY.to_string(),
1507 serde_json::to_string(&run_ids).expect("run index json"),
1508 );
1509 },
1510 )
1511 .await
1512 .expect("replace workflow run index")
1513 .expect("workflow session");
1514 }
1515
1516 async fn wait_for_workflow_run_to_settle(
1517 access: &WorkflowRunAccess,
1518 run_id: &str,
1519 ) -> WorkflowRunSnapshot {
1520 tokio::time::timeout(std::time::Duration::from_secs(2), async {
1521 loop {
1522 let snapshot = access
1523 .engine
1524 .progress(run_id, u64::MAX)
1525 .await
1526 .expect("workflow progress")
1527 .snapshot;
1528 if snapshot.status.is_terminal() && !access.engine.is_run_active(run_id) {
1529 break snapshot;
1530 }
1531 tokio::task::yield_now().await;
1532 }
1533 })
1534 .await
1535 .expect("workflow run settles")
1536 }
1537
1538 fn private_workflow_snapshot(status: WorkflowRunStatus) -> WorkflowRunSnapshot {
1539 let definition = WorkflowRunDefinition {
1540 workflow_schema: 1,
1541 id: "review-flow".to_string(),
1542 revision: 42,
1543 input_schema: json!({"private_schema": "PRIVATE-SCHEMA-SENTINEL"}),
1544 output_schema: Some(json!({"private_output": "PRIVATE-SCHEMA-SENTINEL"})),
1545 steps: vec![bamboo_domain::WorkflowStepDefinition {
1546 id: "inspect".to_string(),
1547 kind: WorkflowStepKind::Tool {
1548 tool: "PRIVATE-TOOL-SENTINEL".to_string(),
1549 args: json!({"credential": "PRIVATE-ARG-SENTINEL"}),
1550 capabilities: vec!["PRIVATE-CAPABILITY-SENTINEL".to_string()],
1551 },
1552 failure: bamboo_domain::FailurePolicy::FailFast,
1553 output_schema: Some(json!({"private": "PRIVATE-STEP-SCHEMA-SENTINEL"})),
1554 }],
1555 plan: WorkflowPlan::Retry {
1556 node: Box::new(WorkflowPlan::Map {
1557 source: bamboo_domain::ValueRef::Literal {
1558 value: json!("PRIVATE-BINDING-SENTINEL"),
1559 },
1560 item: "PRIVATE-ITEM-SENTINEL".to_string(),
1561 body: Box::new(WorkflowPlan::Step {
1562 step: "inspect".to_string(),
1563 }),
1564 }),
1565 max_attempts: 3,
1566 delay_ms: 987_654,
1567 },
1568 budgets: WorkflowBudgets {
1569 max_concurrency: 2,
1570 max_agents: 4,
1571 max_steps: 8,
1572 max_retries: 3,
1573 max_nesting_depth: 2,
1574 wall_time_ms: 10_000,
1575 max_tokens: Some(1_000),
1576 max_cost_micros: Some(2_000),
1577 },
1578 };
1579 let definition_bundle = WorkflowDefinitionBundle {
1580 publication_revision: 7,
1581 root_id: definition.id.clone(),
1582 root_revision: definition.revision,
1583 root_invocation_policy: json!({"private": "PRIVATE-POLICY-SENTINEL"}),
1584 definitions: BTreeMap::from([(
1585 WorkflowDefinitionBundle::key(&definition.id, definition.revision),
1586 definition.clone(),
1587 )]),
1588 };
1589 let now = chrono::Utc::now();
1590 WorkflowRunSnapshot {
1591 run_id: "public-run".to_string(),
1592 parent_run_id: Some("public-parent-run".to_string()),
1593 parent_step_id: Some("public-parent-step".to_string()),
1594 session_id: "public-session".to_string(),
1595 definition,
1596 definition_bundle,
1597 definition_bundle_hash: "public-bundle-hash".to_string(),
1598 validated_args: json!({"password": "PRIVATE-VALIDATED-ARG-SENTINEL"}),
1599 status,
1600 steps: BTreeMap::from([(
1601 "inspect".to_string(),
1602 WorkflowStepSnapshot {
1603 id: "inspect".to_string(),
1604 status: WorkflowStepStatus::Failed,
1605 input_hash: "PRIVATE-INPUT-HASH-SENTINEL".to_string(),
1606 output: Some(json!({"raw_tool_output": "PRIVATE-OUTPUT-SENTINEL"})),
1607 failure: Some(WorkflowFailure {
1608 code: WorkflowFailureCode::Storage,
1609 message: "/private/workspace/PRIVATE-DIAGNOSTIC-SENTINEL".to_string(),
1610 retryable: true,
1611 }),
1612 attempts: 2,
1613 },
1614 )]),
1615 usage: WorkflowBudgetUsage {
1616 steps: 1,
1617 retries: 1,
1618 agents: 3,
1619 tokens: 40,
1620 cost_micros: 50,
1621 },
1622 last_sequence: 9,
1623 output: Some(json!({"raw_run_output": "PRIVATE-RUN-OUTPUT-SENTINEL"})),
1624 failure: Some(WorkflowFailure {
1625 code: WorkflowFailureCode::ExecutionFailed,
1626 message: "credential PRIVATE-RUN-FAILURE-SENTINEL".to_string(),
1627 retryable: false,
1628 }),
1629 suspension: Some(WorkflowSuspensionContext::ToolApproval {
1630 step_id: "inspect".to_string(),
1631 tool: "PRIVATE-SUSPENSION-TOOL-SENTINEL".to_string(),
1632 tool_call_id: "PRIVATE-TOOL-CALL-SENTINEL".to_string(),
1633 }),
1634 created_at: now,
1635 updated_at: now,
1636 }
1637 }
1638
1639 #[test]
1640 fn public_workflow_snapshot_is_stable_metadata_only_for_every_status() {
1641 let statuses = [
1642 (WorkflowRunStatus::Queued, "queued", true),
1643 (WorkflowRunStatus::Running, "running", true),
1644 (WorkflowRunStatus::Suspended, "suspended", true),
1645 (WorkflowRunStatus::Succeeded, "succeeded", false),
1646 (WorkflowRunStatus::Failed, "failed", false),
1647 (WorkflowRunStatus::Cancelled, "cancelled", true),
1648 ];
1649
1650 for (status, wire_status, can_cancel) in statuses {
1651 let public =
1652 serde_json::to_value(public_workflow_snapshot(private_workflow_snapshot(status)))
1653 .expect("public snapshot serializes");
1654 let text = public.to_string();
1655 assert_eq!(public["status"], wire_status);
1656 assert_eq!(public["can_cancel"], can_cancel);
1657 assert_eq!(public["can_restart_as_new_run"], false);
1658 assert_eq!(public["workflow_id"], "review-flow");
1659 assert_eq!(public["workflow_revision"], 42);
1660 assert_eq!(public["definition_bundle_hash"], "public-bundle-hash");
1661 assert_eq!(public["planned_steps"]["inspect"]["kind"], "tool");
1662 assert_eq!(public["plan"]["type"], "retry");
1663 assert_eq!(public["steps"]["inspect"]["attempts"], 2);
1664 assert_eq!(public["budget"]["max_steps"], 8);
1665 assert_eq!(public["usage"]["agents"], 3);
1666 assert_eq!(public["child_agent_count"], 3);
1667 assert_eq!(public["last_sequence"], 9);
1668 assert_eq!(public["failure"]["message"], "Workflow execution failed");
1669 assert_eq!(public["suspension"]["type"], "tool_approval");
1670 for internal_field in [
1671 "definition",
1672 "definition_bundle",
1673 "validated_args",
1674 "output",
1675 ] {
1676 assert!(
1677 public.get(internal_field).is_none(),
1678 "public snapshot exposed internal field {internal_field}: {text}"
1679 );
1680 }
1681
1682 for private in [
1683 "PRIVATE-",
1684 "validated_args",
1685 "input_hash",
1686 "output_schema",
1687 "root_invocation_policy",
1688 "delay_ms",
1689 "tool_call_id",
1690 "raw_tool_output",
1691 "raw_run_output",
1692 ] {
1693 assert!(
1694 !text.contains(private),
1695 "public snapshot leaked {private}: {text}"
1696 );
1697 }
1698 }
1699
1700 let mut recovery = private_workflow_snapshot(WorkflowRunStatus::Suspended);
1701 recovery.suspension = Some(WorkflowSuspensionContext::Recovery {
1702 reason: "PRIVATE-RECOVERY-REASON-SENTINEL".to_string(),
1703 });
1704 let public = serde_json::to_value(public_workflow_snapshot(recovery))
1705 .expect("public recovery snapshot serializes");
1706 let text = public.to_string();
1707 assert_eq!(public["can_cancel"], true);
1708 assert_eq!(public["can_restart_as_new_run"], true);
1709 assert_eq!(public["suspension"]["type"], "recovery");
1710 assert!(!text.contains("PRIVATE-RECOVERY-REASON-SENTINEL"));
1711 }
1712
1713 #[tokio::test]
1714 async fn public_workflow_actions_match_cancel_and_restart_endpoint_acceptance() {
1715 let (access, repo, directory) = workflow_test_access().await;
1716 let workspace = directory.path().join("action-workspace");
1717 std::fs::create_dir_all(&workspace).expect("action workspace");
1718 let session_id = "workflow-action-matrix";
1719 let mut session = Session::new(session_id, "model");
1720 session.workspace = Some(workspace.to_string_lossy().into_owned());
1721 repo.save(&mut session).await.expect("save action session");
1722
1723 let cases = vec![
1724 ("queued", WorkflowRunStatus::Queued, None, true, false),
1725 ("running", WorkflowRunStatus::Running, None, true, false),
1726 (
1727 "suspended_without_context",
1728 WorkflowRunStatus::Suspended,
1729 None,
1730 true,
1731 false,
1732 ),
1733 (
1734 "tool_approval",
1735 WorkflowRunStatus::Suspended,
1736 Some(WorkflowSuspensionContext::ToolApproval {
1737 step_id: "inspect".to_string(),
1738 tool: "Read".to_string(),
1739 tool_call_id: "approval-call".to_string(),
1740 }),
1741 true,
1742 false,
1743 ),
1744 (
1745 "tool_running",
1746 WorkflowRunStatus::Suspended,
1747 Some(WorkflowSuspensionContext::ToolRunning {
1748 step_id: "inspect".to_string(),
1749 tool: "Read".to_string(),
1750 tool_call_id: "running-call".to_string(),
1751 killed: true,
1752 }),
1753 true,
1754 false,
1755 ),
1756 (
1757 "recovery",
1758 WorkflowRunStatus::Suspended,
1759 Some(WorkflowSuspensionContext::Recovery {
1760 reason: "process restarted".to_string(),
1761 }),
1762 true,
1763 true,
1764 ),
1765 (
1766 "succeeded",
1767 WorkflowRunStatus::Succeeded,
1768 None,
1769 false,
1770 false,
1771 ),
1772 ("failed", WorkflowRunStatus::Failed, None, false, false),
1773 ("cancelled", WorkflowRunStatus::Cancelled, None, true, false),
1774 ];
1775
1776 for (name, status, suspension, can_cancel, can_restart_as_new_run) in cases {
1777 let cancel_run_id = format!("action-cancel-{name}");
1778 let cancel_snapshot = seed_durable_workflow_run(
1779 &access,
1780 directory.path(),
1781 &workspace,
1782 session_id,
1783 &cancel_run_id,
1784 status,
1785 suspension.clone(),
1786 )
1787 .await;
1788 let cancel_public = public_workflow_snapshot(cancel_snapshot);
1789 assert_eq!(
1790 cancel_public.can_cancel, can_cancel,
1791 "cancel projection mismatch for {name}"
1792 );
1793 assert_eq!(
1794 access
1795 .cancel_for_session(session_id, &cancel_run_id)
1796 .await
1797 .is_ok(),
1798 can_cancel,
1799 "cancel endpoint mismatch for {name}"
1800 );
1801
1802 let restart_run_id = format!("action-restart-{name}");
1803 let restart_snapshot = seed_durable_workflow_run(
1804 &access,
1805 directory.path(),
1806 &workspace,
1807 session_id,
1808 &restart_run_id,
1809 status,
1810 suspension,
1811 )
1812 .await;
1813 let restart_public = public_workflow_snapshot(restart_snapshot);
1814 assert_eq!(
1815 restart_public.can_restart_as_new_run, can_restart_as_new_run,
1816 "restart projection mismatch for {name}"
1817 );
1818 let restarted = access
1819 .restart_for_session(session_id, &restart_run_id)
1820 .await;
1821 assert_eq!(
1822 restarted.is_ok(),
1823 can_restart_as_new_run,
1824 "restart endpoint mismatch for {name}: {restarted:?}"
1825 );
1826 if let Ok(restarted) = restarted {
1827 let _ = access
1828 .cancel_for_session(session_id, &restarted.run_id)
1829 .await;
1830 }
1831 }
1832 }
1833
1834 #[tokio::test]
1835 async fn restart_as_new_run_is_indexed_isolated_and_survives_reconstruction() {
1836 let (access, repo, directory, storage) =
1837 workflow_test_access_with_tools_and_storage(Arc::new(WorkflowReadTool)).await;
1838 let workspace = directory.path().join("restart-workspace");
1839 std::fs::create_dir_all(&workspace).expect("restart workspace");
1840 let session_id = "restart-owner";
1841 let mut session = Session::new(session_id, "model");
1842 session.workspace = Some(workspace.to_string_lossy().into_owned());
1843 repo.save(&mut session).await.expect("save restart owner");
1844 let mut other = Session::new("restart-other", "model");
1845 other.workspace = Some(workspace.to_string_lossy().into_owned());
1846 repo.save(&mut other).await.expect("save other session");
1847
1848 let original = seed_durable_workflow_run(
1849 &access,
1850 directory.path(),
1851 &workspace,
1852 session_id,
1853 "recovery-original",
1854 WorkflowRunStatus::Suspended,
1855 Some(WorkflowSuspensionContext::Recovery {
1856 reason: "process restarted".to_string(),
1857 }),
1858 )
1859 .await;
1860 replace_workflow_run_index(&repo, session_id, vec![original.run_id.clone()]).await;
1861
1862 let restarted = access
1863 .restart_for_session(session_id, &original.run_id)
1864 .await
1865 .expect("restart recovery suspension as a new run");
1866 assert_ne!(restarted.run_id, original.run_id);
1867 assert_eq!(restarted.session_id, session_id);
1868 assert_eq!(
1869 access
1870 .progress_for_session(session_id, &original.run_id, u64::MAX)
1871 .await
1872 .expect("original progress")
1873 .snapshot,
1874 original,
1875 "restart-as-new must not mutate the original suspended run"
1876 );
1877
1878 let immediate = access
1879 .list_for_session(session_id)
1880 .await
1881 .expect("immediate owner list");
1882 let immediate_ids = immediate
1883 .iter()
1884 .map(|snapshot| snapshot.run_id.as_str())
1885 .collect::<BTreeSet<_>>();
1886 assert_eq!(
1887 immediate_ids,
1888 BTreeSet::from([original.run_id.as_str(), restarted.run_id.as_str()])
1889 );
1890 assert!(access
1891 .list_for_session("restart-other")
1892 .await
1893 .expect("isolated other list")
1894 .is_empty());
1895 assert!(matches!(
1896 access
1897 .progress_for_session("restart-other", &restarted.run_id, u64::MAX)
1898 .await,
1899 Err(WorkflowRunError::NotFound)
1900 ));
1901 assert!(matches!(
1902 access
1903 .restart_for_session("restart-other", &original.run_id)
1904 .await,
1905 Err(WorkflowRunError::NotFound)
1906 ));
1907
1908 let settled = wait_for_workflow_run_to_settle(&access, &restarted.run_id).await;
1909 assert_eq!(settled.status, WorkflowRunStatus::Succeeded);
1910 let skills = access.skills.clone();
1911 drop(access);
1912 drop(repo);
1913
1914 let storage_port: Arc<dyn Storage> = storage;
1915 let reopened_repo = bamboo_engine::SessionRepository::new(
1916 Arc::default(),
1917 storage_port.clone(),
1918 Arc::new(bamboo_storage::LockedSessionStore::new(storage_port)),
1919 );
1920 let reopened = WorkflowRunAccess::new(
1921 directory.path(),
1922 Arc::new(WorkflowReadTool),
1923 skills,
1924 reopened_repo,
1925 )
1926 .await
1927 .expect("reconstruct workflow access");
1928 let reconstructed = reopened
1929 .list_for_session(session_id)
1930 .await
1931 .expect("list after reconstruction");
1932 let reconstructed_ids = reconstructed
1933 .iter()
1934 .map(|snapshot| snapshot.run_id.as_str())
1935 .collect::<BTreeSet<_>>();
1936 assert_eq!(
1937 reconstructed_ids,
1938 BTreeSet::from([original.run_id.as_str(), restarted.run_id.as_str()])
1939 );
1940 assert_eq!(
1941 reopened
1942 .progress_for_session(session_id, &original.run_id, u64::MAX)
1943 .await
1944 .expect("reconstructed original")
1945 .snapshot,
1946 original
1947 );
1948 }
1949
1950 #[tokio::test]
1951 async fn restart_capacity_preflight_creates_no_new_run_or_active_orphan() {
1952 let (access, repo, directory) = workflow_test_access().await;
1953 let workspace = directory.path().join("restart-capacity-workspace");
1954 std::fs::create_dir_all(&workspace).expect("restart capacity workspace");
1955 let session_id = "restart-capacity";
1956 let mut session = Session::new(session_id, "model");
1957 session.workspace = Some(workspace.to_string_lossy().into_owned());
1958 repo.save(&mut session)
1959 .await
1960 .expect("save restart capacity session");
1961
1962 let original = seed_durable_workflow_run(
1963 &access,
1964 directory.path(),
1965 &workspace,
1966 session_id,
1967 "capacity-run-0",
1968 WorkflowRunStatus::Suspended,
1969 Some(WorkflowSuspensionContext::Recovery {
1970 reason: "process restarted".to_string(),
1971 }),
1972 )
1973 .await;
1974 let mut run_ids = vec![original.run_id.clone()];
1975 for index in 1..MAX_WORKFLOW_RUN_IDS_PER_SESSION {
1976 let mut snapshot = original.clone();
1977 snapshot.run_id = format!("capacity-run-{index}");
1978 snapshot.created_at = chrono::Utc::now();
1979 snapshot.updated_at = snapshot.created_at;
1980 persist_durable_workflow_run(directory.path(), &snapshot).await;
1981 run_ids.push(snapshot.run_id);
1982 }
1983 replace_workflow_run_index(&repo, session_id, run_ids).await;
1984 let before = access
1985 .engine
1986 .list_run_ids()
1987 .await
1988 .expect("run ids before capacity rejection")
1989 .into_iter()
1990 .collect::<BTreeSet<_>>();
1991
1992 assert!(matches!(
1993 access
1994 .restart_for_session(session_id, &original.run_id)
1995 .await,
1996 Err(WorkflowRunError::Preflight(message))
1997 if message == "workflow run index is full of active runs"
1998 ));
1999 let after = access
2000 .engine
2001 .list_run_ids()
2002 .await
2003 .expect("run ids after capacity rejection")
2004 .into_iter()
2005 .collect::<BTreeSet<_>>();
2006 assert_eq!(
2007 after, before,
2008 "capacity rejection must precede run creation"
2009 );
2010 assert!(after
2011 .iter()
2012 .all(|run_id| !access.engine.is_run_active(run_id)));
2013 assert_eq!(
2014 access
2015 .progress_for_session(session_id, &original.run_id, u64::MAX)
2016 .await
2017 .expect("original after capacity rejection")
2018 .snapshot,
2019 original
2020 );
2021 }
2022
2023 #[tokio::test]
2024 async fn restart_index_persistence_failure_cancels_the_unindexed_new_run() {
2025 let (access, repo, directory, storage) =
2026 workflow_test_access_with_tools_and_storage(Arc::new(BlockingWorkflowReadTool {
2027 entered: Arc::new(Semaphore::new(0)),
2028 }))
2029 .await;
2030 let workspace = directory.path().join("restart-failure-workspace");
2031 std::fs::create_dir_all(&workspace).expect("restart failure workspace");
2032 let session_id = "restart-index-failure";
2033 let mut session = Session::new(session_id, "model");
2034 session.workspace = Some(workspace.to_string_lossy().into_owned());
2035 repo.save(&mut session)
2036 .await
2037 .expect("save restart failure session");
2038 let original = seed_durable_workflow_run(
2039 &access,
2040 directory.path(),
2041 &workspace,
2042 session_id,
2043 "failure-original",
2044 WorkflowRunStatus::Suspended,
2045 Some(WorkflowSuspensionContext::Recovery {
2046 reason: "process restarted".to_string(),
2047 }),
2048 )
2049 .await;
2050 replace_workflow_run_index(&repo, session_id, vec![original.run_id.clone()]).await;
2051 let before = access
2052 .engine
2053 .list_run_ids()
2054 .await
2055 .expect("run ids before injected failure")
2056 .into_iter()
2057 .collect::<BTreeSet<_>>();
2058
2059 storage.set_fail_saves(true);
2060 let error = access
2061 .restart_for_session(session_id, &original.run_id)
2062 .await
2063 .expect_err("injected index persistence failure");
2064 storage.set_fail_saves(false);
2065 assert!(matches!(error, WorkflowRunError::Storage(_)));
2066
2067 let after = access
2068 .engine
2069 .list_run_ids()
2070 .await
2071 .expect("run ids after injected failure")
2072 .into_iter()
2073 .collect::<BTreeSet<_>>();
2074 let created = after.difference(&before).cloned().collect::<Vec<_>>();
2075 assert_eq!(created.len(), 1, "restart should have created one new run");
2076 let unindexed_run_id = &created[0];
2077 let compensated = wait_for_workflow_run_to_settle(&access, unindexed_run_id).await;
2078 assert_eq!(compensated.status, WorkflowRunStatus::Cancelled);
2079 assert!(!access.engine.is_run_active(unindexed_run_id));
2080 assert_eq!(
2081 access
2082 .progress_for_session(session_id, &original.run_id, u64::MAX)
2083 .await
2084 .expect("original after index failure")
2085 .snapshot,
2086 original
2087 );
2088 let listed = access
2089 .list_for_session(session_id)
2090 .await
2091 .expect("owner list after index failure");
2092 assert_eq!(listed.len(), 1);
2093 assert_eq!(listed[0].run_id, original.run_id);
2094 let durable = storage
2095 .load_session(session_id)
2096 .await
2097 .expect("load durable owner")
2098 .expect("durable owner");
2099 let durable_ids = durable
2100 .metadata
2101 .get(bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY)
2102 .and_then(|raw| serde_json::from_str::<Vec<String>>(raw).ok())
2103 .expect("durable run index");
2104 assert_eq!(durable_ids, vec![original.run_id]);
2105 assert!(!durable_ids.contains(unindexed_run_id));
2106 }
2107
2108 #[test]
2109 fn public_workflow_events_preserve_sequence_and_drop_private_payloads() {
2110 let private = "PRIVATE-EVENT-SENTINEL";
2111 let failure = WorkflowFailure {
2112 code: WorkflowFailureCode::ExecutionFailed,
2113 message: format!("/private/workspace/{private}"),
2114 retryable: true,
2115 };
2116 let kinds = vec![
2117 WorkflowRunEventKind::RunQueued,
2118 WorkflowRunEventKind::RunStarted,
2119 WorkflowRunEventKind::Phase {
2120 name: private.to_string(),
2121 },
2122 WorkflowRunEventKind::StepQueued,
2123 WorkflowRunEventKind::StepStarted,
2124 WorkflowRunEventKind::StepSuspended {
2125 reason: private.to_string(),
2126 },
2127 WorkflowRunEventKind::StepCompleted {
2128 output: json!({"raw": private}),
2129 },
2130 WorkflowRunEventKind::StepFailed {
2131 failure: failure.clone(),
2132 },
2133 WorkflowRunEventKind::StepCancelled,
2134 WorkflowRunEventKind::StepSkipped {
2135 reason: private.to_string(),
2136 },
2137 WorkflowRunEventKind::RunSuspended {
2138 reason: private.to_string(),
2139 },
2140 WorkflowRunEventKind::RunSucceeded {
2141 output: json!({"raw": private}),
2142 },
2143 WorkflowRunEventKind::RunFailed { failure },
2144 WorkflowRunEventKind::RunCancelled,
2145 ];
2146 let at = chrono::Utc::now();
2147 let events = kinds
2148 .into_iter()
2149 .enumerate()
2150 .map(|(index, kind)| {
2151 public_workflow_event(WorkflowRunEvent {
2152 run_id: "public-run".to_string(),
2153 sequence: index as u64 + 1,
2154 at,
2155 step_id: Some("inspect".to_string()),
2156 kind,
2157 })
2158 })
2159 .collect::<Vec<_>>();
2160 let public = serde_json::to_value(&events).expect("public events serialize");
2161 let text = public.to_string();
2162
2163 assert!(!text.contains(private), "public events leaked: {text}");
2164 assert_eq!(public[2]["name"], "workflow_progressed");
2165 assert_eq!(public[6]["type"], "step_completed");
2166 assert!(public[6].get("output").is_none());
2167 assert_eq!(public[7]["failure"]["message"], "Workflow execution failed");
2168 assert_eq!(public[11]["type"], "run_succeeded");
2169 assert!(public[11].get("output").is_none());
2170 assert_eq!(
2171 events
2172 .iter()
2173 .map(|event| event.sequence)
2174 .collect::<Vec<_>>(),
2175 (1..=14).collect::<Vec<_>>()
2176 );
2177 }
2178
2179 #[test]
2180 fn workflow_tool_errors_never_expose_backend_diagnostics() {
2181 let sentinel = "/private/workspace/credentials-PRIVATE-SENTINEL";
2182 let cases = [
2183 WorkflowRunError::Storage(sentinel.to_string()),
2184 WorkflowRunError::InvalidInput(sentinel.to_string()),
2185 WorkflowRunError::Preflight(sentinel.to_string()),
2186 ];
2187 for error in cases {
2188 assert!(!workflow_tool_error(error).to_string().contains(sentinel));
2189 }
2190 assert!(!workflow_tool_error(WorkflowRunError::Compile(
2191 bamboo_domain::WorkflowCompileError::InvalidSchema(sentinel.to_string())
2192 ))
2193 .to_string()
2194 .contains(sentinel));
2195 }
2196
2197 fn canonical_workflow_run_schema() -> Value {
2198 json!({
2199 "type": "object",
2200 "properties": {
2201 "action": {
2202 "type": "string",
2203 "enum": ["start", "list", "get", "events", "cancel", "restart"]
2204 },
2205 "workflow_id": {"type": "string", "minLength": 1},
2206 "revision": {"type": "integer", "minimum": 1},
2207 "args": {"type": "object", "default": {}},
2208 "budget": {
2209 "type": "object",
2210 "properties": {
2211 "max_concurrency": {"type": "integer", "minimum": 1},
2212 "max_agents": {"type": "integer", "minimum": 0},
2213 "max_steps": {"type": "integer", "minimum": 1},
2214 "max_retries": {"type": "integer", "minimum": 0},
2215 "max_nesting_depth": {"type": "integer", "minimum": 1},
2216 "wall_time_ms": {"type": "integer", "minimum": 1},
2217 "max_tokens": {"type": "integer", "minimum": 0},
2218 "max_cost_micros": {"type": "integer", "minimum": 0}
2219 },
2220 "required": [
2221 "max_concurrency",
2222 "max_agents",
2223 "max_steps",
2224 "max_retries",
2225 "max_nesting_depth",
2226 "wall_time_ms"
2227 ],
2228 "additionalProperties": false
2229 },
2230 "run_id": {"type": "string"},
2231 "since": {"type": "integer", "minimum": 0}
2232 },
2233 "required": ["action"],
2234 "additionalProperties": false
2235 })
2236 }
2237
2238 #[tokio::test]
2239 async fn workflow_run_schema_is_flat_complete_and_canonical() {
2240 let (access, _, _) = workflow_test_access().await;
2241 let schema = WorkflowRunTool::new(access).parameters_schema();
2242
2243 for combinator in ["oneOf", "anyOf", "allOf"] {
2244 assert!(
2245 schema.get(combinator).is_none(),
2246 "workflow_run must not advertise root {combinator}"
2247 );
2248 }
2249 assert_eq!(schema, canonical_workflow_run_schema());
2250 }
2251
2252 #[tokio::test]
2253 async fn workflow_run_schema_survives_openai_sanitization_with_all_properties() {
2254 let (access, _, _) = workflow_test_access().await;
2255 let schema = WorkflowRunTool::new(access).parameters_schema();
2256 let sanitized =
2257 bamboo_llm::providers::common::tool_schema::sanitize_openai_function_parameters_schema(
2258 &schema,
2259 );
2260
2261 let properties = sanitized["properties"]
2262 .as_object()
2263 .expect("sanitized workflow_run properties");
2264 assert!(!properties.is_empty());
2265 assert_eq!(properties.len(), 7);
2266 assert_eq!(sanitized, canonical_workflow_run_schema());
2267 }
2268
2269 #[tokio::test]
2270 async fn workflow_run_schema_reaches_gemini_unchanged() {
2271 let (access, _, _) = workflow_test_access().await;
2272 let direct = WorkflowRunTool::new(access).to_schema();
2273 let gemini: GeminiTool = direct.to_provider().expect("Gemini tool conversion");
2274 let declaration = gemini
2275 .function_declarations
2276 .first()
2277 .expect("workflow_run declaration");
2278
2279 assert_eq!(declaration.name, "workflow_run");
2280 assert_eq!(
2281 declaration.parameters_json_schema.as_ref(),
2282 Some(&canonical_workflow_run_schema())
2283 );
2284 assert!(declaration.parameters.is_none());
2285 }
2286
2287 #[tokio::test]
2288 async fn workflow_run_enforces_opt_in_tightens_budget_lists_and_isolates_sessions() {
2289 let (access, repo, directory) = workflow_test_access().await;
2290 let workspace = directory.path().join("workspace");
2291 std::fs::create_dir_all(&workspace).expect("workspace");
2292 let mut session = Session::new("workflow-session", "model");
2293 session.workspace = Some(workspace.to_string_lossy().into_owned());
2294 repo.save(&mut session).await.expect("save session");
2295
2296 let denied = access
2297 .start_from_tool("workflow-session", "review-flow", 42, json!({}), None)
2298 .await
2299 .expect_err("model start defaults off without session opt-in");
2300 assert!(denied.to_string().contains("opt-in"));
2301
2302 session.metadata.insert(
2303 bamboo_skills::WORKFLOW_ORCHESTRATION_OPT_IN_METADATA_KEY.to_string(),
2304 "true".to_string(),
2305 );
2306 repo.save(&mut session).await.expect("save opt-in");
2307 let requested = WorkflowBudgets {
2308 max_concurrency: 1,
2309 max_agents: 0,
2310 max_steps: 2,
2311 max_retries: 0,
2312 max_nesting_depth: 1,
2313 wall_time_ms: 5_000,
2314 max_tokens: Some(500),
2315 max_cost_micros: Some(500),
2316 };
2317 let started = access
2318 .start_from_tool(
2319 "workflow-session",
2320 "review-flow",
2321 42,
2322 json!({}),
2323 Some(requested.clone()),
2324 )
2325 .await
2326 .expect("opted-in model start");
2327 assert_eq!(started.definition.budgets, requested);
2328 let listed = access
2329 .list_for_session("workflow-session")
2330 .await
2331 .expect("session run list");
2332 assert_eq!(listed.len(), 1);
2333 assert_eq!(listed[0].run_id, started.run_id);
2334 let progress = access
2335 .progress_for_session("workflow-session", &started.run_id, 0)
2336 .await
2337 .expect("run events");
2338 assert!(progress
2339 .events
2340 .first()
2341 .is_some_and(|event| event.kind == bamboo_domain::WorkflowRunEventKind::RunQueued));
2342 let completed = tokio::time::timeout(std::time::Duration::from_secs(2), async {
2343 loop {
2344 let progress = access
2345 .progress_for_session("workflow-session", &started.run_id, 0)
2346 .await
2347 .expect("terminal run progress");
2348 if progress.snapshot.status.is_terminal() {
2349 break progress;
2350 }
2351 tokio::time::sleep(std::time::Duration::from_millis(10)).await;
2352 }
2353 })
2354 .await
2355 .expect("workflow reaches terminal state");
2356 assert_eq!(
2357 completed.snapshot.status,
2358 bamboo_domain::WorkflowRunStatus::Succeeded
2359 );
2360 assert_eq!(
2361 completed
2362 .events
2363 .iter()
2364 .map(|event| event.sequence)
2365 .collect::<Vec<_>>(),
2366 (1..=7).collect::<Vec<_>>()
2367 );
2368 assert!(matches!(
2369 completed.events.as_slice(),
2370 [
2371 bamboo_domain::WorkflowRunEvent {
2372 kind: bamboo_domain::WorkflowRunEventKind::RunQueued,
2373 ..
2374 },
2375 bamboo_domain::WorkflowRunEvent {
2376 kind: bamboo_domain::WorkflowRunEventKind::RunStarted,
2377 ..
2378 },
2379 bamboo_domain::WorkflowRunEvent {
2380 kind: bamboo_domain::WorkflowRunEventKind::Phase { ref name },
2381 ..
2382 },
2383 bamboo_domain::WorkflowRunEvent {
2384 kind: bamboo_domain::WorkflowRunEventKind::StepQueued,
2385 ..
2386 },
2387 bamboo_domain::WorkflowRunEvent {
2388 kind: bamboo_domain::WorkflowRunEventKind::StepStarted,
2389 ..
2390 },
2391 bamboo_domain::WorkflowRunEvent {
2392 kind: bamboo_domain::WorkflowRunEventKind::StepCompleted { .. },
2393 ..
2394 },
2395 bamboo_domain::WorkflowRunEvent {
2396 kind: bamboo_domain::WorkflowRunEventKind::RunSucceeded { .. },
2397 ..
2398 }
2399 ] if name == "step_reserved"
2400 ));
2401 assert_eq!(
2402 completed.snapshot.last_sequence,
2403 completed.events.last().expect("terminal event").sequence
2404 );
2405
2406 let invalid_budget = WorkflowBudgets {
2407 max_steps: 0,
2408 ..requested.clone()
2409 };
2410 assert!(matches!(
2411 access
2412 .start_from_tool(
2413 "workflow-session",
2414 "review-flow",
2415 42,
2416 json!({}),
2417 Some(invalid_budget),
2418 )
2419 .await,
2420 Err(WorkflowRunError::InvalidInput(_))
2421 ));
2422
2423 let workflow_path = directory.path().join("skills/review-flow/workflow.yaml");
2426 let original_workflow = std::fs::read_to_string(&workflow_path).expect("workflow yaml");
2427 std::fs::write(
2428 &workflow_path,
2429 original_workflow.replace(
2430 "invocation_policy: {explicit: true, automatic: true}",
2431 "invocation_policy: {explicit: true, automatic: false}",
2432 ),
2433 )
2434 .expect("disable automatic live policy");
2435 access.skills.store().reload().await.expect("reload policy");
2436 let restart = access
2437 .restart_from_tool("workflow-session", &started.run_id)
2438 .await
2439 .expect_err("succeeded runs are terminal");
2440 assert!(matches!(restart, WorkflowRunError::Terminal));
2441
2442 let mut other = Session::new("other-session", "model");
2443 other.workspace = Some(workspace.to_string_lossy().into_owned());
2444 repo.save(&mut other).await.expect("save other session");
2445 assert!(access
2446 .list_for_session("other-session")
2447 .await
2448 .expect("isolated list")
2449 .is_empty());
2450 assert!(matches!(
2451 access
2452 .progress_for_session("other-session", &started.run_id, 0)
2453 .await,
2454 Err(WorkflowRunError::NotFound)
2455 ));
2456
2457 let run_id = started.run_id.clone();
2461 let skills = access.skills.clone();
2462 drop(access);
2463 let reopened = WorkflowRunAccess::new(
2464 directory.path(),
2465 Arc::new(WorkflowReadTool),
2466 skills,
2467 repo.clone(),
2468 )
2469 .await
2470 .expect("reopen workflow adapter");
2471 let reconnected = reopened
2472 .progress_for_session("workflow-session", &run_id, 4)
2473 .await
2474 .expect("durable reconnect");
2475 assert_eq!(reconnected.snapshot.status, WorkflowRunStatus::Succeeded);
2476 assert_eq!(
2477 reconnected
2478 .events
2479 .iter()
2480 .map(|event| event.sequence)
2481 .collect::<Vec<_>>(),
2482 vec![5, 6, 7]
2483 );
2484 assert_eq!(
2485 public_workflow_snapshot(reconnected.snapshot).last_sequence,
2486 7
2487 );
2488 assert!(matches!(
2489 public_workflow_event(reconnected.events.last().cloned().expect("tail event")).kind,
2490 PublicWorkflowRunEventKind::RunSucceeded
2491 ));
2492 assert_eq!(
2493 reopened
2494 .list_for_session("workflow-session")
2495 .await
2496 .expect("reconnected list")
2497 .len(),
2498 1
2499 );
2500 }
2501
2502 #[tokio::test]
2503 async fn workflow_cancel_is_idempotent_and_reconnects_to_one_terminal_event() {
2504 let entered = Arc::new(Semaphore::new(0));
2505 let (access, repo, directory) =
2506 workflow_test_access_with_tools(Arc::new(BlockingWorkflowReadTool {
2507 entered: entered.clone(),
2508 }))
2509 .await;
2510 let workspace = directory.path().join("workspace");
2511 std::fs::create_dir_all(&workspace).expect("workspace");
2512 let mut session = Session::new("cancel-session", "model");
2513 session.workspace = Some(workspace.to_string_lossy().into_owned());
2514 repo.save(&mut session).await.expect("save session");
2515
2516 let started = access
2517 .start("cancel-session", "review-flow", 42, json!({}), None)
2518 .await
2519 .expect("start blocking run");
2520 let _entered = tokio::time::timeout(std::time::Duration::from_secs(2), entered.acquire())
2521 .await
2522 .expect("blocking step entered")
2523 .expect("semaphore open");
2524
2525 let first = access
2526 .cancel_for_session("cancel-session", &started.run_id)
2527 .await
2528 .expect("first cancel");
2529 let second = access
2530 .cancel_for_session("cancel-session", &started.run_id)
2531 .await
2532 .expect("idempotent cancel");
2533 assert_eq!(first.status, WorkflowRunStatus::Cancelled);
2534 assert_eq!(second.status, WorkflowRunStatus::Cancelled);
2535 assert_eq!(first.last_sequence, second.last_sequence);
2536
2537 let progress = access
2538 .progress_for_session("cancel-session", &started.run_id, 0)
2539 .await
2540 .expect("cancelled journal");
2541 assert_eq!(progress.snapshot.status, WorkflowRunStatus::Cancelled);
2542 assert_eq!(progress.snapshot.last_sequence, first.last_sequence);
2543 assert_eq!(
2544 progress
2545 .events
2546 .iter()
2547 .filter(|event| matches!(event.kind, WorkflowRunEventKind::RunCancelled))
2548 .count(),
2549 1
2550 );
2551 assert!(!progress
2552 .events
2553 .iter()
2554 .any(|event| matches!(event.kind, WorkflowRunEventKind::RunSucceeded { .. })));
2555 assert_eq!(
2556 progress
2557 .events
2558 .iter()
2559 .map(|event| event.sequence)
2560 .collect::<Vec<_>>(),
2561 (1..=progress.snapshot.last_sequence).collect::<Vec<_>>()
2562 );
2563 let tail = access
2564 .progress_for_session(
2565 "cancel-session",
2566 &started.run_id,
2567 progress.snapshot.last_sequence,
2568 )
2569 .await
2570 .expect("terminal reconnect tail");
2571 assert!(tail.events.is_empty());
2572 assert!(matches!(
2573 public_workflow_snapshot(second).status,
2574 WorkflowRunStatus::Cancelled
2575 ));
2576 }
2577
2578 #[test]
2579 fn tool_input_rejects_security_context_spoofing() {
2580 let error = serde_json::from_value::<WorkflowToolInput>(json!({
2581 "action": "start",
2582 "workflow_id": "safe",
2583 "revision": 1,
2584 "workspace_trusted": true
2585 }))
2586 .unwrap_err();
2587 assert!(error.to_string().contains("unknown field"));
2588 }
2589
2590 #[test]
2591 fn tool_input_rejects_fields_from_other_actions() {
2592 let invalid = [
2593 json!({"action": "list", "run_id": "run-1"}),
2594 json!({"action": "get", "run_id": "run-1", "since": 1}),
2595 json!({"action": "events", "run_id": "run-1", "workflow_id": "flow"}),
2596 json!({"action": "cancel", "run_id": "run-1", "revision": 1}),
2597 json!({"action": "restart", "run_id": "run-1", "budget": {}}),
2598 json!({
2599 "action": "start",
2600 "workflow_id": "flow",
2601 "revision": 1,
2602 "run_id": "run-1"
2603 }),
2604 ];
2605
2606 for input in invalid {
2607 let error = serde_json::from_value::<WorkflowToolInput>(input.clone())
2608 .expect_err("action-specific fields must remain authoritative at runtime");
2609 assert!(
2610 error.to_string().contains("unknown field"),
2611 "unexpected error for {input}: {error}"
2612 );
2613 }
2614
2615 assert!(matches!(
2616 serde_json::from_value::<WorkflowToolInput>(json!({"action": "list"}))
2617 .expect("fieldless list action"),
2618 WorkflowToolInput::List {}
2619 ));
2620 }
2621
2622 #[tokio::test]
2623 async fn omitted_start_args_and_zero_budgets_match_schema() {
2624 let WorkflowToolInput::Start { args, .. } =
2625 serde_json::from_value::<WorkflowToolInput>(json!({
2626 "action": "start",
2627 "workflow_id": "safe",
2628 "revision": 1
2629 }))
2630 .expect("tool args default")
2631 else {
2632 panic!("start input")
2633 };
2634 assert_eq!(args, json!({}));
2635 let http: crate::handlers::workflow_runs::StartWorkflowRunRequest =
2636 serde_json::from_value(json!({"workflow_id":"safe", "revision":1}))
2637 .expect("http args default");
2638 assert_eq!(http.args, json!({}));
2639
2640 let (access, _, _) = workflow_test_access().await;
2641 let schema = WorkflowRunTool { access }.parameters_schema();
2642 let properties = &schema["properties"];
2643 assert_eq!(properties["args"]["default"], json!({}));
2644 assert_eq!(
2645 properties["budget"]["properties"]["max_agents"]["minimum"],
2646 0
2647 );
2648 assert_eq!(
2649 properties["budget"]["properties"]["max_retries"]["minimum"],
2650 0
2651 );
2652 assert_eq!(
2653 properties["budget"]["properties"]["max_tokens"]["minimum"],
2654 0
2655 );
2656 assert_eq!(
2657 properties["budget"]["properties"]["max_cost_micros"]["minimum"],
2658 0
2659 );
2660 }
2661
2662 #[tokio::test]
2663 async fn run_index_updates_are_concurrent_and_never_evict_at_active_capacity() {
2664 let (access, repo, _) = workflow_test_access().await;
2665 let mut session = Session::new("run-index", "model");
2666 repo.save(&mut session).await.expect("seed session");
2667 let results = futures::future::join_all((0..32).map(|index| {
2668 let access = access.clone();
2669 async move {
2670 let run_id = format!("run-{index}");
2671 access.remember_run_id("run-index", &run_id).await
2672 }
2673 }))
2674 .await;
2675 assert!(results.into_iter().all(|result| result.is_ok()));
2676 let concurrent = repo
2677 .try_load("run-index")
2678 .await
2679 .expect("load")
2680 .expect("session");
2681 let ids = serde_json::from_str::<Vec<String>>(
2682 concurrent
2683 .metadata
2684 .get(bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY)
2685 .expect("run ids"),
2686 )
2687 .expect("ids json");
2688 assert_eq!(ids.len(), 32);
2689 assert_eq!(ids.iter().collect::<BTreeSet<_>>().len(), 32);
2690
2691 let capacity_ids = (0..MAX_WORKFLOW_RUN_IDS_PER_SESSION)
2692 .map(|index| format!("active-{index}"))
2693 .collect::<Vec<_>>();
2694 repo.update_runtime_session(
2695 "run-index",
2696 &[bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY],
2697 {
2698 let capacity_ids = capacity_ids.clone();
2699 move |session| {
2700 session.metadata.insert(
2701 bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY.to_string(),
2702 serde_json::to_string(&capacity_ids).expect("ids json"),
2703 );
2704 }
2705 },
2706 )
2707 .await
2708 .expect("fill index")
2709 .expect("session");
2710 assert!(matches!(
2711 access.remember_run_id("run-index", "new-run").await,
2712 Err(WorkflowRunError::Storage(_))
2713 ));
2714 let retained = repo
2715 .try_load("run-index")
2716 .await
2717 .expect("load")
2718 .expect("session");
2719 let retained = serde_json::from_str::<Vec<String>>(
2720 retained
2721 .metadata
2722 .get(bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY)
2723 .expect("run ids"),
2724 )
2725 .expect("ids json");
2726 assert_eq!(
2727 retained, capacity_ids,
2728 "oldest active id must not be evicted"
2729 );
2730 }
2731
2732 #[tokio::test]
2733 async fn real_model_workflow_run_index_survives_tool_result_and_final_session_save() {
2734 use bamboo_agent_core::storage::AttachmentReader;
2735 use bamboo_engine::{Agent, ExecuteRequestBuilder};
2736 use bamboo_llm::{LLMChunk, LLMProvider, LLMStream};
2737 use futures::stream;
2738 use tokio::sync::Mutex;
2739 use tokio_util::sync::CancellationToken;
2740
2741 struct NoAttachments;
2742 #[async_trait]
2743 impl AttachmentReader for NoAttachments {
2744 async fn read_attachment(
2745 &self,
2746 _session_id: &str,
2747 _attachment_id: &str,
2748 ) -> std::io::Result<Option<(Vec<u8>, String)>> {
2749 Ok(None)
2750 }
2751 }
2752 struct QueueProvider {
2753 queue: Mutex<Vec<Vec<bamboo_llm::provider::Result<LLMChunk>>>>,
2754 }
2755 #[async_trait]
2756 impl LLMProvider for QueueProvider {
2757 async fn chat_stream(
2758 &self,
2759 _messages: &[bamboo_agent_core::Message],
2760 _tools: &[ToolSchema],
2761 _max_output_tokens: Option<u32>,
2762 _model: &str,
2763 ) -> bamboo_llm::provider::Result<LLMStream> {
2764 Ok(Box::pin(stream::iter(self.queue.lock().await.remove(0))))
2765 }
2766 }
2767
2768 let (access, repo, directory) = workflow_test_access().await;
2769 let session_id = "real-model-workflow-run";
2770 let mut session = Session::new(session_id, "test-model");
2771 session.metadata.insert(
2772 bamboo_skills::WORKFLOW_ORCHESTRATION_OPT_IN_METADATA_KEY.to_string(),
2773 "true".to_string(),
2774 );
2775 session
2776 .metadata
2777 .insert("external.metadata".to_string(), "preserve".to_string());
2778 session.add_message(bamboo_agent_core::Message::system("system"));
2779 session.add_message(bamboo_agent_core::Message::user("run review workflow"));
2780 repo.save(&mut session).await.expect("seed session");
2781 let call = ToolCall {
2782 id: "call-workflow-run".to_string(),
2783 tool_type: "function".to_string(),
2784 function: bamboo_agent_core::tools::FunctionCall {
2785 name: "workflow_run".to_string(),
2786 arguments: json!({
2787 "action":"start",
2788 "workflow_id":"review-flow",
2789 "revision":42
2790 })
2791 .to_string(),
2792 },
2793 };
2794 let provider = Arc::new(QueueProvider {
2795 queue: Mutex::new(vec![
2796 vec![Ok(LLMChunk::ToolCalls(vec![call])), Ok(LLMChunk::Done)],
2797 vec![Ok(LLMChunk::Token("done".to_string())), Ok(LLMChunk::Done)],
2798 ]),
2799 });
2800 let tools = Arc::new(
2801 bamboo_tools::BuiltinToolExecutorBuilder::new()
2802 .with_tool(WorkflowRunTool::new(access.clone()))
2803 .expect("workflow tool")
2804 .build(),
2805 );
2806 let metrics = bamboo_metrics::MetricsCollector::spawn(
2807 Arc::new(bamboo_metrics::SqliteMetricsStorage::new(
2808 directory.path().join("runner-metrics.db"),
2809 )),
2810 7,
2811 );
2812 let agent = Agent::builder()
2813 .storage(repo.storage().clone())
2814 .persistence(Arc::new(repo.clone()))
2815 .attachment_reader(Arc::new(NoAttachments))
2816 .skill_manager(access.skills.clone())
2817 .metrics_collector(metrics)
2818 .config(Arc::new(RwLock::new(bamboo_llm::Config::default())))
2819 .provider(provider)
2820 .default_tools(tools)
2821 .build()
2822 .expect("agent");
2823 let (event_tx, _event_rx) = tokio::sync::mpsc::channel(128);
2824 agent
2825 .execute(
2826 &mut session,
2827 ExecuteRequestBuilder::new(
2828 "run review workflow",
2829 event_tx,
2830 CancellationToken::new(),
2831 )
2832 .model("test-model")
2833 .build(),
2834 )
2835 .await
2836 .expect("real model workflow run");
2837
2838 let saved = repo
2839 .storage()
2840 .load_session(session_id)
2841 .await
2842 .expect("load")
2843 .expect("saved");
2844 assert_eq!(
2845 saved.metadata.get("external.metadata").map(String::as_str),
2846 Some("preserve")
2847 );
2848 assert!(saved
2849 .metadata
2850 .contains_key(bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY));
2851 let listed = access
2852 .list_for_session(session_id)
2853 .await
2854 .expect("list after final save");
2855 assert_eq!(listed.len(), 1);
2856 assert!(saved.messages.iter().any(|message| {
2857 message.tool_calls.as_ref().is_some_and(|calls| {
2858 calls
2859 .iter()
2860 .any(|call| call.function.name == "workflow_run")
2861 })
2862 }));
2863 }
2864
2865 #[tokio::test]
2866 async fn http_start_survives_concurrent_stale_runner_final_save_and_server_restart() {
2867 use bamboo_agent_core::storage::AttachmentReader;
2868 use bamboo_engine::{Agent, ExecuteRequestBuilder};
2869 use bamboo_llm::{LLMChunk, LLMProvider, LLMStream};
2870 use futures::stream;
2871 use tokio::sync::{oneshot, Mutex};
2872 use tokio_util::sync::CancellationToken;
2873
2874 struct NoAttachments;
2875 #[async_trait]
2876 impl AttachmentReader for NoAttachments {
2877 async fn read_attachment(
2878 &self,
2879 _session_id: &str,
2880 _attachment_id: &str,
2881 ) -> std::io::Result<Option<(Vec<u8>, String)>> {
2882 Ok(None)
2883 }
2884 }
2885
2886 struct PausingProvider {
2887 entered: Mutex<Option<oneshot::Sender<()>>>,
2888 resume: Mutex<Option<oneshot::Receiver<()>>>,
2889 }
2890 #[async_trait]
2891 impl LLMProvider for PausingProvider {
2892 async fn chat_stream(
2893 &self,
2894 _messages: &[bamboo_agent_core::Message],
2895 _tools: &[ToolSchema],
2896 _max_output_tokens: Option<u32>,
2897 _model: &str,
2898 ) -> bamboo_llm::provider::Result<LLMStream> {
2899 if let Some(entered) = self.entered.lock().await.take() {
2900 let _ = entered.send(());
2901 }
2902 if let Some(resume) = self.resume.lock().await.take() {
2903 let _ = resume.await;
2904 }
2905 Ok(Box::pin(stream::iter(vec![
2906 Ok(LLMChunk::Token("done".to_string())),
2907 Ok(LLMChunk::Done),
2908 ])))
2909 }
2910 }
2911
2912 let (access, repo, directory) = workflow_test_access().await;
2913 let session_id = "http-start-concurrent-runner-save";
2914 let mut session = Session::new(session_id, "test-model");
2915 session.add_message(bamboo_agent_core::Message::system("system"));
2916 session.add_message(bamboo_agent_core::Message::user("keep running"));
2917 repo.save(&mut session).await.expect("seed session");
2918 let mut runner_session = repo
2919 .try_load(session_id)
2920 .await
2921 .expect("load runner session")
2922 .expect("runner session");
2923
2924 let (entered_tx, entered_rx) = oneshot::channel();
2925 let (resume_tx, resume_rx) = oneshot::channel();
2926 let provider = Arc::new(PausingProvider {
2927 entered: Mutex::new(Some(entered_tx)),
2928 resume: Mutex::new(Some(resume_rx)),
2929 });
2930 let metrics = bamboo_metrics::MetricsCollector::spawn(
2931 Arc::new(bamboo_metrics::SqliteMetricsStorage::new(
2932 directory.path().join("http-runner-metrics.db"),
2933 )),
2934 7,
2935 );
2936 let agent = Agent::builder()
2937 .storage(repo.storage().clone())
2938 .persistence(Arc::new(repo.clone()))
2939 .attachment_reader(Arc::new(NoAttachments))
2940 .skill_manager(access.skills.clone())
2941 .metrics_collector(metrics)
2942 .config(Arc::new(RwLock::new(bamboo_llm::Config::default())))
2943 .provider(provider)
2944 .default_tools(Arc::new(
2945 bamboo_tools::BuiltinToolExecutorBuilder::new().build(),
2946 ))
2947 .build()
2948 .expect("agent");
2949 let (event_tx, _event_rx) = tokio::sync::mpsc::channel(128);
2950 let runner = tokio::spawn(async move {
2951 agent
2952 .execute(
2953 &mut runner_session,
2954 ExecuteRequestBuilder::new("keep running", event_tx, CancellationToken::new())
2955 .model("test-model")
2956 .build(),
2957 )
2958 .await
2959 });
2960 tokio::time::timeout(std::time::Duration::from_secs(2), entered_rx)
2961 .await
2962 .expect("runner enters model round")
2963 .expect("runner entry signal");
2964
2965 let started = access
2966 .start(session_id, "review-flow", 42, json!({}), None)
2967 .await
2968 .expect("HTTP-equivalent explicit start");
2969 let durable_during_round = repo
2970 .storage()
2971 .load_session(session_id)
2972 .await
2973 .expect("load during round")
2974 .expect("session during round");
2975 assert!(durable_during_round
2976 .metadata
2977 .get(bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY)
2978 .is_some_and(|raw| raw.contains(&started.run_id)));
2979
2980 resume_tx.send(()).expect("resume runner");
2981 tokio::time::timeout(std::time::Duration::from_secs(2), runner)
2982 .await
2983 .expect("runner completes")
2984 .expect("runner task")
2985 .expect("runner execution");
2986
2987 let durable_after_final_save = repo
2988 .storage()
2989 .load_session(session_id)
2990 .await
2991 .expect("load after final save")
2992 .expect("saved session");
2993 assert!(durable_after_final_save
2994 .metadata
2995 .get(bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY)
2996 .is_some_and(|raw| raw.contains(&started.run_id)));
2997
2998 tokio::time::timeout(std::time::Duration::from_secs(2), async {
2999 loop {
3000 let progress = access
3001 .progress_for_session(session_id, &started.run_id, u64::MAX)
3002 .await
3003 .expect("workflow progress before restart");
3004 if progress.snapshot.status.is_terminal()
3005 && !access.engine.is_run_active(&started.run_id)
3006 {
3007 break;
3008 }
3009 tokio::task::yield_now().await;
3010 }
3011 })
3012 .await
3013 .expect("workflow reaches terminal state before restart");
3014
3015 let skills = access.skills.clone();
3016 drop(access);
3017 let restarted =
3018 WorkflowRunAccess::new(directory.path(), Arc::new(WorkflowReadTool), skills, repo)
3019 .await
3020 .expect("restart workflow access");
3021 let listed = restarted
3022 .list_for_session(session_id)
3023 .await
3024 .expect("list after restart");
3025 assert_eq!(listed.len(), 1);
3026 assert_eq!(listed[0].run_id, started.run_id);
3027 }
3028
3029 #[tokio::test]
3030 async fn production_policy_allows_read_without_fabricating_workspace_trust() {
3031 let read = BTreeSet::from(["read".to_string()]);
3032 assert_eq!(
3033 ServerWorkflowPolicy
3034 .authorize(
3035 "session",
3036 &WorkflowPolicyTarget::Tool("read_file".to_string()),
3037 &read,
3038 false,
3039 )
3040 .await,
3041 PermissionDecision::Allow
3042 );
3043
3044 let write = BTreeSet::from(["write".to_string()]);
3045 assert!(matches!(
3046 ServerWorkflowPolicy
3047 .authorize(
3048 "session",
3049 &WorkflowPolicyTarget::Tool("write_file".to_string()),
3050 &write,
3051 false,
3052 )
3053 .await,
3054 PermissionDecision::Deny(_)
3055 ));
3056
3057 for hostile_target in [
3058 "Write",
3059 "write_file",
3060 "WebFetch",
3061 "mcp::remote_tool",
3062 "Bash",
3063 ] {
3064 for claimed in [BTreeSet::new(), read.clone()] {
3065 assert!(matches!(
3066 ServerWorkflowPolicy
3067 .authorize(
3068 "session",
3069 &WorkflowPolicyTarget::Tool(hostile_target.to_string()),
3070 &claimed,
3071 false,
3072 )
3073 .await,
3074 PermissionDecision::Deny(_)
3075 ));
3076 }
3077 }
3078
3079 assert_eq!(
3080 ServerWorkflowPolicy
3081 .authorize(
3082 "session",
3083 &WorkflowPolicyTarget::Workflow {
3084 id: "nested-review".to_string(),
3085 revision: 1,
3086 },
3087 &BTreeSet::new(),
3088 false,
3089 )
3090 .await,
3091 PermissionDecision::Allow
3092 );
3093 }
3094
3095 #[test]
3096 fn pinned_bundle_limits_reject_oversized_runs_before_engine_start() {
3097 assert!(enforce_pinned_bundle_limits(
3098 MAX_PINNED_DEFINITIONS_PER_RUN,
3099 MAX_PINNED_BUNDLE_BYTES_PER_RUN
3100 )
3101 .is_ok());
3102 assert!(enforce_pinned_bundle_limits(
3103 MAX_PINNED_DEFINITIONS_PER_RUN + 1,
3104 MAX_PINNED_BUNDLE_BYTES_PER_RUN
3105 )
3106 .is_err());
3107 assert!(enforce_pinned_bundle_limits(
3108 MAX_PINNED_DEFINITIONS_PER_RUN,
3109 MAX_PINNED_BUNDLE_BYTES_PER_RUN + 1
3110 )
3111 .is_err());
3112 }
3113}