Skip to main content

bamboo_server/workflow/
run.rs

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