CI / Check (push) Has been skipped
CI / Detect Changes (push) Successful in 3s
CI / Deploy Agent (push) Successful in 3m31s
CI / Deploy Dashboard (push) Successful in 2m38s
CI / Deploy Docs (push) Has been skipped
CI / Deploy MCP (push) Successful in 1m47s
259 lines
7.8 KiB
Rust
259 lines
7.8 KiB
Rust
//! 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> {
|
|
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();
|
|
}
|