feat(werkbank): Mongo-backed job queue with lease + visibility timeout (WB-02)
CI / Check (pull_request) Successful in 5m49s
CI / Detect Changes (pull_request) Has been skipped
CI / Deploy Agent (pull_request) Has been skipped
CI / Deploy Dashboard (pull_request) Has been skipped
CI / Deploy Docs (pull_request) Has been skipped
CI / Deploy MCP (pull_request) Has been skipped
CI / Check (pull_request) Successful in 5m49s
CI / Detect Changes (pull_request) Has been skipped
CI / Deploy Agent (pull_request) Has been skipped
CI / Deploy Dashboard (pull_request) Has been skipped
CI / Deploy Docs (pull_request) Has been skipped
CI / Deploy MCP (pull_request) Has been skipped
The control-plane pull queue behind the Werkbank runner flow (implements sharang/werkbank#3). A JobQueue over a `werkbank_jobs` collection: - enqueue — idempotent by job id (unique index; duplicate is a no-op) - lease — atomic find-and-modify of the oldest queued job the runner can run, matched by executor and by labels (job labels must be a subset of the runner's, empty/absent matches any), returns the job + a lease token, bumps attempts - heartbeat — extends the lease, flips leased→running, surfaces a cancel request; None means the lease was lost (token mismatch / already terminal) - complete — records the terminal result, token-guarded and only from an active state, so it's idempotent - cancel — queued→cancelled outright, in-flight flagged for the next heartbeat - sweep_expired — the visibility timeout: expired leases go back to queued, or to expired once attempts hit max, so a crashed runner's job recovers All transitions are single atomic Mongo updates guarded by the lease token, so two runners can never both own a job. Every op takes an explicit `now` for deterministic tests. Adds JobRecord/LeasedJob/HeartbeatAck to the contract (BSON datetimes so range queries compare correctly) and the werkbank_jobs indexes. Tests: 5 integration against a real Mongo (idempotent enqueue, executor+label matching + FIFO, heartbeat/cancel, token-guarded idempotent complete, sweep requeue→expire; skip cleanly with no Mongo) + 2 unit. clippy + fmt clean. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Fable 5
parent
7b218fffef
commit
25f232774e
@@ -465,6 +465,34 @@ impl Database {
|
||||
)
|
||||
.await?;
|
||||
|
||||
// werkbank_jobs: unique job id (idempotent enqueue by job id)
|
||||
self.werkbank_jobs()
|
||||
.create_index(
|
||||
IndexModel::builder()
|
||||
.keys(doc! { "job.id": 1 })
|
||||
.options(IndexOptions::builder().unique(true).build())
|
||||
.build(),
|
||||
)
|
||||
.await?;
|
||||
|
||||
// werkbank_jobs: lease query — oldest queued job for an executor
|
||||
self.werkbank_jobs()
|
||||
.create_index(
|
||||
IndexModel::builder()
|
||||
.keys(doc! { "status": 1, "job.executor": 1, "created_at": 1 })
|
||||
.build(),
|
||||
)
|
||||
.await?;
|
||||
|
||||
// werkbank_jobs: visibility-timeout sweep of expired leases
|
||||
self.werkbank_jobs()
|
||||
.create_index(
|
||||
IndexModel::builder()
|
||||
.keys(doc! { "status": 1, "lease_expires_at": 1 })
|
||||
.build(),
|
||||
)
|
||||
.await?;
|
||||
|
||||
tracing::info!("Database indexes ensured");
|
||||
Ok(())
|
||||
}
|
||||
@@ -563,6 +591,12 @@ impl Database {
|
||||
self.inner.collection("pentest_messages")
|
||||
}
|
||||
|
||||
/// The Werkbank job queue (WB-02): declarative dynamic-execution jobs the
|
||||
/// control plane enqueues and runners lease.
|
||||
pub fn werkbank_jobs(&self) -> Collection<compliance_core::models::werkbank::JobRecord> {
|
||||
self.inner.collection("werkbank_jobs")
|
||||
}
|
||||
|
||||
#[allow(dead_code)]
|
||||
pub fn raw_collection(&self, name: &str) -> Collection<mongodb::bson::Document> {
|
||||
self.inner.collection(name)
|
||||
|
||||
@@ -16,3 +16,4 @@ pub mod ssh;
|
||||
#[allow(dead_code)]
|
||||
pub mod trackers;
|
||||
pub mod webhooks;
|
||||
pub mod werkbank;
|
||||
|
||||
@@ -0,0 +1,10 @@
|
||||
//! Werkbank control-plane: the dynamic-execution job queue.
|
||||
//!
|
||||
//! The control plane enqueues declarative [`Job`](compliance_core::models::werkbank::Job)s
|
||||
//! and Werkbank runners lease, run, and complete them. [`queue::JobQueue`] is the
|
||||
//! Mongo-backed queue behind that flow (WB-02); the runner-facing HTTP transport
|
||||
//! and the runner itself land in later stories.
|
||||
|
||||
pub mod queue;
|
||||
|
||||
pub use queue::{JobQueue, SweepOutcome};
|
||||
@@ -0,0 +1,309 @@
|
||||
//! The Mongo-backed Werkbank job queue (WB-02).
|
||||
//!
|
||||
//! A pull queue: the control plane [`enqueue`](JobQueue::enqueue)s jobs; a runner
|
||||
//! [`lease`](JobQueue::lease)s the oldest queued job it can run (matched by
|
||||
//! executor + labels), [`heartbeat`](JobQueue::heartbeat)s while it works, and
|
||||
//! [`complete`](JobQueue::complete)s it. Leases carry a visibility timeout: if a
|
||||
//! runner dies mid-job its heartbeats stop, the lease expires, and
|
||||
//! [`sweep_expired`](JobQueue::sweep_expired) returns the job to `queued` (or
|
||||
//! `expired` once it has been retried too many times).
|
||||
//!
|
||||
//! All state transitions are single atomic Mongo updates guarded by the lease
|
||||
//! token, so two runners can never both own a job. Every operation takes an
|
||||
//! explicit `now` so the queue's time-dependent behaviour is deterministically
|
||||
//! testable.
|
||||
|
||||
use std::time::Duration;
|
||||
|
||||
use chrono::{DateTime, Utc};
|
||||
use mongodb::bson::{doc, Bson, DateTime as BsonDateTime};
|
||||
use mongodb::error::{ErrorKind, WriteFailure};
|
||||
use mongodb::options::ReturnDocument;
|
||||
use mongodb::Collection;
|
||||
|
||||
use compliance_core::models::werkbank::{
|
||||
Executor, HeartbeatAck, Job, JobRecord, JobResult, JobStatus, LeasedJob,
|
||||
};
|
||||
|
||||
use crate::database::Database;
|
||||
use crate::error::AgentError;
|
||||
|
||||
/// The non-terminal states a job can be swept or cancelled from.
|
||||
const ACTIVE_STATES: [&str; 2] = ["leased", "running"];
|
||||
/// Every terminal state (no further transitions).
|
||||
const TERMINAL_STATES: [&str; 4] = ["succeeded", "failed", "expired", "cancelled"];
|
||||
|
||||
/// What a visibility-timeout sweep did.
|
||||
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
|
||||
pub struct SweepOutcome {
|
||||
/// Expired-lease jobs returned to `queued` for another runner.
|
||||
pub requeued: u64,
|
||||
/// Jobs that had exhausted their attempts and were marked `expired`.
|
||||
pub expired: u64,
|
||||
}
|
||||
|
||||
/// The Mongo-backed job queue.
|
||||
pub struct JobQueue {
|
||||
coll: Collection<JobRecord>,
|
||||
}
|
||||
|
||||
impl JobQueue {
|
||||
/// Build a queue over a tenant database's `werkbank_jobs` collection.
|
||||
pub fn new(db: &Database) -> Self {
|
||||
Self {
|
||||
coll: db.werkbank_jobs(),
|
||||
}
|
||||
}
|
||||
|
||||
/// Enqueue a job. Idempotent by job id: a job that is already present is a
|
||||
/// no-op. Returns `true` if this call inserted it, `false` if it existed.
|
||||
pub async fn enqueue(&self, job: Job, now: DateTime<Utc>) -> Result<bool, AgentError> {
|
||||
let record = JobRecord::queued(job, now);
|
||||
match self.coll.insert_one(&record).await {
|
||||
Ok(_) => Ok(true),
|
||||
Err(e) if is_duplicate_key(&e) => Ok(false),
|
||||
Err(e) => Err(e.into()),
|
||||
}
|
||||
}
|
||||
|
||||
/// Atomically lease the oldest `queued` job this runner can run — matched by
|
||||
/// executor and by labels (every label the job requires must be one the
|
||||
/// runner advertises). Returns the job plus a lease token, or `None` if
|
||||
/// nothing is runnable.
|
||||
pub async fn lease(
|
||||
&self,
|
||||
runner_id: &str,
|
||||
executor: Executor,
|
||||
runner_labels: &[String],
|
||||
lease_ttl: Duration,
|
||||
now: DateTime<Utc>,
|
||||
) -> Result<Option<LeasedJob>, AgentError> {
|
||||
let token = uuid::Uuid::new_v4().to_string();
|
||||
let expires = bson_dt(now + ttl(lease_ttl));
|
||||
let executor_bson = mongodb::bson::to_bson(&executor).unwrap_or(Bson::Null);
|
||||
|
||||
let filter = doc! {
|
||||
"status": "queued",
|
||||
"cancel_requested": { "$ne": true },
|
||||
"job.executor": executor_bson,
|
||||
// Every label the job requires must be in the runner's set — i.e. the
|
||||
// job has no label that is not offered by the runner. Absent/empty
|
||||
// job labels match any runner.
|
||||
"job.labels": { "$not": { "$elemMatch": { "$nin": runner_labels.to_vec() } } },
|
||||
};
|
||||
let update = doc! {
|
||||
"$set": {
|
||||
"status": "leased",
|
||||
"lease_token": &token,
|
||||
"leased_by": runner_id,
|
||||
"lease_expires_at": expires,
|
||||
"heartbeat_at": bson_dt(now),
|
||||
"updated_at": bson_dt(now),
|
||||
},
|
||||
"$inc": { "attempts": 1 },
|
||||
};
|
||||
|
||||
let record = self
|
||||
.coll
|
||||
.find_one_and_update(filter, update)
|
||||
.sort(doc! { "created_at": 1 }) // FIFO
|
||||
.return_document(ReturnDocument::After)
|
||||
.await?;
|
||||
Ok(record.map(|r| LeasedJob {
|
||||
job: r.job,
|
||||
lease_token: token,
|
||||
}))
|
||||
}
|
||||
|
||||
/// Extend a lease and report whether the job has been asked to cancel.
|
||||
/// Transitions the job to `running` on the first heartbeat. Returns `None`
|
||||
/// when the lease is no longer valid (token mismatch, or the job is already
|
||||
/// terminal) — the runner should then abandon the work.
|
||||
pub async fn heartbeat(
|
||||
&self,
|
||||
job_id: &str,
|
||||
lease_token: &str,
|
||||
lease_ttl: Duration,
|
||||
now: DateTime<Utc>,
|
||||
) -> Result<Option<HeartbeatAck>, AgentError> {
|
||||
let filter = doc! {
|
||||
"job.id": job_id,
|
||||
"lease_token": lease_token,
|
||||
"status": { "$in": ACTIVE_STATES.to_vec() },
|
||||
};
|
||||
let update = doc! {
|
||||
"$set": {
|
||||
"status": "running",
|
||||
"lease_expires_at": bson_dt(now + ttl(lease_ttl)),
|
||||
"heartbeat_at": bson_dt(now),
|
||||
"updated_at": bson_dt(now),
|
||||
},
|
||||
};
|
||||
let record = self
|
||||
.coll
|
||||
.find_one_and_update(filter, update)
|
||||
.return_document(ReturnDocument::After)
|
||||
.await?;
|
||||
Ok(record.map(|r| HeartbeatAck {
|
||||
cancelled: r.cancel_requested,
|
||||
}))
|
||||
}
|
||||
|
||||
/// Record a job's terminal result. Guarded by the lease token and only from
|
||||
/// an active (`leased`/`running`) state, so it is idempotent — a duplicate or
|
||||
/// late submission after the job already finished matches nothing. Returns
|
||||
/// `true` if this call recorded the result.
|
||||
pub async fn complete(
|
||||
&self,
|
||||
job_id: &str,
|
||||
lease_token: &str,
|
||||
result: &JobResult,
|
||||
now: DateTime<Utc>,
|
||||
) -> Result<bool, AgentError> {
|
||||
let status = result.status.unwrap_or(JobStatus::Failed);
|
||||
let status_bson = mongodb::bson::to_bson(&status).unwrap_or(Bson::String("failed".into()));
|
||||
let result_bson =
|
||||
mongodb::bson::to_bson(result).map_err(|e| AgentError::Other(e.to_string()))?;
|
||||
|
||||
let filter = doc! {
|
||||
"job.id": job_id,
|
||||
"lease_token": lease_token,
|
||||
"status": { "$in": ACTIVE_STATES.to_vec() },
|
||||
};
|
||||
let update = doc! {
|
||||
"$set": {
|
||||
"status": status_bson,
|
||||
"result": result_bson,
|
||||
"lease_token": Bson::Null,
|
||||
"lease_expires_at": Bson::Null,
|
||||
"updated_at": bson_dt(now),
|
||||
},
|
||||
};
|
||||
let res = self.coll.update_one(filter, update).await?;
|
||||
Ok(res.modified_count == 1)
|
||||
}
|
||||
|
||||
/// Request cancellation of a job. A still-`queued` job is cancelled outright;
|
||||
/// an in-flight one is flagged so the runner sees it on its next heartbeat and
|
||||
/// tears down. Returns `true` if a non-terminal job matched.
|
||||
pub async fn cancel(&self, job_id: &str, now: DateTime<Utc>) -> Result<bool, AgentError> {
|
||||
let filter = doc! {
|
||||
"job.id": job_id,
|
||||
"status": { "$nin": TERMINAL_STATES.to_vec() },
|
||||
};
|
||||
// Pipeline update: flag cancellation, and if still queued flip straight to
|
||||
// cancelled (nothing is running it).
|
||||
let pipeline = vec![doc! {
|
||||
"$set": {
|
||||
"cancel_requested": true,
|
||||
"status": {
|
||||
"$cond": [ { "$eq": ["$status", "queued"] }, "cancelled", "$status" ]
|
||||
},
|
||||
"updated_at": bson_dt(now),
|
||||
}
|
||||
}];
|
||||
let res = self.coll.update_one(filter, pipeline).await?;
|
||||
Ok(res.matched_count == 1)
|
||||
}
|
||||
|
||||
/// Sweep leases whose visibility timeout has elapsed: return them to `queued`
|
||||
/// for another runner, or mark them `expired` once they have been leased
|
||||
/// `max_attempts` times. This is what makes a crashed runner's job recover.
|
||||
pub async fn sweep_expired(
|
||||
&self,
|
||||
now: DateTime<Utc>,
|
||||
max_attempts: u32,
|
||||
// (kept explicit rather than a const so callers can tune retry policy)
|
||||
) -> Result<SweepOutcome, AgentError> {
|
||||
let now_bson = bson_dt(now);
|
||||
let max = i64::from(max_attempts);
|
||||
|
||||
let requeue = self
|
||||
.coll
|
||||
.update_many(
|
||||
doc! {
|
||||
"status": { "$in": ACTIVE_STATES.to_vec() },
|
||||
"lease_expires_at": { "$lt": &now_bson },
|
||||
"attempts": { "$lt": max },
|
||||
},
|
||||
doc! { "$set": {
|
||||
"status": "queued",
|
||||
"lease_token": Bson::Null,
|
||||
"leased_by": Bson::Null,
|
||||
"lease_expires_at": Bson::Null,
|
||||
"updated_at": &now_bson,
|
||||
} },
|
||||
)
|
||||
.await?;
|
||||
|
||||
let expire = self
|
||||
.coll
|
||||
.update_many(
|
||||
doc! {
|
||||
"status": { "$in": ACTIVE_STATES.to_vec() },
|
||||
"lease_expires_at": { "$lt": &now_bson },
|
||||
"attempts": { "$gte": max },
|
||||
},
|
||||
doc! { "$set": {
|
||||
"status": "expired",
|
||||
"lease_token": Bson::Null,
|
||||
"lease_expires_at": Bson::Null,
|
||||
"updated_at": &now_bson,
|
||||
} },
|
||||
)
|
||||
.await?;
|
||||
|
||||
Ok(SweepOutcome {
|
||||
requeued: requeue.modified_count,
|
||||
expired: expire.modified_count,
|
||||
})
|
||||
}
|
||||
|
||||
/// Fetch a job record by job id (inspection / control-plane reads).
|
||||
pub async fn get(&self, job_id: &str) -> Result<Option<JobRecord>, AgentError> {
|
||||
Ok(self.coll.find_one(doc! { "job.id": job_id }).await?)
|
||||
}
|
||||
}
|
||||
|
||||
/// A `chrono::Duration` for a lease TTL, saturating rather than panicking on an
|
||||
/// absurd input (`chrono::Duration::seconds` panics past its internal bound).
|
||||
fn ttl(d: Duration) -> chrono::Duration {
|
||||
let secs = i64::try_from(d.as_secs()).unwrap_or(i64::MAX);
|
||||
chrono::Duration::try_seconds(secs).unwrap_or(chrono::Duration::MAX)
|
||||
}
|
||||
|
||||
/// A chrono instant as a BSON date (so Mongo stores/compares it as a real date).
|
||||
fn bson_dt(dt: DateTime<Utc>) -> BsonDateTime {
|
||||
BsonDateTime::from_chrono(dt)
|
||||
}
|
||||
|
||||
/// Whether a Mongo error is a duplicate-key (E11000) violation — a job with this
|
||||
/// id is already enqueued.
|
||||
fn is_duplicate_key(e: &mongodb::error::Error) -> bool {
|
||||
match &*e.kind {
|
||||
ErrorKind::Write(WriteFailure::WriteError(we)) => we.code == 11000,
|
||||
_ => false,
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn ttl_saturates_and_converts() {
|
||||
assert_eq!(ttl(Duration::from_secs(30)), chrono::Duration::seconds(30));
|
||||
// An absurd TTL saturates instead of panicking.
|
||||
assert_eq!(ttl(Duration::from_secs(u64::MAX)), chrono::Duration::MAX);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn state_constants_are_disjoint() {
|
||||
for s in ACTIVE_STATES {
|
||||
assert!(
|
||||
!TERMINAL_STATES.contains(&s),
|
||||
"{s} cannot be both active and terminal"
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,258 @@
|
||||
//! 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();
|
||||
}
|
||||
@@ -49,5 +49,6 @@ pub use repository::ScanTrigger;
|
||||
pub use sbom::{SbomEntry, VulnRef};
|
||||
pub use scan::{ScanPhase, ScanRun, ScanRunStatus, ScanType};
|
||||
pub use werkbank::{
|
||||
DastCollect, Executor, InputRef, Job, JobCollect, JobResult, JobRuntime, JobStatus, JobType,
|
||||
DastCollect, Executor, HeartbeatAck, InputRef, Job, JobCollect, JobRecord, JobResult,
|
||||
JobRuntime, JobStatus, JobType, LeasedJob,
|
||||
};
|
||||
|
||||
@@ -266,6 +266,89 @@ impl JobResult {
|
||||
}
|
||||
}
|
||||
|
||||
/// A queued job as persisted by the control plane (WB-02): the [`Job`] contract
|
||||
/// plus the queue bookkeeping — status, lease ownership, attempt count, and the
|
||||
/// eventual result. The runner never sees this record; on lease it receives a
|
||||
/// [`LeasedJob`] (the job plus a token it presents to heartbeat/complete).
|
||||
///
|
||||
/// Timestamps persist as native BSON dates so the queue's range queries (lease
|
||||
/// FIFO by `created_at`, visibility-timeout sweep by `lease_expires_at`) compare
|
||||
/// correctly.
|
||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||
pub struct JobRecord {
|
||||
/// The job to run.
|
||||
pub job: Job,
|
||||
/// Current queue state.
|
||||
pub status: JobStatus,
|
||||
/// The lease token held by the current runner (proves lease ownership).
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub lease_token: Option<String>,
|
||||
/// Id of the runner holding the lease.
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub leased_by: Option<String>,
|
||||
/// When the current lease expires — the visibility timeout after which a
|
||||
/// crashed runner's job is swept back to `queued`.
|
||||
#[serde(default, with = "super::serde_helpers::opt_bson_datetime")]
|
||||
pub lease_expires_at: Option<DateTime<Utc>>,
|
||||
/// Last heartbeat from the runner.
|
||||
#[serde(default, with = "super::serde_helpers::opt_bson_datetime")]
|
||||
pub heartbeat_at: Option<DateTime<Utc>>,
|
||||
/// How many times the job has been leased (incremented on each lease).
|
||||
#[serde(default)]
|
||||
pub attempts: u32,
|
||||
/// Set when the control plane requests cancellation; the runner sees it on
|
||||
/// its next heartbeat and aborts.
|
||||
#[serde(default)]
|
||||
pub cancel_requested: bool,
|
||||
/// The result, once the job reaches a terminal state.
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub result: Option<JobResult>,
|
||||
/// When the job was enqueued.
|
||||
#[serde(with = "super::serde_helpers::bson_datetime")]
|
||||
pub created_at: DateTime<Utc>,
|
||||
/// Last modification.
|
||||
#[serde(with = "super::serde_helpers::bson_datetime")]
|
||||
pub updated_at: DateTime<Utc>,
|
||||
}
|
||||
|
||||
impl JobRecord {
|
||||
/// A freshly-enqueued (`queued`) record for a job.
|
||||
pub fn queued(job: Job, now: DateTime<Utc>) -> Self {
|
||||
Self {
|
||||
job,
|
||||
status: JobStatus::Queued,
|
||||
lease_token: None,
|
||||
leased_by: None,
|
||||
lease_expires_at: None,
|
||||
heartbeat_at: None,
|
||||
attempts: 0,
|
||||
cancel_requested: false,
|
||||
result: None,
|
||||
created_at: now,
|
||||
updated_at: now,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// A job handed to a runner on lease: what to run plus the token the runner must
|
||||
/// present to heartbeat and complete it.
|
||||
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
|
||||
pub struct LeasedJob {
|
||||
/// The job to execute.
|
||||
pub job: Job,
|
||||
/// The lease token proving ownership (opaque to the runner).
|
||||
pub lease_token: String,
|
||||
}
|
||||
|
||||
/// The runner's view of a heartbeat: whether the control plane has asked the job
|
||||
/// to stop. `None` from the queue means the lease was lost (token mismatch or the
|
||||
/// job already terminal) and the runner should abandon the work.
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
|
||||
pub struct HeartbeatAck {
|
||||
/// The control plane requested cancellation — the runner should tear down.
|
||||
pub cancelled: bool,
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
#[allow(clippy::expect_used, clippy::unwrap_used)]
|
||||
mod tests {
|
||||
|
||||
Reference in New Issue
Block a user