Skip to main content

systemprompt_api/services/server/
runner.rs

1//! Server run loop: MCP orchestrator wiring and lifecycle supervision.
2//!
3//! Copyright (c) systemprompt.io — Business Source License 1.1.
4//! See <https://systemprompt.io> for licensing details.
5
6use 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}