| 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 | ®istration.owner.plugin_id, |
| 75 | ®istration.content_hash, |
| 76 | Some(®istration.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(¶ms.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(¶ms, config, generation) |
| 164 | } |
| 165 | .await; |
| 166 | match outcome { |
| 167 | Ok(handle) => RegisterResult::Admitted { handle }, |
| 168 | Err(refused) => { |
| 169 | shared.plugin_diagnostic( |
| 170 | ¶ms.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(¶ms.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(¶ms.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 |