Skip to content

Commit dcc2977

Browse files
committed
fix(gateway): reconnect after config save and hold workspace pin across initialize
Evict-only left enabled servers bound-but-offline after a definition write, and stdio OAuth reconnect still built an HTTP transport. Hold a non-empty X-Mcpmux-Workspace until mcp-session-id exists so multi-root PendingRoots does not win after a gateway reload. Signed-off-by: crimsonsunset <jsangio1@gmail.com>
1 parent 2da2f50 commit dcc2977

12 files changed

Lines changed: 843 additions & 76 deletions

crates/mcpmux-gateway/src/consumers/server_config_handler.rs

Lines changed: 263 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -1,31 +1,38 @@
1-
//! Server Config Event Handler - Evicts pool instances on config changes
1+
//! Server Config Event Handler - Reconnects pool instances on config changes
22
//!
3-
//! Listens for `ServerConfigUpdated` and removes the pooled instance for
4-
//! enabled servers so the next connect rebuilds transport from DB.
3+
//! Listens for `ServerConfigUpdated` and runs `reconnect_fresh` for
4+
//! enabled servers so the next call uses transport rebuilt from DB.
55
66
use mcpmux_core::{DomainEvent, InstalledServerRepository};
7+
use std::path::PathBuf;
78
use std::sync::Arc;
9+
use std::time::Instant;
810
use tokio::sync::broadcast;
911
use tracing::{debug, info, warn};
1012
use uuid::Uuid;
1113

12-
use crate::pool::PoolService;
14+
use crate::pool::transport::resolution::resolve_auto_connection_context;
15+
use crate::pool::{ConnectionResult, PoolService};
1316

14-
/// Evicts pooled server instances when configuration changes.
17+
/// Evicts and reconnects pooled server instances when configuration changes.
1518
pub struct ServerConfigUpdatedHandler {
1619
installed_server_repo: Arc<dyn InstalledServerRepository + Send + Sync>,
1720
pool_service: Arc<PoolService>,
21+
state_dir: Option<PathBuf>,
1822
}
1923

2024
impl ServerConfigUpdatedHandler {
21-
/// Create a handler wired to the installed-server repo and connection pool.
25+
/// Create a handler wired to the installed-server repo, connection pool,
26+
/// and optional MCP state dir used when re-resolving transport.
2227
pub fn new(
2328
installed_server_repo: Arc<dyn InstalledServerRepository + Send + Sync>,
2429
pool_service: Arc<PoolService>,
30+
state_dir: Option<PathBuf>,
2531
) -> Self {
2632
Self {
2733
installed_server_repo,
2834
pool_service,
35+
state_dir,
2936
}
3037
}
3138

@@ -53,7 +60,7 @@ impl ServerConfigUpdatedHandler {
5360
});
5461
}
5562

56-
/// Handle one domain event, evicting the pool instance when applicable.
63+
/// Handle one domain event, reconnecting the pool instance when applicable.
5764
async fn handle_event(&self, event: DomainEvent) -> anyhow::Result<()> {
5865
let DomainEvent::ServerConfigUpdated {
5966
space_id,
@@ -66,7 +73,10 @@ impl ServerConfigUpdatedHandler {
6673
self.handle_config_updated(space_id, &server_id).await
6774
}
6875

69-
/// Remove a pooled instance for an enabled server after a config/definition write.
76+
/// Reconnect an enabled server from current DB config after a write.
77+
///
78+
/// Resolve failure falls back to evict-only so a stale instance cannot
79+
/// keep serving the pre-save transport.
7080
async fn handle_config_updated(&self, space_id: Uuid, server_id: &str) -> anyhow::Result<()> {
7181
let space_id_str = space_id.to_string();
7282
let Some(installed) = self
@@ -87,10 +97,251 @@ impl ServerConfigUpdatedHandler {
8797
return Ok(());
8898
}
8999

90-
info!(
91-
"[ServerConfigHandler] Evicting pool instance for enabled server {space_id}/{server_id} after config update"
92-
);
93-
self.pool_service.remove_instance(space_id, server_id);
100+
let started = Instant::now();
101+
match resolve_auto_connection_context(
102+
self.installed_server_repo.as_ref(),
103+
self.state_dir.as_deref(),
104+
space_id,
105+
server_id,
106+
)
107+
.await
108+
{
109+
Ok(ctx) => {
110+
let result = self.pool_service.reconnect_fresh(&ctx).await;
111+
info!(
112+
server_id = %server_id,
113+
space_id = %space_id,
114+
ok = result.is_connected(),
115+
duration_ms = started.elapsed().as_millis(),
116+
"[ServerConfigHandler] reconnect_fresh after config update"
117+
);
118+
if let ConnectionResult::Failed { error } = result {
119+
warn!(
120+
server_id = %server_id,
121+
space_id = %space_id,
122+
error = %error,
123+
"[ServerConfigHandler] reconnect_fresh failed after config update"
124+
);
125+
}
126+
}
127+
Err(error) => {
128+
warn!(
129+
server_id = %server_id,
130+
space_id = %space_id,
131+
error = %error,
132+
"[ServerConfigHandler] re-resolve failed, evicting only"
133+
);
134+
self.pool_service.remove_instance(space_id, server_id);
135+
}
136+
}
94137
Ok(())
95138
}
96139
}
140+
141+
#[cfg(test)]
142+
mod tests {
143+
use super::*;
144+
use crate::pool::PoolService;
145+
use async_trait::async_trait;
146+
use chrono::Utc;
147+
use mcpmux_core::{InstalledServer, ServerDefinition, TransportConfig};
148+
use std::collections::HashMap;
149+
use std::sync::Mutex;
150+
151+
struct MapInstalledRepo {
152+
servers: Mutex<HashMap<(String, String), InstalledServer>>,
153+
}
154+
155+
impl MapInstalledRepo {
156+
/// Store one installed-server row keyed by `(space_id, server_id)`.
157+
fn with(server: InstalledServer) -> Arc<Self> {
158+
let mut servers = HashMap::new();
159+
servers.insert((server.space_id.clone(), server.server_id.clone()), server);
160+
Arc::new(Self {
161+
servers: Mutex::new(servers),
162+
})
163+
}
164+
}
165+
166+
#[async_trait]
167+
impl InstalledServerRepository for MapInstalledRepo {
168+
async fn list(&self) -> mcpmux_core::repository::RepoResult<Vec<InstalledServer>> {
169+
Ok(self.servers.lock().unwrap().values().cloned().collect())
170+
}
171+
async fn list_for_space(
172+
&self,
173+
_space_id: &str,
174+
) -> mcpmux_core::repository::RepoResult<Vec<InstalledServer>> {
175+
Ok(vec![])
176+
}
177+
async fn list_by_source_file(
178+
&self,
179+
_file_path: &std::path::Path,
180+
) -> mcpmux_core::repository::RepoResult<Vec<InstalledServer>> {
181+
Ok(vec![])
182+
}
183+
async fn get(
184+
&self,
185+
_id: &Uuid,
186+
) -> mcpmux_core::repository::RepoResult<Option<InstalledServer>> {
187+
Ok(None)
188+
}
189+
async fn get_by_server_id(
190+
&self,
191+
space_id: &str,
192+
server_id: &str,
193+
) -> mcpmux_core::repository::RepoResult<Option<InstalledServer>> {
194+
Ok(self
195+
.servers
196+
.lock()
197+
.unwrap()
198+
.get(&(space_id.to_string(), server_id.to_string()))
199+
.cloned())
200+
}
201+
async fn install(
202+
&self,
203+
_server: &InstalledServer,
204+
) -> mcpmux_core::repository::RepoResult<()> {
205+
Ok(())
206+
}
207+
async fn update(
208+
&self,
209+
_server: &InstalledServer,
210+
) -> mcpmux_core::repository::RepoResult<()> {
211+
Ok(())
212+
}
213+
async fn uninstall(&self, _id: &Uuid) -> mcpmux_core::repository::RepoResult<()> {
214+
Ok(())
215+
}
216+
async fn list_enabled(
217+
&self,
218+
_space_id: &str,
219+
) -> mcpmux_core::repository::RepoResult<Vec<InstalledServer>> {
220+
Ok(vec![])
221+
}
222+
async fn list_enabled_all(
223+
&self,
224+
) -> mcpmux_core::repository::RepoResult<Vec<InstalledServer>> {
225+
Ok(vec![])
226+
}
227+
async fn set_enabled(
228+
&self,
229+
_id: &Uuid,
230+
_enabled: bool,
231+
) -> mcpmux_core::repository::RepoResult<()> {
232+
Ok(())
233+
}
234+
async fn set_oauth_connected(
235+
&self,
236+
_id: &Uuid,
237+
_connected: bool,
238+
) -> mcpmux_core::repository::RepoResult<()> {
239+
Ok(())
240+
}
241+
async fn update_inputs(
242+
&self,
243+
_id: &Uuid,
244+
_input_values: HashMap<String, String>,
245+
) -> mcpmux_core::repository::RepoResult<()> {
246+
Ok(())
247+
}
248+
async fn update_cached_definition(
249+
&self,
250+
_id: &Uuid,
251+
_server_name: Option<String>,
252+
_cached_definition: Option<String>,
253+
) -> mcpmux_core::repository::RepoResult<()> {
254+
Ok(())
255+
}
256+
async fn set_display_name_override(
257+
&self,
258+
_id: &Uuid,
259+
_value: Option<String>,
260+
) -> mcpmux_core::repository::RepoResult<()> {
261+
Ok(())
262+
}
263+
async fn update_version_cache(
264+
&self,
265+
_id: &Uuid,
266+
_latest_available_version: Option<String>,
267+
_current_version: Option<String>,
268+
_version_checked_at: chrono::DateTime<Utc>,
269+
) -> mcpmux_core::repository::RepoResult<()> {
270+
Ok(())
271+
}
272+
}
273+
274+
/// Minimal stdio definition whose command does not exist on disk.
275+
fn stdio_definition(server_id: &str) -> ServerDefinition {
276+
ServerDefinition {
277+
id: server_id.to_string(),
278+
name: server_id.to_string(),
279+
description: None,
280+
alias: None,
281+
auth: None,
282+
icon: None,
283+
transport: TransportConfig::Stdio {
284+
command: "/nonexistent/mcpmux-config-handler-test".into(),
285+
args: vec![],
286+
env: HashMap::new(),
287+
metadata: Default::default(),
288+
},
289+
categories: vec![],
290+
publisher: None,
291+
source: Default::default(),
292+
badges: vec![],
293+
hosting_type: Default::default(),
294+
license: None,
295+
license_url: None,
296+
installation: None,
297+
capabilities: None,
298+
sponsored: None,
299+
media: None,
300+
changelog_url: None,
301+
}
302+
}
303+
304+
#[tokio::test]
305+
async fn enabled_server_reconnects_and_evicts_stale_instance() {
306+
let space_id = Uuid::new_v4();
307+
let server_id = "cfg-enabled";
308+
let installed = InstalledServer::new(space_id.to_string(), server_id)
309+
.with_definition(&stdio_definition(server_id))
310+
.with_enabled(true);
311+
let repo = MapInstalledRepo::with(installed);
312+
let pool = Arc::new(PoolService::new_test_with_repo(repo.clone()));
313+
pool.insert_test_instance(space_id, server_id);
314+
assert_eq!(pool.stats().total_instances, 1);
315+
316+
let handler = ServerConfigUpdatedHandler::new(repo, pool.clone(), None);
317+
handler
318+
.handle_config_updated(space_id, server_id)
319+
.await
320+
.expect("handler");
321+
322+
assert!(
323+
pool.get_instance(space_id, server_id).is_none(),
324+
"stale instance must be gone after reconnect_fresh"
325+
);
326+
}
327+
328+
#[tokio::test]
329+
async fn disabled_server_is_left_in_the_pool() {
330+
let space_id = Uuid::new_v4();
331+
let server_id = "cfg-disabled";
332+
let installed = InstalledServer::new(space_id.to_string(), server_id)
333+
.with_definition(&stdio_definition(server_id))
334+
.with_enabled(false);
335+
let repo = MapInstalledRepo::with(installed);
336+
let pool = Arc::new(PoolService::new_test_with_repo(repo.clone()));
337+
pool.insert_test_instance(space_id, server_id);
338+
339+
let handler = ServerConfigUpdatedHandler::new(repo, pool.clone(), None);
340+
handler
341+
.handle_config_updated(space_id, server_id)
342+
.await
343+
.expect("handler");
344+
345+
assert!(pool.get_instance(space_id, server_id).is_some());
346+
}
347+
}

crates/mcpmux-gateway/src/mcp/oauth_middleware.rs

Lines changed: 24 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -271,11 +271,19 @@ pub async fn mcp_oauth_middleware(
271271
services.session_roots.set_pinned(sid, ws);
272272
}
273273
(None, Some(ws)) => {
274-
warn!(
274+
info!(
275275
trace_id = %trace_id,
276276
workspace_header = %ws,
277-
"[SessionRoots] X-Mcpmux-Workspace present without mcp-session-id — pin skipped",
277+
"[SessionRoots] X-Mcpmux-Workspace held until mcp-session-id exists",
278278
);
279+
services
280+
.session_roots
281+
.remember_pending_workspace(&client_id, ws);
282+
}
283+
(Some(sid), None) => {
284+
services
285+
.session_roots
286+
.apply_pending_workspace(&client_id, sid);
279287
}
280288
_ => {}
281289
}
@@ -327,6 +335,20 @@ pub async fn mcp_oauth_middleware(
327335

328336
let response = next.run(request).await;
329337

338+
if let Some(ws) = workspace_header
339+
.as_deref()
340+
.filter(|value| !value.trim().is_empty())
341+
{
342+
if let Some(sid) = response
343+
.headers()
344+
.get("mcp-session-id")
345+
.and_then(|value| value.to_str().ok())
346+
.filter(|value| !value.is_empty())
347+
{
348+
services.session_roots.set_pinned(sid, ws);
349+
}
350+
}
351+
330352
// Log errors only — except two rmcp spec-correct shapes that are not
331353
// gateway problems: a GET without Mcp-Session-Id (client opening the SSE
332354
// stream before initialize) and any request against a session rmcp has

0 commit comments

Comments
 (0)