pub mod a2a;
pub mod bridge;
pub mod context;
pub mod error;
pub mod execution;
pub mod output;
pub mod persistence;
pub mod types;
use chrono::{DateTime, Utc};
use dashmap::DashMap;
use serde::{Deserialize, Serialize};
use std::sync::Arc;
use uuid::Uuid;
use self::error::{SessionError, SessionResult};
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct AutoAcceptConfig {
pub enabled: bool,
pub risk_threshold: u8,
}
use crate::identity::AgentRole;
use crate::resource::{ResourceMonitor, SessionResourceIntegration};
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub enum SessionStatus {
Active,
Paused,
Detached,
Background,
Terminated,
Error(String),
}
impl std::fmt::Display for SessionStatus {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
SessionStatus::Active => write!(f, "Active"),
SessionStatus::Paused => write!(f, "Paused"),
SessionStatus::Detached => write!(f, "Detached"),
SessionStatus::Background => write!(f, "Background"),
SessionStatus::Terminated => write!(f, "Terminated"),
SessionStatus::Error(e) => write!(f, "Error: {}", e),
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct AgentSession {
pub id: String,
pub agent_id: String,
pub agent_role: AgentRole,
pub status: SessionStatus,
pub background_mode: bool,
pub auto_accept: bool,
pub auto_accept_config: Option<AutoAcceptConfig>,
pub created_at: DateTime<Utc>,
pub last_activity: DateTime<Utc>,
pub description: Option<String>,
pub working_directory: String,
pub tasks_processed: usize,
pub tasks_queued: usize,
}
impl AgentSession {
pub fn new(
agent_id: String,
agent_role: AgentRole,
working_directory: String,
description: Option<String>,
) -> Self {
let session_id = Uuid::new_v4().to_string();
Self {
id: session_id,
agent_id,
agent_role,
status: SessionStatus::Active,
background_mode: false,
auto_accept: false,
auto_accept_config: None,
created_at: Utc::now(),
last_activity: Utc::now(),
description,
working_directory,
tasks_processed: 0,
tasks_queued: 0,
}
}
pub fn touch(&mut self) {
self.last_activity = Utc::now();
}
pub fn is_runnable(&self) -> bool {
matches!(
self.status,
SessionStatus::Active | SessionStatus::Background | SessionStatus::Detached
)
}
pub fn increment_tasks_processed(&mut self) {
self.tasks_processed += 1;
self.touch();
}
pub fn enable_auto_accept(&mut self, config: AutoAcceptConfig) {
self.auto_accept = true;
self.auto_accept_config = Some(config);
self.touch();
}
pub fn disable_auto_accept(&mut self) {
self.auto_accept = false;
self.auto_accept_config = None;
self.touch();
}
pub fn update_auto_accept_config(&mut self, config: AutoAcceptConfig) {
if self.auto_accept {
self.auto_accept_config = Some(config);
self.touch();
}
}
pub fn get_auto_accept_config(&self) -> Option<&AutoAcceptConfig> {
self.auto_accept_config.as_ref()
}
pub fn is_auto_accept_ready(&self) -> bool {
self.auto_accept && self.auto_accept_config.is_some()
}
}
pub struct SessionManager {
sessions: DashMap<String, AgentSession>,
resource_monitor: Option<Arc<ResourceMonitor>>,
resource_integration: Option<Arc<SessionResourceIntegration>>,
}
impl SessionManager {
pub async fn new() -> SessionResult<Self> {
Ok(Self {
sessions: DashMap::new(),
resource_monitor: None,
resource_integration: None,
})
}
pub async fn with_resource_monitoring(
resource_limits: crate::resource::ResourceLimits,
) -> SessionResult<Self> {
Ok(Self::with_resource_monitoring_inner(resource_limits))
}
fn with_resource_monitoring_inner(resource_limits: crate::resource::ResourceLimits) -> Self {
let resource_monitor = Arc::new(ResourceMonitor::new(resource_limits));
let resource_integration =
Arc::new(SessionResourceIntegration::new(resource_monitor.clone()));
if let Ok(runtime) = tokio::runtime::Handle::try_current() {
let monitor_clone = resource_monitor.clone();
runtime.spawn(async move {
monitor_clone.start_monitoring_loop().await;
});
}
Self {
sessions: DashMap::new(),
resource_monitor: Some(resource_monitor),
resource_integration: Some(resource_integration),
}
}
pub async fn create_session(
&self,
agent_id: String,
agent_role: AgentRole,
working_directory: String,
description: Option<String>,
auto_start: bool,
) -> SessionResult<AgentSession> {
let session =
AgentSession::new(agent_id, agent_role, working_directory.clone(), description);
self.setup_session_environment(&session).await?;
if auto_start {
self.start_agent_in_session(&session).await?;
}
self.sessions.insert(session.id.clone(), session.clone());
if let Some(ref integration) = self.resource_integration
&& let Err(e) = integration
.on_session_created(&session.id, &session.agent_id, None)
.await
{
tracing::warn!("Failed to start resource monitoring: {}", e);
}
Ok(session)
}
pub async fn pause_session(&self, session_id: &str) -> SessionResult<()> {
self.transition_session(session_id, "pause", |session| {
(session.status == SessionStatus::Active || session.status == SessionStatus::Background)
.then_some(SessionStatus::Paused)
})
}
pub async fn resume_session(&self, session_id: &str) -> SessionResult<()> {
self.transition_session(session_id, "resume", |session| {
(session.status == SessionStatus::Paused).then_some({
if session.background_mode {
SessionStatus::Background
} else {
SessionStatus::Active
}
})
})
}
fn with_session_mut<T>(
&self,
session_id: &str,
operation: impl FnOnce(&mut AgentSession) -> SessionResult<T>,
) -> SessionResult<T> {
let mut session =
self.sessions
.get_mut(session_id)
.ok_or_else(|| SessionError::NotFound {
id: session_id.to_string(),
})?;
operation(&mut session)
}
fn transition_session(
&self,
session_id: &str,
operation: &str,
next_status: impl FnOnce(&AgentSession) -> Option<SessionStatus>,
) -> SessionResult<()> {
self.with_session_mut(session_id, |session| {
let Some(status) = next_status(session) else {
return Err(SessionError::InvalidState {
state: format!("{:?}", session.status),
operation: operation.to_string(),
});
};
session.status = status;
session.touch();
Ok(())
})
}
pub async fn detach_session(&self, session_id: &str) -> SessionResult<()> {
self.transition_session(session_id, "detach", |session| {
(session.status == SessionStatus::Active || session.status == SessionStatus::Background)
.then_some(SessionStatus::Detached)
})
}
pub async fn attach_session(&self, session_id: &str) -> SessionResult<()> {
self.transition_session(session_id, "attach", |session| {
(session.status == SessionStatus::Detached).then_some({
if session.background_mode {
SessionStatus::Background
} else {
SessionStatus::Active
}
})
})
}
pub async fn terminate_session(&self, session_id: &str) -> SessionResult<()> {
let agent_id = {
let session = self
.sessions
.get(session_id)
.ok_or_else(|| SessionError::NotFound {
id: session_id.to_string(),
})?;
session.agent_id.clone()
};
if let Some(ref integration) = self.resource_integration
&& let Err(e) = integration
.on_session_terminated(session_id, &agent_id)
.await
{
tracing::warn!("Failed to stop resource monitoring: {}", e);
}
if let Some(mut session) = self.sessions.get_mut(session_id) {
session.status = SessionStatus::Terminated;
session.touch();
}
Ok(())
}
pub async fn set_background_mode(
&self,
session_id: &str,
auto_accept: bool,
) -> SessionResult<()> {
self.with_session_mut(session_id, |session| {
session.background_mode = true;
if auto_accept && session.auto_accept_config.is_none() {
session.enable_auto_accept(AutoAcceptConfig::default());
} else {
session.auto_accept = auto_accept;
}
if session.status == SessionStatus::Active {
session.status = SessionStatus::Background;
}
session.touch();
Ok(())
})
}
pub async fn enable_auto_accept(
&self,
session_id: &str,
config: AutoAcceptConfig,
) -> SessionResult<()> {
self.with_session_mut(session_id, |session| {
session.enable_auto_accept(config);
Ok(())
})
}
pub async fn disable_auto_accept(&self, session_id: &str) -> SessionResult<()> {
self.with_session_mut(session_id, |session| {
session.disable_auto_accept();
Ok(())
})
}
pub async fn update_auto_accept_config(
&self,
session_id: &str,
config: AutoAcceptConfig,
) -> SessionResult<()> {
self.with_session_mut(session_id, |session| {
session.update_auto_accept_config(config);
Ok(())
})
}
pub async fn emergency_stop_all_auto_accept(&self) -> usize {
let mut count = 0;
for mut entry in self.sessions.iter_mut() {
if entry.value().auto_accept {
entry.value_mut().disable_auto_accept();
count += 1;
}
}
count
}
pub fn get_auto_accept_sessions(&self) -> Vec<AgentSession> {
self.sessions
.iter()
.filter(|entry| entry.value().is_auto_accept_ready())
.map(|entry| entry.value().clone())
.collect()
}
pub fn get_session(&self, session_id: &str) -> Option<AgentSession> {
self.sessions.get(session_id).map(|entry| entry.clone())
}
pub fn list_sessions(&self) -> Vec<AgentSession> {
self.sessions
.iter()
.map(|entry| entry.value().clone())
.collect()
}
pub fn list_active_sessions(&self) -> Vec<AgentSession> {
self.sessions
.iter()
.filter(|entry| entry.value().is_runnable())
.map(|entry| entry.value().clone())
.collect()
}
pub fn get_sessions_by_role(&self, role: AgentRole) -> Vec<AgentSession> {
self.sessions
.iter()
.filter(|entry| entry.value().agent_role == role)
.map(|entry| entry.value().clone())
.collect()
}
pub async fn cleanup_terminated_sessions(&self) -> SessionResult<usize> {
let terminated: Vec<String> = self
.sessions
.iter()
.filter(|entry| entry.value().status == SessionStatus::Terminated)
.map(|entry| entry.key().clone())
.collect();
let count = terminated.len();
for id in terminated {
self.sessions.remove(&id);
}
Ok(count)
}
pub async fn check_and_suspend_idle_agents(&self) -> SessionResult<Vec<String>> {
let mut suspended_agents = Vec::new();
if let Some(ref integration) = self.resource_integration {
let agents_to_check: Vec<(String, String)> = self
.sessions
.iter()
.filter(|entry| entry.value().is_runnable())
.map(|entry| (entry.value().id.clone(), entry.value().agent_id.clone()))
.collect();
for (session_id, agent_id) in agents_to_check {
if integration.check_agent_suspension(&agent_id).await {
if let Err(e) = self.pause_session(&session_id).await {
tracing::warn!("Failed to suspend idle agent {}: {}", agent_id, e);
} else {
suspended_agents.push(agent_id);
}
}
}
}
Ok(suspended_agents)
}
pub fn get_session_resource_usage(
&self,
session_id: &str,
) -> Option<crate::resource::ResourceUsage> {
let session = self.sessions.get(session_id)?;
self.resource_monitor
.as_ref()
.and_then(|monitor| monitor.get_agent_usage(&session.agent_id))
}
pub fn get_resource_efficiency_stats(
&self,
) -> Option<crate::resource::ResourceEfficiencyStats> {
self.resource_monitor
.as_ref()
.map(|monitor| monitor.get_efficiency_stats())
}
async fn setup_session_environment(&self, _session: &AgentSession) -> SessionResult<()> {
Ok(())
}
async fn start_agent_in_session(&self, session: &AgentSession) -> SessionResult<()> {
let _command = format!(
"ccswarm agent start --id {} --role {} --session {} --working-dir {}",
session.agent_id,
session.agent_role.name(),
session.id,
session.working_directory
);
Ok(())
}
}
impl Default for SessionManager {
fn default() -> Self {
Self::with_resource_monitoring_inner(crate::resource::ResourceLimits::default())
}
}
#[cfg(test)]
mod tests {
use super::*;
fn backend_role() -> AgentRole {
AgentRole::Backend {
technologies: Vec::new(),
responsibilities: Vec::new(),
boundaries: Vec::new(),
}
}
#[tokio::test]
async fn session_transitions_preserve_background_mode() {
let manager = SessionManager::new().await.expect("manager");
let session = manager
.create_session(
"backend".to_string(),
backend_role(),
".".to_string(),
None,
false,
)
.await
.expect("session");
manager
.set_background_mode(&session.id, false)
.await
.expect("background");
manager.pause_session(&session.id).await.expect("pause");
manager.resume_session(&session.id).await.expect("resume");
assert_eq!(
manager.get_session(&session.id).expect("stored").status,
SessionStatus::Background
);
manager.detach_session(&session.id).await.expect("detach");
manager.attach_session(&session.id).await.expect("attach");
assert_eq!(
manager.get_session(&session.id).expect("stored").status,
SessionStatus::Background
);
}
#[tokio::test]
async fn session_transition_errors_remain_specific() {
let manager = SessionManager::new().await.expect("manager");
let missing = manager.pause_session("missing").await;
assert!(matches!(missing, Err(SessionError::NotFound { .. })));
let session = manager
.create_session(
"backend".to_string(),
backend_role(),
".".to_string(),
None,
false,
)
.await
.expect("session");
let invalid = manager.resume_session(&session.id).await;
assert!(matches!(
invalid,
Err(SessionError::InvalidState { operation, .. }) if operation == "resume"
));
}
#[test]
fn default_without_runtime_enables_resource_monitoring() {
let manager = SessionManager::default();
assert!(manager.resource_monitor.is_some());
assert!(manager.resource_integration.is_some());
}
#[tokio::test(flavor = "current_thread")]
async fn default_supports_current_thread_runtime() {
let manager = SessionManager::default();
assert!(manager.resource_monitor.is_some());
assert!(manager.resource_integration.is_some());
}
}