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 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 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 #[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 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 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 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 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 "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 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 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}