Skip to main content

bamboo_server_tools/skill_runtime/
load_skill.rs

1use std::sync::Arc;
2use std::time::Duration;
3
4use async_trait::async_trait;
5use serde::Deserialize;
6use serde_json::json;
7use sha2::{Digest, Sha256};
8use tokio::sync::RwLock;
9
10use bamboo_llm::Config;
11use bamboo_skills::access_control;
12use bamboo_skills::runtime_metadata::{
13    LAST_LOADED_SKILL_ID_METADATA_KEY, LAST_LOADED_SKILL_SUMMARY_METADATA_KEY,
14    LOADED_SKILL_IDS_METADATA_KEY,
15};
16use bamboo_skills::SkillManager;
17
18use bamboo_agent_core::tools::{
19    FunctionCall, Tool, ToolCall, ToolCtx, ToolError, ToolExecutionContext,
20    ToolExecutionSessionFlags, ToolExecutor, ToolOutcome, ToolResult,
21};
22
23use super::{
24    skill_access_error_to_tool_error, validate_runtime_activation,
25    validate_runtime_activation_descriptor, SkillToolAccess,
26};
27
28pub(super) fn plan_allows_dynamic_provider(plan_read_only: bool, tool_name: &str) -> bool {
29    !plan_read_only || bamboo_tools::orchestrator::plan_mode_allows_tool(tool_name)
30}
31
32#[derive(Debug, Deserialize)]
33struct LoadSkillArgs {
34    skill_id: String,
35}
36
37const LOAD_SKILL_OWNED_METADATA_KEYS: &[&str] = &[
38    bamboo_skills::runtime_metadata::SKILL_RUNTIME_SELECTED_CATALOG_KEY,
39    bamboo_skills::runtime_metadata::SKILL_RUNTIME_PINNED_SNAPSHOT_KEY,
40    bamboo_skills::runtime_metadata::SKILL_RUNTIME_ACTIVATION_ERROR_KEY,
41    bamboo_skills::ACTIVE_WORKFLOW_METADATA_KEY,
42    bamboo_skills::ACTIVE_WORKFLOW_SNAPSHOT_METADATA_KEY,
43    bamboo_skills::WORKFLOW_ACTIVATION_EVENT_METADATA_KEY,
44    bamboo_skills::WORKFLOW_LAST_DYNAMIC_CONTEXT_METADATA_KEY,
45    bamboo_skills::WORKFLOW_CONTEXT_CACHE_METADATA_KEY,
46    LOADED_SKILL_IDS_METADATA_KEY,
47    LAST_LOADED_SKILL_ID_METADATA_KEY,
48    LAST_LOADED_SKILL_SUMMARY_METADATA_KEY,
49];
50
51pub struct LoadSkillTool {
52    access: SkillToolAccess,
53    context_tools: Option<Arc<dyn ToolExecutor>>,
54    dynamic_context_permission_config: Option<Arc<bamboo_tools::permission::PermissionConfig>>,
55}
56
57impl LoadSkillTool {
58    pub fn new(
59        skill_manager: Arc<SkillManager>,
60        config: Arc<RwLock<Config>>,
61        session_repo: bamboo_engine::SessionRepository,
62    ) -> Self {
63        Self {
64            access: SkillToolAccess::new(skill_manager, config, session_repo),
65            context_tools: None,
66            dynamic_context_permission_config: None,
67        }
68    }
69
70    pub fn with_project_store(mut self, project_store: Arc<bamboo_projects::ProjectStore>) -> Self {
71        self.access = self.access.with_project_store(project_store);
72        self
73    }
74
75    /// Register the permission-wrapped base tool surface used by declared
76    /// dynamic context providers. Keeping this explicit prevents arbitrary
77    /// shell preprocessing and avoids recursive load_skill/workflow_run calls.
78    pub fn with_fail_closed_context_registry(mut self, tools: Arc<dyn ToolExecutor>) -> Self {
79        self.context_tools = Some(tools);
80        self
81    }
82
83    /// Register the same permission-wrapped tool surface and typed policy used
84    /// by normal tool dispatch. A missing typed config deliberately leaves the
85    /// fail-closed registry above inactive.
86    pub fn with_permission_checked_context_registry(
87        mut self,
88        tools: Arc<dyn ToolExecutor>,
89        permission_config: Option<Arc<bamboo_tools::permission::PermissionConfig>>,
90    ) -> Self {
91        self.context_tools = Some(tools);
92        self.dynamic_context_permission_config = permission_config;
93        self
94    }
95
96    /// Test-only injection seam for provider behavior.
97    #[cfg(test)]
98    pub(crate) fn with_test_context_tools(mut self, tools: Arc<dyn ToolExecutor>) -> Self {
99        self.context_tools = Some(tools);
100        self.dynamic_context_permission_config =
101            Some(Arc::new(bamboo_tools::permission::PermissionConfig::new()));
102        self
103    }
104
105    async fn persist_owned_metadata(
106        &self,
107        session_id: &str,
108        source: &bamboo_agent_core::Session,
109        operation: &str,
110    ) -> Result<(), ToolError> {
111        let updates = LOAD_SKILL_OWNED_METADATA_KEYS
112            .iter()
113            .map(|key| ((*key).to_string(), source.metadata.get(*key).cloned()))
114            .collect::<Vec<_>>();
115        self.access
116            .session_repo
117            .update_runtime_session(session_id, LOAD_SKILL_OWNED_METADATA_KEYS, move |latest| {
118                for (key, value) in updates {
119                    if let Some(value) = value {
120                        latest.metadata.insert(key, value);
121                    } else {
122                        latest.metadata.remove(&key);
123                    }
124                }
125            })
126            .await
127            .map_err(|error| ToolError::Execution(format!("{operation}: {error}")))?
128            .ok_or_else(|| ToolError::Execution(format!("Session '{session_id}' not found")))?;
129        Ok(())
130    }
131
132    async fn resolve_dynamic_context(
133        &self,
134        skill: &bamboo_skills::SkillDefinition,
135        session: &mut bamboo_agent_core::Session,
136        ctx: &ToolCtx,
137    ) -> Result<Vec<bamboo_skills::DynamicContextBlock>, ToolError> {
138        const MAX_PROVIDERS: usize = 8;
139        const MAX_PROVIDER_CHARS: usize = 16_384;
140        const MAX_TOTAL_CHARS: usize = 32_768;
141        const MAX_TIMEOUT_MS: u64 = 15_000;
142        const MAX_CACHE_TTL_SECS: u64 = 3_600;
143        const MAX_CACHE_ENTRIES: usize = 64;
144        const MAX_CACHE_SERIALIZED_CHARS: usize = 131_072;
145        const REGISTERED_DYNAMIC_CONTEXT_TOOLS: &[&str] = &[
146            "Read",
147            "read_file",
148            "GetFileInfo",
149            "Glob",
150            "Grep",
151            "list_directory",
152        ];
153        let now = chrono::Utc::now();
154        let degraded_metadata = |message: &str| {
155            vec![bamboo_skills::DynamicContextBlock {
156                provider_id: "metadata".to_string(),
157                tool: "none".to_string(),
158                provenance: "invalid_workflow_metadata".to_string(),
159                generated_at: now,
160                expires_at: None,
161                status: bamboo_skills::WorkflowActivationStatus::Degraded,
162                stop_on_failure: true,
163                content: String::new(),
164                diagnostic: Some(bamboo_skills::WorkflowActivationDiagnostic {
165                    code: bamboo_skills::WorkflowActivationErrorCode::ProviderOutputInvalid,
166                    message: message.to_string(),
167                    recoverable: true,
168                }),
169            }]
170        };
171        let declarations = match skill
172            .metadata
173            .as_ref()
174            .and_then(|metadata| metadata.get("dynamic_context"))
175        {
176            Some(value) => match serde_json::from_value::<
177                Vec<bamboo_skills::DynamicContextDeclaration>,
178            >(value.clone())
179            {
180                Ok(value) => value,
181                Err(_) => {
182                    return Ok(degraded_metadata(
183                        "workflow dynamic context declaration metadata is invalid",
184                    ));
185                }
186            },
187            None => Vec::new(),
188        };
189        if declarations.is_empty() {
190            return Ok(Vec::new());
191        }
192        if declarations.len() > MAX_PROVIDERS {
193            return Ok(degraded_metadata(
194                "workflow declares too many dynamic context providers",
195            ));
196        }
197        let Some(permission_config) = self.dynamic_context_permission_config.as_ref() else {
198            return Ok(declarations
199                .into_iter()
200                .map(|declaration| bamboo_skills::DynamicContextBlock {
201                    provider_id: declaration.id,
202                    tool: declaration.tool,
203                    provenance: "typed_authority_unavailable".to_string(),
204                    generated_at: now,
205                    expires_at: None,
206                    status: bamboo_skills::WorkflowActivationStatus::Degraded,
207                    stop_on_failure: declaration.stop_on_failure,
208                    content: String::new(),
209                    diagnostic: Some(bamboo_skills::WorkflowActivationDiagnostic {
210                        code: bamboo_skills::WorkflowActivationErrorCode::ProviderFailed,
211                        message: "dynamic context provider authority is unavailable until the typed permission dependency is active".to_string(),
212                        recoverable: true,
213                    }),
214                })
215                .collect());
216        };
217        let Some(tools) = self.context_tools.as_ref() else {
218            return Ok(degraded_metadata(
219                "dynamic context provider registry is unavailable",
220            ));
221        };
222        let mut cache = match session
223            .metadata
224            .get(bamboo_skills::WORKFLOW_CONTEXT_CACHE_METADATA_KEY)
225        {
226            Some(raw) => match serde_json::from_str::<bamboo_skills::DynamicContextCache>(raw) {
227                Ok(cache) => cache,
228                Err(_) => {
229                    return Ok(degraded_metadata(
230                        "workflow dynamic context cache metadata is invalid",
231                    ));
232                }
233            },
234            None => bamboo_skills::DynamicContextCache::new(),
235        };
236        cache.retain(|_, block| block.expires_at.is_some_and(|expires| expires > now));
237        let mut blocks = Vec::new();
238        let mut total_chars = 0usize;
239        for declaration in declarations {
240            if declaration.id.trim().is_empty()
241                || declaration.tool.trim().is_empty()
242                || !REGISTERED_DYNAMIC_CONTEXT_TOOLS
243                    .iter()
244                    .any(|tool| tool.eq_ignore_ascii_case(&declaration.tool))
245            {
246                return Ok(degraded_metadata(
247                    "workflow dynamic context declaration is invalid",
248                ));
249            }
250            if bamboo_agent_core::tools::classify_tool(&declaration.tool)
251                != bamboo_agent_core::tools::ToolMutability::ReadOnly
252            {
253                return Ok(degraded_metadata(
254                    "dynamic context provider is not classified read-only",
255                ));
256            }
257            let provider_input = match confine_dynamic_provider_input(
258                session,
259                &declaration.tool,
260                &declaration.input,
261            )
262            .await
263            {
264                Ok(input) => input,
265                Err(diagnostic) => {
266                    blocks.push(bamboo_skills::DynamicContextBlock {
267                        provider_id: declaration.id,
268                        tool: declaration.tool,
269                        provenance: "workspace_confinement_denied".to_string(),
270                        generated_at: now,
271                        expires_at: None,
272                        status: bamboo_skills::WorkflowActivationStatus::Degraded,
273                        stop_on_failure: declaration.stop_on_failure,
274                        content: String::new(),
275                        diagnostic: Some(diagnostic),
276                    });
277                    if declaration.stop_on_failure {
278                        break;
279                    }
280                    continue;
281                }
282            };
283            let catalog_identity = session
284                .metadata
285                .get(bamboo_skills::runtime_metadata::SKILL_RUNTIME_SELECTED_CATALOG_KEY)
286                .and_then(|raw| {
287                    serde_json::from_str::<Vec<bamboo_skills::WorkflowCatalogEntry>>(raw).ok()
288                })
289                .and_then(|entries| entries.into_iter().find(|entry| entry.id == skill.id))
290                .map(|entry| format!("{:?}:{}", entry.source, entry.revision))
291                .unwrap_or_else(|| "unknown:0".to_string());
292            let workspace_scope = session.workspace_path_meta().unwrap_or_default();
293            let cache_material = json!({
294                "workflow": skill.id,
295                "catalog": catalog_identity,
296                "provider": declaration.id,
297                "tool": declaration.tool,
298                "input": provider_input,
299                "workspace": workspace_scope,
300                "permission_scope": "strict_no_bypass",
301                "permission_policy_revision": permission_config.policy_revision(),
302            });
303            let cache_key = hex::encode(Sha256::digest(cache_material.to_string().as_bytes()));
304            // Cached content is never injected without re-running the concrete
305            // tool permission gate. The current ToolExecutor surface has no
306            // authorization-only check or immutable policy revision, so a raw
307            // cache hit could bypass an allow -> deny/approval policy change.
308            let available = tools.list_tools();
309            let canonical_provider = bamboo_tools::exposure::canonical_tool_name(&declaration.tool);
310            if !available.iter().any(|schema| {
311                bamboo_tools::exposure::canonical_tool_name(&schema.function.name)
312                    == canonical_provider
313            }) {
314                let diagnostic = bamboo_skills::WorkflowActivationDiagnostic {
315                    code: bamboo_skills::WorkflowActivationErrorCode::ProviderFailed,
316                    message: format!(
317                        "dynamic context provider '{}' is not registered",
318                        declaration.id
319                    ),
320                    recoverable: true,
321                };
322                blocks.push(bamboo_skills::DynamicContextBlock {
323                    provider_id: declaration.id,
324                    tool: declaration.tool,
325                    provenance: "registered_tool".to_string(),
326                    generated_at: now,
327                    expires_at: None,
328                    status: bamboo_skills::WorkflowActivationStatus::Degraded,
329                    stop_on_failure: declaration.stop_on_failure,
330                    content: String::new(),
331                    diagnostic: Some(diagnostic),
332                });
333                if declaration.stop_on_failure {
334                    break;
335                }
336                continue;
337            }
338            let call_id = format!("dynamic-context-{}-{}", skill.id, declaration.id);
339            let call = ToolCall {
340                id: call_id.clone(),
341                tool_type: "function".to_string(),
342                function: FunctionCall {
343                    name: declaration.tool.clone(),
344                    arguments: provider_input.to_string(),
345                },
346            };
347            let (fallback_tx, _fallback_rx) = tokio::sync::mpsc::channel(1);
348            let event_tx = ctx.event_tx.as_ref().unwrap_or(&fallback_tx);
349            let execution_context = ToolExecutionContext::for_dispatch(
350                session.id.as_str(),
351                &call_id,
352                event_tx,
353                &available,
354                // Dynamic context never inherits legacy Bypass. Every provider
355                // passes the normal permission/workspace gate independently.
356                // Auto's zero-prompt contract and Plan's hard read-only overlay,
357                // however, must survive this nested dispatch.
358                ToolExecutionSessionFlags {
359                    bypass_permissions: false,
360                    auto_approve_permissions: ctx.auto_approve_permissions,
361                    plan_read_only: ctx.plan_read_only,
362                },
363                false,
364                None,
365                Some(&provider_input),
366            );
367            if !plan_allows_dynamic_provider(ctx.plan_read_only, &call.function.name) {
368                blocks.push(bamboo_skills::DynamicContextBlock {
369                    provider_id: declaration.id,
370                    tool: declaration.tool,
371                    provenance: "registered_tool_permission_checked".to_string(),
372                    generated_at: now,
373                    expires_at: None,
374                    status: bamboo_skills::WorkflowActivationStatus::Degraded,
375                    stop_on_failure: declaration.stop_on_failure,
376                    content: String::new(),
377                    diagnostic: Some(bamboo_skills::WorkflowActivationDiagnostic {
378                        code: bamboo_skills::WorkflowActivationErrorCode::ProviderFailed,
379                        message: "Plan mode blocked a mutating dynamic context provider"
380                            .to_string(),
381                        recoverable: true,
382                    }),
383                });
384                if declaration.stop_on_failure {
385                    break;
386                }
387                continue;
388            }
389            let timeout_ms = declaration.timeout_ms.clamp(1, MAX_TIMEOUT_MS);
390            let result = tokio::time::timeout(
391                Duration::from_millis(timeout_ms),
392                tools.execute_with_context_outcome(&call, execution_context),
393            )
394            .await;
395            let block = match result {
396                Ok(Ok(ToolOutcome::Completed(result)))
397                    if result.success && !tool_result_requires_approval(&result) =>
398                {
399                    let max_chars = declaration.max_chars.clamp(1, MAX_PROVIDER_CHARS);
400                    let redacted = redact_dynamic_context(&result.result);
401                    let original_chars = redacted.chars().count();
402                    let remaining_chars = MAX_TOTAL_CHARS.saturating_sub(total_chars);
403                    let output_chars = max_chars.min(remaining_chars);
404                    let content = redacted.chars().take(output_chars).collect::<String>();
405                    let truncated = original_chars > output_chars;
406                    total_chars = total_chars.saturating_add(content.chars().count());
407                    let ttl = declaration.cache_ttl_secs.min(MAX_CACHE_TTL_SECS);
408                    bamboo_skills::DynamicContextBlock {
409                        provider_id: declaration.id,
410                        tool: declaration.tool,
411                        provenance: "registered_tool_permission_checked".to_string(),
412                        generated_at: now,
413                        expires_at: (ttl > 0)
414                            .then(|| now + chrono::Duration::seconds(ttl as i64)),
415                        status: bamboo_skills::WorkflowActivationStatus::Active,
416                        stop_on_failure: declaration.stop_on_failure,
417                        content,
418                        diagnostic: truncated.then(|| bamboo_skills::WorkflowActivationDiagnostic {
419                            code: bamboo_skills::WorkflowActivationErrorCode::ProviderOutputInvalid,
420                            message: format!(
421                                "dynamic context output was truncated from {original_chars} to {output_chars} characters"
422                            ),
423                            recoverable: true,
424                        }),
425                    }
426                }
427                other => {
428                    let reason = match other {
429                        Err(_) => "provider timed out".to_string(),
430                        Ok(Err(error)) => format!("provider execution failed: {error}"),
431                        Ok(Ok(ToolOutcome::Completed(result)))
432                            if tool_result_requires_approval(&result) =>
433                        {
434                            "provider requires permission approval".to_string()
435                        }
436                        Ok(Ok(ToolOutcome::Completed(_))) => {
437                            "provider returned an unsuccessful result".to_string()
438                        }
439                        Ok(Ok(ToolOutcome::NeedsHuman { .. })) => {
440                            "provider requires human approval".to_string()
441                        }
442                        Ok(Ok(ToolOutcome::Running(_))) => {
443                            "provider attempted detached execution".to_string()
444                        }
445                    };
446                    bamboo_skills::DynamicContextBlock {
447                        provider_id: declaration.id,
448                        tool: declaration.tool,
449                        provenance: "registered_tool_permission_checked".to_string(),
450                        generated_at: now,
451                        expires_at: None,
452                        status: bamboo_skills::WorkflowActivationStatus::Degraded,
453                        stop_on_failure: declaration.stop_on_failure,
454                        content: String::new(),
455                        diagnostic: Some(bamboo_skills::WorkflowActivationDiagnostic {
456                            code: bamboo_skills::WorkflowActivationErrorCode::ProviderFailed,
457                            message: reason,
458                            recoverable: true,
459                        }),
460                    }
461                }
462            };
463            tracing::info!(
464                workflow_id = %skill.id,
465                provider_id = %block.provider_id,
466                tool = %block.tool,
467                status = ?block.status,
468                stop_on_failure = block.stop_on_failure,
469                content_chars = block.content.chars().count(),
470                "dynamic workflow context provider completed"
471            );
472            if block.expires_at.is_some() {
473                cache.insert(cache_key, block.clone());
474                while cache.len() > MAX_CACHE_ENTRIES
475                    || serde_json::to_string(&cache)
476                        .map(|raw| raw.chars().count() > MAX_CACHE_SERIALIZED_CHARS)
477                        .unwrap_or(true)
478                {
479                    let Some(oldest_key) = cache
480                        .iter()
481                        .min_by_key(|(_, block)| block.generated_at)
482                        .map(|(key, _)| key.clone())
483                    else {
484                        break;
485                    };
486                    cache.remove(&oldest_key);
487                }
488            }
489            blocks.push(block);
490            if blocks.last().is_some_and(|block| {
491                block.status == bamboo_skills::WorkflowActivationStatus::Degraded
492                    && block.stop_on_failure
493            }) {
494                break;
495            }
496        }
497        session.metadata.insert(
498            bamboo_skills::WORKFLOW_CONTEXT_CACHE_METADATA_KEY.to_string(),
499            serde_json::to_string(&cache).unwrap_or_else(|_| "{}".to_string()),
500        );
501        Ok(blocks)
502    }
503}
504
505fn confinement_diagnostic() -> bamboo_skills::WorkflowActivationDiagnostic {
506    bamboo_skills::WorkflowActivationDiagnostic {
507        code: bamboo_skills::WorkflowActivationErrorCode::ProviderFailed,
508        message:
509            "dynamic context provider input is not confined to the canonical session workspace"
510                .to_string(),
511        recoverable: true,
512    }
513}
514
515fn tool_result_requires_approval(result: &ToolResult) -> bool {
516    result
517        .display_preference
518        .as_deref()
519        .is_some_and(|value| value.eq_ignore_ascii_case("request_permissions"))
520        || serde_json::from_str::<serde_json::Value>(&result.result)
521            .ok()
522            .is_some_and(|value| {
523                value["status"].as_str().is_some_and(|status| {
524                    status.eq_ignore_ascii_case("awaiting_permission_approval")
525                }) || value["request_permissions"].is_object()
526            })
527}
528
529fn path_has_parent_component(path: &std::path::Path) -> bool {
530    path.components()
531        .any(|component| component == std::path::Component::ParentDir)
532}
533
534async fn canonicalize_scoped_path(
535    workspace: &std::path::Path,
536    raw: &str,
537    allow_missing_leaf: bool,
538) -> Result<std::path::PathBuf, bamboo_skills::WorkflowActivationDiagnostic> {
539    let path = std::path::PathBuf::from(raw.trim());
540    if raw.trim().is_empty() || path_has_parent_component(&path) {
541        return Err(confinement_diagnostic());
542    }
543    let candidate = if path.is_absolute() {
544        path
545    } else {
546        workspace.join(path)
547    };
548    match tokio::fs::canonicalize(&candidate).await {
549        Ok(canonical) if canonical.starts_with(workspace) => return Ok(canonical),
550        Ok(_) => return Err(confinement_diagnostic()),
551        Err(error) if allow_missing_leaf && error.kind() == std::io::ErrorKind::NotFound => {}
552        Err(_) => return Err(confinement_diagnostic()),
553    }
554
555    // GetFileInfo intentionally supports `exists:false`. Canonicalize the
556    // nearest existing ancestor (including symlinks), verify it remains in the
557    // workspace, then append only the already parent-free missing suffix.
558    let mut ancestor = candidate.as_path();
559    while !ancestor.exists() {
560        ancestor = ancestor.parent().ok_or_else(confinement_diagnostic)?;
561    }
562    let canonical_ancestor = tokio::fs::canonicalize(ancestor)
563        .await
564        .map_err(|_| confinement_diagnostic())?;
565    if !canonical_ancestor.starts_with(workspace) {
566        return Err(confinement_diagnostic());
567    }
568    let suffix = candidate
569        .strip_prefix(ancestor)
570        .map_err(|_| confinement_diagnostic())?;
571    Ok(canonical_ancestor.join(suffix))
572}
573
574async fn confine_dynamic_provider_input(
575    session: &bamboo_agent_core::Session,
576    tool: &str,
577    input: &serde_json::Value,
578) -> Result<serde_json::Value, bamboo_skills::WorkflowActivationDiagnostic> {
579    let workspace = session
580        .workspace_path_meta()
581        .filter(|path| !path.trim().is_empty())
582        .ok_or_else(confinement_diagnostic)?;
583    let workspace = tokio::fs::canonicalize(workspace)
584        .await
585        .map_err(|_| confinement_diagnostic())?;
586    let mut normalized = input
587        .as_object()
588        .cloned()
589        .ok_or_else(confinement_diagnostic)?;
590    let canonical_tool = tool.to_ascii_lowercase();
591    if matches!(canonical_tool.as_str(), "read" | "read_file") {
592        let raw = normalized
593            .get("file_path")
594            .or_else(|| normalized.get("path"))
595            .and_then(serde_json::Value::as_str)
596            .ok_or_else(confinement_diagnostic)?;
597        let canonical = canonicalize_scoped_path(&workspace, raw, false).await?;
598        normalized.remove("path");
599        normalized.insert(
600            "file_path".to_string(),
601            serde_json::Value::String(canonical.to_string_lossy().into_owned()),
602        );
603    } else if canonical_tool == "getfileinfo" {
604        let raw = normalized
605            .get("path")
606            .and_then(serde_json::Value::as_str)
607            .ok_or_else(confinement_diagnostic)?;
608        let canonical = canonicalize_scoped_path(&workspace, raw, true).await?;
609        normalized.remove("file_path");
610        normalized.insert(
611            "path".to_string(),
612            serde_json::Value::String(canonical.to_string_lossy().into_owned()),
613        );
614    } else if matches!(canonical_tool.as_str(), "glob" | "grep" | "list_directory") {
615        let raw = normalized
616            .get("path")
617            .and_then(serde_json::Value::as_str)
618            .unwrap_or(".");
619        let canonical = canonicalize_scoped_path(&workspace, raw, false).await?;
620        normalized.insert(
621            "path".to_string(),
622            serde_json::Value::String(canonical.to_string_lossy().into_owned()),
623        );
624        for pattern_key in ["pattern", "glob"] {
625            if let Some(pattern) = normalized
626                .get(pattern_key)
627                .and_then(serde_json::Value::as_str)
628            {
629                let pattern_path = std::path::Path::new(pattern);
630                if pattern_path.is_absolute() || path_has_parent_component(pattern_path) {
631                    return Err(confinement_diagnostic());
632                }
633            }
634        }
635    } else {
636        return Err(confinement_diagnostic());
637    }
638    Ok(serde_json::Value::Object(normalized))
639}
640
641fn redact_dynamic_context(raw: &str) -> String {
642    fn redact_value(value: &mut serde_json::Value) {
643        match value {
644            serde_json::Value::Object(map) => {
645                for (key, value) in map {
646                    let sensitive = [
647                        "token",
648                        "secret",
649                        "password",
650                        "authorization",
651                        "cookie",
652                        "api_key",
653                    ]
654                    .iter()
655                    .any(|needle| key.to_ascii_lowercase().contains(needle));
656                    if sensitive {
657                        *value = serde_json::Value::String("[REDACTED]".to_string());
658                    } else {
659                        redact_value(value);
660                    }
661                }
662            }
663            serde_json::Value::Array(values) => values.iter_mut().for_each(redact_value),
664            serde_json::Value::String(text) => {
665                let lower = text.to_ascii_lowercase();
666                if [
667                    "-----begin",
668                    "private key-----",
669                    "bearer ",
670                    "akia",
671                    "github_pat_",
672                    "ghp_",
673                    "sk-",
674                ]
675                .iter()
676                .any(|marker| lower.contains(marker))
677                {
678                    *text = "[REDACTED]".to_string();
679                }
680            }
681            _ => {}
682        }
683    }
684    if let Ok(mut value) = serde_json::from_str::<serde_json::Value>(raw) {
685        redact_value(&mut value);
686        return value.to_string();
687    }
688    let mut in_private_key = false;
689    raw.lines()
690        .map(|line| {
691            let lower = line.to_ascii_lowercase();
692            if lower.contains("-----begin") && lower.contains("private key-----") {
693                in_private_key = true;
694                return "[REDACTED]".to_string();
695            }
696            if in_private_key {
697                if lower.contains("-----end") && lower.contains("private key-----") {
698                    in_private_key = false;
699                }
700                return "[REDACTED]".to_string();
701            }
702            if [
703                "token=",
704                "token:",
705                "secret=",
706                "secret:",
707                "password=",
708                "password:",
709                "authorization:",
710                "api_key=",
711                "api_key:",
712                "bearer ",
713                "akia",
714                "github_pat_",
715                "ghp_",
716                "sk-",
717            ]
718            .iter()
719            .any(|needle| lower.contains(needle))
720            {
721                "[REDACTED]".to_string()
722            } else {
723                line.to_string()
724            }
725        })
726        .collect::<Vec<_>>()
727        .join("\n")
728}
729
730#[cfg(test)]
731#[allow(clippy::items_after_test_module)]
732mod redaction_tests {
733    use super::{confine_dynamic_provider_input, redact_dynamic_context};
734    use bamboo_agent_core::tools::{FunctionCall, ToolCall, ToolExecutor};
735    use bamboo_agent_core::Session;
736    use bamboo_tools::BuiltinToolExecutorBuilder;
737
738    async fn execute(
739        executor: &dyn ToolExecutor,
740        name: &str,
741        input: serde_json::Value,
742    ) -> bamboo_agent_core::tools::ToolResult {
743        let call = ToolCall {
744            id: format!("test-{name}"),
745            tool_type: "function".to_string(),
746            function: FunctionCall {
747                name: name.to_string(),
748                arguments: input.to_string(),
749            },
750        };
751        executor.execute(&call).await.expect("builtin executes")
752    }
753
754    #[test]
755    fn redacts_common_non_json_secret_formats() {
756        let raw = "safe=value\ntoken: tok-value\nAuthorization: Bearer bearer-value\nAWS_ACCESS_KEY_ID=AKIAABCDEFGHIJKLMNOP\nghp_exampletoken\n-----BEGIN PRIVATE KEY-----\nprivate-material\n-----END PRIVATE KEY-----\nafter=safe";
757        let redacted = redact_dynamic_context(raw);
758        assert!(redacted.contains("safe=value"));
759        assert!(redacted.contains("after=safe"));
760        for secret in [
761            "tok-value",
762            "bearer-value",
763            "AKIAABCDEFGHIJKLMNOP",
764            "ghp_exampletoken",
765            "private-material",
766        ] {
767            assert!(!redacted.contains(secret), "leaked {secret}");
768        }
769
770        let json = redact_dynamic_context(
771            r#"{"data":"-----BEGIN PRIVATE KEY----- hidden","value":"ghp_json_secret"}"#,
772        );
773        assert!(!json.contains("hidden"));
774        assert!(!json.contains("ghp_json_secret"));
775    }
776
777    #[tokio::test]
778    async fn real_builtins_are_workspace_confined_with_aliases_and_missing_file_info() {
779        let workspace = tempfile::tempdir().expect("workspace");
780        let outside = tempfile::tempdir().expect("outside");
781        std::fs::create_dir_all(workspace.path().join("src")).expect("src");
782        std::fs::write(
783            workspace.path().join("src/inside.txt"),
784            "WORKSPACE_ONLY_MARKER",
785        )
786        .expect("inside file");
787        std::fs::write(outside.path().join("secret.txt"), "OUTSIDE_SECRET").expect("outside");
788        #[cfg(unix)]
789        std::os::unix::fs::symlink(outside.path(), workspace.path().join("escape"))
790            .expect("escape symlink");
791
792        let mut session = Session::new("confinement", "model");
793        session.set_workspace_path_meta(workspace.path().to_string_lossy().into_owned());
794        let canonical_workspace =
795            std::fs::canonicalize(workspace.path()).expect("canonical workspace");
796        let executor = BuiltinToolExecutorBuilder::new()
797            .with_default_tools()
798            .build();
799
800        for (name, input, expected_field) in [
801            (
802                "Read",
803                serde_json::json!({"file_path": "src/inside.txt"}),
804                "file_path",
805            ),
806            (
807                "read_file",
808                serde_json::json!({"path": "src/inside.txt"}),
809                "file_path",
810            ),
811            (
812                "Glob",
813                serde_json::json!({"path": "src", "pattern": "*.txt"}),
814                "path",
815            ),
816            ("list_directory", serde_json::json!({"path": "src"}), "path"),
817            (
818                "Grep",
819                serde_json::json!({"path": "src", "pattern": "WORKSPACE_ONLY"}),
820                "path",
821            ),
822        ] {
823            let confined = confine_dynamic_provider_input(&session, name, &input)
824                .await
825                .unwrap_or_else(|error| panic!("{name} should be confined: {}", error.message));
826            assert!(
827                confined[expected_field].as_str().is_some_and(
828                    |path| path.starts_with(canonical_workspace.to_string_lossy().as_ref())
829                ),
830                "{name} normalized input: {confined}"
831            );
832            let result = execute(&executor, name, confined).await;
833            assert!(result.success, "{name}: {}", result.result);
834            assert!(!result.result.contains("OUTSIDE_SECRET"));
835        }
836
837        let missing = confine_dynamic_provider_input(
838            &session,
839            "GetFileInfo",
840            &serde_json::json!({"path": "src/missing.txt"}),
841        )
842        .await
843        .expect("missing workspace leaf is safe");
844        assert!(missing.get("file_path").is_none());
845        let result = execute(&executor, "GetFileInfo", missing).await;
846        assert!(result.success);
847        assert_eq!(
848            serde_json::from_str::<serde_json::Value>(&result.result).expect("file info json")
849                ["exists"],
850            false
851        );
852
853        for (name, input) in [
854            (
855                "Read",
856                serde_json::json!({"file_path": outside.path().join("secret.txt")}),
857            ),
858            (
859                "Read",
860                serde_json::json!({"file_path": "../outside/secret.txt"}),
861            ),
862            (
863                "GetFileInfo",
864                serde_json::json!({"path": "escape/missing.txt"}),
865            ),
866        ] {
867            let error = confine_dynamic_provider_input(&session, name, &input)
868                .await
869                .expect_err("escape must be denied");
870            assert!(!error.message.contains("secret.txt"));
871            assert!(!error
872                .message
873                .contains(outside.path().to_string_lossy().as_ref()));
874        }
875
876        let no_workspace = Session::new("no-workspace", "model");
877        assert!(confine_dynamic_provider_input(
878            &no_workspace,
879            "Read",
880            &serde_json::json!({"file_path": "src/inside.txt"})
881        )
882        .await
883        .is_err());
884    }
885}
886
887#[async_trait]
888impl Tool for LoadSkillTool {
889    fn name(&self) -> &str {
890        "load_skill"
891    }
892
893    fn description(&self) -> &str {
894        "Load a skill's detailed SKILL.md instructions by skill_id."
895    }
896
897    fn parameters_schema(&self) -> serde_json::Value {
898        json!({
899            "type": "object",
900            "properties": {
901                "skill_id": {
902                    "type": "string",
903                    "description": "Skill ID from the advertised skill list (for example: skill-creator)."
904                }
905            },
906            "required": ["skill_id"]
907        })
908    }
909
910    async fn invoke(
911        &self,
912        args: serde_json::Value,
913        ctx: ToolCtx,
914    ) -> Result<ToolOutcome, ToolError> {
915        let parsed: LoadSkillArgs = serde_json::from_value(args).map_err(|err| {
916            ToolError::InvalidArguments(format!("Invalid load_skill args: {err}"))
917        })?;
918        let skill_id = parsed.skill_id.trim();
919        if skill_id.is_empty() {
920            return Err(ToolError::InvalidArguments(
921                "skill_id must be a non-empty string".to_string(),
922            ));
923        }
924
925        let session_id = ctx.session_id().ok_or_else(|| {
926            ToolError::Execution("load_skill requires a session_id in tool context".to_string())
927        })?;
928        let store = self.access.skill_store(ctx.session_id()).await?;
929        let reused_published_activation =
930            validate_runtime_activation(&self.access, store.as_ref(), session_id, skill_id).await?;
931        if !reused_published_activation {
932            access_control::ensure_skill_allowed(&self.access, skill_id, ctx.session_id())
933                .await
934                .map_err(skill_access_error_to_tool_error)?;
935            let skill_mode =
936                access_control::selected_skill_mode(&self.access, ctx.session_id()).await;
937            let selected_ids =
938                access_control::selected_skill_allowlist(&self.access, ctx.session_id())
939                    .await
940                    .ok_or_else(|| {
941                        ToolError::Execution(
942                            "load_skill cannot pin a request with no published skill selection"
943                                .to_string(),
944                        )
945                    })?
946                    .into_iter()
947                    .collect::<Vec<_>>();
948            self.access
949                .pin_current_activation(session_id, &selected_ids, skill_mode.as_deref())
950                .await
951                .map_err(|err| {
952                    ToolError::Execution(format!(
953                        "Failed to pin workflow activation for '{skill_id}': {err}"
954                    ))
955                })?;
956        }
957        let (skill, skill_root, revision, resources, payload_descriptor) = store
958            .get_pinned_skill_with_root_and_descriptor(session_id, skill_id)
959            .await
960            .map_err(|err| {
961                ToolError::Execution(format!("Failed to load skill '{skill_id}': {err}"))
962            })?;
963        validate_runtime_activation_descriptor(
964            &self.access,
965            &payload_descriptor,
966            session_id,
967            skill_id,
968        )
969        .await?;
970        let catalog_entry = store
971            .pinned_activation_catalog_entries(session_id)
972            .await
973            .and_then(|entries| entries.into_iter().find(|entry| entry.id == skill_id))
974            .ok_or_else(|| {
975                ToolError::Execution("pinned workflow catalog entry is missing".to_string())
976            })?;
977        if catalog_entry.kind != bamboo_skills::WorkflowKind::Instruction {
978            return Err(ToolError::InvalidArguments(
979                "orchestration workflows cannot be loaded as instruction skills; use workflow_run"
980                    .to_string(),
981            ));
982        }
983        let restored_snapshot = store.activation_was_restored(session_id).await;
984        let canonical_skill_root = if restored_snapshot
985            || catalog_entry.migration_status
986                == Some(bamboo_skills::LegacyWorkflowMigrationStatus::Available)
987        {
988            None
989        } else {
990            Some(
991                tokio::fs::canonicalize(&skill_root)
992                    .await
993                    .unwrap_or(skill_root),
994            )
995        };
996        let mut session = self
997            .access
998            .session_for_context(Some(session_id))
999            .await
1000            .ok_or_else(|| ToolError::Execution(format!("Session '{session_id}' not found")))?;
1001        // Direct SDK/tool integrations may pin through the guarded fallback
1002        // above without running the agent's setup publisher. Materialize the
1003        // same immutable catalog/snapshot metadata from that pin; never rebuild
1004        // it from the live catalog after the fact.
1005        if !reused_published_activation
1006            || !session
1007                .metadata
1008                .contains_key(bamboo_skills::runtime_metadata::SKILL_RUNTIME_SELECTED_CATALOG_KEY)
1009        {
1010            let pinned_catalog = store
1011                .pinned_activation_catalog_entries(session_id)
1012                .await
1013                .ok_or_else(|| {
1014                    ToolError::Execution("pinned workflow catalog is unavailable".to_string())
1015                })?;
1016            session.metadata.insert(
1017                bamboo_skills::runtime_metadata::SKILL_RUNTIME_SELECTED_CATALOG_KEY.to_string(),
1018                serde_json::to_string(&pinned_catalog).map_err(|_| {
1019                    ToolError::Execution("pinned workflow catalog is invalid".to_string())
1020                })?,
1021            );
1022        }
1023        if !reused_published_activation
1024            || !session
1025                .metadata
1026                .contains_key(bamboo_skills::runtime_metadata::SKILL_RUNTIME_PINNED_SNAPSHOT_KEY)
1027        {
1028            let snapshot = store
1029                .export_activation_snapshot(session_id)
1030                .await
1031                .ok_or_else(|| {
1032                    ToolError::Execution("pinned workflow snapshot is unavailable".to_string())
1033                })?;
1034            session.metadata.insert(
1035                bamboo_skills::runtime_metadata::SKILL_RUNTIME_PINNED_SNAPSHOT_KEY.to_string(),
1036                serde_json::to_string(&snapshot).map_err(|_| {
1037                    ToolError::Execution("pinned workflow snapshot is invalid".to_string())
1038                })?,
1039            );
1040        }
1041        if let Some(selection) = session
1042            .metadata
1043            .get(bamboo_skills::WORKFLOW_SELECTION_METADATA_KEY)
1044            .and_then(|raw| serde_json::from_str::<bamboo_skills::WorkflowSelection>(raw).ok())
1045        {
1046            bamboo_domain::validate_schema(&catalog_entry.argument_schema, &selection.args)
1047                .map_err(|error| {
1048                    ToolError::InvalidArguments(format!(
1049                        "workflow arguments do not match the pinned schema: {error}"
1050                    ))
1051                })?;
1052        }
1053        let dynamic_context = self
1054            .resolve_dynamic_context(&skill, &mut session, &ctx)
1055            .await?;
1056        session.metadata.insert(
1057            bamboo_skills::WORKFLOW_LAST_DYNAMIC_CONTEXT_METADATA_KEY.to_string(),
1058            serde_json::to_string(&dynamic_context).unwrap_or_else(|_| "[]".to_string()),
1059        );
1060        let activation_stopped = dynamic_context.iter().any(|block| {
1061            block.status == bamboo_skills::WorkflowActivationStatus::Degraded
1062                && block.stop_on_failure
1063        });
1064        let payload = json!({
1065            "skill_id": skill.id.clone(),
1066            "revision": revision,
1067            "name": skill.name.clone(),
1068            "description": skill.description.clone(),
1069            "license": skill.license.clone(),
1070            "compatibility": skill.compatibility.clone(),
1071            "allowed_tools": skill.tool_refs.clone(),
1072            "instructions": skill.prompt.clone(),
1073            "skill_base_dir": canonical_skill_root
1074                .as_ref()
1075                .map(|root| bamboo_config::paths::path_to_display_string(root)),
1076            "snapshot_provenance": if restored_snapshot { "durable_session_lkg" } else { "live_catalog_pin" },
1077            "resource_files": resources,
1078            "dynamic_context": dynamic_context.clone(),
1079            // Non-stopping provider degradation is carried inside the typed
1080            // runtime blocks while the workflow itself remains active.
1081            "activation_status": if activation_stopped { "degraded" } else { "active" },
1082        });
1083        if activation_stopped {
1084            session.metadata.insert(
1085                bamboo_skills::runtime_metadata::SKILL_RUNTIME_ACTIVATION_ERROR_KEY.to_string(),
1086                serde_json::to_string(&dynamic_context).unwrap_or_else(|_| {
1087                    "dynamic workflow context degraded before activation".to_string()
1088                }),
1089            );
1090            session
1091                .metadata
1092                .remove(bamboo_skills::ACTIVE_WORKFLOW_METADATA_KEY);
1093            session
1094                .metadata
1095                .remove(bamboo_skills::ACTIVE_WORKFLOW_SNAPSHOT_METADATA_KEY);
1096            self.persist_owned_metadata(
1097                session_id,
1098                &session,
1099                "degraded workflow state could not be persisted",
1100            )
1101            .await?;
1102            return Ok(ToolOutcome::Completed(ToolResult {
1103                success: true,
1104                result: payload.to_string(),
1105                display_preference: Some("Collapsible".to_string()),
1106                images: Vec::new(),
1107            }));
1108        }
1109        let canonical_context = json!({
1110            "id": skill.id,
1111            "revision": revision,
1112            "selection": session.metadata.get(bamboo_skills::WORKFLOW_SELECTION_METADATA_KEY),
1113            "instructions": skill.prompt,
1114            "dynamic_context": payload.get("dynamic_context"),
1115        });
1116        let fingerprint = hex::encode(Sha256::digest(canonical_context.to_string().as_bytes()));
1117        let mut loaded_ids = session
1118            .metadata
1119            .get(LOADED_SKILL_IDS_METADATA_KEY)
1120            .map(|raw| access_control::parse_loaded_skill_ids(raw))
1121            .unwrap_or_default();
1122        loaded_ids.insert(skill_id.to_string());
1123        session.metadata.insert(
1124            LOADED_SKILL_IDS_METADATA_KEY.to_string(),
1125            access_control::serialize_loaded_skill_ids(&loaded_ids),
1126        );
1127        session.metadata.insert(
1128            LAST_LOADED_SKILL_ID_METADATA_KEY.to_string(),
1129            skill_id.to_string(),
1130        );
1131        session.metadata.insert(
1132            LAST_LOADED_SKILL_SUMMARY_METADATA_KEY.to_string(),
1133            json!({"skill_id": skill_id, "loaded_count": loaded_ids.len()}).to_string(),
1134        );
1135        let already_active = session
1136            .metadata
1137            .get(bamboo_skills::ACTIVE_WORKFLOW_METADATA_KEY)
1138            .and_then(|raw| serde_json::from_str::<bamboo_skills::ActiveWorkflow>(raw).ok())
1139            .is_some_and(|active| {
1140                active.id == skill_id
1141                    && active.revision == revision
1142                    && active.status == bamboo_skills::WorkflowActivationStatus::Active
1143            });
1144        if already_active {
1145            // Repeated model calls may reload the payload, but they are not a
1146            // second activation and must not duplicate lifecycle events.
1147            self.persist_owned_metadata(
1148                session_id,
1149                &session,
1150                "workflow loaded state could not be persisted",
1151            )
1152            .await?;
1153            return Ok(ToolOutcome::Completed(ToolResult {
1154                success: true,
1155                result: payload.to_string(),
1156                display_preference: Some("Collapsible".to_string()),
1157                images: Vec::new(),
1158            }));
1159        }
1160        let active = match bamboo_skills::record_loaded_workflow_activation(
1161            &mut session.metadata,
1162            skill_id,
1163            fingerprint,
1164        ) {
1165            Ok(active) => active,
1166            Err(diagnostic) => {
1167                session.metadata.insert(
1168                    bamboo_skills::runtime_metadata::SKILL_RUNTIME_ACTIVATION_ERROR_KEY.to_string(),
1169                    serde_json::to_string(&diagnostic)
1170                        .unwrap_or_else(|_| diagnostic.message.clone()),
1171                );
1172                session
1173                    .metadata
1174                    .remove(bamboo_skills::ACTIVE_WORKFLOW_METADATA_KEY);
1175                session
1176                    .metadata
1177                    .remove(bamboo_skills::ACTIVE_WORKFLOW_SNAPSHOT_METADATA_KEY);
1178                session
1179                    .metadata
1180                    .remove(bamboo_skills::WORKFLOW_ACTIVATION_EVENT_METADATA_KEY);
1181                self.persist_owned_metadata(
1182                    session_id,
1183                    &session,
1184                    "degraded workflow activation could not be persisted",
1185                )
1186                .await?;
1187                return Ok(ToolOutcome::Completed(ToolResult {
1188                    success: true,
1189                    result: json!({
1190                        "skill_id": skill_id,
1191                        "revision": revision,
1192                        "activation_status": "degraded",
1193                        "diagnostic": diagnostic,
1194                    })
1195                    .to_string(),
1196                    display_preference: Some("Collapsible".to_string()),
1197                    images: Vec::new(),
1198                }));
1199            }
1200        };
1201        self.persist_owned_metadata(
1202            session_id,
1203            &session,
1204            "workflow activation could not be persisted",
1205        )
1206        .await?;
1207        let pending_event = session
1208            .metadata
1209            .get(bamboo_skills::WORKFLOW_ACTIVATION_EVENT_METADATA_KEY)
1210            .and_then(|raw| serde_json::from_str::<serde_json::Value>(raw).ok())
1211            .ok_or_else(|| {
1212                ToolError::Execution("workflow activation event metadata is invalid".to_string())
1213            })?;
1214        let activation_event = bamboo_agent_core::AgentEvent::WorkflowActivated {
1215            event_id: bamboo_skills::workflow_lifecycle_event_id(session_id, &pending_event),
1216            session_id: session_id.to_string(),
1217            workflow_id: active.id,
1218            revision: active.revision,
1219            invoked_by: format!("{:?}", active.invoked_by).to_ascii_lowercase(),
1220        };
1221        if let Some(sender) = ctx.cloned_sender() {
1222            if sender.send(activation_event).await.is_ok() {
1223                // The persisted pending event is the crash-recovery publication
1224                // record. Remove it only after the live runner accepted it.
1225                session
1226                    .metadata
1227                    .remove(bamboo_skills::WORKFLOW_ACTIVATION_EVENT_METADATA_KEY);
1228                self.persist_owned_metadata(
1229                    session_id,
1230                    &session,
1231                    "workflow activation event acknowledgement could not be persisted",
1232                )
1233                .await?;
1234            }
1235        }
1236
1237        Ok(ToolOutcome::Completed(ToolResult {
1238            success: true,
1239            result: payload.to_string(),
1240            display_preference: Some("Collapsible".to_string()),
1241            images: Vec::new(),
1242        }))
1243    }
1244}