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    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}