Skip to main content

bamboo_server/workflow/
run.rs

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/// Server-owned access boundary for workflow runs. Session, workspace trust and
42/// capabilities are derived here rather than accepted from HTTP/tool callers.
43#[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                    // Server bootstrap registers the workspace-root provider, so this
106                    // yields a persistent session-scoped directory under data_dir when
107                    // the session has no explicit/configured workspace (#217).
108                    Some(
109                        bamboo_agent_core::workspace_state::workspace_or_process_cwd(Some(
110                            session_id,
111                        )),
112                    )
113                });
114        // Path resolution and workspace trust are separate authorities. Until
115        // #601 provides an explicit server-owned trust decision, never infer
116        // trust merely because the resolved path is absolute. Read-only runs
117        // remain available through `ServerWorkflowPolicy`; every stronger
118        // capability stays fail-closed.
119        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                    // #601 is not a production capability authority yet. Grant
240                    // only the server-owned read-only class needed by the review
241                    // dogfood workflow; the concrete base ToolExecutor still
242                    // performs its normal per-tool/per-resource permission gate.
243                    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
567/// The durable snapshot keeps the complete pinned bundle for deterministic
568/// restart, but clients need only the root definition, bundle identity/hash,
569/// step tree and events. Avoid echoing every nested definition over HTTP/tools.
570pub(crate) fn public_workflow_snapshot(mut snapshot: WorkflowRunSnapshot) -> WorkflowRunSnapshot {
571    snapshot.definition_bundle.definitions.clear();
572    snapshot
573}
574
575/// The catalog adapter above always calls `start_pinned`. Generic engine starts
576/// remain unavailable so no future server call site can accidentally re-read a
577/// live definition or mix publications.
578struct 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
590/// #563 named-agent registry integration is not complete. Unknown and named
591/// agents therefore fail preflight rather than falling back to a dynamic agent.
592struct 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        // Definition-declared capabilities are descriptive input, not an
625        // authority. Until #601 provides a server-owned capability registry,
626        // bind authorization to this strict server-owned target allowlist too.
627        // Unknown/MCP/network/script/write targets therefore fail closed even
628        // when hostile YAML claims `capabilities: []` or `[read]`.
629        match target {
630            // The reference itself belongs to the same immutable, validated
631            // bundle. Every concrete step in the nested definition is still
632            // authorized independently during this preflight walk.
633            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        // Restart authority is the original immutable publication, not the
1052        // current catalog policy.
1053        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}