systemprompt_api/services/proxy/
resolver.rs1use 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}