Skip to main content

systemprompt_api/services/proxy/
resolver.rs

1//! Service resolution for the MCP proxy, with restart-on-dead-backend.
2//!
3//! Copyright (c) systemprompt.io — Business Source License 1.1.
4//! See <https://systemprompt.io> for licensing details.
5
6use std::sync::Arc;
7
8use systemprompt_database::ServiceConfig;
9use systemprompt_mcp::services::McpOrchestrator;
10use systemprompt_mcp::services::spawn_target::SpawnTarget;
11use systemprompt_runtime::AppContext;
12
13use super::backend::ProxyError;
14
15#[derive(Debug, Clone, Copy)]
16pub struct ServiceResolver;
17
18pub fn stale_port(db_port: i32, config_port: u16) -> Option<u16> {
19    if db_port == i32::from(config_port) {
20        None
21    } else {
22        Some(config_port)
23    }
24}
25
26impl ServiceResolver {
27    pub async fn resolve(
28        service_name: &str,
29        ctx: &AppContext,
30    ) -> Result<ServiceConfig, ProxyError> {
31        let service_repo = ctx.service_repository();
32
33        let service = match service_repo.find_service_by_name(service_name).await {
34            Ok(svc) => svc,
35            Err(e) => {
36                tracing::error!(service = %service_name, error = %e, "Database error when looking up service");
37                return Err(ProxyError::DatabaseError {
38                    service: service_name.to_owned(),
39                    source: e,
40                });
41            },
42        };
43
44        let Some(service) = service else {
45            tracing::warn!(service = %service_name, "Service not found");
46            return Err(ProxyError::ServiceNotFound {
47                service: service_name.to_owned(),
48            });
49        };
50
51        if service.status != "running" {
52            if service.status == "crashed" {
53                tracing::info!(service = %service_name, "Service crashed, attempting restart");
54
55                if Self::attempt_restart(service_name, ctx).await.is_ok() {
56                    let restarted = service_repo
57                        .find_service_by_name(service_name)
58                        .await
59                        .map_err(|e| ProxyError::DatabaseError {
60                            service: service_name.to_owned(),
61                            source: e,
62                        })?;
63
64                    if let Some(restarted) = restarted
65                        && restarted.status == "running"
66                    {
67                        tracing::info!(service = %service_name, "Service restarted, retrying proxy");
68                        return Ok(restarted);
69                    }
70
71                    tracing::warn!(
72                        service = %service_name,
73                        "Restart reported success but the service is not running"
74                    );
75                }
76            }
77
78            tracing::warn!(service = %service_name, status = %service.status, "Service not running");
79            return Err(ProxyError::ServiceNotRunning {
80                service: service_name.to_owned(),
81                status: service.status.clone(),
82            });
83        }
84
85        Ok(Self::reconcile_internal_port(service_name, service, ctx).await)
86    }
87
88    async fn reconcile_internal_port(
89        service_name: &str,
90        mut service: ServiceConfig,
91        ctx: &AppContext,
92    ) -> ServiceConfig {
93        let config = match ctx.mcp_registry().find_server(service_name) {
94            Ok(Some(config)) => config,
95            Ok(None) => return service,
96            Err(e) => {
97                tracing::debug!(service = %service_name, error = %e, "Registry lookup failed while reconciling service port");
98                return service;
99            },
100        };
101
102        let Ok(config_port) = config.spawn_port() else {
103            return service;
104        };
105
106        if let Some(new_port) = stale_port(service.port, config_port) {
107            tracing::warn!(
108                service = %service_name,
109                db_port = service.port,
110                config_port = new_port,
111                "DB service port is stale; reconciling to the port this instance spawns before proxying"
112            );
113            if let Err(e) = ctx
114                .service_repository()
115                .update_service_port(service_name, new_port)
116                .await
117            {
118                tracing::error!(service = %service_name, error = %e, "Failed to persist reconciled service port");
119            }
120            service.port = i32::from(new_port);
121        }
122
123        service
124    }
125
126    async fn attempt_restart(service_name: &str, ctx: &AppContext) -> Result<(), ProxyError> {
127        let orchestrator = McpOrchestrator::new(
128            (**ctx.service_repository()).clone(),
129            Arc::clone(ctx.app_paths_arc()),
130            ctx.mcp_registry().clone(),
131        )
132        .map_err(|e| ProxyError::ServiceNotRunning {
133            service: service_name.to_owned(),
134            status: format!("Failed to create orchestrator: {e}"),
135        })?;
136
137        match orchestrator
138            .start_services(Some(service_name.to_owned()))
139            .await
140        {
141            Ok(()) => {},
142            Err(e) => {
143                tracing::error!(service = %service_name, error = %e, "Failed to restart service");
144                return Err(ProxyError::ServiceNotRunning {
145                    service: service_name.to_owned(),
146                    status: format!("Restart failed: {e}"),
147                });
148            },
149        }
150
151        tokio::time::sleep(std::time::Duration::from_millis(500)).await;
152        Ok(())
153    }
154}