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