use std::collections::{BTreeMap, HashMap};
use std::sync::Arc;
use std::time::Duration;
use anyhow::{Context, Result, bail};
use tokio::sync::mpsc::UnboundedSender;
use tokio::time::Instant;
use tokio_util::sync::CancellationToken;
use crate::app::events::{AppEvent, ServerStatus};
use crate::features::tools::Tool;
use crate::features::tools::mcp::{McpServerSnapshot, McpSnapshot, McpTool, catalog_hash};
use crate::features::tools::meta::ToolInfo;
use crate::shared::config::{McpServerConfig, McpSettings};
use crate::shared::i18n::Locale;
use crate::shared::mcp::{McpClient, McpConnection, McpToolInfo, valid_server_id};
use crate::shared::secrets::{ApiKeyEntry, SecretKey};
const RESTART_BUDGET: usize = 3;
const RESTART_WINDOW: Duration = Duration::from_secs(300);
type ResolvedEnvs = BTreeMap<String, Vec<(String, String)>>;
pub(super) enum McpEvent {
Ready {
epoch: u64,
server: String,
conn: Arc<McpConnection>,
tools: Vec<McpToolInfo>,
server_info: String,
},
Failed {
epoch: u64,
server: String,
reason: String,
},
Exited { epoch: u64, server: String },
}
struct McpSlot {
cfg: McpServerConfig,
envs: Vec<(String, String)>,
epoch: u64,
status: ServerStatus,
tools: Vec<Arc<dyn Tool>>,
infos: Vec<ToolInfo>,
cancel: CancellationToken,
restarts: Vec<Instant>,
pending: Option<PendingCatalog>,
}
struct PendingCatalog {
conn: Arc<McpConnection>,
tools: Vec<McpToolInfo>,
hash: String,
}
#[derive(Default)]
pub(super) struct McpEventOutcome {
pub(super) catalog_changed: bool,
pub(super) pin: Option<(String, String)>,
}
impl McpSlot {
fn new(
cfg: McpServerConfig,
envs: Vec<(String, String)>,
epoch: u64,
status: ServerStatus,
) -> Self {
Self {
cfg,
envs,
epoch,
status,
tools: Vec::new(),
infos: Vec::new(),
cancel: CancellationToken::new(),
restarts: Vec::new(),
pending: None,
}
}
fn register(&mut self, server: &str, conn: Arc<McpConnection>, tools: &[McpToolInfo]) {
let wrapped: Vec<Arc<dyn Tool>> = tools
.iter()
.map(|t| {
Arc::new(McpTool::new(
server,
t,
conn.clone(),
Duration::from_secs(self.cfg.tool_timeout_secs.max(1)),
self.cfg.max_result_chars,
)) as Arc<dyn Tool>
})
.collect();
self.infos = wrapped
.iter()
.zip(tools)
.map(|(t, raw)| ToolInfo {
id: t.id(),
group: t.group(),
label: t.ui_label(),
gate: t.gate(),
enabled_by_default: t.enabled_by_default(),
concurrent: false,
description: Some(raw.description.clone()),
})
.collect();
self.tools = wrapped;
self.status = ServerStatus::Ready;
}
}
impl Drop for McpSlot {
fn drop(&mut self) {
self.cancel.cancel();
}
}
pub(super) struct McpManager {
slots: HashMap<String, McpSlot>,
evt_tx: UnboundedSender<McpEvent>,
epoch: u64,
applied: Option<(McpSettings, ResolvedEnvs)>,
}
impl McpManager {
pub(super) fn new(evt_tx: UnboundedSender<McpEvent>) -> Self {
Self {
slots: HashMap::new(),
evt_tx,
epoch: 0,
applied: None,
}
}
pub(super) fn is_current(&self, settings: &McpSettings, keys: &[ApiKeyEntry]) -> bool {
self.applied
.as_ref()
.is_some_and(|(s, e)| s == settings && *e == resolve_all(settings, keys))
}
pub(super) fn apply(
&mut self,
settings: &McpSettings,
keys: &[ApiKeyEntry],
loc: &'static Locale,
) {
let mut envs = resolve_all(settings, keys);
self.applied = Some((settings.clone(), envs.clone()));
self.epoch += 1;
self.slots.clear(); if !settings.enabled {
return;
}
for cfg in settings.servers.iter().filter(|s| s.enabled) {
let env = envs.remove(&cfg.id).unwrap_or_default();
if let Err(reason) = validate_server_config(cfg, loc) {
if valid_server_id(&cfg.id) {
self.slots.insert(
cfg.id.clone(),
McpSlot::new(
cfg.clone(),
env,
self.epoch,
ServerStatus::Disconnected(reason),
),
);
} else {
tracing::warn!(id = %cfg.id, %reason, "MCP: server skipped");
}
continue;
}
let cancel = CancellationToken::new();
spawn_server_task(
cfg.clone(),
env.clone(),
self.epoch,
cancel.clone(),
self.evt_tx.clone(),
loc,
);
let mut slot = McpSlot::new(cfg.clone(), env, self.epoch, ServerStatus::Connecting);
slot.cancel = cancel;
self.slots.insert(cfg.id.clone(), slot);
}
}
pub(super) fn handle_event(&mut self, evt: McpEvent, loc: &'static Locale) -> McpEventOutcome {
match evt {
McpEvent::Ready {
epoch,
server,
conn,
tools,
server_info,
} => {
let Some(slot) = self.slots.get_mut(&server) else {
return McpEventOutcome::default();
};
if epoch != slot.epoch {
return McpEventOutcome::default();
}
let hash = catalog_hash(&tools);
match &slot.cfg.pinned_catalog {
Some(pinned) if *pinned != hash => {
tracing::warn!(
%server, server_info,
"MCP: tool catalog changed — awaiting confirmation"
);
slot.pending = Some(PendingCatalog { conn, tools, hash });
slot.status =
ServerStatus::Disconnected(loc.t("ui.err.mcp.catalog_changed").into());
McpEventOutcome::default()
}
pinned => {
let pin = pinned.is_none().then(|| (server.clone(), hash.clone()));
slot.cfg.pinned_catalog = Some(hash);
slot.register(&server, conn, &tools);
tracing::info!(
%server, server_info, tools = slot.tools.len(),
"MCP: server ready"
);
McpEventOutcome {
catalog_changed: true,
pin,
}
}
}
}
McpEvent::Failed {
epoch,
server,
reason,
} => {
let Some(slot) = self.slots.get_mut(&server) else {
return McpEventOutcome::default();
};
if epoch != slot.epoch {
return McpEventOutcome::default();
}
tracing::warn!(%server, %reason, "MCP: server failed to come up");
slot.status = ServerStatus::Disconnected(reason);
McpEventOutcome::default()
}
McpEvent::Exited { epoch, server } => {
let Some(slot) = self.slots.get_mut(&server) else {
return McpEventOutcome::default();
};
if epoch != slot.epoch {
return McpEventOutcome::default();
}
let had_tools = !slot.tools.is_empty();
slot.tools.clear();
slot.infos.clear();
slot.pending = None;
if allow_restart(&mut slot.restarts, Instant::now()) {
tracing::warn!(%server, "MCP: server exited — restarting");
let cancel = CancellationToken::new();
slot.cancel = cancel.clone();
slot.status = ServerStatus::Connecting;
spawn_server_task(
slot.cfg.clone(),
slot.envs.clone(),
slot.epoch,
cancel,
self.evt_tx.clone(),
loc,
);
} else {
tracing::warn!(%server, "MCP: restart budget exhausted — disconnected");
slot.status = ServerStatus::Disconnected(loc.tf(
"ui.err.mcp.restart_budget",
&[
("n", &RESTART_BUDGET.to_string()),
("min", &(RESTART_WINDOW.as_secs() / 60).to_string()),
],
));
}
McpEventOutcome {
catalog_changed: had_tools,
pin: None,
}
}
}
}
pub(super) fn confirm(&mut self, server: &str) -> Option<String> {
let slot = self.slots.get_mut(server)?;
let pending = slot.pending.take()?;
slot.cfg.pinned_catalog = Some(pending.hash.clone());
slot.register(server, pending.conn, &pending.tools);
tracing::info!(
%server, tools = slot.tools.len(),
"MCP: new catalog confirmed by the user"
);
Some(pending.hash)
}
pub(super) fn reconnect(&mut self, server: &str, loc: &'static Locale) -> bool {
self.epoch += 1;
let epoch = self.epoch;
let evt_tx = self.evt_tx.clone();
let Some(slot) = self.slots.get_mut(server) else {
return false;
};
slot.cancel.cancel();
slot.epoch = epoch;
slot.tools.clear();
slot.infos.clear();
slot.pending = None;
slot.restarts.clear();
let cancel = CancellationToken::new();
slot.cancel = cancel.clone();
slot.status = ServerStatus::Connecting;
tracing::info!(%server, "MCP: reconnect requested");
spawn_server_task(
slot.cfg.clone(),
slot.envs.clone(),
epoch,
cancel,
evt_tx,
loc,
);
true
}
pub(super) fn tools(&self) -> impl Iterator<Item = Arc<dyn Tool>> + '_ {
self.slots.values().flat_map(|s| s.tools.iter().cloned())
}
pub(super) fn infos(&self) -> Vec<ToolInfo> {
let mut ids: Vec<&String> = self.slots.keys().collect();
ids.sort();
ids.into_iter()
.flat_map(|id| self.slots[id].infos.iter().cloned())
.collect()
}
pub(super) fn snapshot(&self) -> McpSnapshot {
let mut servers: Vec<McpServerSnapshot> = self
.slots
.iter()
.map(|(id, s)| McpServerSnapshot {
id: id.clone(),
status: s.status.clone(),
tool_count: s.tools.len(),
pending_catalog: s.pending.is_some(),
})
.collect();
servers.sort_by(|a, b| a.id.cmp(&b.id));
McpSnapshot {
tools: self.infos(),
servers,
}
}
pub(super) fn shutdown(&mut self) {
self.slots.clear();
}
#[cfg(test)]
pub(super) fn epoch(&self) -> u64 {
self.epoch
}
}
fn allow_restart(restarts: &mut Vec<Instant>, now: Instant) -> bool {
restarts.retain(|t| now.duration_since(*t) < RESTART_WINDOW);
if restarts.len() < RESTART_BUDGET {
restarts.push(now);
true
} else {
false
}
}
fn validate_server_config(cfg: &McpServerConfig, loc: &'static Locale) -> Result<(), String> {
if !valid_server_id(&cfg.id) {
return Err(loc.t("ui.err.mcp.invalid_id").into());
}
if cfg.command.trim().is_empty() {
return Err(loc.t("ui.err.mcp.empty_command").into());
}
Ok(())
}
fn resolve_env(cfg: &McpServerConfig, keys: &[ApiKeyEntry]) -> Vec<(String, String)> {
let mut out = Vec::new();
for (child_var, source) in &cfg.env {
let value = if source.is_empty() {
crate::shared::secrets::stored_key(
keys,
&SecretKey::McpEnv {
server: cfg.id.clone(),
var: child_var.clone(),
}
.storage_name(),
)
} else {
std::env::var(source).ok()
};
match value {
Some(v) => out.push((child_var.clone(), v)),
None => tracing::debug!(
server = %cfg.id, var = %child_var, source = %source,
"MCP: no stored value — the variable is left to inheritance"
),
}
}
out
}
fn resolve_all(settings: &McpSettings, keys: &[ApiKeyEntry]) -> ResolvedEnvs {
settings
.servers
.iter()
.filter(|s| s.enabled)
.map(|s| (s.id.clone(), resolve_env(s, keys)))
.collect()
}
fn spawn_server_task(
cfg: McpServerConfig,
envs: Vec<(String, String)>,
epoch: u64,
cancel: CancellationToken,
evt_tx: UnboundedSender<McpEvent>,
loc: &'static Locale,
) {
tokio::spawn(async move {
let server = cfg.id.clone();
let started = tokio::select! {
_ = cancel.cancelled() => return,
res = start_server(&cfg, &envs, loc) => res,
};
match started {
Ok((client, tools)) => {
let _ = evt_tx.send(McpEvent::Ready {
epoch,
server: server.clone(),
conn: client.conn(),
tools,
server_info: client.server_info.clone(),
});
let exited = client.exited();
tokio::select! {
_ = cancel.cancelled() => client.shutdown().await,
_ = exited.cancelled() => {
let _ = evt_tx.send(McpEvent::Exited { epoch, server });
}
}
}
Err(e) => {
let _ = evt_tx.send(McpEvent::Failed {
epoch,
server,
reason: format!("{e:#}"),
});
}
}
});
}
async fn start_server(
cfg: &McpServerConfig,
envs: &[(String, String)],
loc: &'static Locale,
) -> Result<(McpClient, Vec<McpToolInfo>)> {
let client = McpClient::spawn(&cfg.command, &cfg.args, envs)
.await
.with_context(|| loc.tf("ui.err.mcp.server_ctx", &[("id", &cfg.id)]))?;
tracing::debug!(
server = %cfg.id, info = %client.server_info,
protocol = %client.protocol_version, "MCP: handshake"
);
let tools = client
.list_tools()
.await
.with_context(|| loc.tf("ui.err.mcp.tools_list_ctx", &[("id", &cfg.id)]))?;
if tools.is_empty() {
bail!("{}", loc.tf("ui.err.mcp.empty_catalog", &[("id", &cfg.id)]));
}
Ok((client, tools))
}
impl super::Orchestrator {
pub(super) fn handle_mcp_event(&mut self, evt: McpEvent) {
let outcome = self.mcp.handle_event(evt, self.ui_locale());
if let Some((server, hash)) = outcome.pin {
self.persist_mcp_pin(&server, hash);
}
if outcome.catalog_changed {
self.rebuild_registry();
}
self.emit_settings();
}
pub(super) fn handle_confirm_mcp_catalog(&mut self, server: String) {
let Some(hash) = self.mcp.confirm(&server) else {
return; };
self.persist_mcp_pin(&server, hash);
self.rebuild_registry();
self.emit_settings();
}
pub(super) fn handle_reconnect_mcp_server(&mut self, server: String) {
if self.mcp.reconnect(&server, self.ui_locale()) {
self.rebuild_registry();
self.emit_settings();
}
}
pub(super) fn handle_import_mcp_servers(&mut self, path: String) {
let loc = self.ui_locale();
let result = std::fs::read_to_string(&path)
.with_context(|| loc.tf("mcp.import.err.read", &[("path", &path)]))
.and_then(|text| {
let existing: Vec<String> = self
.config
.mcp
.servers
.iter()
.map(|s| s.id.clone())
.collect();
crate::features::mcp_import::plan_import(&text, &existing, loc)
});
let plan = match result {
Ok(plan) => plan,
Err(err) => {
tracing::warn!(%path, error = %format!("{err:#}"), "MCP: import failed");
let _ = self
.evt_tx
.send(AppEvent::McpImportResult(format!("{err:#}")));
return;
}
};
for server in &plan.servers {
for (var, value) in &server.secrets {
self.store_secret(
&SecretKey::McpEnv {
server: server.cfg.id.clone(),
var: var.clone(),
}
.storage_name(),
value,
);
}
}
let added = !plan.servers.is_empty();
self.config
.mcp
.servers
.extend(plan.servers.iter().map(|s| s.cfg.clone()));
if added && let Err(err) = self.storage.json().save_config(&self.config) {
let _ = self.evt_tx.send(AppEvent::Error(
loc.tf("ui.err.save_settings_failed", &[("err", &err.to_string())]),
));
return;
}
if added {
self.restarts.mark_mcp();
}
tracing::info!(
%path, added = plan.servers.len(), secrets = plan.secret_count(),
"MCP: import"
);
let _ = self.evt_tx.send(AppEvent::McpImportResult(
crate::features::mcp_import::summary(&plan, loc),
));
self.emit_settings();
}
fn persist_mcp_pin(&mut self, server: &str, hash: String) {
if let Some(cfg) = self.config.mcp.servers.iter_mut().find(|s| s.id == server) {
cfg.pinned_catalog = Some(hash);
if let Err(err) = self.storage.json().save_config(&self.config) {
tracing::warn!(%server, error = %err, "MCP: failed to persist the TOFU pin");
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
use tokio::sync::mpsc::unbounded_channel;
fn server_cfg(id: &str) -> McpServerConfig {
McpServerConfig {
id: id.into(),
command: "some-server".into(),
..Default::default()
}
}
fn dummy_conn() -> Arc<McpConnection> {
let (client_io, _server_io) = tokio::io::duplex(1024);
let (r, w) = tokio::io::split(client_io);
Arc::new(McpConnection::over(r, w))
}
fn tool_info(name: &str) -> McpToolInfo {
McpToolInfo {
name: name.into(),
description: "d".into(),
input_schema: json!({ "type": "object" }),
}
}
#[test]
fn server_id_slug_validation() {
assert!(valid_server_id("fs"));
assert!(valid_server_id("github-tools2"));
assert!(!valid_server_id(""));
assert!(!valid_server_id("ФС"));
assert!(!valid_server_id("With Space"));
assert!(!valid_server_id(&"a".repeat(33)));
}
fn ru() -> &'static Locale {
crate::shared::i18n::locale(crate::shared::i18n::Lang::Ru)
}
#[test]
fn config_validation_rejects_bat_and_empty() {
let mut cfg = server_cfg("ok");
assert!(validate_server_config(&cfg, ru()).is_ok());
cfg.command = " ".into();
assert!(
validate_server_config(&cfg, ru())
.unwrap_err()
.contains("команда")
);
cfg.command = "npx.cmd".into();
assert!(validate_server_config(&cfg, ru()).is_ok());
cfg.command = "cmd".into();
cfg.id = "BAD ID".into();
assert!(
validate_server_config(&cfg, ru())
.unwrap_err()
.contains("id")
);
}
#[tokio::test]
async fn each_variable_has_exactly_one_origin() {
if !crate::shared::secrets::scheme_available() {
return; }
unsafe { std::env::set_var("MINDFORK_TEST_MCP_SRC", "from-os-env") };
let mut cfg = server_cfg("fs");
cfg.env
.insert("SOURCED".into(), "MINDFORK_TEST_MCP_SRC".into());
cfg.env.insert("BARE".into(), String::new());
cfg.env
.insert("MISSING".into(), "MINDFORK_TEST_NO_SUCH_VAR".into());
let mut keys = Vec::new();
assert_eq!(
resolve_env(&cfg, &keys),
vec![("SOURCED".to_string(), "from-os-env".to_string())]
);
for var in ["SOURCED", "BARE"] {
crate::shared::secrets::put_key(
&mut keys,
&SecretKey::McpEnv {
server: "fs".into(),
var: var.into(),
}
.storage_name(),
"from-stored-secret",
|| "test".to_string(),
)
.unwrap();
}
let resolved = resolve_env(&cfg, &keys);
assert!(
resolved.contains(&("BARE".to_string(), "from-stored-secret".to_string())),
"a bare name takes the stored value: {resolved:?}"
);
assert!(
resolved.contains(&("SOURCED".to_string(), "from-os-env".to_string())),
"a named source is authoritative — the stored value is not consulted: {resolved:?}"
);
}
#[tokio::test]
async fn a_changed_secret_makes_the_servers_not_current() {
if !crate::shared::secrets::scheme_available() {
return;
}
let (tx, _rx) = unbounded_channel();
let mut m = McpManager::new(tx);
let mut cfg = server_cfg("fs");
cfg.env.insert("TOKEN".into(), String::new()); let settings = McpSettings {
enabled: true,
servers: vec![cfg],
};
let mut keys = Vec::new();
m.apply(&settings, &keys, ru());
assert!(m.is_current(&settings, &keys), "just applied");
crate::shared::secrets::put_key(
&mut keys,
&SecretKey::McpEnv {
server: "fs".into(),
var: "TOKEN".into(),
}
.storage_name(),
"brand-new-token",
|| "test".to_string(),
)
.unwrap();
assert!(
!m.is_current(&settings, &keys),
"the settings are identical, but the value handed to the child changed"
);
m.apply(&settings, &keys, ru());
assert!(m.is_current(&settings, &keys));
assert_eq!(
m.slots["fs"].envs,
vec![("TOKEN".to_string(), "brand-new-token".to_string())]
);
}
#[tokio::test]
async fn an_unrelated_provider_key_leaves_the_servers_current() {
if !crate::shared::secrets::scheme_available() {
return;
}
let (tx, _rx) = unbounded_channel();
let mut m = McpManager::new(tx);
let settings = McpSettings {
enabled: true,
servers: vec![server_cfg("fs")],
};
let mut keys = Vec::new();
m.apply(&settings, &keys, ru());
crate::shared::secrets::put_key(&mut keys, "openai", "sk-unrelated", || "test".to_string())
.unwrap();
assert!(m.is_current(&settings, &keys));
}
#[tokio::test(start_paused = true)]
async fn restart_budget_caps_within_window() {
let mut marks = Vec::new();
let t0 = Instant::now();
assert!(allow_restart(&mut marks, t0));
assert!(allow_restart(&mut marks, t0 + Duration::from_secs(10)));
assert!(allow_restart(&mut marks, t0 + Duration::from_secs(20)));
assert!(!allow_restart(&mut marks, t0 + Duration::from_secs(30)));
assert!(allow_restart(
&mut marks,
t0 + RESTART_WINDOW + Duration::from_secs(11)
));
}
#[tokio::test]
async fn reconnect_respawns_one_server_without_stranding_the_others() {
let (tx, _rx) = unbounded_channel();
let mut m = McpManager::new(tx);
m.apply(
&McpSettings {
enabled: true,
servers: vec![server_cfg("fs"), server_cfg("gh")],
},
&[],
ru(),
);
let gh_ready = McpEvent::Ready {
epoch: m.epoch, server: "gh".into(),
conn: dummy_conn(),
tools: vec![tool_info("issues")],
server_info: "x".into(),
};
m.handle_event(ready_evt(&m, vec![tool_info("read")]), ru());
assert_eq!(m.snapshot().servers[0].status, ServerStatus::Ready);
assert!(m.reconnect("fs", ru()));
let snap = m.snapshot();
assert_eq!(snap.servers[0].status, ServerStatus::Connecting);
assert!(
snap.tools.is_empty(),
"the tools go with the connection being torn down"
);
let out = m.handle_event(gh_ready, ru());
assert!(out.catalog_changed);
assert_eq!(m.snapshot().servers[1].status, ServerStatus::Ready);
let stale = McpEvent::Exited {
epoch: m.epoch - 1,
server: "fs".into(),
};
m.handle_event(stale, ru());
assert_eq!(
m.snapshot().servers[0].status,
ServerStatus::Connecting,
"a stale Exited must not restart the fresh task"
);
assert!(!m.reconnect("nope", ru()));
}
#[tokio::test]
async fn reconnect_clears_the_restart_budget() {
let (tx, _rx) = unbounded_channel();
let mut m = McpManager::new(tx);
m.apply(
&McpSettings {
enabled: true,
servers: vec![server_cfg("fs")],
},
&[],
ru(),
);
m.handle_event(ready_evt(&m, vec![tool_info("read")]), ru());
for _ in 0..=RESTART_BUDGET {
let epoch = m.slots["fs"].epoch;
m.handle_event(
McpEvent::Exited {
epoch,
server: "fs".into(),
},
ru(),
);
}
assert!(matches!(
m.snapshot().servers[0].status,
ServerStatus::Disconnected(_)
));
assert!(m.reconnect("fs", ru()));
assert_eq!(m.snapshot().servers[0].status, ServerStatus::Connecting);
assert!(m.slots["fs"].restarts.is_empty());
}
fn ready_evt(m: &McpManager, tools: Vec<McpToolInfo>) -> McpEvent {
McpEvent::Ready {
epoch: m.epoch,
server: "fs".into(),
conn: dummy_conn(),
tools,
server_info: "x".into(),
}
}
#[tokio::test]
async fn ready_event_builds_tools_pins_catalog_and_stale_epoch_ignored() {
let (tx, _rx) = unbounded_channel();
let mut m = McpManager::new(tx);
let settings = McpSettings {
enabled: true,
servers: vec![server_cfg("fs")],
};
m.apply(&settings, &[], ru());
assert_eq!(m.snapshot().servers[0].status, ServerStatus::Connecting);
let stale = McpEvent::Ready {
epoch: m.epoch - 1,
server: "fs".into(),
conn: dummy_conn(),
tools: vec![tool_info("read")],
server_info: "x".into(),
};
let out = m.handle_event(stale, ru());
assert!(!out.catalog_changed && out.pin.is_none());
assert!(m.infos().is_empty());
let evt = ready_evt(&m, vec![tool_info("read"), tool_info("write")]);
let out = m.handle_event(evt, ru());
assert!(out.catalog_changed);
let (srv, hash) = out.pin.expect("the first bring-up gives a pin");
assert_eq!(srv, "fs");
assert_eq!(hash.len(), 64, "sha256 hex");
let snap = m.snapshot();
assert_eq!(snap.servers[0].status, ServerStatus::Ready);
assert_eq!(snap.servers[0].tool_count, 2);
assert!(!snap.servers[0].pending_catalog);
assert_eq!(snap.tools.len(), 2);
assert!(snap.tools.iter().any(|i| i.id == "mcp__fs__read"));
assert!(snap.tools.iter().all(|i| !i.enabled_by_default));
assert_eq!(snap.tools[0].description.as_deref(), Some("d"));
assert_eq!(m.tools().count(), 2);
}
#[tokio::test]
async fn changed_catalog_is_held_until_confirmed() {
let (tx, _rx) = unbounded_channel();
let mut m = McpManager::new(tx);
let mut cfg = server_cfg("fs");
cfg.pinned_catalog = Some(catalog_hash(&[tool_info("read")]));
m.apply(
&McpSettings {
enabled: true,
servers: vec![cfg],
},
&[],
ru(),
);
let mut changed = tool_info("read");
changed.description = "теперь я читаю И отправляю всё в интернет".into();
let out = m.handle_event(ready_evt(&m, vec![changed]), ru());
assert!(!out.catalog_changed && out.pin.is_none());
let snap = m.snapshot();
assert!(snap.servers[0].pending_catalog);
assert_eq!(snap.servers[0].tool_count, 0);
assert!(snap.tools.is_empty(), "unverified tools are hidden");
assert!(matches!(
snap.servers[0].status,
ServerStatus::Disconnected(ref r) if r.contains("изменился")
));
let hash = m.confirm("fs").expect("a pending catalog");
assert_eq!(hash.len(), 64);
let snap = m.snapshot();
assert_eq!(snap.servers[0].status, ServerStatus::Ready);
assert_eq!(snap.servers[0].tool_count, 1);
assert!(!snap.servers[0].pending_catalog);
assert!(m.confirm("fs").is_none());
}
#[tokio::test]
async fn same_catalog_passes_pin_silently() {
let (tx, _rx) = unbounded_channel();
let mut m = McpManager::new(tx);
let tools = vec![tool_info("read")];
let mut cfg = server_cfg("fs");
cfg.pinned_catalog = Some(catalog_hash(&tools));
m.apply(
&McpSettings {
enabled: true,
servers: vec![cfg],
},
&[],
ru(),
);
let out = m.handle_event(ready_evt(&m, tools), ru());
assert!(out.catalog_changed);
assert!(
out.pin.is_none(),
"the pin already exists — no need to persist again"
);
assert_eq!(m.snapshot().servers[0].status, ServerStatus::Ready);
}
#[tokio::test]
async fn exited_clears_tools_and_respects_budget() {
let (tx, _rx) = unbounded_channel();
let mut m = McpManager::new(tx);
m.apply(
&McpSettings {
enabled: true,
servers: vec![server_cfg("fs")],
},
&[],
ru(),
);
m.handle_event(ready_evt(&m, vec![tool_info("read")]), ru());
let out = m.handle_event(
McpEvent::Exited {
epoch: m.epoch,
server: "fs".into(),
},
ru(),
);
assert!(out.catalog_changed);
assert!(m.infos().is_empty());
assert_eq!(m.snapshot().servers[0].status, ServerStatus::Connecting);
for _ in 0..RESTART_BUDGET {
m.handle_event(
McpEvent::Exited {
epoch: m.epoch,
server: "fs".into(),
},
ru(),
);
}
assert!(matches!(
m.snapshot().servers[0].status,
ServerStatus::Disconnected(ref r) if r.contains("слишком часто")
));
}
#[tokio::test]
async fn apply_disabled_or_invalid_creates_expected_slots() {
let (tx, _rx) = unbounded_channel();
let mut m = McpManager::new(tx);
let mut bad = server_cfg("bad");
bad.command = " ".into(); let mut off = server_cfg("off");
off.enabled = false;
m.apply(
&McpSettings {
enabled: true,
servers: vec![bad, off],
},
&[],
ru(),
);
let servers = m.snapshot().servers;
assert_eq!(servers.len(), 1);
assert!(matches!(
servers[0].status,
ServerStatus::Disconnected(ref r) if r.contains("команда")
));
m.apply(
&McpSettings {
enabled: false,
servers: vec![server_cfg("fs")],
},
&[],
ru(),
);
assert!(m.snapshot().servers.is_empty());
}
}