From 633f945a1e6a9a96317733cd18a0856646280e1a Mon Sep 17 00:00:00 2001 From: Sharang Parnerkar Date: Fri, 17 Jul 2026 09:15:23 +0000 Subject: [PATCH] feat(werkbank): Mongo-backed job queue with lease + visibility timeout (WB-02) (#206) --- compliance-agent/src/database.rs | 34 +++ compliance-agent/src/lib.rs | 1 + compliance-agent/src/werkbank/mod.rs | 10 + compliance-agent/src/werkbank/queue.rs | 309 +++++++++++++++++++++++ compliance-agent/tests/werkbank_queue.rs | 258 +++++++++++++++++++ compliance-core/src/models/mod.rs | 3 +- compliance-core/src/models/werkbank.rs | 83 ++++++ 7 files changed, 697 insertions(+), 1 deletion(-) create mode 100644 compliance-agent/src/werkbank/mod.rs create mode 100644 compliance-agent/src/werkbank/queue.rs create mode 100644 compliance-agent/tests/werkbank_queue.rs diff --git a/compliance-agent/src/database.rs b/compliance-agent/src/database.rs index e10f3bf..4b50f36 100644 --- a/compliance-agent/src/database.rs +++ b/compliance-agent/src/database.rs @@ -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 { + self.inner.collection("werkbank_jobs") + } + #[allow(dead_code)] pub fn raw_collection(&self, name: &str) -> Collection { self.inner.collection(name) diff --git a/compliance-agent/src/lib.rs b/compliance-agent/src/lib.rs index 2cc5179..a28c50f 100644 --- a/compliance-agent/src/lib.rs +++ b/compliance-agent/src/lib.rs @@ -16,3 +16,4 @@ pub mod ssh; #[allow(dead_code)] pub mod trackers; pub mod webhooks; +pub mod werkbank; diff --git a/compliance-agent/src/werkbank/mod.rs b/compliance-agent/src/werkbank/mod.rs new file mode 100644 index 0000000..a183ec0 --- /dev/null +++ b/compliance-agent/src/werkbank/mod.rs @@ -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}; diff --git a/compliance-agent/src/werkbank/queue.rs b/compliance-agent/src/werkbank/queue.rs new file mode 100644 index 0000000..8c1197e --- /dev/null +++ b/compliance-agent/src/werkbank/queue.rs @@ -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, +} + +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) -> Result { + 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, + ) -> Result, 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, + ) -> Result, 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, + ) -> Result { + 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) -> Result { + 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, + max_attempts: u32, + // (kept explicit rather than a const so callers can tune retry policy) + ) -> Result { + 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, 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) -> 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" + ); + } + } +} diff --git a/compliance-agent/tests/werkbank_queue.rs b/compliance-agent/tests/werkbank_queue.rs new file mode 100644 index 0000000..a5eb4ce --- /dev/null +++ b/compliance-agent/tests/werkbank_queue.rs @@ -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.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(); +} diff --git a/compliance-core/src/models/mod.rs b/compliance-core/src/models/mod.rs index 51c0755..9ef2323 100644 --- a/compliance-core/src/models/mod.rs +++ b/compliance-core/src/models/mod.rs @@ -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, }; diff --git a/compliance-core/src/models/werkbank.rs b/compliance-core/src/models/werkbank.rs index 5969bfd..2be2771 100644 --- a/compliance-core/src/models/werkbank.rs +++ b/compliance-core/src/models/werkbank.rs @@ -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, + /// Id of the runner holding the lease. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub leased_by: Option, + /// 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>, + /// Last heartbeat from the runner. + #[serde(default, with = "super::serde_helpers::opt_bson_datetime")] + pub heartbeat_at: Option>, + /// 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, + /// When the job was enqueued. + #[serde(with = "super::serde_helpers::bson_datetime")] + pub created_at: DateTime, + /// Last modification. + #[serde(with = "super::serde_helpers::bson_datetime")] + pub updated_at: DateTime, +} + +impl JobRecord { + /// A freshly-enqueued (`queued`) record for a job. + pub fn queued(job: Job, now: DateTime) -> 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 {