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