use super::*;
use crate::subprocess::{MockProcessRunner, ProcessRunner};
use crate::worktree::{WorktreeManager, WorktreePool, WorktreeSession};
use serde_json::json;
use std::sync::Arc;
use std::time::Duration;
use tokio::time::timeout;
#[tokio::test]
async fn test_resource_manager_creation() {
let manager = ResourceManager::new(None);
assert!(manager.worktree_pool.is_none());
assert!(manager.active_sessions.read().await.is_empty());
assert_eq!(manager.get_metrics().await.active_sessions, 0);
}
#[tokio::test]
async fn test_resource_manager_with_worktree_pool() {
let mock_runner = MockProcessRunner::new();
let subprocess =
crate::subprocess::SubprocessManager::new(Arc::new(mock_runner) as Arc<dyn ProcessRunner>);
let worktree_manager = WorktreeManager::new(std::path::PathBuf::from("/tmp"), subprocess)
.ok()
.map(Arc::new);
let config = Default::default();
let pool = worktree_manager
.clone()
.map(|manager| Arc::new(WorktreePool::new(config, manager)));
let manager = ResourceManager::new(pool);
assert!(manager.worktree_pool.is_some());
assert_eq!(manager.get_metrics().await.active_sessions, 0);
}
#[tokio::test]
async fn test_session_registration_and_unregistration() {
let manager = ResourceManager::new(None);
let session = WorktreeSession {
name: "test-worktree".to_string(),
branch: "test-branch".to_string(),
path: std::path::PathBuf::from("/tmp/test"),
created_at: chrono::Utc::now(),
};
manager
.register_session("agent-1".to_string(), session.clone())
.await;
let active = manager.get_active_sessions().await;
assert_eq!(active.len(), 1);
assert_eq!(active[0].0, "agent-1");
let metrics = manager.get_metrics().await;
assert_eq!(metrics.active_sessions, 1);
let unregistered = manager.unregister_session("agent-1").await;
assert!(unregistered.is_some());
let active = manager.get_active_sessions().await;
assert_eq!(active.len(), 0);
let metrics = manager.get_metrics().await;
assert_eq!(metrics.active_sessions, 0);
}
#[tokio::test]
async fn test_cleanup_orphaned_resources() {
let manager = ResourceManager::new(None);
manager.cleanup_orphaned_resources(&[]).await;
let worktree_names = vec!["worktree-1".to_string(), "worktree-2".to_string()];
manager.cleanup_orphaned_resources(&worktree_names).await;
}
#[tokio::test]
async fn test_cleanup_all_resources() {
let manager = ResourceManager::new(None);
for i in 0..3 {
let session = WorktreeSession {
name: format!("test-worktree-{}", i),
branch: "test-branch".to_string(),
path: std::path::PathBuf::from(format!("/tmp/test-{}", i)),
created_at: chrono::Utc::now(),
};
manager
.register_session(format!("agent-{}", i), session)
.await;
}
assert_eq!(manager.get_active_sessions().await.len(), 3);
let result = manager.cleanup_all().await;
assert!(result.is_ok());
assert_eq!(manager.get_active_sessions().await.len(), 0);
}
#[tokio::test]
async fn test_resource_metrics_tracking() {
let manager = ResourceManager::new(None);
let metrics = manager.get_metrics().await;
assert_eq!(metrics.active_sessions, 0);
assert_eq!(metrics.total_created, 0);
assert_eq!(metrics.total_reused, 0);
for i in 0..5 {
let session = WorktreeSession {
name: format!("test-worktree-{}", i),
branch: "test-branch".to_string(),
path: std::path::PathBuf::from(format!("/tmp/test-{}", i)),
created_at: chrono::Utc::now(),
};
manager
.register_session(format!("agent-{}", i), session)
.await;
}
let metrics = manager.get_metrics().await;
assert_eq!(metrics.active_sessions, 5);
}
#[tokio::test]
async fn test_resource_guard_cleanup() {
let counter = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let counter_clone = counter.clone();
{
let _guard = ResourceGuard::new(42, move |_value| {
counter_clone.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
});
assert_eq!(counter.load(std::sync::atomic::Ordering::Relaxed), 0);
}
assert_eq!(counter.load(std::sync::atomic::Ordering::Relaxed), 1);
}
#[tokio::test]
async fn test_resource_guard_take() {
let counter = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let counter_clone = counter.clone();
let guard = ResourceGuard::new(42, move |_value| {
counter_clone.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
});
let value = guard.take();
assert_eq!(value, Some(42));
assert_eq!(counter.load(std::sync::atomic::Ordering::Relaxed), 0);
}
#[tokio::test]
async fn test_agent_resource_manager() {
let agent_manager = AgentResourceManager::new();
let context = AgentContext::new(
"agent-1".to_string(),
json!({"test": "data"}),
0,
WorktreeSession {
name: "test-worktree".to_string(),
branch: "test-branch".to_string(),
path: std::path::PathBuf::from("/tmp/test"),
created_at: chrono::Utc::now(),
},
Default::default(),
);
agent_manager
.register_context("agent-1".to_string(), context.clone())
.await;
let retrieved = agent_manager.get_context("agent-1").await;
assert!(retrieved.is_some());
assert_eq!(retrieved.unwrap().agent_id, "agent-1");
let all_contexts = agent_manager.get_active_contexts().await;
assert_eq!(all_contexts.len(), 1);
let unregistered = agent_manager.unregister_context("agent-1").await;
assert!(unregistered.is_some());
let retrieved = agent_manager.get_context("agent-1").await;
assert!(retrieved.is_none());
}
#[tokio::test]
async fn test_agent_context_initialization() {
let agent_manager = AgentResourceManager::new();
let session = WorktreeSession {
name: "test-worktree".to_string(),
branch: "test-branch".to_string(),
path: std::path::PathBuf::from("/tmp/test"),
created_at: chrono::Utc::now(),
};
let item = json!({"key": "value", "number": 42});
let context =
agent_manager.initialize_agent_context("agent-test", &item, 5, &session, "correlation-123");
assert_eq!(context.get("agent_id").unwrap(), &json!("agent-test"));
assert_eq!(context.get("item").unwrap(), &item);
assert_eq!(context.get("item_index").unwrap(), &json!(5));
let map_context = context.get("map").unwrap();
assert!(map_context.is_object());
assert_eq!(map_context["job_id"], json!("correlation-123"));
assert_eq!(map_context["agent"]["id"], json!("agent-test"));
assert_eq!(map_context["agent"]["index"], json!(5));
}
#[tokio::test]
async fn test_cleanup_registry_operations() {
let registry = CleanupRegistry::new();
struct TestCleanupTask {
executed: Arc<std::sync::atomic::AtomicBool>,
}
#[async_trait::async_trait]
impl CleanupTask for TestCleanupTask {
async fn cleanup(&self) -> MapReduceResult<()> {
self.executed
.store(true, std::sync::atomic::Ordering::Relaxed);
Ok(())
}
fn priority(&self) -> CleanupPriority {
CleanupPriority::Normal
}
}
let executed = Arc::new(std::sync::atomic::AtomicBool::new(false));
let task = Box::new(TestCleanupTask {
executed: executed.clone(),
});
registry.register(task).await;
let result = registry.execute_all().await;
assert!(result.is_ok());
assert!(executed.load(std::sync::atomic::Ordering::Relaxed));
}
#[tokio::test]
async fn test_resource_manager_stress() {
let manager = Arc::new(ResourceManager::new(None));
let mut handles = vec![];
for i in 0..10 {
let manager_clone = manager.clone();
let handle = tokio::spawn(async move {
for j in 0..10 {
let session = WorktreeSession {
name: format!("test-worktree-{}-{}", i, j),
branch: "test-branch".to_string(),
path: std::path::PathBuf::from(format!("/tmp/test-{}-{}", i, j)),
created_at: chrono::Utc::now(),
};
let agent_id = format!("agent-{}-{}", i, j);
manager_clone
.register_session(agent_id.clone(), session)
.await;
tokio::time::sleep(Duration::from_millis(1)).await;
manager_clone.unregister_session(&agent_id).await;
}
});
handles.push(handle);
}
for handle in handles {
let _ = timeout(Duration::from_secs(5), handle).await;
}
assert_eq!(manager.get_active_sessions().await.len(), 0);
}
#[tokio::test]
async fn test_worktree_error_creation() {
let agent_manager = AgentResourceManager::new();
let error = agent_manager.create_worktree_error(
"test-agent",
"Failed to create worktree".to_string(),
"correlation-789",
);
match error {
crate::cook::execution::errors::MapReduceError::WorktreeCreationFailed {
agent_id,
reason,
..
} => {
assert_eq!(agent_id, "test-agent");
assert_eq!(reason, "Failed to create worktree");
}
_ => panic!("Expected WorktreeCreationFailed error"),
}
}
#[tokio::test]
async fn test_cleanup_task_priority_ordering() {
let registry = CleanupRegistry::new();
let execution_order = Arc::new(std::sync::Mutex::new(Vec::new()));
struct PriorityTestTask {
name: String,
priority: CleanupPriority,
execution_order: Arc<std::sync::Mutex<Vec<String>>>,
}
#[async_trait::async_trait]
impl CleanupTask for PriorityTestTask {
async fn cleanup(&self) -> MapReduceResult<()> {
let mut order = self.execution_order.lock().unwrap();
order.push(self.name.clone());
Ok(())
}
fn priority(&self) -> CleanupPriority {
self.priority
}
}
let tasks = vec![
("normal", CleanupPriority::Normal),
("high", CleanupPriority::High),
("low", CleanupPriority::Low),
("critical", CleanupPriority::Critical),
];
for (name, priority) in tasks {
let task = Box::new(PriorityTestTask {
name: name.to_string(),
priority,
execution_order: execution_order.clone(),
});
registry.register(task).await;
}
registry.execute_all().await.unwrap();
let order = execution_order.lock().unwrap();
assert_eq!(order[0], "critical");
assert_eq!(order[1], "high");
assert_eq!(order[2], "normal");
assert_eq!(order[3], "low");
}