systemprompt_api/services/server/
runner.rs1use anyhow::Result;
7use std::sync::Arc;
8use systemprompt_runtime::AppContext;
9use systemprompt_scheduler::services::SchedulerHandle;
10use systemprompt_traits::{Phase, StartupEvent, StartupEventExt, StartupEventSender};
11
12use super::lifecycle::{
13 initialize_scheduler, reconcile_agents, reconcile_system_services, start_event_bridge,
14 start_registry_heartbeat,
15};
16
17pub async fn run_server(
18 ctx: AppContext,
19 events: Option<StartupEventSender>,
20 early: super::startup::EarlyServer,
21) -> Result<()> {
22 let start_time = std::time::Instant::now();
23
24 let mcp_orchestrator = create_mcp_orchestrator(&ctx)?;
25
26 start_event_bridge(&ctx);
27 let heartbeat = start_registry_heartbeat(&ctx);
28 reconcile_system_services(&ctx, &mcp_orchestrator, events.as_ref()).await?;
29
30 run_agents_phase(&ctx, events.as_ref()).await?;
31 let scheduler_handle = run_scheduler_phase(&ctx, events.as_ref()).await?;
32
33 if let Some(ref tx) = events {
34 tx.phase_started(Phase::ApiServer);
35 }
36 let router = crate::services::server::setup_api_server(&ctx, events.as_ref())?;
37 let accounting_recovery = start_accounting_recovery(&ctx).await?;
38 let addr = ctx.server_address();
39
40 early.activate(router);
41 let metrics_listener = start_metrics_listener(&ctx).await?;
42 let reporting = systemprompt_runtime::reporting::spawn(ctx.db_pool())?;
43 super::readiness::signal_ready();
44
45 if let Some(ref tx) = events {
46 tx.phase_completed(Phase::ApiServer);
47 }
48
49 if let Some(ref tx) = events {
50 tx.startup_complete(start_time.elapsed(), format!("http://{}", addr), vec![]);
51 }
52
53 systemprompt_logging::set_startup_mode(false);
54
55 let serve_result = super::shutdown::join_within_drain_grace(early.join()).await;
56
57 super::shutdown::arm_forced_exit();
58 heartbeat.abort();
59 reporting.abort();
60 if let Err(error) = reporting.await
61 && !error.is_cancelled()
62 {
63 tracing::warn!(error = %error, "Analytics projection worker ended abnormally");
64 }
65 if let Some(recovery) = accounting_recovery {
66 recovery.abort();
67 }
68 if let Some(listener) = metrics_listener {
69 listener.abort();
70 }
71 super::shutdown::drain(&ctx, scheduler_handle).await;
72
73 serve_result
74}
75
76async fn run_agents_phase(ctx: &AppContext, events: Option<&StartupEventSender>) -> Result<()> {
77 if let Some(tx) = events {
78 tx.phase_started(Phase::Agents);
79 }
80 match reconcile_agents(ctx, events).await {
81 Ok(started_count) => {
82 if let Some(tx) = events {
83 send_startup_event(
84 tx,
85 StartupEvent::AgentReconciliationComplete {
86 running: started_count,
87 total: started_count,
88 },
89 );
90 tx.phase_completed(Phase::Agents);
91 }
92 Ok(())
93 },
94 Err(e) => Err(fail_phase(
95 events,
96 Phase::Agents,
97 format!("Agent reconciliation failed: {e}"),
98 e,
99 )),
100 }
101}
102
103async fn run_scheduler_phase(
104 ctx: &AppContext,
105 events: Option<&StartupEventSender>,
106) -> Result<Option<SchedulerHandle>> {
107 if let Some(tx) = events {
108 tx.phase_started(Phase::Scheduler);
109 }
110 match initialize_scheduler(ctx, events).await {
111 Ok(handle) => {
112 if let Some(tx) = events {
113 tx.phase_completed(Phase::Scheduler);
114 }
115 Ok(handle)
116 },
117 Err(e) => Err(fail_phase(
118 events,
119 Phase::Scheduler,
120 format!("Scheduler initialization failed: {e}"),
121 e,
122 )),
123 }
124}
125
126fn fail_phase(
127 events: Option<&StartupEventSender>,
128 phase: Phase,
129 message: String,
130 error: anyhow::Error,
131) -> anyhow::Error {
132 if let Some(tx) = events {
133 tx.phase_failed(phase, error.to_string());
134 send_startup_event(
135 tx,
136 StartupEvent::Error {
137 message,
138 fatal: true,
139 },
140 );
141 }
142 error
143}
144
145fn send_startup_event(tx: &StartupEventSender, event: StartupEvent) {
146 if tx.unbounded_send(event).is_err() {
147 tracing::debug!("Startup event receiver dropped");
148 }
149}
150
151fn create_mcp_orchestrator(
152 ctx: &AppContext,
153) -> Result<Arc<systemprompt_mcp::services::McpOrchestrator>> {
154 use systemprompt_mcp::services::McpOrchestrator;
155 let manager = McpOrchestrator::new(
156 (**ctx.service_repository()).clone(),
157 Arc::clone(ctx.app_paths_arc()),
158 ctx.mcp_registry().clone(),
159 )?;
160 Ok(Arc::new(manager))
161}
162
163async fn start_metrics_listener(ctx: &AppContext) -> Result<Option<tokio::task::JoinHandle<()>>> {
164 let Some(port) = ctx.config().metrics_port else {
165 return Ok(None);
166 };
167 let handle = super::metrics::install_recorder(&ctx.config().instance_id)?;
168 let addr = std::net::SocketAddr::new(ctx.config().host.parse()?, port);
169 Ok(Some(
170 super::metrics::serve_metrics_listener(addr, handle).await?,
171 ))
172}
173
174async fn start_accounting_recovery(
175 ctx: &AppContext,
176) -> Result<Option<tokio::task::JoinHandle<()>>> {
177 if !crate::routes::gateway::gateway_enabled(ctx) {
178 return Ok(None);
179 }
180 let settlement = crate::routes::gateway::gateway_repositories(ctx)?.settlement();
181 let settled = crate::services::gateway::audit::journal::recover(&settlement).await?;
182 if settled > 0 {
183 tracing::info!(settled, "Gateway accounting receipts recovered at startup");
184 }
185 Ok(Some(
186 crate::services::gateway::audit::journal::spawn_recovery(settlement),
187 ))
188}