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