systemprompt_api/services/server/
runner.rs1use 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};
15
16pub async fn run_server(
17 ctx: AppContext,
18 events: Option<StartupEventSender>,
19 early: super::startup::EarlyServer,
20) -> Result<()> {
21 let start_time = std::time::Instant::now();
22
23 let mcp_orchestrator = create_mcp_orchestrator(&ctx)?;
24
25 start_event_bridge(&ctx);
26 reconcile_system_services(&ctx, &mcp_orchestrator, events.as_ref()).await?;
27
28 run_agents_phase(&ctx, events.as_ref()).await?;
29 let scheduler_handle = run_scheduler_phase(&ctx, events.as_ref()).await?;
30
31 if let Some(ref tx) = events {
32 tx.phase_started(Phase::ApiServer);
33 }
34 let router = crate::services::server::setup_api_server(&ctx, events.as_ref())?;
35 let addr = ctx.server_address();
36
37 early.activate(router);
38 super::readiness::signal_ready();
39
40 if let Some(ref tx) = events {
41 tx.phase_completed(Phase::ApiServer);
42 }
43
44 if let Some(ref tx) = events {
45 tx.startup_complete(start_time.elapsed(), format!("http://{}", addr), vec![]);
46 }
47
48 systemprompt_logging::set_startup_mode(false);
49
50 let serve_result = super::shutdown::join_within_drain_grace(early.join()).await;
51
52 super::shutdown::arm_forced_exit();
53 super::shutdown::drain(&ctx, scheduler_handle).await;
54
55 serve_result
56}
57
58async fn run_agents_phase(ctx: &AppContext, events: Option<&StartupEventSender>) -> Result<()> {
59 if let Some(tx) = events {
60 tx.phase_started(Phase::Agents);
61 }
62 match reconcile_agents(ctx, events).await {
63 Ok(started_count) => {
64 if let Some(tx) = events {
65 send_startup_event(
66 tx,
67 StartupEvent::AgentReconciliationComplete {
68 running: started_count,
69 total: started_count,
70 },
71 );
72 tx.phase_completed(Phase::Agents);
73 }
74 Ok(())
75 },
76 Err(e) => Err(fail_phase(
77 events,
78 Phase::Agents,
79 format!("Agent reconciliation failed: {e}"),
80 e,
81 )),
82 }
83}
84
85async fn run_scheduler_phase(
86 ctx: &AppContext,
87 events: Option<&StartupEventSender>,
88) -> Result<Option<SchedulerHandle>> {
89 if let Some(tx) = events {
90 tx.phase_started(Phase::Scheduler);
91 }
92 match initialize_scheduler(ctx, events).await {
93 Ok(handle) => {
94 if let Some(tx) = events {
95 tx.phase_completed(Phase::Scheduler);
96 }
97 Ok(handle)
98 },
99 Err(e) => Err(fail_phase(
100 events,
101 Phase::Scheduler,
102 format!("Scheduler initialization failed: {e}"),
103 e,
104 )),
105 }
106}
107
108fn fail_phase(
109 events: Option<&StartupEventSender>,
110 phase: Phase,
111 message: String,
112 error: anyhow::Error,
113) -> anyhow::Error {
114 if let Some(tx) = events {
115 tx.phase_failed(phase, error.to_string());
116 send_startup_event(
117 tx,
118 StartupEvent::Error {
119 message,
120 fatal: true,
121 },
122 );
123 }
124 error
125}
126
127fn send_startup_event(tx: &StartupEventSender, event: StartupEvent) {
128 if tx.unbounded_send(event).is_err() {
129 tracing::debug!("Startup event receiver dropped");
130 }
131}
132
133fn create_mcp_orchestrator(
134 ctx: &AppContext,
135) -> Result<Arc<systemprompt_mcp::services::McpOrchestrator>> {
136 use systemprompt_mcp::services::McpOrchestrator;
137 let manager = McpOrchestrator::new(
138 Arc::clone(ctx.db_pool()),
139 (**ctx.service_repository()).clone(),
140 Arc::clone(ctx.app_paths_arc()),
141 ctx.mcp_registry().clone(),
142 )?;
143 Ok(Arc::new(manager))
144}