Compare commits
1
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
463d4f6eb2 |
@@ -14,7 +14,6 @@ pub mod pentest_handlers;
|
|||||||
pub use pentest_handlers as pentest;
|
pub use pentest_handlers as pentest;
|
||||||
pub mod sbom;
|
pub mod sbom;
|
||||||
pub mod scans;
|
pub mod scans;
|
||||||
pub mod werkbank_jobs;
|
|
||||||
|
|
||||||
// Re-export all handler functions so routes.rs can use `handlers::function_name`
|
// Re-export all handler functions so routes.rs can use `handlers::function_name`
|
||||||
pub use dto::*;
|
pub use dto::*;
|
||||||
|
|||||||
@@ -1,193 +0,0 @@
|
|||||||
//! Werkbank runner endpoints (`/api/v1/werkbank/jobs/*`).
|
|
||||||
//!
|
|
||||||
//! The pull API a Werkbank runner talks to: lease a job, heartbeat while it runs,
|
|
||||||
//! and post the result back. Machine auth is a **static bearer token**
|
|
||||||
//! (`WERKBANK_RUNNER_TOKEN`) — not a Keycloak JWT, because a runner acts across
|
|
||||||
//! tenants (each request names its `tenant`). Routes are only mounted when the
|
|
||||||
//! token is configured; with none set they don't exist (404).
|
|
||||||
//!
|
|
||||||
//! On completion the runner's findings are persisted against the job's target,
|
|
||||||
//! so a job run by a remote runner lands the same findings an in-process run
|
|
||||||
//! would (WB-05, the control-plane cut-over).
|
|
||||||
|
|
||||||
use axum::extract::{Extension, Request};
|
|
||||||
use axum::http::{header, StatusCode};
|
|
||||||
use axum::middleware::Next;
|
|
||||||
use axum::response::{IntoResponse, Response};
|
|
||||||
use axum::Json;
|
|
||||||
use mongodb::bson::doc;
|
|
||||||
use secrecy::ExposeSecret;
|
|
||||||
use std::time::Duration;
|
|
||||||
|
|
||||||
use compliance_core::models::werkbank::{
|
|
||||||
CompleteRequest, CompleteResponse, HeartbeatRequest, JobResult, LeaseRequest,
|
|
||||||
};
|
|
||||||
|
|
||||||
use super::dto::AgentExt;
|
|
||||||
use crate::database::Database;
|
|
||||||
use crate::werkbank::JobQueue;
|
|
||||||
|
|
||||||
/// Gate the runner endpoints behind the static runner bearer token.
|
|
||||||
pub async fn require_runner_token(
|
|
||||||
Extension(agent): AgentExt,
|
|
||||||
request: Request,
|
|
||||||
next: Next,
|
|
||||||
) -> Response {
|
|
||||||
let Some(expected) = agent.config.werkbank_runner_token.as_ref() else {
|
|
||||||
return (StatusCode::NOT_FOUND, "werkbank runner API disabled").into_response();
|
|
||||||
};
|
|
||||||
let presented = request
|
|
||||||
.headers()
|
|
||||||
.get(header::AUTHORIZATION)
|
|
||||||
.and_then(|v| v.to_str().ok())
|
|
||||||
.and_then(|s| s.strip_prefix("Bearer "))
|
|
||||||
.map(str::trim)
|
|
||||||
.filter(|s| !s.is_empty());
|
|
||||||
let Some(presented) = presented else {
|
|
||||||
return (StatusCode::UNAUTHORIZED, "Missing bearer token").into_response();
|
|
||||||
};
|
|
||||||
if !constant_time_eq(presented, expected.expose_secret()) {
|
|
||||||
return (StatusCode::UNAUTHORIZED, "Invalid runner token").into_response();
|
|
||||||
}
|
|
||||||
next.run(request).await
|
|
||||||
}
|
|
||||||
|
|
||||||
/// `POST /api/v1/werkbank/jobs/lease` — lease the oldest runnable job, or `204`.
|
|
||||||
#[tracing::instrument(skip_all, fields(tenant = %req.tenant, runner = %req.runner_id))]
|
|
||||||
pub async fn lease(
|
|
||||||
Extension(agent): AgentExt,
|
|
||||||
Json(req): Json<LeaseRequest>,
|
|
||||||
) -> Result<Response, StatusCode> {
|
|
||||||
let queue = JobQueue::new(&tenant_db(&agent, &req.tenant).await?);
|
|
||||||
let leased = queue
|
|
||||||
.lease(
|
|
||||||
&req.runner_id,
|
|
||||||
req.executor,
|
|
||||||
&req.labels,
|
|
||||||
Duration::from_secs(req.lease_ttl_secs),
|
|
||||||
chrono::Utc::now(),
|
|
||||||
)
|
|
||||||
.await
|
|
||||||
.map_err(internal)?;
|
|
||||||
Ok(match leased {
|
|
||||||
Some(job) => Json(job).into_response(),
|
|
||||||
None => StatusCode::NO_CONTENT.into_response(),
|
|
||||||
})
|
|
||||||
}
|
|
||||||
|
|
||||||
/// `POST /api/v1/werkbank/jobs/heartbeat` — extend the lease; `409` if it's lost.
|
|
||||||
#[tracing::instrument(skip_all, fields(tenant = %req.tenant, job = %req.job_id))]
|
|
||||||
pub async fn heartbeat(
|
|
||||||
Extension(agent): AgentExt,
|
|
||||||
Json(req): Json<HeartbeatRequest>,
|
|
||||||
) -> Result<Response, StatusCode> {
|
|
||||||
let queue = JobQueue::new(&tenant_db(&agent, &req.tenant).await?);
|
|
||||||
let ack = queue
|
|
||||||
.heartbeat(
|
|
||||||
&req.job_id,
|
|
||||||
&req.lease_token,
|
|
||||||
Duration::from_secs(req.lease_ttl_secs),
|
|
||||||
chrono::Utc::now(),
|
|
||||||
)
|
|
||||||
.await
|
|
||||||
.map_err(internal)?;
|
|
||||||
Ok(match ack {
|
|
||||||
Some(ack) => Json(ack).into_response(),
|
|
||||||
// Lease lost — the runner should abandon the job.
|
|
||||||
None => StatusCode::CONFLICT.into_response(),
|
|
||||||
})
|
|
||||||
}
|
|
||||||
|
|
||||||
/// `POST /api/v1/werkbank/jobs/complete` — record the result and persist findings.
|
|
||||||
#[tracing::instrument(skip_all, fields(tenant = %req.tenant, job = %req.job_id))]
|
|
||||||
pub async fn complete(
|
|
||||||
Extension(agent): AgentExt,
|
|
||||||
Json(req): Json<CompleteRequest>,
|
|
||||||
) -> Result<Json<CompleteResponse>, StatusCode> {
|
|
||||||
let db = tenant_db(&agent, &req.tenant).await?;
|
|
||||||
let queue = JobQueue::new(&db);
|
|
||||||
let now = chrono::Utc::now();
|
|
||||||
let recorded = queue
|
|
||||||
.complete(&req.job_id, &req.lease_token, &req.result, now)
|
|
||||||
.await
|
|
||||||
.map_err(internal)?;
|
|
||||||
|
|
||||||
// Only persist findings for the run that actually recorded the result, so a
|
|
||||||
// duplicate/late completion can't double-insert.
|
|
||||||
if recorded {
|
|
||||||
if let Some(record) = queue.get(&req.job_id).await.map_err(internal)? {
|
|
||||||
persist_findings(&db, &record.job.target_id, &req.result).await;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
Ok(Json(CompleteResponse { recorded }))
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Persist a job result's findings against its target: general findings
|
|
||||||
/// (dedup'd by fingerprint) and DAST findings. Best-effort — a persistence hiccup
|
|
||||||
/// is logged, not surfaced to the runner (its result is already recorded).
|
|
||||||
async fn persist_findings(db: &Database, target_id: &str, result: &JobResult) {
|
|
||||||
for finding in &result.findings {
|
|
||||||
let exists = db
|
|
||||||
.findings()
|
|
||||||
.find_one(doc! { "fingerprint": &finding.fingerprint })
|
|
||||||
.await
|
|
||||||
.ok()
|
|
||||||
.flatten()
|
|
||||||
.is_some();
|
|
||||||
if !exists {
|
|
||||||
if let Err(e) = db.findings().insert_one(finding).await {
|
|
||||||
tracing::warn!(target_id, error = %e, "werkbank: persist finding failed");
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
for finding in &result.dast_findings {
|
|
||||||
if let Err(e) = db.dast_findings().insert_one(finding).await {
|
|
||||||
tracing::warn!(target_id, error = %e, "werkbank: persist DAST finding failed");
|
|
||||||
}
|
|
||||||
}
|
|
||||||
tracing::info!(
|
|
||||||
target_id,
|
|
||||||
findings = result.findings.len(),
|
|
||||||
dast = result.dast_findings.len(),
|
|
||||||
"werkbank: persisted runner results"
|
|
||||||
);
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Resolve the tenant-scoped database for a request.
|
|
||||||
async fn tenant_db(
|
|
||||||
agent: &crate::agent::ComplianceAgent,
|
|
||||||
tenant: &str,
|
|
||||||
) -> Result<Database, StatusCode> {
|
|
||||||
agent.db_pool.for_tenant_id(tenant).await.map_err(internal)
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Map any internal error to a 500.
|
|
||||||
fn internal<E: std::fmt::Display>(e: E) -> StatusCode {
|
|
||||||
tracing::error!("werkbank endpoint error: {e}");
|
|
||||||
StatusCode::INTERNAL_SERVER_ERROR
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Length-checked, constant-time-ish token comparison.
|
|
||||||
fn constant_time_eq(a: &str, b: &str) -> bool {
|
|
||||||
if a.len() != b.len() {
|
|
||||||
return false;
|
|
||||||
}
|
|
||||||
let mut diff = 0u8;
|
|
||||||
for (x, y) in a.bytes().zip(b.bytes()) {
|
|
||||||
diff |= x ^ y;
|
|
||||||
}
|
|
||||||
diff == 0
|
|
||||||
}
|
|
||||||
|
|
||||||
#[cfg(test)]
|
|
||||||
mod tests {
|
|
||||||
use super::constant_time_eq;
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn token_compare() {
|
|
||||||
assert!(constant_time_eq("secret", "secret"));
|
|
||||||
assert!(!constant_time_eq("secret", "secrex"));
|
|
||||||
assert!(!constant_time_eq("secret", "secretx"));
|
|
||||||
assert!(!constant_time_eq("", "x"));
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -4,7 +4,7 @@ use axum::extract::{DefaultBodyLimit, Request};
|
|||||||
use axum::http::HeaderValue;
|
use axum::http::HeaderValue;
|
||||||
use axum::middleware::Next;
|
use axum::middleware::Next;
|
||||||
use axum::response::Response;
|
use axum::response::Response;
|
||||||
use axum::routing::{delete, get, post};
|
use axum::routing::{delete, get};
|
||||||
use axum::{middleware, Extension, Router};
|
use axum::{middleware, Extension, Router};
|
||||||
use tokio::sync::RwLock;
|
use tokio::sync::RwLock;
|
||||||
use tower_http::cors::CorsLayer;
|
use tower_http::cors::CorsLayer;
|
||||||
@@ -72,35 +72,8 @@ pub async fn start_api_server(agent: ComplianceAgent, port: u16) -> Result<(), A
|
|||||||
Router::new()
|
Router::new()
|
||||||
};
|
};
|
||||||
|
|
||||||
// Werkbank runner API. Like admin, only mounted when its bearer token is
|
|
||||||
// configured; runners authenticate with WERKBANK_RUNNER_TOKEN (not a JWT).
|
|
||||||
let werkbank_router: Router = if agent.config.werkbank_runner_token.is_some() {
|
|
||||||
tracing::info!(
|
|
||||||
"Werkbank runner API enabled — /api/v1/werkbank/jobs/* behind WERKBANK_RUNNER_TOKEN"
|
|
||||||
);
|
|
||||||
Router::new()
|
|
||||||
.route(
|
|
||||||
"/api/v1/werkbank/jobs/lease",
|
|
||||||
post(handlers::werkbank_jobs::lease),
|
|
||||||
)
|
|
||||||
.route(
|
|
||||||
"/api/v1/werkbank/jobs/heartbeat",
|
|
||||||
post(handlers::werkbank_jobs::heartbeat),
|
|
||||||
)
|
|
||||||
.route(
|
|
||||||
"/api/v1/werkbank/jobs/complete",
|
|
||||||
post(handlers::werkbank_jobs::complete),
|
|
||||||
)
|
|
||||||
.layer(middleware::from_fn(
|
|
||||||
handlers::werkbank_jobs::require_runner_token,
|
|
||||||
))
|
|
||||||
} else {
|
|
||||||
Router::new()
|
|
||||||
};
|
|
||||||
|
|
||||||
let mut app = routes::build_router()
|
let mut app = routes::build_router()
|
||||||
.merge(admin_router)
|
.merge(admin_router)
|
||||||
.merge(werkbank_router)
|
|
||||||
// Allow large artifact uploads (PLC .projectarchive, firmware images,
|
// Allow large artifact uploads (PLC .projectarchive, firmware images,
|
||||||
// mobile packages) — axum's default request-body limit is only 2 MiB.
|
// mobile packages) — axum's default request-body limit is only 2 MiB.
|
||||||
.layer(DefaultBodyLimit::max(512 * 1024 * 1024))
|
.layer(DefaultBodyLimit::max(512 * 1024 * 1024))
|
||||||
|
|||||||
@@ -65,7 +65,6 @@ pub fn load_config() -> Result<AgentConfig, AgentError> {
|
|||||||
admin_api_token: env_secret_opt("ADMIN_API_TOKEN"),
|
admin_api_token: env_secret_opt("ADMIN_API_TOKEN"),
|
||||||
tenant_registry_url: env_var_opt("TENANT_REGISTRY_URL"),
|
tenant_registry_url: env_var_opt("TENANT_REGISTRY_URL"),
|
||||||
plc_runtime: load_plc_runtime_config(),
|
plc_runtime: load_plc_runtime_config(),
|
||||||
werkbank_runner_token: env_secret_opt("WERKBANK_RUNNER_TOKEN"),
|
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -465,34 +465,6 @@ impl Database {
|
|||||||
)
|
)
|
||||||
.await?;
|
.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");
|
tracing::info!("Database indexes ensured");
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
@@ -591,12 +563,6 @@ impl Database {
|
|||||||
self.inner.collection("pentest_messages")
|
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)]
|
#[allow(dead_code)]
|
||||||
pub fn raw_collection(&self, name: &str) -> Collection<mongodb::bson::Document> {
|
pub fn raw_collection(&self, name: &str) -> Collection<mongodb::bson::Document> {
|
||||||
self.inner.collection(name)
|
self.inner.collection(name)
|
||||||
|
|||||||
@@ -16,4 +16,3 @@ pub mod ssh;
|
|||||||
#[allow(dead_code)]
|
#[allow(dead_code)]
|
||||||
pub mod trackers;
|
pub mod trackers;
|
||||||
pub mod webhooks;
|
pub mod webhooks;
|
||||||
pub mod werkbank;
|
|
||||||
|
|||||||
@@ -343,7 +343,6 @@ mod tests {
|
|||||||
admin_api_token: None,
|
admin_api_token: None,
|
||||||
tenant_registry_url: None,
|
tenant_registry_url: None,
|
||||||
plc_runtime: compliance_core::PlcRuntimeConfig::default(),
|
plc_runtime: compliance_core::PlcRuntimeConfig::default(),
|
||||||
werkbank_runner_token: None,
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -1,10 +0,0 @@
|
|||||||
//! 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};
|
|
||||||
@@ -1,309 +0,0 @@
|
|||||||
//! 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"
|
|
||||||
);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -2,10 +2,6 @@
|
|||||||
//
|
//
|
||||||
// Spins up the agent API server on a random port with an isolated test
|
// Spins up the agent API server on a random port with an isolated test
|
||||||
// database. Each test gets a fresh database that is dropped on cleanup.
|
// database. Each test gets a fresh database that is dropped on cleanup.
|
||||||
//
|
|
||||||
// Included via `mod common;` in several test binaries; not every binary uses
|
|
||||||
// every helper, so allow dead code here.
|
|
||||||
#![allow(dead_code)]
|
|
||||||
|
|
||||||
use std::sync::Arc;
|
use std::sync::Arc;
|
||||||
|
|
||||||
@@ -15,54 +11,6 @@ use compliance_agent::database::DatabasePool;
|
|||||||
use compliance_core::AgentConfig;
|
use compliance_core::AgentConfig;
|
||||||
use secrecy::SecretString;
|
use secrecy::SecretString;
|
||||||
|
|
||||||
/// The runner bearer token wired into the test config.
|
|
||||||
pub const TEST_RUNNER_TOKEN: &str = "test-runner-token";
|
|
||||||
|
|
||||||
/// A minimal dev [`AgentConfig`] for tests: unauthenticated (no Keycloak), the
|
|
||||||
/// Werkbank runner API enabled with [`TEST_RUNNER_TOKEN`].
|
|
||||||
pub fn dev_config(mongodb_uri: String, db_name: String) -> AgentConfig {
|
|
||||||
AgentConfig {
|
|
||||||
mongodb_uri,
|
|
||||||
mongodb_database: db_name,
|
|
||||||
litellm_url: std::env::var("TEST_LITELLM_URL")
|
|
||||||
.unwrap_or_else(|_| "http://localhost:4000".into()),
|
|
||||||
litellm_api_key: SecretString::from(String::new()),
|
|
||||||
litellm_model: "gpt-4o".into(),
|
|
||||||
litellm_embed_model: "text-embedding-3-small".into(),
|
|
||||||
agent_port: 0, // not used — we bind ourselves
|
|
||||||
scan_schedule: String::new(),
|
|
||||||
cve_monitor_schedule: String::new(),
|
|
||||||
git_clone_base_path: "/tmp/compliance-scanner-tests/repos".into(),
|
|
||||||
artifact_store_base_path: "/tmp/compliance-scanner-tests/artifacts".into(),
|
|
||||||
ssh_key_path: "/tmp/compliance-scanner-tests/ssh/id_ed25519".into(),
|
|
||||||
github_token: None,
|
|
||||||
github_webhook_secret: None,
|
|
||||||
gitlab_url: None,
|
|
||||||
gitlab_token: None,
|
|
||||||
gitlab_webhook_secret: None,
|
|
||||||
jira_url: None,
|
|
||||||
jira_email: None,
|
|
||||||
jira_api_token: None,
|
|
||||||
jira_project_key: None,
|
|
||||||
searxng_url: None,
|
|
||||||
nvd_api_key: None,
|
|
||||||
keycloak_url: None,
|
|
||||||
keycloak_realm: None,
|
|
||||||
keycloak_admin_username: None,
|
|
||||||
keycloak_admin_password: None,
|
|
||||||
pentest_verification_email: None,
|
|
||||||
pentest_imap_host: None,
|
|
||||||
pentest_imap_port: None,
|
|
||||||
pentest_imap_tls: false,
|
|
||||||
pentest_imap_username: None,
|
|
||||||
pentest_imap_password: None,
|
|
||||||
admin_api_token: None,
|
|
||||||
tenant_registry_url: None,
|
|
||||||
plc_runtime: compliance_core::PlcRuntimeConfig::default(),
|
|
||||||
werkbank_runner_token: Some(SecretString::from(TEST_RUNNER_TOKEN.to_string())),
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/// A running test server with a unique database.
|
/// A running test server with a unique database.
|
||||||
pub struct TestServer {
|
pub struct TestServer {
|
||||||
pub base_url: String,
|
pub base_url: String,
|
||||||
@@ -85,7 +33,45 @@ impl TestServer {
|
|||||||
.await
|
.await
|
||||||
.expect("Failed to build DatabasePool");
|
.expect("Failed to build DatabasePool");
|
||||||
|
|
||||||
let config = dev_config(mongodb_uri.clone(), db_name.clone());
|
let config = AgentConfig {
|
||||||
|
mongodb_uri: mongodb_uri.clone(),
|
||||||
|
mongodb_database: db_name.clone(),
|
||||||
|
litellm_url: std::env::var("TEST_LITELLM_URL")
|
||||||
|
.unwrap_or_else(|_| "http://localhost:4000".into()),
|
||||||
|
litellm_api_key: SecretString::from(String::new()),
|
||||||
|
litellm_model: "gpt-4o".into(),
|
||||||
|
litellm_embed_model: "text-embedding-3-small".into(),
|
||||||
|
agent_port: 0, // not used — we bind ourselves
|
||||||
|
scan_schedule: String::new(),
|
||||||
|
cve_monitor_schedule: String::new(),
|
||||||
|
git_clone_base_path: "/tmp/compliance-scanner-tests/repos".into(),
|
||||||
|
artifact_store_base_path: "/tmp/compliance-scanner-tests/artifacts".into(),
|
||||||
|
ssh_key_path: "/tmp/compliance-scanner-tests/ssh/id_ed25519".into(),
|
||||||
|
github_token: None,
|
||||||
|
github_webhook_secret: None,
|
||||||
|
gitlab_url: None,
|
||||||
|
gitlab_token: None,
|
||||||
|
gitlab_webhook_secret: None,
|
||||||
|
jira_url: None,
|
||||||
|
jira_email: None,
|
||||||
|
jira_api_token: None,
|
||||||
|
jira_project_key: None,
|
||||||
|
searxng_url: None,
|
||||||
|
nvd_api_key: None,
|
||||||
|
keycloak_url: None,
|
||||||
|
keycloak_realm: None,
|
||||||
|
keycloak_admin_username: None,
|
||||||
|
keycloak_admin_password: None,
|
||||||
|
pentest_verification_email: None,
|
||||||
|
pentest_imap_host: None,
|
||||||
|
pentest_imap_port: None,
|
||||||
|
pentest_imap_tls: false,
|
||||||
|
pentest_imap_username: None,
|
||||||
|
pentest_imap_password: None,
|
||||||
|
admin_api_token: None,
|
||||||
|
tenant_registry_url: None,
|
||||||
|
plc_runtime: compliance_core::PlcRuntimeConfig::default(),
|
||||||
|
};
|
||||||
|
|
||||||
let agent = ComplianceAgent::new(config, db_pool);
|
let agent = ComplianceAgent::new(config, db_pool);
|
||||||
|
|
||||||
|
|||||||
@@ -1,222 +0,0 @@
|
|||||||
//! Integration tests for the Werkbank runner endpoints (WB-05).
|
|
||||||
//!
|
|
||||||
//! Drives the real HTTP handlers (lease/heartbeat/complete) against a live Mongo:
|
|
||||||
//! a runner leases a seeded job, completes it, and the result's findings are
|
|
||||||
//! persisted against the job's target. Also checks the bearer-token gate. Skips
|
|
||||||
//! cleanly when no Mongo is reachable.
|
|
||||||
|
|
||||||
#![allow(clippy::expect_used, clippy::unwrap_used)]
|
|
||||||
|
|
||||||
mod common;
|
|
||||||
|
|
||||||
use std::sync::Arc;
|
|
||||||
|
|
||||||
use axum::routing::post;
|
|
||||||
use axum::{middleware, Extension, Router};
|
|
||||||
|
|
||||||
use compliance_agent::agent::ComplianceAgent;
|
|
||||||
use compliance_agent::api::handlers::werkbank_jobs;
|
|
||||||
use compliance_agent::database::DatabasePool;
|
|
||||||
use compliance_agent::werkbank::JobQueue;
|
|
||||||
use compliance_core::models::werkbank::{InputRef, Job, JobResult, JobStatus, LeasedJob};
|
|
||||||
use compliance_core::models::{Finding, ScanType, Severity};
|
|
||||||
|
|
||||||
use common::{dev_config, TEST_RUNNER_TOKEN};
|
|
||||||
|
|
||||||
const TENANT: &str = "dev";
|
|
||||||
|
|
||||||
/// A running werkbank API on a random port, or `None` if no Mongo.
|
|
||||||
struct Harness {
|
|
||||||
base_url: String,
|
|
||||||
client: reqwest::Client,
|
|
||||||
pool: DatabasePool,
|
|
||||||
db_name: String,
|
|
||||||
}
|
|
||||||
|
|
||||||
async fn start() -> Option<Harness> {
|
|
||||||
let uri = std::env::var("TEST_MONGODB_URI")
|
|
||||||
.unwrap_or_else(|_| "mongodb://root:example@localhost:27017/?authSource=admin".into());
|
|
||||||
let db_name = format!("wba_{}", &uuid::Uuid::new_v4().simple().to_string()[..12]);
|
|
||||||
let pool = match DatabasePool::connect(&uri, &db_name).await {
|
|
||||||
Ok(p) => p,
|
|
||||||
Err(_) => {
|
|
||||||
eprintln!("SKIP werkbank_api: no MongoDB reachable at {uri}");
|
|
||||||
return None;
|
|
||||||
}
|
|
||||||
};
|
|
||||||
// Touch the tenant DB so indexes are ensured before the queue is used.
|
|
||||||
pool.for_tenant_id(TENANT).await.expect("tenant db");
|
|
||||||
|
|
||||||
let agent = ComplianceAgent::new(dev_config(uri, db_name.clone()), pool.clone());
|
|
||||||
let app = Router::new()
|
|
||||||
.route("/api/v1/werkbank/jobs/lease", post(werkbank_jobs::lease))
|
|
||||||
.route(
|
|
||||||
"/api/v1/werkbank/jobs/heartbeat",
|
|
||||||
post(werkbank_jobs::heartbeat),
|
|
||||||
)
|
|
||||||
.route(
|
|
||||||
"/api/v1/werkbank/jobs/complete",
|
|
||||||
post(werkbank_jobs::complete),
|
|
||||||
)
|
|
||||||
.layer(middleware::from_fn(werkbank_jobs::require_runner_token))
|
|
||||||
.layer(Extension(Arc::new(agent)));
|
|
||||||
|
|
||||||
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
|
|
||||||
let port = listener.local_addr().unwrap().port();
|
|
||||||
tokio::spawn(async move {
|
|
||||||
axum::serve(listener, app).await.ok();
|
|
||||||
});
|
|
||||||
|
|
||||||
Some(Harness {
|
|
||||||
base_url: format!("http://127.0.0.1:{port}"),
|
|
||||||
client: reqwest::Client::new(),
|
|
||||||
pool,
|
|
||||||
db_name,
|
|
||||||
})
|
|
||||||
}
|
|
||||||
|
|
||||||
impl Harness {
|
|
||||||
fn post(
|
|
||||||
&self,
|
|
||||||
path: &str,
|
|
||||||
token: Option<&str>,
|
|
||||||
body: serde_json::Value,
|
|
||||||
) -> reqwest::RequestBuilder {
|
|
||||||
let mut r = self
|
|
||||||
.client
|
|
||||||
.post(format!("{}{path}", self.base_url))
|
|
||||||
.json(&body);
|
|
||||||
if let Some(t) = token {
|
|
||||||
r = r.bearer_auth(t);
|
|
||||||
}
|
|
||||||
r
|
|
||||||
}
|
|
||||||
async fn cleanup(&self) {
|
|
||||||
let _ = self
|
|
||||||
.pool
|
|
||||||
.client()
|
|
||||||
.database(&format!("{}_{TENANT}", self.db_name))
|
|
||||||
.drop()
|
|
||||||
.await;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
fn finding_for(target: &str, fp: &str) -> Finding {
|
|
||||||
let mut f = Finding::new(
|
|
||||||
target.to_string(),
|
|
||||||
fp.to_string(),
|
|
||||||
"ics-probe".to_string(),
|
|
||||||
ScanType::IcsProbe,
|
|
||||||
"Modbus exposed".to_string(),
|
|
||||||
"unauthenticated".to_string(),
|
|
||||||
Severity::Critical,
|
|
||||||
);
|
|
||||||
f.rule_id = Some("ics-modbus-exposed".to_string());
|
|
||||||
f
|
|
||||||
}
|
|
||||||
|
|
||||||
#[tokio::test]
|
|
||||||
async fn lease_complete_persists_findings_against_the_target() {
|
|
||||||
let Some(h) = start().await else { return };
|
|
||||||
let db = h.pool.for_tenant_id(TENANT).await.unwrap();
|
|
||||||
let queue = JobQueue::new(&db);
|
|
||||||
|
|
||||||
// Seed a queued job.
|
|
||||||
let job = Job::plc_provision("job-1", TENANT, "target-1", InputRef::blob("sha256:x"), 180);
|
|
||||||
assert!(queue.enqueue(job, chrono::Utc::now()).await.unwrap());
|
|
||||||
|
|
||||||
// Lease it over HTTP.
|
|
||||||
let resp = h
|
|
||||||
.post(
|
|
||||||
"/api/v1/werkbank/jobs/lease",
|
|
||||||
Some(TEST_RUNNER_TOKEN),
|
|
||||||
serde_json::json!({
|
|
||||||
"tenant": TENANT, "runner_id": "r1", "executor": "docker",
|
|
||||||
"labels": [], "lease_ttl_secs": 60
|
|
||||||
}),
|
|
||||||
)
|
|
||||||
.send()
|
|
||||||
.await
|
|
||||||
.unwrap();
|
|
||||||
assert_eq!(resp.status(), 200, "lease should return a job");
|
|
||||||
let leased: LeasedJob = resp.json().await.unwrap();
|
|
||||||
assert_eq!(leased.job.id, "job-1");
|
|
||||||
|
|
||||||
// Complete it with a finding.
|
|
||||||
let mut result = JobResult::succeeded("job-1");
|
|
||||||
result.findings = vec![finding_for("target-1", "fp-abc")];
|
|
||||||
let resp = h
|
|
||||||
.post(
|
|
||||||
"/api/v1/werkbank/jobs/complete",
|
|
||||||
Some(TEST_RUNNER_TOKEN),
|
|
||||||
serde_json::json!({
|
|
||||||
"tenant": TENANT, "job_id": "job-1",
|
|
||||||
"lease_token": leased.lease_token, "result": result
|
|
||||||
}),
|
|
||||||
)
|
|
||||||
.send()
|
|
||||||
.await
|
|
||||||
.unwrap();
|
|
||||||
assert_eq!(resp.status(), 200);
|
|
||||||
assert!(resp.json::<serde_json::Value>().await.unwrap()["recorded"]
|
|
||||||
.as_bool()
|
|
||||||
.unwrap());
|
|
||||||
|
|
||||||
// The job is now succeeded, and the finding was persisted to the target.
|
|
||||||
assert_eq!(
|
|
||||||
queue.get("job-1").await.unwrap().unwrap().status,
|
|
||||||
JobStatus::Succeeded
|
|
||||||
);
|
|
||||||
let stored = db
|
|
||||||
.findings()
|
|
||||||
.find_one(mongodb::bson::doc! { "fingerprint": "fp-abc" })
|
|
||||||
.await
|
|
||||||
.unwrap();
|
|
||||||
assert!(stored.is_some(), "finding should be persisted");
|
|
||||||
|
|
||||||
h.cleanup().await;
|
|
||||||
}
|
|
||||||
|
|
||||||
#[tokio::test]
|
|
||||||
async fn empty_queue_leases_nothing() {
|
|
||||||
let Some(h) = start().await else { return };
|
|
||||||
let resp = h
|
|
||||||
.post(
|
|
||||||
"/api/v1/werkbank/jobs/lease",
|
|
||||||
Some(TEST_RUNNER_TOKEN),
|
|
||||||
serde_json::json!({
|
|
||||||
"tenant": TENANT, "runner_id": "r1", "executor": "docker",
|
|
||||||
"labels": [], "lease_ttl_secs": 60
|
|
||||||
}),
|
|
||||||
)
|
|
||||||
.send()
|
|
||||||
.await
|
|
||||||
.unwrap();
|
|
||||||
assert_eq!(resp.status(), 204, "no job → 204");
|
|
||||||
h.cleanup().await;
|
|
||||||
}
|
|
||||||
|
|
||||||
#[tokio::test]
|
|
||||||
async fn runner_endpoints_require_the_bearer_token() {
|
|
||||||
let Some(h) = start().await else { return };
|
|
||||||
let body = serde_json::json!({
|
|
||||||
"tenant": TENANT, "runner_id": "r1", "executor": "docker",
|
|
||||||
"labels": [], "lease_ttl_secs": 60
|
|
||||||
});
|
|
||||||
|
|
||||||
let no_token = h
|
|
||||||
.post("/api/v1/werkbank/jobs/lease", None, body.clone())
|
|
||||||
.send()
|
|
||||||
.await
|
|
||||||
.unwrap();
|
|
||||||
assert_eq!(no_token.status(), 401, "missing token → 401");
|
|
||||||
|
|
||||||
let bad_token = h
|
|
||||||
.post("/api/v1/werkbank/jobs/lease", Some("wrong"), body)
|
|
||||||
.send()
|
|
||||||
.await
|
|
||||||
.unwrap();
|
|
||||||
assert_eq!(bad_token.status(), 401, "wrong token → 401");
|
|
||||||
|
|
||||||
h.cleanup().await;
|
|
||||||
}
|
|
||||||
@@ -1,258 +0,0 @@
|
|||||||
//! 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();
|
|
||||||
}
|
|
||||||
@@ -64,11 +64,11 @@ struct Claims {
|
|||||||
const PUBLIC_ENDPOINTS: &[&str] = &["/api/v1/health"];
|
const PUBLIC_ENDPOINTS: &[&str] = &["/api/v1/health"];
|
||||||
|
|
||||||
/// Path prefixes that bypass JWT validation. The admin sub-router
|
/// Path prefixes that bypass JWT validation. The admin sub-router
|
||||||
/// (`/api/v1/admin/*`) and the Werkbank runner API (`/api/v1/werkbank/*`)
|
/// (`/api/v1/admin/*`) has its own static-bearer middleware and must
|
||||||
/// have their own static-bearer middleware and must not be routed through the
|
/// not be routed through the customer-JWT path — a Keycloak token
|
||||||
/// customer-JWT path — a Keycloak token always carries a single tenant_id and
|
/// always carries a single tenant_id and would semantically conflict
|
||||||
/// would semantically conflict with these cross-tenant / machine operations.
|
/// with cross-tenant admin operations.
|
||||||
const PUBLIC_PREFIXES: &[&str] = &["/api/v1/admin/", "/api/v1/werkbank/"];
|
const PUBLIC_PREFIXES: &[&str] = &["/api/v1/admin/"];
|
||||||
|
|
||||||
/// Middleware that validates Bearer JWT tokens against Keycloak's JWKS
|
/// Middleware that validates Bearer JWT tokens against Keycloak's JWKS
|
||||||
/// and attaches a `TenantContext` extension on success.
|
/// and attaches a `TenantContext` extension on success.
|
||||||
|
|||||||
@@ -53,11 +53,6 @@ pub struct AgentConfig {
|
|||||||
/// default: it needs Docker access in the agent's runtime, which is a
|
/// default: it needs Docker access in the agent's runtime, which is a
|
||||||
/// deployment opt-in.
|
/// deployment opt-in.
|
||||||
pub plc_runtime: PlcRuntimeConfig,
|
pub plc_runtime: PlcRuntimeConfig,
|
||||||
/// Static bearer for the Werkbank runner endpoints
|
|
||||||
/// (`/api/v1/werkbank/jobs/*`). Machine auth for runners leasing/completing
|
|
||||||
/// jobs — NOT a Keycloak JWT, since a runner acts across tenants. When
|
|
||||||
/// `None`, those endpoints are not mounted at all.
|
|
||||||
pub werkbank_runner_token: Option<SecretString>,
|
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Configuration for the ephemeral soft-PLC "provision-and-test" path (#183).
|
/// Configuration for the ephemeral soft-PLC "provision-and-test" path (#183).
|
||||||
|
|||||||
@@ -49,7 +49,5 @@ pub use repository::ScanTrigger;
|
|||||||
pub use sbom::{SbomEntry, VulnRef};
|
pub use sbom::{SbomEntry, VulnRef};
|
||||||
pub use scan::{ScanPhase, ScanRun, ScanRunStatus, ScanType};
|
pub use scan::{ScanPhase, ScanRun, ScanRunStatus, ScanType};
|
||||||
pub use werkbank::{
|
pub use werkbank::{
|
||||||
CompleteRequest, CompleteResponse, DastCollect, Executor, HeartbeatAck, HeartbeatRequest,
|
DastCollect, Executor, InputRef, Job, JobCollect, JobResult, JobRuntime, JobStatus, JobType,
|
||||||
InputRef, Job, JobCollect, JobRecord, JobResult, JobRuntime, JobStatus, JobType, LeaseRequest,
|
|
||||||
LeasedJob,
|
|
||||||
};
|
};
|
||||||
|
|||||||
@@ -266,142 +266,6 @@ 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,
|
|
||||||
}
|
|
||||||
|
|
||||||
// --- Runner ↔ control-plane transport (the pull API wire types) ---------------
|
|
||||||
// Shared so the runner (client) and the control plane (server) agree on shapes.
|
|
||||||
|
|
||||||
/// Runner → control plane: lease the oldest runnable job for this runner.
|
|
||||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
|
||||||
pub struct LeaseRequest {
|
|
||||||
/// The tenant queue to lease from.
|
|
||||||
pub tenant: String,
|
|
||||||
/// The runner id (advertised for attribution).
|
|
||||||
pub runner_id: String,
|
|
||||||
/// The executor this runner provides.
|
|
||||||
pub executor: Executor,
|
|
||||||
/// The capability labels this runner advertises.
|
|
||||||
#[serde(default)]
|
|
||||||
pub labels: Vec<String>,
|
|
||||||
/// Requested lease lifetime (the visibility timeout), in seconds.
|
|
||||||
pub lease_ttl_secs: u64,
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Runner → control plane: prove lease ownership and extend it.
|
|
||||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
|
||||||
pub struct HeartbeatRequest {
|
|
||||||
/// The tenant queue.
|
|
||||||
pub tenant: String,
|
|
||||||
/// The job being worked.
|
|
||||||
pub job_id: String,
|
|
||||||
/// The lease token from the [`LeasedJob`].
|
|
||||||
pub lease_token: String,
|
|
||||||
/// Lease lifetime to extend to, in seconds.
|
|
||||||
pub lease_ttl_secs: u64,
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Runner → control plane: record a job's terminal result.
|
|
||||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
|
||||||
pub struct CompleteRequest {
|
|
||||||
/// The tenant queue.
|
|
||||||
pub tenant: String,
|
|
||||||
/// The job being completed.
|
|
||||||
pub job_id: String,
|
|
||||||
/// The lease token proving ownership.
|
|
||||||
pub lease_token: String,
|
|
||||||
/// The result to record.
|
|
||||||
pub result: JobResult,
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Control plane → runner: whether the completion was recorded (false if the
|
|
||||||
/// lease was already lost — token mismatch or the job had become terminal).
|
|
||||||
#[derive(Debug, Clone, Copy, Serialize, Deserialize)]
|
|
||||||
pub struct CompleteResponse {
|
|
||||||
/// Whether the result was recorded.
|
|
||||||
pub recorded: bool,
|
|
||||||
}
|
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
#[allow(clippy::expect_used, clippy::unwrap_used)]
|
#[allow(clippy::expect_used, clippy::unwrap_used)]
|
||||||
mod tests {
|
mod tests {
|
||||||
|
|||||||
Reference in New Issue
Block a user