Skip to main content

systemprompt_agent/services/agent_orchestration/
mod.rs

1//! Supervision of agent worker processes: lifecycle, monitoring, and
2//! reconciliation.
3//!
4//! This module groups the services that keep the database's view of running
5//! agents consistent with the OS process table. [`AgentOrchestrator`] is the
6//! top-level facade; the submodules cover process lifecycle, health
7//! checks, drift reconciliation, port allocation, and the low-level process
8//! primitives. [`AgentStatus`] is the shared status model and
9//! [`OrchestrationError`] the unified error type.
10//!
11//! Copyright (c) systemprompt.io — Business Source License 1.1.
12//! See <https://systemprompt.io> for licensing details.
13
14pub mod database;
15pub mod lifecycle;
16pub mod monitor;
17pub mod orchestrator;
18pub mod port_service;
19pub mod process;
20pub mod reconciler;
21
22pub use orchestrator::{AgentInfo, AgentOrchestrator};
23pub use port_service::PortService;
24
25#[derive(Debug, Clone, PartialEq, Eq)]
26pub enum AgentStatus {
27    Running {
28        pid: u32,
29        port: u16,
30    },
31    Failed {
32        reason: String,
33        last_attempt: Option<String>,
34        retry_count: u32,
35    },
36}
37
38#[derive(Debug, Clone)]
39pub struct ValidationReport {
40    pub valid: bool,
41    pub issues: Vec<String>,
42}
43
44impl Default for ValidationReport {
45    fn default() -> Self {
46        Self::new()
47    }
48}
49
50impl ValidationReport {
51    pub const fn new() -> Self {
52        Self {
53            valid: true,
54            issues: Vec::new(),
55        }
56    }
57
58    pub fn with_issue(issue: String) -> Self {
59        Self {
60            valid: false,
61            issues: vec![issue],
62        }
63    }
64
65    pub fn add_issue(&mut self, issue: String) {
66        self.valid = false;
67        self.issues.push(issue);
68    }
69}
70
71use crate::services::shared::AgentServiceError;
72use systemprompt_identifiers::AgentName;
73use systemprompt_traits::{BoxedSource, RepositoryError};
74use thiserror::Error;
75
76#[derive(Error, Debug)]
77pub enum OrchestrationError {
78    #[error("Agent {0} not found")]
79    AgentNotFound(String),
80
81    #[error("Agent {0} already running")]
82    AgentAlreadyRunning(String),
83
84    #[error("Process spawn failed: {0}")]
85    ProcessSpawnFailed(String),
86
87    #[error("Process spawn failed: {context}: {source}")]
88    Spawn {
89        context: String,
90        #[source]
91        source: BoxedSource,
92    },
93
94    #[error("repository: {0}")]
95    Repository(#[from] RepositoryError),
96
97    #[error("IO error: {0}")]
98    IoError(#[from] std::io::Error),
99
100    #[error("Health check timeout for agent {0}")]
101    HealthCheckTimeout(String),
102
103    #[error("Failed to load agent registry: {0}")]
104    Registry(#[source] crate::error::AgentError),
105
106    #[error("agent: {0}")]
107    Agent(#[from] crate::error::AgentError),
108
109    #[error("Service error: {0}")]
110    AgentService(#[from] AgentServiceError),
111
112    #[error("process supervision: {0}")]
113    Supervision(#[from] systemprompt_loader::subprocess::SupervisionError),
114
115    #[error(
116        "port {port} for agent {agent} is held by process {pid}, which this installation did \
117         not spawn; stop it or choose a different port"
118    )]
119    PortHeldByForeignProcess {
120        port: u16,
121        pid: u32,
122        agent: AgentName,
123    },
124}
125
126impl OrchestrationError {
127    pub fn spawn<E>(context: impl Into<String>, source: E) -> Self
128    where
129        E: std::error::Error + Send + Sync + 'static,
130    {
131        Self::Spawn {
132            context: context.into(),
133            source: Box::new(source),
134        }
135    }
136}
137
138impl From<sqlx::Error> for OrchestrationError {
139    fn from(err: sqlx::Error) -> Self {
140        Self::Repository(err.into())
141    }
142}
143
144pub type OrchestrationResult<T> = Result<T, OrchestrationError>;