/// Manages MCP sessions across multiple `(user, server)` pairs.
///
/// Server names are typed via [`McpServerName`] so a free-form string can't
/// bypass allowlist validation at the boundary. Callers convert raw strings
/// via `McpServerName::new` (validating) or `McpServerName::from_trusted`
/// (for names the caller already validated). This makes identity-confusion
/// bugs — matching the shape described in `.claude/rules/types.md` — a
/// compile error rather than a runtime surprise.
pub struct McpSessionManager {
/// Active sessions keyed by `(user_id, server_name)`.
sessions: RwLock<HashMap<McpSessionKey, McpSession>>,
/// Maximum idle time before a session is considered stale (in seconds).
max_idle_secs: u64,30min
}
/// Session state for a single `(user, server)` MCP connection.
#[derive(Debug, Clone)]
pub struct McpSession {
/// Session ID returned by the server (via Mcp-Session-Id header).
pub session_id: Option<String>,
/// Last activity timestamp for this session.
pub last_activity: Instant,
/// Server URL this session is connected to.
pub server_url: String,
/// Whether initialization has completed.
pub initialized: bool,
}
2. MCPProcessManager(管理本地server)
/// Manages stdio MCP server processes.
///
/// Handles spawning, tracking, and shutdown of child processes. Keyed
/// by `(user_id, server_name)` so that multiple tenants activating
/// the same server name end up with distinct, independently tracked
/// child processes — see `McpProcessKey` for the rationale.
pub struct McpProcessManager {
transports: RwLock<HashMap<McpProcessKey, Arc<StdioMcpTransport>>>,
configs: RwLock<HashMap<McpProcessKey, StdioSpawnConfig>>,
}
/// Composite key for a stdio MCP child process: the activating user
/// plus the server name. Both fields participate in `Hash` / `Eq` so
/// two users activating the same server name each get — and keep —
/// their own child process instead of one silently overwriting the
/// other's transport handle.
///
/// Stdio MCP servers receive credentials via their spawn `env` map, so
/// sharing a single child across users would leak one tenant's
/// credentials to the other's dispatches. Per-user children are
/// required; the process manager must track them independently.
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
pub struct McpProcessKey {
pub user_id: String,
pub server_name: String,
}
1. "复合键"为啥是 (user, server_name) 而不是 server_name
直觉上 stdio MCP 是个全局的本地服务,似乎一个 server 跑一个进程就够了。但 stdio 模式下,每个用户激活时,凭据是通过 env 注入到子进程的环境变量里的(而不是像 HTTP 那样每次请求带 header)。
所以会出现这种场景:
┌───────┬─────────────────────────────────────────┐
│ 用户 │ 同一个 MCP server github-mcp │
├───────┼─────────────────────────────────────────┤
│ Alice │ spawn 时 env: { GH_TOKEN: alice_token } │
├───────┼─────────────────────────────────────────┤
│ Bob │ spawn 时 env: { GH_TOKEN: bob_token } │
└───────┴─────────────────────────────────────────┘
如果用 server_name 当 key,后到的 Bob 会覆盖 Alice 的子进程:
- Alice 的 child 被 kill
- Bob 的 child 占位
- 之后 Alice 的请求走到 Bob 的进程 → Alice 的请求用 Bob 的 token 调 GitHub
→ 凭据跨租户泄露,而且因为是静默覆盖,审计日志里也看不出来。
2. 为什么不通过请求时切换 env 复用同一个子进程
stdio 是单向字节流(stdin 写、stdout 读),不是按"租户"分流的。一个 child 进程内部,server 端的 process.env.GH_TOKEN 在启动那一刻就固定了——不是按请求切。
要支持多租户,只能多 child。
3. 这跟 HTTP / SSE 模式形成对比
┌───────────┬─────────────────────────────────┬─────────────────────────────────┐
│ Transport │ 凭据怎么传 │ 是否需要每用户一个 child │
├───────────┼─────────────────────────────────┼─────────────────────────────────┤
│ stdio │ spawn 时 env 注入,固定在进程里 │ ✅ 需要 │
├───────────┼─────────────────────────────────┼─────────────────────────────────┤
│ http │ 每次请求带 Authorization header │ ❌ 一个 server 进程服务所有用户 │
├───────────┼─────────────────────────────────┼─────────────────────────────────┤
│ sse │ 同 http │ ❌ │
└───────────┴─────────────────────────────────┴─────────────────────────────────┘
这是为啥前面讲 McpProcessManager 在 stdio 模式才用、http/sse 模式闲置——只有 stdio 有"凭据绑死进程"这个问题。
4. Hash / Eq 参与的意义
McpProcessManager 内部应该是 HashMap<ProcessKey, ChildHandle> 之类:
struct ProcessKey {
owner_id: String,
server_name: String,
}
只有两个字段都参与哈希,两个用户激活同名 server时才会被识别为不同的 key,各自拿到独立的 child,互不污染。
如果 Hash / Eq 只看 server_name(字段存在但没参与哈希),HashMap 行为就 undefined——这是 Rust 里经典的"字段没参与等价性"陷阱,注释把这点点出来是为了显式约束而不是默认行为。
5. 一句话
▎ stdio 模式下,凭据绑死子进程的 env 启动参数,所以必须按 (user, server_name) 复合键管理子进程——否则会跨用户泄露凭据,而且是静默的。HTTP/SSE 模式没这个问题,所以那个模式根本不需要这个 manager。
● 这次问题问得对 —— 这段代码确实只创建了 owner 的 client,多租户确实没在这一步里体现。我先把上下文补齐让你看到全貌,然后解释清楚为什么这"看起来"是漏洞,但其实是有意分层的设计。
Read 2 files (ctrl+o to expand)
● 你说得对 —— 启动期只创建 owner 的 client,多用户 client 不是在这里建的
那段 app.rs:843 的代码确实是 owner 一次性预加载,只为 owner_id 创建 N 个 client(N = enabled_servers 数)。多用户的 client 不在这里建。
那"多用户多 client"到底在哪建?在每个用户运行时调 activate_mcp_server 时:
第二个创建点:extensions/manager.rs:5964
pub async fn activate_mcp_server(&self, name: &str, user_id: &str) -> Result<...> {
...
let client = crate::tools::mcp::create_client_from_config(
server.clone(),
&self.mcp_session_manager,
&self.mcp_process_manager,
Some(Arc::clone(&self.secrets)),
user_id, // ← 烧死调用方 user_id
).await?;
...
}
调用方是谁?任意用户在聊天/CLI 里触发"激活某个 MCP server"。流程大致是:
1. user B 在 web UI 里点 "Activate github-mcp"
2. 后端 handler 调 manager.activate_mcp_server("github-mcp", user_B)
3. 第 5964 行 create_client_from_config(server, ..., user_B) 给 user B 独立烧一个 client
4. 第 6046-6049 行 self.mcp_clients.insert(user_B, "github-mcp", client, ...) 插进 store
5. Tool wrappers 在这一步才被注册到全局 ToolRegistry(如果还没有 —— 因为 wrappers 是 per-server 不是 per-user)
关键架构:wrappers 是 per-server,不重复注册
这是为啥"启动期只为 owner 建"够用的原因。看 client.rs:888-895 那条注释我之前贴过:
▎ Holds a reference to the shared McpClientStore instead of a specific Arc<McpClient> so the same registered wrapper serves every user.
所以:
启动期:
owner 的 client → store[(owner, server)] ✓
wrappers 注册到 ToolRegistry(per-server,一份)✓
运行时 user B 激活:
user B 的 client → store[(user_B, server)] ✓ ← 新增这条 entry
wrappers 已经在 registry 里了 → 不再注册 ← 跳过这一步
dispatch 时:
McpToolWrapper.execute(args, ctx)
→ store.get(ctx.user_id, server_name) ← 这里才按 ctx.user_id 路由
→ ctx.user_id == user_B → user_B 的 client(user_B 的 token、user_B 的 secrets handle)
没有重复创建 wrapper,也没有重复创建 client —— wrappers 是 server 级共享一份,client 是 user 级独立一份。
那 app.rs:843 的 owner 加载有什么意义?
不是"占住位置防别人插",是启动时让 owner 提前可用(不用等到 owner 第一次激活请求)。其他用户仍然需要在各自第一次激活时走 manager.rs:5964 那条路径,没人能"白嫖" owner 那份 client —— 因为
McpClientStore.get(user_B, server_name) 查的是 (user_B, server_name),不会返回 (owner, server_name) 那条 entry。
一个相关的设计护栏:check_surface_conflict
manager.rs:5995-6026 这个保护就是为了防御 ToolRegistry 是全局带来的跨租户隐患:
// Before registering any tool wrappers for this user, fingerprint
// the tool surface the server reported and reject activation if
// another user already has the same `name` active with a
// DIFFERENT surface. The `ToolRegistry` keys wrappers by tool
// name only, so without this check user B's incoming schemas
// would silently shadow user A's
意思是:如果 github-mcp 这个名字已经被 user A 激活了,user B 再去激活同名 server 时会比对其工具表面,如果发现 user B 的 server 报上来的工具集和 user A 的不一样(schema 不同),直接拒绝 —— 因为
ToolRegistry 里 wrapper 是按工具名存的,user B 的 schema 会顶掉 user A 的,导致 user A 看到的是 user B 的 schema。
如果 user B 的 server 工具表面和 user A 完全一致(同一个 MCP server 实现的同一组工具),就放行 —— wrappers 不重复注册,store 里多一条 (user_B, server_name) entry,dispatch 时各自拿到各自的 client。
修正后的完整多租户图景
┌─────────────────────────────────┐
│ McpClientStore (单例) │
│ key=(user_id, server_name) │
│ │
│ (owner, github-mcp) → C₁ │ ← app.rs:843 启动期
│ (alice, github-mcp) → C₂ │ ← manager.rs:5964 运行时
│ (bob, github-mcp) → C₃ │ ← manager.rs:5964 运行时
│ (alice, notion-mcp) → C₄ │ ← manager.rs:5964 运行时
└─────────────────────────────────┘
↑
store.get(ctx.user_id, server_name)
↑
┌─────────────────────────────────┐
│ McpToolWrapper (per-server) │
│ github-mcp_search ─┐ │
│ github-mcp_create ─┼── 共享 │ ← ToolRegistry 按工具名 key
│ notion-mcp_query ─┘ │ 全局只一份 wrapper
└─────────────────────────────────┘
↑
Tool::execute(args, ctx)
↑
ctx.user_id
为什么是这样两层结构(不直接 store[(user, server)] → wrapper)
因为 ToolRegistry 是全局的、keyed by 工具名(不是 (user, tool_name))。如果在 wrapper 里绑死 client,就会出现:
- wrapper github-mcp_search 内部存 Arc<McpClient> → 只能服务一个 user
- user B 激活时,要么 reuse 同一个 wrapper 但改 client(user A 的 wrapper 改了 → user A 后续调用拿 user B 的 client)→ 租户串
- 要么重新注册一个 wrapper github-mcp_search(user B 的)→ ToolRegistry 不允许同名
所以wrapper 必须和 client 解耦 —— wrapper 只持 store,运行时按 ctx 路由 client。这是被 ToolRegistry 的全局性逼出来的设计,不是随便选的。
一句话
▎ 启动期只给 owner 建 client 是设计如此(owner 预热 + wrappers 预注册),多用户的 client 在每个用户运行时调 activate_mcp_server 时建,插进共享 McpClientStore。wrappers 是 per-server
▎ 共享一份(按工具名注册到全局 ToolRegistry),client 是 per-user 独立一份(按 (user_id, server_name) 进 store),dispatch 时 wrapper 用 ctx.user_id 去 store 拿对应的 client —— 这就是"看起来 owner
▎ 一份就够、实际每个用户都自己一份"的原因。
2. 加载注册owner mcp
let mcp_servers_future = {
let secrets_store = self.secrets_store.clone();
let db = self.db.clone();
let mcp_sm = Arc::clone(&mcp_session_manager);
let pm = Arc::clone(&mcp_process_manager);
let owner_id = self.config.owner_id.clone();
async move {
let servers_result =
crate::tools::mcp::config::load_mcp_servers_ready(db.as_deref(), &owner_id)
.await;//得到要mcpserver
match servers_result {
Ok(servers) => {
let enabled: Vec<_> = servers.enabled_servers().cloned().collect();
if !enabled.is_empty() {
tracing::debug!(
"Loading {} configured MCP server(s)...",
enabled.len()
);
}
let mut join_set = tokio::task::JoinSet::new();
for server in enabled {
let mcp_sm = Arc::clone(&mcp_sm);
let secrets = secrets_store.clone();
let pm = Arc::clone(&pm);
let owner_id = owner_id.clone();
join_set.spawn(async move {
let server_name = server.name.clone();
let has_custom_auth_header = server.has_custom_auth_header();
//客户端是每个用户一个
let client = match crate::tools::mcp::create_client_from_config(//根据不同uid、server创建对应的client
server,
&mcp_sm,
&pm,
secrets,
&owner_id,
)
/// Create an `McpClient` from a server configuration, dispatching on the
/// effective transport type.
pub async fn create_client_from_config(
mut server: McpServerConfig,
session_manager: &Arc<McpSessionManager>,
process_manager: &Arc<McpProcessManager>,
secrets: Option<Arc<dyn SecretsStore + Send + Sync>>,
user_id: &str,
) -> Result<McpClient, McpFactoryError> {
match server.effective_transport() {
EffectiveTransport::Stdio { command, args, env } => {//根据不同类型创建不同transport, 然后给客户端
let transport = process_manager
.spawn_stdio(//通过管理器起本地子进程,然后过程在章节一
user_id,
validated_name.as_str(),
command,
args.to_vec(),
env.clone(),
)
.await
.map_err(|e| McpFactoryError::StdioSpawn {
name: server_name.clone(),
reason: e.to_string(),
})?;
Ok(McpClient::new_with_transport(
validated_name.as_str(),
transport as Arc<dyn McpTransport>,
None,
secrets,
user_id,
Some(server),
))
}
EffectiveTransport::Http => {//远程mcpServer需要session_manager管理sessionId
// Authenticated (OAuth) path: tokens exist or server requires auth.
if let Some(ref secrets) = secrets {
let has_tokens =
crate::tools::mcp::is_authenticated(&server, secrets, user_id).await;
if has_tokens || server.requires_auth() {
return Ok(McpClient::new_authenticated(
server,
Arc::clone(session_manager),
Arc::clone(secrets),
user_id,
));
}
}
// Non-OAuth HTTP: wire the session manager into the *transport* so
// it captures `Mcp-Session-Id` from responses. Passing it only to
// the client (via `with_session_manager`) is not enough — the
// transport must know about it to read/write the header.
let transport = Arc::new(
HttpMcpTransport::new(server.url.clone(), validated_name.as_str())
.with_session_manager(Arc::clone(session_manager), user_id),
);
Ok(McpClient::new_with_transport(
validated_name.as_str(),
transport,
Some(Arc::clone(session_manager)),
secrets,
user_id,
Some(server),
))
}
/// MCP transport that communicates with a server over HTTP.
///
/// Sends JSON-RPC requests as HTTP POST with `Content-Type: application/json`
/// and accepts either JSON or SSE (`text/event-stream`) responses. Optionally
/// manages session IDs via [`McpSessionManager`] and supports custom headers.
pub struct HttpMcpTransport {
server_url: String,
/// Typed name so session-manager lookups cannot accidentally be keyed
/// by a free-form string. See `ironclaw_common::identity`.
server_name: McpServerName,
http_client: reqwest::Client,//通过客户端发送接受请求
session_manager: Option<Arc<McpSessionManager>>,
session_user_id: Option<String>,//烧录死,uid-serverName层次单独一个Client
custom_headers: HashMap<String, String>,
}
async fn send(
&self,
request: &McpRequest,
headers: &HashMap<String, String>,
) -> Result<McpResponse, ToolError> {//发送方法
// Build the HTTP request.
let mut req_builder = self
.http_client
.post(&self.server_url)
.header("Content-Type", "application/json")
.header("Accept", "application/json, text/event-stream")
.json(request);
// Apply custom headers configured on the transport.
for (key, value) in &self.custom_headers {
req_builder = req_builder.header(key.as_str(), value.as_str());
}
// Apply per-request headers (e.g. Authorization, Mcp-Session-Id).
for (key, value) in headers {
req_builder = req_builder.header(key.as_str(), value.as_str());
}
// Send the request.
let response = req_builder.send().await.map_err(|e| {
let mut chain = format!("[{}] MCP HTTP request failed: {}", self.server_name, e);
let mut source = std::error::Error::source(&e);
while let Some(cause) = source {
chain.push_str(&format!(" -> {}", cause));
source = cause.source();
}
ToolError::ExternalService(chain)
})?;
// Handle error status codes before accepting any session state from the response.
if !response.status().is_success() {
let status = response.status();
let body = response.text().await.unwrap_or_default();
let sanitized = sanitize_error_body(&body);
return Err(ToolError::ExternalService(format!(
"[{}] MCP server returned status: {} - {}",
self.server_name, status, sanitized
)));
}
// Extract session ID from successful response headers before consuming the body.
// Scope by `(session_user_id, server_name)` so a second user's
// initialize handshake can't overwrite the first user's stored
// session ID and silently redirect their subsequent requests to
// the wrong server-side session.
if let Some(ref session_manager) = self.session_manager
&& let Some(ref user_id) = self.session_user_id
&& let Some(session_id) = response
.headers()
.get("Mcp-Session-Id")
.and_then(|v| v.to_str().ok())
{
let session_id = session_id.trim();
if !is_safe_mcp_session_id(session_id) {
return Err(ToolError::ExternalService(format!(
"[{}] MCP server returned invalid session id",
self.server_name
)));
}
session_manager
.update_session_id(user_id, &self.server_name, Some(session_id.to_string()))
.await;
}
// MCP notifications commonly acknowledge with 202 Accepted and no body.
if response.status() == reqwest::StatusCode::ACCEPTED {
return Ok(McpResponse {
jsonrpc: "2.0".to_string(),
id: request.id,
result: None,
error: None,
});
}
// Determine response format from Content-Type.
let content_type = response
.headers()
.get("content-type")
.and_then(|v| v.to_str().ok())
.unwrap_or("")
.to_string();
if content_type.contains("text/event-stream") {
self.parse_sse_response(response, request.id).await
} else {
response.json().await.map_err(|e| {
ToolError::ExternalService(format!(
"[{}] Failed to parse MCP response: {}",
self.server_name, e
))
})
}
}
pub struct ExtensionManager {
registry: ExtensionRegistry,
discovery: OnlineDiscovery,
// MCP infrastructure
mcp_session_manager: Arc<McpSessionManager>,
mcp_process_manager: Arc<crate::tools::mcp::process::McpProcessManager>,
/// Active MCP clients keyed by `(user, server)`. Shared as `Arc` with
/// every registered `McpToolWrapper` so tool dispatch can resolve the
/// caller's per-user client at execute time instead of embedding a
/// specific client in the globally-registered wrapper (which would
/// let the second activating user's credentials shadow the first).
mcp_clients: Arc<crate::tools::mcp::McpClientStore>, //缓存MCPCLient
/// Per-server async mutex that serialises `activate_mcp` and the
/// `McpServer` arm of `remove` on the same server name. Without this,
/// user B's `remove` (which unregisters the server's global tool
/// wrappers once it's the last user out) can interleave with user C's
/// `activate` (which re-registers the wrappers and inserts C's
/// client), leaving the store with C's client but the registry with
/// C's wrappers already unregistered. Parallelism across *different*
/// servers is preserved.
mcp_lifecycle_locks: RwLock<HashMap<String, Arc<tokio::sync::Mutex<()>>>>,
// WASM tool infrastructure
wasm_tool_runtime: Option<Arc<WasmToolRuntime>>,
wasm_tools_dir: PathBuf,
wasm_channels_dir: PathBuf,
latent_wasm_provider_actions: RwLock<HashMap<String, Vec<LatentProviderAction>>>,
/// Per-server URL cache for `mcp_supports_auth` metadata discovery.
/// Avoids re-issuing a network probe on every `list()` call.
mcp_auth_support_cache: RwLock<HashMap<String, bool>>,
// WASM channel hot-activation infrastructure (set post-construction)
channel_runtime: RwLock<Option<ChannelRuntimeState>>,
/// Channel manager for hot-adding relay channels (set independently of WASM runtime).
relay_channel_manager: RwLock<Option<Arc<ChannelManager>>>,
// Shared
secrets: Arc<dyn SecretsStore + Send + Sync>,
tool_registry: Arc<ToolRegistry>,
hooks: Option<Arc<HookRegistry>>,
pending_auth: RwLock<HashMap<PendingAuthKey, PendingAuth>>,
/// Tunnel URL for webhook configuration and remote OAuth callbacks.
tunnel_url: Option<String>,
user_id: String,
/// Optional database store for DB-backed MCP config.
store: Option<Arc<dyn crate::db::Database>>,
/// When set, settings reads/writes go through this cache-backed store
/// instead of the raw `Database`. Populated via `with_settings_store()`.
settings_override: Option<Arc<dyn crate::db::SettingsStore + Send + Sync>>,
/// Names of WASM/relay channels that are actively running in this process.
///
/// This is runtime state, not install/discovery state. Installed WASM
/// channels are still discovered from disk in `list()`, but only channels
/// present in this set are reported as active.
active_channel_names: RwLock<HashSet<String>>,
/// Installed channel-relay extensions (no on-disk artifact, tracked in memory).
installed_relay_extensions: RwLock<HashSet<String>>,
/// Last activation error for each WASM channel (ephemeral, cleared on success).
activation_errors: RwLock<HashMap<String, String>>,
/// SSE broadcast manager (set post-construction via `set_sse_sender()`).
sse_manager: RwLock<Option<Arc<crate::channels::web::sse::SseManager>>>,
/// Shared registry of pending OAuth flows for gateway-routed callbacks.
///
/// Keyed by CSRF `state` parameter. Populated in `start_wasm_oauth()`
/// when running in gateway mode, consumed by the web gateway's
/// `/oauth/callback` handler.
pending_oauth_flows: crate::auth::oauth::PendingOAuthRegistry,
/// OAuth proxy auth token for authenticating with the hosted token exchange proxy.
/// Resolved once at construction from `IRONCLAW_OAUTH_PROXY_AUTH_TOKEN`,
/// then `GATEWAY_AUTH_TOKEN` as a backward-compatible fallback.
oauth_proxy_auth_token: Option<String>,
/// Relay config captured at startup. Used by `auth_channel_relay` and
/// `activate_channel_relay` instead of re-reading env vars.
relay_config: Option<crate::config::RelayConfig>,
/// Shared event sender for the relay webhook endpoint.
/// Populated by `activate_channel_relay`, consumed by the web gateway's
/// `/relay/events` handler.
relay_event_tx: Arc<
tokio::sync::Mutex<
Option<tokio::sync::mpsc::Sender<crate::channels::relay::client::ChannelEvent>>,
>,
>,
/// Per-instance callback signing secret fetched from channel-relay at activation.
/// Stored here so the web gateway can verify incoming callbacks without
/// any env var or shared secret.
relay_signing_secret_cache: Arc<std::sync::Mutex<Option<Vec<u8>>>>,
/// PairingStore for multi-tenant relay identity resolution.
pairing_store: Option<Arc<crate::pairing::PairingStore>>,
/// When `true`, OAuth flows always return an auth URL to the caller
/// instead of opening a browser on the server via `open::that()`.
/// Set by the web gateway at startup via `enable_gateway_mode()`.
gateway_mode: std::sync::atomic::AtomicBool,
/// Reborn Telegram v2 ProductAdapter (issue #3285) feature flag.
///
/// When `true`, [`Self::activate_wasm_channel`] fails closed on the
/// legacy `telegram` WASM channel — both paths must not handle the
/// same Telegram webhook installation. The runtime-tier startup
/// guard in `main.rs` rejects the same conflict at boot; this
/// post-startup flag closes the hot-activation bypass Henry flagged
/// on PR #3356. Set by the host at startup via
/// [`Self::set_reborn_telegram_v2_enabled`].
reborn_telegram_v2_enabled: std::sync::atomic::AtomicBool,
/// The gateway's own base URL for building OAuth redirect URIs.
/// Set by the web gateway at startup via `enable_gateway_mode()`.
gateway_base_url: RwLock<Option<String>>,
pending_wechat_logins: RwLock<HashMap<String, PendingWechatLogin>>,
channel_activation_locks: RwLock<HashMap<String, Arc<tokio::sync::Mutex<()>>>>,
#[cfg(test)]
test_wasm_channel_loader: RwLock<Option<TestWasmChannelLoader>>,
#[cfg(test)]
test_wechat_login_starter: RwLock<Option<TestWechatLoginStarter>>,
#[cfg(test)]
test_wechat_login_poller: RwLock<Option<TestWechatLoginPoller>>,
}
***************************
/// Per-user MCP client registry. Typically held as `Arc<McpClientStore>`
/// by both `ExtensionManager` (for lifecycle) and every `McpToolWrapper`
/// (for dispatch-time lookup).
#[derive(Default)]
pub struct McpClientStore {uid-sreverName
clients: RwLock<HashMap<McpClientKey, McpClientEntry>>,
}
init_ext内部会调这个 :主要是将工具统一注册(注册前会先初始化、获取工具列表)
let registered = manager
.inject_mcp_client(normalized_name.clone(), &self.config.owner_id, client)
McpClient
pub async fn create_tools_with_store(
&self,
store: Arc<super::McpClientStore>,
) -> Result<Vec<Arc<dyn Tool>>, ToolError> {
let mcp_tools = self.list_tools().await?;//去拉列表
每个tool都包装成一个mcp_wrapper 返回用于toolRegistry的注册
**************
/// List available tools from the MCP server.
pub async fn list_tools(&self) -> Result<Vec<McpTool>, ToolError> {
if let Some(tools) = self.tools_cache.read().await.as_ref() {//看看有没有缓存
return Ok(tools.clone());
}
self.initialize().await?;//初始化链接
let request = McpRequest::list_tools(self.next_request_id());//发送实际机请求,拉取tool_list
let response = self.send_request(request).await?;
if let Some(error) = response.error {
return Err(ToolError::ExternalService(format!(
"MCP error: {} (code {})",
error.message, error.code
)));
}
let result: ListToolsResult = response
.result
.ok_or_else(|| ToolError::ExternalService("No result in MCP response".to_string()))
.and_then(|r| {
serde_json::from_value(r)
.map_err(|e| ToolError::ExternalService(format!("Invalid tools list: {}", e)))
})?;
*self.tools_cache.write().await = Some(result.tools.clone());
Ok(result.tools)
}
**************
/// Initialize the connection to the MCP server.
///
/// Uses `OnceCell` to guarantee that exactly one caller performs the
/// handshake, even under concurrent access. Subsequent calls return
/// immediately.
pub async fn initialize(&self) -> Result<InitializeResult, ToolError> {
let result = self
.initialized
.get_or_try_init(|| async {//仅会初始化一次
if let Some(ref session_manager) = self.session_manager
&& session_manager
.is_initialized(&self.user_id, &self.server_name)
.await
{
return Ok(InitializeResult::default());
}
self.reinitialize_session().await
})
.await?;
Ok(result.clone())
}
********************
重连
/// Re-run the MCP initialize handshake outside the OnceCell cache.
///
/// This is used for recoverable session-expiry failures when an MCP server
/// reports that the current session ID is no longer valid.
async fn reinitialize_session(&self) -> Result<InitializeResult, ToolError> {
if let Some(ref session_manager) = self.session_manager {
session_manager
.terminate(&self.user_id, &self.server_name)
.await;
session_manager
.get_or_create(&self.user_id, &self.server_name, &self.server_url)
.await;
}
let request = McpRequest::initialize(self.next_request_id());
let response = self
.transport
.send(&request, &self.build_request_headers().await?)//借助tranport去发连接信息
.await?;
if let Some(error) = response.error {
return Err(ToolError::ExternalService(format!(
"MCP initialization error: {} (code {})",
error.message, error.code
)));
}
let init_result: InitializeResult = response
.result
.ok_or_else(|| {
ToolError::ExternalService("No result in initialize response".to_string())
})
.and_then(|r| {
serde_json::from_value(r).map_err(|e| {
ToolError::ExternalService(format!("Invalid initialize result: {}", e))
})
})?;
if let Some(ref session_manager) = self.session_manager {
session_manager
.mark_initialized(&self.user_id, &self.server_name)
.await;
}
let notification = McpRequest::initialized_notification();
if let Err(e) = self
.transport
.send(¬ification, &self.build_request_headers().await?)
.await
{
tracing::debug!(
"Failed to send initialized notification to '{}': {}",
self.server_name,
e
);
}
Ok(init_result)
}
*****
async fn build_request_headers(&self) -> Result<HashMap<String, String>, ToolError> {//每次client发送时会读sessionId
let mut headers = self.custom_headers.clone();
// Only inject OAuth token if the user hasn't set a custom Authorization header.
let has_custom_auth = self
.custom_headers
.keys()
.any(|k| k.eq_ignore_ascii_case("authorization"));
if !has_custom_auth && let Some(token) = self.get_access_token().await? {
let trimmed = token.trim();
if !trimmed.is_empty() {
headers.insert("Authorization".to_string(), format!("Bearer {}", trimmed));
}
}
if let Some(ref session_manager) = self.session_manager
&& let Some(session_id) = session_manager
.get_session_id(&self.user_id, &self.server_name)
.await
{
headers.insert("Mcp-Session-Id".to_string(), session_id);
}
Ok(headers)
}
****************
client.initialize()
└─ reinitialize_session()
├─ session_manager.terminate + get_or_create ← 清旧建新(空)
├─ transport.send(initialize_request) ← ① 发出
│ ↓ HTTP 响应回来
│ transport.send() 内部:
│ ├─ 检查 response.headers["Mcp-Session-Id"] ← ② 抓 session ID
│ └─ session_manager.update_session_id(...) ← ③ 存进 manager
├─ 解析 response.result (InitializeResult) ← ④ 只读 JSON-RPC 响应体
├─ session_manager.mark_initialized(...) ← ⑤ 标记握手完成
└─ transport.send(notifications/initialized) ← ⑥ 同样会抓一遍 header
↓
response.headers["Mcp-Session-Id"] ← ② 同上路径
session_manager.update_session_id(...)
*****************client本身有重试,发送,底层还是transport的send,被list_tool和call_tool调用
/// Send a request to the MCP server with auth and session headers.
/// Automatically attempts token refresh on 401 errors (HTTP transports only).
async fn send_request(&self, request: McpRequest) -> Result<McpResponse, ToolError> {
// For non-HTTP transports, just send directly without retry logic
if !self.transport.supports_http_features() {
let headers = self.build_request_headers().await?;
return self.transport.send(&request, &headers).await;
}
// HTTP transport: try up to 2 times (first attempt, then retry after token refresh
// or recoverable session reinitialization).
for attempt in 0..2 {
let headers = self.build_request_headers().await?;
let result = self.transport.send(&request, &headers).await;
**************************too_call
/// Call a tool on the MCP server.
pub async fn call_tool(//在wraper的execute里调用这个,真正的执行方法
&self,
name: &str,
arguments: serde_json::Value,
) -> Result<CallToolResult, ToolError> {
self.initialize().await?;
let request = McpRequest::call_tool(self.next_request_id(), name, arguments);
let response = self.send_request(request).await?;
/// Session state for a single `(user, server)` MCP connection.
#[derive(Debug, Clone)]
pub struct McpSession {
/// Session ID returned by the server (via Mcp-Session-Id header).//注意sessionId
pub session_id: Option<String>,
/// Last activity timestamp for this session.
pub last_activity: Instant,
/// Server URL this session is connected to.
pub server_url: String,
/// Whether initialization has completed.
pub initialized: bool,
}
/// Check if a session is initialized.
pub async fn is_initialized(&self, user_id: &str, server_name: &McpServerName) -> bool {检查是否初始化
let sessions = self.sessions.read().await;
sessions
.get(&Self::key(user_id, server_name))
.map(|s| s.initialized)
.unwrap_or(false)
}
5. McpWrapper
一个 McpToolWrapper 只包装一个 McpTool,一个 MCP Server 的多个 tool 会生成多个 wrapper。
struct McpToolWrapper {
tool: McpTool,//工具信息
prefixed_name: String,
provider_extension: String,
server_name: String,
client_store: Arc<super::McpClientStore>,实际借助Client去发起调用
}
async fn execute(
&self,
params: serde_json::Value,
ctx: &JobContext,
) -> Result<ToolOutput, ToolError> {
let start = std::time::Instant::now();
// Strip top-level null values before forwarding — LLMs often emit
// `"field": null` for optional params, but many MCP servers reject
// explicit nulls for fields that should simply be absent.
let params = strip_top_level_nulls(params);
let client = self
.client_store
.get(&ctx.user_id, &self.server_name)
.await
.ok_or_else(|| {
ToolError::ExternalService(format!(
"MCP server '{}' is not active for this user",
self.server_name
))
})?;
let result = client.call_tool(&self.tool.name, params).await?;//这里发起调用
let content: String = result
.content
.iter()
.filter_map(|b| b.as_text())
.collect::<Vec<_>>()
.join("\n");
if result.is_error {
return Err(ToolError::ExecutionFailed(content));
}
Ok(ToolOutput::text(content, start.elapsed()))
}
/// MCP client for communicating with MCP servers.
///
/// Supports multiple transport types:
/// - HTTP: For remote MCP servers (created via `new`, `new_with_name`, `new_authenticated`)
/// - Stdio/Unix: Via `new_with_transport` with a custom `McpTransport` implementation
pub struct McpClient {
/// Transport for sending requests.
transport: Arc<dyn McpTransport>,
/// Server URL (kept for accessor compatibility).
server_url: String,
/// Server name (for logging and session management).
///
/// Typed via `McpServerName` so session-manager lookups and logging
/// are both compile-time-gated behind the allowlist validation that
/// lives in `ironclaw_common::identity`.
server_name: McpServerName,
/// Request ID counter.
next_id: AtomicU64,
/// Cached tools.
tools_cache: RwLock<Option<Vec<McpTool>>>,
/// Session manager (shared across clients).
session_manager: Option<Arc<McpSessionManager>>,
/// Secrets store for retrieving access tokens.
secrets: Option<Arc<dyn SecretsStore + Send + Sync>>,
/// User ID for secrets lookup.
user_id: String,
/// Server configuration (for token secret name lookup).
server_config: Option<McpServerConfig>,
/// Custom headers to include in every request.
custom_headers: HashMap<String, String>,
/// Ensures the MCP initialize handshake runs exactly once.
/// Uses `OnceCell` to serialize concurrent callers so only one
/// actually sends the request; subsequent calls return immediately.
initialized: tokio::sync::OnceCell<InitializeResult>,
/// Test-only marker recording which constructor produced this client.
/// Used by caller-level tests to assert the factory chose the correct path.
#[cfg(test)]
constructor_kind: McpClientConstructor,
}