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 if let Err(error) = self.remember_run_id(session_id, &snapshot.run_id).await {
302 return match self.engine.cancel(&snapshot.run_id).await {
303 Ok(cancelled) if cancelled.status.is_terminal() => Err(WorkflowRunError::Storage(
304 format!(
305 "run index persistence failed; run {} reached terminal {:?}: {error}",
306 snapshot.run_id, cancelled.status
307 ),
308 )),
309 Ok(cancelled) => Err(WorkflowRunError::Storage(format!(
310 "run index persistence failed; orphan run {} remains {:?}; repair with this run id: {error}",
311 snapshot.run_id, cancelled.status
312 ))),
313 Err(cancel_error) => Err(WorkflowRunError::Storage(format!(
314 "run index persistence failed; orphan run {} could not be cancelled ({cancel_error}); repair with this run id: {error}",
315 snapshot.run_id
316 ))),
317 };
318 }
319 Ok(snapshot)
320 }
321
322 async fn ensure_run_index_capacity(&self, session_id: &str) -> Result<(), WorkflowRunError> {
323 let session = self
324 .sessions
325 .try_load(session_id)
326 .await
327 .map_err(|_| WorkflowRunError::Storage("session state unavailable".to_string()))?
328 .ok_or(WorkflowRunError::NotFound)?;
329 let ids = session
330 .metadata
331 .get(bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY)
332 .and_then(|raw| serde_json::from_str::<Vec<String>>(raw).ok())
333 .unwrap_or_default();
334 if ids.len() < MAX_WORKFLOW_RUN_IDS_PER_SESSION {
335 return Ok(());
336 }
337 let mut evictable = BTreeSet::new();
338 for run_id in &ids {
339 match self.engine.progress(run_id, u64::MAX).await {
340 Ok(progress) if progress.snapshot.status.is_terminal() => {
341 evictable.insert(run_id.clone());
342 }
343 Err(WorkflowRunError::NotFound) => {
344 evictable.insert(run_id.clone());
345 }
346 Ok(_) => {}
347 Err(error) => return Err(error),
348 }
349 }
350 if !evictable.is_empty() {
351 self.sessions
352 .update_runtime_session(
353 session_id,
354 &[bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY],
355 move |session| {
356 let mut ids = session
357 .metadata
358 .get(bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY)
359 .and_then(|raw| serde_json::from_str::<Vec<String>>(raw).ok())
360 .unwrap_or_default();
361 ids.retain(|id| !evictable.contains(id));
362 session.metadata.insert(
363 bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY.to_string(),
364 serde_json::to_string(&ids)
365 .expect("string vector serialization cannot fail"),
366 );
367 },
368 )
369 .await
370 .map_err(|_| WorkflowRunError::Storage("run index pruning failed".to_string()))?
371 .ok_or(WorkflowRunError::NotFound)?;
372 }
373 let remaining = self
374 .sessions
375 .try_load(session_id)
376 .await
377 .map_err(|_| WorkflowRunError::Storage("session state unavailable".to_string()))?
378 .and_then(|session| {
379 session
380 .metadata
381 .get(bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY)
382 .and_then(|raw| serde_json::from_str::<Vec<String>>(raw).ok())
383 })
384 .unwrap_or_default()
385 .len();
386 if remaining >= MAX_WORKFLOW_RUN_IDS_PER_SESSION {
387 return Err(WorkflowRunError::Preflight(
388 "workflow run index is full of active runs".to_string(),
389 ));
390 }
391 Ok(())
392 }
393
394 async fn remember_run_id(
395 &self,
396 session_id: &str,
397 run_id: &str,
398 ) -> Result<(), WorkflowRunError> {
399 let run_id = run_id.to_string();
400 let index_full = Arc::new(AtomicBool::new(false));
401 let index_full_in_transaction = index_full.clone();
402 self.sessions
403 .update_runtime_session(
404 session_id,
405 &[bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY],
406 move |session| {
407 let mut ids = session
408 .metadata
409 .get(bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY)
410 .and_then(|raw| serde_json::from_str::<Vec<String>>(raw).ok())
411 .unwrap_or_default();
412 ids.retain(|existing| existing != &run_id);
413 if ids.len() >= MAX_WORKFLOW_RUN_IDS_PER_SESSION {
414 index_full_in_transaction.store(true, Ordering::SeqCst);
415 return;
416 }
417 ids.push(run_id);
418 session.metadata.insert(
419 bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY.to_string(),
420 serde_json::to_string(&ids)
421 .expect("string vector serialization cannot fail"),
422 );
423 },
424 )
425 .await
426 .map_err(|_| WorkflowRunError::Storage("run index persistence failed".to_string()))?
427 .ok_or(WorkflowRunError::NotFound)
428 .and_then(|_| {
429 if index_full.load(Ordering::SeqCst) {
430 Err(WorkflowRunError::Storage(
431 "workflow run index reached its active-run capacity".to_string(),
432 ))
433 } else {
434 Ok(())
435 }
436 })
437 }
438
439 pub async fn list_for_session(
440 &self,
441 session_id: &str,
442 ) -> Result<Vec<WorkflowRunSnapshot>, WorkflowRunError> {
443 self.session_context(session_id).await?;
444 let session = self
445 .sessions
446 .try_load(session_id)
447 .await
448 .map_err(|_| WorkflowRunError::Storage("session state unavailable".to_string()))?
449 .ok_or(WorkflowRunError::NotFound)?;
450 let run_ids = session
451 .metadata
452 .get(bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY)
453 .and_then(|raw| serde_json::from_str::<Vec<String>>(raw).ok())
454 .unwrap_or_default();
455 let mut snapshots = Vec::new();
456 let mut stale = BTreeSet::new();
457 for run_id in run_ids {
458 match self.engine.progress(&run_id, u64::MAX).await {
459 Ok(progress) if progress.snapshot.session_id == session_id => {
460 snapshots.push(progress.snapshot);
461 }
462 Ok(_) | Err(WorkflowRunError::NotFound) => {
463 stale.insert(run_id);
464 }
465 Err(error) => return Err(error),
466 }
467 }
468 if !stale.is_empty() {
469 let _ = self
470 .sessions
471 .update_runtime_session(
472 session_id,
473 &[bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY],
474 move |session| {
475 let mut ids = session
476 .metadata
477 .get(bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY)
478 .and_then(|raw| serde_json::from_str::<Vec<String>>(raw).ok())
479 .unwrap_or_default();
480 ids.retain(|id| !stale.contains(id));
481 session.metadata.insert(
482 bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY.to_string(),
483 serde_json::to_string(&ids)
484 .expect("string vector serialization cannot fail"),
485 );
486 },
487 )
488 .await;
489 }
490 snapshots.sort_by_key(|snapshot| std::cmp::Reverse(snapshot.created_at));
491 Ok(snapshots)
492 }
493
494 pub async fn progress_for_session(
495 &self,
496 session_id: &str,
497 run_id: &str,
498 since: u64,
499 ) -> Result<WorkflowProgress, WorkflowRunError> {
500 let progress = self.engine.progress(run_id, since).await?;
501 if progress.snapshot.session_id != session_id {
502 return Err(WorkflowRunError::NotFound);
503 }
504 Ok(progress)
505 }
506
507 pub async fn cancel_for_session(
508 &self,
509 session_id: &str,
510 run_id: &str,
511 ) -> Result<WorkflowRunSnapshot, WorkflowRunError> {
512 self.progress_for_session(session_id, run_id, u64::MAX)
513 .await?;
514 self.engine.cancel(run_id).await
515 }
516
517 pub async fn restart_for_session(
518 &self,
519 session_id: &str,
520 run_id: &str,
521 ) -> Result<WorkflowRunSnapshot, WorkflowRunError> {
522 self.progress_for_session(session_id, run_id, u64::MAX)
523 .await?;
524 let (_, workspace_trusted) = self.session_context(session_id).await?;
525 self.engine
526 .restart(run_id, workspace_trusted, vec!["read".to_string()])
527 .await
528 }
529
530 pub async fn restart_from_tool(
531 &self,
532 session_id: &str,
533 run_id: &str,
534 ) -> Result<WorkflowRunSnapshot, WorkflowRunError> {
535 let progress = self
536 .progress_for_session(session_id, run_id, u64::MAX)
537 .await?;
538 if progress.snapshot.definition_bundle.root_invocation_policy["automatic"].as_bool()
539 != Some(true)
540 {
541 return Err(WorkflowRunError::Preflight(
542 "pinned workflow invocation policy denies automatic restart".to_string(),
543 ));
544 }
545 let session = self
546 .sessions
547 .try_load(session_id)
548 .await
549 .map_err(|_| WorkflowRunError::Preflight("session state is unavailable".to_string()))?
550 .ok_or(WorkflowRunError::NotFound)?;
551 let opted_in = session
552 .metadata
553 .get(bamboo_skills::WORKFLOW_ORCHESTRATION_OPT_IN_METADATA_KEY)
554 .is_some_and(|value| value.eq_ignore_ascii_case("true"));
555 if !opted_in {
556 return Err(WorkflowRunError::Preflight(
557 "model-started orchestration restart requires explicit session opt-in".to_string(),
558 ));
559 }
560 self.restart_for_session(session_id, run_id).await
561 }
562}
563
564fn tighten_workflow_budget(
565 definition: &WorkflowBudgets,
566 requested: &WorkflowBudgets,
567) -> WorkflowBudgets {
568 WorkflowBudgets {
569 max_concurrency: definition.max_concurrency.min(requested.max_concurrency),
570 max_agents: definition.max_agents.min(requested.max_agents),
571 max_steps: definition.max_steps.min(requested.max_steps),
572 max_retries: definition.max_retries.min(requested.max_retries),
573 max_nesting_depth: definition
574 .max_nesting_depth
575 .min(requested.max_nesting_depth),
576 wall_time_ms: definition.wall_time_ms.min(requested.wall_time_ms),
577 max_tokens: match (definition.max_tokens, requested.max_tokens) {
578 (Some(left), Some(right)) => Some(left.min(right)),
579 (Some(value), None) | (None, Some(value)) => Some(value),
580 (None, None) => None,
581 },
582 max_cost_micros: match (definition.max_cost_micros, requested.max_cost_micros) {
583 (Some(left), Some(right)) => Some(left.min(right)),
584 (Some(value), None) | (None, Some(value)) => Some(value),
585 (None, None) => None,
586 },
587 }
588}
589
590fn validate_requested_budget(requested: &WorkflowBudgets) -> Result<(), WorkflowRunError> {
591 if requested.max_concurrency == 0
592 || requested.max_steps == 0
593 || requested.max_nesting_depth == 0
594 || requested.wall_time_ms == 0
595 {
596 return Err(WorkflowRunError::InvalidInput(
597 "workflow execution limits must be positive".to_string(),
598 ));
599 }
600 Ok(())
601}
602
603fn enforce_pinned_bundle_limits(
604 definition_count: usize,
605 serialized_bytes: usize,
606) -> Result<(), WorkflowRunError> {
607 if definition_count > MAX_PINNED_DEFINITIONS_PER_RUN {
608 return Err(WorkflowRunError::Preflight(
609 "workflow dependency count exceeds the server limit".to_string(),
610 ));
611 }
612 if serialized_bytes > MAX_PINNED_BUNDLE_BYTES_PER_RUN {
613 return Err(WorkflowRunError::Preflight(
614 "workflow definition bundle exceeds the server size limit".to_string(),
615 ));
616 }
617 Ok(())
618}
619
620#[derive(Debug, Clone, Serialize, PartialEq)]
628pub(crate) struct PublicWorkflowRunSnapshot {
629 pub run_id: String,
630 #[serde(skip_serializing_if = "Option::is_none")]
631 pub parent_run_id: Option<String>,
632 #[serde(skip_serializing_if = "Option::is_none")]
633 pub parent_step_id: Option<String>,
634 pub session_id: String,
635 pub workflow_id: String,
636 pub workflow_revision: u64,
637 pub definition_bundle_hash: String,
638 pub status: WorkflowRunStatus,
639 pub planned_steps: BTreeMap<String, PublicWorkflowPlannedStep>,
640 pub plan: PublicWorkflowPlan,
641 pub steps: BTreeMap<String, PublicWorkflowStepSnapshot>,
642 pub budget: WorkflowBudgets,
643 pub usage: WorkflowBudgetUsage,
644 pub child_agent_count: u32,
645 pub last_sequence: u64,
646 #[serde(skip_serializing_if = "Option::is_none")]
647 pub failure: Option<PublicWorkflowFailure>,
648 #[serde(skip_serializing_if = "Option::is_none")]
649 pub suspension: Option<PublicWorkflowSuspension>,
650 pub created_at: chrono::DateTime<chrono::Utc>,
651 pub updated_at: chrono::DateTime<chrono::Utc>,
652}
653
654#[derive(Debug, Clone, Serialize, PartialEq)]
655pub(crate) struct PublicWorkflowStepSnapshot {
656 pub id: String,
657 pub status: WorkflowStepStatus,
658 #[serde(skip_serializing_if = "Option::is_none")]
659 pub failure: Option<PublicWorkflowFailure>,
660 pub attempts: u32,
661}
662
663#[derive(Debug, Clone, Serialize, PartialEq, Eq)]
664pub(crate) struct PublicWorkflowPlannedStep {
665 pub id: String,
666 pub kind: PublicWorkflowPlannedStepKind,
667}
668
669#[derive(Debug, Clone, Copy, Serialize, PartialEq, Eq)]
670#[serde(rename_all = "snake_case")]
671pub(crate) enum PublicWorkflowPlannedStepKind {
672 Tool,
673 Agent,
674 Workflow,
675}
676
677#[derive(Debug, Clone, Serialize, PartialEq, Eq)]
680#[serde(tag = "type", rename_all = "snake_case")]
681pub(crate) enum PublicWorkflowPlan {
682 Step {
683 step: String,
684 },
685 Sequence {
686 nodes: Vec<PublicWorkflowPlan>,
687 },
688 Parallel {
689 nodes: Vec<PublicWorkflowPlan>,
690 },
691 Map {
692 body: Box<PublicWorkflowPlan>,
693 },
694 Retry {
695 node: Box<PublicWorkflowPlan>,
696 max_attempts: u32,
697 },
698}
699
700#[derive(Debug, Clone, Serialize, PartialEq, Eq)]
701pub(crate) struct PublicWorkflowFailure {
702 pub code: WorkflowFailureCode,
703 pub message: &'static str,
704 pub retryable: bool,
705}
706
707#[derive(Debug, Clone, Serialize, PartialEq, Eq)]
708#[serde(tag = "type", rename_all = "snake_case")]
709pub(crate) enum PublicWorkflowSuspension {
710 ToolApproval { step_id: String },
711 ToolRunning { step_id: String, killed: bool },
712 Recovery,
713}
714
715#[derive(Debug, Clone, Serialize, PartialEq)]
719pub(crate) struct PublicWorkflowRunEvent {
720 pub run_id: String,
721 pub sequence: u64,
722 pub at: chrono::DateTime<chrono::Utc>,
723 #[serde(skip_serializing_if = "Option::is_none")]
724 pub step_id: Option<String>,
725 #[serde(flatten)]
726 pub kind: PublicWorkflowRunEventKind,
727}
728
729#[derive(Debug, Clone, Serialize, PartialEq)]
730#[serde(tag = "type", rename_all = "snake_case")]
731pub(crate) enum PublicWorkflowRunEventKind {
732 RunQueued,
733 RunStarted,
734 Phase { name: &'static str },
735 StepQueued,
736 StepStarted,
737 StepSuspended,
738 StepCompleted,
739 StepFailed { failure: PublicWorkflowFailure },
740 StepCancelled,
741 StepSkipped,
742 RunSuspended,
743 RunSucceeded,
744 RunFailed { failure: PublicWorkflowFailure },
745 RunCancelled,
746}
747
748fn public_workflow_plan(plan: &WorkflowPlan) -> PublicWorkflowPlan {
749 match plan {
750 WorkflowPlan::Step { step } => PublicWorkflowPlan::Step { step: step.clone() },
751 WorkflowPlan::Sequence { nodes } => PublicWorkflowPlan::Sequence {
752 nodes: nodes.iter().map(public_workflow_plan).collect(),
753 },
754 WorkflowPlan::Parallel { nodes } => PublicWorkflowPlan::Parallel {
755 nodes: nodes.iter().map(public_workflow_plan).collect(),
756 },
757 WorkflowPlan::Map { body, .. } => PublicWorkflowPlan::Map {
758 body: Box::new(public_workflow_plan(body)),
759 },
760 WorkflowPlan::Retry {
761 node, max_attempts, ..
762 } => PublicWorkflowPlan::Retry {
763 node: Box::new(public_workflow_plan(node)),
764 max_attempts: *max_attempts,
765 },
766 }
767}
768
769fn public_planned_steps(
770 definition: &WorkflowRunDefinition,
771) -> BTreeMap<String, PublicWorkflowPlannedStep> {
772 definition
773 .steps
774 .iter()
775 .map(|step| {
776 let kind = match &step.kind {
777 WorkflowStepKind::Tool { .. } => PublicWorkflowPlannedStepKind::Tool,
778 WorkflowStepKind::Agent { .. } => PublicWorkflowPlannedStepKind::Agent,
779 WorkflowStepKind::Workflow { .. } => PublicWorkflowPlannedStepKind::Workflow,
780 };
781 (
782 step.id.clone(),
783 PublicWorkflowPlannedStep {
784 id: step.id.clone(),
785 kind,
786 },
787 )
788 })
789 .collect()
790}
791
792fn public_workflow_failure(failure: WorkflowFailure) -> PublicWorkflowFailure {
793 let message = match failure.code {
794 WorkflowFailureCode::InvalidDefinition => "Workflow definition is invalid",
795 WorkflowFailureCode::InvalidInput => "Workflow input is invalid",
796 WorkflowFailureCode::InvalidOutput => "Workflow output is invalid",
797 WorkflowFailureCode::UnknownReference => "Workflow reference is unavailable",
798 WorkflowFailureCode::PermissionDenied => "Workflow permission was denied",
799 WorkflowFailureCode::UntrustedWorkspace => "Workflow workspace is not trusted",
800 WorkflowFailureCode::BudgetExceeded => "Workflow execution budget was exceeded",
801 WorkflowFailureCode::RetryExhausted => "Workflow retry budget was exhausted",
802 WorkflowFailureCode::ExecutionFailed => "Workflow execution failed",
803 WorkflowFailureCode::Cancelled => "Workflow execution was cancelled",
804 WorkflowFailureCode::RecoverySuspended => "Workflow recovery requires attention",
805 WorkflowFailureCode::Suspended => "Workflow execution is suspended",
806 WorkflowFailureCode::DependencySkipped => "Workflow dependency was skipped",
807 WorkflowFailureCode::Storage => "Workflow storage is unavailable",
808 };
809 PublicWorkflowFailure {
810 code: failure.code,
811 message,
812 retryable: failure.retryable,
813 }
814}
815
816fn public_workflow_step(step: WorkflowStepSnapshot) -> PublicWorkflowStepSnapshot {
817 PublicWorkflowStepSnapshot {
818 id: step.id,
819 status: step.status,
820 failure: step.failure.map(public_workflow_failure),
821 attempts: step.attempts,
822 }
823}
824
825fn public_workflow_suspension(suspension: WorkflowSuspensionContext) -> PublicWorkflowSuspension {
826 match suspension {
827 WorkflowSuspensionContext::ToolApproval { step_id, .. } => {
828 PublicWorkflowSuspension::ToolApproval { step_id }
829 }
830 WorkflowSuspensionContext::ToolRunning {
831 step_id, killed, ..
832 } => PublicWorkflowSuspension::ToolRunning { step_id, killed },
833 WorkflowSuspensionContext::Recovery { .. } => PublicWorkflowSuspension::Recovery,
834 }
835}
836
837fn public_workflow_phase(name: &str) -> &'static str {
838 match name {
839 "retry_reserved" => "retry_reserved",
840 "step_reserved" => "step_reserved",
841 "agent_reserved" => "agent_reserved",
842 "agent_usage_recorded" => "agent_usage_recorded",
843 "suspension_context_persisted" => "suspension_context_persisted",
844 _ => "workflow_progressed",
845 }
846}
847
848pub(crate) fn public_workflow_snapshot(snapshot: WorkflowRunSnapshot) -> PublicWorkflowRunSnapshot {
849 let planned_steps = public_planned_steps(&snapshot.definition);
850 let plan = public_workflow_plan(&snapshot.definition.plan);
851 let child_agent_count = snapshot.usage.agents;
852 PublicWorkflowRunSnapshot {
853 run_id: snapshot.run_id,
854 parent_run_id: snapshot.parent_run_id,
855 parent_step_id: snapshot.parent_step_id,
856 session_id: snapshot.session_id,
857 workflow_id: snapshot.definition.id,
858 workflow_revision: snapshot.definition.revision,
859 definition_bundle_hash: snapshot.definition_bundle_hash,
860 status: snapshot.status,
861 planned_steps,
862 plan,
863 steps: snapshot
864 .steps
865 .into_iter()
866 .map(|(id, step)| (id, public_workflow_step(step)))
867 .collect(),
868 budget: snapshot.definition.budgets,
869 usage: snapshot.usage,
870 child_agent_count,
871 last_sequence: snapshot.last_sequence,
872 failure: snapshot.failure.map(public_workflow_failure),
873 suspension: snapshot.suspension.map(public_workflow_suspension),
874 created_at: snapshot.created_at,
875 updated_at: snapshot.updated_at,
876 }
877}
878
879pub(crate) fn public_workflow_event(event: WorkflowRunEvent) -> PublicWorkflowRunEvent {
880 let kind = match event.kind {
881 WorkflowRunEventKind::RunQueued => PublicWorkflowRunEventKind::RunQueued,
882 WorkflowRunEventKind::RunStarted => PublicWorkflowRunEventKind::RunStarted,
883 WorkflowRunEventKind::Phase { name } => PublicWorkflowRunEventKind::Phase {
884 name: public_workflow_phase(&name),
885 },
886 WorkflowRunEventKind::StepQueued => PublicWorkflowRunEventKind::StepQueued,
887 WorkflowRunEventKind::StepStarted => PublicWorkflowRunEventKind::StepStarted,
888 WorkflowRunEventKind::StepSuspended { .. } => PublicWorkflowRunEventKind::StepSuspended,
889 WorkflowRunEventKind::StepCompleted { .. } => PublicWorkflowRunEventKind::StepCompleted,
890 WorkflowRunEventKind::StepFailed { failure } => PublicWorkflowRunEventKind::StepFailed {
891 failure: public_workflow_failure(failure),
892 },
893 WorkflowRunEventKind::StepCancelled => PublicWorkflowRunEventKind::StepCancelled,
894 WorkflowRunEventKind::StepSkipped { .. } => PublicWorkflowRunEventKind::StepSkipped,
895 WorkflowRunEventKind::RunSuspended { .. } => PublicWorkflowRunEventKind::RunSuspended,
896 WorkflowRunEventKind::RunSucceeded { .. } => PublicWorkflowRunEventKind::RunSucceeded,
897 WorkflowRunEventKind::RunFailed { failure } => PublicWorkflowRunEventKind::RunFailed {
898 failure: public_workflow_failure(failure),
899 },
900 WorkflowRunEventKind::RunCancelled => PublicWorkflowRunEventKind::RunCancelled,
901 };
902 PublicWorkflowRunEvent {
903 run_id: event.run_id,
904 sequence: event.sequence,
905 at: event.at,
906 step_id: event.step_id,
907 kind,
908 }
909}
910
911struct ExternallyPinnedDefinitions;
915
916#[async_trait]
917impl WorkflowDefinitionPort for ExternallyPinnedDefinitions {
918 async fn pin_bundle(
919 &self,
920 _root: &WorkflowRunDefinition,
921 ) -> Result<WorkflowDefinitionBundle, String> {
922 Err("server workflows must be pinned through SkillManager".to_string())
923 }
924}
925
926struct UnavailableAgentPort;
929
930#[async_trait]
931impl AgentStepPort for UnavailableAgentPort {
932 async fn resolve(&self, _name: &str) -> Result<Option<NamedAgentSpec>, String> {
933 Ok(None)
934 }
935
936 async fn execute(
937 &self,
938 _spec: &NamedAgentSpec,
939 _prompt: Value,
940 _model: Option<&str>,
941 _effort: Option<&str>,
942 _capabilities: &BTreeSet<String>,
943 _session_id: &str,
944 ) -> Result<AgentStepResult, String> {
945 Err("named-agent execution is not available".to_string())
946 }
947}
948
949struct ServerWorkflowPolicy;
950
951#[async_trait]
952impl WorkflowPolicyPort for ServerWorkflowPolicy {
953 async fn authorize(
954 &self,
955 _session_id: &str,
956 target: &WorkflowPolicyTarget,
957 requested: &BTreeSet<String>,
958 _workspace_trusted: bool,
959 ) -> PermissionDecision {
960 match target {
966 WorkflowPolicyTarget::Workflow { .. } if requested.is_empty() => {
970 PermissionDecision::Allow
971 }
972 WorkflowPolicyTarget::Tool(name) => {
973 let safe_target = SAFE_UNTRUSTED_WORKFLOW_TOOLS
974 .iter()
975 .any(|candidate| candidate.eq_ignore_ascii_case(name));
976 if safe_target && requested.iter().all(|capability| capability == "read") {
977 PermissionDecision::Allow
978 } else {
979 PermissionDecision::Deny(
980 "workflow capability authority is not available".to_string(),
981 )
982 }
983 }
984 WorkflowPolicyTarget::Agent(_) | WorkflowPolicyTarget::Workflow { .. } => {
985 PermissionDecision::Deny(
986 "workflow capability authority is not available".to_string(),
987 )
988 }
989 }
990 }
991}
992
993struct UnavailableSecretResolver;
994
995#[async_trait]
996impl WorkflowSecretResolverPort for UnavailableSecretResolver {
997 async fn resolve(
998 &self,
999 _session_id: &str,
1000 _capability: &str,
1001 ) -> Result<WorkflowSecretMaterial, String> {
1002 Err("workflow secret capability resolver is not available".to_string())
1003 }
1004}
1005
1006#[derive(Debug, Deserialize)]
1007#[serde(tag = "action", rename_all = "snake_case", deny_unknown_fields)]
1008enum WorkflowToolInput {
1009 Start {
1010 workflow_id: String,
1011 revision: u64,
1012 #[serde(default = "empty_object")]
1013 args: Value,
1014 #[serde(default)]
1015 budget: Option<WorkflowBudgets>,
1016 },
1017 List {},
1018 Get {
1019 run_id: String,
1020 },
1021 Events {
1022 run_id: String,
1023 #[serde(default)]
1024 since: u64,
1025 },
1026 Cancel {
1027 run_id: String,
1028 },
1029 Restart {
1030 run_id: String,
1031 },
1032}
1033
1034fn empty_object() -> Value {
1035 json!({})
1036}
1037
1038pub struct WorkflowRunTool {
1039 access: WorkflowRunAccess,
1040}
1041
1042impl WorkflowRunTool {
1043 pub fn new(access: WorkflowRunAccess) -> Self {
1044 Self { access }
1045 }
1046}
1047
1048#[async_trait]
1049impl Tool for WorkflowRunTool {
1050 fn name(&self) -> &str {
1051 "workflow_run"
1052 }
1053
1054 fn description(&self) -> &str {
1055 "Start, inspect, cancel, or safely restart a catalog-pinned workflow run"
1056 }
1057
1058 fn parameters_schema(&self) -> Value {
1059 json!({
1060 "type": "object",
1061 "properties": {
1062 "action": {
1063 "type": "string",
1064 "enum": ["start", "list", "get", "events", "cancel", "restart"]
1065 },
1066 "workflow_id": {"type": "string", "minLength": 1},
1067 "revision": {"type": "integer", "minimum": 1},
1068 "args": {"type": "object", "default": {}},
1069 "budget": {
1070 "type": "object",
1071 "properties": {
1072 "max_concurrency": {"type": "integer", "minimum": 1},
1073 "max_agents": {"type": "integer", "minimum": 0},
1074 "max_steps": {"type": "integer", "minimum": 1},
1075 "max_retries": {"type": "integer", "minimum": 0},
1076 "max_nesting_depth": {"type": "integer", "minimum": 1},
1077 "wall_time_ms": {"type": "integer", "minimum": 1},
1078 "max_tokens": {"type": "integer", "minimum": 0},
1079 "max_cost_micros": {"type": "integer", "minimum": 0}
1080 },
1081 "required": [
1082 "max_concurrency",
1083 "max_agents",
1084 "max_steps",
1085 "max_retries",
1086 "max_nesting_depth",
1087 "wall_time_ms"
1088 ],
1089 "additionalProperties": false
1090 },
1091 "run_id": {"type": "string"},
1092 "since": {"type": "integer", "minimum": 0}
1093 },
1094 "required": ["action"],
1095 "additionalProperties": false
1096 })
1097 }
1098
1099 fn classify(&self, args: &Value) -> ToolClass {
1100 match args.get("action").and_then(Value::as_str) {
1101 Some("get" | "list" | "events") => ToolClass::READONLY_PARALLEL,
1102 _ => ToolClass::MUTATING_SERIAL,
1103 }
1104 }
1105
1106 async fn invoke(&self, args: Value, ctx: ToolCtx) -> Result<ToolOutcome, ToolError> {
1107 let input: WorkflowToolInput = serde_json::from_value(args)
1108 .map_err(|error| ToolError::InvalidArguments(error.to_string()))?;
1109 let session_id = ctx.session_id().ok_or_else(|| {
1110 ToolError::InvalidArguments("workflow_run requires a session".to_string())
1111 })?;
1112 let result = match input {
1113 WorkflowToolInput::Start {
1114 workflow_id,
1115 revision,
1116 args,
1117 budget,
1118 } => serde_json::to_value(public_workflow_snapshot(
1119 self.access
1120 .start_from_tool(session_id, &workflow_id, revision, args, budget)
1121 .await
1122 .map_err(workflow_tool_error)?,
1123 )),
1124 WorkflowToolInput::List {} => serde_json::to_value(
1125 self.access
1126 .list_for_session(session_id)
1127 .await
1128 .map_err(workflow_tool_error)?
1129 .into_iter()
1130 .map(public_workflow_snapshot)
1131 .collect::<Vec<_>>(),
1132 ),
1133 WorkflowToolInput::Get { run_id } => {
1134 let progress = self
1135 .access
1136 .progress_for_session(session_id, &run_id, u64::MAX)
1137 .await
1138 .map_err(workflow_tool_error)?;
1139 serde_json::to_value(public_workflow_snapshot(progress.snapshot))
1140 }
1141 WorkflowToolInput::Events { run_id, since } => {
1142 let progress = self
1143 .access
1144 .progress_for_session(session_id, &run_id, since)
1145 .await
1146 .map_err(workflow_tool_error)?;
1147 serde_json::to_value(
1148 progress
1149 .events
1150 .into_iter()
1151 .map(public_workflow_event)
1152 .collect::<Vec<_>>(),
1153 )
1154 }
1155 WorkflowToolInput::Cancel { run_id } => serde_json::to_value(public_workflow_snapshot(
1156 self.access
1157 .cancel_for_session(session_id, &run_id)
1158 .await
1159 .map_err(workflow_tool_error)?,
1160 )),
1161 WorkflowToolInput::Restart { run_id } => {
1162 serde_json::to_value(public_workflow_snapshot(
1163 self.access
1164 .restart_from_tool(session_id, &run_id)
1165 .await
1166 .map_err(workflow_tool_error)?,
1167 ))
1168 }
1169 }
1170 .map_err(|error| ToolError::Execution(error.to_string()))?;
1171 Ok(ToolOutcome::Completed(ToolResult::text(
1172 true,
1173 serde_json::to_string(&result)
1174 .map_err(|error| ToolError::Execution(error.to_string()))?,
1175 )))
1176 }
1177}
1178
1179fn workflow_tool_error(error: WorkflowRunError) -> ToolError {
1180 match error {
1181 WorkflowRunError::InvalidInput(_) => {
1182 ToolError::InvalidArguments("workflow input is invalid".to_string())
1183 }
1184 WorkflowRunError::Compile(_) => {
1185 ToolError::InvalidArguments("workflow definition is invalid".to_string())
1186 }
1187 WorkflowRunError::Preflight(_) => {
1188 ToolError::InvalidArguments("workflow preflight failed".to_string())
1189 }
1190 WorkflowRunError::Storage(details) => {
1191 tracing::error!(%details, "workflow tool storage unavailable");
1192 ToolError::Execution("workflow storage unavailable".to_string())
1193 }
1194 WorkflowRunError::NotFound => ToolError::Execution("workflow run not found".to_string()),
1195 WorkflowRunError::Terminal => {
1196 ToolError::Execution("workflow run is already terminal".to_string())
1197 }
1198 }
1199}
1200
1201#[cfg(test)]
1202mod tests {
1203 use super::*;
1204 use bamboo_agent_core::storage::Storage;
1205 use bamboo_agent_core::tools::{
1206 FunctionSchema, ToolCall, ToolExecutor, ToolResult, ToolSchema,
1207 };
1208 use bamboo_agent_core::Session;
1209 use bamboo_llm::protocol::{gemini::GeminiTool, ToProvider};
1210 use std::collections::HashMap;
1211 use tokio::sync::{RwLock, Semaphore};
1212
1213 #[derive(Default)]
1214 struct WorkflowTestStorage {
1215 sessions: RwLock<HashMap<String, Session>>,
1216 }
1217
1218 #[async_trait]
1219 impl Storage for WorkflowTestStorage {
1220 async fn save_session(&self, session: &Session) -> std::io::Result<()> {
1221 self.sessions
1222 .write()
1223 .await
1224 .insert(session.id.clone(), session.clone());
1225 Ok(())
1226 }
1227
1228 async fn load_session(&self, session_id: &str) -> std::io::Result<Option<Session>> {
1229 Ok(self.sessions.read().await.get(session_id).cloned())
1230 }
1231
1232 async fn delete_session(&self, session_id: &str) -> std::io::Result<bool> {
1233 Ok(self.sessions.write().await.remove(session_id).is_some())
1234 }
1235 }
1236
1237 struct WorkflowReadTool;
1238
1239 #[async_trait]
1240 impl ToolExecutor for WorkflowReadTool {
1241 async fn execute(
1242 &self,
1243 _call: &ToolCall,
1244 ) -> bamboo_agent_core::tools::executor::Result<ToolResult> {
1245 Ok(ToolResult::text(true, r#"{"ok":true}"#))
1246 }
1247
1248 fn list_tools(&self) -> Vec<ToolSchema> {
1249 vec![ToolSchema {
1250 schema_type: "function".to_string(),
1251 function: FunctionSchema {
1252 name: "Read".to_string(),
1253 description: "read".to_string(),
1254 parameters: serde_json::json!({"type":"object"}),
1255 },
1256 }]
1257 }
1258 }
1259
1260 struct BlockingWorkflowReadTool {
1261 entered: Arc<Semaphore>,
1262 }
1263
1264 #[async_trait]
1265 impl ToolExecutor for BlockingWorkflowReadTool {
1266 async fn execute(
1267 &self,
1268 _call: &ToolCall,
1269 ) -> bamboo_agent_core::tools::executor::Result<ToolResult> {
1270 self.entered.add_permits(1);
1271 std::future::pending().await
1272 }
1273
1274 fn list_tools(&self) -> Vec<ToolSchema> {
1275 WorkflowReadTool.list_tools()
1276 }
1277 }
1278
1279 async fn workflow_test_access_with_tools(
1280 tools: Arc<dyn ToolExecutor>,
1281 ) -> (
1282 WorkflowRunAccess,
1283 bamboo_engine::SessionRepository,
1284 tempfile::TempDir,
1285 ) {
1286 let directory = tempfile::tempdir().expect("tempdir");
1287 let skills_dir = directory.path().join("skills");
1288 let root = skills_dir.join("review-flow");
1289 std::fs::create_dir_all(&root).expect("workflow dir");
1290 std::fs::write(
1291 root.join("SKILL.md"),
1292 "---\nname: review-flow\ndescription: Review flow\n---\nRun review flow.\n",
1293 )
1294 .expect("skill");
1295 std::fs::write(
1296 root.join("workflow.yaml"),
1297 "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",
1298 )
1299 .expect("workflow");
1300 let skills = Arc::new(SkillManager::with_config(bamboo_skills::SkillStoreConfig {
1301 skills_dir,
1302 ..Default::default()
1303 }));
1304 skills.initialize().await.expect("skills initialize");
1305 let storage: Arc<dyn Storage> = Arc::new(WorkflowTestStorage::default());
1306 let persistence = Arc::new(bamboo_storage::LockedSessionStore::new(storage.clone()));
1307 let cache = Arc::new(dashmap::DashMap::new());
1308 let repo = bamboo_engine::SessionRepository::new(cache, storage, persistence);
1309 let access = WorkflowRunAccess::new(directory.path(), tools, skills, repo.clone())
1310 .await
1311 .expect("workflow access");
1312 (access, repo, directory)
1313 }
1314
1315 async fn workflow_test_access() -> (
1316 WorkflowRunAccess,
1317 bamboo_engine::SessionRepository,
1318 tempfile::TempDir,
1319 ) {
1320 workflow_test_access_with_tools(Arc::new(WorkflowReadTool)).await
1321 }
1322
1323 fn private_workflow_snapshot(status: WorkflowRunStatus) -> WorkflowRunSnapshot {
1324 let definition = WorkflowRunDefinition {
1325 workflow_schema: 1,
1326 id: "review-flow".to_string(),
1327 revision: 42,
1328 input_schema: json!({"private_schema": "PRIVATE-SCHEMA-SENTINEL"}),
1329 output_schema: Some(json!({"private_output": "PRIVATE-SCHEMA-SENTINEL"})),
1330 steps: vec![bamboo_domain::WorkflowStepDefinition {
1331 id: "inspect".to_string(),
1332 kind: WorkflowStepKind::Tool {
1333 tool: "PRIVATE-TOOL-SENTINEL".to_string(),
1334 args: json!({"credential": "PRIVATE-ARG-SENTINEL"}),
1335 capabilities: vec!["PRIVATE-CAPABILITY-SENTINEL".to_string()],
1336 },
1337 failure: bamboo_domain::FailurePolicy::FailFast,
1338 output_schema: Some(json!({"private": "PRIVATE-STEP-SCHEMA-SENTINEL"})),
1339 }],
1340 plan: WorkflowPlan::Retry {
1341 node: Box::new(WorkflowPlan::Map {
1342 source: bamboo_domain::ValueRef::Literal {
1343 value: json!("PRIVATE-BINDING-SENTINEL"),
1344 },
1345 item: "PRIVATE-ITEM-SENTINEL".to_string(),
1346 body: Box::new(WorkflowPlan::Step {
1347 step: "inspect".to_string(),
1348 }),
1349 }),
1350 max_attempts: 3,
1351 delay_ms: 987_654,
1352 },
1353 budgets: WorkflowBudgets {
1354 max_concurrency: 2,
1355 max_agents: 4,
1356 max_steps: 8,
1357 max_retries: 3,
1358 max_nesting_depth: 2,
1359 wall_time_ms: 10_000,
1360 max_tokens: Some(1_000),
1361 max_cost_micros: Some(2_000),
1362 },
1363 };
1364 let definition_bundle = WorkflowDefinitionBundle {
1365 publication_revision: 7,
1366 root_id: definition.id.clone(),
1367 root_revision: definition.revision,
1368 root_invocation_policy: json!({"private": "PRIVATE-POLICY-SENTINEL"}),
1369 definitions: BTreeMap::from([(
1370 WorkflowDefinitionBundle::key(&definition.id, definition.revision),
1371 definition.clone(),
1372 )]),
1373 };
1374 let now = chrono::Utc::now();
1375 WorkflowRunSnapshot {
1376 run_id: "public-run".to_string(),
1377 parent_run_id: Some("public-parent-run".to_string()),
1378 parent_step_id: Some("public-parent-step".to_string()),
1379 session_id: "public-session".to_string(),
1380 definition,
1381 definition_bundle,
1382 definition_bundle_hash: "public-bundle-hash".to_string(),
1383 validated_args: json!({"password": "PRIVATE-VALIDATED-ARG-SENTINEL"}),
1384 status,
1385 steps: BTreeMap::from([(
1386 "inspect".to_string(),
1387 WorkflowStepSnapshot {
1388 id: "inspect".to_string(),
1389 status: WorkflowStepStatus::Failed,
1390 input_hash: "PRIVATE-INPUT-HASH-SENTINEL".to_string(),
1391 output: Some(json!({"raw_tool_output": "PRIVATE-OUTPUT-SENTINEL"})),
1392 failure: Some(WorkflowFailure {
1393 code: WorkflowFailureCode::Storage,
1394 message: "/private/workspace/PRIVATE-DIAGNOSTIC-SENTINEL".to_string(),
1395 retryable: true,
1396 }),
1397 attempts: 2,
1398 },
1399 )]),
1400 usage: WorkflowBudgetUsage {
1401 steps: 1,
1402 retries: 1,
1403 agents: 3,
1404 tokens: 40,
1405 cost_micros: 50,
1406 },
1407 last_sequence: 9,
1408 output: Some(json!({"raw_run_output": "PRIVATE-RUN-OUTPUT-SENTINEL"})),
1409 failure: Some(WorkflowFailure {
1410 code: WorkflowFailureCode::ExecutionFailed,
1411 message: "credential PRIVATE-RUN-FAILURE-SENTINEL".to_string(),
1412 retryable: false,
1413 }),
1414 suspension: Some(WorkflowSuspensionContext::ToolApproval {
1415 step_id: "inspect".to_string(),
1416 tool: "PRIVATE-SUSPENSION-TOOL-SENTINEL".to_string(),
1417 tool_call_id: "PRIVATE-TOOL-CALL-SENTINEL".to_string(),
1418 }),
1419 created_at: now,
1420 updated_at: now,
1421 }
1422 }
1423
1424 #[test]
1425 fn public_workflow_snapshot_is_stable_metadata_only_for_every_status() {
1426 let statuses = [
1427 (WorkflowRunStatus::Queued, "queued"),
1428 (WorkflowRunStatus::Running, "running"),
1429 (WorkflowRunStatus::Suspended, "suspended"),
1430 (WorkflowRunStatus::Succeeded, "succeeded"),
1431 (WorkflowRunStatus::Failed, "failed"),
1432 (WorkflowRunStatus::Cancelled, "cancelled"),
1433 ];
1434
1435 for (status, wire_status) in statuses {
1436 let public =
1437 serde_json::to_value(public_workflow_snapshot(private_workflow_snapshot(status)))
1438 .expect("public snapshot serializes");
1439 let text = public.to_string();
1440 assert_eq!(public["status"], wire_status);
1441 assert_eq!(public["workflow_id"], "review-flow");
1442 assert_eq!(public["workflow_revision"], 42);
1443 assert_eq!(public["definition_bundle_hash"], "public-bundle-hash");
1444 assert_eq!(public["planned_steps"]["inspect"]["kind"], "tool");
1445 assert_eq!(public["plan"]["type"], "retry");
1446 assert_eq!(public["steps"]["inspect"]["attempts"], 2);
1447 assert_eq!(public["budget"]["max_steps"], 8);
1448 assert_eq!(public["usage"]["agents"], 3);
1449 assert_eq!(public["child_agent_count"], 3);
1450 assert_eq!(public["last_sequence"], 9);
1451 assert_eq!(public["failure"]["message"], "Workflow execution failed");
1452 assert_eq!(public["suspension"]["type"], "tool_approval");
1453 for internal_field in [
1454 "definition",
1455 "definition_bundle",
1456 "validated_args",
1457 "output",
1458 ] {
1459 assert!(
1460 public.get(internal_field).is_none(),
1461 "public snapshot exposed internal field {internal_field}: {text}"
1462 );
1463 }
1464
1465 for private in [
1466 "PRIVATE-",
1467 "validated_args",
1468 "input_hash",
1469 "output_schema",
1470 "root_invocation_policy",
1471 "delay_ms",
1472 "tool_call_id",
1473 "raw_tool_output",
1474 "raw_run_output",
1475 ] {
1476 assert!(
1477 !text.contains(private),
1478 "public snapshot leaked {private}: {text}"
1479 );
1480 }
1481 }
1482 }
1483
1484 #[test]
1485 fn public_workflow_events_preserve_sequence_and_drop_private_payloads() {
1486 let private = "PRIVATE-EVENT-SENTINEL";
1487 let failure = WorkflowFailure {
1488 code: WorkflowFailureCode::ExecutionFailed,
1489 message: format!("/private/workspace/{private}"),
1490 retryable: true,
1491 };
1492 let kinds = vec![
1493 WorkflowRunEventKind::RunQueued,
1494 WorkflowRunEventKind::RunStarted,
1495 WorkflowRunEventKind::Phase {
1496 name: private.to_string(),
1497 },
1498 WorkflowRunEventKind::StepQueued,
1499 WorkflowRunEventKind::StepStarted,
1500 WorkflowRunEventKind::StepSuspended {
1501 reason: private.to_string(),
1502 },
1503 WorkflowRunEventKind::StepCompleted {
1504 output: json!({"raw": private}),
1505 },
1506 WorkflowRunEventKind::StepFailed {
1507 failure: failure.clone(),
1508 },
1509 WorkflowRunEventKind::StepCancelled,
1510 WorkflowRunEventKind::StepSkipped {
1511 reason: private.to_string(),
1512 },
1513 WorkflowRunEventKind::RunSuspended {
1514 reason: private.to_string(),
1515 },
1516 WorkflowRunEventKind::RunSucceeded {
1517 output: json!({"raw": private}),
1518 },
1519 WorkflowRunEventKind::RunFailed { failure },
1520 WorkflowRunEventKind::RunCancelled,
1521 ];
1522 let at = chrono::Utc::now();
1523 let events = kinds
1524 .into_iter()
1525 .enumerate()
1526 .map(|(index, kind)| {
1527 public_workflow_event(WorkflowRunEvent {
1528 run_id: "public-run".to_string(),
1529 sequence: index as u64 + 1,
1530 at,
1531 step_id: Some("inspect".to_string()),
1532 kind,
1533 })
1534 })
1535 .collect::<Vec<_>>();
1536 let public = serde_json::to_value(&events).expect("public events serialize");
1537 let text = public.to_string();
1538
1539 assert!(!text.contains(private), "public events leaked: {text}");
1540 assert_eq!(public[2]["name"], "workflow_progressed");
1541 assert_eq!(public[6]["type"], "step_completed");
1542 assert!(public[6].get("output").is_none());
1543 assert_eq!(public[7]["failure"]["message"], "Workflow execution failed");
1544 assert_eq!(public[11]["type"], "run_succeeded");
1545 assert!(public[11].get("output").is_none());
1546 assert_eq!(
1547 events
1548 .iter()
1549 .map(|event| event.sequence)
1550 .collect::<Vec<_>>(),
1551 (1..=14).collect::<Vec<_>>()
1552 );
1553 }
1554
1555 #[test]
1556 fn workflow_tool_errors_never_expose_backend_diagnostics() {
1557 let sentinel = "/private/workspace/credentials-PRIVATE-SENTINEL";
1558 let cases = [
1559 WorkflowRunError::Storage(sentinel.to_string()),
1560 WorkflowRunError::InvalidInput(sentinel.to_string()),
1561 WorkflowRunError::Preflight(sentinel.to_string()),
1562 ];
1563 for error in cases {
1564 assert!(!workflow_tool_error(error).to_string().contains(sentinel));
1565 }
1566 assert!(!workflow_tool_error(WorkflowRunError::Compile(
1567 bamboo_domain::WorkflowCompileError::InvalidSchema(sentinel.to_string())
1568 ))
1569 .to_string()
1570 .contains(sentinel));
1571 }
1572
1573 fn canonical_workflow_run_schema() -> Value {
1574 json!({
1575 "type": "object",
1576 "properties": {
1577 "action": {
1578 "type": "string",
1579 "enum": ["start", "list", "get", "events", "cancel", "restart"]
1580 },
1581 "workflow_id": {"type": "string", "minLength": 1},
1582 "revision": {"type": "integer", "minimum": 1},
1583 "args": {"type": "object", "default": {}},
1584 "budget": {
1585 "type": "object",
1586 "properties": {
1587 "max_concurrency": {"type": "integer", "minimum": 1},
1588 "max_agents": {"type": "integer", "minimum": 0},
1589 "max_steps": {"type": "integer", "minimum": 1},
1590 "max_retries": {"type": "integer", "minimum": 0},
1591 "max_nesting_depth": {"type": "integer", "minimum": 1},
1592 "wall_time_ms": {"type": "integer", "minimum": 1},
1593 "max_tokens": {"type": "integer", "minimum": 0},
1594 "max_cost_micros": {"type": "integer", "minimum": 0}
1595 },
1596 "required": [
1597 "max_concurrency",
1598 "max_agents",
1599 "max_steps",
1600 "max_retries",
1601 "max_nesting_depth",
1602 "wall_time_ms"
1603 ],
1604 "additionalProperties": false
1605 },
1606 "run_id": {"type": "string"},
1607 "since": {"type": "integer", "minimum": 0}
1608 },
1609 "required": ["action"],
1610 "additionalProperties": false
1611 })
1612 }
1613
1614 #[tokio::test]
1615 async fn workflow_run_schema_is_flat_complete_and_canonical() {
1616 let (access, _, _) = workflow_test_access().await;
1617 let schema = WorkflowRunTool::new(access).parameters_schema();
1618
1619 for combinator in ["oneOf", "anyOf", "allOf"] {
1620 assert!(
1621 schema.get(combinator).is_none(),
1622 "workflow_run must not advertise root {combinator}"
1623 );
1624 }
1625 assert_eq!(schema, canonical_workflow_run_schema());
1626 }
1627
1628 #[tokio::test]
1629 async fn workflow_run_schema_survives_openai_sanitization_with_all_properties() {
1630 let (access, _, _) = workflow_test_access().await;
1631 let schema = WorkflowRunTool::new(access).parameters_schema();
1632 let sanitized =
1633 bamboo_llm::providers::common::tool_schema::sanitize_openai_function_parameters_schema(
1634 &schema,
1635 );
1636
1637 let properties = sanitized["properties"]
1638 .as_object()
1639 .expect("sanitized workflow_run properties");
1640 assert!(!properties.is_empty());
1641 assert_eq!(properties.len(), 7);
1642 assert_eq!(sanitized, canonical_workflow_run_schema());
1643 }
1644
1645 #[tokio::test]
1646 async fn workflow_run_schema_reaches_gemini_unchanged() {
1647 let (access, _, _) = workflow_test_access().await;
1648 let direct = WorkflowRunTool::new(access).to_schema();
1649 let gemini: GeminiTool = direct.to_provider().expect("Gemini tool conversion");
1650 let declaration = gemini
1651 .function_declarations
1652 .first()
1653 .expect("workflow_run declaration");
1654
1655 assert_eq!(declaration.name, "workflow_run");
1656 assert_eq!(
1657 declaration.parameters_json_schema.as_ref(),
1658 Some(&canonical_workflow_run_schema())
1659 );
1660 assert!(declaration.parameters.is_none());
1661 }
1662
1663 #[tokio::test]
1664 async fn workflow_run_enforces_opt_in_tightens_budget_lists_and_isolates_sessions() {
1665 let (access, repo, directory) = workflow_test_access().await;
1666 let workspace = directory.path().join("workspace");
1667 std::fs::create_dir_all(&workspace).expect("workspace");
1668 let mut session = Session::new("workflow-session", "model");
1669 session.workspace = Some(workspace.to_string_lossy().into_owned());
1670 repo.save(&mut session).await.expect("save session");
1671
1672 let denied = access
1673 .start_from_tool("workflow-session", "review-flow", 42, json!({}), None)
1674 .await
1675 .expect_err("model start defaults off without session opt-in");
1676 assert!(denied.to_string().contains("opt-in"));
1677
1678 session.metadata.insert(
1679 bamboo_skills::WORKFLOW_ORCHESTRATION_OPT_IN_METADATA_KEY.to_string(),
1680 "true".to_string(),
1681 );
1682 repo.save(&mut session).await.expect("save opt-in");
1683 let requested = WorkflowBudgets {
1684 max_concurrency: 1,
1685 max_agents: 0,
1686 max_steps: 2,
1687 max_retries: 0,
1688 max_nesting_depth: 1,
1689 wall_time_ms: 5_000,
1690 max_tokens: Some(500),
1691 max_cost_micros: Some(500),
1692 };
1693 let started = access
1694 .start_from_tool(
1695 "workflow-session",
1696 "review-flow",
1697 42,
1698 json!({}),
1699 Some(requested.clone()),
1700 )
1701 .await
1702 .expect("opted-in model start");
1703 assert_eq!(started.definition.budgets, requested);
1704 let listed = access
1705 .list_for_session("workflow-session")
1706 .await
1707 .expect("session run list");
1708 assert_eq!(listed.len(), 1);
1709 assert_eq!(listed[0].run_id, started.run_id);
1710 let progress = access
1711 .progress_for_session("workflow-session", &started.run_id, 0)
1712 .await
1713 .expect("run events");
1714 assert!(progress
1715 .events
1716 .first()
1717 .is_some_and(|event| event.kind == bamboo_domain::WorkflowRunEventKind::RunQueued));
1718 let completed = tokio::time::timeout(std::time::Duration::from_secs(2), async {
1719 loop {
1720 let progress = access
1721 .progress_for_session("workflow-session", &started.run_id, 0)
1722 .await
1723 .expect("terminal run progress");
1724 if progress.snapshot.status.is_terminal() {
1725 break progress;
1726 }
1727 tokio::time::sleep(std::time::Duration::from_millis(10)).await;
1728 }
1729 })
1730 .await
1731 .expect("workflow reaches terminal state");
1732 assert_eq!(
1733 completed.snapshot.status,
1734 bamboo_domain::WorkflowRunStatus::Succeeded
1735 );
1736 assert_eq!(
1737 completed
1738 .events
1739 .iter()
1740 .map(|event| event.sequence)
1741 .collect::<Vec<_>>(),
1742 (1..=7).collect::<Vec<_>>()
1743 );
1744 assert!(matches!(
1745 completed.events.as_slice(),
1746 [
1747 bamboo_domain::WorkflowRunEvent {
1748 kind: bamboo_domain::WorkflowRunEventKind::RunQueued,
1749 ..
1750 },
1751 bamboo_domain::WorkflowRunEvent {
1752 kind: bamboo_domain::WorkflowRunEventKind::RunStarted,
1753 ..
1754 },
1755 bamboo_domain::WorkflowRunEvent {
1756 kind: bamboo_domain::WorkflowRunEventKind::Phase { ref name },
1757 ..
1758 },
1759 bamboo_domain::WorkflowRunEvent {
1760 kind: bamboo_domain::WorkflowRunEventKind::StepQueued,
1761 ..
1762 },
1763 bamboo_domain::WorkflowRunEvent {
1764 kind: bamboo_domain::WorkflowRunEventKind::StepStarted,
1765 ..
1766 },
1767 bamboo_domain::WorkflowRunEvent {
1768 kind: bamboo_domain::WorkflowRunEventKind::StepCompleted { .. },
1769 ..
1770 },
1771 bamboo_domain::WorkflowRunEvent {
1772 kind: bamboo_domain::WorkflowRunEventKind::RunSucceeded { .. },
1773 ..
1774 }
1775 ] if name == "step_reserved"
1776 ));
1777 assert_eq!(
1778 completed.snapshot.last_sequence,
1779 completed.events.last().expect("terminal event").sequence
1780 );
1781
1782 let invalid_budget = WorkflowBudgets {
1783 max_steps: 0,
1784 ..requested.clone()
1785 };
1786 assert!(matches!(
1787 access
1788 .start_from_tool(
1789 "workflow-session",
1790 "review-flow",
1791 42,
1792 json!({}),
1793 Some(invalid_budget),
1794 )
1795 .await,
1796 Err(WorkflowRunError::InvalidInput(_))
1797 ));
1798
1799 let workflow_path = directory.path().join("skills/review-flow/workflow.yaml");
1802 let original_workflow = std::fs::read_to_string(&workflow_path).expect("workflow yaml");
1803 std::fs::write(
1804 &workflow_path,
1805 original_workflow.replace(
1806 "invocation_policy: {explicit: true, automatic: true}",
1807 "invocation_policy: {explicit: true, automatic: false}",
1808 ),
1809 )
1810 .expect("disable automatic live policy");
1811 access.skills.store().reload().await.expect("reload policy");
1812 let restart = access
1813 .restart_from_tool("workflow-session", &started.run_id)
1814 .await
1815 .expect_err("succeeded runs are terminal");
1816 assert!(matches!(restart, WorkflowRunError::Terminal));
1817
1818 let mut other = Session::new("other-session", "model");
1819 other.workspace = Some(workspace.to_string_lossy().into_owned());
1820 repo.save(&mut other).await.expect("save other session");
1821 assert!(access
1822 .list_for_session("other-session")
1823 .await
1824 .expect("isolated list")
1825 .is_empty());
1826 assert!(matches!(
1827 access
1828 .progress_for_session("other-session", &started.run_id, 0)
1829 .await,
1830 Err(WorkflowRunError::NotFound)
1831 ));
1832
1833 let run_id = started.run_id.clone();
1837 let skills = access.skills.clone();
1838 drop(access);
1839 let reopened = WorkflowRunAccess::new(
1840 directory.path(),
1841 Arc::new(WorkflowReadTool),
1842 skills,
1843 repo.clone(),
1844 )
1845 .await
1846 .expect("reopen workflow adapter");
1847 let reconnected = reopened
1848 .progress_for_session("workflow-session", &run_id, 4)
1849 .await
1850 .expect("durable reconnect");
1851 assert_eq!(reconnected.snapshot.status, WorkflowRunStatus::Succeeded);
1852 assert_eq!(
1853 reconnected
1854 .events
1855 .iter()
1856 .map(|event| event.sequence)
1857 .collect::<Vec<_>>(),
1858 vec![5, 6, 7]
1859 );
1860 assert_eq!(
1861 public_workflow_snapshot(reconnected.snapshot).last_sequence,
1862 7
1863 );
1864 assert!(matches!(
1865 public_workflow_event(reconnected.events.last().cloned().expect("tail event")).kind,
1866 PublicWorkflowRunEventKind::RunSucceeded
1867 ));
1868 assert_eq!(
1869 reopened
1870 .list_for_session("workflow-session")
1871 .await
1872 .expect("reconnected list")
1873 .len(),
1874 1
1875 );
1876 }
1877
1878 #[tokio::test]
1879 async fn workflow_cancel_is_idempotent_and_reconnects_to_one_terminal_event() {
1880 let entered = Arc::new(Semaphore::new(0));
1881 let (access, repo, directory) =
1882 workflow_test_access_with_tools(Arc::new(BlockingWorkflowReadTool {
1883 entered: entered.clone(),
1884 }))
1885 .await;
1886 let workspace = directory.path().join("workspace");
1887 std::fs::create_dir_all(&workspace).expect("workspace");
1888 let mut session = Session::new("cancel-session", "model");
1889 session.workspace = Some(workspace.to_string_lossy().into_owned());
1890 repo.save(&mut session).await.expect("save session");
1891
1892 let started = access
1893 .start("cancel-session", "review-flow", 42, json!({}), None)
1894 .await
1895 .expect("start blocking run");
1896 let _entered = tokio::time::timeout(std::time::Duration::from_secs(2), entered.acquire())
1897 .await
1898 .expect("blocking step entered")
1899 .expect("semaphore open");
1900
1901 let first = access
1902 .cancel_for_session("cancel-session", &started.run_id)
1903 .await
1904 .expect("first cancel");
1905 let second = access
1906 .cancel_for_session("cancel-session", &started.run_id)
1907 .await
1908 .expect("idempotent cancel");
1909 assert_eq!(first.status, WorkflowRunStatus::Cancelled);
1910 assert_eq!(second.status, WorkflowRunStatus::Cancelled);
1911 assert_eq!(first.last_sequence, second.last_sequence);
1912
1913 let progress = access
1914 .progress_for_session("cancel-session", &started.run_id, 0)
1915 .await
1916 .expect("cancelled journal");
1917 assert_eq!(progress.snapshot.status, WorkflowRunStatus::Cancelled);
1918 assert_eq!(progress.snapshot.last_sequence, first.last_sequence);
1919 assert_eq!(
1920 progress
1921 .events
1922 .iter()
1923 .filter(|event| matches!(event.kind, WorkflowRunEventKind::RunCancelled))
1924 .count(),
1925 1
1926 );
1927 assert!(!progress
1928 .events
1929 .iter()
1930 .any(|event| matches!(event.kind, WorkflowRunEventKind::RunSucceeded { .. })));
1931 assert_eq!(
1932 progress
1933 .events
1934 .iter()
1935 .map(|event| event.sequence)
1936 .collect::<Vec<_>>(),
1937 (1..=progress.snapshot.last_sequence).collect::<Vec<_>>()
1938 );
1939 let tail = access
1940 .progress_for_session(
1941 "cancel-session",
1942 &started.run_id,
1943 progress.snapshot.last_sequence,
1944 )
1945 .await
1946 .expect("terminal reconnect tail");
1947 assert!(tail.events.is_empty());
1948 assert!(matches!(
1949 public_workflow_snapshot(second).status,
1950 WorkflowRunStatus::Cancelled
1951 ));
1952 }
1953
1954 #[test]
1955 fn tool_input_rejects_security_context_spoofing() {
1956 let error = serde_json::from_value::<WorkflowToolInput>(json!({
1957 "action": "start",
1958 "workflow_id": "safe",
1959 "revision": 1,
1960 "workspace_trusted": true
1961 }))
1962 .unwrap_err();
1963 assert!(error.to_string().contains("unknown field"));
1964 }
1965
1966 #[test]
1967 fn tool_input_rejects_fields_from_other_actions() {
1968 let invalid = [
1969 json!({"action": "list", "run_id": "run-1"}),
1970 json!({"action": "get", "run_id": "run-1", "since": 1}),
1971 json!({"action": "events", "run_id": "run-1", "workflow_id": "flow"}),
1972 json!({"action": "cancel", "run_id": "run-1", "revision": 1}),
1973 json!({"action": "restart", "run_id": "run-1", "budget": {}}),
1974 json!({
1975 "action": "start",
1976 "workflow_id": "flow",
1977 "revision": 1,
1978 "run_id": "run-1"
1979 }),
1980 ];
1981
1982 for input in invalid {
1983 let error = serde_json::from_value::<WorkflowToolInput>(input.clone())
1984 .expect_err("action-specific fields must remain authoritative at runtime");
1985 assert!(
1986 error.to_string().contains("unknown field"),
1987 "unexpected error for {input}: {error}"
1988 );
1989 }
1990
1991 assert!(matches!(
1992 serde_json::from_value::<WorkflowToolInput>(json!({"action": "list"}))
1993 .expect("fieldless list action"),
1994 WorkflowToolInput::List {}
1995 ));
1996 }
1997
1998 #[tokio::test]
1999 async fn omitted_start_args_and_zero_budgets_match_schema() {
2000 let WorkflowToolInput::Start { args, .. } =
2001 serde_json::from_value::<WorkflowToolInput>(json!({
2002 "action": "start",
2003 "workflow_id": "safe",
2004 "revision": 1
2005 }))
2006 .expect("tool args default")
2007 else {
2008 panic!("start input")
2009 };
2010 assert_eq!(args, json!({}));
2011 let http: crate::handlers::workflow_runs::StartWorkflowRunRequest =
2012 serde_json::from_value(json!({"workflow_id":"safe", "revision":1}))
2013 .expect("http args default");
2014 assert_eq!(http.args, json!({}));
2015
2016 let (access, _, _) = workflow_test_access().await;
2017 let schema = WorkflowRunTool { access }.parameters_schema();
2018 let properties = &schema["properties"];
2019 assert_eq!(properties["args"]["default"], json!({}));
2020 assert_eq!(
2021 properties["budget"]["properties"]["max_agents"]["minimum"],
2022 0
2023 );
2024 assert_eq!(
2025 properties["budget"]["properties"]["max_retries"]["minimum"],
2026 0
2027 );
2028 assert_eq!(
2029 properties["budget"]["properties"]["max_tokens"]["minimum"],
2030 0
2031 );
2032 assert_eq!(
2033 properties["budget"]["properties"]["max_cost_micros"]["minimum"],
2034 0
2035 );
2036 }
2037
2038 #[tokio::test]
2039 async fn run_index_updates_are_concurrent_and_never_evict_at_active_capacity() {
2040 let (access, repo, _) = workflow_test_access().await;
2041 let mut session = Session::new("run-index", "model");
2042 repo.save(&mut session).await.expect("seed session");
2043 let results = futures::future::join_all((0..32).map(|index| {
2044 let access = access.clone();
2045 async move {
2046 let run_id = format!("run-{index}");
2047 access.remember_run_id("run-index", &run_id).await
2048 }
2049 }))
2050 .await;
2051 assert!(results.into_iter().all(|result| result.is_ok()));
2052 let concurrent = repo
2053 .try_load("run-index")
2054 .await
2055 .expect("load")
2056 .expect("session");
2057 let ids = serde_json::from_str::<Vec<String>>(
2058 concurrent
2059 .metadata
2060 .get(bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY)
2061 .expect("run ids"),
2062 )
2063 .expect("ids json");
2064 assert_eq!(ids.len(), 32);
2065 assert_eq!(ids.iter().collect::<BTreeSet<_>>().len(), 32);
2066
2067 let capacity_ids = (0..MAX_WORKFLOW_RUN_IDS_PER_SESSION)
2068 .map(|index| format!("active-{index}"))
2069 .collect::<Vec<_>>();
2070 repo.update_runtime_session(
2071 "run-index",
2072 &[bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY],
2073 {
2074 let capacity_ids = capacity_ids.clone();
2075 move |session| {
2076 session.metadata.insert(
2077 bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY.to_string(),
2078 serde_json::to_string(&capacity_ids).expect("ids json"),
2079 );
2080 }
2081 },
2082 )
2083 .await
2084 .expect("fill index")
2085 .expect("session");
2086 assert!(matches!(
2087 access.remember_run_id("run-index", "new-run").await,
2088 Err(WorkflowRunError::Storage(_))
2089 ));
2090 let retained = repo
2091 .try_load("run-index")
2092 .await
2093 .expect("load")
2094 .expect("session");
2095 let retained = serde_json::from_str::<Vec<String>>(
2096 retained
2097 .metadata
2098 .get(bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY)
2099 .expect("run ids"),
2100 )
2101 .expect("ids json");
2102 assert_eq!(
2103 retained, capacity_ids,
2104 "oldest active id must not be evicted"
2105 );
2106 }
2107
2108 #[tokio::test]
2109 async fn real_model_workflow_run_index_survives_tool_result_and_final_session_save() {
2110 use bamboo_agent_core::storage::AttachmentReader;
2111 use bamboo_engine::{Agent, ExecuteRequestBuilder};
2112 use bamboo_llm::{LLMChunk, LLMProvider, LLMStream};
2113 use futures::stream;
2114 use tokio::sync::Mutex;
2115 use tokio_util::sync::CancellationToken;
2116
2117 struct NoAttachments;
2118 #[async_trait]
2119 impl AttachmentReader for NoAttachments {
2120 async fn read_attachment(
2121 &self,
2122 _session_id: &str,
2123 _attachment_id: &str,
2124 ) -> std::io::Result<Option<(Vec<u8>, String)>> {
2125 Ok(None)
2126 }
2127 }
2128 struct QueueProvider {
2129 queue: Mutex<Vec<Vec<bamboo_llm::provider::Result<LLMChunk>>>>,
2130 }
2131 #[async_trait]
2132 impl LLMProvider for QueueProvider {
2133 async fn chat_stream(
2134 &self,
2135 _messages: &[bamboo_agent_core::Message],
2136 _tools: &[ToolSchema],
2137 _max_output_tokens: Option<u32>,
2138 _model: &str,
2139 ) -> bamboo_llm::provider::Result<LLMStream> {
2140 Ok(Box::pin(stream::iter(self.queue.lock().await.remove(0))))
2141 }
2142 }
2143
2144 let (access, repo, directory) = workflow_test_access().await;
2145 let session_id = "real-model-workflow-run";
2146 let mut session = Session::new(session_id, "test-model");
2147 session.metadata.insert(
2148 bamboo_skills::WORKFLOW_ORCHESTRATION_OPT_IN_METADATA_KEY.to_string(),
2149 "true".to_string(),
2150 );
2151 session
2152 .metadata
2153 .insert("external.metadata".to_string(), "preserve".to_string());
2154 session.add_message(bamboo_agent_core::Message::system("system"));
2155 session.add_message(bamboo_agent_core::Message::user("run review workflow"));
2156 repo.save(&mut session).await.expect("seed session");
2157 let call = ToolCall {
2158 id: "call-workflow-run".to_string(),
2159 tool_type: "function".to_string(),
2160 function: bamboo_agent_core::tools::FunctionCall {
2161 name: "workflow_run".to_string(),
2162 arguments: json!({
2163 "action":"start",
2164 "workflow_id":"review-flow",
2165 "revision":42
2166 })
2167 .to_string(),
2168 },
2169 };
2170 let provider = Arc::new(QueueProvider {
2171 queue: Mutex::new(vec![
2172 vec![Ok(LLMChunk::ToolCalls(vec![call])), Ok(LLMChunk::Done)],
2173 vec![Ok(LLMChunk::Token("done".to_string())), Ok(LLMChunk::Done)],
2174 ]),
2175 });
2176 let tools = Arc::new(
2177 bamboo_tools::BuiltinToolExecutorBuilder::new()
2178 .with_tool(WorkflowRunTool::new(access.clone()))
2179 .expect("workflow tool")
2180 .build(),
2181 );
2182 let metrics = bamboo_metrics::MetricsCollector::spawn(
2183 Arc::new(bamboo_metrics::SqliteMetricsStorage::new(
2184 directory.path().join("runner-metrics.db"),
2185 )),
2186 7,
2187 );
2188 let agent = Agent::builder()
2189 .storage(repo.storage().clone())
2190 .persistence(Arc::new(repo.clone()))
2191 .attachment_reader(Arc::new(NoAttachments))
2192 .skill_manager(access.skills.clone())
2193 .metrics_collector(metrics)
2194 .config(Arc::new(RwLock::new(bamboo_llm::Config::default())))
2195 .provider(provider)
2196 .default_tools(tools)
2197 .build()
2198 .expect("agent");
2199 let (event_tx, _event_rx) = tokio::sync::mpsc::channel(128);
2200 agent
2201 .execute(
2202 &mut session,
2203 ExecuteRequestBuilder::new(
2204 "run review workflow",
2205 event_tx,
2206 CancellationToken::new(),
2207 )
2208 .model("test-model")
2209 .build(),
2210 )
2211 .await
2212 .expect("real model workflow run");
2213
2214 let saved = repo
2215 .storage()
2216 .load_session(session_id)
2217 .await
2218 .expect("load")
2219 .expect("saved");
2220 assert_eq!(
2221 saved.metadata.get("external.metadata").map(String::as_str),
2222 Some("preserve")
2223 );
2224 assert!(saved
2225 .metadata
2226 .contains_key(bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY));
2227 let listed = access
2228 .list_for_session(session_id)
2229 .await
2230 .expect("list after final save");
2231 assert_eq!(listed.len(), 1);
2232 assert!(saved.messages.iter().any(|message| {
2233 message.tool_calls.as_ref().is_some_and(|calls| {
2234 calls
2235 .iter()
2236 .any(|call| call.function.name == "workflow_run")
2237 })
2238 }));
2239 }
2240
2241 #[tokio::test]
2242 async fn http_start_survives_concurrent_stale_runner_final_save_and_server_restart() {
2243 use bamboo_agent_core::storage::AttachmentReader;
2244 use bamboo_engine::{Agent, ExecuteRequestBuilder};
2245 use bamboo_llm::{LLMChunk, LLMProvider, LLMStream};
2246 use futures::stream;
2247 use tokio::sync::{oneshot, Mutex};
2248 use tokio_util::sync::CancellationToken;
2249
2250 struct NoAttachments;
2251 #[async_trait]
2252 impl AttachmentReader for NoAttachments {
2253 async fn read_attachment(
2254 &self,
2255 _session_id: &str,
2256 _attachment_id: &str,
2257 ) -> std::io::Result<Option<(Vec<u8>, String)>> {
2258 Ok(None)
2259 }
2260 }
2261
2262 struct PausingProvider {
2263 entered: Mutex<Option<oneshot::Sender<()>>>,
2264 resume: Mutex<Option<oneshot::Receiver<()>>>,
2265 }
2266 #[async_trait]
2267 impl LLMProvider for PausingProvider {
2268 async fn chat_stream(
2269 &self,
2270 _messages: &[bamboo_agent_core::Message],
2271 _tools: &[ToolSchema],
2272 _max_output_tokens: Option<u32>,
2273 _model: &str,
2274 ) -> bamboo_llm::provider::Result<LLMStream> {
2275 if let Some(entered) = self.entered.lock().await.take() {
2276 let _ = entered.send(());
2277 }
2278 if let Some(resume) = self.resume.lock().await.take() {
2279 let _ = resume.await;
2280 }
2281 Ok(Box::pin(stream::iter(vec![
2282 Ok(LLMChunk::Token("done".to_string())),
2283 Ok(LLMChunk::Done),
2284 ])))
2285 }
2286 }
2287
2288 let (access, repo, directory) = workflow_test_access().await;
2289 let session_id = "http-start-concurrent-runner-save";
2290 let mut session = Session::new(session_id, "test-model");
2291 session.add_message(bamboo_agent_core::Message::system("system"));
2292 session.add_message(bamboo_agent_core::Message::user("keep running"));
2293 repo.save(&mut session).await.expect("seed session");
2294 let mut runner_session = repo
2295 .try_load(session_id)
2296 .await
2297 .expect("load runner session")
2298 .expect("runner session");
2299
2300 let (entered_tx, entered_rx) = oneshot::channel();
2301 let (resume_tx, resume_rx) = oneshot::channel();
2302 let provider = Arc::new(PausingProvider {
2303 entered: Mutex::new(Some(entered_tx)),
2304 resume: Mutex::new(Some(resume_rx)),
2305 });
2306 let metrics = bamboo_metrics::MetricsCollector::spawn(
2307 Arc::new(bamboo_metrics::SqliteMetricsStorage::new(
2308 directory.path().join("http-runner-metrics.db"),
2309 )),
2310 7,
2311 );
2312 let agent = Agent::builder()
2313 .storage(repo.storage().clone())
2314 .persistence(Arc::new(repo.clone()))
2315 .attachment_reader(Arc::new(NoAttachments))
2316 .skill_manager(access.skills.clone())
2317 .metrics_collector(metrics)
2318 .config(Arc::new(RwLock::new(bamboo_llm::Config::default())))
2319 .provider(provider)
2320 .default_tools(Arc::new(
2321 bamboo_tools::BuiltinToolExecutorBuilder::new().build(),
2322 ))
2323 .build()
2324 .expect("agent");
2325 let (event_tx, _event_rx) = tokio::sync::mpsc::channel(128);
2326 let runner = tokio::spawn(async move {
2327 agent
2328 .execute(
2329 &mut runner_session,
2330 ExecuteRequestBuilder::new("keep running", event_tx, CancellationToken::new())
2331 .model("test-model")
2332 .build(),
2333 )
2334 .await
2335 });
2336 tokio::time::timeout(std::time::Duration::from_secs(2), entered_rx)
2337 .await
2338 .expect("runner enters model round")
2339 .expect("runner entry signal");
2340
2341 let started = access
2342 .start(session_id, "review-flow", 42, json!({}), None)
2343 .await
2344 .expect("HTTP-equivalent explicit start");
2345 let durable_during_round = repo
2346 .storage()
2347 .load_session(session_id)
2348 .await
2349 .expect("load during round")
2350 .expect("session during round");
2351 assert!(durable_during_round
2352 .metadata
2353 .get(bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY)
2354 .is_some_and(|raw| raw.contains(&started.run_id)));
2355
2356 resume_tx.send(()).expect("resume runner");
2357 tokio::time::timeout(std::time::Duration::from_secs(2), runner)
2358 .await
2359 .expect("runner completes")
2360 .expect("runner task")
2361 .expect("runner execution");
2362
2363 let durable_after_final_save = repo
2364 .storage()
2365 .load_session(session_id)
2366 .await
2367 .expect("load after final save")
2368 .expect("saved session");
2369 assert!(durable_after_final_save
2370 .metadata
2371 .get(bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY)
2372 .is_some_and(|raw| raw.contains(&started.run_id)));
2373
2374 tokio::time::timeout(std::time::Duration::from_secs(2), async {
2375 loop {
2376 let progress = access
2377 .progress_for_session(session_id, &started.run_id, u64::MAX)
2378 .await
2379 .expect("workflow progress before restart");
2380 if progress.snapshot.status.is_terminal()
2381 && !access.engine.is_run_active(&started.run_id)
2382 {
2383 break;
2384 }
2385 tokio::task::yield_now().await;
2386 }
2387 })
2388 .await
2389 .expect("workflow reaches terminal state before restart");
2390
2391 let skills = access.skills.clone();
2392 drop(access);
2393 let restarted =
2394 WorkflowRunAccess::new(directory.path(), Arc::new(WorkflowReadTool), skills, repo)
2395 .await
2396 .expect("restart workflow access");
2397 let listed = restarted
2398 .list_for_session(session_id)
2399 .await
2400 .expect("list after restart");
2401 assert_eq!(listed.len(), 1);
2402 assert_eq!(listed[0].run_id, started.run_id);
2403 }
2404
2405 #[tokio::test]
2406 async fn production_policy_allows_read_without_fabricating_workspace_trust() {
2407 let read = BTreeSet::from(["read".to_string()]);
2408 assert_eq!(
2409 ServerWorkflowPolicy
2410 .authorize(
2411 "session",
2412 &WorkflowPolicyTarget::Tool("read_file".to_string()),
2413 &read,
2414 false,
2415 )
2416 .await,
2417 PermissionDecision::Allow
2418 );
2419
2420 let write = BTreeSet::from(["write".to_string()]);
2421 assert!(matches!(
2422 ServerWorkflowPolicy
2423 .authorize(
2424 "session",
2425 &WorkflowPolicyTarget::Tool("write_file".to_string()),
2426 &write,
2427 false,
2428 )
2429 .await,
2430 PermissionDecision::Deny(_)
2431 ));
2432
2433 for hostile_target in [
2434 "Write",
2435 "write_file",
2436 "WebFetch",
2437 "mcp::remote_tool",
2438 "Bash",
2439 ] {
2440 for claimed in [BTreeSet::new(), read.clone()] {
2441 assert!(matches!(
2442 ServerWorkflowPolicy
2443 .authorize(
2444 "session",
2445 &WorkflowPolicyTarget::Tool(hostile_target.to_string()),
2446 &claimed,
2447 false,
2448 )
2449 .await,
2450 PermissionDecision::Deny(_)
2451 ));
2452 }
2453 }
2454
2455 assert_eq!(
2456 ServerWorkflowPolicy
2457 .authorize(
2458 "session",
2459 &WorkflowPolicyTarget::Workflow {
2460 id: "nested-review".to_string(),
2461 revision: 1,
2462 },
2463 &BTreeSet::new(),
2464 false,
2465 )
2466 .await,
2467 PermissionDecision::Allow
2468 );
2469 }
2470
2471 #[test]
2472 fn pinned_bundle_limits_reject_oversized_runs_before_engine_start() {
2473 assert!(enforce_pinned_bundle_limits(
2474 MAX_PINNED_DEFINITIONS_PER_RUN,
2475 MAX_PINNED_BUNDLE_BYTES_PER_RUN
2476 )
2477 .is_ok());
2478 assert!(enforce_pinned_bundle_limits(
2479 MAX_PINNED_DEFINITIONS_PER_RUN + 1,
2480 MAX_PINNED_BUNDLE_BYTES_PER_RUN
2481 )
2482 .is_err());
2483 assert!(enforce_pinned_bundle_limits(
2484 MAX_PINNED_DEFINITIONS_PER_RUN,
2485 MAX_PINNED_BUNDLE_BYTES_PER_RUN + 1
2486 )
2487 .is_err());
2488 }
2489}