Skip to content

Commit c09e569

Browse files
committed
fix(gateway): retry call_tool after a dead backend connection
Widen reconnect beyond auth errors so a silently closed transport gets one evict+re-resolve+retry, and reject missing feature_set_ids before the bind FK fires. Signed-off-by: crimsonsunset <jsangio1@gmail.com>
1 parent 1367048 commit c09e569

12 files changed

Lines changed: 456 additions & 92 deletions

File tree

apps/desktop/src-tauri/src/lib.rs

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -12,11 +12,11 @@ mod commands;
1212
mod macos_dock;
1313
mod macos_permissions;
1414
mod main_window;
15-
#[cfg(unix)]
16-
mod unix_signal;
1715
mod services;
1816
mod state;
1917
mod tray;
18+
#[cfg(unix)]
19+
mod unix_signal;
2020

2121
// Re-export deep link handler
2222
use commands::oauth::{route_or_buffer_deep_link, PendingInitialDeepLink};

apps/desktop/src-tauri/src/unix_signal.rs

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -43,7 +43,11 @@ pub async fn wait_for_term() {
4343
let sender = SENDER_PID.load(Ordering::SeqCst);
4444
let sig = SIGNAL_NO.load(Ordering::SeqCst);
4545
let code = SIGNAL_CODE.load(Ordering::SeqCst);
46-
let name = if sig == libc::SIGINT { "SIGINT" } else { "SIGTERM" };
46+
let name = if sig == libc::SIGINT {
47+
"SIGINT"
48+
} else {
49+
"SIGTERM"
50+
};
4751

4852
info!(
4953
pid,

crates/mcpmux-gateway/src/admin/write_runtime.rs

Lines changed: 12 additions & 22 deletions
Original file line numberDiff line numberDiff line change
@@ -11,7 +11,9 @@ use tokio::sync::RwLock;
1111
use tracing::warn;
1212
use uuid::Uuid;
1313

14-
use crate::pool::transport::resolution::{build_transport_config, TransportResolutionOptions};
14+
use crate::pool::transport::resolution::{
15+
build_transport_config, resolve_auto_connection_context, TransportResolutionOptions,
16+
};
1517
use crate::pool::{
1618
ConnectionContext, ConnectionResult, FeatureService, PoolService, ServerKey, ServerManager,
1719
};
@@ -180,31 +182,19 @@ impl GatewayWriteRuntime for LiveGatewayWriteRuntime {
180182

181183
async fn retry_connection(&self, space_id: String, server_id: String) -> Result<Value> {
182184
let space_uuid = Uuid::parse_str(&space_id)?;
183-
184-
self.pool_service.remove_instance(space_uuid, &server_id);
185-
186-
let installed = self
187-
.installed_server_repo
188-
.get_by_server_id(&space_id, &server_id)
189-
.await?
190-
.ok_or_else(|| anyhow!("Server not found: {space_id}/{server_id}"))?;
191-
192-
let server_definition = installed
193-
.get_definition()
194-
.ok_or_else(|| anyhow!("Server {server_id} has no cached definition"))?;
185+
let ctx = resolve_auto_connection_context(
186+
self.installed_server_repo.as_ref(),
187+
Some(&self.data_dir),
188+
space_uuid,
189+
&server_id,
190+
)
191+
.await
192+
.map_err(|e| anyhow!(e))?;
195193

196194
let key = ServerKey::new(space_uuid, &server_id);
197195
self.server_manager.set_connecting(&key).await;
198196

199-
let transport = build_transport_config(
200-
&server_definition.transport,
201-
&installed,
202-
Some(&self.data_dir),
203-
TransportResolutionOptions::default(),
204-
);
205-
206-
let ctx = ConnectionContext::auto(space_uuid, server_id.clone(), transport);
207-
let result = self.pool_service.connect_server(&ctx).await;
197+
let result = self.pool_service.reconnect_fresh(&ctx).await;
208198

209199
match result {
210200
ConnectionResult::Connected { features, .. } => {

crates/mcpmux-gateway/src/pool/connection.rs

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -48,6 +48,13 @@ pub enum ConnectionResult {
4848
},
4949
}
5050

51+
impl ConnectionResult {
52+
/// True when the attempt produced a live connection (new or reused).
53+
pub fn is_connected(&self) -> bool {
54+
matches!(self, Self::Connected { .. })
55+
}
56+
}
57+
5158
/// Connection Service handles server connection lifecycle
5259
pub struct ConnectionService {
5360
token_service: Arc<TokenService>,

crates/mcpmux-gateway/src/pool/instance.rs

Lines changed: 38 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -484,3 +484,41 @@ impl ServerInstance {
484484
}
485485
}
486486
}
487+
488+
#[cfg(test)]
489+
mod record_stats_tests {
490+
use super::{InstanceKey, ServerInstance, TransportType};
491+
use std::collections::HashMap;
492+
use uuid::Uuid;
493+
494+
fn test_instance() -> ServerInstance {
495+
ServerInstance::new(
496+
InstanceKey::stdio(Uuid::new_v4(), "echo", &[], &HashMap::new()),
497+
"echo".into(),
498+
TransportType::Stdio,
499+
)
500+
}
501+
502+
#[test]
503+
fn record_failure_increments_consecutive_failures_without_flipping_state() {
504+
let instance = test_instance();
505+
let before = instance.state();
506+
instance.record_failure("mcp error -32000: connection closed");
507+
let stats = instance.stats.read();
508+
assert_eq!(stats.consecutive_failures, 1);
509+
assert_eq!(
510+
stats.last_error.as_deref(),
511+
Some("mcp error -32000: connection closed")
512+
);
513+
drop(stats);
514+
assert_eq!(instance.state(), before);
515+
}
516+
517+
#[test]
518+
fn record_success_increments_requests_served() {
519+
let instance = test_instance();
520+
instance.record_success();
521+
instance.record_success();
522+
assert_eq!(instance.stats.read().requests_served, 2);
523+
}
524+
}

0 commit comments

Comments
 (0)