返回 CodeWhale
native_mcp.rs
根目录 / crates / tui / src / extension_host / native_mcp.rs
1 //! Native definition admission joins Rust's existing MCP pool; no transport here.
2 use super::ManagerShared;
3 use super::composition_scope::SelectionRevision;
4 use super::protocol::{EntryRef, OwnerRef, RegisterKind, RegisterParams, RegisterResult};
5 use super::supervisor::HostRequestContext;
6 use super::tier::HostTier;
7 use crate::mcp::McpServerConfig;
8 use crate::plugins::{
9 PluginRegistry, activation::PluginActivationCapability, types::PluginAuthority,
10 };
11 use std::sync::atomic::{AtomicU64, Ordering};
12 use std::sync::{Arc, Weak};
13
14 pub(super) const MAX_PER_OWNER: usize = 64;
15 pub(super) const MAX_PER_HOST: usize = 256;
16 pub(super) const MAX_DEFINITION_BYTES: usize = 64 * 1024;
17 static EPOCH: AtomicU64 = AtomicU64::new(1);
18 pub(crate) fn epoch() -> u64 {
19 EPOCH.load(Ordering::SeqCst)
20 }
21 pub(super) fn changed() {
22 EPOCH.fetch_add(1, Ordering::SeqCst);
23 }
24
25 #[derive(Debug, Clone)]
26 pub(super) struct McpRegistration {
27 pub handle: u64,
28 pub owner: OwnerRef,
29 pub scope: EntryRef,
30 pub content_hash: String,
31 pub host_generation: u64,
32 pub name: String,
33 pub config: McpServerConfig,
34 pub cancel: tokio_util::sync::CancellationToken,
35 }
36 /// In-memory only: exact generation/handle/caller receipt. No host token is persisted.
37 #[derive(Debug, Clone)]
38 pub(crate) struct NativeMcpRef {
39 shared: Weak<dyn NativeMcpAuthority>,
40 registration: Arc<McpRegistration>,
41 selection: SelectionRevision,
42 caller_cancel: tokio_util::sync::CancellationToken,
43 }
44 // Keep MCP configs independent of the concrete manager's Engine/task graph.
45 // The weak port still targets the original manager and owns no runtime state.
46 trait NativeMcpAuthority: Send + Sync {
47 fn validate_mcp(
48 &self,
49 selection: SelectionRevision,
50 registration: &McpRegistration,
51 ) -> Result<(), String>;
52 }
53 impl NativeMcpAuthority for ManagerShared {
54 fn validate_mcp(
55 &self,
56 selection: SelectionRevision,
57 registration: &McpRegistration,
58 ) -> Result<(), String> {
59 if !crate::plugins::activation::extension_host_policy_enabled() {
60 return Err("Native MCP is disabled by current policy".into());
61 }
62 self.ready_host(HostTier::Plugin)
63 .map_err(|status| status.to_string())?;
64 if self
65 .tier_runtime(HostTier::Plugin)
66 .host_generation
67 .load(Ordering::SeqCst)
68 != registration.host_generation
69 {
70 return Err("Native MCP host generation changed".into());
71 }
72 if !self.selection_current(
73 selection,
74 &registration.owner.plugin_id,
75 &registration.content_hash,
76 Some(&registration.scope),
77 ) {
78 return Err("Native MCP entry is no longer selected for this caller".into());
79 }
80 let registry = self.registry.lock().expect("registry lock");
81 if !registry.is_live_mcp(registration) {
82 return Err("Native MCP definition was withdrawn".into());
83 }
84 Ok(())
85 }
86 }
87 impl NativeMcpRef {
88 pub(crate) fn validate(&self) -> Result<(), String> {
89 self.shared
90 .upgrade()
91 .ok_or("Native MCP manager exited")?
92 .validate_mcp(self.selection, &self.registration)
93 }
94 pub(crate) fn selection(&self) -> SelectionRevision {
95 self.selection
96 }
97 pub(crate) fn catalog_identity(&self) -> String {
98 format!(
99 "{}:{}:{}:{}:{}:{}",
100 self.registration.host_generation,
101 self.registration.owner.generation,
102 self.registration.handle,
103 self.selection.attachment_id,
104 self.selection.revision,
105 self.registration.scope.sha256
106 )
107 }
108 pub(crate) async fn withdrawn(&self) {
109 tokio::select! {
110 biased;
111 _ = self.caller_cancel.cancelled() => {},
112 _ = self.registration.cancel.cancelled() => {},
113 }
114 }
115 }
116
117 pub(super) async fn admit(
118 shared: &Arc<ManagerShared>,
119 tier: HostTier,
120 generation: u64,
121 params: RegisterParams,
122 cx: &HostRequestContext,
123 ) -> RegisterResult {
124 let outcome = async {
125 params.check_spec()?;
126 if params.kind != RegisterKind::McpServer
127 || tier != HostTier::Plugin
128 || params.scope.is_none()
129 {
130 return Err("MCP definitions require a reviewed Native entry scope".into());
131 }
132 if params.spec.description.len() > MAX_DEFINITION_BYTES {
133 return Err("MCP definition exceeds 64 KiB".into());
134 }
135 let authority = shared
136 .live_owner_authority(tier, |registry| {
137 registry
138 .authority_for(&params.owner)
139 .ok_or_else(|| "stale MCP owner".to_string())?;
140 Ok(params.owner.clone())
141 })?
142 .ok_or("MCP definition has no Native authority")?;
143 let proposal = params.clone();
144 let config = super::skills::bounded_review_check(
145 Arc::clone(&shared.skill_admission),
146 &cx.cancel,
147 move || validate_definition(&authority, &proposal),
148 )
149 .await?;
150 shared
151 .ready_host(tier)
152 .map_err(|status| status.to_string())?;
153 let runtime = shared.tier_runtime(tier);
154 let _host = runtime.host.lock().expect("host lock");
155 if cx.cancel.is_cancelled() || runtime.host_generation.load(Ordering::SeqCst) != generation
156 {
157 return Err("MCP admission cancelled or host changed".into());
158 }
159 shared
160 .registry
161 .lock()
162 .expect("registry lock")
163 .register_mcp(&params, config, generation)
164 }
165 .await;
166 match outcome {
167 Ok(handle) => RegisterResult::Admitted { handle },
168 Err(refused) => {
169 shared.plugin_diagnostic(
170 &params.owner.plugin_id,
171 format!("MCP definition refused: {refused}"),
172 );
173 RegisterResult::Refused { refused }
174 }
175 }
176 }
177 fn validate_definition(
178 authority: &PluginAuthority,
179 params: &RegisterParams,
180 ) -> Result<McpServerConfig, String> {
181 crate::plugins::registry::verify_plugin_component_authority(
182 authority,
183 PluginActivationCapability::Native,
184 )?;
185 let declaration: serde_json::Value = serde_json::from_str(&params.spec.description)
186 .map_err(|_| "MCP proposal is not a literal server definition")?;
187 let text = serde_json::to_string(
188 &serde_json::json!({"mcpServers":{(params.spec.name.clone()):declaration}}),
189 )
190 .map_err(|_| "MCP definition encoding failed")?;
191 let mut servers = crate::plugins::agent_plugin::parse_mcp_json(&text)
192 .map_err(|_| "MCP proposal does not match the existing transport schema")?;
193 let config = servers
194 .remove(&params.spec.name)
195 .ok_or("MCP proposal omitted its server")?;
196 // Native rows cannot transport credential values to Core. Authentication
197 // remains the existing OAuth/environment-key authority at connection time.
198 if !config.env.is_empty() || !config.headers.is_empty() {
199 return Err("MCP credential values must remain under Core authentication authority".into());
200 }
201 let mut validated =
202 crate::plugins::manifest::PluginManifest::validate_from_path(&authority.staged_manifest)?;
203 let root = &validated.canonical_root;
204 validated.manifest.mcp_servers = Some(servers);
205 validated
206 .manifest
207 .mcp_servers
208 .as_mut()
209 .expect("set servers")
210 .insert(params.spec.name.clone(), config.clone());
211 validated.manifest.validate_mcp_servers(root)?;
212 crate::plugins::registry::verify_plugin_component_authority(
213 authority,
214 PluginActivationCapability::Native,
215 )?;
216 Ok(config)
217 }
218
219 pub(crate) fn for_plugins(
220 plugins: &PluginRegistry,
221 ) -> Result<Vec<(String, McpServerConfig, PluginAuthority, NativeMcpRef)>, String> {
222 if !crate::plugins::activation::extension_host_policy_enabled() {
223 return Ok(Vec::new());
224 }
225 let Some(selection) = plugins.caller_selection() else {
226 return Ok(Vec::new());
227 };
228 let manager = super::manager();
229 let shared = &manager.shared;
230 if shared.ready_host(HostTier::Plugin).is_err() {
231 return Ok(Vec::new());
232 }
233 let Some(caller_cancel) = shared.selection_cancellation(selection) else {
234 return Ok(Vec::new());
235 };
236 let live_authority: Arc<dyn NativeMcpAuthority> = shared.clone();
237 let definitions = shared.registry.lock().expect("registry lock").live_mcp();
238 let mut result = Vec::new();
239 let mut names = std::collections::HashSet::new();
240 for definition in definitions {
241 let Some(plugin) = plugins.get(&definition.owner.plugin_id) else {
242 continue;
243 };
244 if !plugin.active()
245 || plugin.content_hash != definition.content_hash
246 || !shared.selection_current(
247 selection,
248 &definition.owner.plugin_id,
249 &definition.content_hash,
250 Some(&definition.scope),
251 )
252 {
253 continue;
254 }
255 let Some(authority) = plugins.authority_for(&definition.owner.plugin_id) else {
256 continue;
257 };
258 let receipt = NativeMcpRef {
259 shared: Arc::downgrade(&live_authority),
260 registration: Arc::new(definition.clone()),
261 selection,
262 caller_cancel: caller_cancel.clone(),
263 };
264 if receipt.validate().is_err()
265 || crate::plugins::registry::verify_plugin_component_authority(
266 &authority,
267 PluginActivationCapability::Native,
268 )
269 .is_err()
270 {
271 continue;
272 }
273 let qualified =
274 crate::mcp::qualified_plugin_server_name(&authority.plugin_name, &definition.name);
275 if !names.insert(qualified) {
276 return Err("Selected Native entries duplicate an MCP server namespace".into());
277 }
278 result.push((definition.name, definition.config, authority, receipt));
279 }
280 Ok(result)
281 }
282
283 #[cfg(test)]
284 mod tests;
285
285 lines RUST