//! Integration tests for the Werkbank job queue (WB-02). //! //! Exercises the atomic lease/heartbeat/complete/sweep flow against a real //! MongoDB — the guarantees (idempotent enqueue, single-owner lease, visibility //! timeout) are Mongo-semantics-dependent and can't be unit-tested in isolation. //! Skips cleanly when no Mongo is reachable (set `TEST_MONGODB_URI` to point at //! one; defaults to the local dev cluster). #![allow(clippy::expect_used, clippy::unwrap_used)] use std::time::Duration; use chrono::{DateTime, TimeZone, Utc}; use compliance_agent::database::Database; use compliance_agent::werkbank::JobQueue; use compliance_core::models::werkbank::{Executor, InputRef, Job, JobResult}; /// Connect + ensure indexes on a throwaway database, or `None` if no Mongo. async fn setup() -> Option<(JobQueue, mongodb::Database)> { let uri = std::env::var("TEST_MONGODB_URI") .unwrap_or_else(|_| "mongodb://root:example@localhost:27017/?authSource=admin".into()); let db_name = format!("wbq_{}", &uuid::Uuid::new_v4().simple().to_string()[..12]); let db = match Database::connect(&uri, &db_name).await { Ok(d) => d, Err(_) => { eprintln!("SKIP werkbank_queue: no MongoDB reachable at {uri}"); return None; } }; db.ensure_indexes().await.expect("ensure indexes"); let queue = JobQueue::new(&db); Some((queue, db.inner().clone())) } fn base_time() -> DateTime { Utc.timestamp_opt(1_700_000_000, 0).unwrap() } fn job(id: &str) -> Job { Job::plc_provision(id, "acme", "target-1", InputRef::blob("sha256:abc"), 180) } fn job_with_labels(id: &str, labels: &[&str]) -> Job { let mut j = job(id); j.labels = labels.iter().map(|s| s.to_string()).collect(); j } macro_rules! skip_if_no_mongo { () => { match setup().await { Some(v) => v, None => return, } }; } #[tokio::test] async fn enqueue_is_idempotent() { let (q, db) = skip_if_no_mongo!(); let now = base_time(); assert!(q.enqueue(job("j1"), now).await.expect("enqueue")); // Same id again — no duplicate row, reports "already present". assert!(!q.enqueue(job("j1"), now).await.expect("enqueue2")); let rec = q.get("j1").await.expect("get").expect("exists"); assert_eq!( rec.status, compliance_core::models::werkbank::JobStatus::Queued ); assert_eq!(rec.attempts, 0); db.drop().await.ok(); } #[tokio::test] async fn lease_matches_executor_and_labels_and_is_fifo() { let (q, db) = skip_if_no_mongo!(); let t0 = base_time(); // Two docker jobs (j_old older than j_new) + one requiring a kvm label. q.enqueue(job("j_old"), t0).await.unwrap(); q.enqueue(job("j_new"), t0 + chrono::Duration::seconds(5)) .await .unwrap(); q.enqueue(job_with_labels("j_kvm", &["kvm=true"]), t0) .await .unwrap(); // Wrong executor: a shell runner leases nothing. assert!(q .lease("r-shell", Executor::Shell, &[], Duration::from_secs(30), t0) .await .unwrap() .is_none()); // A docker runner without the kvm label gets the oldest label-free job (FIFO). let leased = q .lease("r1", Executor::Docker, &[], Duration::from_secs(30), t0) .await .unwrap() .expect("leased"); assert_eq!(leased.job.id, "j_old", "oldest matching job first"); assert!(!leased.lease_token.is_empty()); // The kvm job stays unleased for that runner (missing label)... let none = q .lease("r1", Executor::Docker, &[], Duration::from_secs(30), t0) .await .unwrap() .expect("next"); assert_eq!(none.job.id, "j_new", "label-free job, not the kvm one"); // ...but a runner advertising kvm can take it. let kvm = q .lease( "r2", Executor::Docker, &["kvm=true".to_string(), "arch=amd64".to_string()], Duration::from_secs(30), t0, ) .await .unwrap() .expect("kvm leased"); assert_eq!(kvm.job.id, "j_kvm"); // A leased job increments attempts and is no longer queued. let rec = q.get("j_old").await.unwrap().unwrap(); assert_eq!(rec.attempts, 1); assert_eq!(rec.leased_by.as_deref(), Some("r1")); db.drop().await.ok(); } #[tokio::test] async fn heartbeat_extends_lease_and_surfaces_cancel() { let (q, db) = skip_if_no_mongo!(); let now = base_time(); q.enqueue(job("j1"), now).await.unwrap(); let leased = q .lease("r1", Executor::Docker, &[], Duration::from_secs(30), now) .await .unwrap() .unwrap(); // A valid heartbeat moves it to running and reports not-cancelled. let ack = q .heartbeat("j1", &leased.lease_token, Duration::from_secs(30), now) .await .unwrap() .expect("valid lease"); assert!(!ack.cancelled); assert_eq!( q.get("j1").await.unwrap().unwrap().status, compliance_core::models::werkbank::JobStatus::Running ); // A wrong token is a lost lease. assert!(q .heartbeat("j1", "wrong-token", Duration::from_secs(30), now) .await .unwrap() .is_none()); // Cancelling an in-flight job flags it; the next heartbeat reports cancelled. assert!(q.cancel("j1", now).await.unwrap()); let ack = q .heartbeat("j1", &leased.lease_token, Duration::from_secs(30), now) .await .unwrap() .expect("still leased"); assert!(ack.cancelled); db.drop().await.ok(); } #[tokio::test] async fn complete_is_idempotent_and_token_guarded() { let (q, db) = skip_if_no_mongo!(); let now = base_time(); q.enqueue(job("j1"), now).await.unwrap(); let leased = q .lease("r1", Executor::Docker, &[], Duration::from_secs(30), now) .await .unwrap() .unwrap(); // Wrong token cannot complete. let mut result = JobResult::succeeded("j1"); result.findings = Vec::new(); assert!(!q.complete("j1", "nope", &result, now).await.unwrap()); // The lease holder completes it once... assert!(q .complete("j1", &leased.lease_token, &result, now) .await .unwrap()); let rec = q.get("j1").await.unwrap().unwrap(); assert_eq!( rec.status, compliance_core::models::werkbank::JobStatus::Succeeded ); assert!(rec.result.is_some()); assert!(rec.lease_token.is_none(), "lease cleared on completion"); // ...and a second (duplicate) completion is a no-op. assert!(!q .complete("j1", &leased.lease_token, &result, now) .await .unwrap()); db.drop().await.ok(); } #[tokio::test] async fn sweep_requeues_expired_then_expires_after_max_attempts() { let (q, db) = skip_if_no_mongo!(); let t0 = base_time(); q.enqueue(job("j1"), t0).await.unwrap(); // Lease #1 with a 10s TTL; then time jumps past expiry. q.lease("r1", Executor::Docker, &[], Duration::from_secs(10), t0) .await .unwrap() .unwrap(); let past = t0 + chrono::Duration::seconds(60); // attempts=1 < max=2 → requeued. let swept = q.sweep_expired(past, 2).await.unwrap(); assert_eq!(swept.requeued, 1); assert_eq!(swept.expired, 0); assert_eq!( q.get("j1").await.unwrap().unwrap().status, compliance_core::models::werkbank::JobStatus::Queued ); // Lease #2 (attempts=2), let it expire again → now expired (>= max). q.lease("r2", Executor::Docker, &[], Duration::from_secs(10), past) .await .unwrap() .unwrap(); let later = past + chrono::Duration::seconds(60); let swept = q.sweep_expired(later, 2).await.unwrap(); assert_eq!(swept.requeued, 0); assert_eq!(swept.expired, 1); assert_eq!( q.get("j1").await.unwrap().unwrap().status, compliance_core::models::werkbank::JobStatus::Expired ); db.drop().await.ok(); }