systemprompt-agent 0.63.0

Agent-to-Agent (A2A) protocol for systemprompt.io AI governance: streaming, JSON-RPC models, task lifecycle, .well-known discovery, and governed agent orchestration.
Documentation
//! The agent supervision facade tying lifecycle, monitoring, and reconciliation
//! together.
//!
//! [`AgentOrchestrator`] is constructed once with shared [`AgentState`] and
//! [`AppPaths`], runs startup reconciliation, and exposes the operator-facing
//! verbs (start/stop/restart, status, health checks, bulk start/disable) by
//! delegating to the owned lifecycle, monitor, and reconciler services. The
//! `cleanup`, `reconciliation`, and `status` submodules carry agent deletion,
//! startup reconciliation, and status aggregation respectively.
//!
//! Copyright (c) systemprompt.io — Business Source License 1.1.
//! See <https://systemprompt.io> for licensing details.

mod cleanup;
mod reconciliation;
mod status;

use std::sync::Arc;
use systemprompt_traits::{Phase, StartupEvent, StartupEventExt, StartupEventSender};

use crate::services::agent_orchestration::database::AgentDatabaseService;
use crate::services::agent_orchestration::lifecycle::AgentLifecycle;
use crate::services::agent_orchestration::monitor::AgentMonitor;
use crate::services::agent_orchestration::reconciler::AgentReconciler;
use crate::services::agent_orchestration::{AgentStatus, OrchestrationResult, monitor};
use crate::state::AgentState;
use systemprompt_config::paths::AppPaths;
use systemprompt_identifiers::AgentName;

#[derive(Debug, Clone)]
pub struct AgentInfo {
    pub name: AgentName,
    pub status: AgentStatus,
    pub port: u16,
}

#[derive(Debug)]
pub struct AgentOrchestrator {
    pub(super) db_service: AgentDatabaseService,
    pub(super) lifecycle: AgentLifecycle,
    pub(super) reconciler: AgentReconciler,
    monitor: AgentMonitor,
}

impl AgentOrchestrator {
    pub async fn new(
        agent_state: Arc<AgentState>,
        app_paths: Arc<AppPaths>,
        events: Option<&StartupEventSender>,
    ) -> OrchestrationResult<Self> {
        tracing::debug!("Initializing Agent Orchestrator");

        let agent_repo = agent_state.repositories().agent_services.clone();

        let db_service = AgentDatabaseService::new(agent_repo.clone())?;
        let lifecycle = AgentLifecycle::new(agent_repo.clone(), app_paths)?;
        let reconciler = AgentReconciler::new(agent_repo.clone())?;
        let monitor = AgentMonitor::new(agent_repo)?;

        let orchestrator = Self {
            db_service,
            lifecycle,
            reconciler,
            monitor,
        };

        orchestrator.startup_reconciliation(events).await?;

        tracing::debug!("Agent Orchestrator initialized");
        Ok(orchestrator)
    }

    pub fn set_registry(&mut self, registry: crate::services::registry::AgentRegistry) {
        self.db_service.registry = registry.clone();
        self.lifecycle.db_service.registry = registry;
    }

    pub async fn start_agent(
        &self,
        agent_name: &AgentName,
        events: Option<&StartupEventSender>,
    ) -> OrchestrationResult<String> {
        self.lifecycle.start_agent(agent_name, events).await
    }

    pub async fn enable_agent(
        &self,
        agent_name: &AgentName,
        events: Option<&StartupEventSender>,
    ) -> OrchestrationResult<String> {
        self.lifecycle.enable_agent(agent_name, events).await
    }

    pub async fn disable_agent(&self, agent_name: &AgentName) -> OrchestrationResult<()> {
        self.lifecycle.disable_agent(agent_name).await
    }

    pub async fn restart_agent(
        &self,
        agent_name: &AgentName,
        events: Option<&StartupEventSender>,
    ) -> OrchestrationResult<String> {
        self.lifecycle.restart_agent(agent_name, events).await
    }

    pub async fn get_status(&self, agent_name: &AgentName) -> OrchestrationResult<AgentStatus> {
        self.db_service.get_status(agent_name).await
    }

    pub async fn list_agents(&self) -> OrchestrationResult<Vec<(AgentName, AgentStatus)>> {
        self.db_service.list_all_agents().await
    }

    pub async fn health_check(
        &self,
        agent_name: &AgentName,
    ) -> OrchestrationResult<monitor::HealthCheckResult> {
        self.monitor.comprehensive_health_check(agent_name).await
    }

    pub async fn disable_all(&self) -> OrchestrationResult<()> {
        let agents = self.db_service.list_all_agents().await?;

        for (agent_name, _) in agents {
            if let Err(e) = self.disable_agent(&agent_name).await {
                tracing::error!(agent_name = %agent_name, error = %e, "Failed to disable agent");
            }
        }

        Ok(())
    }

    pub async fn reconcile(&self, events: Option<&StartupEventSender>) -> OrchestrationResult<()> {
        if let Some(tx) = events {
            tx.phase_started(Phase::Agents);
        }

        self.startup_reconciliation(events).await?;

        let agents = self.db_service.list_all_agents().await?;
        let running = agents
            .iter()
            .filter(|(_, s)| matches!(s, AgentStatus::Running { .. }))
            .count();
        let total = agents.len();

        if let Some(tx) = events {
            if tx
                .unbounded_send(StartupEvent::AgentReconciliationComplete { running, total })
                .is_err()
            {
                tracing::trace!("No receivers for agent reconciliation event");
            }
            tx.phase_completed(Phase::Agents);
        }

        Ok(())
    }

    pub async fn update_agent_running(
        &self,
        agent_name: &AgentName,
        pid: u32,
        port: u16,
    ) -> OrchestrationResult<()> {
        self.db_service
            .update_agent_running(agent_name, pid, port)
            .await
    }

    pub async fn update_agent_stopped(&self, agent_name: &AgentName) -> OrchestrationResult<()> {
        self.db_service.update_agent_stopped(agent_name).await
    }
}