Compare commits

...
Author SHA1 Message Date
Sharang ParnerkarandClaude Fable 5 25f232774e 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
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>
2026-07-17 11:07:31 +02:00
sharang 7b218fffef feat(werkbank): job/result contract in compliance-core (WB-01) (#205)
CI / Check (push) Has been skipped
CI / Deploy Dashboard (push) Successful in 2m38s
CI / Deploy Docs (push) Has been skipped
CI / Deploy MCP (push) Successful in 1m47s
CI / Detect Changes (push) Successful in 3s
CI / Deploy Agent (push) Successful in 3m51s
2026-07-17 08:54:11 +00:00
9 changed files with 1083 additions and 0 deletions
Generated
+1
View File
@@ -723,6 +723,7 @@ dependencies = [
"sha2",
"thiserror 2.0.18",
"tokio",
"toml",
"tracing",
"tracing-opentelemetry",
"tracing-subscriber",
+34
View File
@@ -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)
+1
View File
@@ -16,3 +16,4 @@ pub mod ssh;
#[allow(dead_code)]
pub mod trackers;
pub mod webhooks;
pub mod werkbank;
+10
View File
@@ -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};
+309
View File
@@ -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"
);
}
}
}
+258
View File
@@ -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();
}
+4
View File
@@ -50,3 +50,7 @@ axum = { version = "0.8", optional = true }
jsonwebtoken = { version = "9", optional = true }
reqwest = { workspace = true, optional = true }
tokio = { workspace = true, optional = true }
[dev-dependencies]
# Parse the declarative TOML job specs in the Werkbank contract tests.
toml = "0.8"
+5
View File
@@ -15,6 +15,7 @@ pub mod repository;
pub mod sbom;
pub mod scan;
pub(crate) mod serde_helpers;
pub mod werkbank;
pub use auth::AuthInfo;
pub use chat::{ChatMessage, ChatRequest, ChatResponse, SourceReference};
@@ -47,3 +48,7 @@ pub use pentest::{
pub use repository::ScanTrigger;
pub use sbom::{SbomEntry, VulnRef};
pub use scan::{ScanPhase, ScanRun, ScanRunStatus, ScanType};
pub use werkbank::{
DastCollect, Executor, HeartbeatAck, InputRef, Job, JobCollect, JobRecord, JobResult,
JobRuntime, JobStatus, JobType, LeasedJob,
};
+461
View File
@@ -0,0 +1,461 @@
//! The Werkbank job/result contract (WB-01).
//!
//! The shared, dependency-free vocabulary the control plane and the Werkbank
//! execution runner agree on: what a [`Job`] is, which [`Executor`] runs it, how
//! it moves through the queue ([`JobStatus`]), and what a [`JobResult`] carries
//! back. Jobs are declarative — TOML on disk, JSON on the wire — and results
//! reuse the existing scanner result types ([`Finding`], [`DastFinding`],
//! [`SbomEntry`]) so the runner produces exactly what the control plane persists.
//!
//! This module is intentionally free of the `mongodb`/`axum` features so the
//! runner can depend on `compliance-core` without pulling the server stack.
use std::collections::BTreeMap;
use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use super::dast::DastFinding;
use super::finding::Finding;
use super::sbom::SbomEntry;
/// The kind of dynamic-execution job.
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "kebab-case")]
pub enum JobType {
/// Instantiate control logic on an ephemeral soft-PLC and probe it.
PlcProvision,
/// Boot a firmware image under QEMU and run dynamic checks.
QemuBoot,
/// Crawl and dynamically test a running web endpoint.
Dast,
/// Run an active penetration test against a running target.
Pentest,
}
/// How a runner executes a job — the CI-runner-style classification. A runner
/// advertises exactly one; a job requires one.
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "lowercase")]
pub enum Executor {
/// A subprocess on the runner host (dev / trusted single-node).
Shell,
/// One or more containers on the runner's Docker (default; QEMU runs here).
Docker,
/// A Pod/Job in a Kubernetes cluster (scale-out / multi-tenant).
K8s,
}
/// Lifecycle state of a job in the queue.
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "lowercase")]
pub enum JobStatus {
/// Waiting to be leased.
Queued,
/// Leased by a runner but not yet started.
Leased,
/// Executing on a runner.
Running,
/// Completed successfully.
Succeeded,
/// Completed with an error.
Failed,
/// The lease/lifetime deadline elapsed before completion.
Expired,
/// Cancelled by the control plane.
Cancelled,
}
impl JobStatus {
/// Whether the job has reached a terminal state (no further transitions).
pub fn is_terminal(self) -> bool {
matches!(
self,
JobStatus::Succeeded | JobStatus::Failed | JobStatus::Expired | JobStatus::Cancelled
)
}
}
/// A reference to an input artifact. Resolved by the runner from a source it can
/// reach; the blob itself never flows through the control plane (so an on-prem
/// runner keeps customer data local). Exactly one of `blob`/`url` should be set.
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct InputRef {
/// Content-addressed blob (e.g. `sha256:…`) the runner fetches from its store.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub blob: Option<String>,
/// A URL the runner can reach (git repo, internal artifact store, …).
#[serde(default, skip_serializing_if = "Option::is_none")]
pub url: Option<String>,
}
impl InputRef {
/// A content-addressed blob reference.
pub fn blob(id: impl Into<String>) -> Self {
Self {
blob: Some(id.into()),
url: None,
}
}
}
/// Sandbox runtime knobs. Fields are executor/job-type specific and all optional;
/// `extra` carries anything not modelled explicitly.
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct JobRuntime {
/// Container image (Docker executor).
#[serde(default, skip_serializing_if = "Option::is_none")]
pub image: Option<String>,
/// Memory cap (e.g. `512m`).
#[serde(default, skip_serializing_if = "Option::is_none")]
pub memory: Option<String>,
/// CPU cap (e.g. `0.5`).
#[serde(default, skip_serializing_if = "Option::is_none")]
pub cpus: Option<String>,
/// Network to join (e.g. `isolated`).
#[serde(default, skip_serializing_if = "Option::is_none")]
pub network: Option<String>,
/// QEMU machine type (qemu-boot).
#[serde(default, skip_serializing_if = "Option::is_none")]
pub machine: Option<String>,
/// QEMU target architecture (qemu-boot).
#[serde(default, skip_serializing_if = "Option::is_none")]
pub arch: Option<String>,
/// Executor-specific extras not modelled above.
#[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
pub extra: BTreeMap<String, String>,
}
/// DAST collection settings for jobs that scan a web endpoint.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct DastCollect {
/// Maximum crawl depth (kept shallow for ephemeral instances).
pub max_crawl_depth: u32,
}
/// What to collect from a run.
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct JobCollect {
/// Run the industrial-protocol probe (Modbus/OPC-UA/EtherNet-IP).
#[serde(default)]
pub ics_probe: bool,
/// Run DAST against the provisioned/booted web endpoint.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub dast: Option<DastCollect>,
/// Run an active pentest.
#[serde(default)]
pub pentest: bool,
/// Collect an SBOM.
#[serde(default)]
pub sbom: bool,
}
/// A declarative dynamic-execution job the control plane enqueues and a Werkbank
/// runner leases and executes.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct Job {
/// Unique job id (assigned by the control plane on enqueue).
pub id: String,
/// What kind of job this is.
#[serde(rename = "type")]
pub job_type: JobType,
/// Owning tenant.
pub tenant: String,
/// The onboarded target this job tests.
pub target_id: String,
/// The executor a runner must provide to run this job.
pub executor: Executor,
/// Runner capabilities this job requires (e.g. `arch=amd64`, `kvm=true`).
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub labels: Vec<String>,
/// Hard lifetime deadline for the whole job.
pub timeout_secs: u64,
/// Named input artifacts (e.g. `program`, `firmware`), by reference.
#[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
pub inputs: BTreeMap<String, InputRef>,
/// Sandbox runtime knobs.
#[serde(default)]
pub runtime: JobRuntime,
/// What to collect from the run.
#[serde(default)]
pub collect: JobCollect,
}
impl Job {
/// A `plc-provision` job: instantiate the control logic named `program` on an
/// ephemeral soft-PLC (Docker executor) and collect the ICS probe + DAST.
pub fn plc_provision(
id: impl Into<String>,
tenant: impl Into<String>,
target_id: impl Into<String>,
program: InputRef,
timeout_secs: u64,
) -> Self {
let mut inputs = BTreeMap::new();
inputs.insert("program".to_string(), program);
Self {
id: id.into(),
job_type: JobType::PlcProvision,
tenant: tenant.into(),
target_id: target_id.into(),
executor: Executor::Docker,
labels: Vec::new(),
timeout_secs,
inputs,
runtime: JobRuntime::default(),
collect: JobCollect {
ics_probe: true,
dast: Some(DastCollect { max_crawl_depth: 2 }),
pentest: false,
sbom: false,
},
}
}
}
/// The outcome of running a job, posted back to the control plane. Findings and
/// SBOM reuse the shared scanner types, so the control plane persists them
/// unchanged. Submission is idempotent — keyed by [`JobResult::job_id`].
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct JobResult {
/// The job this result is for.
pub job_id: String,
/// Terminal status of the job.
pub status: Option<JobStatus>,
/// General scanner findings (e.g. ICS-probe findings).
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub findings: Vec<Finding>,
/// DAST findings from a web-endpoint scan.
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub dast_findings: Vec<DastFinding>,
/// SBOM components collected from the run.
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub sbom: Vec<SbomEntry>,
/// Error message when the job failed.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub error: Option<String>,
/// Captured execution log (truncated by the runner).
#[serde(default, skip_serializing_if = "Option::is_none")]
pub logs: Option<String>,
/// When execution started on the runner.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub started_at: Option<DateTime<Utc>>,
/// When execution finished.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub finished_at: Option<DateTime<Utc>>,
}
impl JobResult {
/// A successful result for a job.
pub fn succeeded(job_id: impl Into<String>) -> Self {
Self {
job_id: job_id.into(),
status: Some(JobStatus::Succeeded),
..Default::default()
}
}
/// A failed result carrying an error message.
pub fn failed(job_id: impl Into<String>, error: impl Into<String>) -> Self {
Self {
job_id: job_id.into(),
status: Some(JobStatus::Failed),
error: Some(error.into()),
..Default::default()
}
}
}
/// 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 {
use super::*;
#[test]
fn job_round_trips_through_json() {
let job = Job::plc_provision("job_1", "acme", "64f0aa", InputRef::blob("sha256:abc"), 180);
let json = serde_json::to_string(&job).expect("serialize");
let back: Job = serde_json::from_str(&json).expect("deserialize");
assert_eq!(job, back);
// Enum wire forms are the kebab/lowercase the contract documents.
assert!(json.contains("\"type\":\"plc-provision\""));
assert!(json.contains("\"executor\":\"docker\""));
}
#[test]
fn parses_the_design_doc_plc_provision_toml() {
// The exact shape from docs/DESIGN.md §5 (wrapped in a [job] table).
#[derive(Deserialize)]
struct JobFile {
job: Job,
}
let src = r#"
[job]
id = "job_01H"
type = "plc-provision"
tenant = "acme"
target_id = "64f0"
executor = "docker"
labels = ["arch=amd64"]
timeout_secs = 180
[job.inputs]
program = { blob = "sha256:deadbeef" }
[job.runtime]
image = "openplc:latest"
memory = "512m"
cpus = "0.5"
network = "isolated"
[job.collect]
ics_probe = true
dast = { max_crawl_depth = 2 }
"#;
let file: JobFile = toml::from_str(src).expect("parse job toml");
let job = file.job;
assert_eq!(job.job_type, JobType::PlcProvision);
assert_eq!(job.executor, Executor::Docker);
assert_eq!(job.labels, vec!["arch=amd64".to_string()]);
assert_eq!(
job.inputs.get("program").and_then(|i| i.blob.as_deref()),
Some("sha256:deadbeef")
);
assert_eq!(job.runtime.image.as_deref(), Some("openplc:latest"));
assert!(job.collect.ics_probe);
assert_eq!(job.collect.dast.map(|d| d.max_crawl_depth), Some(2));
}
#[test]
fn qemu_boot_runtime_fields_parse() {
#[derive(Deserialize)]
struct JobFile {
job: Job,
}
let src = r#"
[job]
id = "j2"
type = "qemu-boot"
tenant = "acme"
target_id = "t"
executor = "docker"
labels = ["kvm=true"]
timeout_secs = 600
[job.inputs]
firmware = { blob = "sha256:cafe" }
[job.runtime]
machine = "virt"
arch = "arm"
memory = "1g"
"#;
let file: JobFile = toml::from_str(src).expect("parse");
assert_eq!(file.job.job_type, JobType::QemuBoot);
assert_eq!(file.job.runtime.arch.as_deref(), Some("arm"));
assert_eq!(
file.job
.inputs
.get("firmware")
.and_then(|i| i.blob.as_deref()),
Some("sha256:cafe")
);
}
#[test]
fn status_terminality() {
assert!(JobStatus::Succeeded.is_terminal());
assert!(JobStatus::Expired.is_terminal());
assert!(!JobStatus::Queued.is_terminal());
assert!(!JobStatus::Running.is_terminal());
}
#[test]
fn result_constructors() {
assert_eq!(JobResult::succeeded("j").status, Some(JobStatus::Succeeded));
let f = JobResult::failed("j", "boom");
assert_eq!(f.status, Some(JobStatus::Failed));
assert_eq!(f.error.as_deref(), Some("boom"));
}
}