Skip to main content

systemprompt_api/services/server/
runner.rs

1//! Server run loop: MCP orchestrator wiring and lifecycle supervision.
2//!
3//! The startup phases a node runs follow its `server.role`
4//! ([`super::routes::role::lifecycle_plan`]): a gateway node spawns no MCP
5//! servers or agents and runs no scheduler.
6//!
7//! Copyright (c) systemprompt.io — Business Source License 1.1.
8//! See <https://systemprompt.io> for licensing details.
9
10use 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}