use anyhow::Result;
use prodigy::cook::execution::ResumeLockManager;
use std::sync::Arc;
use std::time::Duration;
use tempfile::TempDir;
use tokio::time::sleep;
#[tokio::test]
async fn test_concurrent_resume_attempts_blocked() -> Result<()> {
let temp_dir = TempDir::new()?;
let manager = Arc::new(ResumeLockManager::new(temp_dir.path().to_path_buf())?);
let job_id = "test-job-concurrent";
let manager1 = manager.clone();
let job_id1 = job_id.to_string();
let handle1 = tokio::spawn(async move {
let lock = manager1.acquire_lock(&job_id1).await;
if lock.is_ok() {
sleep(Duration::from_millis(100)).await;
}
lock
});
let manager2 = manager.clone();
let job_id2 = job_id.to_string();
let handle2 = tokio::spawn(async move {
sleep(Duration::from_millis(10)).await;
manager2.acquire_lock(&job_id2).await
});
let result1 = handle1.await?;
let result2 = handle2.await?;
assert!(
(result1.is_ok() && result2.is_err()) || (result1.is_err() && result2.is_ok()),
"Exactly one lock acquisition should succeed"
);
let error = if result1.is_err() {
result1.unwrap_err()
} else {
result2.unwrap_err()
};
assert!(error.to_string().contains("already in progress"));
Ok(())
}
#[tokio::test]
async fn test_sequential_resume_succeeds() -> Result<()> {
let temp_dir = TempDir::new()?;
let manager = ResumeLockManager::new(temp_dir.path().to_path_buf())?;
let job_id = "test-job-sequential";
{
let _lock = manager.acquire_lock(job_id).await?;
sleep(Duration::from_millis(50)).await;
}
let lock2 = manager.acquire_lock(job_id).await;
assert!(lock2.is_ok(), "Second resume should succeed");
Ok(())
}
#[tokio::test]
async fn test_resume_after_crash_cleans_stale_lock() -> Result<()> {
let temp_dir = TempDir::new()?;
let manager = ResumeLockManager::new(temp_dir.path().to_path_buf())?;
let job_id = "test-job-stale";
let lock_path = temp_dir
.path()
.join("resume_locks")
.join(format!("{}.lock", job_id));
std::fs::create_dir_all(lock_path.parent().unwrap())?;
let stale_lock_data = serde_json::json!({
"job_id": job_id,
"process_id": 999999, "hostname": "test-host",
"acquired_at": chrono::Utc::now().to_rfc3339()
});
std::fs::write(&lock_path, serde_json::to_string(&stale_lock_data)?)?;
let lock = manager.acquire_lock(job_id).await;
assert!(
lock.is_ok(),
"Resume should clean up stale lock and succeed"
);
Ok(())
}
#[tokio::test]
async fn test_lock_error_message_includes_details() -> Result<()> {
let temp_dir = TempDir::new()?;
let manager = ResumeLockManager::new(temp_dir.path().to_path_buf())?;
let job_id = "test-job-error-msg";
let _lock1 = manager.acquire_lock(job_id).await?;
let result = manager.acquire_lock(job_id).await;
assert!(result.is_err());
let error = result.unwrap_err().to_string();
assert!(error.contains("already in progress"));
assert!(error.contains("PID"));
assert!(error.contains(job_id));
Ok(())
}
#[tokio::test]
async fn test_lock_released_on_task_panic() -> Result<()> {
let temp_dir = TempDir::new()?;
let manager = Arc::new(ResumeLockManager::new(temp_dir.path().to_path_buf())?);
let job_id = "test-job-panic";
let manager_clone = manager.clone();
let job_id_clone = job_id.to_string();
let handle = tokio::spawn(async move {
let _lock = manager_clone.acquire_lock(&job_id_clone).await.unwrap();
sleep(Duration::from_millis(50)).await;
});
let _ = handle.await;
let lock = manager.acquire_lock(job_id).await;
assert!(
lock.is_ok(),
"Lock should be available after task completes"
);
Ok(())
}
#[tokio::test]
async fn test_multiple_jobs_independent_locks() -> Result<()> {
let temp_dir = TempDir::new()?;
let manager = Arc::new(ResumeLockManager::new(temp_dir.path().to_path_buf())?);
let mut handles = vec![];
for i in 0..5 {
let manager_clone = manager.clone();
let handle = tokio::spawn(async move {
let job_id = format!("job-{}", i);
manager_clone.acquire_lock(&job_id).await
});
handles.push(handle);
}
let results = futures::future::join_all(handles).await;
for result in results {
assert!(result.is_ok());
assert!(result.unwrap().is_ok());
}
Ok(())
}