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        self.index_started_run_or_compensate(session_id, snapshot)
302            .await
303    }
304
305    async fn index_started_run_or_compensate(
306        &self,
307        session_id: &str,
308        snapshot: WorkflowRunSnapshot,
309    ) -> Result<WorkflowRunSnapshot, WorkflowRunError> {
310        if let Err(error) = self.remember_run_id(session_id, &snapshot.run_id).await {
311            return match self.engine.cancel(&snapshot.run_id).await {
312                Ok(cancelled) if cancelled.status.is_terminal() => Err(WorkflowRunError::Storage(
313                    format!(
314                        "run index persistence failed; run {} reached terminal {:?}: {error}",
315                        snapshot.run_id, cancelled.status
316                    ),
317                )),
318                Ok(cancelled) => Err(WorkflowRunError::Storage(format!(
319                    "run index persistence failed; orphan run {} remains {:?}; repair with this run id: {error}",
320                    snapshot.run_id, cancelled.status
321                ))),
322                Err(cancel_error) => Err(WorkflowRunError::Storage(format!(
323                    "run index persistence failed; orphan run {} could not be cancelled ({cancel_error}); repair with this run id: {error}",
324                    snapshot.run_id
325                ))),
326            };
327        }
328        Ok(snapshot)
329    }
330
331    async fn ensure_run_index_capacity(&self, session_id: &str) -> Result<(), WorkflowRunError> {
332        let session = self
333            .sessions
334            .try_load(session_id)
335            .await
336            .map_err(|_| WorkflowRunError::Storage("session state unavailable".to_string()))?
337            .ok_or(WorkflowRunError::NotFound)?;
338        let ids = session
339            .metadata
340            .get(bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY)
341            .and_then(|raw| serde_json::from_str::<Vec<String>>(raw).ok())
342            .unwrap_or_default();
343        if ids.len() < MAX_WORKFLOW_RUN_IDS_PER_SESSION {
344            return Ok(());
345        }
346        let mut evictable = BTreeSet::new();
347        for run_id in &ids {
348            match self.engine.progress(run_id, u64::MAX).await {
349                Ok(progress) if progress.snapshot.status.is_terminal() => {
350                    evictable.insert(run_id.clone());
351                }
352                Err(WorkflowRunError::NotFound) => {
353                    evictable.insert(run_id.clone());
354                }
355                Ok(_) => {}
356                Err(error) => return Err(error),
357            }
358        }
359        if !evictable.is_empty() {
360            self.sessions
361                .update_runtime_session(
362                    session_id,
363                    &[bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY],
364                    move |session| {
365                        let mut ids = session
366                            .metadata
367                            .get(bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY)
368                            .and_then(|raw| serde_json::from_str::<Vec<String>>(raw).ok())
369                            .unwrap_or_default();
370                        ids.retain(|id| !evictable.contains(id));
371                        session.metadata.insert(
372                            bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY.to_string(),
373                            serde_json::to_string(&ids)
374                                .expect("string vector serialization cannot fail"),
375                        );
376                    },
377                )
378                .await
379                .map_err(|_| WorkflowRunError::Storage("run index pruning failed".to_string()))?
380                .ok_or(WorkflowRunError::NotFound)?;
381        }
382        let remaining = self
383            .sessions
384            .try_load(session_id)
385            .await
386            .map_err(|_| WorkflowRunError::Storage("session state unavailable".to_string()))?
387            .and_then(|session| {
388                session
389                    .metadata
390                    .get(bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY)
391                    .and_then(|raw| serde_json::from_str::<Vec<String>>(raw).ok())
392            })
393            .unwrap_or_default()
394            .len();
395        if remaining >= MAX_WORKFLOW_RUN_IDS_PER_SESSION {
396            return Err(WorkflowRunError::Preflight(
397                "workflow run index is full of active runs".to_string(),
398            ));
399        }
400        Ok(())
401    }
402
403    async fn remember_run_id(
404        &self,
405        session_id: &str,
406        run_id: &str,
407    ) -> Result<(), WorkflowRunError> {
408        let run_id = run_id.to_string();
409        let index_full = Arc::new(AtomicBool::new(false));
410        let index_full_in_transaction = index_full.clone();
411        self.sessions
412            .update_runtime_session(
413                session_id,
414                &[bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY],
415                move |session| {
416                    let mut ids = session
417                        .metadata
418                        .get(bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY)
419                        .and_then(|raw| serde_json::from_str::<Vec<String>>(raw).ok())
420                        .unwrap_or_default();
421                    ids.retain(|existing| existing != &run_id);
422                    if ids.len() >= MAX_WORKFLOW_RUN_IDS_PER_SESSION {
423                        index_full_in_transaction.store(true, Ordering::SeqCst);
424                        return;
425                    }
426                    ids.push(run_id);
427                    session.metadata.insert(
428                        bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY.to_string(),
429                        serde_json::to_string(&ids)
430                            .expect("string vector serialization cannot fail"),
431                    );
432                },
433            )
434            .await
435            .map_err(|_| WorkflowRunError::Storage("run index persistence failed".to_string()))?
436            .ok_or(WorkflowRunError::NotFound)
437            .and_then(|_| {
438                if index_full.load(Ordering::SeqCst) {
439                    Err(WorkflowRunError::Storage(
440                        "workflow run index reached its active-run capacity".to_string(),
441                    ))
442                } else {
443                    Ok(())
444                }
445            })
446    }
447
448    pub async fn list_for_session(
449        &self,
450        session_id: &str,
451    ) -> Result<Vec<WorkflowRunSnapshot>, WorkflowRunError> {
452        self.session_context(session_id).await?;
453        let session = self
454            .sessions
455            .try_load(session_id)
456            .await
457            .map_err(|_| WorkflowRunError::Storage("session state unavailable".to_string()))?
458            .ok_or(WorkflowRunError::NotFound)?;
459        let run_ids = session
460            .metadata
461            .get(bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY)
462            .and_then(|raw| serde_json::from_str::<Vec<String>>(raw).ok())
463            .unwrap_or_default();
464        let mut snapshots = Vec::new();
465        let mut stale = BTreeSet::new();
466        for run_id in run_ids {
467            match self.engine.progress(&run_id, u64::MAX).await {
468                Ok(progress) if progress.snapshot.session_id == session_id => {
469                    snapshots.push(progress.snapshot);
470                }
471                Ok(_) | Err(WorkflowRunError::NotFound) => {
472                    stale.insert(run_id);
473                }
474                Err(error) => return Err(error),
475            }
476        }
477        if !stale.is_empty() {
478            let _ = self
479                .sessions
480                .update_runtime_session(
481                    session_id,
482                    &[bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY],
483                    move |session| {
484                        let mut ids = session
485                            .metadata
486                            .get(bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY)
487                            .and_then(|raw| serde_json::from_str::<Vec<String>>(raw).ok())
488                            .unwrap_or_default();
489                        ids.retain(|id| !stale.contains(id));
490                        session.metadata.insert(
491                            bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY.to_string(),
492                            serde_json::to_string(&ids)
493                                .expect("string vector serialization cannot fail"),
494                        );
495                    },
496                )
497                .await;
498        }
499        snapshots.sort_by_key(|snapshot| std::cmp::Reverse(snapshot.created_at));
500        Ok(snapshots)
501    }
502
503    pub async fn progress_for_session(
504        &self,
505        session_id: &str,
506        run_id: &str,
507        since: u64,
508    ) -> Result<WorkflowProgress, WorkflowRunError> {
509        let progress = self.engine.progress(run_id, since).await?;
510        if progress.snapshot.session_id != session_id {
511            return Err(WorkflowRunError::NotFound);
512        }
513        Ok(progress)
514    }
515
516    pub async fn cancel_for_session(
517        &self,
518        session_id: &str,
519        run_id: &str,
520    ) -> Result<WorkflowRunSnapshot, WorkflowRunError> {
521        let progress = self
522            .progress_for_session(session_id, run_id, u64::MAX)
523            .await?;
524        ensure_workflow_cancel_allowed(&progress.snapshot)?;
525        self.engine.cancel(run_id).await
526    }
527
528    pub async fn restart_for_session(
529        &self,
530        session_id: &str,
531        run_id: &str,
532    ) -> Result<WorkflowRunSnapshot, WorkflowRunError> {
533        let progress = self
534            .progress_for_session(session_id, run_id, u64::MAX)
535            .await?;
536        ensure_workflow_restart_as_new_run_allowed(&progress.snapshot)?;
537        self.ensure_run_index_capacity(session_id).await?;
538        let (_, workspace_trusted) = self.session_context(session_id).await?;
539        let snapshot = self
540            .engine
541            .restart(run_id, workspace_trusted, vec!["read".to_string()])
542            .await?;
543        self.index_started_run_or_compensate(session_id, snapshot)
544            .await
545    }
546
547    pub async fn restart_from_tool(
548        &self,
549        session_id: &str,
550        run_id: &str,
551    ) -> Result<WorkflowRunSnapshot, WorkflowRunError> {
552        let progress = self
553            .progress_for_session(session_id, run_id, u64::MAX)
554            .await?;
555        if progress.snapshot.definition_bundle.root_invocation_policy["automatic"].as_bool()
556            != Some(true)
557        {
558            return Err(WorkflowRunError::Preflight(
559                "pinned workflow invocation policy denies automatic restart".to_string(),
560            ));
561        }
562        let session = self
563            .sessions
564            .try_load(session_id)
565            .await
566            .map_err(|_| WorkflowRunError::Preflight("session state is unavailable".to_string()))?
567            .ok_or(WorkflowRunError::NotFound)?;
568        let opted_in = session
569            .metadata
570            .get(bamboo_skills::WORKFLOW_ORCHESTRATION_OPT_IN_METADATA_KEY)
571            .is_some_and(|value| value.eq_ignore_ascii_case("true"));
572        if !opted_in {
573            return Err(WorkflowRunError::Preflight(
574                "model-started orchestration restart requires explicit session opt-in".to_string(),
575            ));
576        }
577        self.restart_for_session(session_id, run_id).await
578    }
579}
580
581fn ensure_workflow_cancel_allowed(snapshot: &WorkflowRunSnapshot) -> Result<(), WorkflowRunError> {
582    match snapshot.status {
583        WorkflowRunStatus::Succeeded | WorkflowRunStatus::Failed => Err(WorkflowRunError::Terminal),
584        WorkflowRunStatus::Queued
585        | WorkflowRunStatus::Running
586        | WorkflowRunStatus::Suspended
587        | WorkflowRunStatus::Cancelled => Ok(()),
588    }
589}
590
591fn ensure_workflow_restart_as_new_run_allowed(
592    snapshot: &WorkflowRunSnapshot,
593) -> Result<(), WorkflowRunError> {
594    match (snapshot.status, snapshot.suspension.as_ref()) {
595        (WorkflowRunStatus::Suspended, Some(WorkflowSuspensionContext::Recovery { .. })) => Ok(()),
596        (status, _) if status.is_terminal() => Err(WorkflowRunError::Terminal),
597        _ => Err(WorkflowRunError::Preflight(
598            "only recovery-suspended workflows can restart as a new run".to_string(),
599        )),
600    }
601}
602
603fn tighten_workflow_budget(
604    definition: &WorkflowBudgets,
605    requested: &WorkflowBudgets,
606) -> WorkflowBudgets {
607    WorkflowBudgets {
608        max_concurrency: definition.max_concurrency.min(requested.max_concurrency),
609        max_agents: definition.max_agents.min(requested.max_agents),
610        max_steps: definition.max_steps.min(requested.max_steps),
611        max_retries: definition.max_retries.min(requested.max_retries),
612        max_nesting_depth: definition
613            .max_nesting_depth
614            .min(requested.max_nesting_depth),
615        wall_time_ms: definition.wall_time_ms.min(requested.wall_time_ms),
616        max_tokens: match (definition.max_tokens, requested.max_tokens) {
617            (Some(left), Some(right)) => Some(left.min(right)),
618            (Some(value), None) | (None, Some(value)) => Some(value),
619            (None, None) => None,
620        },
621        max_cost_micros: match (definition.max_cost_micros, requested.max_cost_micros) {
622            (Some(left), Some(right)) => Some(left.min(right)),
623            (Some(value), None) | (None, Some(value)) => Some(value),
624            (None, None) => None,
625        },
626    }
627}
628
629fn validate_requested_budget(requested: &WorkflowBudgets) -> Result<(), WorkflowRunError> {
630    if requested.max_concurrency == 0
631        || requested.max_steps == 0
632        || requested.max_nesting_depth == 0
633        || requested.wall_time_ms == 0
634    {
635        return Err(WorkflowRunError::InvalidInput(
636            "workflow execution limits must be positive".to_string(),
637        ));
638    }
639    Ok(())
640}
641
642fn enforce_pinned_bundle_limits(
643    definition_count: usize,
644    serialized_bytes: usize,
645) -> Result<(), WorkflowRunError> {
646    if definition_count > MAX_PINNED_DEFINITIONS_PER_RUN {
647        return Err(WorkflowRunError::Preflight(
648            "workflow dependency count exceeds the server limit".to_string(),
649        ));
650    }
651    if serialized_bytes > MAX_PINNED_BUNDLE_BYTES_PER_RUN {
652        return Err(WorkflowRunError::Preflight(
653            "workflow definition bundle exceeds the server size limit".to_string(),
654        ));
655    }
656    Ok(())
657}
658
659/// Stable metadata-only projection of a durable workflow run.
660///
661/// The internal snapshot deliberately persists the complete definition bundle,
662/// validated arguments, input hashes, tool outputs, and diagnostic strings for
663/// deterministic recovery. None of those values are public API data. HTTP and
664/// model-facing tools must serialize this projection instead of the durable
665/// snapshot itself.
666#[derive(Debug, Clone, Serialize, PartialEq)]
667pub(crate) struct PublicWorkflowRunSnapshot {
668    pub run_id: String,
669    #[serde(skip_serializing_if = "Option::is_none")]
670    pub parent_run_id: Option<String>,
671    #[serde(skip_serializing_if = "Option::is_none")]
672    pub parent_step_id: Option<String>,
673    pub session_id: String,
674    pub workflow_id: String,
675    pub workflow_revision: u64,
676    pub definition_bundle_hash: String,
677    pub status: WorkflowRunStatus,
678    pub can_cancel: bool,
679    pub can_restart_as_new_run: bool,
680    pub planned_steps: BTreeMap<String, PublicWorkflowPlannedStep>,
681    pub plan: PublicWorkflowPlan,
682    pub steps: BTreeMap<String, PublicWorkflowStepSnapshot>,
683    pub budget: WorkflowBudgets,
684    pub usage: WorkflowBudgetUsage,
685    pub child_agent_count: u32,
686    pub last_sequence: u64,
687    #[serde(skip_serializing_if = "Option::is_none")]
688    pub failure: Option<PublicWorkflowFailure>,
689    #[serde(skip_serializing_if = "Option::is_none")]
690    pub suspension: Option<PublicWorkflowSuspension>,
691    pub created_at: chrono::DateTime<chrono::Utc>,
692    pub updated_at: chrono::DateTime<chrono::Utc>,
693}
694
695#[derive(Debug, Clone, Serialize, PartialEq)]
696pub(crate) struct PublicWorkflowStepSnapshot {
697    pub id: String,
698    pub status: WorkflowStepStatus,
699    #[serde(skip_serializing_if = "Option::is_none")]
700    pub failure: Option<PublicWorkflowFailure>,
701    pub attempts: u32,
702}
703
704#[derive(Debug, Clone, Serialize, PartialEq, Eq)]
705pub(crate) struct PublicWorkflowPlannedStep {
706    pub id: String,
707    pub kind: PublicWorkflowPlannedStepKind,
708}
709
710#[derive(Debug, Clone, Copy, Serialize, PartialEq, Eq)]
711#[serde(rename_all = "snake_case")]
712pub(crate) enum PublicWorkflowPlannedStepKind {
713    Tool,
714    Agent,
715    Workflow,
716}
717
718/// Safe plan topology. Data bindings, map item names, delays, tool arguments,
719/// prompts, schemas, and capabilities remain internal.
720#[derive(Debug, Clone, Serialize, PartialEq, Eq)]
721#[serde(tag = "type", rename_all = "snake_case")]
722pub(crate) enum PublicWorkflowPlan {
723    Step {
724        step: String,
725    },
726    Sequence {
727        nodes: Vec<PublicWorkflowPlan>,
728    },
729    Parallel {
730        nodes: Vec<PublicWorkflowPlan>,
731    },
732    Map {
733        body: Box<PublicWorkflowPlan>,
734    },
735    Retry {
736        node: Box<PublicWorkflowPlan>,
737        max_attempts: u32,
738    },
739}
740
741#[derive(Debug, Clone, Serialize, PartialEq, Eq)]
742pub(crate) struct PublicWorkflowFailure {
743    pub code: WorkflowFailureCode,
744    pub message: &'static str,
745    pub retryable: bool,
746}
747
748#[derive(Debug, Clone, Serialize, PartialEq, Eq)]
749#[serde(tag = "type", rename_all = "snake_case")]
750pub(crate) enum PublicWorkflowSuspension {
751    ToolApproval { step_id: String },
752    ToolRunning { step_id: String, killed: bool },
753    Recovery,
754}
755
756/// Durable reconnect events use the same metadata-only boundary as snapshots.
757/// Raw step/run outputs, suspension reasons, and backend diagnostics are never
758/// serialized to HTTP, SSE consumers, or model-facing tools.
759#[derive(Debug, Clone, Serialize, PartialEq)]
760pub(crate) struct PublicWorkflowRunEvent {
761    pub run_id: String,
762    pub sequence: u64,
763    pub at: chrono::DateTime<chrono::Utc>,
764    #[serde(skip_serializing_if = "Option::is_none")]
765    pub step_id: Option<String>,
766    #[serde(flatten)]
767    pub kind: PublicWorkflowRunEventKind,
768}
769
770#[derive(Debug, Clone, Serialize, PartialEq)]
771#[serde(tag = "type", rename_all = "snake_case")]
772pub(crate) enum PublicWorkflowRunEventKind {
773    RunQueued,
774    RunStarted,
775    Phase { name: &'static str },
776    StepQueued,
777    StepStarted,
778    StepSuspended,
779    StepCompleted,
780    StepFailed { failure: PublicWorkflowFailure },
781    StepCancelled,
782    StepSkipped,
783    RunSuspended,
784    RunSucceeded,
785    RunFailed { failure: PublicWorkflowFailure },
786    RunCancelled,
787}
788
789fn public_workflow_plan(plan: &WorkflowPlan) -> PublicWorkflowPlan {
790    match plan {
791        WorkflowPlan::Step { step } => PublicWorkflowPlan::Step { step: step.clone() },
792        WorkflowPlan::Sequence { nodes } => PublicWorkflowPlan::Sequence {
793            nodes: nodes.iter().map(public_workflow_plan).collect(),
794        },
795        WorkflowPlan::Parallel { nodes } => PublicWorkflowPlan::Parallel {
796            nodes: nodes.iter().map(public_workflow_plan).collect(),
797        },
798        WorkflowPlan::Map { body, .. } => PublicWorkflowPlan::Map {
799            body: Box::new(public_workflow_plan(body)),
800        },
801        WorkflowPlan::Retry {
802            node, max_attempts, ..
803        } => PublicWorkflowPlan::Retry {
804            node: Box::new(public_workflow_plan(node)),
805            max_attempts: *max_attempts,
806        },
807    }
808}
809
810fn public_planned_steps(
811    definition: &WorkflowRunDefinition,
812) -> BTreeMap<String, PublicWorkflowPlannedStep> {
813    definition
814        .steps
815        .iter()
816        .map(|step| {
817            let kind = match &step.kind {
818                WorkflowStepKind::Tool { .. } => PublicWorkflowPlannedStepKind::Tool,
819                WorkflowStepKind::Agent { .. } => PublicWorkflowPlannedStepKind::Agent,
820                WorkflowStepKind::Workflow { .. } => PublicWorkflowPlannedStepKind::Workflow,
821            };
822            (
823                step.id.clone(),
824                PublicWorkflowPlannedStep {
825                    id: step.id.clone(),
826                    kind,
827                },
828            )
829        })
830        .collect()
831}
832
833fn public_workflow_failure(failure: WorkflowFailure) -> PublicWorkflowFailure {
834    let message = match failure.code {
835        WorkflowFailureCode::InvalidDefinition => "Workflow definition is invalid",
836        WorkflowFailureCode::InvalidInput => "Workflow input is invalid",
837        WorkflowFailureCode::InvalidOutput => "Workflow output is invalid",
838        WorkflowFailureCode::UnknownReference => "Workflow reference is unavailable",
839        WorkflowFailureCode::PermissionDenied => "Workflow permission was denied",
840        WorkflowFailureCode::UntrustedWorkspace => "Workflow workspace is not trusted",
841        WorkflowFailureCode::BudgetExceeded => "Workflow execution budget was exceeded",
842        WorkflowFailureCode::RetryExhausted => "Workflow retry budget was exhausted",
843        WorkflowFailureCode::ExecutionFailed => "Workflow execution failed",
844        WorkflowFailureCode::Cancelled => "Workflow execution was cancelled",
845        WorkflowFailureCode::RecoverySuspended => "Workflow recovery requires attention",
846        WorkflowFailureCode::Suspended => "Workflow execution is suspended",
847        WorkflowFailureCode::DependencySkipped => "Workflow dependency was skipped",
848        WorkflowFailureCode::Storage => "Workflow storage is unavailable",
849    };
850    PublicWorkflowFailure {
851        code: failure.code,
852        message,
853        retryable: failure.retryable,
854    }
855}
856
857fn public_workflow_step(step: WorkflowStepSnapshot) -> PublicWorkflowStepSnapshot {
858    PublicWorkflowStepSnapshot {
859        id: step.id,
860        status: step.status,
861        failure: step.failure.map(public_workflow_failure),
862        attempts: step.attempts,
863    }
864}
865
866fn public_workflow_suspension(suspension: WorkflowSuspensionContext) -> PublicWorkflowSuspension {
867    match suspension {
868        WorkflowSuspensionContext::ToolApproval { step_id, .. } => {
869            PublicWorkflowSuspension::ToolApproval { step_id }
870        }
871        WorkflowSuspensionContext::ToolRunning {
872            step_id, killed, ..
873        } => PublicWorkflowSuspension::ToolRunning { step_id, killed },
874        WorkflowSuspensionContext::Recovery { .. } => PublicWorkflowSuspension::Recovery,
875    }
876}
877
878fn public_workflow_phase(name: &str) -> &'static str {
879    match name {
880        "retry_reserved" => "retry_reserved",
881        "step_reserved" => "step_reserved",
882        "agent_reserved" => "agent_reserved",
883        "agent_usage_recorded" => "agent_usage_recorded",
884        "suspension_context_persisted" => "suspension_context_persisted",
885        _ => "workflow_progressed",
886    }
887}
888
889pub(crate) fn public_workflow_snapshot(snapshot: WorkflowRunSnapshot) -> PublicWorkflowRunSnapshot {
890    let can_cancel = ensure_workflow_cancel_allowed(&snapshot).is_ok();
891    let can_restart_as_new_run = ensure_workflow_restart_as_new_run_allowed(&snapshot).is_ok();
892    let planned_steps = public_planned_steps(&snapshot.definition);
893    let plan = public_workflow_plan(&snapshot.definition.plan);
894    let child_agent_count = snapshot.usage.agents;
895    PublicWorkflowRunSnapshot {
896        run_id: snapshot.run_id,
897        parent_run_id: snapshot.parent_run_id,
898        parent_step_id: snapshot.parent_step_id,
899        session_id: snapshot.session_id,
900        workflow_id: snapshot.definition.id,
901        workflow_revision: snapshot.definition.revision,
902        definition_bundle_hash: snapshot.definition_bundle_hash,
903        status: snapshot.status,
904        can_cancel,
905        can_restart_as_new_run,
906        planned_steps,
907        plan,
908        steps: snapshot
909            .steps
910            .into_iter()
911            .map(|(id, step)| (id, public_workflow_step(step)))
912            .collect(),
913        budget: snapshot.definition.budgets,
914        usage: snapshot.usage,
915        child_agent_count,
916        last_sequence: snapshot.last_sequence,
917        failure: snapshot.failure.map(public_workflow_failure),
918        suspension: snapshot.suspension.map(public_workflow_suspension),
919        created_at: snapshot.created_at,
920        updated_at: snapshot.updated_at,
921    }
922}
923
924pub(crate) fn public_workflow_event(event: WorkflowRunEvent) -> PublicWorkflowRunEvent {
925    let kind = match event.kind {
926        WorkflowRunEventKind::RunQueued => PublicWorkflowRunEventKind::RunQueued,
927        WorkflowRunEventKind::RunStarted => PublicWorkflowRunEventKind::RunStarted,
928        WorkflowRunEventKind::Phase { name } => PublicWorkflowRunEventKind::Phase {
929            name: public_workflow_phase(&name),
930        },
931        WorkflowRunEventKind::StepQueued => PublicWorkflowRunEventKind::StepQueued,
932        WorkflowRunEventKind::StepStarted => PublicWorkflowRunEventKind::StepStarted,
933        WorkflowRunEventKind::StepSuspended { .. } => PublicWorkflowRunEventKind::StepSuspended,
934        WorkflowRunEventKind::StepCompleted { .. } => PublicWorkflowRunEventKind::StepCompleted,
935        WorkflowRunEventKind::StepFailed { failure } => PublicWorkflowRunEventKind::StepFailed {
936            failure: public_workflow_failure(failure),
937        },
938        WorkflowRunEventKind::StepCancelled => PublicWorkflowRunEventKind::StepCancelled,
939        WorkflowRunEventKind::StepSkipped { .. } => PublicWorkflowRunEventKind::StepSkipped,
940        WorkflowRunEventKind::RunSuspended { .. } => PublicWorkflowRunEventKind::RunSuspended,
941        WorkflowRunEventKind::RunSucceeded { .. } => PublicWorkflowRunEventKind::RunSucceeded,
942        WorkflowRunEventKind::RunFailed { failure } => PublicWorkflowRunEventKind::RunFailed {
943            failure: public_workflow_failure(failure),
944        },
945        WorkflowRunEventKind::RunCancelled => PublicWorkflowRunEventKind::RunCancelled,
946    };
947    PublicWorkflowRunEvent {
948        run_id: event.run_id,
949        sequence: event.sequence,
950        at: event.at,
951        step_id: event.step_id,
952        kind,
953    }
954}
955
956/// The catalog adapter above always calls `start_pinned`. Generic engine starts
957/// remain unavailable so no future server call site can accidentally re-read a
958/// live definition or mix publications.
959struct ExternallyPinnedDefinitions;
960
961#[async_trait]
962impl WorkflowDefinitionPort for ExternallyPinnedDefinitions {
963    async fn pin_bundle(
964        &self,
965        _root: &WorkflowRunDefinition,
966    ) -> Result<WorkflowDefinitionBundle, String> {
967        Err("server workflows must be pinned through SkillManager".to_string())
968    }
969}
970
971/// #563 named-agent registry integration is not complete. Unknown and named
972/// agents therefore fail preflight rather than falling back to a dynamic agent.
973struct UnavailableAgentPort;
974
975#[async_trait]
976impl AgentStepPort for UnavailableAgentPort {
977    async fn resolve(&self, _name: &str) -> Result<Option<NamedAgentSpec>, String> {
978        Ok(None)
979    }
980
981    async fn execute(
982        &self,
983        _spec: &NamedAgentSpec,
984        _prompt: Value,
985        _model: Option<&str>,
986        _effort: Option<&str>,
987        _capabilities: &BTreeSet<String>,
988        _session_id: &str,
989    ) -> Result<AgentStepResult, String> {
990        Err("named-agent execution is not available".to_string())
991    }
992}
993
994struct ServerWorkflowPolicy;
995
996#[async_trait]
997impl WorkflowPolicyPort for ServerWorkflowPolicy {
998    async fn authorize(
999        &self,
1000        _session_id: &str,
1001        target: &WorkflowPolicyTarget,
1002        requested: &BTreeSet<String>,
1003        _workspace_trusted: bool,
1004    ) -> PermissionDecision {
1005        // Definition-declared capabilities are descriptive input, not an
1006        // authority. Until #601 provides a server-owned capability registry,
1007        // bind authorization to this strict server-owned target allowlist too.
1008        // Unknown/MCP/network/script/write targets therefore fail closed even
1009        // when hostile YAML claims `capabilities: []` or `[read]`.
1010        match target {
1011            // The reference itself belongs to the same immutable, validated
1012            // bundle. Every concrete step in the nested definition is still
1013            // authorized independently during this preflight walk.
1014            WorkflowPolicyTarget::Workflow { .. } if requested.is_empty() => {
1015                PermissionDecision::Allow
1016            }
1017            WorkflowPolicyTarget::Tool(name) => {
1018                let safe_target = SAFE_UNTRUSTED_WORKFLOW_TOOLS
1019                    .iter()
1020                    .any(|candidate| candidate.eq_ignore_ascii_case(name));
1021                if safe_target && requested.iter().all(|capability| capability == "read") {
1022                    PermissionDecision::Allow
1023                } else {
1024                    PermissionDecision::Deny(
1025                        "workflow capability authority is not available".to_string(),
1026                    )
1027                }
1028            }
1029            WorkflowPolicyTarget::Agent(_) | WorkflowPolicyTarget::Workflow { .. } => {
1030                PermissionDecision::Deny(
1031                    "workflow capability authority is not available".to_string(),
1032                )
1033            }
1034        }
1035    }
1036}
1037
1038struct UnavailableSecretResolver;
1039
1040#[async_trait]
1041impl WorkflowSecretResolverPort for UnavailableSecretResolver {
1042    async fn resolve(
1043        &self,
1044        _session_id: &str,
1045        _capability: &str,
1046    ) -> Result<WorkflowSecretMaterial, String> {
1047        Err("workflow secret capability resolver is not available".to_string())
1048    }
1049}
1050
1051#[derive(Debug, Deserialize)]
1052#[serde(tag = "action", rename_all = "snake_case", deny_unknown_fields)]
1053enum WorkflowToolInput {
1054    Start {
1055        workflow_id: String,
1056        revision: u64,
1057        #[serde(default = "empty_object")]
1058        args: Value,
1059        #[serde(default)]
1060        budget: Option<WorkflowBudgets>,
1061    },
1062    List {},
1063    Get {
1064        run_id: String,
1065    },
1066    Events {
1067        run_id: String,
1068        #[serde(default)]
1069        since: u64,
1070    },
1071    Cancel {
1072        run_id: String,
1073    },
1074    Restart {
1075        run_id: String,
1076    },
1077}
1078
1079fn empty_object() -> Value {
1080    json!({})
1081}
1082
1083pub struct WorkflowRunTool {
1084    access: WorkflowRunAccess,
1085}
1086
1087impl WorkflowRunTool {
1088    pub fn new(access: WorkflowRunAccess) -> Self {
1089        Self { access }
1090    }
1091}
1092
1093#[async_trait]
1094impl Tool for WorkflowRunTool {
1095    fn name(&self) -> &str {
1096        "workflow_run"
1097    }
1098
1099    fn description(&self) -> &str {
1100        "Start, inspect, cancel, or safely restart a catalog-pinned workflow run"
1101    }
1102
1103    fn parameters_schema(&self) -> Value {
1104        json!({
1105            "type": "object",
1106            "properties": {
1107                "action": {
1108                    "type": "string",
1109                    "enum": ["start", "list", "get", "events", "cancel", "restart"]
1110                },
1111                "workflow_id": {"type": "string", "minLength": 1},
1112                "revision": {"type": "integer", "minimum": 1},
1113                "args": {"type": "object", "default": {}},
1114                "budget": {
1115                    "type": "object",
1116                    "properties": {
1117                        "max_concurrency": {"type": "integer", "minimum": 1},
1118                        "max_agents": {"type": "integer", "minimum": 0},
1119                        "max_steps": {"type": "integer", "minimum": 1},
1120                        "max_retries": {"type": "integer", "minimum": 0},
1121                        "max_nesting_depth": {"type": "integer", "minimum": 1},
1122                        "wall_time_ms": {"type": "integer", "minimum": 1},
1123                        "max_tokens": {"type": "integer", "minimum": 0},
1124                        "max_cost_micros": {"type": "integer", "minimum": 0}
1125                    },
1126                    "required": [
1127                        "max_concurrency",
1128                        "max_agents",
1129                        "max_steps",
1130                        "max_retries",
1131                        "max_nesting_depth",
1132                        "wall_time_ms"
1133                    ],
1134                    "additionalProperties": false
1135                },
1136                "run_id": {"type": "string"},
1137                "since": {"type": "integer", "minimum": 0}
1138            },
1139            "required": ["action"],
1140            "additionalProperties": false
1141        })
1142    }
1143
1144    fn classify(&self, args: &Value) -> ToolClass {
1145        match args.get("action").and_then(Value::as_str) {
1146            Some("get" | "list" | "events") => ToolClass::READONLY_PARALLEL,
1147            _ => ToolClass::MUTATING_SERIAL,
1148        }
1149    }
1150
1151    async fn invoke(&self, args: Value, ctx: ToolCtx) -> Result<ToolOutcome, ToolError> {
1152        let input: WorkflowToolInput = serde_json::from_value(args)
1153            .map_err(|error| ToolError::InvalidArguments(error.to_string()))?;
1154        let session_id = ctx.session_id().ok_or_else(|| {
1155            ToolError::InvalidArguments("workflow_run requires a session".to_string())
1156        })?;
1157        let result = match input {
1158            WorkflowToolInput::Start {
1159                workflow_id,
1160                revision,
1161                args,
1162                budget,
1163            } => serde_json::to_value(public_workflow_snapshot(
1164                self.access
1165                    .start_from_tool(session_id, &workflow_id, revision, args, budget)
1166                    .await
1167                    .map_err(workflow_tool_error)?,
1168            )),
1169            WorkflowToolInput::List {} => serde_json::to_value(
1170                self.access
1171                    .list_for_session(session_id)
1172                    .await
1173                    .map_err(workflow_tool_error)?
1174                    .into_iter()
1175                    .map(public_workflow_snapshot)
1176                    .collect::<Vec<_>>(),
1177            ),
1178            WorkflowToolInput::Get { run_id } => {
1179                let progress = self
1180                    .access
1181                    .progress_for_session(session_id, &run_id, u64::MAX)
1182                    .await
1183                    .map_err(workflow_tool_error)?;
1184                serde_json::to_value(public_workflow_snapshot(progress.snapshot))
1185            }
1186            WorkflowToolInput::Events { run_id, since } => {
1187                let progress = self
1188                    .access
1189                    .progress_for_session(session_id, &run_id, since)
1190                    .await
1191                    .map_err(workflow_tool_error)?;
1192                serde_json::to_value(
1193                    progress
1194                        .events
1195                        .into_iter()
1196                        .map(public_workflow_event)
1197                        .collect::<Vec<_>>(),
1198                )
1199            }
1200            WorkflowToolInput::Cancel { run_id } => serde_json::to_value(public_workflow_snapshot(
1201                self.access
1202                    .cancel_for_session(session_id, &run_id)
1203                    .await
1204                    .map_err(workflow_tool_error)?,
1205            )),
1206            WorkflowToolInput::Restart { run_id } => {
1207                serde_json::to_value(public_workflow_snapshot(
1208                    self.access
1209                        .restart_from_tool(session_id, &run_id)
1210                        .await
1211                        .map_err(workflow_tool_error)?,
1212                ))
1213            }
1214        }
1215        .map_err(|error| ToolError::Execution(error.to_string()))?;
1216        Ok(ToolOutcome::Completed(ToolResult::text(
1217            true,
1218            serde_json::to_string(&result)
1219                .map_err(|error| ToolError::Execution(error.to_string()))?,
1220        )))
1221    }
1222}
1223
1224fn workflow_tool_error(error: WorkflowRunError) -> ToolError {
1225    match error {
1226        WorkflowRunError::InvalidInput(_) => {
1227            ToolError::InvalidArguments("workflow input is invalid".to_string())
1228        }
1229        WorkflowRunError::Compile(_) => {
1230            ToolError::InvalidArguments("workflow definition is invalid".to_string())
1231        }
1232        WorkflowRunError::Preflight(_) => {
1233            ToolError::InvalidArguments("workflow preflight failed".to_string())
1234        }
1235        WorkflowRunError::Storage(details) => {
1236            tracing::error!(%details, "workflow tool storage unavailable");
1237            ToolError::Execution("workflow storage unavailable".to_string())
1238        }
1239        WorkflowRunError::NotFound => ToolError::Execution("workflow run not found".to_string()),
1240        WorkflowRunError::Terminal => {
1241            ToolError::Execution("workflow run is already terminal".to_string())
1242        }
1243    }
1244}
1245
1246#[cfg(test)]
1247mod tests {
1248    use super::*;
1249    use bamboo_agent_core::storage::Storage;
1250    use bamboo_agent_core::tools::{
1251        FunctionSchema, ToolCall, ToolExecutor, ToolResult, ToolSchema,
1252    };
1253    use bamboo_agent_core::Session;
1254    use bamboo_engine::WorkflowRunRepository;
1255    use bamboo_llm::protocol::{gemini::GeminiTool, ToProvider};
1256    use std::collections::HashMap;
1257    use tokio::sync::{RwLock, Semaphore};
1258
1259    #[derive(Default)]
1260    struct WorkflowTestStorage {
1261        sessions: RwLock<HashMap<String, Session>>,
1262        fail_saves: AtomicBool,
1263    }
1264
1265    impl WorkflowTestStorage {
1266        fn set_fail_saves(&self, fail: bool) {
1267            self.fail_saves.store(fail, Ordering::SeqCst);
1268        }
1269    }
1270
1271    #[async_trait]
1272    impl Storage for WorkflowTestStorage {
1273        async fn save_session(&self, session: &Session) -> std::io::Result<()> {
1274            if self.fail_saves.load(Ordering::SeqCst) {
1275                return Err(std::io::Error::other(
1276                    "injected workflow session persistence failure",
1277                ));
1278            }
1279            self.sessions
1280                .write()
1281                .await
1282                .insert(session.id.clone(), session.clone());
1283            Ok(())
1284        }
1285
1286        async fn load_session(&self, session_id: &str) -> std::io::Result<Option<Session>> {
1287            Ok(self.sessions.read().await.get(session_id).cloned())
1288        }
1289
1290        async fn delete_session(&self, session_id: &str) -> std::io::Result<bool> {
1291            Ok(self.sessions.write().await.remove(session_id).is_some())
1292        }
1293    }
1294
1295    struct WorkflowReadTool;
1296
1297    #[async_trait]
1298    impl ToolExecutor for WorkflowReadTool {
1299        async fn execute(
1300            &self,
1301            _call: &ToolCall,
1302        ) -> bamboo_agent_core::tools::executor::Result<ToolResult> {
1303            Ok(ToolResult::text(true, r#"{"ok":true}"#))
1304        }
1305
1306        fn list_tools(&self) -> Vec<ToolSchema> {
1307            vec![ToolSchema {
1308                schema_type: "function".to_string(),
1309                function: FunctionSchema {
1310                    name: "Read".to_string(),
1311                    description: "read".to_string(),
1312                    parameters: serde_json::json!({"type":"object"}),
1313                },
1314            }]
1315        }
1316    }
1317
1318    struct BlockingWorkflowReadTool {
1319        entered: Arc<Semaphore>,
1320    }
1321
1322    #[async_trait]
1323    impl ToolExecutor for BlockingWorkflowReadTool {
1324        async fn execute(
1325            &self,
1326            _call: &ToolCall,
1327        ) -> bamboo_agent_core::tools::executor::Result<ToolResult> {
1328            self.entered.add_permits(1);
1329            std::future::pending().await
1330        }
1331
1332        fn list_tools(&self) -> Vec<ToolSchema> {
1333            WorkflowReadTool.list_tools()
1334        }
1335    }
1336
1337    async fn workflow_test_access_with_tools(
1338        tools: Arc<dyn ToolExecutor>,
1339    ) -> (
1340        WorkflowRunAccess,
1341        bamboo_engine::SessionRepository,
1342        tempfile::TempDir,
1343    ) {
1344        let (access, repo, directory, _) = workflow_test_access_with_tools_and_storage(tools).await;
1345        (access, repo, directory)
1346    }
1347
1348    async fn workflow_test_access_with_tools_and_storage(
1349        tools: Arc<dyn ToolExecutor>,
1350    ) -> (
1351        WorkflowRunAccess,
1352        bamboo_engine::SessionRepository,
1353        tempfile::TempDir,
1354        Arc<WorkflowTestStorage>,
1355    ) {
1356        let directory = tempfile::tempdir().expect("tempdir");
1357        let skills_dir = directory.path().join("skills");
1358        let root = skills_dir.join("review-flow");
1359        std::fs::create_dir_all(&root).expect("workflow dir");
1360        std::fs::write(
1361            root.join("SKILL.md"),
1362            "---\nname: review-flow\ndescription: Review flow\n---\nRun review flow.\n",
1363        )
1364        .expect("skill");
1365        std::fs::write(
1366            root.join("workflow.yaml"),
1367            "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",
1368        )
1369        .expect("workflow");
1370        let skills = Arc::new(SkillManager::with_config(bamboo_skills::SkillStoreConfig {
1371            skills_dir,
1372            ..Default::default()
1373        }));
1374        skills.initialize().await.expect("skills initialize");
1375        let storage = Arc::new(WorkflowTestStorage::default());
1376        let storage_port: Arc<dyn Storage> = storage.clone();
1377        let persistence = Arc::new(bamboo_storage::LockedSessionStore::new(
1378            storage_port.clone(),
1379        ));
1380        let cache = Arc::default();
1381        let repo = bamboo_engine::SessionRepository::new(cache, storage_port, persistence);
1382        let access = WorkflowRunAccess::new(directory.path(), tools, skills, repo.clone())
1383            .await
1384            .expect("workflow access");
1385        (access, repo, directory, storage)
1386    }
1387
1388    async fn workflow_test_access() -> (
1389        WorkflowRunAccess,
1390        bamboo_engine::SessionRepository,
1391        tempfile::TempDir,
1392    ) {
1393        workflow_test_access_with_tools(Arc::new(WorkflowReadTool)).await
1394    }
1395
1396    async fn seed_durable_workflow_run(
1397        access: &WorkflowRunAccess,
1398        directory: &Path,
1399        workspace: &Path,
1400        session_id: &str,
1401        run_id: &str,
1402        status: WorkflowRunStatus,
1403        suspension: Option<WorkflowSuspensionContext>,
1404    ) -> WorkflowRunSnapshot {
1405        let bundle = access
1406            .skills
1407            .pin_workflow_definition_bundle(Some(workspace), "review-flow", 42)
1408            .await
1409            .expect("pin workflow test bundle");
1410        let definition = bundle.root().cloned().expect("pinned root definition");
1411        let step_status = match status {
1412            WorkflowRunStatus::Queued => WorkflowStepStatus::Queued,
1413            WorkflowRunStatus::Running => WorkflowStepStatus::Running,
1414            WorkflowRunStatus::Suspended => WorkflowStepStatus::Suspended,
1415            WorkflowRunStatus::Succeeded => WorkflowStepStatus::Succeeded,
1416            WorkflowRunStatus::Failed => WorkflowStepStatus::Failed,
1417            WorkflowRunStatus::Cancelled => WorkflowStepStatus::Cancelled,
1418        };
1419        let now = chrono::Utc::now();
1420        let failure = (status == WorkflowRunStatus::Failed).then(|| WorkflowFailure {
1421            code: WorkflowFailureCode::ExecutionFailed,
1422            message: "seeded workflow failure".to_string(),
1423            retryable: false,
1424        });
1425        let snapshot = WorkflowRunSnapshot {
1426            run_id: run_id.to_string(),
1427            parent_run_id: None,
1428            parent_step_id: None,
1429            session_id: session_id.to_string(),
1430            definition: definition.clone(),
1431            definition_bundle: bundle,
1432            definition_bundle_hash: "seeded-public-bundle-hash".to_string(),
1433            validated_args: json!({}),
1434            status,
1435            steps: definition
1436                .steps
1437                .iter()
1438                .map(|step| {
1439                    (
1440                        step.id.clone(),
1441                        WorkflowStepSnapshot {
1442                            id: step.id.clone(),
1443                            status: step_status,
1444                            input_hash: "seeded-input-hash".to_string(),
1445                            output: None,
1446                            failure: failure.clone(),
1447                            attempts: 0,
1448                        },
1449                    )
1450                })
1451                .collect(),
1452            usage: WorkflowBudgetUsage::default(),
1453            last_sequence: 1,
1454            output: None,
1455            failure: failure.clone(),
1456            suspension,
1457            created_at: now,
1458            updated_at: now,
1459        };
1460        persist_durable_workflow_run(directory, &snapshot).await;
1461        snapshot
1462    }
1463
1464    async fn persist_durable_workflow_run(directory: &Path, snapshot: &WorkflowRunSnapshot) {
1465        let kind = match snapshot.status {
1466            WorkflowRunStatus::Queued => WorkflowRunEventKind::RunQueued,
1467            WorkflowRunStatus::Running => WorkflowRunEventKind::RunStarted,
1468            WorkflowRunStatus::Suspended => WorkflowRunEventKind::RunSuspended {
1469                reason: "seeded suspension".to_string(),
1470            },
1471            WorkflowRunStatus::Succeeded => WorkflowRunEventKind::RunSucceeded {
1472                output: Value::Null,
1473            },
1474            WorkflowRunStatus::Failed => WorkflowRunEventKind::RunFailed {
1475                failure: snapshot.failure.clone().expect("failed run has failure"),
1476            },
1477            WorkflowRunStatus::Cancelled => WorkflowRunEventKind::RunCancelled,
1478        };
1479        let repository = FileWorkflowRunRepository::new(directory.join("workflow-runs"))
1480            .expect("open workflow test repository");
1481        repository
1482            .create(
1483                snapshot,
1484                &WorkflowRunEvent {
1485                    run_id: snapshot.run_id.clone(),
1486                    sequence: snapshot.last_sequence,
1487                    at: snapshot.updated_at,
1488                    step_id: None,
1489                    kind,
1490                },
1491            )
1492            .await
1493            .expect("seed durable workflow run");
1494    }
1495
1496    async fn replace_workflow_run_index(
1497        repo: &bamboo_engine::SessionRepository,
1498        session_id: &str,
1499        run_ids: Vec<String>,
1500    ) {
1501        repo.update_runtime_session(
1502            session_id,
1503            &[bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY],
1504            move |session| {
1505                session.metadata.insert(
1506                    bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY.to_string(),
1507                    serde_json::to_string(&run_ids).expect("run index json"),
1508                );
1509            },
1510        )
1511        .await
1512        .expect("replace workflow run index")
1513        .expect("workflow session");
1514    }
1515
1516    async fn wait_for_workflow_run_to_settle(
1517        access: &WorkflowRunAccess,
1518        run_id: &str,
1519    ) -> WorkflowRunSnapshot {
1520        tokio::time::timeout(std::time::Duration::from_secs(2), async {
1521            loop {
1522                let snapshot = access
1523                    .engine
1524                    .progress(run_id, u64::MAX)
1525                    .await
1526                    .expect("workflow progress")
1527                    .snapshot;
1528                if snapshot.status.is_terminal() && !access.engine.is_run_active(run_id) {
1529                    break snapshot;
1530                }
1531                tokio::task::yield_now().await;
1532            }
1533        })
1534        .await
1535        .expect("workflow run settles")
1536    }
1537
1538    fn private_workflow_snapshot(status: WorkflowRunStatus) -> WorkflowRunSnapshot {
1539        let definition = WorkflowRunDefinition {
1540            workflow_schema: 1,
1541            id: "review-flow".to_string(),
1542            revision: 42,
1543            input_schema: json!({"private_schema": "PRIVATE-SCHEMA-SENTINEL"}),
1544            output_schema: Some(json!({"private_output": "PRIVATE-SCHEMA-SENTINEL"})),
1545            steps: vec![bamboo_domain::WorkflowStepDefinition {
1546                id: "inspect".to_string(),
1547                kind: WorkflowStepKind::Tool {
1548                    tool: "PRIVATE-TOOL-SENTINEL".to_string(),
1549                    args: json!({"credential": "PRIVATE-ARG-SENTINEL"}),
1550                    capabilities: vec!["PRIVATE-CAPABILITY-SENTINEL".to_string()],
1551                },
1552                failure: bamboo_domain::FailurePolicy::FailFast,
1553                output_schema: Some(json!({"private": "PRIVATE-STEP-SCHEMA-SENTINEL"})),
1554            }],
1555            plan: WorkflowPlan::Retry {
1556                node: Box::new(WorkflowPlan::Map {
1557                    source: bamboo_domain::ValueRef::Literal {
1558                        value: json!("PRIVATE-BINDING-SENTINEL"),
1559                    },
1560                    item: "PRIVATE-ITEM-SENTINEL".to_string(),
1561                    body: Box::new(WorkflowPlan::Step {
1562                        step: "inspect".to_string(),
1563                    }),
1564                }),
1565                max_attempts: 3,
1566                delay_ms: 987_654,
1567            },
1568            budgets: WorkflowBudgets {
1569                max_concurrency: 2,
1570                max_agents: 4,
1571                max_steps: 8,
1572                max_retries: 3,
1573                max_nesting_depth: 2,
1574                wall_time_ms: 10_000,
1575                max_tokens: Some(1_000),
1576                max_cost_micros: Some(2_000),
1577            },
1578        };
1579        let definition_bundle = WorkflowDefinitionBundle {
1580            publication_revision: 7,
1581            root_id: definition.id.clone(),
1582            root_revision: definition.revision,
1583            root_invocation_policy: json!({"private": "PRIVATE-POLICY-SENTINEL"}),
1584            definitions: BTreeMap::from([(
1585                WorkflowDefinitionBundle::key(&definition.id, definition.revision),
1586                definition.clone(),
1587            )]),
1588        };
1589        let now = chrono::Utc::now();
1590        WorkflowRunSnapshot {
1591            run_id: "public-run".to_string(),
1592            parent_run_id: Some("public-parent-run".to_string()),
1593            parent_step_id: Some("public-parent-step".to_string()),
1594            session_id: "public-session".to_string(),
1595            definition,
1596            definition_bundle,
1597            definition_bundle_hash: "public-bundle-hash".to_string(),
1598            validated_args: json!({"password": "PRIVATE-VALIDATED-ARG-SENTINEL"}),
1599            status,
1600            steps: BTreeMap::from([(
1601                "inspect".to_string(),
1602                WorkflowStepSnapshot {
1603                    id: "inspect".to_string(),
1604                    status: WorkflowStepStatus::Failed,
1605                    input_hash: "PRIVATE-INPUT-HASH-SENTINEL".to_string(),
1606                    output: Some(json!({"raw_tool_output": "PRIVATE-OUTPUT-SENTINEL"})),
1607                    failure: Some(WorkflowFailure {
1608                        code: WorkflowFailureCode::Storage,
1609                        message: "/private/workspace/PRIVATE-DIAGNOSTIC-SENTINEL".to_string(),
1610                        retryable: true,
1611                    }),
1612                    attempts: 2,
1613                },
1614            )]),
1615            usage: WorkflowBudgetUsage {
1616                steps: 1,
1617                retries: 1,
1618                agents: 3,
1619                tokens: 40,
1620                cost_micros: 50,
1621            },
1622            last_sequence: 9,
1623            output: Some(json!({"raw_run_output": "PRIVATE-RUN-OUTPUT-SENTINEL"})),
1624            failure: Some(WorkflowFailure {
1625                code: WorkflowFailureCode::ExecutionFailed,
1626                message: "credential PRIVATE-RUN-FAILURE-SENTINEL".to_string(),
1627                retryable: false,
1628            }),
1629            suspension: Some(WorkflowSuspensionContext::ToolApproval {
1630                step_id: "inspect".to_string(),
1631                tool: "PRIVATE-SUSPENSION-TOOL-SENTINEL".to_string(),
1632                tool_call_id: "PRIVATE-TOOL-CALL-SENTINEL".to_string(),
1633            }),
1634            created_at: now,
1635            updated_at: now,
1636        }
1637    }
1638
1639    #[test]
1640    fn public_workflow_snapshot_is_stable_metadata_only_for_every_status() {
1641        let statuses = [
1642            (WorkflowRunStatus::Queued, "queued", true),
1643            (WorkflowRunStatus::Running, "running", true),
1644            (WorkflowRunStatus::Suspended, "suspended", true),
1645            (WorkflowRunStatus::Succeeded, "succeeded", false),
1646            (WorkflowRunStatus::Failed, "failed", false),
1647            (WorkflowRunStatus::Cancelled, "cancelled", true),
1648        ];
1649
1650        for (status, wire_status, can_cancel) in statuses {
1651            let public =
1652                serde_json::to_value(public_workflow_snapshot(private_workflow_snapshot(status)))
1653                    .expect("public snapshot serializes");
1654            let text = public.to_string();
1655            assert_eq!(public["status"], wire_status);
1656            assert_eq!(public["can_cancel"], can_cancel);
1657            assert_eq!(public["can_restart_as_new_run"], false);
1658            assert_eq!(public["workflow_id"], "review-flow");
1659            assert_eq!(public["workflow_revision"], 42);
1660            assert_eq!(public["definition_bundle_hash"], "public-bundle-hash");
1661            assert_eq!(public["planned_steps"]["inspect"]["kind"], "tool");
1662            assert_eq!(public["plan"]["type"], "retry");
1663            assert_eq!(public["steps"]["inspect"]["attempts"], 2);
1664            assert_eq!(public["budget"]["max_steps"], 8);
1665            assert_eq!(public["usage"]["agents"], 3);
1666            assert_eq!(public["child_agent_count"], 3);
1667            assert_eq!(public["last_sequence"], 9);
1668            assert_eq!(public["failure"]["message"], "Workflow execution failed");
1669            assert_eq!(public["suspension"]["type"], "tool_approval");
1670            for internal_field in [
1671                "definition",
1672                "definition_bundle",
1673                "validated_args",
1674                "output",
1675            ] {
1676                assert!(
1677                    public.get(internal_field).is_none(),
1678                    "public snapshot exposed internal field {internal_field}: {text}"
1679                );
1680            }
1681
1682            for private in [
1683                "PRIVATE-",
1684                "validated_args",
1685                "input_hash",
1686                "output_schema",
1687                "root_invocation_policy",
1688                "delay_ms",
1689                "tool_call_id",
1690                "raw_tool_output",
1691                "raw_run_output",
1692            ] {
1693                assert!(
1694                    !text.contains(private),
1695                    "public snapshot leaked {private}: {text}"
1696                );
1697            }
1698        }
1699
1700        let mut recovery = private_workflow_snapshot(WorkflowRunStatus::Suspended);
1701        recovery.suspension = Some(WorkflowSuspensionContext::Recovery {
1702            reason: "PRIVATE-RECOVERY-REASON-SENTINEL".to_string(),
1703        });
1704        let public = serde_json::to_value(public_workflow_snapshot(recovery))
1705            .expect("public recovery snapshot serializes");
1706        let text = public.to_string();
1707        assert_eq!(public["can_cancel"], true);
1708        assert_eq!(public["can_restart_as_new_run"], true);
1709        assert_eq!(public["suspension"]["type"], "recovery");
1710        assert!(!text.contains("PRIVATE-RECOVERY-REASON-SENTINEL"));
1711    }
1712
1713    #[tokio::test]
1714    async fn public_workflow_actions_match_cancel_and_restart_endpoint_acceptance() {
1715        let (access, repo, directory) = workflow_test_access().await;
1716        let workspace = directory.path().join("action-workspace");
1717        std::fs::create_dir_all(&workspace).expect("action workspace");
1718        let session_id = "workflow-action-matrix";
1719        let mut session = Session::new(session_id, "model");
1720        session.workspace = Some(workspace.to_string_lossy().into_owned());
1721        repo.save(&mut session).await.expect("save action session");
1722
1723        let cases = vec![
1724            ("queued", WorkflowRunStatus::Queued, None, true, false),
1725            ("running", WorkflowRunStatus::Running, None, true, false),
1726            (
1727                "suspended_without_context",
1728                WorkflowRunStatus::Suspended,
1729                None,
1730                true,
1731                false,
1732            ),
1733            (
1734                "tool_approval",
1735                WorkflowRunStatus::Suspended,
1736                Some(WorkflowSuspensionContext::ToolApproval {
1737                    step_id: "inspect".to_string(),
1738                    tool: "Read".to_string(),
1739                    tool_call_id: "approval-call".to_string(),
1740                }),
1741                true,
1742                false,
1743            ),
1744            (
1745                "tool_running",
1746                WorkflowRunStatus::Suspended,
1747                Some(WorkflowSuspensionContext::ToolRunning {
1748                    step_id: "inspect".to_string(),
1749                    tool: "Read".to_string(),
1750                    tool_call_id: "running-call".to_string(),
1751                    killed: true,
1752                }),
1753                true,
1754                false,
1755            ),
1756            (
1757                "recovery",
1758                WorkflowRunStatus::Suspended,
1759                Some(WorkflowSuspensionContext::Recovery {
1760                    reason: "process restarted".to_string(),
1761                }),
1762                true,
1763                true,
1764            ),
1765            (
1766                "succeeded",
1767                WorkflowRunStatus::Succeeded,
1768                None,
1769                false,
1770                false,
1771            ),
1772            ("failed", WorkflowRunStatus::Failed, None, false, false),
1773            ("cancelled", WorkflowRunStatus::Cancelled, None, true, false),
1774        ];
1775
1776        for (name, status, suspension, can_cancel, can_restart_as_new_run) in cases {
1777            let cancel_run_id = format!("action-cancel-{name}");
1778            let cancel_snapshot = seed_durable_workflow_run(
1779                &access,
1780                directory.path(),
1781                &workspace,
1782                session_id,
1783                &cancel_run_id,
1784                status,
1785                suspension.clone(),
1786            )
1787            .await;
1788            let cancel_public = public_workflow_snapshot(cancel_snapshot);
1789            assert_eq!(
1790                cancel_public.can_cancel, can_cancel,
1791                "cancel projection mismatch for {name}"
1792            );
1793            assert_eq!(
1794                access
1795                    .cancel_for_session(session_id, &cancel_run_id)
1796                    .await
1797                    .is_ok(),
1798                can_cancel,
1799                "cancel endpoint mismatch for {name}"
1800            );
1801
1802            let restart_run_id = format!("action-restart-{name}");
1803            let restart_snapshot = seed_durable_workflow_run(
1804                &access,
1805                directory.path(),
1806                &workspace,
1807                session_id,
1808                &restart_run_id,
1809                status,
1810                suspension,
1811            )
1812            .await;
1813            let restart_public = public_workflow_snapshot(restart_snapshot);
1814            assert_eq!(
1815                restart_public.can_restart_as_new_run, can_restart_as_new_run,
1816                "restart projection mismatch for {name}"
1817            );
1818            let restarted = access
1819                .restart_for_session(session_id, &restart_run_id)
1820                .await;
1821            assert_eq!(
1822                restarted.is_ok(),
1823                can_restart_as_new_run,
1824                "restart endpoint mismatch for {name}: {restarted:?}"
1825            );
1826            if let Ok(restarted) = restarted {
1827                let _ = access
1828                    .cancel_for_session(session_id, &restarted.run_id)
1829                    .await;
1830            }
1831        }
1832    }
1833
1834    #[tokio::test]
1835    async fn restart_as_new_run_is_indexed_isolated_and_survives_reconstruction() {
1836        let (access, repo, directory, storage) =
1837            workflow_test_access_with_tools_and_storage(Arc::new(WorkflowReadTool)).await;
1838        let workspace = directory.path().join("restart-workspace");
1839        std::fs::create_dir_all(&workspace).expect("restart workspace");
1840        let session_id = "restart-owner";
1841        let mut session = Session::new(session_id, "model");
1842        session.workspace = Some(workspace.to_string_lossy().into_owned());
1843        repo.save(&mut session).await.expect("save restart owner");
1844        let mut other = Session::new("restart-other", "model");
1845        other.workspace = Some(workspace.to_string_lossy().into_owned());
1846        repo.save(&mut other).await.expect("save other session");
1847
1848        let original = seed_durable_workflow_run(
1849            &access,
1850            directory.path(),
1851            &workspace,
1852            session_id,
1853            "recovery-original",
1854            WorkflowRunStatus::Suspended,
1855            Some(WorkflowSuspensionContext::Recovery {
1856                reason: "process restarted".to_string(),
1857            }),
1858        )
1859        .await;
1860        replace_workflow_run_index(&repo, session_id, vec![original.run_id.clone()]).await;
1861
1862        let restarted = access
1863            .restart_for_session(session_id, &original.run_id)
1864            .await
1865            .expect("restart recovery suspension as a new run");
1866        assert_ne!(restarted.run_id, original.run_id);
1867        assert_eq!(restarted.session_id, session_id);
1868        assert_eq!(
1869            access
1870                .progress_for_session(session_id, &original.run_id, u64::MAX)
1871                .await
1872                .expect("original progress")
1873                .snapshot,
1874            original,
1875            "restart-as-new must not mutate the original suspended run"
1876        );
1877
1878        let immediate = access
1879            .list_for_session(session_id)
1880            .await
1881            .expect("immediate owner list");
1882        let immediate_ids = immediate
1883            .iter()
1884            .map(|snapshot| snapshot.run_id.as_str())
1885            .collect::<BTreeSet<_>>();
1886        assert_eq!(
1887            immediate_ids,
1888            BTreeSet::from([original.run_id.as_str(), restarted.run_id.as_str()])
1889        );
1890        assert!(access
1891            .list_for_session("restart-other")
1892            .await
1893            .expect("isolated other list")
1894            .is_empty());
1895        assert!(matches!(
1896            access
1897                .progress_for_session("restart-other", &restarted.run_id, u64::MAX)
1898                .await,
1899            Err(WorkflowRunError::NotFound)
1900        ));
1901        assert!(matches!(
1902            access
1903                .restart_for_session("restart-other", &original.run_id)
1904                .await,
1905            Err(WorkflowRunError::NotFound)
1906        ));
1907
1908        let settled = wait_for_workflow_run_to_settle(&access, &restarted.run_id).await;
1909        assert_eq!(settled.status, WorkflowRunStatus::Succeeded);
1910        let skills = access.skills.clone();
1911        drop(access);
1912        drop(repo);
1913
1914        let storage_port: Arc<dyn Storage> = storage;
1915        let reopened_repo = bamboo_engine::SessionRepository::new(
1916            Arc::default(),
1917            storage_port.clone(),
1918            Arc::new(bamboo_storage::LockedSessionStore::new(storage_port)),
1919        );
1920        let reopened = WorkflowRunAccess::new(
1921            directory.path(),
1922            Arc::new(WorkflowReadTool),
1923            skills,
1924            reopened_repo,
1925        )
1926        .await
1927        .expect("reconstruct workflow access");
1928        let reconstructed = reopened
1929            .list_for_session(session_id)
1930            .await
1931            .expect("list after reconstruction");
1932        let reconstructed_ids = reconstructed
1933            .iter()
1934            .map(|snapshot| snapshot.run_id.as_str())
1935            .collect::<BTreeSet<_>>();
1936        assert_eq!(
1937            reconstructed_ids,
1938            BTreeSet::from([original.run_id.as_str(), restarted.run_id.as_str()])
1939        );
1940        assert_eq!(
1941            reopened
1942                .progress_for_session(session_id, &original.run_id, u64::MAX)
1943                .await
1944                .expect("reconstructed original")
1945                .snapshot,
1946            original
1947        );
1948    }
1949
1950    #[tokio::test]
1951    async fn restart_capacity_preflight_creates_no_new_run_or_active_orphan() {
1952        let (access, repo, directory) = workflow_test_access().await;
1953        let workspace = directory.path().join("restart-capacity-workspace");
1954        std::fs::create_dir_all(&workspace).expect("restart capacity workspace");
1955        let session_id = "restart-capacity";
1956        let mut session = Session::new(session_id, "model");
1957        session.workspace = Some(workspace.to_string_lossy().into_owned());
1958        repo.save(&mut session)
1959            .await
1960            .expect("save restart capacity session");
1961
1962        let original = seed_durable_workflow_run(
1963            &access,
1964            directory.path(),
1965            &workspace,
1966            session_id,
1967            "capacity-run-0",
1968            WorkflowRunStatus::Suspended,
1969            Some(WorkflowSuspensionContext::Recovery {
1970                reason: "process restarted".to_string(),
1971            }),
1972        )
1973        .await;
1974        let mut run_ids = vec![original.run_id.clone()];
1975        for index in 1..MAX_WORKFLOW_RUN_IDS_PER_SESSION {
1976            let mut snapshot = original.clone();
1977            snapshot.run_id = format!("capacity-run-{index}");
1978            snapshot.created_at = chrono::Utc::now();
1979            snapshot.updated_at = snapshot.created_at;
1980            persist_durable_workflow_run(directory.path(), &snapshot).await;
1981            run_ids.push(snapshot.run_id);
1982        }
1983        replace_workflow_run_index(&repo, session_id, run_ids).await;
1984        let before = access
1985            .engine
1986            .list_run_ids()
1987            .await
1988            .expect("run ids before capacity rejection")
1989            .into_iter()
1990            .collect::<BTreeSet<_>>();
1991
1992        assert!(matches!(
1993            access
1994                .restart_for_session(session_id, &original.run_id)
1995                .await,
1996            Err(WorkflowRunError::Preflight(message))
1997                if message == "workflow run index is full of active runs"
1998        ));
1999        let after = access
2000            .engine
2001            .list_run_ids()
2002            .await
2003            .expect("run ids after capacity rejection")
2004            .into_iter()
2005            .collect::<BTreeSet<_>>();
2006        assert_eq!(
2007            after, before,
2008            "capacity rejection must precede run creation"
2009        );
2010        assert!(after
2011            .iter()
2012            .all(|run_id| !access.engine.is_run_active(run_id)));
2013        assert_eq!(
2014            access
2015                .progress_for_session(session_id, &original.run_id, u64::MAX)
2016                .await
2017                .expect("original after capacity rejection")
2018                .snapshot,
2019            original
2020        );
2021    }
2022
2023    #[tokio::test]
2024    async fn restart_index_persistence_failure_cancels_the_unindexed_new_run() {
2025        let (access, repo, directory, storage) =
2026            workflow_test_access_with_tools_and_storage(Arc::new(BlockingWorkflowReadTool {
2027                entered: Arc::new(Semaphore::new(0)),
2028            }))
2029            .await;
2030        let workspace = directory.path().join("restart-failure-workspace");
2031        std::fs::create_dir_all(&workspace).expect("restart failure workspace");
2032        let session_id = "restart-index-failure";
2033        let mut session = Session::new(session_id, "model");
2034        session.workspace = Some(workspace.to_string_lossy().into_owned());
2035        repo.save(&mut session)
2036            .await
2037            .expect("save restart failure session");
2038        let original = seed_durable_workflow_run(
2039            &access,
2040            directory.path(),
2041            &workspace,
2042            session_id,
2043            "failure-original",
2044            WorkflowRunStatus::Suspended,
2045            Some(WorkflowSuspensionContext::Recovery {
2046                reason: "process restarted".to_string(),
2047            }),
2048        )
2049        .await;
2050        replace_workflow_run_index(&repo, session_id, vec![original.run_id.clone()]).await;
2051        let before = access
2052            .engine
2053            .list_run_ids()
2054            .await
2055            .expect("run ids before injected failure")
2056            .into_iter()
2057            .collect::<BTreeSet<_>>();
2058
2059        storage.set_fail_saves(true);
2060        let error = access
2061            .restart_for_session(session_id, &original.run_id)
2062            .await
2063            .expect_err("injected index persistence failure");
2064        storage.set_fail_saves(false);
2065        assert!(matches!(error, WorkflowRunError::Storage(_)));
2066
2067        let after = access
2068            .engine
2069            .list_run_ids()
2070            .await
2071            .expect("run ids after injected failure")
2072            .into_iter()
2073            .collect::<BTreeSet<_>>();
2074        let created = after.difference(&before).cloned().collect::<Vec<_>>();
2075        assert_eq!(created.len(), 1, "restart should have created one new run");
2076        let unindexed_run_id = &created[0];
2077        let compensated = wait_for_workflow_run_to_settle(&access, unindexed_run_id).await;
2078        assert_eq!(compensated.status, WorkflowRunStatus::Cancelled);
2079        assert!(!access.engine.is_run_active(unindexed_run_id));
2080        assert_eq!(
2081            access
2082                .progress_for_session(session_id, &original.run_id, u64::MAX)
2083                .await
2084                .expect("original after index failure")
2085                .snapshot,
2086            original
2087        );
2088        let listed = access
2089            .list_for_session(session_id)
2090            .await
2091            .expect("owner list after index failure");
2092        assert_eq!(listed.len(), 1);
2093        assert_eq!(listed[0].run_id, original.run_id);
2094        let durable = storage
2095            .load_session(session_id)
2096            .await
2097            .expect("load durable owner")
2098            .expect("durable owner");
2099        let durable_ids = durable
2100            .metadata
2101            .get(bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY)
2102            .and_then(|raw| serde_json::from_str::<Vec<String>>(raw).ok())
2103            .expect("durable run index");
2104        assert_eq!(durable_ids, vec![original.run_id]);
2105        assert!(!durable_ids.contains(unindexed_run_id));
2106    }
2107
2108    #[test]
2109    fn public_workflow_events_preserve_sequence_and_drop_private_payloads() {
2110        let private = "PRIVATE-EVENT-SENTINEL";
2111        let failure = WorkflowFailure {
2112            code: WorkflowFailureCode::ExecutionFailed,
2113            message: format!("/private/workspace/{private}"),
2114            retryable: true,
2115        };
2116        let kinds = vec![
2117            WorkflowRunEventKind::RunQueued,
2118            WorkflowRunEventKind::RunStarted,
2119            WorkflowRunEventKind::Phase {
2120                name: private.to_string(),
2121            },
2122            WorkflowRunEventKind::StepQueued,
2123            WorkflowRunEventKind::StepStarted,
2124            WorkflowRunEventKind::StepSuspended {
2125                reason: private.to_string(),
2126            },
2127            WorkflowRunEventKind::StepCompleted {
2128                output: json!({"raw": private}),
2129            },
2130            WorkflowRunEventKind::StepFailed {
2131                failure: failure.clone(),
2132            },
2133            WorkflowRunEventKind::StepCancelled,
2134            WorkflowRunEventKind::StepSkipped {
2135                reason: private.to_string(),
2136            },
2137            WorkflowRunEventKind::RunSuspended {
2138                reason: private.to_string(),
2139            },
2140            WorkflowRunEventKind::RunSucceeded {
2141                output: json!({"raw": private}),
2142            },
2143            WorkflowRunEventKind::RunFailed { failure },
2144            WorkflowRunEventKind::RunCancelled,
2145        ];
2146        let at = chrono::Utc::now();
2147        let events = kinds
2148            .into_iter()
2149            .enumerate()
2150            .map(|(index, kind)| {
2151                public_workflow_event(WorkflowRunEvent {
2152                    run_id: "public-run".to_string(),
2153                    sequence: index as u64 + 1,
2154                    at,
2155                    step_id: Some("inspect".to_string()),
2156                    kind,
2157                })
2158            })
2159            .collect::<Vec<_>>();
2160        let public = serde_json::to_value(&events).expect("public events serialize");
2161        let text = public.to_string();
2162
2163        assert!(!text.contains(private), "public events leaked: {text}");
2164        assert_eq!(public[2]["name"], "workflow_progressed");
2165        assert_eq!(public[6]["type"], "step_completed");
2166        assert!(public[6].get("output").is_none());
2167        assert_eq!(public[7]["failure"]["message"], "Workflow execution failed");
2168        assert_eq!(public[11]["type"], "run_succeeded");
2169        assert!(public[11].get("output").is_none());
2170        assert_eq!(
2171            events
2172                .iter()
2173                .map(|event| event.sequence)
2174                .collect::<Vec<_>>(),
2175            (1..=14).collect::<Vec<_>>()
2176        );
2177    }
2178
2179    #[test]
2180    fn workflow_tool_errors_never_expose_backend_diagnostics() {
2181        let sentinel = "/private/workspace/credentials-PRIVATE-SENTINEL";
2182        let cases = [
2183            WorkflowRunError::Storage(sentinel.to_string()),
2184            WorkflowRunError::InvalidInput(sentinel.to_string()),
2185            WorkflowRunError::Preflight(sentinel.to_string()),
2186        ];
2187        for error in cases {
2188            assert!(!workflow_tool_error(error).to_string().contains(sentinel));
2189        }
2190        assert!(!workflow_tool_error(WorkflowRunError::Compile(
2191            bamboo_domain::WorkflowCompileError::InvalidSchema(sentinel.to_string())
2192        ))
2193        .to_string()
2194        .contains(sentinel));
2195    }
2196
2197    fn canonical_workflow_run_schema() -> Value {
2198        json!({
2199            "type": "object",
2200            "properties": {
2201                "action": {
2202                    "type": "string",
2203                    "enum": ["start", "list", "get", "events", "cancel", "restart"]
2204                },
2205                "workflow_id": {"type": "string", "minLength": 1},
2206                "revision": {"type": "integer", "minimum": 1},
2207                "args": {"type": "object", "default": {}},
2208                "budget": {
2209                    "type": "object",
2210                    "properties": {
2211                        "max_concurrency": {"type": "integer", "minimum": 1},
2212                        "max_agents": {"type": "integer", "minimum": 0},
2213                        "max_steps": {"type": "integer", "minimum": 1},
2214                        "max_retries": {"type": "integer", "minimum": 0},
2215                        "max_nesting_depth": {"type": "integer", "minimum": 1},
2216                        "wall_time_ms": {"type": "integer", "minimum": 1},
2217                        "max_tokens": {"type": "integer", "minimum": 0},
2218                        "max_cost_micros": {"type": "integer", "minimum": 0}
2219                    },
2220                    "required": [
2221                        "max_concurrency",
2222                        "max_agents",
2223                        "max_steps",
2224                        "max_retries",
2225                        "max_nesting_depth",
2226                        "wall_time_ms"
2227                    ],
2228                    "additionalProperties": false
2229                },
2230                "run_id": {"type": "string"},
2231                "since": {"type": "integer", "minimum": 0}
2232            },
2233            "required": ["action"],
2234            "additionalProperties": false
2235        })
2236    }
2237
2238    #[tokio::test]
2239    async fn workflow_run_schema_is_flat_complete_and_canonical() {
2240        let (access, _, _) = workflow_test_access().await;
2241        let schema = WorkflowRunTool::new(access).parameters_schema();
2242
2243        for combinator in ["oneOf", "anyOf", "allOf"] {
2244            assert!(
2245                schema.get(combinator).is_none(),
2246                "workflow_run must not advertise root {combinator}"
2247            );
2248        }
2249        assert_eq!(schema, canonical_workflow_run_schema());
2250    }
2251
2252    #[tokio::test]
2253    async fn workflow_run_schema_survives_openai_sanitization_with_all_properties() {
2254        let (access, _, _) = workflow_test_access().await;
2255        let schema = WorkflowRunTool::new(access).parameters_schema();
2256        let sanitized =
2257            bamboo_llm::providers::common::tool_schema::sanitize_openai_function_parameters_schema(
2258                &schema,
2259            );
2260
2261        let properties = sanitized["properties"]
2262            .as_object()
2263            .expect("sanitized workflow_run properties");
2264        assert!(!properties.is_empty());
2265        assert_eq!(properties.len(), 7);
2266        assert_eq!(sanitized, canonical_workflow_run_schema());
2267    }
2268
2269    #[tokio::test]
2270    async fn workflow_run_schema_reaches_gemini_unchanged() {
2271        let (access, _, _) = workflow_test_access().await;
2272        let direct = WorkflowRunTool::new(access).to_schema();
2273        let gemini: GeminiTool = direct.to_provider().expect("Gemini tool conversion");
2274        let declaration = gemini
2275            .function_declarations
2276            .first()
2277            .expect("workflow_run declaration");
2278
2279        assert_eq!(declaration.name, "workflow_run");
2280        assert_eq!(
2281            declaration.parameters_json_schema.as_ref(),
2282            Some(&canonical_workflow_run_schema())
2283        );
2284        assert!(declaration.parameters.is_none());
2285    }
2286
2287    #[tokio::test]
2288    async fn workflow_run_enforces_opt_in_tightens_budget_lists_and_isolates_sessions() {
2289        let (access, repo, directory) = workflow_test_access().await;
2290        let workspace = directory.path().join("workspace");
2291        std::fs::create_dir_all(&workspace).expect("workspace");
2292        let mut session = Session::new("workflow-session", "model");
2293        session.workspace = Some(workspace.to_string_lossy().into_owned());
2294        repo.save(&mut session).await.expect("save session");
2295
2296        let denied = access
2297            .start_from_tool("workflow-session", "review-flow", 42, json!({}), None)
2298            .await
2299            .expect_err("model start defaults off without session opt-in");
2300        assert!(denied.to_string().contains("opt-in"));
2301
2302        session.metadata.insert(
2303            bamboo_skills::WORKFLOW_ORCHESTRATION_OPT_IN_METADATA_KEY.to_string(),
2304            "true".to_string(),
2305        );
2306        repo.save(&mut session).await.expect("save opt-in");
2307        let requested = WorkflowBudgets {
2308            max_concurrency: 1,
2309            max_agents: 0,
2310            max_steps: 2,
2311            max_retries: 0,
2312            max_nesting_depth: 1,
2313            wall_time_ms: 5_000,
2314            max_tokens: Some(500),
2315            max_cost_micros: Some(500),
2316        };
2317        let started = access
2318            .start_from_tool(
2319                "workflow-session",
2320                "review-flow",
2321                42,
2322                json!({}),
2323                Some(requested.clone()),
2324            )
2325            .await
2326            .expect("opted-in model start");
2327        assert_eq!(started.definition.budgets, requested);
2328        let listed = access
2329            .list_for_session("workflow-session")
2330            .await
2331            .expect("session run list");
2332        assert_eq!(listed.len(), 1);
2333        assert_eq!(listed[0].run_id, started.run_id);
2334        let progress = access
2335            .progress_for_session("workflow-session", &started.run_id, 0)
2336            .await
2337            .expect("run events");
2338        assert!(progress
2339            .events
2340            .first()
2341            .is_some_and(|event| event.kind == bamboo_domain::WorkflowRunEventKind::RunQueued));
2342        let completed = tokio::time::timeout(std::time::Duration::from_secs(2), async {
2343            loop {
2344                let progress = access
2345                    .progress_for_session("workflow-session", &started.run_id, 0)
2346                    .await
2347                    .expect("terminal run progress");
2348                if progress.snapshot.status.is_terminal() {
2349                    break progress;
2350                }
2351                tokio::time::sleep(std::time::Duration::from_millis(10)).await;
2352            }
2353        })
2354        .await
2355        .expect("workflow reaches terminal state");
2356        assert_eq!(
2357            completed.snapshot.status,
2358            bamboo_domain::WorkflowRunStatus::Succeeded
2359        );
2360        assert_eq!(
2361            completed
2362                .events
2363                .iter()
2364                .map(|event| event.sequence)
2365                .collect::<Vec<_>>(),
2366            (1..=7).collect::<Vec<_>>()
2367        );
2368        assert!(matches!(
2369            completed.events.as_slice(),
2370            [
2371                bamboo_domain::WorkflowRunEvent {
2372                    kind: bamboo_domain::WorkflowRunEventKind::RunQueued,
2373                    ..
2374                },
2375                bamboo_domain::WorkflowRunEvent {
2376                    kind: bamboo_domain::WorkflowRunEventKind::RunStarted,
2377                    ..
2378                },
2379                bamboo_domain::WorkflowRunEvent {
2380                    kind: bamboo_domain::WorkflowRunEventKind::Phase { ref name },
2381                    ..
2382                },
2383                bamboo_domain::WorkflowRunEvent {
2384                    kind: bamboo_domain::WorkflowRunEventKind::StepQueued,
2385                    ..
2386                },
2387                bamboo_domain::WorkflowRunEvent {
2388                    kind: bamboo_domain::WorkflowRunEventKind::StepStarted,
2389                    ..
2390                },
2391                bamboo_domain::WorkflowRunEvent {
2392                    kind: bamboo_domain::WorkflowRunEventKind::StepCompleted { .. },
2393                    ..
2394                },
2395                bamboo_domain::WorkflowRunEvent {
2396                    kind: bamboo_domain::WorkflowRunEventKind::RunSucceeded { .. },
2397                    ..
2398                }
2399            ] if name == "step_reserved"
2400        ));
2401        assert_eq!(
2402            completed.snapshot.last_sequence,
2403            completed.events.last().expect("terminal event").sequence
2404        );
2405
2406        let invalid_budget = WorkflowBudgets {
2407            max_steps: 0,
2408            ..requested.clone()
2409        };
2410        assert!(matches!(
2411            access
2412                .start_from_tool(
2413                    "workflow-session",
2414                    "review-flow",
2415                    42,
2416                    json!({}),
2417                    Some(invalid_budget),
2418                )
2419                .await,
2420            Err(WorkflowRunError::InvalidInput(_))
2421        ));
2422
2423        // Restart authority is the original immutable publication, not the
2424        // current catalog policy.
2425        let workflow_path = directory.path().join("skills/review-flow/workflow.yaml");
2426        let original_workflow = std::fs::read_to_string(&workflow_path).expect("workflow yaml");
2427        std::fs::write(
2428            &workflow_path,
2429            original_workflow.replace(
2430                "invocation_policy: {explicit: true, automatic: true}",
2431                "invocation_policy: {explicit: true, automatic: false}",
2432            ),
2433        )
2434        .expect("disable automatic live policy");
2435        access.skills.store().reload().await.expect("reload policy");
2436        let restart = access
2437            .restart_from_tool("workflow-session", &started.run_id)
2438            .await
2439            .expect_err("succeeded runs are terminal");
2440        assert!(matches!(restart, WorkflowRunError::Terminal));
2441
2442        let mut other = Session::new("other-session", "model");
2443        other.workspace = Some(workspace.to_string_lossy().into_owned());
2444        repo.save(&mut other).await.expect("save other session");
2445        assert!(access
2446            .list_for_session("other-session")
2447            .await
2448            .expect("isolated list")
2449            .is_empty());
2450        assert!(matches!(
2451            access
2452                .progress_for_session("other-session", &started.run_id, 0)
2453                .await,
2454            Err(WorkflowRunError::NotFound)
2455        ));
2456
2457        // Recreate the adapter over the same durable journal, then combine the
2458        // snapshot with only the tail after sequence 4. This is the reconnect
2459        // contract used by Lotus after a process or transport restart.
2460        let run_id = started.run_id.clone();
2461        let skills = access.skills.clone();
2462        drop(access);
2463        let reopened = WorkflowRunAccess::new(
2464            directory.path(),
2465            Arc::new(WorkflowReadTool),
2466            skills,
2467            repo.clone(),
2468        )
2469        .await
2470        .expect("reopen workflow adapter");
2471        let reconnected = reopened
2472            .progress_for_session("workflow-session", &run_id, 4)
2473            .await
2474            .expect("durable reconnect");
2475        assert_eq!(reconnected.snapshot.status, WorkflowRunStatus::Succeeded);
2476        assert_eq!(
2477            reconnected
2478                .events
2479                .iter()
2480                .map(|event| event.sequence)
2481                .collect::<Vec<_>>(),
2482            vec![5, 6, 7]
2483        );
2484        assert_eq!(
2485            public_workflow_snapshot(reconnected.snapshot).last_sequence,
2486            7
2487        );
2488        assert!(matches!(
2489            public_workflow_event(reconnected.events.last().cloned().expect("tail event")).kind,
2490            PublicWorkflowRunEventKind::RunSucceeded
2491        ));
2492        assert_eq!(
2493            reopened
2494                .list_for_session("workflow-session")
2495                .await
2496                .expect("reconnected list")
2497                .len(),
2498            1
2499        );
2500    }
2501
2502    #[tokio::test]
2503    async fn workflow_cancel_is_idempotent_and_reconnects_to_one_terminal_event() {
2504        let entered = Arc::new(Semaphore::new(0));
2505        let (access, repo, directory) =
2506            workflow_test_access_with_tools(Arc::new(BlockingWorkflowReadTool {
2507                entered: entered.clone(),
2508            }))
2509            .await;
2510        let workspace = directory.path().join("workspace");
2511        std::fs::create_dir_all(&workspace).expect("workspace");
2512        let mut session = Session::new("cancel-session", "model");
2513        session.workspace = Some(workspace.to_string_lossy().into_owned());
2514        repo.save(&mut session).await.expect("save session");
2515
2516        let started = access
2517            .start("cancel-session", "review-flow", 42, json!({}), None)
2518            .await
2519            .expect("start blocking run");
2520        let _entered = tokio::time::timeout(std::time::Duration::from_secs(2), entered.acquire())
2521            .await
2522            .expect("blocking step entered")
2523            .expect("semaphore open");
2524
2525        let first = access
2526            .cancel_for_session("cancel-session", &started.run_id)
2527            .await
2528            .expect("first cancel");
2529        let second = access
2530            .cancel_for_session("cancel-session", &started.run_id)
2531            .await
2532            .expect("idempotent cancel");
2533        assert_eq!(first.status, WorkflowRunStatus::Cancelled);
2534        assert_eq!(second.status, WorkflowRunStatus::Cancelled);
2535        assert_eq!(first.last_sequence, second.last_sequence);
2536
2537        let progress = access
2538            .progress_for_session("cancel-session", &started.run_id, 0)
2539            .await
2540            .expect("cancelled journal");
2541        assert_eq!(progress.snapshot.status, WorkflowRunStatus::Cancelled);
2542        assert_eq!(progress.snapshot.last_sequence, first.last_sequence);
2543        assert_eq!(
2544            progress
2545                .events
2546                .iter()
2547                .filter(|event| matches!(event.kind, WorkflowRunEventKind::RunCancelled))
2548                .count(),
2549            1
2550        );
2551        assert!(!progress
2552            .events
2553            .iter()
2554            .any(|event| matches!(event.kind, WorkflowRunEventKind::RunSucceeded { .. })));
2555        assert_eq!(
2556            progress
2557                .events
2558                .iter()
2559                .map(|event| event.sequence)
2560                .collect::<Vec<_>>(),
2561            (1..=progress.snapshot.last_sequence).collect::<Vec<_>>()
2562        );
2563        let tail = access
2564            .progress_for_session(
2565                "cancel-session",
2566                &started.run_id,
2567                progress.snapshot.last_sequence,
2568            )
2569            .await
2570            .expect("terminal reconnect tail");
2571        assert!(tail.events.is_empty());
2572        assert!(matches!(
2573            public_workflow_snapshot(second).status,
2574            WorkflowRunStatus::Cancelled
2575        ));
2576    }
2577
2578    #[test]
2579    fn tool_input_rejects_security_context_spoofing() {
2580        let error = serde_json::from_value::<WorkflowToolInput>(json!({
2581            "action": "start",
2582            "workflow_id": "safe",
2583            "revision": 1,
2584            "workspace_trusted": true
2585        }))
2586        .unwrap_err();
2587        assert!(error.to_string().contains("unknown field"));
2588    }
2589
2590    #[test]
2591    fn tool_input_rejects_fields_from_other_actions() {
2592        let invalid = [
2593            json!({"action": "list", "run_id": "run-1"}),
2594            json!({"action": "get", "run_id": "run-1", "since": 1}),
2595            json!({"action": "events", "run_id": "run-1", "workflow_id": "flow"}),
2596            json!({"action": "cancel", "run_id": "run-1", "revision": 1}),
2597            json!({"action": "restart", "run_id": "run-1", "budget": {}}),
2598            json!({
2599                "action": "start",
2600                "workflow_id": "flow",
2601                "revision": 1,
2602                "run_id": "run-1"
2603            }),
2604        ];
2605
2606        for input in invalid {
2607            let error = serde_json::from_value::<WorkflowToolInput>(input.clone())
2608                .expect_err("action-specific fields must remain authoritative at runtime");
2609            assert!(
2610                error.to_string().contains("unknown field"),
2611                "unexpected error for {input}: {error}"
2612            );
2613        }
2614
2615        assert!(matches!(
2616            serde_json::from_value::<WorkflowToolInput>(json!({"action": "list"}))
2617                .expect("fieldless list action"),
2618            WorkflowToolInput::List {}
2619        ));
2620    }
2621
2622    #[tokio::test]
2623    async fn omitted_start_args_and_zero_budgets_match_schema() {
2624        let WorkflowToolInput::Start { args, .. } =
2625            serde_json::from_value::<WorkflowToolInput>(json!({
2626                "action": "start",
2627                "workflow_id": "safe",
2628                "revision": 1
2629            }))
2630            .expect("tool args default")
2631        else {
2632            panic!("start input")
2633        };
2634        assert_eq!(args, json!({}));
2635        let http: crate::handlers::workflow_runs::StartWorkflowRunRequest =
2636            serde_json::from_value(json!({"workflow_id":"safe", "revision":1}))
2637                .expect("http args default");
2638        assert_eq!(http.args, json!({}));
2639
2640        let (access, _, _) = workflow_test_access().await;
2641        let schema = WorkflowRunTool { access }.parameters_schema();
2642        let properties = &schema["properties"];
2643        assert_eq!(properties["args"]["default"], json!({}));
2644        assert_eq!(
2645            properties["budget"]["properties"]["max_agents"]["minimum"],
2646            0
2647        );
2648        assert_eq!(
2649            properties["budget"]["properties"]["max_retries"]["minimum"],
2650            0
2651        );
2652        assert_eq!(
2653            properties["budget"]["properties"]["max_tokens"]["minimum"],
2654            0
2655        );
2656        assert_eq!(
2657            properties["budget"]["properties"]["max_cost_micros"]["minimum"],
2658            0
2659        );
2660    }
2661
2662    #[tokio::test]
2663    async fn run_index_updates_are_concurrent_and_never_evict_at_active_capacity() {
2664        let (access, repo, _) = workflow_test_access().await;
2665        let mut session = Session::new("run-index", "model");
2666        repo.save(&mut session).await.expect("seed session");
2667        let results = futures::future::join_all((0..32).map(|index| {
2668            let access = access.clone();
2669            async move {
2670                let run_id = format!("run-{index}");
2671                access.remember_run_id("run-index", &run_id).await
2672            }
2673        }))
2674        .await;
2675        assert!(results.into_iter().all(|result| result.is_ok()));
2676        let concurrent = repo
2677            .try_load("run-index")
2678            .await
2679            .expect("load")
2680            .expect("session");
2681        let ids = serde_json::from_str::<Vec<String>>(
2682            concurrent
2683                .metadata
2684                .get(bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY)
2685                .expect("run ids"),
2686        )
2687        .expect("ids json");
2688        assert_eq!(ids.len(), 32);
2689        assert_eq!(ids.iter().collect::<BTreeSet<_>>().len(), 32);
2690
2691        let capacity_ids = (0..MAX_WORKFLOW_RUN_IDS_PER_SESSION)
2692            .map(|index| format!("active-{index}"))
2693            .collect::<Vec<_>>();
2694        repo.update_runtime_session(
2695            "run-index",
2696            &[bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY],
2697            {
2698                let capacity_ids = capacity_ids.clone();
2699                move |session| {
2700                    session.metadata.insert(
2701                        bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY.to_string(),
2702                        serde_json::to_string(&capacity_ids).expect("ids json"),
2703                    );
2704                }
2705            },
2706        )
2707        .await
2708        .expect("fill index")
2709        .expect("session");
2710        assert!(matches!(
2711            access.remember_run_id("run-index", "new-run").await,
2712            Err(WorkflowRunError::Storage(_))
2713        ));
2714        let retained = repo
2715            .try_load("run-index")
2716            .await
2717            .expect("load")
2718            .expect("session");
2719        let retained = serde_json::from_str::<Vec<String>>(
2720            retained
2721                .metadata
2722                .get(bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY)
2723                .expect("run ids"),
2724        )
2725        .expect("ids json");
2726        assert_eq!(
2727            retained, capacity_ids,
2728            "oldest active id must not be evicted"
2729        );
2730    }
2731
2732    #[tokio::test]
2733    async fn real_model_workflow_run_index_survives_tool_result_and_final_session_save() {
2734        use bamboo_agent_core::storage::AttachmentReader;
2735        use bamboo_engine::{Agent, ExecuteRequestBuilder};
2736        use bamboo_llm::{LLMChunk, LLMProvider, LLMStream};
2737        use futures::stream;
2738        use tokio::sync::Mutex;
2739        use tokio_util::sync::CancellationToken;
2740
2741        struct NoAttachments;
2742        #[async_trait]
2743        impl AttachmentReader for NoAttachments {
2744            async fn read_attachment(
2745                &self,
2746                _session_id: &str,
2747                _attachment_id: &str,
2748            ) -> std::io::Result<Option<(Vec<u8>, String)>> {
2749                Ok(None)
2750            }
2751        }
2752        struct QueueProvider {
2753            queue: Mutex<Vec<Vec<bamboo_llm::provider::Result<LLMChunk>>>>,
2754        }
2755        #[async_trait]
2756        impl LLMProvider for QueueProvider {
2757            async fn chat_stream(
2758                &self,
2759                _messages: &[bamboo_agent_core::Message],
2760                _tools: &[ToolSchema],
2761                _max_output_tokens: Option<u32>,
2762                _model: &str,
2763            ) -> bamboo_llm::provider::Result<LLMStream> {
2764                Ok(Box::pin(stream::iter(self.queue.lock().await.remove(0))))
2765            }
2766        }
2767
2768        let (access, repo, directory) = workflow_test_access().await;
2769        let session_id = "real-model-workflow-run";
2770        let mut session = Session::new(session_id, "test-model");
2771        session.metadata.insert(
2772            bamboo_skills::WORKFLOW_ORCHESTRATION_OPT_IN_METADATA_KEY.to_string(),
2773            "true".to_string(),
2774        );
2775        session
2776            .metadata
2777            .insert("external.metadata".to_string(), "preserve".to_string());
2778        session.add_message(bamboo_agent_core::Message::system("system"));
2779        session.add_message(bamboo_agent_core::Message::user("run review workflow"));
2780        repo.save(&mut session).await.expect("seed session");
2781        let call = ToolCall {
2782            id: "call-workflow-run".to_string(),
2783            tool_type: "function".to_string(),
2784            function: bamboo_agent_core::tools::FunctionCall {
2785                name: "workflow_run".to_string(),
2786                arguments: json!({
2787                    "action":"start",
2788                    "workflow_id":"review-flow",
2789                    "revision":42
2790                })
2791                .to_string(),
2792            },
2793        };
2794        let provider = Arc::new(QueueProvider {
2795            queue: Mutex::new(vec![
2796                vec![Ok(LLMChunk::ToolCalls(vec![call])), Ok(LLMChunk::Done)],
2797                vec![Ok(LLMChunk::Token("done".to_string())), Ok(LLMChunk::Done)],
2798            ]),
2799        });
2800        let tools = Arc::new(
2801            bamboo_tools::BuiltinToolExecutorBuilder::new()
2802                .with_tool(WorkflowRunTool::new(access.clone()))
2803                .expect("workflow tool")
2804                .build(),
2805        );
2806        let metrics = bamboo_metrics::MetricsCollector::spawn(
2807            Arc::new(bamboo_metrics::SqliteMetricsStorage::new(
2808                directory.path().join("runner-metrics.db"),
2809            )),
2810            7,
2811        );
2812        let agent = Agent::builder()
2813            .storage(repo.storage().clone())
2814            .persistence(Arc::new(repo.clone()))
2815            .attachment_reader(Arc::new(NoAttachments))
2816            .skill_manager(access.skills.clone())
2817            .metrics_collector(metrics)
2818            .config(Arc::new(RwLock::new(bamboo_llm::Config::default())))
2819            .provider(provider)
2820            .default_tools(tools)
2821            .build()
2822            .expect("agent");
2823        let (event_tx, _event_rx) = tokio::sync::mpsc::channel(128);
2824        agent
2825            .execute(
2826                &mut session,
2827                ExecuteRequestBuilder::new(
2828                    "run review workflow",
2829                    event_tx,
2830                    CancellationToken::new(),
2831                )
2832                .model("test-model")
2833                .build(),
2834            )
2835            .await
2836            .expect("real model workflow run");
2837
2838        let saved = repo
2839            .storage()
2840            .load_session(session_id)
2841            .await
2842            .expect("load")
2843            .expect("saved");
2844        assert_eq!(
2845            saved.metadata.get("external.metadata").map(String::as_str),
2846            Some("preserve")
2847        );
2848        assert!(saved
2849            .metadata
2850            .contains_key(bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY));
2851        let listed = access
2852            .list_for_session(session_id)
2853            .await
2854            .expect("list after final save");
2855        assert_eq!(listed.len(), 1);
2856        assert!(saved.messages.iter().any(|message| {
2857            message.tool_calls.as_ref().is_some_and(|calls| {
2858                calls
2859                    .iter()
2860                    .any(|call| call.function.name == "workflow_run")
2861            })
2862        }));
2863    }
2864
2865    #[tokio::test]
2866    async fn http_start_survives_concurrent_stale_runner_final_save_and_server_restart() {
2867        use bamboo_agent_core::storage::AttachmentReader;
2868        use bamboo_engine::{Agent, ExecuteRequestBuilder};
2869        use bamboo_llm::{LLMChunk, LLMProvider, LLMStream};
2870        use futures::stream;
2871        use tokio::sync::{oneshot, Mutex};
2872        use tokio_util::sync::CancellationToken;
2873
2874        struct NoAttachments;
2875        #[async_trait]
2876        impl AttachmentReader for NoAttachments {
2877            async fn read_attachment(
2878                &self,
2879                _session_id: &str,
2880                _attachment_id: &str,
2881            ) -> std::io::Result<Option<(Vec<u8>, String)>> {
2882                Ok(None)
2883            }
2884        }
2885
2886        struct PausingProvider {
2887            entered: Mutex<Option<oneshot::Sender<()>>>,
2888            resume: Mutex<Option<oneshot::Receiver<()>>>,
2889        }
2890        #[async_trait]
2891        impl LLMProvider for PausingProvider {
2892            async fn chat_stream(
2893                &self,
2894                _messages: &[bamboo_agent_core::Message],
2895                _tools: &[ToolSchema],
2896                _max_output_tokens: Option<u32>,
2897                _model: &str,
2898            ) -> bamboo_llm::provider::Result<LLMStream> {
2899                if let Some(entered) = self.entered.lock().await.take() {
2900                    let _ = entered.send(());
2901                }
2902                if let Some(resume) = self.resume.lock().await.take() {
2903                    let _ = resume.await;
2904                }
2905                Ok(Box::pin(stream::iter(vec![
2906                    Ok(LLMChunk::Token("done".to_string())),
2907                    Ok(LLMChunk::Done),
2908                ])))
2909            }
2910        }
2911
2912        let (access, repo, directory) = workflow_test_access().await;
2913        let session_id = "http-start-concurrent-runner-save";
2914        let mut session = Session::new(session_id, "test-model");
2915        session.add_message(bamboo_agent_core::Message::system("system"));
2916        session.add_message(bamboo_agent_core::Message::user("keep running"));
2917        repo.save(&mut session).await.expect("seed session");
2918        let mut runner_session = repo
2919            .try_load(session_id)
2920            .await
2921            .expect("load runner session")
2922            .expect("runner session");
2923
2924        let (entered_tx, entered_rx) = oneshot::channel();
2925        let (resume_tx, resume_rx) = oneshot::channel();
2926        let provider = Arc::new(PausingProvider {
2927            entered: Mutex::new(Some(entered_tx)),
2928            resume: Mutex::new(Some(resume_rx)),
2929        });
2930        let metrics = bamboo_metrics::MetricsCollector::spawn(
2931            Arc::new(bamboo_metrics::SqliteMetricsStorage::new(
2932                directory.path().join("http-runner-metrics.db"),
2933            )),
2934            7,
2935        );
2936        let agent = Agent::builder()
2937            .storage(repo.storage().clone())
2938            .persistence(Arc::new(repo.clone()))
2939            .attachment_reader(Arc::new(NoAttachments))
2940            .skill_manager(access.skills.clone())
2941            .metrics_collector(metrics)
2942            .config(Arc::new(RwLock::new(bamboo_llm::Config::default())))
2943            .provider(provider)
2944            .default_tools(Arc::new(
2945                bamboo_tools::BuiltinToolExecutorBuilder::new().build(),
2946            ))
2947            .build()
2948            .expect("agent");
2949        let (event_tx, _event_rx) = tokio::sync::mpsc::channel(128);
2950        let runner = tokio::spawn(async move {
2951            agent
2952                .execute(
2953                    &mut runner_session,
2954                    ExecuteRequestBuilder::new("keep running", event_tx, CancellationToken::new())
2955                        .model("test-model")
2956                        .build(),
2957                )
2958                .await
2959        });
2960        tokio::time::timeout(std::time::Duration::from_secs(2), entered_rx)
2961            .await
2962            .expect("runner enters model round")
2963            .expect("runner entry signal");
2964
2965        let started = access
2966            .start(session_id, "review-flow", 42, json!({}), None)
2967            .await
2968            .expect("HTTP-equivalent explicit start");
2969        let durable_during_round = repo
2970            .storage()
2971            .load_session(session_id)
2972            .await
2973            .expect("load during round")
2974            .expect("session during round");
2975        assert!(durable_during_round
2976            .metadata
2977            .get(bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY)
2978            .is_some_and(|raw| raw.contains(&started.run_id)));
2979
2980        resume_tx.send(()).expect("resume runner");
2981        tokio::time::timeout(std::time::Duration::from_secs(2), runner)
2982            .await
2983            .expect("runner completes")
2984            .expect("runner task")
2985            .expect("runner execution");
2986
2987        let durable_after_final_save = repo
2988            .storage()
2989            .load_session(session_id)
2990            .await
2991            .expect("load after final save")
2992            .expect("saved session");
2993        assert!(durable_after_final_save
2994            .metadata
2995            .get(bamboo_skills::WORKFLOW_RUN_IDS_METADATA_KEY)
2996            .is_some_and(|raw| raw.contains(&started.run_id)));
2997
2998        tokio::time::timeout(std::time::Duration::from_secs(2), async {
2999            loop {
3000                let progress = access
3001                    .progress_for_session(session_id, &started.run_id, u64::MAX)
3002                    .await
3003                    .expect("workflow progress before restart");
3004                if progress.snapshot.status.is_terminal()
3005                    && !access.engine.is_run_active(&started.run_id)
3006                {
3007                    break;
3008                }
3009                tokio::task::yield_now().await;
3010            }
3011        })
3012        .await
3013        .expect("workflow reaches terminal state before restart");
3014
3015        let skills = access.skills.clone();
3016        drop(access);
3017        let restarted =
3018            WorkflowRunAccess::new(directory.path(), Arc::new(WorkflowReadTool), skills, repo)
3019                .await
3020                .expect("restart workflow access");
3021        let listed = restarted
3022            .list_for_session(session_id)
3023            .await
3024            .expect("list after restart");
3025        assert_eq!(listed.len(), 1);
3026        assert_eq!(listed[0].run_id, started.run_id);
3027    }
3028
3029    #[tokio::test]
3030    async fn production_policy_allows_read_without_fabricating_workspace_trust() {
3031        let read = BTreeSet::from(["read".to_string()]);
3032        assert_eq!(
3033            ServerWorkflowPolicy
3034                .authorize(
3035                    "session",
3036                    &WorkflowPolicyTarget::Tool("read_file".to_string()),
3037                    &read,
3038                    false,
3039                )
3040                .await,
3041            PermissionDecision::Allow
3042        );
3043
3044        let write = BTreeSet::from(["write".to_string()]);
3045        assert!(matches!(
3046            ServerWorkflowPolicy
3047                .authorize(
3048                    "session",
3049                    &WorkflowPolicyTarget::Tool("write_file".to_string()),
3050                    &write,
3051                    false,
3052                )
3053                .await,
3054            PermissionDecision::Deny(_)
3055        ));
3056
3057        for hostile_target in [
3058            "Write",
3059            "write_file",
3060            "WebFetch",
3061            "mcp::remote_tool",
3062            "Bash",
3063        ] {
3064            for claimed in [BTreeSet::new(), read.clone()] {
3065                assert!(matches!(
3066                    ServerWorkflowPolicy
3067                        .authorize(
3068                            "session",
3069                            &WorkflowPolicyTarget::Tool(hostile_target.to_string()),
3070                            &claimed,
3071                            false,
3072                        )
3073                        .await,
3074                    PermissionDecision::Deny(_)
3075                ));
3076            }
3077        }
3078
3079        assert_eq!(
3080            ServerWorkflowPolicy
3081                .authorize(
3082                    "session",
3083                    &WorkflowPolicyTarget::Workflow {
3084                        id: "nested-review".to_string(),
3085                        revision: 1,
3086                    },
3087                    &BTreeSet::new(),
3088                    false,
3089                )
3090                .await,
3091            PermissionDecision::Allow
3092        );
3093    }
3094
3095    #[test]
3096    fn pinned_bundle_limits_reject_oversized_runs_before_engine_start() {
3097        assert!(enforce_pinned_bundle_limits(
3098            MAX_PINNED_DEFINITIONS_PER_RUN,
3099            MAX_PINNED_BUNDLE_BYTES_PER_RUN
3100        )
3101        .is_ok());
3102        assert!(enforce_pinned_bundle_limits(
3103            MAX_PINNED_DEFINITIONS_PER_RUN + 1,
3104            MAX_PINNED_BUNDLE_BYTES_PER_RUN
3105        )
3106        .is_err());
3107        assert!(enforce_pinned_bundle_limits(
3108            MAX_PINNED_DEFINITIONS_PER_RUN,
3109            MAX_PINNED_BUNDLE_BYTES_PER_RUN + 1
3110        )
3111        .is_err());
3112    }
3113}