use crate::cli_settings::CliConfig;
use crate::interactive::{Prompter, confirm_optional};
use anyhow::{Context, Result};
use std::sync::Arc;
use std::time::Duration;
use systemprompt_config::ProfileBootstrap;
use systemprompt_loader::subprocess::{self, ChildKind};
use systemprompt_logging::CliService;
use systemprompt_runtime::{AppContext, ShutdownRequest, validate_system};
use systemprompt_scheduler::{port_holders, wait_for_port_free};
use systemprompt_traits::{Phase, StartupEvent, StartupEventExt, StartupEventSender};
use super::{get_api_addr, get_api_port};
const CONFIRMED_HOLDER_GRACE: Duration = Duration::from_secs(2);
const PORT_RELEASE: Duration = Duration::from_secs(5);
#[derive(Debug, Clone, Copy)]
pub struct ServeOptions {
pub foreground: bool,
pub kill_port_process: bool,
pub run_migrations: bool,
}
pub async fn execute_with_events(
prompter: &dyn Prompter,
options: ServeOptions,
config: &CliConfig,
events: Option<&StartupEventSender>,
) -> Result<String> {
let ServeOptions {
foreground,
kill_port_process,
run_migrations,
} = options;
let port = get_api_port();
if events.is_none() {
CliService::startup_banner(Some("Starting services..."));
}
ensure_port_free(prompter, port, kill_port_process, config, events).await?;
let shutdown = ShutdownRequest::default();
let early = bind_early(foreground, events, shutdown.clone()).await?;
let ctx = Arc::new(
AppContext::builder()
.with_startup_warnings(true)
.with_shutdown(shutdown)
.with_migrations(run_migrations)
.with_schema_verification(!run_migrations)
.build()
.await
.context("Failed to initialize application context")?,
);
if events.is_none() {
if run_migrations {
CliService::phase_success("Database schemas installed", None);
} else {
CliService::phase_success("Schema current (migrations skipped)", None);
}
}
if let Some(tx) = events {
tx.phase_started(Phase::Database);
if let Err(e) = tx.unbounded_send(StartupEvent::DatabaseValidated) {
tracing::debug!(error = %e, "startup event channel closed: DatabaseValidated");
}
tx.phase_completed(Phase::Database);
} else {
CliService::phase("Validation");
CliService::phase_info("Running system validation...", None);
}
validate_system(&ctx)
.await
.context("System validation failed")?;
if events.is_none() {
CliService::phase_success("System validation complete", None);
CliService::phase_info(&format!("Node role: {}", ctx.config().role), None);
} else if let Some(tx) = events {
tx.info(format!("Node role: {}", ctx.config().role));
}
if events.is_none() {
CliService::phase("Server");
if !foreground {
CliService::phase_warning("Daemon mode not supported", Some("running in foreground"));
}
} else if let Some(tx) = events {
tx.phase_started(Phase::ApiServer);
if !foreground {
tx.warning("Daemon mode not supported, running in foreground");
}
}
if let Some(early) = early {
systemprompt_api::services::server::run_server(
Arc::unwrap_or_clone(ctx),
events.cloned(),
early,
)
.await?;
}
Ok(format!("http://127.0.0.1:{}", port))
}
#[derive(Debug, Clone, Copy)]
pub struct ServeFlags {
pub foreground: bool,
pub kill_port_process: bool,
pub skip_migrate: bool,
}
pub const fn effective_run_migrations(skip_flag: bool, migrate_on_boot: bool) -> bool {
!skip_flag && migrate_on_boot
}
pub async fn execute(prompter: &dyn Prompter, flags: ServeFlags, config: &CliConfig) -> Result<()> {
let migrate_on_boot = ProfileBootstrap::get()
.context("Profile not initialized")?
.database
.migrate_on_boot;
execute_with_events(
prompter,
ServeOptions {
foreground: flags.foreground,
kill_port_process: flags.kill_port_process,
run_migrations: effective_run_migrations(flags.skip_migrate, migrate_on_boot),
},
config,
None,
)
.await
.map(|_| ())
}
async fn ensure_port_free(
prompter: &dyn Prompter,
port: u16,
kill_port_process: bool,
config: &CliConfig,
events: Option<&StartupEventSender>,
) -> Result<()> {
if let Some(pid) = port_holder(port).await? {
if let Some(tx) = events
&& let Err(e) = tx.unbounded_send(StartupEvent::PortConflict { port, pid })
{
tracing::debug!(error = %e, "startup event channel closed: PortConflict");
}
handle_port_conflict(
prompter,
PortConflict { port, pid },
kill_port_process,
config,
events,
)
.await?;
if let Some(tx) = events
&& let Err(e) = tx.unbounded_send(StartupEvent::PortConflictResolved { port })
{
tracing::debug!(error = %e, "startup event channel closed: PortConflictResolved");
}
} else if let Some(tx) = events {
tx.port_available(port);
} else {
CliService::phase_success(&format!("Port {} available", port), None);
}
Ok(())
}
async fn bind_early(
foreground: bool,
events: Option<&StartupEventSender>,
shutdown: ShutdownRequest,
) -> Result<Option<systemprompt_api::services::server::EarlyServer>> {
if !foreground {
return Ok(None);
}
let addr = get_api_addr().context("Profile not initialized; cannot determine bind address")?;
let early =
systemprompt_api::services::server::bind_and_serve(&addr, events.cloned(), shutdown)
.await
.context("Failed to bind API listener")?;
Ok(Some(early))
}
#[derive(Debug, Clone, Copy)]
struct PortConflict {
port: u16,
pid: u32,
}
async fn port_holder(port: u16) -> Result<Option<u32>> {
Ok(port_holders(port).await?.first().copied())
}
async fn stop_confirmed_holder(port: u16, pid: u32) -> Result<()> {
if port_holders(port).await?.contains(&pid) {
subprocess::terminate_gracefully(pid, CONFIRMED_HOLDER_GRACE)
.await
.with_context(|| format!("Failed to stop PID {pid} holding port {port}"))?;
}
wait_for_port_free(port, PORT_RELEASE)
.await
.with_context(|| format!("Failed to free port {port} after stopping PID {pid}"))?;
Ok(())
}
async fn handle_port_conflict(
prompter: &dyn Prompter,
conflict: PortConflict,
kill_port_process: bool,
config: &CliConfig,
events: Option<&StartupEventSender>,
) -> Result<()> {
let PortConflict { port, pid } = conflict;
if events.is_none() {
CliService::warning(&format!("Port {} is already in use by PID {}", port, pid));
}
let verified = subprocess::owns(pid, ChildKind::Api, &subprocess::api_server_service()).await;
let should_kill = if verified {
kill_port_process
|| confirm_optional(
prompter,
&format!("Stop the running API server (PID {pid}) and restart?"),
false,
config,
)?
} else if kill_port_process {
confirm_optional(
prompter,
&format!(
"PID {pid} holding port {port} is not a verified systemprompt API server. Signal \
PID {pid} anyway?"
),
false,
config,
)?
} else {
false
};
if should_kill {
if events.is_none() {
CliService::info(&format!("Stopping process {}...", pid));
}
stop_confirmed_holder(port, pid).await?;
if events.is_none() {
CliService::success(&format!("Port {} is now available", port));
}
return Ok(());
}
if !verified {
return Err(anyhow::anyhow!(
"Port {port} is held by PID {pid}, which is not a verified systemprompt API server; \
it was not signalled. Stop it by hand, or rerun interactively with \
--kill-port-process and confirm."
));
}
if config.is_interactive() {
return Err(anyhow::anyhow!(
"Port {} is occupied by PID {}. Aborted by user.",
port,
pid
));
}
if events.is_none() {
CliService::error(&format!("Port {} is already in use by PID {}", port, pid));
CliService::info("Use --kill-port-process to stop it, or:");
CliService::info(" - systemprompt infra services restart api");
CliService::info(" - systemprompt infra services stop --all --force");
CliService::info(&format!(
" - kill {} (manually kill the process)",
pid
));
}
Err(anyhow::anyhow!(
"Port {} is occupied by PID {}. Use --kill-port-process to terminate.",
port,
pid
))
}