systemprompt-api 0.63.0

Axum-based HTTP server and API gateway for systemprompt.io AI governance infrastructure. Exposes governed agents, MCP, A2A, and admin endpoints with rate limiting and RBAC.
Documentation
//! MCP service reconciliation during server startup.
//!
//! `reconcile_system_services` cleans stale service rows, reconciles the MCP
//! orchestrator to the required set of enabled servers, and verifies each is
//! registered and running in the database — failing server startup loudly if
//! any required MCP server is missing, since agents depend on their tools.
//!
//! Copyright (c) systemprompt.io — Business Source License 1.1.
//! See <https://systemprompt.io> for licensing details.

use anyhow::Result;
use std::sync::Arc;
use systemprompt_identifiers::ServiceName;
use systemprompt_loader::subprocess::{self, ChildKind};
use systemprompt_manifest::services::ServiceStatus;
use systemprompt_runtime::AppContext;
use systemprompt_traits::{Phase, StartupEventExt, StartupEventSender};

struct ReconcileSuccessParams<'a> {
    running_count: usize,
    required_count: usize,
    required_servers: &'a [systemprompt_mcp::McpServerConfig],
    mcp_orchestrator: &'a Arc<systemprompt_mcp::services::McpOrchestrator>,
    ctx: &'a AppContext,
    events: Option<&'a StartupEventSender>,
}

pub(in crate::services::server) async fn reconcile_system_services(
    ctx: &AppContext,
    mcp_orchestrator: &Arc<systemprompt_mcp::services::McpOrchestrator>,
    events: Option<&StartupEventSender>,
) -> Result<()> {
    events.phase_started(Phase::McpServers);

    match cleanup_stale_service_entries(ctx, events).await {
        Ok(count) => {
            if count > 0 {
                events.mcp_service_cleanup(format!("{count} services"), "Stale entries removed");
            }
        },
        Err(e) => {
            events.warning(format!("Could not clean stale entries: {e}"));
        },
    }

    let required_servers = ctx.mcp_registry().get_managed_servers()?;
    let required_count = required_servers.len();

    match mcp_orchestrator.reconcile().await {
        Ok(running_count) => {
            handle_reconcile_success(ReconcileSuccessParams {
                running_count,
                required_count,
                required_servers: &required_servers,
                mcp_orchestrator,
                ctx,
                events,
            })
            .await?;
        },
        Err(e) => {
            events.phase_failed(Phase::McpServers, e.to_string());
            return Err(anyhow::anyhow!(
                "FATAL: MCP reconciliation failed: {}\n\nCannot start API without MCP servers.",
                e
            ));
        },
    }

    events.phase_completed(Phase::McpServers);
    Ok(())
}

async fn handle_reconcile_success(params: ReconcileSuccessParams<'_>) -> Result<()> {
    if params.running_count < params.required_count {
        return handle_missing_servers(
            params.required_servers,
            params.mcp_orchestrator,
            params.events,
        )
        .await;
    }

    if params.running_count > 0 {
        verify_database_registration(params.required_servers, params.ctx, params.events).await?;
    }

    params
        .events
        .mcp_reconciliation_complete(params.running_count, params.required_count);

    Ok(())
}

#[expect(
    clippy::collection_is_never_read,
    reason = "`events` is consumed by `StartupEventExt` trait methods \
              (`events.error(...)`); clippy's `collection_is_never_read` heuristic does not \
              recognise those calls as reads of the `Option`"
)]
pub async fn handle_missing_servers(
    required_servers: &[systemprompt_mcp::McpServerConfig],
    mcp_orchestrator: &Arc<systemprompt_mcp::services::McpOrchestrator>,
    events: Option<&StartupEventSender>,
) -> Result<()> {
    let running_servers = mcp_orchestrator.get_running_servers().await?;
    let running_names: std::collections::HashSet<String> =
        running_servers.iter().map(|s| s.name.clone()).collect();

    let missing: Vec<String> = required_servers
        .iter()
        .map(|s| s.name.clone())
        .filter(|name| !running_names.contains(name))
        .collect();

    events.error(
        format!(
            "Server status mismatch: {} servers failed to start: {}",
            missing.len(),
            missing.join(", ")
        ),
        true,
    );

    Err(anyhow::anyhow!(
        "FATAL: {} required MCP server(s) failed to start: {}\n\nsystemprompt.io OS cannot \
         operate without MCP servers.\nAgents need tools to function.\n\nBuild missing binaries \
         with:\n  cargo build --bin {}\n\nOr build all MCP servers:\n  systemprompt build mcp",
        missing.len(),
        missing.join(", "),
        missing.join(" --bin ")
    ))
}

pub const VERIFY_ATTEMPTS: u32 = 5;
pub const VERIFY_BACKOFF: std::time::Duration = std::time::Duration::from_millis(250);

pub async fn verify_database_registration(
    required_servers: &[systemprompt_mcp::McpServerConfig],
    ctx: &AppContext,
    events: Option<&StartupEventSender>,
) -> Result<()> {
    let mut pending: Vec<&systemprompt_mcp::McpServerConfig> = required_servers.iter().collect();
    let mut failures = Vec::new();

    for attempt in 1..=VERIFY_ATTEMPTS {
        let (still_pending, attempt_failures) = verify_once(&pending, ctx, events).await;
        pending = still_pending;
        failures = attempt_failures;
        if pending.is_empty() {
            return Ok(());
        }
        if attempt < VERIFY_ATTEMPTS {
            tokio::time::sleep(VERIFY_BACKOFF).await;
        }
    }

    events.error(
        format!(
            "Database verification failed for {} service(s): {}",
            failures.len(),
            failures.join(", ")
        ),
        true,
    );
    Err(anyhow::anyhow!(
        "FATAL: MCP services running but not properly registered in database after {} \
         attempts\n\nThis indicates a race condition or database synchronization \
         issue.\nFailed services: {}",
        VERIFY_ATTEMPTS,
        failures.join(", ")
    ))
}

#[expect(
    clippy::collection_is_never_read,
    reason = "`events` is consumed by StartupEventExt trait methods that clippy does not \
              recognise as reads"
)]
async fn verify_once<'a>(
    servers: &[&'a systemprompt_mcp::McpServerConfig],
    ctx: &AppContext,
    events: Option<&StartupEventSender>,
) -> (Vec<&'a systemprompt_mcp::McpServerConfig>, Vec<String>) {
    let service_repo = ctx.service_repository();
    let mut pending = Vec::new();
    let mut failures = Vec::new();

    for &server in servers {
        let failure = match service_repo
            .find_service_by_name(&ServiceName::new(server.name.as_str()))
            .await
        {
            Ok(Some(service)) if service.status == ServiceStatus::Running => {
                events.mcp_ready(
                    server.name.clone(),
                    service.port as u16,
                    std::time::Duration::ZERO,
                    None,
                );
                continue;
            },
            Ok(Some(service)) => format!("{} (status: {})", server.name, service.status),
            Ok(None) => format!("{} (not in database)", server.name),
            Err(e) => format!("{} (db error: {})", server.name, e),
        };
        pending.push(server);
        failures.push(failure);
    }

    (pending, failures)
}

#[expect(
    clippy::collection_is_never_read,
    reason = "`events` is consumed by StartupEventExt trait methods that clippy does not \
              recognise as reads"
)]
pub async fn cleanup_stale_service_entries(
    ctx: &AppContext,
    events: Option<&StartupEventSender>,
) -> Result<u64> {
    let repo = ctx.service_repository();
    let mut deleted_count = 0u64;

    let mcp_services = repo.list_mcp_services().await?;
    for service in mcp_services {
        if !service_row_is_stale(service.status, service.pid, ChildKind::Mcp, &service.name).await {
            continue;
        }
        if repo.delete_service(&service.name).await.is_ok() {
            deleted_count += 1;
            events.mcp_service_cleanup(
                service.name.as_str(),
                format!(
                    "Stale entry (status: {}, pid: {:?})",
                    service.status, service.pid
                ),
            );
        }
    }

    let agent_service_names = repo.list_all_agent_service_names().await?;
    for service_name in agent_service_names {
        if let Ok(Some(service)) = repo.find_service_by_name(&service_name).await {
            if !service_row_is_stale(service.status, service.pid, ChildKind::Agent, &service_name)
                .await
            {
                continue;
            }
            if repo.delete_service(&service_name).await.is_ok() {
                deleted_count += 1;
                events.agent_cleanup(
                    service_name.as_str(),
                    format!(
                        "Stale entry (status: {}, pid: {:?})",
                        service.status, service.pid
                    ),
                );
            }
        }
    }

    Ok(deleted_count)
}

pub async fn service_row_is_stale(
    status: ServiceStatus,
    pid: Option<i32>,
    kind: ChildKind,
    name: &ServiceName,
) -> bool {
    match status {
        ServiceStatus::Running => {
            let Some(pid) = pid.and_then(|p| u32::try_from(p).ok()) else {
                return true;
            };
            !subprocess::owns(pid, kind, name).await
        },
        ServiceStatus::Error | ServiceStatus::Stopped => true,
        ServiceStatus::Starting | ServiceStatus::Stopping => false,
    }
}