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