/
kleidinc
/
brain
Обзор
Документация
Войти
/
kleidinc
/
brain
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
master
src/session.rs
378 строк
12 KB
Your Name
docs: complete rewrite of documentation suite
20 мар 2026, 14:16
20 мар 2026, 14:16
64a640d
Код
Авторство
О чём код?
use anyhow::Result; use serde::{Deserialize, Serialize}; use std::collections::HashMap; use std::sync::Arc; use tokio::sync::RwLock; use crate::storage::{ConversationMemoryStore, SessionMetadata, SessionStatus}; /// Session manager for conversation lifecycle and context pub struct SessionManager { sessions: Arc<RwLock<HashMap<String, Arc<RwLock<SessionData>>>>>, memory_store: Arc<ConversationMemoryStore>, config: SessionConfig, } #[derive(Debug, Clone)] pub struct SessionConfig { pub max_sessions_per_agent: usize, pub session_timeout_hours: u64, pub cleanup_interval_minutes: u64, pub max_context_messages: usize, } impl Default for SessionConfig { fn default() -> Self { Self { max_sessions_per_agent: 10, session_timeout_hours: 24, cleanup_interval_minutes: 15, max_context_messages: 20, } } } #[derive(Debug, Clone, Serialize, Deserialize)] pub struct SessionData { pub metadata: SessionMetadata, pub last_access: std::time::Instant, pub is_active: bool, } impl SessionData { pub fn new(agent_id: &str, client_type: &str, project_context: &str) -> Self { let session_id = uuid::Uuid::new_v4().to_string(); let now = chrono::Utc::now(); let timestamp = now.to_rfc3339(); Self { metadata: SessionMetadata { session_id: session_id.clone(), agent_id: agent_id.to_string(), client_type: client_type.to_string(), project_context: project_context.to_string(), preferences: serde_json::Value::Object(serde_json::Map::new()), start_time: timestamp.clone(), last_active: timestamp, message_count: 0, status: SessionStatus::Active, }, last_access: std::time::Instant::now(), is_active: true, } } pub fn update_activity(&mut self) { self.last_access = std::time::Instant::now(); self.metadata.last_active = chrono::Utc::now().to_rfc3339(); } pub fn add_message(&mut self) { self.metadata.message_count += 1; self.update_activity(); } pub fn is_expired(&self, timeout_hours: u64) -> bool { self.last_access.elapsed().as_secs() > (timeout_hours * 3600) } pub fn pause(&mut self) { self.metadata.status = SessionStatus::Paused; self.is_active = false; } pub fn resume(&mut self) { self.metadata.status = SessionStatus::Active; self.is_active = true; self.update_activity(); } pub fn complete(&mut self) { self.metadata.status = SessionStatus::Completed; self.is_active = false; } } impl SessionManager { pub fn new(memory_store: Arc<ConversationMemoryStore>, config: SessionConfig) -> Self { let manager = Self { sessions: Arc::new(RwLock::new(HashMap::new())), memory_store, config, }; // Start cleanup task let manager_clone = manager.clone(); tokio::spawn(async move { manager_clone.cleanup_task().await; }); manager } /// Create a new session for an agent pub async fn create_session( &self, agent_id: &str, client_type: &str, project_context: &str, ) -> Result<String> { let mut sessions = self.sessions.write().await; // Clean up expired sessions for this agent self.cleanup_agent_sessions(&mut sessions, agent_id).await; // Check session limit let agent_sessions: Vec<_> = sessions .iter() .filter(|(_, data_lock)| { tokio::task::block_in_place(|| { futures::executor::block_on(async { data_lock.read().await.metadata.agent_id == agent_id }) }) }) .collect(); if agent_sessions.len() >= self.config.max_sessions_per_agent { // Complete oldest session let oldest = agent_sessions .iter() .min_by_key(|(_, data_lock)| { tokio::task::block_in_place(|| { futures::executor::block_on(async { data_lock.read().await.last_access }) }) }); if let Some((session_id, _)) = oldest { if let Some(session_data) = sessions.remove(*session_id) { let mut data = session_data.write().await; data.complete(); let _ = self.memory_store.upsert_session(&data.metadata).await; } } } // Create new session let session_data = Arc::new(RwLock::new(SessionData::new(agent_id, client_type, project_context))); let session_id = session_data.read().await.metadata.session_id.clone(); // Persist to database self.memory_store.upsert_session(&session_data.read().await.metadata).await?; sessions.insert(session_id.clone(), session_data); tracing::info!("Created conversation session: {} for agent: {}", session_id, agent_id); Ok(session_id) } /// Get or create active session for agent pub async fn get_or_create_session( &self, agent_id: &str, client_type: &str, project_context: &str, ) -> Result<String> { // First try to find existing active session if let Some(session_id) = self.find_active_session(agent_id, project_context).await { return Ok(session_id); } // Create new session if none found self.create_session(agent_id, client_type, project_context).await } /// Find active session for agent and project context pub async fn find_active_session(&self, agent_id: &str, project_context: &str) -> Option<String> { let sessions = self.sessions.read().await; for (session_id, data_lock) in sessions.iter() { let data = data_lock.read().await; if data.is_active && data.metadata.agent_id == agent_id && data.metadata.project_context == project_context && data.metadata.status == SessionStatus::Active { return Some(session_id.clone()); } } None } /// Get session data by ID pub async fn get_session(&self, session_id: &str) -> Option<Arc<RwLock<SessionData>>> { let sessions = self.sessions.read().await; sessions.get(session_id).cloned() } /// Update session activity and persist pub async fn touch_session(&self, session_id: &str) -> Result<()> { if let Some(session_data) = self.get_session(session_id).await { let mut data = session_data.write().await; data.update_activity(); self.memory_store.upsert_session(&data.metadata).await?; } Ok(()) } /// Record a message in the session pub async fn record_message(&self, session_id: &str) -> Result<()> { if let Some(session_data) = self.get_session(session_id).await { let mut data = session_data.write().await; data.add_message(); self.memory_store.upsert_session(&data.metadata).await?; } Ok(()) } /// Pause a session pub async fn pause_session(&self, session_id: &str) -> Result<()> { if let Some(session_data) = self.get_session(session_id).await { let mut data = session_data.write().await; data.pause(); self.memory_store.upsert_session(&data.metadata).await?; } Ok(()) } /// Resume a paused session pub async fn resume_session(&self, session_id: &str) -> Result<()> { if let Some(session_data) = self.get_session(session_id).await { let mut data = session_data.write().await; data.resume(); self.memory_store.upsert_session(&data.metadata).await?; } Ok(()) } /// Complete a session pub async fn complete_session(&self, session_id: &str) -> Result<()> { if let Some(session_data) = self.sessions.write().await.remove(session_id) { let mut data = session_data.write().await; data.complete(); self.memory_store.upsert_session(&data.metadata).await?; } Ok(()) } /// Get all active sessions for an agent pub async fn get_agent_sessions(&self, agent_id: &str) -> Vec<SessionMetadata> { let sessions = self.sessions.read().await; let mut result = Vec::new(); for data_lock in sessions.values() { let data = data_lock.read().await; if data.metadata.agent_id == agent_id { result.push(data.metadata.clone()); } } result } /// Get session statistics pub async fn get_stats(&self) -> SessionStats { let sessions = self.sessions.read().await; let mut active = 0; let mut paused = 0; let mut completed = 0; let mut total_messages = 0; for data_lock in sessions.values() { let data = data_lock.read().await; total_messages += data.metadata.message_count; match data.metadata.status { SessionStatus::Active => active += 1, SessionStatus::Paused => paused += 1, SessionStatus::Completed => completed += 1, SessionStatus::Expired => {} // Don't count expired sessions } } SessionStats { active_sessions: active, paused_sessions: paused, completed_sessions: completed, total_sessions: sessions.len(), total_messages, } } /// Periodically clean up expired sessions async fn cleanup_task(self: Arc<Self>) { let mut interval = tokio::time::interval( std::time::Duration::from_secs(self.config.cleanup_interval_minutes * 60) ); loop { interval.tick().await; let mut sessions = self.sessions.write().await; let mut to_remove = Vec::new(); for (session_id, data_lock) in sessions.iter() { let data = data_lock.read().await; if data.is_expired(self.config.session_timeout_hours) { to_remove.push(session_id.clone()); } } for session_id in to_remove { if let Some(session_data) = sessions.remove(&session_id) { let mut data = session_data.write().await; data.metadata.status = SessionStatus::Expired; let _ = self.memory_store.upsert_session(&data.metadata).await; } } if !to_remove.is_empty() { tracing::info!("Cleaned up {} expired sessions", to_remove.len()); } } } /// Clean up expired sessions for a specific agent async fn cleanup_agent_sessions( &self, sessions: &mut HashMap<String, Arc<RwLock<SessionData>>>, agent_id: &str, ) { let mut to_remove = Vec::new(); for (session_id, data_lock) in sessions.iter() { let data = data_lock.read().await; if data.metadata.agent_id == agent_id && data.is_expired(self.config.session_timeout_hours) { to_remove.push(session_id.clone()); } } for session_id in to_remove { if let Some(session_data) = sessions.remove(&session_id) { let mut data = session_data.write().await; data.metadata.status = SessionStatus::Expired; let _ = self.memory_store.upsert_session(&data.metadata).await; } } } fn clone(&self) -> Self { Self { sessions: self.sessions.clone(), memory_store: self.memory_store.clone(), config: self.config.clone(), } } } #[derive(Debug, Clone, Serialize, Deserialize)] pub struct SessionStats { pub active_sessions: usize, pub paused_sessions: usize, pub completed_sessions: usize, pub total_sessions: usize, pub total_messages: usize, }