1use std::collections::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, WorkflowBudgets, WorkflowDefinitionBundle, WorkflowProgress,
12 WorkflowRunDefinition, WorkflowRunSnapshot,
13};
14use bamboo_engine::{
15 AgentStepPort, AgentStepResult, FileWorkflowRunRepository, NamedAgentSpec, PermissionDecision,
16 WorkflowDefinitionPort, WorkflowPolicyPort, WorkflowPolicyTarget, WorkflowRunEngine,
17 WorkflowRunError, WorkflowSecretMaterial, WorkflowSecretResolverPort,
18 WorkflowSessionPermissionPort,
19};
20use bamboo_skills::SkillManager;
21use serde::Deserialize;
22use serde_json::{json, Value};
23
24const MAX_CONCURRENCY: usize = 8;
25const MAX_AGENTS: u32 = 16;
26const MAX_STEPS: u32 = 512;
27const MAX_RETRIES: u32 = 16;
28const MAX_NESTING_DEPTH: u32 = 8;
29const MAX_WALL_TIME_MS: u64 = 60 * 60 * 1000;
30const MAX_TOKENS: u64 = 2_000_000;
31const MAX_COST_MICROS: u64 = 100_000_000;
32const MAX_PINNED_DEFINITIONS_PER_RUN: usize = 32;
33const MAX_PINNED_BUNDLE_BYTES_PER_RUN: usize = 512 * 1024;
34const MAX_WORKFLOW_RUN_IDS_PER_SESSION: usize = 256;
35const SAFE_UNTRUSTED_WORKFLOW_TOOLS: &[&str] = &[
36 "Read",
37 "read_file",
38 "GetFileInfo",
39 "Glob",
40 "list_directory",
41 "Grep",
42];
43
44#[derive(Clone)]
47pub struct WorkflowRunAccess {
48 engine: Arc<WorkflowRunEngine>,
49 skills: Arc<SkillManager>,
50 sessions: bamboo_engine::SessionRepository,
51}
52
53struct ServerWorkflowSessionPermissions {
54 sessions: bamboo_engine::SessionRepository,
55 permission_config: Arc<bamboo_tools::permission::PermissionConfig>,
56}
57
58#[async_trait]
59impl WorkflowSessionPermissionPort for ServerWorkflowSessionPermissions {
60 async fn flags_for_session(
61 &self,
62 session_id: &str,
63 ) -> Result<ToolExecutionSessionFlags, String> {
64 let session = self
65 .sessions
66 .try_load(session_id)
67 .await
68 .map_err(|error| error.to_string())?
69 .ok_or_else(|| format!("workflow session '{session_id}' does not exist"))?;
70 let configured = if session
71 .agent_runtime_state
72 .as_ref()
73 .is_some_and(|state| state.plan_mode.is_some())
74 {
75 bamboo_domain::PermissionMode::Plan
76 } else {
77 self.permission_config.mode()
78 };
79 Ok(ToolExecutionSessionFlags::from_session_and_configured_mode(
80 &session, configured,
81 ))
82 }
83}
84
85impl WorkflowRunAccess {
86 pub async fn new(
87 data_dir: &Path,
88 tools: Arc<dyn bamboo_agent_core::tools::ToolExecutor>,
89 skills: Arc<SkillManager>,
90 sessions: bamboo_engine::SessionRepository,
91 ) -> Result<Self, String> {
92 Self::new_with_permission_config(data_dir, tools, skills, sessions, None).await
93 }
94
95 pub async fn new_with_permission_config(
96 data_dir: &Path,
97 tools: Arc<dyn bamboo_agent_core::tools::ToolExecutor>,
98 skills: Arc<SkillManager>,
99 sessions: bamboo_engine::SessionRepository,
100 permission_config: Option<Arc<bamboo_tools::permission::PermissionConfig>>,
101 ) -> Result<Self, String> {
102 let repository = Arc::new(
103 FileWorkflowRunRepository::new(data_dir.join("workflow-runs"))
104 .map_err(|error| format!("failed to initialize workflow journal: {error}"))?,
105 );
106 let engine = WorkflowRunEngine::new(
107 repository,
108 tools,
109 Arc::new(UnavailableAgentPort),
110 Arc::new(ExternallyPinnedDefinitions),
111 Arc::new(ServerWorkflowPolicy),
112 Arc::new(UnavailableSecretResolver),
113 WorkflowBudgets {
114 max_concurrency: MAX_CONCURRENCY,
115 max_agents: MAX_AGENTS,
116 max_steps: MAX_STEPS,
117 max_retries: MAX_RETRIES,
118 max_nesting_depth: MAX_NESTING_DEPTH,
119 wall_time_ms: MAX_WALL_TIME_MS,
120 max_tokens: Some(MAX_TOKENS),
121 max_cost_micros: Some(MAX_COST_MICROS),
122 },
123 );
124 if let Some(permission_config) = permission_config {
125 engine.set_session_permission_port(Arc::new(ServerWorkflowSessionPermissions {
126 sessions: sessions.clone(),
127 permission_config,
128 }));
129 }
130 engine
131 .recover()
132 .await
133 .map_err(|error| format!("failed to recover workflow journal: {error}"))?;
134 Ok(Self {
135 engine,
136 skills,
137 sessions,
138 })
139 }
140
141 async fn session_context(
142 &self,
143 session_id: &str,
144 ) -> Result<(Option<PathBuf>, bool), WorkflowRunError> {
145 let session =
146 self.sessions.try_load(session_id).await.map_err(|_| {
147 WorkflowRunError::Preflight("session state is unavailable".to_string())
148 })?;
149 let session = session.ok_or_else(|| {
150 WorkflowRunError::Preflight("workflow session does not exist".to_string())
151 })?;
152 let preferred = session.workspace.map(PathBuf::from);
153 let workspace =
154 bamboo_agent_core::workspace_state::ensure_session_workspace(session_id, preferred)
155 .or_else(|| {
156 Some(
160 bamboo_agent_core::workspace_state::workspace_or_process_cwd(Some(
161 session_id,
162 )),
163 )
164 });
165 Ok((workspace, false))
171 }
172
173 pub async fn start(
174 &self,
175 session_id: &str,
176 workflow_id: &str,
177 revision: u64,
178 args: Value,
179 budget: Option<WorkflowBudgets>,
180 ) -> Result<WorkflowRunSnapshot, WorkflowRunError> {
181 self.start_for_invoker(session_id, workflow_id, revision, args, budget, false)
182 .await
183 }
184
185 pub async fn start_from_tool(
186 &self,
187 session_id: &str,
188 workflow_id: &str,
189 revision: u64,
190 args: Value,
191 budget: Option<WorkflowBudgets>,
192 ) -> Result<WorkflowRunSnapshot, WorkflowRunError> {
193 self.start_for_invoker(session_id, workflow_id, revision, args, budget, true)
194 .await
195 }
196
197 async fn start_for_invoker(
198 &self,
199 session_id: &str,
200 workflow_id: &str,
201 revision: u64,
202 args: Value,
203 budget: Option<WorkflowBudgets>,
204 model_started: bool,
205 ) -> Result<WorkflowRunSnapshot, WorkflowRunError> {
206 self.ensure_run_index_capacity(session_id).await?;
207 let (workspace, workspace_trusted) = self.session_context(session_id).await?;
208 let store = self
209 .skills
210 .store_for_workspace(workspace.as_deref())
211 .await
212 .map_err(|_| {
213 WorkflowRunError::Preflight("workflow catalog is unavailable".to_string())
214 })?;
215 let catalog = store.workflow_catalog_snapshot().await;
216 let entry = catalog
217 .entries
218 .iter()
219 .find(|entry| {
220 entry.winner
221 && entry.id == workflow_id
222 && entry.revision == revision
223 && entry.status == bamboo_skills::WorkflowStatus::Valid
224 })
225 .ok_or_else(|| {
226 WorkflowRunError::Preflight(
227 "requested workflow revision is unavailable".to_string(),
228 )
229 })?;
230 if entry.kind != bamboo_skills::WorkflowKind::Orchestration {
231 return Err(WorkflowRunError::Preflight(
232 "instruction workflows must be activated with load_skill".to_string(),
233 ));
234 }
235 if model_started {
236 let session = self
237 .sessions
238 .try_load(session_id)
239 .await
240 .map_err(|_| {
241 WorkflowRunError::Preflight("session state is unavailable".to_string())
242 })?
243 .ok_or(WorkflowRunError::NotFound)?;
244 let opted_in = session
245 .metadata
246 .get(bamboo_skills::WORKFLOW_ORCHESTRATION_OPT_IN_METADATA_KEY)
247 .is_some_and(|value| value.eq_ignore_ascii_case("true"));
248 if !opted_in {
249 return Err(WorkflowRunError::Preflight(
250 "model-started orchestration requires explicit session opt-in".to_string(),
251 ));
252 }
253 }
254 let mut bundle = self
255 .skills
256 .pin_workflow_definition_bundle(workspace.as_deref(), workflow_id, revision)
257 .await
258 .map_err(|_| WorkflowRunError::Preflight("workflow catalog pin failed".to_string()))?;
259 let policy = if model_started {
260 "automatic"
261 } else {
262 "explicit"
263 };
264 if bundle.root_invocation_policy[policy].as_bool() != Some(true) {
265 return Err(WorkflowRunError::Preflight(format!(
266 "pinned workflow invocation policy denies {policy} start"
267 )));
268 }
269 let bundle_bytes = serde_json::to_vec(&bundle)
270 .map_err(|_| WorkflowRunError::Preflight("workflow bundle is invalid".to_string()))?
271 .len();
272 enforce_pinned_bundle_limits(bundle.definitions.len(), bundle_bytes)?;
273 let mut definition = bundle.root().cloned().ok_or_else(|| {
274 WorkflowRunError::Preflight("pinned workflow root is missing".to_string())
275 })?;
276 if let Some(requested) = budget {
277 validate_requested_budget(&requested)?;
278 definition.budgets = tighten_workflow_budget(&definition.budgets, &requested);
279 let root_key = WorkflowDefinitionBundle::key(&definition.id, definition.revision);
280 bundle.definitions.insert(root_key, definition.clone());
281 }
282 let snapshot = self
283 .engine
284 .start_pinned(
285 StartWorkflowRun {
286 definition,
287 args,
288 session_id: session_id.to_string(),
289 workspace_trusted,
290 allowed_capabilities: vec!["read".to_string()],
295 },
296 bundle,
297 )
298 .await?;
299 if let Err(error) = self.remember_run_id(session_id, &snapshot.run_id).await {
300 return match self.engine.cancel(&snapshot.run_id).await {
301 Ok(cancelled) if cancelled.status.is_terminal() => Err(WorkflowRunError::Storage(
302 format!(
303 "run index persistence failed; run {} reached terminal {:?}: {error}",
304 snapshot.run_id, cancelled.status
305 ),
306 )),
307 Ok(cancelled) => Err(WorkflowRunError::Storage(format!(
308 "run index persistence failed; orphan run {} remains {:?}; repair with this run id: {error}",
309 snapshot.run_id, cancelled.status
310 ))),
311 Err(cancel_error) => Err(WorkflowRunError::Storage(format!(
312 "run index persistence failed; orphan run {} could not be cancelled ({cancel_error}); repair with this run id: {error}",
313 snapshot.run_id
314 ))),
315 };
316 }
317 Ok(snapshot)
318 }
319
320 async fn ensure_run_index_capacity(&self, session_id: &str) -> Result<(), WorkflowRunError> {
321 let session = self
322 .sessions
323 .try_load(session_id)
324 .await
325 .map_err(|_| WorkflowRunError::Storage("session state unavailable".to_string()))?
326 .ok_or(WorkflowRunError::NotFound)?;
327 let ids = session
328 .metadata
329 .get(bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY)
330 .and_then(|raw| serde_json::from_str::<Vec<String>>(raw).ok())
331 .unwrap_or_default();
332 if ids.len() < MAX_WORKFLOW_RUN_IDS_PER_SESSION {
333 return Ok(());
334 }
335 let mut evictable = BTreeSet::new();
336 for run_id in &ids {
337 match self.engine.progress(run_id, u64::MAX).await {
338 Ok(progress) if progress.snapshot.status.is_terminal() => {
339 evictable.insert(run_id.clone());
340 }
341 Err(WorkflowRunError::NotFound) => {
342 evictable.insert(run_id.clone());
343 }
344 Ok(_) => {}
345 Err(error) => return Err(error),
346 }
347 }
348 if !evictable.is_empty() {
349 self.sessions
350 .update_runtime_session(
351 session_id,
352 &[bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY],
353 move |session| {
354 let mut ids = session
355 .metadata
356 .get(bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY)
357 .and_then(|raw| serde_json::from_str::<Vec<String>>(raw).ok())
358 .unwrap_or_default();
359 ids.retain(|id| !evictable.contains(id));
360 session.metadata.insert(
361 bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY.to_string(),
362 serde_json::to_string(&ids)
363 .expect("string vector serialization cannot fail"),
364 );
365 },
366 )
367 .await
368 .map_err(|_| WorkflowRunError::Storage("run index pruning failed".to_string()))?
369 .ok_or(WorkflowRunError::NotFound)?;
370 }
371 let remaining = self
372 .sessions
373 .try_load(session_id)
374 .await
375 .map_err(|_| WorkflowRunError::Storage("session state unavailable".to_string()))?
376 .and_then(|session| {
377 session
378 .metadata
379 .get(bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY)
380 .and_then(|raw| serde_json::from_str::<Vec<String>>(raw).ok())
381 })
382 .unwrap_or_default()
383 .len();
384 if remaining >= MAX_WORKFLOW_RUN_IDS_PER_SESSION {
385 return Err(WorkflowRunError::Preflight(
386 "workflow run index is full of active runs".to_string(),
387 ));
388 }
389 Ok(())
390 }
391
392 async fn remember_run_id(
393 &self,
394 session_id: &str,
395 run_id: &str,
396 ) -> Result<(), WorkflowRunError> {
397 let run_id = run_id.to_string();
398 let index_full = Arc::new(AtomicBool::new(false));
399 let index_full_in_transaction = index_full.clone();
400 self.sessions
401 .update_runtime_session(
402 session_id,
403 &[bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY],
404 move |session| {
405 let mut ids = session
406 .metadata
407 .get(bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY)
408 .and_then(|raw| serde_json::from_str::<Vec<String>>(raw).ok())
409 .unwrap_or_default();
410 ids.retain(|existing| existing != &run_id);
411 if ids.len() >= MAX_WORKFLOW_RUN_IDS_PER_SESSION {
412 index_full_in_transaction.store(true, Ordering::SeqCst);
413 return;
414 }
415 ids.push(run_id);
416 session.metadata.insert(
417 bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY.to_string(),
418 serde_json::to_string(&ids)
419 .expect("string vector serialization cannot fail"),
420 );
421 },
422 )
423 .await
424 .map_err(|_| WorkflowRunError::Storage("run index persistence failed".to_string()))?
425 .ok_or(WorkflowRunError::NotFound)
426 .and_then(|_| {
427 if index_full.load(Ordering::SeqCst) {
428 Err(WorkflowRunError::Storage(
429 "workflow run index reached its active-run capacity".to_string(),
430 ))
431 } else {
432 Ok(())
433 }
434 })
435 }
436
437 pub async fn list_for_session(
438 &self,
439 session_id: &str,
440 ) -> Result<Vec<WorkflowRunSnapshot>, WorkflowRunError> {
441 self.session_context(session_id).await?;
442 let session = self
443 .sessions
444 .try_load(session_id)
445 .await
446 .map_err(|_| WorkflowRunError::Storage("session state unavailable".to_string()))?
447 .ok_or(WorkflowRunError::NotFound)?;
448 let run_ids = session
449 .metadata
450 .get(bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY)
451 .and_then(|raw| serde_json::from_str::<Vec<String>>(raw).ok())
452 .unwrap_or_default();
453 let mut snapshots = Vec::new();
454 let mut stale = BTreeSet::new();
455 for run_id in run_ids {
456 match self.engine.progress(&run_id, u64::MAX).await {
457 Ok(progress) if progress.snapshot.session_id == session_id => {
458 snapshots.push(progress.snapshot);
459 }
460 Ok(_) | Err(WorkflowRunError::NotFound) => {
461 stale.insert(run_id);
462 }
463 Err(error) => return Err(error),
464 }
465 }
466 if !stale.is_empty() {
467 let _ = self
468 .sessions
469 .update_runtime_session(
470 session_id,
471 &[bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY],
472 move |session| {
473 let mut ids = session
474 .metadata
475 .get(bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY)
476 .and_then(|raw| serde_json::from_str::<Vec<String>>(raw).ok())
477 .unwrap_or_default();
478 ids.retain(|id| !stale.contains(id));
479 session.metadata.insert(
480 bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY.to_string(),
481 serde_json::to_string(&ids)
482 .expect("string vector serialization cannot fail"),
483 );
484 },
485 )
486 .await;
487 }
488 snapshots.sort_by_key(|snapshot| std::cmp::Reverse(snapshot.created_at));
489 Ok(snapshots)
490 }
491
492 pub async fn progress_for_session(
493 &self,
494 session_id: &str,
495 run_id: &str,
496 since: u64,
497 ) -> Result<WorkflowProgress, WorkflowRunError> {
498 let progress = self.engine.progress(run_id, since).await?;
499 if progress.snapshot.session_id != session_id {
500 return Err(WorkflowRunError::NotFound);
501 }
502 Ok(progress)
503 }
504
505 pub async fn cancel_for_session(
506 &self,
507 session_id: &str,
508 run_id: &str,
509 ) -> Result<WorkflowRunSnapshot, WorkflowRunError> {
510 self.progress_for_session(session_id, run_id, u64::MAX)
511 .await?;
512 self.engine.cancel(run_id).await
513 }
514
515 pub async fn restart_for_session(
516 &self,
517 session_id: &str,
518 run_id: &str,
519 ) -> Result<WorkflowRunSnapshot, WorkflowRunError> {
520 self.progress_for_session(session_id, run_id, u64::MAX)
521 .await?;
522 let (_, workspace_trusted) = self.session_context(session_id).await?;
523 self.engine
524 .restart(run_id, workspace_trusted, vec!["read".to_string()])
525 .await
526 }
527
528 pub async fn restart_from_tool(
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 if progress.snapshot.definition_bundle.root_invocation_policy["automatic"].as_bool()
537 != Some(true)
538 {
539 return Err(WorkflowRunError::Preflight(
540 "pinned workflow invocation policy denies automatic restart".to_string(),
541 ));
542 }
543 let session = self
544 .sessions
545 .try_load(session_id)
546 .await
547 .map_err(|_| WorkflowRunError::Preflight("session state is unavailable".to_string()))?
548 .ok_or(WorkflowRunError::NotFound)?;
549 let opted_in = session
550 .metadata
551 .get(bamboo_skills::WORKFLOW_ORCHESTRATION_OPT_IN_METADATA_KEY)
552 .is_some_and(|value| value.eq_ignore_ascii_case("true"));
553 if !opted_in {
554 return Err(WorkflowRunError::Preflight(
555 "model-started orchestration restart requires explicit session opt-in".to_string(),
556 ));
557 }
558 self.restart_for_session(session_id, run_id).await
559 }
560}
561
562fn tighten_workflow_budget(
563 definition: &WorkflowBudgets,
564 requested: &WorkflowBudgets,
565) -> WorkflowBudgets {
566 WorkflowBudgets {
567 max_concurrency: definition.max_concurrency.min(requested.max_concurrency),
568 max_agents: definition.max_agents.min(requested.max_agents),
569 max_steps: definition.max_steps.min(requested.max_steps),
570 max_retries: definition.max_retries.min(requested.max_retries),
571 max_nesting_depth: definition
572 .max_nesting_depth
573 .min(requested.max_nesting_depth),
574 wall_time_ms: definition.wall_time_ms.min(requested.wall_time_ms),
575 max_tokens: match (definition.max_tokens, requested.max_tokens) {
576 (Some(left), Some(right)) => Some(left.min(right)),
577 (Some(value), None) | (None, Some(value)) => Some(value),
578 (None, None) => None,
579 },
580 max_cost_micros: match (definition.max_cost_micros, requested.max_cost_micros) {
581 (Some(left), Some(right)) => Some(left.min(right)),
582 (Some(value), None) | (None, Some(value)) => Some(value),
583 (None, None) => None,
584 },
585 }
586}
587
588fn validate_requested_budget(requested: &WorkflowBudgets) -> Result<(), WorkflowRunError> {
589 if requested.max_concurrency == 0
590 || requested.max_steps == 0
591 || requested.max_nesting_depth == 0
592 || requested.wall_time_ms == 0
593 {
594 return Err(WorkflowRunError::InvalidInput(
595 "workflow execution limits must be positive".to_string(),
596 ));
597 }
598 Ok(())
599}
600
601fn enforce_pinned_bundle_limits(
602 definition_count: usize,
603 serialized_bytes: usize,
604) -> Result<(), WorkflowRunError> {
605 if definition_count > MAX_PINNED_DEFINITIONS_PER_RUN {
606 return Err(WorkflowRunError::Preflight(
607 "workflow dependency count exceeds the server limit".to_string(),
608 ));
609 }
610 if serialized_bytes > MAX_PINNED_BUNDLE_BYTES_PER_RUN {
611 return Err(WorkflowRunError::Preflight(
612 "workflow definition bundle exceeds the server size limit".to_string(),
613 ));
614 }
615 Ok(())
616}
617
618pub(crate) fn public_workflow_snapshot(mut snapshot: WorkflowRunSnapshot) -> WorkflowRunSnapshot {
622 snapshot.definition_bundle.definitions.clear();
623 snapshot
624}
625
626struct ExternallyPinnedDefinitions;
630
631#[async_trait]
632impl WorkflowDefinitionPort for ExternallyPinnedDefinitions {
633 async fn pin_bundle(
634 &self,
635 _root: &WorkflowRunDefinition,
636 ) -> Result<WorkflowDefinitionBundle, String> {
637 Err("server workflows must be pinned through SkillManager".to_string())
638 }
639}
640
641struct UnavailableAgentPort;
644
645#[async_trait]
646impl AgentStepPort for UnavailableAgentPort {
647 async fn resolve(&self, _name: &str) -> Result<Option<NamedAgentSpec>, String> {
648 Ok(None)
649 }
650
651 async fn execute(
652 &self,
653 _spec: &NamedAgentSpec,
654 _prompt: Value,
655 _model: Option<&str>,
656 _effort: Option<&str>,
657 _capabilities: &BTreeSet<String>,
658 _session_id: &str,
659 ) -> Result<AgentStepResult, String> {
660 Err("named-agent execution is not available".to_string())
661 }
662}
663
664struct ServerWorkflowPolicy;
665
666#[async_trait]
667impl WorkflowPolicyPort for ServerWorkflowPolicy {
668 async fn authorize(
669 &self,
670 _session_id: &str,
671 target: &WorkflowPolicyTarget,
672 requested: &BTreeSet<String>,
673 _workspace_trusted: bool,
674 ) -> PermissionDecision {
675 match target {
681 WorkflowPolicyTarget::Workflow { .. } if requested.is_empty() => {
685 PermissionDecision::Allow
686 }
687 WorkflowPolicyTarget::Tool(name) => {
688 let safe_target = SAFE_UNTRUSTED_WORKFLOW_TOOLS
689 .iter()
690 .any(|candidate| candidate.eq_ignore_ascii_case(name));
691 if safe_target && requested.iter().all(|capability| capability == "read") {
692 PermissionDecision::Allow
693 } else {
694 PermissionDecision::Deny(
695 "workflow capability authority is not available".to_string(),
696 )
697 }
698 }
699 WorkflowPolicyTarget::Agent(_) | WorkflowPolicyTarget::Workflow { .. } => {
700 PermissionDecision::Deny(
701 "workflow capability authority is not available".to_string(),
702 )
703 }
704 }
705 }
706}
707
708struct UnavailableSecretResolver;
709
710#[async_trait]
711impl WorkflowSecretResolverPort for UnavailableSecretResolver {
712 async fn resolve(
713 &self,
714 _session_id: &str,
715 _capability: &str,
716 ) -> Result<WorkflowSecretMaterial, String> {
717 Err("workflow secret capability resolver is not available".to_string())
718 }
719}
720
721#[derive(Debug, Deserialize)]
722#[serde(tag = "action", rename_all = "snake_case", deny_unknown_fields)]
723enum WorkflowToolInput {
724 Start {
725 workflow_id: String,
726 revision: u64,
727 #[serde(default = "empty_object")]
728 args: Value,
729 #[serde(default)]
730 budget: Option<WorkflowBudgets>,
731 },
732 List {},
733 Get {
734 run_id: String,
735 },
736 Events {
737 run_id: String,
738 #[serde(default)]
739 since: u64,
740 },
741 Cancel {
742 run_id: String,
743 },
744 Restart {
745 run_id: String,
746 },
747}
748
749fn empty_object() -> Value {
750 json!({})
751}
752
753pub struct WorkflowRunTool {
754 access: WorkflowRunAccess,
755}
756
757impl WorkflowRunTool {
758 pub fn new(access: WorkflowRunAccess) -> Self {
759 Self { access }
760 }
761}
762
763#[async_trait]
764impl Tool for WorkflowRunTool {
765 fn name(&self) -> &str {
766 "workflow_run"
767 }
768
769 fn description(&self) -> &str {
770 "Start, inspect, cancel, or safely restart a catalog-pinned workflow run"
771 }
772
773 fn parameters_schema(&self) -> Value {
774 json!({
775 "type": "object",
776 "properties": {
777 "action": {
778 "type": "string",
779 "enum": ["start", "list", "get", "events", "cancel", "restart"]
780 },
781 "workflow_id": {"type": "string", "minLength": 1},
782 "revision": {"type": "integer", "minimum": 1},
783 "args": {"type": "object", "default": {}},
784 "budget": {
785 "type": "object",
786 "properties": {
787 "max_concurrency": {"type": "integer", "minimum": 1},
788 "max_agents": {"type": "integer", "minimum": 0},
789 "max_steps": {"type": "integer", "minimum": 1},
790 "max_retries": {"type": "integer", "minimum": 0},
791 "max_nesting_depth": {"type": "integer", "minimum": 1},
792 "wall_time_ms": {"type": "integer", "minimum": 1},
793 "max_tokens": {"type": "integer", "minimum": 0},
794 "max_cost_micros": {"type": "integer", "minimum": 0}
795 },
796 "required": [
797 "max_concurrency",
798 "max_agents",
799 "max_steps",
800 "max_retries",
801 "max_nesting_depth",
802 "wall_time_ms"
803 ],
804 "additionalProperties": false
805 },
806 "run_id": {"type": "string"},
807 "since": {"type": "integer", "minimum": 0}
808 },
809 "required": ["action"],
810 "additionalProperties": false
811 })
812 }
813
814 fn classify(&self, args: &Value) -> ToolClass {
815 match args.get("action").and_then(Value::as_str) {
816 Some("get" | "list" | "events") => ToolClass::READONLY_PARALLEL,
817 _ => ToolClass::MUTATING_SERIAL,
818 }
819 }
820
821 async fn invoke(&self, args: Value, ctx: ToolCtx) -> Result<ToolOutcome, ToolError> {
822 let input: WorkflowToolInput = serde_json::from_value(args)
823 .map_err(|error| ToolError::InvalidArguments(error.to_string()))?;
824 let session_id = ctx.session_id().ok_or_else(|| {
825 ToolError::InvalidArguments("workflow_run requires a session".to_string())
826 })?;
827 let result = match input {
828 WorkflowToolInput::Start {
829 workflow_id,
830 revision,
831 args,
832 budget,
833 } => serde_json::to_value(public_workflow_snapshot(
834 self.access
835 .start_from_tool(session_id, &workflow_id, revision, args, budget)
836 .await
837 .map_err(workflow_tool_error)?,
838 )),
839 WorkflowToolInput::List {} => serde_json::to_value(
840 self.access
841 .list_for_session(session_id)
842 .await
843 .map_err(workflow_tool_error)?
844 .into_iter()
845 .map(public_workflow_snapshot)
846 .collect::<Vec<_>>(),
847 ),
848 WorkflowToolInput::Get { run_id } => {
849 let progress = self
850 .access
851 .progress_for_session(session_id, &run_id, u64::MAX)
852 .await
853 .map_err(workflow_tool_error)?;
854 serde_json::to_value(public_workflow_snapshot(progress.snapshot))
855 }
856 WorkflowToolInput::Events { run_id, since } => {
857 let progress = self
858 .access
859 .progress_for_session(session_id, &run_id, since)
860 .await
861 .map_err(workflow_tool_error)?;
862 serde_json::to_value(progress.events)
863 }
864 WorkflowToolInput::Cancel { run_id } => serde_json::to_value(public_workflow_snapshot(
865 self.access
866 .cancel_for_session(session_id, &run_id)
867 .await
868 .map_err(workflow_tool_error)?,
869 )),
870 WorkflowToolInput::Restart { run_id } => {
871 serde_json::to_value(public_workflow_snapshot(
872 self.access
873 .restart_from_tool(session_id, &run_id)
874 .await
875 .map_err(workflow_tool_error)?,
876 ))
877 }
878 }
879 .map_err(|error| ToolError::Execution(error.to_string()))?;
880 Ok(ToolOutcome::Completed(ToolResult::text(
881 true,
882 serde_json::to_string(&result)
883 .map_err(|error| ToolError::Execution(error.to_string()))?,
884 )))
885 }
886}
887
888fn workflow_tool_error(error: WorkflowRunError) -> ToolError {
889 match error {
890 WorkflowRunError::InvalidInput(message) => ToolError::InvalidArguments(message),
891 WorkflowRunError::Compile(error) => ToolError::InvalidArguments(error.to_string()),
892 other => ToolError::Execution(other.to_string()),
893 }
894}
895
896#[cfg(test)]
897mod tests {
898 use super::*;
899 use bamboo_agent_core::storage::Storage;
900 use bamboo_agent_core::tools::{
901 FunctionSchema, ToolCall, ToolExecutor, ToolResult, ToolSchema,
902 };
903 use bamboo_agent_core::Session;
904 use bamboo_llm::protocol::{gemini::GeminiTool, ToProvider};
905 use std::collections::HashMap;
906 use tokio::sync::RwLock;
907
908 #[derive(Default)]
909 struct WorkflowTestStorage {
910 sessions: RwLock<HashMap<String, Session>>,
911 }
912
913 #[async_trait]
914 impl Storage for WorkflowTestStorage {
915 async fn save_session(&self, session: &Session) -> std::io::Result<()> {
916 self.sessions
917 .write()
918 .await
919 .insert(session.id.clone(), session.clone());
920 Ok(())
921 }
922
923 async fn load_session(&self, session_id: &str) -> std::io::Result<Option<Session>> {
924 Ok(self.sessions.read().await.get(session_id).cloned())
925 }
926
927 async fn delete_session(&self, session_id: &str) -> std::io::Result<bool> {
928 Ok(self.sessions.write().await.remove(session_id).is_some())
929 }
930 }
931
932 struct WorkflowReadTool;
933
934 #[async_trait]
935 impl ToolExecutor for WorkflowReadTool {
936 async fn execute(
937 &self,
938 _call: &ToolCall,
939 ) -> bamboo_agent_core::tools::executor::Result<ToolResult> {
940 Ok(ToolResult::text(true, r#"{"ok":true}"#))
941 }
942
943 fn list_tools(&self) -> Vec<ToolSchema> {
944 vec![ToolSchema {
945 schema_type: "function".to_string(),
946 function: FunctionSchema {
947 name: "Read".to_string(),
948 description: "read".to_string(),
949 parameters: serde_json::json!({"type":"object"}),
950 },
951 }]
952 }
953 }
954
955 async fn workflow_test_access() -> (
956 WorkflowRunAccess,
957 bamboo_engine::SessionRepository,
958 tempfile::TempDir,
959 ) {
960 let directory = tempfile::tempdir().expect("tempdir");
961 let skills_dir = directory.path().join("skills");
962 let root = skills_dir.join("review-flow");
963 std::fs::create_dir_all(&root).expect("workflow dir");
964 std::fs::write(
965 root.join("SKILL.md"),
966 "---\nname: review-flow\ndescription: Review flow\n---\nRun review flow.\n",
967 )
968 .expect("skill");
969 std::fs::write(
970 root.join("workflow.yaml"),
971 "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",
972 )
973 .expect("workflow");
974 let skills = Arc::new(SkillManager::with_config(bamboo_skills::SkillStoreConfig {
975 skills_dir,
976 ..Default::default()
977 }));
978 skills.initialize().await.expect("skills initialize");
979 let storage: Arc<dyn Storage> = Arc::new(WorkflowTestStorage::default());
980 let persistence = Arc::new(bamboo_storage::LockedSessionStore::new(storage.clone()));
981 let cache = Arc::new(dashmap::DashMap::new());
982 let repo = bamboo_engine::SessionRepository::new(cache, storage, persistence);
983 let access = WorkflowRunAccess::new(
984 directory.path(),
985 Arc::new(WorkflowReadTool),
986 skills,
987 repo.clone(),
988 )
989 .await
990 .expect("workflow access");
991 (access, repo, directory)
992 }
993
994 fn canonical_workflow_run_schema() -> Value {
995 json!({
996 "type": "object",
997 "properties": {
998 "action": {
999 "type": "string",
1000 "enum": ["start", "list", "get", "events", "cancel", "restart"]
1001 },
1002 "workflow_id": {"type": "string", "minLength": 1},
1003 "revision": {"type": "integer", "minimum": 1},
1004 "args": {"type": "object", "default": {}},
1005 "budget": {
1006 "type": "object",
1007 "properties": {
1008 "max_concurrency": {"type": "integer", "minimum": 1},
1009 "max_agents": {"type": "integer", "minimum": 0},
1010 "max_steps": {"type": "integer", "minimum": 1},
1011 "max_retries": {"type": "integer", "minimum": 0},
1012 "max_nesting_depth": {"type": "integer", "minimum": 1},
1013 "wall_time_ms": {"type": "integer", "minimum": 1},
1014 "max_tokens": {"type": "integer", "minimum": 0},
1015 "max_cost_micros": {"type": "integer", "minimum": 0}
1016 },
1017 "required": [
1018 "max_concurrency",
1019 "max_agents",
1020 "max_steps",
1021 "max_retries",
1022 "max_nesting_depth",
1023 "wall_time_ms"
1024 ],
1025 "additionalProperties": false
1026 },
1027 "run_id": {"type": "string"},
1028 "since": {"type": "integer", "minimum": 0}
1029 },
1030 "required": ["action"],
1031 "additionalProperties": false
1032 })
1033 }
1034
1035 #[tokio::test]
1036 async fn workflow_run_schema_is_flat_complete_and_canonical() {
1037 let (access, _, _) = workflow_test_access().await;
1038 let schema = WorkflowRunTool::new(access).parameters_schema();
1039
1040 for combinator in ["oneOf", "anyOf", "allOf"] {
1041 assert!(
1042 schema.get(combinator).is_none(),
1043 "workflow_run must not advertise root {combinator}"
1044 );
1045 }
1046 assert_eq!(schema, canonical_workflow_run_schema());
1047 }
1048
1049 #[tokio::test]
1050 async fn workflow_run_schema_survives_openai_sanitization_with_all_properties() {
1051 let (access, _, _) = workflow_test_access().await;
1052 let schema = WorkflowRunTool::new(access).parameters_schema();
1053 let sanitized =
1054 bamboo_llm::providers::common::tool_schema::sanitize_openai_function_parameters_schema(
1055 &schema,
1056 );
1057
1058 let properties = sanitized["properties"]
1059 .as_object()
1060 .expect("sanitized workflow_run properties");
1061 assert!(!properties.is_empty());
1062 assert_eq!(properties.len(), 7);
1063 assert_eq!(sanitized, canonical_workflow_run_schema());
1064 }
1065
1066 #[tokio::test]
1067 async fn workflow_run_schema_reaches_gemini_unchanged() {
1068 let (access, _, _) = workflow_test_access().await;
1069 let direct = WorkflowRunTool::new(access).to_schema();
1070 let gemini: GeminiTool = direct.to_provider().expect("Gemini tool conversion");
1071 let declaration = gemini
1072 .function_declarations
1073 .first()
1074 .expect("workflow_run declaration");
1075
1076 assert_eq!(declaration.name, "workflow_run");
1077 assert_eq!(
1078 declaration.parameters_json_schema.as_ref(),
1079 Some(&canonical_workflow_run_schema())
1080 );
1081 assert!(declaration.parameters.is_none());
1082 }
1083
1084 #[tokio::test]
1085 async fn workflow_run_enforces_opt_in_tightens_budget_lists_and_isolates_sessions() {
1086 let (access, repo, directory) = workflow_test_access().await;
1087 let workspace = directory.path().join("workspace");
1088 std::fs::create_dir_all(&workspace).expect("workspace");
1089 let mut session = Session::new("workflow-session", "model");
1090 session.workspace = Some(workspace.to_string_lossy().into_owned());
1091 repo.save(&mut session).await.expect("save session");
1092
1093 let denied = access
1094 .start_from_tool("workflow-session", "review-flow", 42, json!({}), None)
1095 .await
1096 .expect_err("model start defaults off without session opt-in");
1097 assert!(denied.to_string().contains("opt-in"));
1098
1099 session.metadata.insert(
1100 bamboo_skills::WORKFLOW_ORCHESTRATION_OPT_IN_METADATA_KEY.to_string(),
1101 "true".to_string(),
1102 );
1103 repo.save(&mut session).await.expect("save opt-in");
1104 let requested = WorkflowBudgets {
1105 max_concurrency: 1,
1106 max_agents: 0,
1107 max_steps: 2,
1108 max_retries: 0,
1109 max_nesting_depth: 1,
1110 wall_time_ms: 5_000,
1111 max_tokens: Some(500),
1112 max_cost_micros: Some(500),
1113 };
1114 let started = access
1115 .start_from_tool(
1116 "workflow-session",
1117 "review-flow",
1118 42,
1119 json!({}),
1120 Some(requested.clone()),
1121 )
1122 .await
1123 .expect("opted-in model start");
1124 assert_eq!(started.definition.budgets, requested);
1125 let listed = access
1126 .list_for_session("workflow-session")
1127 .await
1128 .expect("session run list");
1129 assert_eq!(listed.len(), 1);
1130 assert_eq!(listed[0].run_id, started.run_id);
1131 let progress = access
1132 .progress_for_session("workflow-session", &started.run_id, 0)
1133 .await
1134 .expect("run events");
1135 assert!(progress
1136 .events
1137 .first()
1138 .is_some_and(|event| event.kind == bamboo_domain::WorkflowRunEventKind::RunQueued));
1139 let completed = tokio::time::timeout(std::time::Duration::from_secs(2), async {
1140 loop {
1141 let progress = access
1142 .progress_for_session("workflow-session", &started.run_id, 0)
1143 .await
1144 .expect("terminal run progress");
1145 if progress.snapshot.status.is_terminal() {
1146 break progress;
1147 }
1148 tokio::time::sleep(std::time::Duration::from_millis(10)).await;
1149 }
1150 })
1151 .await
1152 .expect("workflow reaches terminal state");
1153 assert_eq!(
1154 completed.snapshot.status,
1155 bamboo_domain::WorkflowRunStatus::Succeeded
1156 );
1157 assert_eq!(
1158 completed
1159 .events
1160 .iter()
1161 .map(|event| event.sequence)
1162 .collect::<Vec<_>>(),
1163 (1..=7).collect::<Vec<_>>()
1164 );
1165 assert!(matches!(
1166 completed.events.as_slice(),
1167 [
1168 bamboo_domain::WorkflowRunEvent {
1169 kind: bamboo_domain::WorkflowRunEventKind::RunQueued,
1170 ..
1171 },
1172 bamboo_domain::WorkflowRunEvent {
1173 kind: bamboo_domain::WorkflowRunEventKind::RunStarted,
1174 ..
1175 },
1176 bamboo_domain::WorkflowRunEvent {
1177 kind: bamboo_domain::WorkflowRunEventKind::Phase { ref name },
1178 ..
1179 },
1180 bamboo_domain::WorkflowRunEvent {
1181 kind: bamboo_domain::WorkflowRunEventKind::StepQueued,
1182 ..
1183 },
1184 bamboo_domain::WorkflowRunEvent {
1185 kind: bamboo_domain::WorkflowRunEventKind::StepStarted,
1186 ..
1187 },
1188 bamboo_domain::WorkflowRunEvent {
1189 kind: bamboo_domain::WorkflowRunEventKind::StepCompleted { .. },
1190 ..
1191 },
1192 bamboo_domain::WorkflowRunEvent {
1193 kind: bamboo_domain::WorkflowRunEventKind::RunSucceeded { .. },
1194 ..
1195 }
1196 ] if name == "step_reserved"
1197 ));
1198 assert_eq!(
1199 completed.snapshot.last_sequence,
1200 completed.events.last().expect("terminal event").sequence
1201 );
1202
1203 let invalid_budget = WorkflowBudgets {
1204 max_steps: 0,
1205 ..requested.clone()
1206 };
1207 assert!(matches!(
1208 access
1209 .start_from_tool(
1210 "workflow-session",
1211 "review-flow",
1212 42,
1213 json!({}),
1214 Some(invalid_budget),
1215 )
1216 .await,
1217 Err(WorkflowRunError::InvalidInput(_))
1218 ));
1219
1220 let workflow_path = directory.path().join("skills/review-flow/workflow.yaml");
1223 let original_workflow = std::fs::read_to_string(&workflow_path).expect("workflow yaml");
1224 std::fs::write(
1225 &workflow_path,
1226 original_workflow.replace(
1227 "invocation_policy: {explicit: true, automatic: true}",
1228 "invocation_policy: {explicit: true, automatic: false}",
1229 ),
1230 )
1231 .expect("disable automatic live policy");
1232 access.skills.store().reload().await.expect("reload policy");
1233 let restart = access
1234 .restart_from_tool("workflow-session", &started.run_id)
1235 .await
1236 .expect_err("succeeded runs are terminal");
1237 assert!(matches!(restart, WorkflowRunError::Terminal));
1238
1239 let mut other = Session::new("other-session", "model");
1240 other.workspace = Some(workspace.to_string_lossy().into_owned());
1241 repo.save(&mut other).await.expect("save other session");
1242 assert!(access
1243 .list_for_session("other-session")
1244 .await
1245 .expect("isolated list")
1246 .is_empty());
1247 assert!(matches!(
1248 access
1249 .progress_for_session("other-session", &started.run_id, 0)
1250 .await,
1251 Err(WorkflowRunError::NotFound)
1252 ));
1253 }
1254
1255 #[test]
1256 fn tool_input_rejects_security_context_spoofing() {
1257 let error = serde_json::from_value::<WorkflowToolInput>(json!({
1258 "action": "start",
1259 "workflow_id": "safe",
1260 "revision": 1,
1261 "workspace_trusted": true
1262 }))
1263 .unwrap_err();
1264 assert!(error.to_string().contains("unknown field"));
1265 }
1266
1267 #[test]
1268 fn tool_input_rejects_fields_from_other_actions() {
1269 let invalid = [
1270 json!({"action": "list", "run_id": "run-1"}),
1271 json!({"action": "get", "run_id": "run-1", "since": 1}),
1272 json!({"action": "events", "run_id": "run-1", "workflow_id": "flow"}),
1273 json!({"action": "cancel", "run_id": "run-1", "revision": 1}),
1274 json!({"action": "restart", "run_id": "run-1", "budget": {}}),
1275 json!({
1276 "action": "start",
1277 "workflow_id": "flow",
1278 "revision": 1,
1279 "run_id": "run-1"
1280 }),
1281 ];
1282
1283 for input in invalid {
1284 let error = serde_json::from_value::<WorkflowToolInput>(input.clone())
1285 .expect_err("action-specific fields must remain authoritative at runtime");
1286 assert!(
1287 error.to_string().contains("unknown field"),
1288 "unexpected error for {input}: {error}"
1289 );
1290 }
1291
1292 assert!(matches!(
1293 serde_json::from_value::<WorkflowToolInput>(json!({"action": "list"}))
1294 .expect("fieldless list action"),
1295 WorkflowToolInput::List {}
1296 ));
1297 }
1298
1299 #[tokio::test]
1300 async fn omitted_start_args_and_zero_budgets_match_schema() {
1301 let WorkflowToolInput::Start { args, .. } =
1302 serde_json::from_value::<WorkflowToolInput>(json!({
1303 "action": "start",
1304 "workflow_id": "safe",
1305 "revision": 1
1306 }))
1307 .expect("tool args default")
1308 else {
1309 panic!("start input")
1310 };
1311 assert_eq!(args, json!({}));
1312 let http: crate::handlers::workflow_runs::StartWorkflowRunRequest =
1313 serde_json::from_value(json!({"workflow_id":"safe", "revision":1}))
1314 .expect("http args default");
1315 assert_eq!(http.args, json!({}));
1316
1317 let (access, _, _) = workflow_test_access().await;
1318 let schema = WorkflowRunTool { access }.parameters_schema();
1319 let properties = &schema["properties"];
1320 assert_eq!(properties["args"]["default"], json!({}));
1321 assert_eq!(
1322 properties["budget"]["properties"]["max_agents"]["minimum"],
1323 0
1324 );
1325 assert_eq!(
1326 properties["budget"]["properties"]["max_retries"]["minimum"],
1327 0
1328 );
1329 assert_eq!(
1330 properties["budget"]["properties"]["max_tokens"]["minimum"],
1331 0
1332 );
1333 assert_eq!(
1334 properties["budget"]["properties"]["max_cost_micros"]["minimum"],
1335 0
1336 );
1337 }
1338
1339 #[tokio::test]
1340 async fn run_index_updates_are_concurrent_and_never_evict_at_active_capacity() {
1341 let (access, repo, _) = workflow_test_access().await;
1342 let mut session = Session::new("run-index", "model");
1343 repo.save(&mut session).await.expect("seed session");
1344 let results = futures::future::join_all((0..32).map(|index| {
1345 let access = access.clone();
1346 async move {
1347 let run_id = format!("run-{index}");
1348 access.remember_run_id("run-index", &run_id).await
1349 }
1350 }))
1351 .await;
1352 assert!(results.into_iter().all(|result| result.is_ok()));
1353 let concurrent = repo
1354 .try_load("run-index")
1355 .await
1356 .expect("load")
1357 .expect("session");
1358 let ids = serde_json::from_str::<Vec<String>>(
1359 concurrent
1360 .metadata
1361 .get(bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY)
1362 .expect("run ids"),
1363 )
1364 .expect("ids json");
1365 assert_eq!(ids.len(), 32);
1366 assert_eq!(ids.iter().collect::<BTreeSet<_>>().len(), 32);
1367
1368 let capacity_ids = (0..MAX_WORKFLOW_RUN_IDS_PER_SESSION)
1369 .map(|index| format!("active-{index}"))
1370 .collect::<Vec<_>>();
1371 repo.update_runtime_session(
1372 "run-index",
1373 &[bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY],
1374 {
1375 let capacity_ids = capacity_ids.clone();
1376 move |session| {
1377 session.metadata.insert(
1378 bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY.to_string(),
1379 serde_json::to_string(&capacity_ids).expect("ids json"),
1380 );
1381 }
1382 },
1383 )
1384 .await
1385 .expect("fill index")
1386 .expect("session");
1387 assert!(matches!(
1388 access.remember_run_id("run-index", "new-run").await,
1389 Err(WorkflowRunError::Storage(_))
1390 ));
1391 let retained = repo
1392 .try_load("run-index")
1393 .await
1394 .expect("load")
1395 .expect("session");
1396 let retained = serde_json::from_str::<Vec<String>>(
1397 retained
1398 .metadata
1399 .get(bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY)
1400 .expect("run ids"),
1401 )
1402 .expect("ids json");
1403 assert_eq!(
1404 retained, capacity_ids,
1405 "oldest active id must not be evicted"
1406 );
1407 }
1408
1409 #[tokio::test]
1410 async fn real_model_workflow_run_index_survives_tool_result_and_final_session_save() {
1411 use bamboo_agent_core::storage::AttachmentReader;
1412 use bamboo_engine::{Agent, ExecuteRequestBuilder};
1413 use bamboo_llm::{LLMChunk, LLMProvider, LLMStream};
1414 use futures::stream;
1415 use tokio::sync::Mutex;
1416 use tokio_util::sync::CancellationToken;
1417
1418 struct NoAttachments;
1419 #[async_trait]
1420 impl AttachmentReader for NoAttachments {
1421 async fn read_attachment(
1422 &self,
1423 _session_id: &str,
1424 _attachment_id: &str,
1425 ) -> std::io::Result<Option<(Vec<u8>, String)>> {
1426 Ok(None)
1427 }
1428 }
1429 struct QueueProvider {
1430 queue: Mutex<Vec<Vec<bamboo_llm::provider::Result<LLMChunk>>>>,
1431 }
1432 #[async_trait]
1433 impl LLMProvider for QueueProvider {
1434 async fn chat_stream(
1435 &self,
1436 _messages: &[bamboo_agent_core::Message],
1437 _tools: &[ToolSchema],
1438 _max_output_tokens: Option<u32>,
1439 _model: &str,
1440 ) -> bamboo_llm::provider::Result<LLMStream> {
1441 Ok(Box::pin(stream::iter(self.queue.lock().await.remove(0))))
1442 }
1443 }
1444
1445 let (access, repo, directory) = workflow_test_access().await;
1446 let session_id = "real-model-workflow-run";
1447 let mut session = Session::new(session_id, "test-model");
1448 session.metadata.insert(
1449 bamboo_skills::WORKFLOW_ORCHESTRATION_OPT_IN_METADATA_KEY.to_string(),
1450 "true".to_string(),
1451 );
1452 session
1453 .metadata
1454 .insert("external.metadata".to_string(), "preserve".to_string());
1455 session.add_message(bamboo_agent_core::Message::system("system"));
1456 session.add_message(bamboo_agent_core::Message::user("run review workflow"));
1457 repo.save(&mut session).await.expect("seed session");
1458 let call = ToolCall {
1459 id: "call-workflow-run".to_string(),
1460 tool_type: "function".to_string(),
1461 function: bamboo_agent_core::tools::FunctionCall {
1462 name: "workflow_run".to_string(),
1463 arguments: json!({
1464 "action":"start",
1465 "workflow_id":"review-flow",
1466 "revision":42
1467 })
1468 .to_string(),
1469 },
1470 };
1471 let provider = Arc::new(QueueProvider {
1472 queue: Mutex::new(vec![
1473 vec![Ok(LLMChunk::ToolCalls(vec![call])), Ok(LLMChunk::Done)],
1474 vec![Ok(LLMChunk::Token("done".to_string())), Ok(LLMChunk::Done)],
1475 ]),
1476 });
1477 let tools = Arc::new(
1478 bamboo_tools::BuiltinToolExecutorBuilder::new()
1479 .with_tool(WorkflowRunTool::new(access.clone()))
1480 .expect("workflow tool")
1481 .build(),
1482 );
1483 let metrics = bamboo_metrics::MetricsCollector::spawn(
1484 Arc::new(bamboo_metrics::SqliteMetricsStorage::new(
1485 directory.path().join("runner-metrics.db"),
1486 )),
1487 7,
1488 );
1489 let agent = Agent::builder()
1490 .storage(repo.storage().clone())
1491 .persistence(Arc::new(repo.clone()))
1492 .attachment_reader(Arc::new(NoAttachments))
1493 .skill_manager(access.skills.clone())
1494 .metrics_collector(metrics)
1495 .config(Arc::new(RwLock::new(bamboo_llm::Config::default())))
1496 .provider(provider)
1497 .default_tools(tools)
1498 .build()
1499 .expect("agent");
1500 let (event_tx, _event_rx) = tokio::sync::mpsc::channel(128);
1501 agent
1502 .execute(
1503 &mut session,
1504 ExecuteRequestBuilder::new(
1505 "run review workflow",
1506 event_tx,
1507 CancellationToken::new(),
1508 )
1509 .model("test-model")
1510 .build(),
1511 )
1512 .await
1513 .expect("real model workflow run");
1514
1515 let saved = repo
1516 .storage()
1517 .load_session(session_id)
1518 .await
1519 .expect("load")
1520 .expect("saved");
1521 assert_eq!(
1522 saved.metadata.get("external.metadata").map(String::as_str),
1523 Some("preserve")
1524 );
1525 assert!(saved
1526 .metadata
1527 .contains_key(bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY));
1528 let listed = access
1529 .list_for_session(session_id)
1530 .await
1531 .expect("list after final save");
1532 assert_eq!(listed.len(), 1);
1533 assert!(saved.messages.iter().any(|message| {
1534 message.tool_calls.as_ref().is_some_and(|calls| {
1535 calls
1536 .iter()
1537 .any(|call| call.function.name == "workflow_run")
1538 })
1539 }));
1540 }
1541
1542 #[tokio::test]
1543 async fn http_start_survives_concurrent_stale_runner_final_save_and_server_restart() {
1544 use bamboo_agent_core::storage::AttachmentReader;
1545 use bamboo_engine::{Agent, ExecuteRequestBuilder};
1546 use bamboo_llm::{LLMChunk, LLMProvider, LLMStream};
1547 use futures::stream;
1548 use tokio::sync::{oneshot, Mutex};
1549 use tokio_util::sync::CancellationToken;
1550
1551 struct NoAttachments;
1552 #[async_trait]
1553 impl AttachmentReader for NoAttachments {
1554 async fn read_attachment(
1555 &self,
1556 _session_id: &str,
1557 _attachment_id: &str,
1558 ) -> std::io::Result<Option<(Vec<u8>, String)>> {
1559 Ok(None)
1560 }
1561 }
1562
1563 struct PausingProvider {
1564 entered: Mutex<Option<oneshot::Sender<()>>>,
1565 resume: Mutex<Option<oneshot::Receiver<()>>>,
1566 }
1567 #[async_trait]
1568 impl LLMProvider for PausingProvider {
1569 async fn chat_stream(
1570 &self,
1571 _messages: &[bamboo_agent_core::Message],
1572 _tools: &[ToolSchema],
1573 _max_output_tokens: Option<u32>,
1574 _model: &str,
1575 ) -> bamboo_llm::provider::Result<LLMStream> {
1576 if let Some(entered) = self.entered.lock().await.take() {
1577 let _ = entered.send(());
1578 }
1579 if let Some(resume) = self.resume.lock().await.take() {
1580 let _ = resume.await;
1581 }
1582 Ok(Box::pin(stream::iter(vec![
1583 Ok(LLMChunk::Token("done".to_string())),
1584 Ok(LLMChunk::Done),
1585 ])))
1586 }
1587 }
1588
1589 let (access, repo, directory) = workflow_test_access().await;
1590 let session_id = "http-start-concurrent-runner-save";
1591 let mut session = Session::new(session_id, "test-model");
1592 session.add_message(bamboo_agent_core::Message::system("system"));
1593 session.add_message(bamboo_agent_core::Message::user("keep running"));
1594 repo.save(&mut session).await.expect("seed session");
1595 let mut runner_session = repo
1596 .try_load(session_id)
1597 .await
1598 .expect("load runner session")
1599 .expect("runner session");
1600
1601 let (entered_tx, entered_rx) = oneshot::channel();
1602 let (resume_tx, resume_rx) = oneshot::channel();
1603 let provider = Arc::new(PausingProvider {
1604 entered: Mutex::new(Some(entered_tx)),
1605 resume: Mutex::new(Some(resume_rx)),
1606 });
1607 let metrics = bamboo_metrics::MetricsCollector::spawn(
1608 Arc::new(bamboo_metrics::SqliteMetricsStorage::new(
1609 directory.path().join("http-runner-metrics.db"),
1610 )),
1611 7,
1612 );
1613 let agent = Agent::builder()
1614 .storage(repo.storage().clone())
1615 .persistence(Arc::new(repo.clone()))
1616 .attachment_reader(Arc::new(NoAttachments))
1617 .skill_manager(access.skills.clone())
1618 .metrics_collector(metrics)
1619 .config(Arc::new(RwLock::new(bamboo_llm::Config::default())))
1620 .provider(provider)
1621 .default_tools(Arc::new(
1622 bamboo_tools::BuiltinToolExecutorBuilder::new().build(),
1623 ))
1624 .build()
1625 .expect("agent");
1626 let (event_tx, _event_rx) = tokio::sync::mpsc::channel(128);
1627 let runner = tokio::spawn(async move {
1628 agent
1629 .execute(
1630 &mut runner_session,
1631 ExecuteRequestBuilder::new("keep running", event_tx, CancellationToken::new())
1632 .model("test-model")
1633 .build(),
1634 )
1635 .await
1636 });
1637 tokio::time::timeout(std::time::Duration::from_secs(2), entered_rx)
1638 .await
1639 .expect("runner enters model round")
1640 .expect("runner entry signal");
1641
1642 let started = access
1643 .start(session_id, "review-flow", 42, json!({}), None)
1644 .await
1645 .expect("HTTP-equivalent explicit start");
1646 let durable_during_round = repo
1647 .storage()
1648 .load_session(session_id)
1649 .await
1650 .expect("load during round")
1651 .expect("session during round");
1652 assert!(durable_during_round
1653 .metadata
1654 .get(bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY)
1655 .is_some_and(|raw| raw.contains(&started.run_id)));
1656
1657 resume_tx.send(()).expect("resume runner");
1658 tokio::time::timeout(std::time::Duration::from_secs(2), runner)
1659 .await
1660 .expect("runner completes")
1661 .expect("runner task")
1662 .expect("runner execution");
1663
1664 let durable_after_final_save = repo
1665 .storage()
1666 .load_session(session_id)
1667 .await
1668 .expect("load after final save")
1669 .expect("saved session");
1670 assert!(durable_after_final_save
1671 .metadata
1672 .get(bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY)
1673 .is_some_and(|raw| raw.contains(&started.run_id)));
1674
1675 tokio::time::timeout(std::time::Duration::from_secs(2), async {
1676 loop {
1677 let progress = access
1678 .progress_for_session(session_id, &started.run_id, u64::MAX)
1679 .await
1680 .expect("workflow progress before restart");
1681 if progress.snapshot.status.is_terminal()
1682 && !access.engine.is_run_active(&started.run_id)
1683 {
1684 break;
1685 }
1686 tokio::task::yield_now().await;
1687 }
1688 })
1689 .await
1690 .expect("workflow reaches terminal state before restart");
1691
1692 let skills = access.skills.clone();
1693 drop(access);
1694 let restarted =
1695 WorkflowRunAccess::new(directory.path(), Arc::new(WorkflowReadTool), skills, repo)
1696 .await
1697 .expect("restart workflow access");
1698 let listed = restarted
1699 .list_for_session(session_id)
1700 .await
1701 .expect("list after restart");
1702 assert_eq!(listed.len(), 1);
1703 assert_eq!(listed[0].run_id, started.run_id);
1704 }
1705
1706 #[tokio::test]
1707 async fn production_policy_allows_read_without_fabricating_workspace_trust() {
1708 let read = BTreeSet::from(["read".to_string()]);
1709 assert_eq!(
1710 ServerWorkflowPolicy
1711 .authorize(
1712 "session",
1713 &WorkflowPolicyTarget::Tool("read_file".to_string()),
1714 &read,
1715 false,
1716 )
1717 .await,
1718 PermissionDecision::Allow
1719 );
1720
1721 let write = BTreeSet::from(["write".to_string()]);
1722 assert!(matches!(
1723 ServerWorkflowPolicy
1724 .authorize(
1725 "session",
1726 &WorkflowPolicyTarget::Tool("write_file".to_string()),
1727 &write,
1728 false,
1729 )
1730 .await,
1731 PermissionDecision::Deny(_)
1732 ));
1733
1734 for hostile_target in [
1735 "Write",
1736 "write_file",
1737 "WebFetch",
1738 "mcp::remote_tool",
1739 "Bash",
1740 ] {
1741 for claimed in [BTreeSet::new(), read.clone()] {
1742 assert!(matches!(
1743 ServerWorkflowPolicy
1744 .authorize(
1745 "session",
1746 &WorkflowPolicyTarget::Tool(hostile_target.to_string()),
1747 &claimed,
1748 false,
1749 )
1750 .await,
1751 PermissionDecision::Deny(_)
1752 ));
1753 }
1754 }
1755
1756 assert_eq!(
1757 ServerWorkflowPolicy
1758 .authorize(
1759 "session",
1760 &WorkflowPolicyTarget::Workflow {
1761 id: "nested-review".to_string(),
1762 revision: 1,
1763 },
1764 &BTreeSet::new(),
1765 false,
1766 )
1767 .await,
1768 PermissionDecision::Allow
1769 );
1770 }
1771
1772 #[test]
1773 fn pinned_bundle_limits_reject_oversized_runs_before_engine_start() {
1774 assert!(enforce_pinned_bundle_limits(
1775 MAX_PINNED_DEFINITIONS_PER_RUN,
1776 MAX_PINNED_BUNDLE_BYTES_PER_RUN
1777 )
1778 .is_ok());
1779 assert!(enforce_pinned_bundle_limits(
1780 MAX_PINNED_DEFINITIONS_PER_RUN + 1,
1781 MAX_PINNED_BUNDLE_BYTES_PER_RUN
1782 )
1783 .is_err());
1784 assert!(enforce_pinned_bundle_limits(
1785 MAX_PINNED_DEFINITIONS_PER_RUN,
1786 MAX_PINNED_BUNDLE_BYTES_PER_RUN + 1
1787 )
1788 .is_err());
1789 }
1790}