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 super::readiness::signal_ready();
43
44 if let Some(ref tx) = events {
45 tx.phase_completed(Phase::ApiServer);
46 }
47
48 if let Some(ref tx) = events {
49 tx.startup_complete(start_time.elapsed(), format!("http://{}", addr), vec![]);
50 }
51
52 systemprompt_logging::set_startup_mode(false);
53
54 let serve_result = super::shutdown::join_within_drain_grace(early.join()).await;
55
56 super::shutdown::arm_forced_exit();
57 heartbeat.abort();
58 if let Some(recovery) = accounting_recovery {
59 recovery.abort();
60 }
61 if let Some(listener) = metrics_listener {
62 listener.abort();
63 }
64 super::shutdown::drain(&ctx, scheduler_handle).await;
65
66 serve_result
67}
68
69async fn run_agents_phase(ctx: &AppContext, events: Option<&StartupEventSender>) -> Result<()> {
70 if let Some(tx) = events {
71 tx.phase_started(Phase::Agents);
72 }
73 match reconcile_agents(ctx, events).await {
74 Ok(started_count) => {
75 if let Some(tx) = events {
76 send_startup_event(
77 tx,
78 StartupEvent::AgentReconciliationComplete {
79 running: started_count,
80 total: started_count,
81 },
82 );
83 tx.phase_completed(Phase::Agents);
84 }
85 Ok(())
86 },
87 Err(e) => Err(fail_phase(
88 events,
89 Phase::Agents,
90 format!("Agent reconciliation failed: {e}"),
91 e,
92 )),
93 }
94}
95
96async fn run_scheduler_phase(
97 ctx: &AppContext,
98 events: Option<&StartupEventSender>,
99) -> Result<Option<SchedulerHandle>> {
100 if let Some(tx) = events {
101 tx.phase_started(Phase::Scheduler);
102 }
103 match initialize_scheduler(ctx, events).await {
104 Ok(handle) => {
105 if let Some(tx) = events {
106 tx.phase_completed(Phase::Scheduler);
107 }
108 Ok(handle)
109 },
110 Err(e) => Err(fail_phase(
111 events,
112 Phase::Scheduler,
113 format!("Scheduler initialization failed: {e}"),
114 e,
115 )),
116 }
117}
118
119fn fail_phase(
120 events: Option<&StartupEventSender>,
121 phase: Phase,
122 message: String,
123 error: anyhow::Error,
124) -> anyhow::Error {
125 if let Some(tx) = events {
126 tx.phase_failed(phase, error.to_string());
127 send_startup_event(
128 tx,
129 StartupEvent::Error {
130 message,
131 fatal: true,
132 },
133 );
134 }
135 error
136}
137
138fn send_startup_event(tx: &StartupEventSender, event: StartupEvent) {
139 if tx.unbounded_send(event).is_err() {
140 tracing::debug!("Startup event receiver dropped");
141 }
142}
143
144fn create_mcp_orchestrator(
145 ctx: &AppContext,
146) -> Result<Arc<systemprompt_mcp::services::McpOrchestrator>> {
147 use systemprompt_mcp::services::McpOrchestrator;
148 let manager = McpOrchestrator::new(
149 Arc::clone(ctx.db_pool()),
150 (**ctx.service_repository()).clone(),
151 Arc::clone(ctx.app_paths_arc()),
152 ctx.mcp_registry().clone(),
153 )?;
154 Ok(Arc::new(manager))
155}
156
157async fn start_metrics_listener(ctx: &AppContext) -> Result<Option<tokio::task::JoinHandle<()>>> {
158 let Some(port) = ctx.config().metrics_port else {
159 return Ok(None);
160 };
161 let handle = super::metrics::install_recorder(&ctx.config().instance_id)?;
162 let addr = std::net::SocketAddr::new(ctx.config().host.parse()?, port);
163 Ok(Some(
164 super::metrics::serve_metrics_listener(addr, handle).await?,
165 ))
166}
167
168async fn start_accounting_recovery(
169 ctx: &AppContext,
170) -> Result<Option<tokio::task::JoinHandle<()>>> {
171 if !crate::routes::gateway::gateway_enabled(ctx) {
172 return Ok(None);
173 }
174 let settlement = crate::routes::gateway::gateway_repositories(ctx)?.settlement();
175 let settled = crate::services::gateway::audit::journal::recover(&settlement).await?;
176 if settled > 0 {
177 tracing::info!(settled, "Gateway accounting receipts recovered at startup");
178 }
179 Ok(Some(
180 crate::services::gateway::audit::journal::spawn_recovery(settlement),
181 ))
182}