Skip to content

Commit 5efd85d

Browse files
committed
fix(gateway): Phase 3 — evict pool instance on ServerConfigUpdated
Autonomous decisions: - UserSpaceSyncService uses optional EventSender via with_event_sender() — keeps existing new() call sites compiling; file-watcher now wired to the app-wide EventBus sender created alongside ServerAppService - Split regression coverage: user_space_sync event emission test + pool remove_instance unit test — full handler integration test skipped due to heavy mock surface in gateway lib tests - Added cfg(test) insert_test_instance on PoolService — minimal test seam to assert eviction clears a pooled entry without standing up a live MCP connection - Amended (not separate commit) — branch is ahead of origin by one unpushed Phase 3 commit; lib.rs wiring is the same logical change Signed-off-by: crimsonsunset <jsangio1@gmail.com>
1 parent 51209c5 commit 5efd85d

7 files changed

Lines changed: 539 additions & 6 deletions

File tree

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

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -643,6 +643,7 @@ pub fn run() {
643643
let app_state: tauri::State<'_, AppState> = app.state();
644644
let spaces_dir = app_state.spaces_dir().to_path_buf();
645645
let installed_repo = app_state.installed_server_repository.clone();
646+
let event_sender_for_sync = event_bus.sender();
646647
let app_handle_for_watcher = app.handle().clone();
647648

648649
// Use the well-known default space UUID
@@ -656,7 +657,10 @@ pub fn run() {
656657
// Create file watcher with UI event emitter
657658
match services::SpaceFileWatcher::new(
658659
spaces_dir.clone(),
659-
Arc::new(mcpmux_core::application::UserSpaceSyncService::new(installed_repo)),
660+
Arc::new(
661+
mcpmux_core::application::UserSpaceSyncService::new(installed_repo)
662+
.with_event_sender(event_sender_for_sync),
663+
),
660664
default_space_id,
661665
services::SpaceFileWatcherEmitters {
662666
on_success: Some(Arc::new(

crates/mcpmux-core/src/application/user_space_sync.rs

Lines changed: 79 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -9,9 +9,11 @@ use std::sync::Arc;
99

1010
use anyhow::{Context, Result};
1111
use tracing::{debug, info};
12+
use uuid::Uuid;
1213

1314
use crate::domain::config::UserSpaceConfig;
14-
use crate::domain::{InstallationSource, InstalledServer, ServerDefinition};
15+
use crate::domain::{DomainEvent, InstallationSource, InstalledServer, ServerDefinition};
16+
use crate::event_bus::EventSender;
1517
use crate::repository::InstalledServerRepository;
1618

1719
/// Result of a sync operation
@@ -54,12 +56,40 @@ enum SyncOutcome {
5456
/// Service for syncing user space JSON config files to InstalledServer records
5557
pub struct UserSpaceSyncService {
5658
installed_repo: Arc<dyn InstalledServerRepository>,
59+
event_sender: Option<EventSender>,
5760
}
5861

5962
impl UserSpaceSyncService {
6063
/// Create a new sync service
6164
pub fn new(installed_repo: Arc<dyn InstalledServerRepository>) -> Self {
62-
Self { installed_repo }
65+
Self {
66+
installed_repo,
67+
event_sender: None,
68+
}
69+
}
70+
71+
/// Attach an event sender so definition updates emit `ServerConfigUpdated`.
72+
pub fn with_event_sender(mut self, sender: EventSender) -> Self {
73+
self.event_sender = Some(sender);
74+
self
75+
}
76+
77+
/// Emit `ServerConfigUpdated` when a wired event sender is present.
78+
fn emit_server_config_updated(&self, space_id: &str, server_id: &str) {
79+
let Some(sender) = &self.event_sender else {
80+
return;
81+
};
82+
let Ok(space_uuid) = Uuid::parse_str(space_id) else {
83+
debug!(
84+
"Skipping ServerConfigUpdated emit for {space_id}/{server_id}: invalid space UUID"
85+
);
86+
return;
87+
};
88+
89+
sender.emit(DomainEvent::ServerConfigUpdated {
90+
space_id: space_uuid,
91+
server_id: server_id.to_string(),
92+
});
6393
}
6494

6595
/// Ensure no two user-config entries normalize to the same MCP server id.
@@ -197,7 +227,8 @@ impl UserSpaceSyncService {
197227
}
198228
Ok(SyncOutcome::Updated) => {
199229
debug!("Updated server: {}", server_id);
200-
result.updated.push(server_id);
230+
result.updated.push(server_id.clone());
231+
self.emit_server_config_updated(space_id, &server_id);
201232
}
202233
Ok(SyncOutcome::Adopted) => {
203234
info!("Adopted server from another source: {}", server_id);
@@ -687,4 +718,49 @@ mod tests {
687718
"unexpected error: {err}"
688719
);
689720
}
721+
722+
#[tokio::test]
723+
async fn sync_from_file_emits_server_config_updated_on_definition_update() {
724+
use crate::EventBus;
725+
726+
let space_id = "00000000-0000-0000-0000-000000000001";
727+
let bus = EventBus::new();
728+
let mut receiver = bus.subscribe();
729+
let repo = Arc::new(InMemoryInstalledServerRepository::new());
730+
let (_dir, file_path) = write_config_file(
731+
r#"{ "mcpServers": {
732+
"alpha": { "command": "echo", "name": "Alpha v1" }
733+
} }"#,
734+
)
735+
.await;
736+
737+
repo.seed(
738+
InstalledServer::new(space_id, "alpha")
739+
.with_source(InstallationSource::UserConfig {
740+
file_path: file_path.clone(),
741+
})
742+
.with_enabled(true),
743+
)
744+
.await;
745+
746+
let service = UserSpaceSyncService::new(repo).with_event_sender(bus.sender());
747+
tokio::fs::write(
748+
&file_path,
749+
r#"{ "mcpServers": {
750+
"alpha": { "command": "echo", "name": "Alpha v2" }
751+
} }"#,
752+
)
753+
.await
754+
.expect("rewrite config file");
755+
756+
let result = service
757+
.sync_from_file(space_id, &file_path)
758+
.await
759+
.expect("sync should update existing server");
760+
assert_eq!(result.updated, vec!["alpha".to_string()]);
761+
762+
let event = receiver.recv().await.expect("ServerConfigUpdated should emit");
763+
assert_eq!(event.type_name(), "server_config_updated");
764+
assert_eq!(event.server_id(), Some("alpha"));
765+
}
690766
}

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

Lines changed: 58 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -178,8 +178,64 @@ impl GatewayWriteRuntime for LiveGatewayWriteRuntime {
178178
Err(gateway_write_unavailable())
179179
}
180180

181-
async fn retry_connection(&self, _space_id: String, _server_id: String) -> Result<Value> {
182-
Err(gateway_write_unavailable())
181+
async fn retry_connection(&self, space_id: String, server_id: String) -> Result<Value> {
182+
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"))?;
195+
196+
let key = ServerKey::new(space_uuid, &server_id);
197+
self.server_manager.set_connecting(&key).await;
198+
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;
208+
209+
match result {
210+
ConnectionResult::Connected { features, .. } => {
211+
self.server_manager.set_connected(&key, features).await;
212+
}
213+
ConnectionResult::OAuthRequired { .. } => {
214+
self.server_manager.set_auth_required(&key, None).await;
215+
if let Err(error) = self
216+
.feature_service
217+
.mark_unavailable(&space_id, &server_id)
218+
.await
219+
{
220+
warn!("[LiveGatewayWriteRuntime] Failed to mark features unavailable: {error}");
221+
}
222+
}
223+
ConnectionResult::Failed { error } => {
224+
self.server_manager.set_error(&key, error.clone()).await;
225+
if let Err(mark_error) = self
226+
.feature_service
227+
.mark_unavailable(&space_id, &server_id)
228+
.await
229+
{
230+
warn!(
231+
"[LiveGatewayWriteRuntime] Failed to mark features unavailable: {mark_error}"
232+
);
233+
}
234+
return Err(anyhow!(error));
235+
}
236+
}
237+
238+
Ok(json!({ "ok": true }))
183239
}
184240

185241
async fn update_server_package(&self, space_id: String, server_id: String) -> Result<Value> {

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

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -31,6 +31,8 @@
3131
3232
mod mcp_notifier;
3333
mod oauth_handler;
34+
mod server_config_handler;
3435

3536
pub use mcp_notifier::MCPNotifier;
3637
pub use oauth_handler::OAuthEventHandler;
38+
pub use server_config_handler::ServerConfigUpdatedHandler;
Lines changed: 92 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,92 @@
1+
//! Server Config Event Handler - Evicts pool instances on config changes
2+
//!
3+
//! Listens for `ServerConfigUpdated` and removes the pooled instance for
4+
//! enabled servers so the next connect rebuilds transport from DB.
5+
6+
use mcpmux_core::{DomainEvent, InstalledServerRepository};
7+
use std::sync::Arc;
8+
use tokio::sync::broadcast;
9+
use tracing::{debug, info, warn};
10+
use uuid::Uuid;
11+
12+
use crate::pool::PoolService;
13+
14+
/// Evicts pooled server instances when configuration changes.
15+
pub struct ServerConfigUpdatedHandler {
16+
installed_server_repo: Arc<dyn InstalledServerRepository + Send + Sync>,
17+
pool_service: Arc<PoolService>,
18+
}
19+
20+
impl ServerConfigUpdatedHandler {
21+
/// Create a handler wired to the installed-server repo and connection pool.
22+
pub fn new(
23+
installed_server_repo: Arc<dyn InstalledServerRepository + Send + Sync>,
24+
pool_service: Arc<PoolService>,
25+
) -> Self {
26+
Self {
27+
installed_server_repo,
28+
pool_service,
29+
}
30+
}
31+
32+
/// Start listening to domain events on a background task.
33+
pub fn start(self: Arc<Self>, mut event_rx: broadcast::Receiver<DomainEvent>) {
34+
tokio::spawn(async move {
35+
info!("[ServerConfigHandler] Started listening for ServerConfigUpdated events");
36+
37+
loop {
38+
match event_rx.recv().await {
39+
Ok(event) => {
40+
if let Err(error) = self.handle_event(event).await {
41+
warn!("[ServerConfigHandler] Failed to handle event: {error}");
42+
}
43+
}
44+
Err(broadcast::error::RecvError::Lagged(skipped)) => {
45+
warn!("[ServerConfigHandler] Lagged behind, skipped {skipped} events");
46+
}
47+
Err(broadcast::error::RecvError::Closed) => {
48+
warn!("[ServerConfigHandler] Event channel closed");
49+
break;
50+
}
51+
}
52+
}
53+
});
54+
}
55+
56+
/// Handle one domain event, evicting the pool instance when applicable.
57+
async fn handle_event(&self, event: DomainEvent) -> anyhow::Result<()> {
58+
let DomainEvent::ServerConfigUpdated { space_id, server_id } = event else {
59+
return Ok(());
60+
};
61+
62+
self.handle_config_updated(space_id, &server_id).await
63+
}
64+
65+
/// Remove a pooled instance for an enabled server after a config/definition write.
66+
async fn handle_config_updated(&self, space_id: Uuid, server_id: &str) -> anyhow::Result<()> {
67+
let space_id_str = space_id.to_string();
68+
let Some(installed) = self
69+
.installed_server_repo
70+
.get_by_server_id(&space_id_str, server_id)
71+
.await?
72+
else {
73+
debug!(
74+
"[ServerConfigHandler] Server {space_id}/{server_id} not found, skipping eviction"
75+
);
76+
return Ok(());
77+
};
78+
79+
if !installed.enabled {
80+
debug!(
81+
"[ServerConfigHandler] Skipping eviction for disabled server {space_id}/{server_id}"
82+
);
83+
return Ok(());
84+
}
85+
86+
info!(
87+
"[ServerConfigHandler] Evicting pool instance for enabled server {space_id}/{server_id} after config update"
88+
);
89+
self.pool_service.remove_instance(space_id, server_id);
90+
Ok(())
91+
}
92+
}

0 commit comments

Comments
 (0)