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