From d33fecb17c3e81404f021ee944aa75673d1b4f97 Mon Sep 17 00:00:00 2001 From: Henry Park Date: Wed, 22 Apr 2026 11:36:25 -0700 Subject: [PATCH] engine-v2: centralize action vs capability surface policy (#2827) * Add canonical engine capability status enum * Add bridge tool surface assignment policy * fix(engine): tighten scoped surface assignment * fix(bridge): remove premature approval_gated field, surface ReadyScoped in capabilities - Remove approval_gated from SurfacePolicyInput (YAGNI until policy uses it) - Change ReadyScoped fallback from neither() to capabilities_only() so scoped subjects remain visible in background context - Update tests to match new behavior Co-Authored-By: Claude Opus 4.6 (1M context) * fix(ci): unblock section 2 policy PR * engine-v2: add capability projection and two-surface prompt baseline (#2826) * Add capability projection and two-surface prompt baseline * Reduce step-context args for clippy-clean two-surface stack * fix(engine): address two-surface review follow-ups * fix(engine): normalize alias-aware capability projection * fix(bridge): share extension fetch between projectors, preserve NeedsAuth in actions - Fetch list_capability_extensions once in EffectBridgeAdapter and pass to both ActionProjector and CapabilityProjector via prefetched_extensions - Keep NeedsAuth provider tools in available_actions so the LLM can trigger auth gates by attempting to call them - Add unit tests for NeedsAuth preservation and latent tool omission at the ActionProjector level where extension maps can be controlled Co-Authored-By: Claude Opus 4.6 (1M context) --------- Co-authored-by: serrrfirat Co-authored-by: Claude Opus 4.6 (1M context) --------- Co-authored-by: serrrfirat Co-authored-by: Claude Opus 4.6 (1M context) --- .../ironclaw_engine/src/executor/context.rs | 67 ++- .../src/executor/loop_engine.rs | 149 +++++- crates/ironclaw_engine/src/executor/mod.rs | 1 + .../src/executor/orchestrator.rs | 103 +--- crates/ironclaw_engine/src/executor/prompt.rs | 130 ++++- .../ironclaw_engine/src/executor/scripting.rs | 13 +- .../src/executor/structured.rs | 11 +- .../src/executor/thread_context.rs | 31 ++ crates/ironclaw_engine/src/lib.rs | 5 +- .../src/runtime/conversation.rs | 9 + .../src/runtime/lease_refresh.rs | 5 +- crates/ironclaw_engine/src/runtime/manager.rs | 18 + crates/ironclaw_engine/src/runtime/mission.rs | 18 + crates/ironclaw_engine/src/traits/effect.rs | 10 +- .../ironclaw_engine/src/types/capability.rs | 96 ++++ deny.toml | 1 - src/auth/extension.rs | 31 +- src/bridge/action_projector.rs | 390 +++++++++++++++ src/bridge/capability_projector.rs | 463 +++++++++++++++++ src/bridge/effect_adapter.rs | 470 ++++++++++++++---- src/bridge/mod.rs | 3 + src/bridge/router.rs | 39 ++ src/bridge/tool_surface.rs | 289 +++++++++++ tests/engine_v2_gate_integration.rs | 18 + tests/engine_v2_skill_codeact.rs | 9 + 25 files changed, 2148 insertions(+), 231 deletions(-) create mode 100644 crates/ironclaw_engine/src/executor/thread_context.rs create mode 100644 src/bridge/action_projector.rs create mode 100644 src/bridge/capability_projector.rs create mode 100644 src/bridge/tool_surface.rs diff --git a/crates/ironclaw_engine/src/executor/context.rs b/crates/ironclaw_engine/src/executor/context.rs index 9a48b204f91..85f38d669c1 100644 --- a/crates/ironclaw_engine/src/executor/context.rs +++ b/crates/ironclaw_engine/src/executor/context.rs @@ -6,12 +6,11 @@ use std::sync::Arc; use crate::memory::RetrievalEngine; -use crate::traits::effect::EffectExecutor; +use crate::traits::effect::{EffectExecutor, ThreadExecutionContext}; use crate::types::capability::{ActionDef, CapabilityLease}; use crate::types::error::EngineError; use crate::types::memory::MemoryDoc; use crate::types::message::ThreadMessage; -use crate::types::project::ProjectId; /// Maximum number of memory docs to inject into context. const MAX_CONTEXT_DOCS: usize = 5; @@ -26,16 +25,19 @@ pub async fn build_step_context( leases: &[CapabilityLease], effects: &Arc, retrieval: Option<&RetrievalEngine>, - project_id: ProjectId, - user_id: &str, - goal: &str, + context: &ThreadExecutionContext, ) -> Result<(Vec, Vec), EngineError> { // Fetch actions and memory docs in parallel — they are independent. - let actions_fut = effects.available_actions(leases); + let actions_fut = effects.available_actions(leases, context); let docs_fut = async { if let Some(engine) = retrieval { engine - .retrieve_context(project_id, user_id, goal, MAX_CONTEXT_DOCS) + .retrieve_context( + context.project_id, + &context.user_id, + context.thread_goal.as_deref().unwrap_or(""), + MAX_CONTEXT_DOCS, + ) .await } else { Ok(Vec::new()) @@ -130,9 +132,18 @@ mod tests { async fn available_actions( &self, _: &[CapabilityLease], + _: &crate::traits::effect::ThreadExecutionContext, ) -> Result, EngineError> { Ok(vec![]) } + + async fn available_capabilities( + &self, + _: &[CapabilityLease], + _: &crate::traits::effect::ThreadExecutionContext, + ) -> Result, EngineError> { + Ok(vec![]) + } } #[tokio::test] @@ -161,9 +172,17 @@ mod tests { &[], &effects, Some(&retrieval), - project, - "test-user", - "search the web", + &crate::traits::effect::ThreadExecutionContext { + thread_id: crate::types::thread::ThreadId::new(), + thread_type: crate::types::thread::ThreadType::Foreground, + project_id: project, + user_id: "test-user".into(), + step_id: crate::types::step::StepId::new(), + current_call_id: None, + source_channel: None, + user_timezone: None, + thread_goal: Some("search the web".into()), + }, ) .await .unwrap(); @@ -191,9 +210,17 @@ mod tests { &[], &effects, None, - ProjectId::new(), - "test-user", - "hello", + &crate::traits::effect::ThreadExecutionContext { + thread_id: crate::types::thread::ThreadId::new(), + thread_type: crate::types::thread::ThreadType::Foreground, + project_id: ProjectId::new(), + user_id: "test-user".into(), + step_id: crate::types::step::StepId::new(), + current_call_id: None, + source_channel: None, + user_timezone: None, + thread_goal: Some("hello".into()), + }, ) .await .unwrap(); @@ -217,9 +244,17 @@ mod tests { &[], &effects, Some(&retrieval), - project, - "test-user", - "hello", + &crate::traits::effect::ThreadExecutionContext { + thread_id: crate::types::thread::ThreadId::new(), + thread_type: crate::types::thread::ThreadType::Foreground, + project_id: project, + user_id: "test-user".into(), + step_id: crate::types::step::StepId::new(), + current_call_id: None, + source_channel: None, + user_timezone: None, + thread_goal: Some("hello".into()), + }, ) .await .unwrap(); diff --git a/crates/ironclaw_engine/src/executor/loop_engine.rs b/crates/ironclaw_engine/src/executor/loop_engine.rs index 175bf4bb2bb..a6f200e003c 100644 --- a/crates/ironclaw_engine/src/executor/loop_engine.rs +++ b/crates/ironclaw_engine/src/executor/loop_engine.rs @@ -18,7 +18,7 @@ use crate::traits::llm::LlmBackend; use crate::types::error::EngineError; use crate::types::event::EventKind; use crate::types::message::ThreadMessage; -use crate::types::step::Step; +use crate::types::step::{Step, StepId}; use crate::types::thread::{Thread, ThreadState}; const RUNTIME_CHECKPOINT_METADATA_KEY: &str = "runtime_checkpoint"; @@ -217,16 +217,37 @@ impl ExecutionLoop { { // Fetch active leases (needed for action list) let active_leases = self.leases.active_for_thread(self.thread.id).await; - let actions = match self.effects.available_actions(&active_leases).await { + let prompt_context = crate::executor::thread_context::thread_execution_context( + &self.thread, + StepId::new(), + None, + ); + let actions = match self + .effects + .available_actions(&active_leases, &prompt_context) + .await + { Ok(a) => a, Err(e) => { debug!(thread_id = %self.thread.id, "failed to load actions for system prompt: {e}"); Vec::new() } }; + let capabilities = match self + .effects + .available_capabilities(&active_leases, &prompt_context) + .await + { + Ok(c) => c, + Err(e) => { + debug!(thread_id = %self.thread.id, "failed to load capabilities for system prompt: {e}"); + Vec::new() + } + }; // Build prompt using pre-fetched docs (no extra Store query) let system_prompt = crate::executor::prompt::build_codeact_system_prompt_with_docs( &actions, + &capabilities, &system_docs, self.platform_info.as_ref(), ); @@ -412,7 +433,10 @@ mod tests { use crate::runtime::messaging::ThreadSignal; use crate::traits::effect::ThreadExecutionContext; use crate::traits::llm::{LlmCallConfig, LlmOutput}; - use crate::types::capability::{ActionDef, CapabilityLease, EffectType, GrantedActions}; + use crate::types::capability::{ + ActionDef, CapabilityLease, CapabilityStatus, CapabilitySummary, CapabilitySummaryKind, + EffectType, GrantedActions, + }; use crate::types::project::ProjectId; use crate::types::step::LlmResponse; use crate::types::step::{ActionResult, TokenUsage}; @@ -501,9 +525,18 @@ mod tests { async fn available_actions( &self, _leases: &[CapabilityLease], + _context: &ThreadExecutionContext, ) -> Result, EngineError> { Ok(self.actions.clone()) } + + async fn available_capabilities( + &self, + _: &[CapabilityLease], + _: &ThreadExecutionContext, + ) -> Result, EngineError> { + Ok(vec![]) + } } // ── Helpers ───────────────────────────────────────────── @@ -597,6 +630,116 @@ mod tests { assert!(exec.thread.total_tokens_used > 0); } + #[tokio::test] + async fn system_prompt_includes_capability_background() { + struct CapturingLlm { + seen_messages: Mutex>>, + } + + #[async_trait::async_trait] + impl LlmBackend for CapturingLlm { + async fn complete( + &self, + messages: &[ThreadMessage], + _actions: &[ActionDef], + _config: &LlmCallConfig, + ) -> Result { + self.seen_messages.lock().unwrap().push(messages.to_vec()); + Ok(text_response("done")) + } + + fn model_name(&self) -> &str { + "capturing" + } + } + + struct CapabilityEffects; + + #[async_trait::async_trait] + impl EffectExecutor for CapabilityEffects { + async fn execute_action( + &self, + _action_name: &str, + _parameters: serde_json::Value, + _lease: &CapabilityLease, + _context: &ThreadExecutionContext, + ) -> Result { + Ok(ActionResult { + call_id: String::new(), + action_name: String::new(), + output: serde_json::json!({}), + is_error: false, + duration: Duration::from_millis(1), + }) + } + + async fn available_actions( + &self, + _leases: &[CapabilityLease], + _context: &ThreadExecutionContext, + ) -> Result, EngineError> { + Ok(vec![test_action()]) + } + + async fn available_capabilities( + &self, + _: &[CapabilityLease], + _: &ThreadExecutionContext, + ) -> Result, EngineError> { + Ok(vec![CapabilitySummary { + name: "telegram".into(), + display_name: Some("Telegram".into()), + kind: CapabilitySummaryKind::Channel, + status: CapabilityStatus::ReadyScoped, + description: Some("Telegram notifications".into()), + routing_hint: Some("Usable through message".into()), + }]) + } + } + + let project_id = ProjectId::new(); + let thread = Thread::new( + "test goal", + ThreadType::Foreground, + project_id, + "test-user", + ThreadConfig::default(), + ); + let tid = thread.id; + + let llm = Arc::new(CapturingLlm { + seen_messages: Mutex::new(Vec::new()), + }); + let effects: Arc = Arc::new(CapabilityEffects); + let leases = Arc::new(LeaseManager::new()); + let policy = Arc::new(PolicyEngine::new()); + leases + .grant(tid, "tools", GrantedActions::All, None, None) + .await + .unwrap(); + let (_tx, rx) = crate::runtime::messaging::signal_channel(16); + + let mut exec = ExecutionLoop::new( + thread, + llm.clone(), + effects, + leases, + policy, + rx, + "test-user".into(), + ); + + let outcome = exec.run().await.unwrap(); + assert!(matches!(outcome, ThreadOutcome::Completed { response: Some(r) } if r == "done")); + + let seen = llm.seen_messages.lock().unwrap(); + let system_prompt = &seen[0][0].content; + assert!(system_prompt.contains("## Available capabilities (background status)")); + assert!(system_prompt.contains("`telegram`")); + assert!(system_prompt.contains("ready_scoped")); + assert!(system_prompt.contains("Usable through message")); + } + #[tokio::test] async fn action_then_text() { let (mut exec, _tx) = make_loop( diff --git a/crates/ironclaw_engine/src/executor/mod.rs b/crates/ironclaw_engine/src/executor/mod.rs index 718ebcf1b7a..8f688bc2ba3 100644 --- a/crates/ironclaw_engine/src/executor/mod.rs +++ b/crates/ironclaw_engine/src/executor/mod.rs @@ -11,6 +11,7 @@ pub mod orchestrator; pub mod prompt; pub mod scripting; pub mod structured; +pub(crate) mod thread_context; pub mod trace; pub use loop_engine::ExecutionLoop; diff --git a/crates/ironclaw_engine/src/executor/orchestrator.rs b/crates/ironclaw_engine/src/executor/orchestrator.rs index 3c173634491..95426fcc78e 100644 --- a/crates/ironclaw_engine/src/executor/orchestrator.rs +++ b/crates/ironclaw_engine/src/executor/orchestrator.rs @@ -30,6 +30,8 @@ use monty::{ }; use tracing::{debug, warn}; +use super::scripting::{execute_code, json_to_monty, monty_to_json, monty_to_string}; +use super::thread_context::thread_execution_context; use crate::capability::lease::LeaseManager; use crate::capability::policy::PolicyEngine; use crate::memory::RetrievalEngine; @@ -45,9 +47,6 @@ use crate::types::project::ProjectId; use crate::types::shared_owner_id; use crate::types::step::{ActionCall, StepId, TokenUsage}; use crate::types::thread::{ActiveSkillProvenance, Thread, ThreadState}; -use ironclaw_common::ValidTimezone; - -use super::scripting::{execute_code, json_to_monty, monty_to_json, monty_to_string}; /// The compiled-in default orchestrator (v0). pub(crate) const DEFAULT_ORCHESTRATOR: &str = include_str!("../../orchestrator/default.py"); @@ -66,24 +65,6 @@ pub struct OrchestratorResult { pub tokens_used: TokenUsage, } -/// Extract source_channel from thread metadata (set by ConversationManager). -fn thread_source_channel(thread: &Thread) -> Option { - thread - .metadata - .get("source_channel") - .and_then(|v| v.as_str()) - .map(String::from) -} - -/// Extract and validate user_timezone from thread metadata (set by bridge router). -fn thread_user_timezone(thread: &Thread) -> Option { - thread - .metadata - .get("user_timezone") - .and_then(|v| v.as_str()) - .and_then(ValidTimezone::parse) -} - fn normalize_pause_outcome( thread: &mut Thread, outcome: &ThreadOutcome, @@ -731,9 +712,10 @@ async fn handle_llm_complete( } let active_leases = deps.leases.active_for_thread(thread.id).await; + let actions_context = thread_execution_context(thread, StepId::new(), None); let actions = deps .effects - .available_actions(&active_leases) + .available_actions(&active_leases, &actions_context) .await .unwrap_or_default(); @@ -833,17 +815,7 @@ async fn handle_execute_code_step( .map(monty_to_json) .unwrap_or(serde_json::json!({})); - let exec_ctx = ThreadExecutionContext { - thread_id: thread.id, - thread_type: thread.thread_type, - project_id: thread.project_id, - user_id: thread.user_id.clone(), - step_id: StepId::new(), - current_call_id: None, - source_channel: thread_source_channel(thread), - user_timezone: thread_user_timezone(thread), - thread_goal: Some(thread.goal.clone()), - }; + let exec_ctx = thread_execution_context(thread, StepId::new(), None); // Run user code in a nested Monty VM (same pattern as rlm_query) let code_start = std::time::Instant::now(); @@ -1008,17 +980,7 @@ async fn handle_execute_action( let call_id = extract_string_kwarg(kwargs, "call_id").unwrap_or_default(); - let exec_ctx = ThreadExecutionContext { - thread_id: thread.id, - thread_type: thread.thread_type, - project_id: thread.project_id, - user_id: thread.user_id.clone(), - step_id: StepId::new(), - current_call_id: Some(call_id.clone()), - source_channel: thread_source_channel(thread), - user_timezone: thread_user_timezone(thread), - thread_goal: Some(thread.goal.clone()), - }; + let exec_ctx = thread_execution_context(thread, StepId::new(), Some(call_id.clone())); // Helper: emit event only. The orchestrator owns transcript recording. let emit_and_record = |thread: &mut Thread, @@ -1066,7 +1028,7 @@ async fn handle_execute_action( // 2. Check policy let action_def = effects - .available_actions(std::slice::from_ref(&lease)) + .available_actions(std::slice::from_ref(&lease), &exec_ctx) .await .ok() .and_then(|actions| actions.into_iter().find(|a| a.name == name)); @@ -1407,8 +1369,9 @@ async fn handle_execute_actions_parallel( }; // Check policy + let exec_ctx = thread_execution_context(thread, step_id, Some(pc.call_id.clone())); let action_def = effects - .available_actions(std::slice::from_ref(&lease)) + .available_actions(std::slice::from_ref(&lease), &exec_ctx) .await .ok() .and_then(|actions| actions.into_iter().find(|a| a.name == pc.name)); @@ -1567,22 +1530,7 @@ async fn handle_execute_actions_parallel( // Single call: execute directly let (idx, lease) = runnable.into_iter().next().unwrap(); // safety: len()==1 checked above let pc = &parsed[idx]; - let exec_ctx = ThreadExecutionContext { - thread_id: thread.id, - thread_type: thread.thread_type, - project_id: thread.project_id, - user_id: thread.user_id.clone(), - step_id, - current_call_id: Some(pc.call_id.clone()), - // Read source_channel from thread metadata so downstream tools - // (e.g. mission_create) can default notify_channels to the - // originating channel. Hardcoding `None` here was a bug — it - // silently dropped the gateway routing for any tool dispatched - // through the parallel batch path. - source_channel: thread_source_channel(thread), - user_timezone: thread_user_timezone(thread), - thread_goal: Some(thread.goal.clone()), - }; + let exec_ctx = thread_execution_context(thread, step_id, Some(pc.call_id.clone())); let ps = summarize_params(&pc.name, &pc.params); let (result_json, event, output) = execute_single_action( effects, @@ -1606,27 +1554,13 @@ async fn handle_execute_actions_parallel( let effects = effects.clone(); // Capture once outside the loop — the thread's metadata is stable // for the duration of the parallel batch. - let parallel_source_channel = thread_source_channel(thread); - let parallel_user_timezone = thread_user_timezone(thread); - for (idx, lease) in runnable { let pc_name = parsed[idx].name.clone(); let pc_params = parsed[idx].params.clone(); let pc_call_id = parsed[idx].call_id.clone(); let effects = effects.clone(); let lease = lease.clone(); - let exec_ctx = ThreadExecutionContext { - thread_id: thread.id, - thread_type: thread.thread_type, - project_id: thread.project_id, - user_id: thread.user_id.clone(), - step_id, - current_call_id: Some(pc_call_id.clone()), - // See comment above — read from thread metadata, not None. - source_channel: parallel_source_channel.clone(), - user_timezone: parallel_user_timezone, - thread_goal: Some(thread.goal.clone()), - }; + let exec_ctx = thread_execution_context(thread, step_id, Some(pc_call_id.clone())); let ps = summarize_params(&pc_name, &pc_params); join_set.spawn(async move { @@ -2057,7 +1991,11 @@ async fn handle_get_actions( } let active_leases = leases.active_for_thread(thread.id).await; - match effects.available_actions(&active_leases).await { + let actions_context = thread_execution_context(thread, StepId::new(), None); + match effects + .available_actions(&active_leases, &actions_context) + .await + { Ok(actions) => { let actions_json: Vec = actions .iter() @@ -3638,9 +3576,18 @@ mod tests { async fn available_actions( &self, _: &[crate::types::capability::CapabilityLease], + _: &ThreadExecutionContext, ) -> Result, EngineError> { Ok(vec![]) } + + async fn available_capabilities( + &self, + _: &[crate::types::capability::CapabilityLease], + _: &ThreadExecutionContext, + ) -> Result, EngineError> { + Ok(vec![]) + } } #[tokio::test] diff --git a/crates/ironclaw_engine/src/executor/prompt.rs b/crates/ironclaw_engine/src/executor/prompt.rs index f54c92376de..dc607bcaaaf 100644 --- a/crates/ironclaw_engine/src/executor/prompt.rs +++ b/crates/ironclaw_engine/src/executor/prompt.rs @@ -11,7 +11,9 @@ use std::sync::Arc; use crate::traits::store::Store; -use crate::types::capability::ActionDef; +use crate::types::capability::{ + ActionDef, CapabilityStatus, CapabilitySummary, CapabilitySummaryKind, +}; use crate::types::project::ProjectId; /// Runtime platform metadata injected into system prompts for self-awareness. @@ -102,6 +104,7 @@ const MAX_PROMPT_OVERLAY_CHARS: usize = 4000; /// mission to evolve the system prompt at runtime. pub async fn build_codeact_system_prompt( actions: &[ActionDef], + capabilities: &[CapabilitySummary], store: Option<&Arc>, project_id: ProjectId, platform: Option<&PlatformInfo>, @@ -111,7 +114,7 @@ pub async fn build_codeact_system_prompt( } else { None }; - build_codeact_system_prompt_inner(actions, overlay.as_deref(), platform) + build_codeact_system_prompt_inner(actions, capabilities, overlay.as_deref(), platform) } /// Build the system prompt using pre-fetched memory docs. @@ -121,16 +124,18 @@ pub async fn build_codeact_system_prompt( /// Store query. pub fn build_codeact_system_prompt_with_docs( actions: &[ActionDef], + capabilities: &[CapabilitySummary], system_docs: &[crate::types::memory::MemoryDoc], platform: Option<&PlatformInfo>, ) -> String { let overlay = extract_prompt_overlay(system_docs); - build_codeact_system_prompt_inner(actions, overlay.as_deref(), platform) + build_codeact_system_prompt_inner(actions, capabilities, overlay.as_deref(), platform) } /// Shared prompt builder used by both the async and pre-fetched-docs variants. fn build_codeact_system_prompt_inner( actions: &[ActionDef], + capabilities: &[CapabilitySummary], overlay: Option<&str>, platform: Option<&PlatformInfo>, ) -> String { @@ -163,10 +168,55 @@ fn build_codeact_system_prompt_inner( } } + if !capabilities.is_empty() { + prompt.push_str("\n## Available capabilities (background status)\n\n"); + for capability in capabilities { + prompt.push_str(&format!( + "- `{}` [{}] — {}", + capability.name, + capability_kind_label(capability.kind), + capability_status_label(capability.status) + )); + if let Some(display_name) = &capability.display_name + && display_name != &capability.name + { + prompt.push_str(&format!(" ({display_name})")); + } + if let Some(routing_hint) = &capability.routing_hint { + prompt.push_str(&format!(". {routing_hint}")); + } + if let Some(description) = &capability.description { + prompt.push_str(&format!(". {description}")); + } + prompt.push('\n'); + } + } + prompt.push_str(CODEACT_POSTAMBLE); prompt } +const fn capability_status_label(status: CapabilityStatus) -> &'static str { + match status { + CapabilityStatus::Ready => "ready", + CapabilityStatus::ReadyScoped => "ready_scoped", + CapabilityStatus::NeedsAuth => "needs_auth", + CapabilityStatus::NeedsSetup => "needs_setup", + CapabilityStatus::Inactive => "inactive", + CapabilityStatus::Latent => "latent", + CapabilityStatus::Error => "error", + CapabilityStatus::AvailableNotInstalled => "available_not_installed", + } +} + +const fn capability_kind_label(kind: CapabilitySummaryKind) -> &'static str { + match kind { + CapabilitySummaryKind::Channel => "channel", + CapabilitySummaryKind::Provider => "provider", + CapabilitySummaryKind::Runtime => "runtime", + } +} + /// Load the prompt overlay from the Store, if one exists for this project. async fn load_prompt_overlay(store: &Arc, project_id: ProjectId) -> Option { let docs = store.list_shared_memory_docs(project_id).await.ok()?; @@ -199,7 +249,7 @@ mod tests { #[tokio::test] async fn prompt_without_store_uses_compiled_preamble() { let prompt = - build_codeact_system_prompt(&[], None, ProjectId(uuid::Uuid::nil()), None).await; + build_codeact_system_prompt(&[], &[], None, ProjectId(uuid::Uuid::nil()), None).await; assert!(prompt.contains("Python REPL environment")); assert!(prompt.contains("Strategy")); assert!(!prompt.contains("Learned Rules")); @@ -223,9 +273,14 @@ mod tests { }; let store = Arc::new(crate::tests::InMemoryStore::with_docs(vec![overlay])); - let prompt = - build_codeact_system_prompt(&[], Some(&(store as Arc)), project_id, None) - .await; + let prompt = build_codeact_system_prompt( + &[], + &[], + Some(&(store as Arc)), + project_id, + None, + ) + .await; assert!(prompt.contains("Learned Rules")); assert!(prompt.contains("Never call web_fetch")); } @@ -251,9 +306,14 @@ mod tests { }; let store = Arc::new(crate::tests::InMemoryStore::with_docs(vec![overlay])); - let prompt = - build_codeact_system_prompt(&[], Some(&(store as Arc)), project_id, None) - .await; + let prompt = build_codeact_system_prompt( + &[], + &[], + Some(&(store as Arc)), + project_id, + None, + ) + .await; let snowman_count = prompt.chars().filter(|c| *c == '\u{2603}').count(); assert_eq!(snowman_count, MAX_PROMPT_OVERLAY_CHARS); @@ -278,9 +338,14 @@ mod tests { }; let store = Arc::new(crate::tests::InMemoryStore::with_docs(vec![overlay])); - let prompt = - build_codeact_system_prompt(&[], Some(&(store as Arc)), project_id, None) - .await; + let prompt = build_codeact_system_prompt( + &[], + &[], + Some(&(store as Arc)), + project_id, + None, + ) + .await; assert!(!prompt.contains("Should not appear")); assert!(!prompt.contains("Learned Rules")); } @@ -297,7 +362,8 @@ mod tests { repo_url: Some("https://github.com/nearai/ironclaw".into()), }; let prompt = - build_codeact_system_prompt(&[], None, ProjectId(uuid::Uuid::nil()), Some(&info)).await; + build_codeact_system_prompt(&[], &[], None, ProjectId(uuid::Uuid::nil()), Some(&info)) + .await; assert!(prompt.contains("IronClaw")); assert!(prompt.contains("1.2.3")); assert!(prompt.contains("nearai")); @@ -311,7 +377,41 @@ mod tests { #[tokio::test] async fn prompt_without_platform_info_has_no_platform_section() { let prompt = - build_codeact_system_prompt(&[], None, ProjectId(uuid::Uuid::nil()), None).await; + build_codeact_system_prompt(&[], &[], None, ProjectId(uuid::Uuid::nil()), None).await; assert!(!prompt.contains("## Platform")); } + + #[test] + fn prompt_with_capabilities_includes_background_statuses() { + let prompt = build_codeact_system_prompt_with_docs( + &[], + &[ + CapabilitySummary { + name: "telegram".into(), + display_name: Some("Telegram".into()), + kind: crate::types::capability::CapabilitySummaryKind::Channel, + status: CapabilityStatus::ReadyScoped, + description: Some("Telegram notifications".into()), + routing_hint: Some("Usable through message".into()), + }, + CapabilitySummary { + name: "slack".into(), + display_name: None, + kind: crate::types::capability::CapabilitySummaryKind::Provider, + status: CapabilityStatus::NeedsAuth, + description: Some("Slack workspace integration".into()), + routing_hint: None, + }, + ], + &[], + None, + ); + + assert!(prompt.contains("## Available capabilities (background status)")); + assert!(prompt.contains("`telegram` [channel]")); + assert!(prompt.contains("ready_scoped")); + assert!(prompt.contains("Usable through message")); + assert!(prompt.contains("`slack` [provider]")); + assert!(prompt.contains("needs_auth")); + } } diff --git a/crates/ironclaw_engine/src/executor/scripting.rs b/crates/ironclaw_engine/src/executor/scripting.rs index 0dcd72da70f..d56a46394a7 100644 --- a/crates/ironclaw_engine/src/executor/scripting.rs +++ b/crates/ironclaw_engine/src/executor/scripting.rs @@ -432,7 +432,7 @@ pub async fn execute_code_with_skills( // resolves the name before calling it, and Undefined → NameError. let active_leases = leases.active_for_thread(thread.id).await; let mut known_actions: std::collections::HashSet = effects - .available_actions(&active_leases) + .available_actions(&active_leases, context) .await .unwrap_or_default() .into_iter() @@ -1147,7 +1147,7 @@ async fn preflight_action( }; let action_def = effects - .available_actions(std::slice::from_ref(&lease)) + .available_actions(std::slice::from_ref(&lease), context) .await .ok() .and_then(|actions| actions.into_iter().find(|a| a.name == action_name)); @@ -1980,9 +1980,18 @@ mod tests { async fn available_actions( &self, _leases: &[CapabilityLease], + _context: &ThreadExecutionContext, ) -> Result, EngineError> { Ok(self.actions.clone()) } + + async fn available_capabilities( + &self, + _: &[CapabilityLease], + _: &ThreadExecutionContext, + ) -> Result, EngineError> { + Ok(vec![]) + } } fn test_action(name: &str) -> ActionDef { diff --git a/crates/ironclaw_engine/src/executor/structured.rs b/crates/ironclaw_engine/src/executor/structured.rs index 58b7e4fba82..18dc1fa5a8f 100644 --- a/crates/ironclaw_engine/src/executor/structured.rs +++ b/crates/ironclaw_engine/src/executor/structured.rs @@ -105,7 +105,7 @@ pub async fn execute_action_calls( // 2. Find the action definition and check policy let action_def = effects - .available_actions(std::slice::from_ref(&lease)) + .available_actions(std::slice::from_ref(&lease), context) .await? .into_iter() .find(|a| action_name_matches(&a.name, &call.action_name)); @@ -514,9 +514,18 @@ mod tests { async fn available_actions( &self, _leases: &[CapabilityLease], + _context: &ThreadExecutionContext, ) -> Result, EngineError> { Ok(self.actions.clone()) } + + async fn available_capabilities( + &self, + _: &[CapabilityLease], + _: &ThreadExecutionContext, + ) -> Result, EngineError> { + Ok(vec![]) + } } fn test_action(name: &str) -> ActionDef { diff --git a/crates/ironclaw_engine/src/executor/thread_context.rs b/crates/ironclaw_engine/src/executor/thread_context.rs new file mode 100644 index 00000000000..bf1a64c80f4 --- /dev/null +++ b/crates/ironclaw_engine/src/executor/thread_context.rs @@ -0,0 +1,31 @@ +use crate::traits::effect::ThreadExecutionContext; +use crate::types::step::StepId; +use crate::types::thread::Thread; +use ironclaw_common::ValidTimezone; + +/// Build an execution context from the current thread state. +pub(crate) fn thread_execution_context( + thread: &Thread, + step_id: StepId, + current_call_id: Option, +) -> ThreadExecutionContext { + ThreadExecutionContext { + thread_id: thread.id, + thread_type: thread.thread_type, + project_id: thread.project_id, + user_id: thread.user_id.clone(), + step_id, + current_call_id, + source_channel: thread + .metadata + .get("source_channel") + .and_then(|v| v.as_str()) + .map(str::to_string), + user_timezone: thread + .metadata + .get("user_timezone") + .and_then(|v| v.as_str()) + .and_then(ValidTimezone::parse), + thread_goal: Some(thread.goal.clone()), + } +} diff --git a/crates/ironclaw_engine/src/lib.rs b/crates/ironclaw_engine/src/lib.rs index fd2fc3b3080..d7154749f41 100644 --- a/crates/ironclaw_engine/src/lib.rs +++ b/crates/ironclaw_engine/src/lib.rs @@ -36,8 +36,9 @@ pub mod workspace; // ── Re-exports: types ─────────────────────────────────────── pub use types::capability::{ - ActionDef, Capability, CapabilityLease, CapabilityStatus, EffectType, GrantedActions, LeaseId, - PolicyCondition, PolicyEffect, PolicyRule, + ActionDef, Capability, CapabilityLease, CapabilityStatus, CapabilitySummary, + CapabilitySummaryKind, EffectType, GrantedActions, LeaseId, PolicyCondition, PolicyEffect, + PolicyRule, }; pub use types::error::{CapabilityError, EngineError, StepError, ThreadError}; pub use types::event::{EventId, EventKind, ThreadEvent}; diff --git a/crates/ironclaw_engine/src/runtime/conversation.rs b/crates/ironclaw_engine/src/runtime/conversation.rs index 14d43901022..74b42e90be2 100644 --- a/crates/ironclaw_engine/src/runtime/conversation.rs +++ b/crates/ironclaw_engine/src/runtime/conversation.rs @@ -636,9 +636,18 @@ mod tests { async fn available_actions( &self, _: &[CapabilityLease], + _: &crate::traits::effect::ThreadExecutionContext, ) -> Result, EngineError> { Ok(vec![]) } + + async fn available_capabilities( + &self, + _: &[CapabilityLease], + _: &crate::traits::effect::ThreadExecutionContext, + ) -> Result, EngineError> { + Ok(vec![]) + } } struct MockStore { diff --git a/crates/ironclaw_engine/src/runtime/lease_refresh.rs b/crates/ironclaw_engine/src/runtime/lease_refresh.rs index a3ffd6421d6..324cf859830 100644 --- a/crates/ironclaw_engine/src/runtime/lease_refresh.rs +++ b/crates/ironclaw_engine/src/runtime/lease_refresh.rs @@ -9,6 +9,7 @@ use crate::traits::effect::EffectExecutor; use crate::traits::store::Store; use crate::types::capability::GrantedActions; use crate::types::error::EngineError; +use crate::types::step::StepId; use crate::types::thread::Thread; pub(crate) async fn reconcile_dynamic_tool_lease( @@ -19,7 +20,9 @@ pub(crate) async fn reconcile_dynamic_tool_lease( lease_planner: &LeasePlanner, ) -> Result<(), EngineError> { let active_leases = leases.active_for_thread(thread.id).await; - let actions = effects.available_actions(&active_leases).await?; + let context = + crate::executor::thread_context::thread_execution_context(thread, StepId::new(), None); + let actions = effects.available_actions(&active_leases, &context).await?; if actions.is_empty() { return Ok(()); } diff --git a/crates/ironclaw_engine/src/runtime/manager.rs b/crates/ironclaw_engine/src/runtime/manager.rs index 9dc811ff87e..3a900ae85b7 100644 --- a/crates/ironclaw_engine/src/runtime/manager.rs +++ b/crates/ironclaw_engine/src/runtime/manager.rs @@ -778,9 +778,18 @@ mod tests { async fn available_actions( &self, _: &[CapabilityLease], + _: &crate::traits::effect::ThreadExecutionContext, ) -> Result, EngineError> { Ok(vec![]) } + + async fn available_capabilities( + &self, + _: &[CapabilityLease], + _: &crate::traits::effect::ThreadExecutionContext, + ) -> Result, EngineError> { + Ok(vec![]) + } } #[async_trait::async_trait] @@ -810,9 +819,18 @@ mod tests { async fn available_actions( &self, _: &[CapabilityLease], + _: &crate::traits::effect::ThreadExecutionContext, ) -> Result, EngineError> { Ok(self.actions.read().await.clone()) } + + async fn available_capabilities( + &self, + _: &[CapabilityLease], + _: &crate::traits::effect::ThreadExecutionContext, + ) -> Result, EngineError> { + Ok(vec![]) + } } struct MockStore { diff --git a/crates/ironclaw_engine/src/runtime/mission.rs b/crates/ironclaw_engine/src/runtime/mission.rs index 56834d1f78b..7b4f15f0492 100644 --- a/crates/ironclaw_engine/src/runtime/mission.rs +++ b/crates/ironclaw_engine/src/runtime/mission.rs @@ -3446,9 +3446,18 @@ mod tests { async fn available_actions( &self, _: &[CapabilityLease], + _: &crate::traits::effect::ThreadExecutionContext, ) -> Result, EngineError> { Ok(vec![]) } + + async fn available_capabilities( + &self, + _: &[CapabilityLease], + _: &crate::traits::effect::ThreadExecutionContext, + ) -> Result, EngineError> { + Ok(vec![]) + } } // ── Helper to build a MissionManager with its dependencies ── @@ -4317,9 +4326,18 @@ mod tests { async fn available_actions( &self, _: &[CapabilityLease], + _: &crate::traits::effect::ThreadExecutionContext, ) -> Result, EngineError> { Ok(vec![]) } + + async fn available_capabilities( + &self, + _: &[CapabilityLease], + _: &crate::traits::effect::ThreadExecutionContext, + ) -> Result, EngineError> { + Ok(vec![]) + } } #[tokio::test] diff --git a/crates/ironclaw_engine/src/traits/effect.rs b/crates/ironclaw_engine/src/traits/effect.rs index ae9f3761c8e..693a30fc24f 100644 --- a/crates/ironclaw_engine/src/traits/effect.rs +++ b/crates/ironclaw_engine/src/traits/effect.rs @@ -4,7 +4,7 @@ //! trait. The main crate implements it by wrapping `ToolRegistry` and //! `SafetyLayer` — the engine itself has no knowledge of specific tools. -use crate::types::capability::{ActionDef, CapabilityLease}; +use crate::types::capability::{ActionDef, CapabilityLease, CapabilitySummary}; use crate::types::error::EngineError; use crate::types::project::ProjectId; use crate::types::step::{ActionResult, StepId}; @@ -65,5 +65,13 @@ pub trait EffectExecutor: Send + Sync { async fn available_actions( &self, leases: &[CapabilityLease], + context: &ThreadExecutionContext, ) -> Result, EngineError>; + + /// List capability background summaries given the current runtime state. + async fn available_capabilities( + &self, + leases: &[CapabilityLease], + context: &ThreadExecutionContext, + ) -> Result, EngineError>; } diff --git a/crates/ironclaw_engine/src/types/capability.rs b/crates/ironclaw_engine/src/types/capability.rs index 28f39b1b75e..f7e30e85012 100644 --- a/crates/ironclaw_engine/src/types/capability.rs +++ b/crates/ironclaw_engine/src/types/capability.rs @@ -160,6 +160,38 @@ pub enum CapabilityStatus { AvailableNotInstalled, } +/// High-level category for capability background summaries. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum CapabilitySummaryKind { + /// Messaging or notification route usable through a bridge action. + Channel, + /// Extension-backed provider or integration. + Provider, + /// Engine-native runtime capability background. + Runtime, +} + +/// Background summary for a non-callable or indirectly-callable capability. +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct CapabilitySummary { + /// Stable capability identifier (for example `telegram` or `slack`). + pub name: String, + /// Human-readable display name when available. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub display_name: Option, + /// High-level category used by prompt/UI renderers. + pub kind: CapabilitySummaryKind, + /// Canonical normalized status. + pub status: CapabilityStatus, + /// Optional human-readable description. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub description: Option, + /// Optional routing guidance such as `Usable through message`. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub routing_hint: Option, +} + // ── Capability ────────────────────────────────────────────── /// A capability — bundles actions, knowledge, and policies. @@ -428,4 +460,68 @@ mod tests { assert_eq!(parsed, expected); } } + + #[test] + fn capability_summary_kind_serializes_as_snake_case() { + let cases = [ + (CapabilitySummaryKind::Channel, json!("channel")), + (CapabilitySummaryKind::Provider, json!("provider")), + (CapabilitySummaryKind::Runtime, json!("runtime")), + ]; + + for (kind, expected) in cases { + assert_eq!(serde_json::to_value(kind).unwrap(), expected); + } + } + + #[test] + fn capability_summary_round_trips_with_optional_fields() { + let summary = CapabilitySummary { + name: "telegram".to_string(), + display_name: Some("Telegram".to_string()), + kind: CapabilitySummaryKind::Channel, + status: CapabilityStatus::ReadyScoped, + description: Some("Telegram messaging".to_string()), + routing_hint: Some("Usable through message".to_string()), + }; + + let json = serde_json::to_value(&summary).unwrap(); + assert_eq!(json["name"], "telegram"); + assert_eq!(json["display_name"], "Telegram"); + assert_eq!(json["kind"], "channel"); + assert_eq!(json["status"], "ready_scoped"); + assert_eq!(json["description"], "Telegram messaging"); + assert_eq!(json["routing_hint"], "Usable through message"); + + let parsed: CapabilitySummary = serde_json::from_value(json).unwrap(); + assert_eq!(parsed, summary); + } + + #[test] + fn capability_summary_allows_minimal_payload_and_omits_none_fields() { + let summary = CapabilitySummary { + name: "notion".to_string(), + display_name: None, + kind: CapabilitySummaryKind::Provider, + status: CapabilityStatus::NeedsAuth, + description: None, + routing_hint: None, + }; + + let json = serde_json::to_value(&summary).unwrap(); + assert_eq!(json["name"], "notion"); + assert_eq!(json["kind"], "provider"); + assert_eq!(json["status"], "needs_auth"); + assert!(json.get("display_name").is_none()); + assert!(json.get("description").is_none()); + assert!(json.get("routing_hint").is_none()); + + let parsed: CapabilitySummary = serde_json::from_value(serde_json::json!({ + "name": "notion", + "kind": "provider", + "status": "needs_auth" + })) + .unwrap(); + assert_eq!(parsed, summary); + } } diff --git a/deny.toml b/deny.toml index b5d2afabaec..ce6a35d5b01 100644 --- a/deny.toml +++ b/deny.toml @@ -63,4 +63,3 @@ allow-git = [ "https://github.com/pydantic/monty.git", "https://github.com/astral-sh/ruff.git", ] - diff --git a/src/auth/extension.rs b/src/auth/extension.rs index fef7bbba783..477d158b762 100644 --- a/src/auth/extension.rs +++ b/src/auth/extension.rs @@ -17,7 +17,9 @@ use crate::auth::{ upsert_auth_descriptor, }; use crate::extensions::naming::canonicalize_extension_name; -use crate::extensions::{ConfigureResult, ExtensionError}; +use crate::extensions::{ + ConfigureResult, ExtensionError, InstalledExtension, LatentProviderAction, +}; use crate::secrets::SecretsStore; use crate::tools::ToolRegistry; use crate::tools::builtin::extract_host_from_params; @@ -519,6 +521,33 @@ impl AuthManager { .collect() } + /// Thin bridge accessor for capability projection inventory. + pub(crate) async fn list_capability_extensions( + &self, + user_id: &str, + ) -> Result, ExtensionError> { + let Some(ext_mgr) = self.extension_manager.as_ref() else { + return Ok(Vec::new()); + }; + + ext_mgr.list(None, true, user_id).await + } + + /// Thin bridge accessor for user-scoped latent provider projection. + pub(crate) async fn latent_provider_actions(&self, user_id: &str) -> Vec { + let Some(ext_mgr) = self.extension_manager.as_ref() else { + return Vec::new(); + }; + + ext_mgr.latent_provider_actions(user_id).await + } + + /// Thin bridge accessor for routed channel capability hints. + pub(crate) async fn notification_target_for_channel(&self, name: &str) -> Option { + let ext_mgr = self.extension_manager.as_ref()?; + ext_mgr.notification_target_for_channel(name).await + } + pub async fn execute_latent_extension_action( &self, action_name: &str, diff --git a/src/bridge/action_projector.rs b/src/bridge/action_projector.rs new file mode 100644 index 00000000000..4c8fd2133e3 --- /dev/null +++ b/src/bridge/action_projector.rs @@ -0,0 +1,390 @@ +use std::collections::{HashMap, HashSet}; +use std::sync::Arc; + +use tracing::debug; + +use ironclaw_engine::{ + ActionDef, CapabilityLease, CapabilityRegistry, CapabilityStatus, EngineError, + ThreadExecutionContext, +}; + +use crate::auth::extension::AuthManager; +use crate::bridge::capability_projector::{ + capability_status_for_extension, capability_surface_subject_for_extension, +}; +use crate::bridge::tool_surface::{ + InvocationMode, SurfacePolicyInput, SurfaceSubjectKind, assign_surface, +}; +use crate::extensions::InstalledExtension; +use crate::extensions::naming::extension_name_candidates; +use crate::tools::ToolRegistry; + +pub(crate) struct ActionProjector; + +impl ActionProjector { + /// Project the set of available actions from the tool registry and + /// capability registry. + /// + /// When `prefetched_extensions` is `Some`, the projector uses that map + /// instead of fetching from `auth_manager`. This allows the caller + /// (typically `EffectBridgeAdapter`) to share a single fetch across + /// both `ActionProjector` and `CapabilityProjector`. + pub(crate) async fn project( + tools: &ToolRegistry, + auth_manager: Option<&AuthManager>, + capability_registry: Option>, + leases: &[CapabilityLease], + context: &ThreadExecutionContext, + prefetched_extensions: Option<&HashMap>, + ) -> Result, EngineError> { + let tool_defs = tools.tool_definitions().await; + let owned_statuses; + let extension_statuses: Option<&HashMap> = if let Some( + prefetched, + ) = + prefetched_extensions + { + Some(prefetched) + } else if let Some(auth_manager) = auth_manager { + match auth_manager + .list_capability_extensions(&context.user_id) + .await + { + Ok(extensions) => { + owned_statuses = extensions + .into_iter() + .map(|extension| (extension.name.clone(), extension)) + .collect::>(); + Some(&owned_statuses) + } + Err(error) => { + debug!( + user_id = %context.user_id, + error = %error, + "failed to load extension inventory for available_actions; omitting extension-backed actions" + ); + owned_statuses = HashMap::new(); + Some(&owned_statuses) + } + } + } else { + None + }; + + let mut actions = Vec::with_capacity(tool_defs.len()); + for td in tool_defs { + if crate::bridge::effect_adapter::is_v1_only_tool(&td.name) { + continue; + } + if crate::bridge::effect_adapter::is_v1_auth_tool(&td.name) { + continue; + } + + if let Some(provider_extension) = tools.provider_extension_for_tool(&td.name).await { + let Some(extension_statuses) = extension_statuses else { + continue; + }; + let Some(extension) = + provider_extension_status(extension_statuses, &provider_extension) + else { + continue; + }; + let status = capability_status_for_extension(extension, false); + // NeedsAuth tools remain in available_actions so the LLM can + // trigger the auth-on-first-call gate by attempting to call them. + if status == CapabilityStatus::NeedsAuth { + // Keep in actions — execution will trigger the auth gate. + } else { + let (kind, invocation_mode) = + capability_surface_subject_for_extension(extension); + let assignment = assign_surface(SurfacePolicyInput { + kind, + status, + invocation_mode, + leased_and_callable: false, + }); + if !assignment.available_actions { + continue; + } + } + } + + actions.push(ActionDef { + name: td.name.replace('-', "_"), + description: td.description, + parameters_schema: td.parameters, + effects: vec![], + requires_approval: false, + }); + } + + if let Some(registry) = capability_registry.as_ref() { + let mut seen: HashSet = actions.iter().map(|a| a.name.clone()).collect(); + for lease in leases { + if lease.capability_name == "tools" { + continue; + } + let Some(cap) = registry.get(&lease.capability_name) else { + continue; + }; + for action in &cap.actions { + if !lease.granted_actions.covers(&action.name) { + continue; + } + if crate::bridge::effect_adapter::is_v1_only_tool(&action.name) + || crate::bridge::effect_adapter::is_v1_auth_tool(&action.name) + { + continue; + } + let assignment = assign_surface(SurfacePolicyInput { + kind: SurfaceSubjectKind::EngineNativeDirectAction, + status: CapabilityStatus::Ready, + invocation_mode: InvocationMode::Direct, + leased_and_callable: true, + }); + if !assignment.available_actions || !seen.insert(action.name.clone()) { + continue; + } + actions.push(action.clone()); + } + } + } + + actions.sort_by(|a, b| a.name.cmp(&b.name)); + Ok(actions) + } +} + +fn provider_extension_status<'a>( + extension_statuses: &'a HashMap, + provider_extension: &str, +) -> Option<&'a InstalledExtension> { + extension_name_candidates(provider_extension) + .into_iter() + .filter_map(|candidate| extension_statuses.get(&candidate)) + .max_by_key(|extension| provider_extension_rank(extension)) +} + +fn provider_extension_rank(extension: &InstalledExtension) -> u8 { + match capability_status_for_extension(extension, false) { + CapabilityStatus::Ready => 5, + CapabilityStatus::Inactive => 4, + CapabilityStatus::NeedsAuth => 3, + CapabilityStatus::NeedsSetup => 2, + CapabilityStatus::Error => 1, + CapabilityStatus::AvailableNotInstalled => 0, + CapabilityStatus::ReadyScoped | CapabilityStatus::Latent => 0, + } +} + +#[cfg(test)] +mod tests { + use std::collections::HashMap; + + use super::{ActionProjector, provider_extension_status}; + use crate::extensions::{ExtensionKind, InstalledExtension}; + use crate::tools::ToolRegistry; + + fn installed_extension(name: &str) -> InstalledExtension { + InstalledExtension { + name: name.to_string(), + kind: ExtensionKind::McpServer, + display_name: Some(name.to_string()), + description: Some(format!("{name} description")), + url: None, + authenticated: true, + active: true, + tools: vec![format!("{name}_search")], + needs_setup: false, + has_auth: true, + installed: true, + activation_error: None, + version: None, + } + } + + fn needs_auth_extension(name: &str) -> InstalledExtension { + InstalledExtension { + authenticated: false, + ..installed_extension(name) + } + } + + #[test] + fn provider_extension_lookup_accepts_legacy_hyphen_alias() { + let extension = installed_extension("linear-server"); + let statuses = HashMap::from([(extension.name.clone(), extension)]); + + let resolved = provider_extension_status(&statuses, "linear_server") + .expect("legacy hyphen alias should resolve"); + + assert_eq!(resolved.name, "linear-server"); + } + + #[test] + fn provider_extension_lookup_prefers_installed_alias_over_registry_only_entry() { + let installed = installed_extension("linear-server"); + let registry_only = InstalledExtension { + installed: false, + active: false, + authenticated: false, + has_auth: true, + tools: Vec::new(), + ..installed_extension("linear_server") + }; + let statuses = HashMap::from([ + (installed.name.clone(), installed), + (registry_only.name.clone(), registry_only), + ]); + + let resolved = provider_extension_status(&statuses, "linear_server") + .expect("installed alias should win over registry-only canonical entry"); + + assert_eq!(resolved.name, "linear-server"); + assert!(resolved.installed); + } + + /// NeedsAuth provider tools must remain in available_actions so the LLM + /// can trigger the auth-on-first-call gate by attempting to call them. + #[tokio::test] + async fn needs_auth_provider_tools_remain_in_available_actions() { + use async_trait::async_trait; + + struct GmailSendTool; + + #[async_trait] + impl crate::tools::Tool for GmailSendTool { + fn name(&self) -> &str { + "gmail_send" + } + fn description(&self) -> &str { + "Send a Gmail message" + } + fn parameters_schema(&self) -> serde_json::Value { + serde_json::json!({"type": "object"}) + } + async fn execute( + &self, + _: serde_json::Value, + _: &crate::context::JobContext, + ) -> Result { + Ok(crate::tools::ToolOutput::success( + serde_json::json!({}), + std::time::Duration::from_millis(1), + )) + } + fn provider_extension(&self) -> Option<&str> { + Some("gmail") + } + } + + let tools = std::sync::Arc::new(ToolRegistry::new()); + tools.register(std::sync::Arc::new(GmailSendTool)).await; + + let ext = needs_auth_extension("gmail"); + let extension_map = HashMap::from([(ext.name.clone(), ext)]); + + let context = ironclaw_engine::ThreadExecutionContext { + thread_id: ironclaw_engine::ThreadId::new(), + thread_type: ironclaw_engine::types::thread::ThreadType::Foreground, + project_id: ironclaw_engine::ProjectId::new(), + user_id: "test_user".to_string(), + step_id: ironclaw_engine::StepId::new(), + current_call_id: None, + source_channel: None, + user_timezone: None, + thread_goal: None, + }; + + let actions = ActionProjector::project( + tools.as_ref(), + None, + None, + &[], + &context, + Some(&extension_map), + ) + .await + .expect("project should succeed"); + + assert!( + actions.iter().any(|a| a.name == "gmail_send"), + "NeedsAuth provider tool should remain in available_actions, got: {:?}", + actions.iter().map(|a| &a.name).collect::>() + ); + } + + /// Latent/not-installed provider tools should still be omitted. + #[tokio::test] + async fn latent_provider_tools_omitted_from_available_actions() { + use async_trait::async_trait; + + struct LatentTool; + + #[async_trait] + impl crate::tools::Tool for LatentTool { + fn name(&self) -> &str { + "latent_send" + } + fn description(&self) -> &str { + "Send via latent provider" + } + fn parameters_schema(&self) -> serde_json::Value { + serde_json::json!({"type": "object"}) + } + async fn execute( + &self, + _: serde_json::Value, + _: &crate::context::JobContext, + ) -> Result { + Ok(crate::tools::ToolOutput::success( + serde_json::json!({}), + std::time::Duration::from_millis(1), + )) + } + fn provider_extension(&self) -> Option<&str> { + Some("latent_provider") + } + } + + let tools = std::sync::Arc::new(ToolRegistry::new()); + tools.register(std::sync::Arc::new(LatentTool)).await; + + // Extension is not installed — only in registry. + let ext = InstalledExtension { + installed: false, + active: false, + authenticated: false, + ..installed_extension("latent_provider") + }; + let extension_map = HashMap::from([(ext.name.clone(), ext)]); + + let context = ironclaw_engine::ThreadExecutionContext { + thread_id: ironclaw_engine::ThreadId::new(), + thread_type: ironclaw_engine::types::thread::ThreadType::Foreground, + project_id: ironclaw_engine::ProjectId::new(), + user_id: "test_user".to_string(), + step_id: ironclaw_engine::StepId::new(), + current_call_id: None, + source_channel: None, + user_timezone: None, + thread_goal: None, + }; + + let actions = ActionProjector::project( + tools.as_ref(), + None, + None, + &[], + &context, + Some(&extension_map), + ) + .await + .expect("project should succeed"); + + assert!( + !actions.iter().any(|a| a.name == "latent_send"), + "Not-installed provider tool should be omitted from available_actions" + ); + } +} diff --git a/src/bridge/capability_projector.rs b/src/bridge/capability_projector.rs new file mode 100644 index 00000000000..6803c87e04c --- /dev/null +++ b/src/bridge/capability_projector.rs @@ -0,0 +1,463 @@ +use std::collections::{BTreeMap, HashMap, HashSet}; + +use ironclaw_engine::{ + CapabilityLease, CapabilityStatus, CapabilitySummary, CapabilitySummaryKind, EngineError, + ThreadExecutionContext, +}; + +use crate::auth::extension::AuthManager; +use crate::bridge::tool_surface::{ + InvocationMode, SurfacePolicyInput, SurfaceSubjectKind, assign_surface, +}; +use crate::extensions::naming::extension_name_candidates; +use crate::extensions::{ExtensionKind, InstalledExtension, LatentProviderAction}; + +pub(crate) struct CapabilityProjector; + +struct CapabilityRuntimeSnapshot { + extensions: Vec, + latent_actions: Vec, + channel_routes: HashMap, +} + +impl CapabilityProjector { + /// Project the set of capability summaries from the runtime extension state. + /// + /// When `prefetched_extensions` is `Some`, the projector uses that list + /// instead of fetching from `auth_manager`. This allows the caller + /// (typically `EffectBridgeAdapter`) to share a single fetch across + /// both `ActionProjector` and `CapabilityProjector`. + pub(crate) async fn project( + auth_manager: Option<&AuthManager>, + leases: &[CapabilityLease], + context: &ThreadExecutionContext, + prefetched_extensions: Option>, + ) -> Result, EngineError> { + let Some(auth_manager) = auth_manager else { + return Ok(Vec::new()); + }; + + let extensions = if let Some(prefetched) = prefetched_extensions { + prefetched + } else { + auth_manager + .list_capability_extensions(&context.user_id) + .await + .map_err(|error| EngineError::Effect { + reason: format!("Failed to list extensions for capability projection: {error}"), + })? + }; + + let latent_actions = auth_manager.latent_provider_actions(&context.user_id).await; + let mut channel_routes = HashMap::new(); + let channel_route_lookups = extensions + .iter() + .filter(|extension| is_channel_extension_kind(extension.kind)) + .map(|extension| { + let name = extension.name.clone(); + async { + let target = auth_manager.notification_target_for_channel(&name).await; + (name, target) + } + }); + + for (name, target) in futures::future::join_all(channel_route_lookups).await { + if let Some(target) = target { + channel_routes.insert(name, target); + } + } + + Ok(Self::project_snapshot( + CapabilityRuntimeSnapshot { + extensions, + latent_actions, + channel_routes, + }, + leases, + )) + } + + fn project_snapshot( + snapshot: CapabilityRuntimeSnapshot, + _leases: &[CapabilityLease], + ) -> Vec { + let mut summaries = BTreeMap::::new(); + let mut registry_only = Vec::new(); + let mut installed_keys = HashSet::new(); + + for extension in snapshot.extensions { + if extension.installed { + installed_keys.insert(normalized_capability_key(&extension.name)); + if let Some(summary) = summarize_extension(&extension, &snapshot.channel_routes) { + upsert_summary(&mut summaries, summary, ProjectionSource::InstalledRuntime); + } + } else { + registry_only.push(extension); + } + } + + for latent in unique_latent_providers(snapshot.latent_actions) { + let normalized_key = normalized_capability_key(&latent.provider_extension); + if installed_keys.contains(&normalized_key) || summaries.contains_key(&normalized_key) { + continue; + } + + let assignment = assign_surface(SurfacePolicyInput { + kind: SurfaceSubjectKind::LatentProviderAction, + status: CapabilityStatus::Latent, + invocation_mode: InvocationMode::Direct, + + leased_and_callable: false, + }); + if !assignment.available_capabilities { + continue; + } + + upsert_summary( + &mut summaries, + CapabilitySummary { + name: latent.provider_extension.clone(), + display_name: None, + kind: CapabilitySummaryKind::Provider, + status: CapabilityStatus::Latent, + description: Some(latent.description), + routing_hint: None, + }, + ProjectionSource::LatentRuntime, + ); + } + + for extension in registry_only { + let normalized_key = normalized_capability_key(&extension.name); + if installed_keys.contains(&normalized_key) || summaries.contains_key(&normalized_key) { + continue; + } + + if let Some(summary) = summarize_extension(&extension, &snapshot.channel_routes) { + upsert_summary(&mut summaries, summary, ProjectionSource::RegistryOnly); + } + } + + summaries + .into_values() + .map(|summary| summary.summary) + .collect() + } +} + +#[derive(Clone, Copy, PartialEq, Eq, PartialOrd, Ord)] +enum ProjectionSource { + RegistryOnly, + LatentRuntime, + InstalledRuntime, +} + +struct PrioritizedSummary { + source: ProjectionSource, + summary: CapabilitySummary, +} + +fn upsert_summary( + summaries: &mut BTreeMap, + summary: CapabilitySummary, + source: ProjectionSource, +) { + let key = normalized_capability_key(&summary.name); + match summaries.get(&key) { + Some(existing) if existing.source >= source => {} + _ => { + summaries.insert(key, PrioritizedSummary { source, summary }); + } + } +} + +fn normalized_capability_key(name: &str) -> String { + extension_name_candidates(name) + .into_iter() + .next() + .unwrap_or_else(|| name.to_string()) +} + +fn summarize_extension( + extension: &InstalledExtension, + channel_routes: &HashMap, +) -> Option { + let status = + capability_status_for_extension(extension, channel_routes.contains_key(&extension.name)); + let (subject_kind, invocation_mode) = capability_surface_subject_for_extension(extension); + let assignment = assign_surface(SurfacePolicyInput { + kind: subject_kind, + status, + invocation_mode, + leased_and_callable: false, + }); + if !assignment.available_capabilities { + return None; + } + + let routing_hint = if matches!(status, CapabilityStatus::ReadyScoped) { + Some("Usable through message".to_string()) + } else { + None + }; + + Some(CapabilitySummary { + name: extension.name.clone(), + display_name: extension.display_name.clone(), + kind: if is_channel_extension_kind(extension.kind) { + CapabilitySummaryKind::Channel + } else { + CapabilitySummaryKind::Provider + }, + status, + description: extension.description.clone(), + routing_hint, + }) +} + +pub(crate) fn capability_surface_subject_for_extension( + extension: &InstalledExtension, +) -> (SurfaceSubjectKind, InvocationMode) { + if is_channel_extension_kind(extension.kind) { + return (SurfaceSubjectKind::Channel, InvocationMode::RoutedOnly); + } + + if extension.installed { + ( + SurfaceSubjectKind::ExtensionDirectAction, + InvocationMode::Direct, + ) + } else { + ( + SurfaceSubjectKind::AvailableNotInstalledProviderEntry, + InvocationMode::Direct, + ) + } +} + +pub(crate) fn capability_status_for_extension( + extension: &InstalledExtension, + route_exists: bool, +) -> CapabilityStatus { + if !extension.installed { + return CapabilityStatus::AvailableNotInstalled; + } + if extension.activation_error.is_some() { + return CapabilityStatus::Error; + } + if extension.needs_setup { + return CapabilityStatus::NeedsSetup; + } + if extension.has_auth && !extension.authenticated { + return CapabilityStatus::NeedsAuth; + } + if !extension.active { + return CapabilityStatus::Inactive; + } + if is_channel_extension_kind(extension.kind) { + return if route_exists { + CapabilityStatus::ReadyScoped + } else { + CapabilityStatus::Inactive + }; + } + + CapabilityStatus::Ready +} + +fn unique_latent_providers( + latent_actions: Vec, +) -> impl Iterator { + let mut providers = BTreeMap::::new(); + for latent in latent_actions { + providers + .entry(latent.provider_extension.clone()) + .or_insert(latent); + } + providers.into_values() +} + +pub(crate) const fn is_channel_extension_kind(kind: ExtensionKind) -> bool { + matches!( + kind, + ExtensionKind::WasmChannel | ExtensionKind::ChannelRelay + ) +} + +#[cfg(test)] +mod tests { + use std::collections::HashMap; + + use ironclaw_engine::{ + CapabilityLease, CapabilityStatus, CapabilitySummaryKind, GrantedActions, + }; + + use super::{CapabilityProjector, CapabilityRuntimeSnapshot}; + use crate::extensions::{ExtensionKind, InstalledExtension, LatentProviderAction}; + + fn make_lease() -> CapabilityLease { + CapabilityLease { + id: ironclaw_engine::LeaseId::new(), + thread_id: ironclaw_engine::ThreadId::new(), + capability_name: "tools".to_string(), + granted_actions: GrantedActions::All, + granted_at: chrono::Utc::now(), + expires_at: None, + max_uses: None, + uses_remaining: None, + revoked: false, + revoked_reason: None, + } + } + + fn installed_extension(name: &str, kind: ExtensionKind) -> InstalledExtension { + InstalledExtension { + name: name.to_string(), + kind, + display_name: Some(name.to_string()), + description: Some(format!("{name} description")), + url: None, + authenticated: true, + active: true, + tools: vec![format!("{name}_send")], + needs_setup: false, + has_auth: true, + installed: true, + activation_error: None, + version: None, + } + } + + fn available_extension(name: &str) -> InstalledExtension { + InstalledExtension { + installed: false, + active: false, + authenticated: false, + has_auth: true, + tools: Vec::new(), + ..installed_extension(name, ExtensionKind::WasmTool) + } + } + + fn latent_action(provider_extension: &str) -> LatentProviderAction { + LatentProviderAction { + action_name: format!("{provider_extension}_send"), + provider_extension: provider_extension.to_string(), + description: format!("{provider_extension} latent action"), + parameters_schema: serde_json::json!({"type": "object"}), + } + } + + #[test] + fn projects_normalized_capability_statuses() { + let mut telegram = installed_extension("telegram", ExtensionKind::WasmChannel); + telegram.tools.clear(); + + let mut slack = installed_extension("slack", ExtensionKind::ChannelRelay); + slack.authenticated = false; + + let mut github = installed_extension("github", ExtensionKind::McpServer); + github.needs_setup = true; + github.authenticated = false; + + let mut notion = installed_extension("notion", ExtensionKind::WasmTool); + notion.active = false; + + let mut broken = installed_extension("broken", ExtensionKind::WasmChannel); + broken.active = false; + broken.activation_error = Some("activation failed".to_string()); + + let unpaired = installed_extension("discord", ExtensionKind::WasmChannel); + + let snapshot = CapabilityRuntimeSnapshot { + extensions: vec![ + telegram, + slack, + github, + notion, + broken, + unpaired, + available_extension("linear"), + ], + latent_actions: vec![latent_action("gmail")], + channel_routes: HashMap::from([("telegram".to_string(), "actor".to_string())]), + }; + + let projected = CapabilityProjector::project_snapshot(snapshot, &[make_lease()]); + let by_name = projected + .into_iter() + .map(|summary| (summary.name.clone(), summary)) + .collect::>(); + + assert_eq!(by_name["telegram"].kind, CapabilitySummaryKind::Channel); + assert_eq!(by_name["telegram"].status, CapabilityStatus::ReadyScoped); + assert_eq!( + by_name["telegram"].routing_hint.as_deref(), + Some("Usable through message") + ); + assert_eq!(by_name["slack"].status, CapabilityStatus::NeedsAuth); + assert_eq!(by_name["github"].status, CapabilityStatus::NeedsSetup); + assert_eq!(by_name["notion"].status, CapabilityStatus::Inactive); + assert_eq!(by_name["gmail"].status, CapabilityStatus::Latent); + assert_eq!( + by_name["linear"].status, + CapabilityStatus::AvailableNotInstalled + ); + // "broken" has Error status which is excluded from all surfaces + // (early return in assign_surface), so it should not appear. + assert!(!by_name.contains_key("broken")); + assert_eq!(by_name["discord"].status, CapabilityStatus::Inactive); + assert_eq!(by_name["discord"].routing_hint, None); + } + + #[test] + fn installed_and_latent_entries_beat_registry_only_duplicates() { + let mut inactive = installed_extension("gmail", ExtensionKind::WasmTool); + inactive.active = false; + + let snapshot = CapabilityRuntimeSnapshot { + extensions: vec![inactive, available_extension("gmail")], + latent_actions: vec![latent_action("gmail")], + channel_routes: HashMap::new(), + }; + + let projected = CapabilityProjector::project_snapshot(snapshot, &[]); + assert_eq!(projected.len(), 1); + assert_eq!(projected[0].name, "gmail"); + assert_eq!(projected[0].status, CapabilityStatus::Inactive); + } + + #[test] + fn installed_alias_suppresses_registry_only_canonical_duplicate() { + let installed = installed_extension("linear-server", ExtensionKind::McpServer); + let registry_only = InstalledExtension { + installed: false, + active: false, + authenticated: false, + has_auth: true, + tools: Vec::new(), + ..installed_extension("linear_server", ExtensionKind::McpServer) + }; + + let snapshot = CapabilityRuntimeSnapshot { + extensions: vec![installed, registry_only], + latent_actions: Vec::new(), + channel_routes: HashMap::new(), + }; + + let projected = CapabilityProjector::project_snapshot(snapshot, &[]); + assert!(projected.is_empty()); + } + + #[test] + fn ready_direct_provider_actions_stay_out_of_capabilities() { + let snapshot = CapabilityRuntimeSnapshot { + extensions: vec![installed_extension("drive", ExtensionKind::WasmTool)], + latent_actions: Vec::new(), + channel_routes: HashMap::new(), + }; + + let projected = CapabilityProjector::project_snapshot(snapshot, &[]); + assert!(projected.is_empty()); + } +} diff --git a/src/bridge/effect_adapter.rs b/src/bridge/effect_adapter.rs index 842ec61c521..3c8f9505e3f 100644 --- a/src/bridge/effect_adapter.rs +++ b/src/bridge/effect_adapter.rs @@ -8,7 +8,7 @@ //! - Sensitive parameter redaction //! - Rate limiting (per-user, per-tool) -use std::collections::HashSet; +use std::collections::{HashMap, HashSet}; use std::sync::Arc; use std::time::Instant; @@ -16,16 +16,19 @@ use tokio::sync::RwLock; use tracing::debug; use ironclaw_engine::{ - ActionDef, ActionResult, CapabilityLease, CapabilityRegistry, EffectExecutor, EngineError, - MountError, Store, ThreadExecutionContext, WorkspaceMounts, + ActionDef, ActionResult, CapabilityLease, CapabilityRegistry, CapabilitySummary, + EffectExecutor, EngineError, MountError, Store, ThreadExecutionContext, WorkspaceMounts, }; use ironclaw_skills::SkillRegistry; use crate::auth::extension::{AuthCheckResult, AuthManager, LatentActionExecution, ToolReadiness}; use crate::auth::oauth::sanitize_auth_url; +use crate::bridge::action_projector::ActionProjector; +use crate::bridge::capability_projector::CapabilityProjector; use crate::bridge::router::synthetic_action_call_id; use crate::bridge::sandbox::{InterceptOutcome, maybe_intercept}; use crate::context::JobContext; +use crate::extensions::InstalledExtension; use crate::hooks::{HookEvent, HookOutcome, HookRegistry}; use crate::tools::permissions::{PermissionState, effective_permission}; use crate::tools::rate_limiter::RateLimiter; @@ -75,6 +78,30 @@ pub struct EffectBridgeAdapter { /// capabilities like `missions` are registered here in `router.rs` and /// would otherwise be invisible to the LLM despite having active leases. capability_registry: RwLock>>, + /// Short-lived cache for `list_capability_extensions` results. Keyed by + /// `user_id` with an expiry timestamp. Both `available_actions` and + /// `available_capabilities` are called in quick succession from the engine + /// loop; this cache eliminates the duplicate fetch. + extension_cache: RwLock>, +} + +/// Cached extension list with a short TTL to deduplicate the fetch across +/// `available_actions` and `available_capabilities` when called in sequence. +struct ExtensionCacheEntry { + user_id: String, + fetched_at: Instant, + extensions: Vec, +} + +impl ExtensionCacheEntry { + /// Cache entries expire after 500ms — long enough to cover back-to-back + /// calls within a single system-prompt build, short enough to never serve + /// stale data across separate engine iterations. + const TTL_MS: u128 = 500; + + fn is_valid_for(&self, user_id: &str) -> bool { + self.user_id == user_id && self.fetched_at.elapsed().as_millis() < Self::TTL_MS + } } impl EffectBridgeAdapter { @@ -98,6 +125,7 @@ impl EffectBridgeAdapter { skill_registry: RwLock::new(None), workspace_mounts: RwLock::new(None), capability_registry: RwLock::new(None), + extension_cache: RwLock::new(None), } } @@ -181,6 +209,72 @@ impl EffectBridgeAdapter { self.mission_manager.read().await.clone() } + /// Fetch the extension list as a `Vec`, using the short-lived cache when + /// available. Returns `Some(vec)` when an auth_manager is present, `None` + /// otherwise. + async fn fetch_extension_list( + &self, + auth_manager: Option<&AuthManager>, + context: &ThreadExecutionContext, + ) -> Option> { + let auth_manager = auth_manager?; + + // Check cache first. + { + let cache = self.extension_cache.read().await; + if let Some(entry) = cache.as_ref() + && entry.is_valid_for(&context.user_id) + { + return Some(entry.extensions.clone()); + } + } + + // Cache miss — fetch from auth_manager. + let extensions = match auth_manager + .list_capability_extensions(&context.user_id) + .await + { + Ok(exts) => exts, + Err(error) => { + debug!( + user_id = %context.user_id, + error = %error, + "failed to load extension inventory; returning empty list" + ); + Vec::new() + } + }; + + // Populate cache for the sibling projector call. + { + let mut cache = self.extension_cache.write().await; + *cache = Some(ExtensionCacheEntry { + user_id: context.user_id.clone(), + fetched_at: Instant::now(), + extensions: extensions.clone(), + }); + } + + Some(extensions) + } + + /// Fetch the extension list as a `HashMap`, using + /// the short-lived cache when available. Returns `Some(map)` when an + /// auth_manager is present, `None` otherwise. + async fn fetch_extension_map( + &self, + auth_manager: Option<&AuthManager>, + context: &ThreadExecutionContext, + ) -> Option> { + let extensions = self.fetch_extension_list(auth_manager, context).await?; + Some( + extensions + .into_iter() + .map(|ext| (ext.name.clone(), ext)) + .collect(), + ) + } + async fn sync_skill_install_result( &self, output_value: &serde_json::Value, @@ -1480,97 +1574,34 @@ impl EffectExecutor for EffectBridgeAdapter { async fn available_actions( &self, leases: &[CapabilityLease], + context: &ThreadExecutionContext, ) -> Result, EngineError> { - let tool_defs = self.tools.tool_definitions().await; - - // Build action defs, excluding v1-only tools and v1 auth tools - let mut actions = Vec::with_capacity(tool_defs.len()); - for td in tool_defs { - // Skip tools that can't work in engine v2 - if is_v1_only_tool(&td.name) { - continue; - } - - // Skip v1 auth management tools — auth is kernel-level in v2 - if is_v1_auth_tool(&td.name) { - continue; - } - - let python_name = td.name.replace('-', "_"); - - actions.push(ActionDef { - name: python_name, - description: td.description, - parameters_schema: td.parameters, - effects: vec![], - // Approval is enforced at execute-time inside this adapter so - // thread-scoped one-shot approvals and auth-aware bypasses can - // participate. Advertising approval here would cause the engine - // policy preflight to interrupt before the adapter can apply - // those runtime checks. - requires_approval: false, - }); - } - - if let Some(auth_mgr) = self.auth_manager.read().await.as_ref() { - for latent in auth_mgr.latent_extension_actions().await { - if actions - .iter() - .any(|action| action.name == latent.action_name) - { - continue; - } - actions.push(ActionDef { - name: latent.action_name, - description: latent.description, - parameters_schema: latent.parameters_schema, - effects: vec![], - requires_approval: false, - }); - } - } - - // Surface actions from engine-native capabilities (e.g. `missions`). - // The v1 `ToolRegistry` path above only covers built-in + extension - // tools; capabilities registered directly against the engine - // (`CapabilityRegistry`) would otherwise be invisible to the LLM - // even though the thread holds active leases for them. Iterate - // leases so we only advertise what the current thread actually has - // access to, and skip the `"tools"` capability — that lease is - // reconciled dynamically from the v1 path already covered above. - if let Some(registry) = self.capability_registry.read().await.as_ref() { - let mut seen: HashSet = actions.iter().map(|a| a.name.clone()).collect(); - for lease in leases { - if lease.capability_name == "tools" { - continue; - } - let Some(cap) = registry.get(&lease.capability_name) else { - continue; - }; - for action in &cap.actions { - if !lease.granted_actions.covers(&action.name) { - continue; - } - // Defensive: apply the same v1-isolation filters we run - // on v1 tools. If a future engine capability registers - // an action under a v1-denylisted name (`create_job`, - // `tool_auth`, ...), the v1 filters above would have - // hidden it — the engine path must not become a - // silent bypass. - if is_v1_only_tool(&action.name) || is_v1_auth_tool(&action.name) { - continue; - } - if !seen.insert(action.name.clone()) { - continue; - } - actions.push(action.clone()); - } - } - } - - actions.sort_by(|a, b| a.name.cmp(&b.name)); + let auth_manager = self.auth_manager.read().await.clone(); + let capability_registry = self.capability_registry.read().await.clone(); + let extensions = self + .fetch_extension_map(auth_manager.as_deref(), context) + .await; + ActionProjector::project( + self.tools.as_ref(), + auth_manager.as_deref(), + capability_registry, + leases, + context, + extensions.as_ref(), + ) + .await + } - Ok(actions) + async fn available_capabilities( + &self, + leases: &[CapabilityLease], + context: &ThreadExecutionContext, + ) -> Result, EngineError> { + let auth_manager = self.auth_manager.read().await.clone(); + let extensions = self + .fetch_extension_list(auth_manager.as_deref(), context) + .await; + CapabilityProjector::project(auth_manager.as_deref(), leases, context, extensions).await } } @@ -2260,7 +2291,7 @@ fn extract_credential_name(error_msg: &str) -> Option { None } -fn is_v1_only_tool(name: &str) -> bool { +pub(crate) fn is_v1_only_tool(name: &str) -> bool { // routine_* tools are surfaced in v2 too, but are intercepted by // `handle_mission_call`'s routine alias path *before* this check fires — // they get translated into mission_* dispatches via the existing @@ -2281,7 +2312,7 @@ fn is_v1_only_tool(name: &str) -> bool { /// Auth management tools from v1 that are now kernel-internal in v2. /// The LLM should not see or call these — auth is handled automatically. -fn is_v1_auth_tool(name: &str) -> bool { +pub(crate) fn is_v1_auth_tool(name: &str) -> bool { matches!(name, "tool_auth" | "tool-auth") } @@ -3683,9 +3714,18 @@ mod tests { async fn available_actions( &self, _: &[ironclaw_engine::CapabilityLease], + _: &ironclaw_engine::ThreadExecutionContext, ) -> Result, EngineError> { Ok(vec![]) } + + async fn available_capabilities( + &self, + _: &[ironclaw_engine::CapabilityLease], + _: &ironclaw_engine::ThreadExecutionContext, + ) -> Result, EngineError> { + Ok(vec![]) + } } // ── Helpers ────────────────────────────────────────── @@ -4229,7 +4269,7 @@ mod tests { } #[tokio::test] - async fn available_actions_include_latent_inactive_provider_actions() { + async fn available_actions_omit_latent_inactive_provider_actions() { use crate::secrets::InMemorySecretsStore; use crate::secrets::SecretsCrypto; use crate::tools::mcp::process::McpProcessManager; @@ -4288,8 +4328,196 @@ mod tests { ))) .await; - let actions = adapter.available_actions(&[]).await.expect("actions"); - assert!(actions.iter().any(|action| action.name == "latent_tool")); + let actions = adapter + .available_actions(&[], &exec_ctx(ironclaw_engine::ThreadId::new(), None)) + .await + .expect("actions"); + assert!(!actions.iter().any(|action| action.name == "latent_tool")); + } + + #[tokio::test] + async fn available_actions_omit_registered_tool_when_provider_is_not_installed() { + use crate::secrets::InMemorySecretsStore; + use crate::secrets::SecretsCrypto; + use crate::tools::mcp::process::McpProcessManager; + use crate::tools::mcp::session::McpSessionManager; + + struct LinearSearchTool; + + #[async_trait] + impl crate::tools::Tool for LinearSearchTool { + fn name(&self) -> &str { + "linear_search" + } + + fn description(&self) -> &str { + "Search Linear issues" + } + + fn parameters_schema(&self) -> serde_json::Value { + serde_json::json!({"type": "object"}) + } + + async fn execute( + &self, + _: serde_json::Value, + _: &crate::context::JobContext, + ) -> Result { + Ok(crate::tools::ToolOutput::success( + serde_json::json!({}), + std::time::Duration::from_millis(1), + )) + } + + fn provider_extension(&self) -> Option<&str> { + Some("linear") + } + } + + let dir = tempfile::tempdir().expect("temp dir"); + std::fs::create_dir_all(dir.path().join("tools")).expect("tools dir"); + std::fs::write(dir.path().join("tools").join("linear.wasm"), b"fake-wasm") + .expect("write wasm"); + std::fs::write( + dir.path().join("tools").join("linear.capabilities.json"), + r#"{"description":"linear provider"}"#, + ) + .expect("write capabilities"); + + let key = secrecy::SecretString::from(crate::secrets::keychain::generate_master_key_hex()); + let crypto = Arc::new(SecretsCrypto::new(key).expect("crypto")); + let secrets: Arc = + Arc::new(InMemorySecretsStore::new(crypto)); + + let tools = Arc::new(ToolRegistry::new()); + tools.register(Arc::new(LinearSearchTool)).await; + let ext_mgr = Arc::new(crate::extensions::ExtensionManager::new( + Arc::new(McpSessionManager::new()), + Arc::new(McpProcessManager::new()), + Arc::clone(&secrets), + Arc::clone(&tools), + None, + None, + dir.path().join("tools"), + dir.path().join("channels"), + None, + "test_user".to_string(), + None, + vec![], + )); + + let adapter = EffectBridgeAdapter::new( + Arc::clone(&tools), + Arc::new(SafetyLayer::new(&ironclaw_safety::SafetyConfig { + max_output_length: 10_000, + injection_check_enabled: false, + })), + Arc::new(HookRegistry::default()), + ); + adapter + .set_auth_manager(Arc::new(AuthManager::new( + secrets, + None, + Some(ext_mgr), + Some(Arc::clone(&tools)), + ))) + .await; + + let actions = adapter + .available_actions(&[], &exec_ctx(ironclaw_engine::ThreadId::new(), None)) + .await + .expect("actions"); + assert!(!actions.iter().any(|action| action.name == "linear_search")); + } + + // NeedsAuth provider tool preservation is tested directly in + // action_projector::tests::needs_auth_provider_tools_remain_in_available_actions + // where the extension map can be constructed with correct NeedsAuth status. + // An EffectBridgeAdapter-level test would require a real WASM module to + // produce installed=true + authenticated=false; fake-wasm files produce + // installed=false (AvailableNotInstalled), making the test unreliable. + + #[tokio::test] + async fn available_capabilities_projects_inactive_provider_background() { + use crate::secrets::InMemorySecretsStore; + use crate::secrets::SecretsCrypto; + use crate::tools::mcp::process::McpProcessManager; + use crate::tools::mcp::session::McpSessionManager; + + let dir = tempfile::tempdir().expect("temp dir"); + std::fs::create_dir_all(dir.path().join("tools")).expect("tools dir"); + std::fs::write( + dir.path().join("tools").join("latent_tool.wasm"), + b"fake-wasm", + ) + .expect("write wasm"); + std::fs::write( + dir.path() + .join("tools") + .join("latent_tool.capabilities.json"), + r#"{"description":"latent capability test"}"#, + ) + .expect("write capabilities"); + + let key = secrecy::SecretString::from(crate::secrets::keychain::generate_master_key_hex()); + let crypto = Arc::new(SecretsCrypto::new(key).expect("crypto")); + let secrets: Arc = + Arc::new(InMemorySecretsStore::new(crypto)); + + let tools = Arc::new(ToolRegistry::new()); + let ext_mgr = Arc::new(crate::extensions::ExtensionManager::new( + Arc::new(McpSessionManager::new()), + Arc::new(McpProcessManager::new()), + Arc::clone(&secrets), + Arc::clone(&tools), + None, + None, + dir.path().join("tools"), + dir.path().join("channels"), + None, + "test_user".to_string(), + None, + vec![], + )); + + let adapter = EffectBridgeAdapter::new( + Arc::clone(&tools), + Arc::new(SafetyLayer::new(&ironclaw_safety::SafetyConfig { + max_output_length: 10_000, + injection_check_enabled: false, + })), + Arc::new(HookRegistry::default()), + ); + adapter + .set_auth_manager(Arc::new(AuthManager::new( + secrets, + None, + Some(ext_mgr), + Some(Arc::clone(&tools)), + ))) + .await; + + let context = ironclaw_engine::ThreadExecutionContext { + thread_id: ironclaw_engine::ThreadId::new(), + thread_type: ironclaw_engine::types::thread::ThreadType::Foreground, + project_id: ironclaw_engine::ProjectId::new(), + user_id: "test_user".to_string(), + step_id: ironclaw_engine::StepId::new(), + current_call_id: None, + source_channel: None, + user_timezone: None, + thread_goal: None, + }; + + let capabilities = adapter + .available_capabilities(&[], &context) + .await + .expect("capabilities"); + + assert!(capabilities.iter().any(|summary| { + summary.name == "latent_tool" + && summary.status == ironclaw_engine::CapabilityStatus::Inactive + })); } #[tokio::test] @@ -5012,9 +5240,19 @@ Use this skill to set up a Pika meeting. async fn available_actions( &self, _: &[ironclaw_engine::CapabilityLease], + _: &ironclaw_engine::ThreadExecutionContext, ) -> Result, ironclaw_engine::EngineError> { Ok(vec![]) } + + async fn available_capabilities( + &self, + _: &[ironclaw_engine::CapabilityLease], + _: &ironclaw_engine::ThreadExecutionContext, + ) -> Result, ironclaw_engine::EngineError> + { + Ok(vec![]) + } } let concrete_store = Arc::new(mission_store::TestStore::new()); @@ -5578,11 +5816,14 @@ Use this skill to set up a Pika meeting. adapter.set_capability_registry(Arc::new(registry)).await; let actions = adapter - .available_actions(&[mission_lease(&[ - "mission_create", - "mission_list", - "mission_complete", - ])]) + .available_actions( + &[mission_lease(&[ + "mission_create", + "mission_list", + "mission_complete", + ])], + &exec_ctx(ironclaw_engine::ThreadId::new(), None), + ) .await .expect("available_actions should succeed"); @@ -5606,7 +5847,10 @@ Use this skill to set up a Pika meeting. // must NOT be advertised to the LLM even though they exist in the // capability registry. let actions = adapter - .available_actions(&[mission_lease(&["mission_list"])]) + .available_actions( + &[mission_lease(&["mission_list"])], + &exec_ctx(ironclaw_engine::ThreadId::new(), None), + ) .await .expect("available_actions should succeed"); @@ -5635,7 +5879,7 @@ Use this skill to set up a Pika meeting. // No leases passed — no capability actions should surface even // though the registry has them. let actions = adapter - .available_actions(&[]) + .available_actions(&[], &exec_ctx(ironclaw_engine::ThreadId::new(), None)) .await .expect("available_actions should succeed"); @@ -5700,11 +5944,14 @@ Use this skill to set up a Pika meeting. adapter.set_capability_registry(Arc::new(registry)).await; let actions = adapter - .available_actions(&[mission_lease(&[ - "mission_create", - "mission_list", - "mission_complete", - ])]) + .available_actions( + &[mission_lease(&[ + "mission_create", + "mission_list", + "mission_complete", + ])], + &exec_ctx(ironclaw_engine::ThreadId::new(), None), + ) .await .expect("available_actions should succeed"); @@ -5777,7 +6024,10 @@ Use this skill to set up a Pika meeting. }; let actions = adapter - .available_actions(&[rogue_lease]) + .available_actions( + &[rogue_lease], + &exec_ctx(ironclaw_engine::ThreadId::new(), None), + ) .await .expect("available_actions should succeed"); diff --git a/src/bridge/mod.rs b/src/bridge/mod.rs index fb1a2e80ab6..221e048a2b8 100644 --- a/src/bridge/mod.rs +++ b/src/bridge/mod.rs @@ -4,6 +4,8 @@ //! route through the engine instead of the existing agentic loop. All //! existing behavior is unchanged when the flag is off. +mod action_projector; +mod capability_projector; mod cost_guard_gate; mod effect_adapter; mod llm_adapter; @@ -11,6 +13,7 @@ mod router; pub mod sandbox; pub mod skill_migration; mod store_adapter; +mod tool_surface; mod user_facing_errors; mod workspace_reader; diff --git a/src/bridge/router.rs b/src/bridge/router.rs index fb265bd7031..28c69489866 100644 --- a/src/bridge/router.rs +++ b/src/bridge/router.rs @@ -5955,9 +5955,18 @@ pub(crate) mod test_support { async fn available_actions( &self, _: &[CapabilityLease], + _: &ironclaw_engine::ThreadExecutionContext, ) -> Result, EngineError> { Ok(vec![]) } + + async fn available_capabilities( + &self, + _: &[CapabilityLease], + _: &ironclaw_engine::ThreadExecutionContext, + ) -> Result, EngineError> { + Ok(vec![]) + } } let store = Arc::new(ThreadTestStore::new()); @@ -7990,9 +7999,19 @@ mod tests { async fn available_actions( &self, _: &[ironclaw_engine::CapabilityLease], + _: &ironclaw_engine::ThreadExecutionContext, ) -> Result, ironclaw_engine::EngineError> { Ok(vec![]) } + + async fn available_capabilities( + &self, + _: &[ironclaw_engine::CapabilityLease], + _: &ironclaw_engine::ThreadExecutionContext, + ) -> Result, ironclaw_engine::EngineError> + { + Ok(vec![]) + } } let store_dyn: Arc = store; @@ -8130,9 +8149,19 @@ mod tests { async fn available_actions( &self, _: &[ironclaw_engine::CapabilityLease], + _: &ironclaw_engine::ThreadExecutionContext, ) -> Result, ironclaw_engine::EngineError> { Ok(vec![]) } + + async fn available_capabilities( + &self, + _: &[ironclaw_engine::CapabilityLease], + _: &ironclaw_engine::ThreadExecutionContext, + ) -> Result, ironclaw_engine::EngineError> + { + Ok(vec![]) + } } let store_dyn: Arc = store; @@ -9455,9 +9484,19 @@ mod tests { async fn available_actions( &self, _: &[ironclaw_engine::CapabilityLease], + _: &ironclaw_engine::ThreadExecutionContext, ) -> Result, ironclaw_engine::EngineError> { Ok(vec![]) } + + async fn available_capabilities( + &self, + _: &[ironclaw_engine::CapabilityLease], + _: &ironclaw_engine::ThreadExecutionContext, + ) -> Result, ironclaw_engine::EngineError> + { + Ok(vec![]) + } } let store: Arc = Arc::new(TestStore::new()); diff --git a/src/bridge/tool_surface.rs b/src/bridge/tool_surface.rs new file mode 100644 index 00000000000..ccf2e763a14 --- /dev/null +++ b/src/bridge/tool_surface.rs @@ -0,0 +1,289 @@ +#![cfg_attr(not(test), allow(dead_code))] + +use ironclaw_engine::CapabilityStatus; + +/// How the subject can be invoked from the model/runtime boundary. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub(crate) enum InvocationMode { + Direct, + RoutedOnly, +} + +/// Bridge-owned subject categories for surface placement policy. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub(crate) enum SurfaceSubjectKind { + BuiltinDirectTool, + ExtensionDirectAction, + EngineNativeDirectAction, + Channel, + LatentProviderAction, + AvailableNotInstalledProviderEntry, +} + +/// Pure input to the bridge-owned surface assignment policy. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub(crate) struct SurfacePolicyInput { + pub(crate) kind: SurfaceSubjectKind, + pub(crate) status: CapabilityStatus, + pub(crate) invocation_mode: InvocationMode, + /// Engine-native direct actions also need a current callable lease before + /// they belong in `available_actions`. + pub(crate) leased_and_callable: bool, +} + +/// Pure result describing which bridge surfaces should include the subject. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub(crate) struct SurfaceAssignment { + pub(crate) available_actions: bool, + pub(crate) available_capabilities: bool, +} + +impl SurfaceAssignment { + const fn actions_only() -> Self { + Self { + available_actions: true, + available_capabilities: false, + } + } + + const fn capabilities_only() -> Self { + Self { + available_actions: false, + available_capabilities: true, + } + } + + const fn neither() -> Self { + Self { + available_actions: false, + available_capabilities: false, + } + } +} + +pub(crate) fn assign_surface(subject: SurfacePolicyInput) -> SurfaceAssignment { + if matches!(subject.status, CapabilityStatus::Error) { + return SurfaceAssignment::neither(); + } + + if matches!(subject.invocation_mode, InvocationMode::RoutedOnly) { + return SurfaceAssignment::capabilities_only(); + } + + match subject.kind { + SurfaceSubjectKind::LatentProviderAction + | SurfaceSubjectKind::AvailableNotInstalledProviderEntry + | SurfaceSubjectKind::Channel => SurfaceAssignment::capabilities_only(), + SurfaceSubjectKind::BuiltinDirectTool | SurfaceSubjectKind::ExtensionDirectAction => { + // Execute-time approval is orthogonal to surface placement: direct + // actions stay on the callable surface as long as they are ready. + if is_direct_ready(subject.status) { + SurfaceAssignment::actions_only() + } else { + fallback_assignment(subject.status) + } + } + SurfaceSubjectKind::EngineNativeDirectAction => { + if is_direct_ready(subject.status) && subject.leased_and_callable { + SurfaceAssignment::actions_only() + } else if is_direct_ready(subject.status) { + SurfaceAssignment::capabilities_only() + } else { + fallback_assignment(subject.status) + } + } + } +} + +const fn is_direct_ready(status: CapabilityStatus) -> bool { + matches!(status, CapabilityStatus::Ready) +} + +const fn fallback_assignment(status: CapabilityStatus) -> SurfaceAssignment { + match status { + // ReadyScoped subjects are not directly callable but should remain + // visible in the capabilities surface so scoped functionality is + // discoverable in background context. + CapabilityStatus::NeedsAuth + | CapabilityStatus::NeedsSetup + | CapabilityStatus::Inactive + | CapabilityStatus::Latent + | CapabilityStatus::AvailableNotInstalled + | CapabilityStatus::ReadyScoped => SurfaceAssignment::capabilities_only(), + CapabilityStatus::Ready | CapabilityStatus::Error => SurfaceAssignment::neither(), + } +} + +#[cfg(test)] +mod tests { + use super::{ + InvocationMode, SurfaceAssignment, SurfacePolicyInput, SurfaceSubjectKind, assign_surface, + }; + use ironclaw_engine::CapabilityStatus; + + #[test] + fn assigns_surface_matrix_rows() { + struct Case { + name: &'static str, + subject: SurfacePolicyInput, + expected: SurfaceAssignment, + } + + let cases = [ + Case { + name: "ready built-in direct tool", + subject: SurfacePolicyInput { + kind: SurfaceSubjectKind::BuiltinDirectTool, + status: CapabilityStatus::Ready, + invocation_mode: InvocationMode::Direct, + leased_and_callable: false, + }, + expected: SurfaceAssignment::actions_only(), + }, + Case { + name: "approval-gated built-in direct tool still stays in available_actions", + subject: SurfacePolicyInput { + kind: SurfaceSubjectKind::BuiltinDirectTool, + status: CapabilityStatus::Ready, + invocation_mode: InvocationMode::Direct, + leased_and_callable: false, + }, + expected: SurfaceAssignment::actions_only(), + }, + Case { + name: "ready extension direct action", + subject: SurfacePolicyInput { + kind: SurfaceSubjectKind::ExtensionDirectAction, + status: CapabilityStatus::Ready, + invocation_mode: InvocationMode::Direct, + leased_and_callable: false, + }, + expected: SurfaceAssignment::actions_only(), + }, + Case { + name: "approval-gated extension direct action still stays in available_actions", + subject: SurfacePolicyInput { + kind: SurfaceSubjectKind::ExtensionDirectAction, + status: CapabilityStatus::Ready, + invocation_mode: InvocationMode::Direct, + leased_and_callable: false, + }, + expected: SurfaceAssignment::actions_only(), + }, + Case { + name: "needs-auth extension direct action", + subject: SurfacePolicyInput { + kind: SurfaceSubjectKind::ExtensionDirectAction, + status: CapabilityStatus::NeedsAuth, + invocation_mode: InvocationMode::Direct, + leased_and_callable: false, + }, + expected: SurfaceAssignment::capabilities_only(), + }, + Case { + name: "needs-setup extension direct action", + subject: SurfacePolicyInput { + kind: SurfaceSubjectKind::ExtensionDirectAction, + status: CapabilityStatus::NeedsSetup, + invocation_mode: InvocationMode::Direct, + leased_and_callable: false, + }, + expected: SurfaceAssignment::capabilities_only(), + }, + Case { + name: "inactive extension direct action", + subject: SurfacePolicyInput { + kind: SurfaceSubjectKind::ExtensionDirectAction, + status: CapabilityStatus::Inactive, + invocation_mode: InvocationMode::Direct, + leased_and_callable: false, + }, + expected: SurfaceAssignment::capabilities_only(), + }, + Case { + name: "error extension direct action", + subject: SurfacePolicyInput { + kind: SurfaceSubjectKind::ExtensionDirectAction, + status: CapabilityStatus::Error, + invocation_mode: InvocationMode::Direct, + leased_and_callable: false, + }, + expected: SurfaceAssignment::neither(), + }, + Case { + name: "latent provider action", + subject: SurfacePolicyInput { + kind: SurfaceSubjectKind::LatentProviderAction, + status: CapabilityStatus::Latent, + invocation_mode: InvocationMode::Direct, + leased_and_callable: false, + }, + expected: SurfaceAssignment::capabilities_only(), + }, + Case { + name: "available-not-installed provider entry", + subject: SurfacePolicyInput { + kind: SurfaceSubjectKind::AvailableNotInstalledProviderEntry, + status: CapabilityStatus::AvailableNotInstalled, + invocation_mode: InvocationMode::Direct, + leased_and_callable: false, + }, + expected: SurfaceAssignment::capabilities_only(), + }, + Case { + name: "routed-only channel", + subject: SurfacePolicyInput { + kind: SurfaceSubjectKind::Channel, + status: CapabilityStatus::ReadyScoped, + invocation_mode: InvocationMode::RoutedOnly, + leased_and_callable: false, + }, + expected: SurfaceAssignment::capabilities_only(), + }, + Case { + name: "ready-scoped extension direct action is not callable but visible in capabilities", + subject: SurfacePolicyInput { + kind: SurfaceSubjectKind::ExtensionDirectAction, + status: CapabilityStatus::ReadyScoped, + invocation_mode: InvocationMode::Direct, + leased_and_callable: false, + }, + expected: SurfaceAssignment::capabilities_only(), + }, + Case { + name: "ready leased engine-native direct action", + subject: SurfacePolicyInput { + kind: SurfaceSubjectKind::EngineNativeDirectAction, + status: CapabilityStatus::Ready, + invocation_mode: InvocationMode::Direct, + leased_and_callable: true, + }, + expected: SurfaceAssignment::actions_only(), + }, + Case { + name: "engine-native direct action without current callable lease", + subject: SurfacePolicyInput { + kind: SurfaceSubjectKind::EngineNativeDirectAction, + status: CapabilityStatus::Ready, + invocation_mode: InvocationMode::Direct, + leased_and_callable: false, + }, + expected: SurfaceAssignment::capabilities_only(), + }, + Case { + name: "error channel stays off all surfaces", + subject: SurfacePolicyInput { + kind: SurfaceSubjectKind::Channel, + status: CapabilityStatus::Error, + invocation_mode: InvocationMode::RoutedOnly, + leased_and_callable: false, + }, + expected: SurfaceAssignment::neither(), + }, + ]; + + for case in cases { + assert_eq!(assign_surface(case.subject), case.expected, "{}", case.name); + } + } +} diff --git a/tests/engine_v2_gate_integration.rs b/tests/engine_v2_gate_integration.rs index 2d9df8cf71f..eff1a0746df 100644 --- a/tests/engine_v2_gate_integration.rs +++ b/tests/engine_v2_gate_integration.rs @@ -236,6 +236,7 @@ impl EffectExecutor for InstallThenAliasEffects { async fn available_actions( &self, _leases: &[CapabilityLease], + _context: &ironclaw_engine::ThreadExecutionContext, ) -> Result, EngineError> { Ok(vec![ ActionDef { @@ -254,6 +255,14 @@ impl EffectExecutor for InstallThenAliasEffects { }, ]) } + + async fn available_capabilities( + &self, + _leases: &[CapabilityLease], + _context: &ironclaw_engine::ThreadExecutionContext, + ) -> Result, EngineError> { + Ok(vec![]) + } } #[async_trait::async_trait] @@ -348,6 +357,7 @@ impl EffectExecutor for GateMockEffects { async fn available_actions( &self, _leases: &[CapabilityLease], + _context: &ironclaw_engine::ThreadExecutionContext, ) -> Result, EngineError> { // requires_approval: false — the gate check is done by the mock's // execute_action() returning GatePaused, not by the PolicyEngine. @@ -375,6 +385,14 @@ impl EffectExecutor for GateMockEffects { }, ]) } + + async fn available_capabilities( + &self, + _leases: &[CapabilityLease], + _context: &ironclaw_engine::ThreadExecutionContext, + ) -> Result, EngineError> { + Ok(vec![]) + } } // ── In-Memory Store (same as engine_v2_skill_codeact) ──────── diff --git a/tests/engine_v2_skill_codeact.rs b/tests/engine_v2_skill_codeact.rs index 3369af0260c..27c809d8417 100644 --- a/tests/engine_v2_skill_codeact.rs +++ b/tests/engine_v2_skill_codeact.rs @@ -130,6 +130,7 @@ impl EffectExecutor for HttpMockEffects { async fn available_actions( &self, _leases: &[CapabilityLease], + _context: &ironclaw_engine::ThreadExecutionContext, ) -> Result, EngineError> { Ok(vec![ActionDef { name: "http".into(), @@ -148,6 +149,14 @@ impl EffectExecutor for HttpMockEffects { requires_approval: false, }]) } + + async fn available_capabilities( + &self, + _leases: &[CapabilityLease], + _context: &ironclaw_engine::ThreadExecutionContext, + ) -> Result, EngineError> { + Ok(vec![]) + } } // ── In-Memory Store ──────────────────────────────────────────