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
}
}