Compare commits
1
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
463d4f6eb2 |
@@ -47,11 +47,6 @@ PLC_RUNTIME_MAX_LIFETIME_SECS=180
|
||||
PLC_RUNTIME_OPENPLC_USER=openplc
|
||||
PLC_RUNTIME_OPENPLC_PASSWORD=openplc
|
||||
|
||||
# Werkbank runner API (/api/v1/werkbank/jobs/*, /api/v1/werkbank/artifacts/*).
|
||||
# When set, mounts the runner-facing queue + artifact endpoints behind this
|
||||
# bearer token; runners present the same token. Unset = endpoints not mounted.
|
||||
WERKBANK_RUNNER_TOKEN=
|
||||
|
||||
# Dashboard
|
||||
DASHBOARD_PORT=8080
|
||||
AGENT_API_URL=http://localhost:3001
|
||||
|
||||
@@ -107,8 +107,6 @@ jobs:
|
||||
run: cargo clippy -p compliance-dashboard --features web --no-default-features -- -D warnings
|
||||
- name: Clippy (mcp)
|
||||
run: cargo clippy -p compliance-mcp -- -D warnings
|
||||
- name: Clippy (werkbank-exec)
|
||||
run: cargo clippy -p werkbank-exec -- -D warnings
|
||||
|
||||
# Security audit
|
||||
- name: Security Audit
|
||||
@@ -117,8 +115,8 @@ jobs:
|
||||
RUSTC_WRAPPER: ""
|
||||
|
||||
# Tests (reuses compilation artifacts from clippy)
|
||||
- name: Tests (core + agent + werkbank-exec)
|
||||
run: cargo test -p compliance-core -p compliance-agent -p werkbank-exec --lib
|
||||
- name: Tests (core + agent)
|
||||
run: cargo test -p compliance-core -p compliance-agent --lib
|
||||
- name: Tests (dashboard server)
|
||||
run: cargo test -p compliance-dashboard --features server --no-default-features
|
||||
- name: Tests (dashboard web)
|
||||
|
||||
Generated
-20
@@ -699,7 +699,6 @@ dependencies = [
|
||||
"urlencoding",
|
||||
"uuid",
|
||||
"walkdir",
|
||||
"werkbank-exec",
|
||||
"zip",
|
||||
]
|
||||
|
||||
@@ -6715,25 +6714,6 @@ dependencies = [
|
||||
"rustls-pki-types",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "werkbank-exec"
|
||||
version = "0.1.0"
|
||||
dependencies = [
|
||||
"compliance-core",
|
||||
"compliance-dast",
|
||||
"futures-util",
|
||||
"hex",
|
||||
"regex",
|
||||
"reqwest",
|
||||
"secrecy",
|
||||
"sha2",
|
||||
"thiserror 2.0.18",
|
||||
"tokio",
|
||||
"tracing",
|
||||
"uuid",
|
||||
"walkdir",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "which"
|
||||
version = "6.0.3"
|
||||
|
||||
@@ -7,7 +7,6 @@ members = [
|
||||
"compliance-dast",
|
||||
"compliance-mcp",
|
||||
"compliance-smoke",
|
||||
"werkbank-exec",
|
||||
]
|
||||
resolver = "2"
|
||||
|
||||
|
||||
@@ -10,9 +10,6 @@ workspace = true
|
||||
compliance-core = { workspace = true, features = ["mongodb", "telemetry", "axum"] }
|
||||
compliance-graph = { path = "../compliance-graph" }
|
||||
compliance-dast = { path = "../compliance-dast" }
|
||||
# Shared dynamic-execution logic (soft-PLC provisioning + ICS probing), also
|
||||
# used by the Werkbank runner.
|
||||
werkbank-exec = { path = "../werkbank-exec" }
|
||||
# Native firmware build/target detection for bare-metal & RTOS artifacts.
|
||||
# Same-company IP, used directly (not via CLI) so the whole tramiton suite is
|
||||
# available to the onboarding classifier. NOTE: CI must be able to fetch this
|
||||
|
||||
@@ -14,7 +14,6 @@ pub mod pentest_handlers;
|
||||
pub use pentest_handlers as pentest;
|
||||
pub mod sbom;
|
||||
pub mod scans;
|
||||
pub mod werkbank_jobs;
|
||||
|
||||
// Re-export all handler functions so routes.rs can use `handlers::function_name`
|
||||
pub use dto::*;
|
||||
|
||||
@@ -1,289 +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, Path, Request};
|
||||
use axum::http::{header, StatusCode};
|
||||
use axum::middleware::Next;
|
||||
use axum::response::{IntoResponse, Response};
|
||||
use axum::Json;
|
||||
use mongodb::bson::{doc, oid::ObjectId};
|
||||
use secrecy::ExposeSecret;
|
||||
use serde::{Deserialize, Serialize};
|
||||
use std::time::Duration;
|
||||
|
||||
use compliance_core::models::werkbank::{
|
||||
CompleteRequest, CompleteResponse, HeartbeatRequest, InputRef, Job, JobResult, LeaseRequest,
|
||||
};
|
||||
use compliance_core::models::ArtifactKind;
|
||||
|
||||
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 }))
|
||||
}
|
||||
|
||||
/// `GET /api/v1/werkbank/artifacts/{hash}` — serve a content-addressed blob (the
|
||||
/// program a runner needs to load). The hash is validated against traversal by
|
||||
/// [`crate::ingest::blob::read_blob`]; a runner fetches this for a job's `blob`
|
||||
/// input.
|
||||
#[tracing::instrument(skip_all, fields(hash = %hash))]
|
||||
pub async fn serve_artifact(
|
||||
Extension(agent): AgentExt,
|
||||
Path(hash): Path<String>,
|
||||
) -> Result<Response, StatusCode> {
|
||||
let base = std::path::Path::new(&agent.config.artifact_store_base_path);
|
||||
match crate::ingest::blob::read_blob(base, &hash) {
|
||||
Ok(bytes) => {
|
||||
Ok(([(header::CONTENT_TYPE, "application/octet-stream")], bytes).into_response())
|
||||
}
|
||||
Err(_) => Err(StatusCode::NOT_FOUND),
|
||||
}
|
||||
}
|
||||
|
||||
/// Enqueue a `plc-provision` job for a target: extract its control-logic program,
|
||||
/// stash it as a content-addressed blob (which the runner fetches via
|
||||
/// [`serve_artifact`]), and queue the job. This is the control-plane "enqueue"
|
||||
/// half of the loop — a runner then leases it, provisions, and posts results.
|
||||
#[derive(Debug, Deserialize)]
|
||||
pub struct EnqueueRequest {
|
||||
/// The tenant whose queue to enqueue into.
|
||||
pub tenant: String,
|
||||
/// The onboarded target to test.
|
||||
pub target_id: String,
|
||||
}
|
||||
|
||||
/// The enqueued job's id.
|
||||
#[derive(Debug, Serialize)]
|
||||
pub struct EnqueueResponse {
|
||||
/// The new job id.
|
||||
pub job_id: String,
|
||||
/// Whether this call inserted it (false = already queued).
|
||||
pub enqueued: bool,
|
||||
}
|
||||
|
||||
#[tracing::instrument(skip_all, fields(tenant = %req.tenant, target = %req.target_id))]
|
||||
pub async fn enqueue(
|
||||
Extension(agent): AgentExt,
|
||||
Json(req): Json<EnqueueRequest>,
|
||||
) -> Result<Json<EnqueueResponse>, StatusCode> {
|
||||
let db = tenant_db(&agent, &req.tenant).await?;
|
||||
let oid = ObjectId::parse_str(&req.target_id).map_err(|_| StatusCode::BAD_REQUEST)?;
|
||||
let target = db
|
||||
.onboarded_targets()
|
||||
.find_one(doc! { "_id": oid })
|
||||
.await
|
||||
.map_err(internal)?
|
||||
.ok_or(StatusCode::NOT_FOUND)?;
|
||||
|
||||
// Extract the control-logic program from the target's PLC-source artifacts
|
||||
// (same selection as the in-process PLC scan).
|
||||
let ctx = crate::ingest::IngestContext::from_config(&agent.config, &req.target_id);
|
||||
let ingest_set = crate::ingest::ingest_all(&target, &ctx).map_err(internal)?;
|
||||
let program = target
|
||||
.artifacts
|
||||
.iter()
|
||||
.filter(|a| {
|
||||
matches!(
|
||||
a.kind,
|
||||
ArtifactKind::PlcProject | ArtifactKind::GitRepo | ArtifactKind::SourceArchive
|
||||
)
|
||||
})
|
||||
.find_map(|a| {
|
||||
let path = ingest_set
|
||||
.get(&a.id)
|
||||
.and_then(|ia| ia.working_path.clone())?;
|
||||
werkbank_exec::plc::extract_program(&path)
|
||||
})
|
||||
.ok_or(StatusCode::UNPROCESSABLE_ENTITY)?;
|
||||
|
||||
// Stash the program source so the runner can fetch it by hash.
|
||||
let base = std::path::Path::new(&agent.config.artifact_store_base_path);
|
||||
let hash =
|
||||
crate::ingest::blob::store_bytes(base, program.source.as_bytes()).map_err(internal)?;
|
||||
|
||||
let job_id = format!("job_{}", uuid::Uuid::new_v4().simple());
|
||||
let job = Job::plc_provision(
|
||||
&job_id,
|
||||
&req.tenant,
|
||||
&req.target_id,
|
||||
InputRef::blob(hash),
|
||||
agent.config.plc_runtime.max_lifetime_secs,
|
||||
);
|
||||
let enqueued = JobQueue::new(&db)
|
||||
.enqueue(job, chrono::Utc::now())
|
||||
.await
|
||||
.map_err(internal)?;
|
||||
Ok(Json(EnqueueResponse { job_id, enqueued }))
|
||||
}
|
||||
|
||||
/// 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::middleware::Next;
|
||||
use axum::response::Response;
|
||||
use axum::routing::{delete, get, post};
|
||||
use axum::routing::{delete, get};
|
||||
use axum::{middleware, Extension, Router};
|
||||
use tokio::sync::RwLock;
|
||||
use tower_http::cors::CorsLayer;
|
||||
@@ -72,43 +72,8 @@ pub async fn start_api_server(agent: ComplianceAgent, port: u16) -> Result<(), A
|
||||
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),
|
||||
)
|
||||
.route(
|
||||
"/api/v1/werkbank/jobs/enqueue",
|
||||
post(handlers::werkbank_jobs::enqueue),
|
||||
)
|
||||
.route(
|
||||
"/api/v1/werkbank/artifacts/{hash}",
|
||||
get(handlers::werkbank_jobs::serve_artifact),
|
||||
)
|
||||
.layer(middleware::from_fn(
|
||||
handlers::werkbank_jobs::require_runner_token,
|
||||
))
|
||||
} else {
|
||||
Router::new()
|
||||
};
|
||||
|
||||
let mut app = routes::build_router()
|
||||
.merge(admin_router)
|
||||
.merge(werkbank_router)
|
||||
// Allow large artifact uploads (PLC .projectarchive, firmware images,
|
||||
// mobile packages) — axum's default request-body limit is only 2 MiB.
|
||||
.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"),
|
||||
tenant_registry_url: env_var_opt("TENANT_REGISTRY_URL"),
|
||||
plc_runtime: load_plc_runtime_config(),
|
||||
werkbank_runner_token: env_secret_opt("WERKBANK_RUNNER_TOKEN"),
|
||||
})
|
||||
}
|
||||
|
||||
|
||||
@@ -465,34 +465,6 @@ 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(())
|
||||
}
|
||||
@@ -591,12 +563,6 @@ 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)
|
||||
|
||||
@@ -27,9 +27,6 @@ pub enum AgentError {
|
||||
#[error("Configuration error: {0}")]
|
||||
Config(String),
|
||||
|
||||
#[error("Dynamic-execution error: {0}")]
|
||||
Exec(#[from] werkbank_exec::ExecError),
|
||||
|
||||
#[error("{0}")]
|
||||
Other(String),
|
||||
}
|
||||
|
||||
@@ -32,30 +32,6 @@ pub fn hash_file(path: &Path) -> Result<(String, u64), AgentError> {
|
||||
Ok((hex::encode(hasher.finalize()), total))
|
||||
}
|
||||
|
||||
/// Store raw bytes in the content-addressed blob store under `base`, returning
|
||||
/// the SHA-256 digest. Used to stash a small derived artifact (e.g. the extracted
|
||||
/// PLC program source) so a Werkbank runner can fetch it by hash. Idempotent.
|
||||
pub fn store_bytes(base: &Path, bytes: &[u8]) -> Result<String, AgentError> {
|
||||
let sha = hex::encode(Sha256::digest(bytes));
|
||||
let dir = base.join("blobs").join(&sha[0..2]);
|
||||
fs::create_dir_all(&dir)?;
|
||||
let dest = dir.join(&sha);
|
||||
if !dest.exists() {
|
||||
fs::write(&dest, bytes)?;
|
||||
}
|
||||
Ok(sha)
|
||||
}
|
||||
|
||||
/// Read a blob's bytes by its SHA-256 digest. Rejects a non-hex/wrong-length hash
|
||||
/// so a request can't traverse outside the blob store.
|
||||
pub fn read_blob(base: &Path, sha: &str) -> Result<Vec<u8>, AgentError> {
|
||||
if sha.len() != 64 || !sha.bytes().all(|b| b.is_ascii_hexdigit()) {
|
||||
return Err(AgentError::Other(format!("invalid content hash '{sha}'")));
|
||||
}
|
||||
let path = base.join("blobs").join(&sha[0..2]).join(sha);
|
||||
Ok(fs::read(path)?)
|
||||
}
|
||||
|
||||
/// Copy `src` into the content-addressed blob store under `base`, returning the
|
||||
/// stored path. Idempotent: an already-present blob is not rewritten.
|
||||
pub fn store_file(base: &Path, src: &Path, sha: &str) -> Result<PathBuf, AgentError> {
|
||||
|
||||
@@ -6,7 +6,7 @@
|
||||
//! is also the reconciliation key against sibling products (a firmware sha256
|
||||
//! matches tramiton's `Artifact.sha256`).
|
||||
|
||||
pub(crate) mod blob;
|
||||
mod blob;
|
||||
|
||||
use std::collections::HashMap;
|
||||
use std::path::{Path, PathBuf};
|
||||
|
||||
@@ -16,4 +16,3 @@ pub mod ssh;
|
||||
#[allow(dead_code)]
|
||||
pub mod trackers;
|
||||
pub mod webhooks;
|
||||
pub mod werkbank;
|
||||
|
||||
@@ -343,7 +343,6 @@ mod tests {
|
||||
admin_api_token: None,
|
||||
tenant_registry_url: None,
|
||||
plc_runtime: compliance_core::PlcRuntimeConfig::default(),
|
||||
werkbank_runner_token: None,
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -14,7 +14,7 @@ use std::time::Duration;
|
||||
|
||||
use compliance_core::models::{Finding, ScanType, Severity};
|
||||
|
||||
use crate::fingerprint as dedup;
|
||||
use crate::pipeline::dedup;
|
||||
|
||||
/// Well-known deep-probe ports (each independent of any WebVisu HTTP port).
|
||||
const MODBUS_PORT: u16 = 502;
|
||||
@@ -5,6 +5,7 @@ pub mod firmware_sbom;
|
||||
pub mod git;
|
||||
pub mod gitleaks;
|
||||
mod graph_build;
|
||||
pub mod ics;
|
||||
mod issue_creation;
|
||||
pub mod lint;
|
||||
pub mod orchestrator;
|
||||
|
||||
@@ -587,7 +587,7 @@ impl PipelineOrchestrator {
|
||||
let path = ingest_set
|
||||
.get(&a.id)
|
||||
.and_then(|ia| ia.working_path.clone())?;
|
||||
werkbank_exec::plc::extract_program(&path)
|
||||
crate::pipeline::plc::runtime::extract_program(&path)
|
||||
});
|
||||
let Some(program) = program else {
|
||||
tracing::info!(
|
||||
@@ -597,9 +597,10 @@ impl PipelineOrchestrator {
|
||||
return Ok(0);
|
||||
};
|
||||
|
||||
let http = werkbank_exec::plc::http_client()?;
|
||||
let provisioner = werkbank_exec::plc::DockerSoftPlc::new(self.config.plc_runtime.clone());
|
||||
let outcome = werkbank_exec::plc::provision_and_test(
|
||||
let http = crate::pipeline::plc::runtime::http_client()?;
|
||||
let provisioner =
|
||||
crate::pipeline::plc::runtime::DockerSoftPlc::new(self.config.plc_runtime.clone());
|
||||
let outcome = crate::pipeline::plc::runtime::provision_and_test(
|
||||
&provisioner,
|
||||
&http,
|
||||
&self.config.plc_runtime,
|
||||
@@ -662,7 +663,7 @@ impl PipelineOrchestrator {
|
||||
};
|
||||
// Short per-request budget so an unreachable device doesn't stall the scan.
|
||||
let budget = std::time::Duration::from_secs(5);
|
||||
let findings = werkbank_exec::ics::probe_target(&endpoint, target_id, budget).await;
|
||||
let findings = crate::pipeline::ics::probe_target(&endpoint, target_id, budget).await;
|
||||
tracing::info!(
|
||||
target_id,
|
||||
endpoint = %endpoint,
|
||||
|
||||
@@ -9,6 +9,7 @@ pub mod lexer;
|
||||
pub mod parser;
|
||||
pub mod plcopen;
|
||||
pub mod rules;
|
||||
pub mod runtime;
|
||||
pub mod sbom;
|
||||
|
||||
use std::path::Path;
|
||||
|
||||
@@ -24,7 +24,7 @@ use compliance_core::models::dast::{DastFinding, DastScanRun, DastTarget, DastTa
|
||||
use compliance_core::models::Finding;
|
||||
use compliance_core::PlcRuntimeConfig;
|
||||
|
||||
use crate::error::ExecError;
|
||||
use crate::error::AgentError;
|
||||
|
||||
pub use provision::{DockerSoftPlc, ProvisionedRuntime, SoftPlc};
|
||||
|
||||
@@ -61,12 +61,12 @@ pub struct PlcProgram {
|
||||
|
||||
/// A cookie-aware HTTP client for the OpenPLC web UI. A fresh client per scan
|
||||
/// isolates the OpenPLC session (its Flask login cookie) from every other scan.
|
||||
pub fn http_client() -> Result<reqwest::Client, ExecError> {
|
||||
pub fn http_client() -> Result<reqwest::Client, AgentError> {
|
||||
reqwest::Client::builder()
|
||||
.cookie_store(true)
|
||||
.timeout(Duration::from_secs(30))
|
||||
.build()
|
||||
.map_err(ExecError::Http)
|
||||
.map_err(AgentError::Http)
|
||||
}
|
||||
|
||||
/// Pick the control-logic program to run from an ingested PLC source tree.
|
||||
@@ -151,7 +151,7 @@ pub async fn provision_and_test<P: SoftPlc>(
|
||||
cfg: &PlcRuntimeConfig,
|
||||
program: &PlcProgram,
|
||||
target_id: &str,
|
||||
) -> Result<ProvisionOutcome, ExecError> {
|
||||
) -> Result<ProvisionOutcome, AgentError> {
|
||||
let handle = provisioner.provision(target_id).await?;
|
||||
tracing::info!(
|
||||
target_id,
|
||||
@@ -193,7 +193,7 @@ async fn run_dynamic_test(
|
||||
program: &PlcProgram,
|
||||
target_id: &str,
|
||||
handle: &ProvisionedRuntime,
|
||||
) -> Result<ProvisionOutcome, ExecError> {
|
||||
) -> Result<ProvisionOutcome, AgentError> {
|
||||
let ready_budget = Duration::from_secs((cfg.max_lifetime_secs / 3).clamp(10, 60));
|
||||
openplc::wait_ready(http, &handle.webvisu_url, ready_budget).await?;
|
||||
|
||||
@@ -212,7 +212,8 @@ async fn run_dynamic_test(
|
||||
tokio::time::sleep(Duration::from_secs(3)).await;
|
||||
|
||||
let probe_budget = Duration::from_secs(5);
|
||||
let findings = crate::ics::probe_target(&handle.modbus_endpoint, target_id, probe_budget).await;
|
||||
let findings =
|
||||
crate::pipeline::ics::probe_target(&handle.modbus_endpoint, target_id, probe_budget).await;
|
||||
tracing::info!(
|
||||
target_id,
|
||||
instance = %handle.name,
|
||||
@@ -347,10 +348,10 @@ mod tests {
|
||||
}
|
||||
|
||||
impl SoftPlc for FakeSoftPlc {
|
||||
async fn provision(&self, _target_id: &str) -> Result<ProvisionedRuntime, ExecError> {
|
||||
async fn provision(&self, _target_id: &str) -> Result<ProvisionedRuntime, AgentError> {
|
||||
self.provisions.fetch_add(1, Ordering::SeqCst);
|
||||
if self.fail_provision {
|
||||
return Err(ExecError::Other("provision failed".into()));
|
||||
return Err(AgentError::Other("provision failed".into()));
|
||||
}
|
||||
// Unreachable address so run_dynamic_test blocks on readiness until the
|
||||
// deadline fires — exercising the teardown-on-deadline path.
|
||||
+15
-15
@@ -10,7 +10,7 @@
|
||||
|
||||
use std::time::Duration;
|
||||
|
||||
use crate::error::ExecError;
|
||||
use crate::error::AgentError;
|
||||
|
||||
use super::PlcProgram;
|
||||
|
||||
@@ -27,7 +27,7 @@ pub async fn wait_ready(
|
||||
http: &reqwest::Client,
|
||||
base_url: &str,
|
||||
budget: Duration,
|
||||
) -> Result<(), ExecError> {
|
||||
) -> Result<(), AgentError> {
|
||||
let login = format!("{base_url}/login");
|
||||
let outcome = tokio::time::timeout(budget, async {
|
||||
loop {
|
||||
@@ -40,7 +40,7 @@ pub async fn wait_ready(
|
||||
}
|
||||
})
|
||||
.await;
|
||||
outcome.map_err(|_| ExecError::Other(format!("OpenPLC at {base_url} did not become ready")))
|
||||
outcome.map_err(|_| AgentError::Other(format!("OpenPLC at {base_url} did not become ready")))
|
||||
}
|
||||
|
||||
/// Log in, upload the program, compile it, and start the runtime. On success the
|
||||
@@ -52,7 +52,7 @@ pub async fn load_and_start(
|
||||
password: &str,
|
||||
program: &PlcProgram,
|
||||
compile_budget: Duration,
|
||||
) -> Result<(), ExecError> {
|
||||
) -> Result<(), AgentError> {
|
||||
login(http, base_url, user, password).await?;
|
||||
let prog_file = upload_program(http, base_url, program).await?;
|
||||
save_program(http, base_url, &prog_file).await?;
|
||||
@@ -68,14 +68,14 @@ async fn login(
|
||||
base_url: &str,
|
||||
user: &str,
|
||||
password: &str,
|
||||
) -> Result<(), ExecError> {
|
||||
) -> Result<(), AgentError> {
|
||||
let resp = http
|
||||
.post(format!("{base_url}/login"))
|
||||
.form(&[("username", user), ("password", password)])
|
||||
.send()
|
||||
.await?;
|
||||
if resp.status().is_server_error() {
|
||||
return Err(ExecError::Other(format!(
|
||||
return Err(AgentError::Other(format!(
|
||||
"OpenPLC login failed: HTTP {}",
|
||||
resp.status()
|
||||
)));
|
||||
@@ -90,7 +90,7 @@ async fn upload_program(
|
||||
http: &reqwest::Client,
|
||||
base_url: &str,
|
||||
program: &PlcProgram,
|
||||
) -> Result<String, ExecError> {
|
||||
) -> Result<String, AgentError> {
|
||||
let part = reqwest::multipart::Part::text(program.source.clone())
|
||||
.file_name(program.file_name.clone())
|
||||
.mime_str("application/octet-stream")?;
|
||||
@@ -102,7 +102,7 @@ async fn upload_program(
|
||||
.await?;
|
||||
let html = resp.text().await?;
|
||||
parse_prog_file(&html).ok_or_else(|| {
|
||||
ExecError::Other("OpenPLC upload did not return a prog_file handle".to_string())
|
||||
AgentError::Other("OpenPLC upload did not return a prog_file handle".to_string())
|
||||
})
|
||||
}
|
||||
|
||||
@@ -113,7 +113,7 @@ async fn save_program(
|
||||
http: &reqwest::Client,
|
||||
base_url: &str,
|
||||
prog_file: &str,
|
||||
) -> Result<(), ExecError> {
|
||||
) -> Result<(), AgentError> {
|
||||
let epoch = std::time::SystemTime::now()
|
||||
.duration_since(std::time::UNIX_EPOCH)
|
||||
.map(|d| d.as_secs())
|
||||
@@ -130,7 +130,7 @@ async fn save_program(
|
||||
.send()
|
||||
.await?;
|
||||
if resp.status().is_server_error() {
|
||||
return Err(ExecError::Other(format!(
|
||||
return Err(AgentError::Other(format!(
|
||||
"OpenPLC save-program failed: HTTP {}",
|
||||
resp.status()
|
||||
)));
|
||||
@@ -146,7 +146,7 @@ async fn compile(
|
||||
base_url: &str,
|
||||
prog_file: &str,
|
||||
budget: Duration,
|
||||
) -> Result<(), ExecError> {
|
||||
) -> Result<(), AgentError> {
|
||||
http.get(format!("{base_url}/compile-program"))
|
||||
.query(&[("file", prog_file)])
|
||||
.send()
|
||||
@@ -168,20 +168,20 @@ async fn compile(
|
||||
.await;
|
||||
match outcome {
|
||||
Ok(true) => Ok(()),
|
||||
Ok(false) => Err(ExecError::Other(
|
||||
Ok(false) => Err(AgentError::Other(
|
||||
"OpenPLC compilation finished with errors".to_string(),
|
||||
)),
|
||||
Err(_) => Err(ExecError::Other(
|
||||
Err(_) => Err(AgentError::Other(
|
||||
"OpenPLC compilation did not finish in time".to_string(),
|
||||
)),
|
||||
}
|
||||
}
|
||||
|
||||
/// `GET /start_plc` — starts the runtime, opening Modbus/TCP on 502.
|
||||
async fn start(http: &reqwest::Client, base_url: &str) -> Result<(), ExecError> {
|
||||
async fn start(http: &reqwest::Client, base_url: &str) -> Result<(), AgentError> {
|
||||
let resp = http.get(format!("{base_url}/start_plc")).send().await?;
|
||||
if resp.status().is_server_error() {
|
||||
return Err(ExecError::Other(format!(
|
||||
return Err(AgentError::Other(format!(
|
||||
"OpenPLC start_plc failed: HTTP {}",
|
||||
resp.status()
|
||||
)));
|
||||
+6
-6
@@ -15,7 +15,7 @@ use std::time::{SystemTime, UNIX_EPOCH};
|
||||
|
||||
use compliance_core::PlcRuntimeConfig;
|
||||
|
||||
use crate::error::ExecError;
|
||||
use crate::error::AgentError;
|
||||
|
||||
/// The Modbus/TCP port an OpenPLC instance opens once a program is running.
|
||||
const MODBUS_PORT: u16 = 502;
|
||||
@@ -45,7 +45,7 @@ pub trait SoftPlc {
|
||||
fn provision(
|
||||
&self,
|
||||
target_id: &str,
|
||||
) -> impl std::future::Future<Output = Result<ProvisionedRuntime, ExecError>> + Send;
|
||||
) -> impl std::future::Future<Output = Result<ProvisionedRuntime, AgentError>> + Send;
|
||||
|
||||
/// Tear an instance down. Best-effort and idempotent — never fails the scan.
|
||||
fn teardown(&self, handle: &ProvisionedRuntime)
|
||||
@@ -65,7 +65,7 @@ impl DockerSoftPlc {
|
||||
}
|
||||
|
||||
impl SoftPlc for DockerSoftPlc {
|
||||
async fn provision(&self, target_id: &str) -> Result<ProvisionedRuntime, ExecError> {
|
||||
async fn provision(&self, target_id: &str) -> Result<ProvisionedRuntime, AgentError> {
|
||||
// Best-effort sweep of any container leaked by a crashed earlier run
|
||||
// before we add another. Only removes instances past their max lifetime,
|
||||
// so it can never disturb a concurrent run.
|
||||
@@ -75,7 +75,7 @@ impl SoftPlc for DockerSoftPlc {
|
||||
let args = run_args(&self.cfg, &name, target_id);
|
||||
let out = run_docker(&args).await?;
|
||||
if !out.status.success() {
|
||||
return Err(ExecError::Other(format!(
|
||||
return Err(AgentError::Other(format!(
|
||||
"docker run for soft-PLC {name} failed: {}",
|
||||
String::from_utf8_lossy(&out.stderr).trim()
|
||||
)));
|
||||
@@ -215,12 +215,12 @@ async fn reap_stale(cfg: &PlcRuntimeConfig, now: u64) {
|
||||
}
|
||||
|
||||
/// Run a `docker` subcommand, capturing its output.
|
||||
async fn run_docker(args: &[String]) -> Result<std::process::Output, ExecError> {
|
||||
async fn run_docker(args: &[String]) -> Result<std::process::Output, AgentError> {
|
||||
tokio::process::Command::new("docker")
|
||||
.args(args)
|
||||
.output()
|
||||
.await
|
||||
.map_err(ExecError::Io)
|
||||
.map_err(AgentError::Io)
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
@@ -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
|
||||
// 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;
|
||||
|
||||
@@ -15,54 +11,6 @@ use compliance_agent::database::DatabasePool;
|
||||
use compliance_core::AgentConfig;
|
||||
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.
|
||||
pub struct TestServer {
|
||||
pub base_url: String,
|
||||
@@ -85,7 +33,45 @@ impl TestServer {
|
||||
.await
|
||||
.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);
|
||||
|
||||
|
||||
@@ -1,291 +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::{get, 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::{
|
||||
Artifact, Finding, OnboardedTarget, PlcFormat, ScanType, Severity, TargetType,
|
||||
};
|
||||
|
||||
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),
|
||||
)
|
||||
.route(
|
||||
"/api/v1/werkbank/jobs/enqueue",
|
||||
post(werkbank_jobs::enqueue),
|
||||
)
|
||||
.route(
|
||||
"/api/v1/werkbank/artifacts/{hash}",
|
||||
get(werkbank_jobs::serve_artifact),
|
||||
)
|
||||
.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 enqueue_extracts_program_stores_a_blob_and_serves_it() {
|
||||
let Some(h) = start().await else { return };
|
||||
let db = h.pool.for_tenant_id(TENANT).await.unwrap();
|
||||
|
||||
// A PlcSps target with a single complete ST program uploaded.
|
||||
let dir = std::env::temp_dir().join(format!("wbq-prog-{}", uuid::Uuid::new_v4()));
|
||||
std::fs::create_dir_all(&dir).unwrap();
|
||||
let st = dir.join("main.st");
|
||||
std::fs::write(
|
||||
&st,
|
||||
"PROGRAM Main\nEND_PROGRAM\nCONFIGURATION C\n RESOURCE R\nEND_CONFIGURATION\n",
|
||||
)
|
||||
.unwrap();
|
||||
let mut target = OnboardedTarget::new("plc".into(), TargetType::PlcSps);
|
||||
let mut art = Artifact::plc_project("main.st", PlcFormat::StructuredText);
|
||||
art.stored_path = Some(st.to_string_lossy().to_string());
|
||||
target.artifacts.push(art);
|
||||
let ins = db.onboarded_targets().insert_one(&target).await.unwrap();
|
||||
let target_id = ins.inserted_id.as_object_id().unwrap().to_hex();
|
||||
|
||||
// Enqueue → a plc-provision job whose program is a content-addressed blob.
|
||||
let resp = h
|
||||
.post(
|
||||
"/api/v1/werkbank/jobs/enqueue",
|
||||
Some(TEST_RUNNER_TOKEN),
|
||||
serde_json::json!({ "tenant": TENANT, "target_id": target_id }),
|
||||
)
|
||||
.send()
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(resp.status(), 200, "enqueue should succeed");
|
||||
let body: serde_json::Value = resp.json().await.unwrap();
|
||||
let job_id = body["job_id"].as_str().unwrap().to_string();
|
||||
|
||||
let rec = JobQueue::new(&db).get(&job_id).await.unwrap().unwrap();
|
||||
let hash = rec
|
||||
.job
|
||||
.inputs
|
||||
.get("program")
|
||||
.and_then(|i| i.blob.clone())
|
||||
.expect("program blob");
|
||||
|
||||
// Serve the blob back and confirm it's the program source (what the runner
|
||||
// would fetch).
|
||||
let served = h
|
||||
.client
|
||||
.get(format!("{}/api/v1/werkbank/artifacts/{hash}", h.base_url))
|
||||
.bearer_auth(TEST_RUNNER_TOKEN)
|
||||
.send()
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(served.status(), 200);
|
||||
assert!(served.text().await.unwrap().contains("CONFIGURATION"));
|
||||
|
||||
h.cleanup().await;
|
||||
let _ = std::fs::remove_dir_all(&dir);
|
||||
}
|
||||
|
||||
#[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"];
|
||||
|
||||
/// Path prefixes that bypass JWT validation. The admin sub-router
|
||||
/// (`/api/v1/admin/*`) and the Werkbank runner API (`/api/v1/werkbank/*`)
|
||||
/// have their own static-bearer middleware and must not be routed through the
|
||||
/// customer-JWT path — a Keycloak token always carries a single tenant_id and
|
||||
/// would semantically conflict with these cross-tenant / machine operations.
|
||||
const PUBLIC_PREFIXES: &[&str] = &["/api/v1/admin/", "/api/v1/werkbank/"];
|
||||
/// (`/api/v1/admin/*`) has its own static-bearer middleware and must
|
||||
/// not be routed through the customer-JWT path — a Keycloak token
|
||||
/// always carries a single tenant_id and would semantically conflict
|
||||
/// with cross-tenant admin operations.
|
||||
const PUBLIC_PREFIXES: &[&str] = &["/api/v1/admin/"];
|
||||
|
||||
/// Middleware that validates Bearer JWT tokens against Keycloak's JWKS
|
||||
/// 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
|
||||
/// deployment opt-in.
|
||||
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).
|
||||
|
||||
@@ -49,7 +49,5 @@ pub use repository::ScanTrigger;
|
||||
pub use sbom::{SbomEntry, VulnRef};
|
||||
pub use scan::{ScanPhase, ScanRun, ScanRunStatus, ScanType};
|
||||
pub use werkbank::{
|
||||
CompleteRequest, CompleteResponse, DastCollect, Executor, HeartbeatAck, HeartbeatRequest,
|
||||
InputRef, Job, JobCollect, JobRecord, JobResult, JobRuntime, JobStatus, JobType, LeaseRequest,
|
||||
LeasedJob,
|
||||
DastCollect, Executor, InputRef, Job, JobCollect, JobResult, JobRuntime, JobStatus, JobType,
|
||||
};
|
||||
|
||||
@@ -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)]
|
||||
#[allow(clippy::expect_used, clippy::unwrap_used)]
|
||||
mod tests {
|
||||
|
||||
@@ -1,23 +0,0 @@
|
||||
[package]
|
||||
name = "werkbank-exec"
|
||||
version = "0.1.0"
|
||||
edition = "2021"
|
||||
description = "Shared dynamic-execution logic: soft-PLC provisioning + industrial-protocol probing, used by the compliance agent and the Werkbank runner."
|
||||
|
||||
[lints]
|
||||
workspace = true
|
||||
|
||||
[dependencies]
|
||||
compliance-core = { workspace = true }
|
||||
compliance-dast = { path = "../compliance-dast" }
|
||||
tokio = { workspace = true }
|
||||
reqwest = { workspace = true }
|
||||
uuid = { workspace = true }
|
||||
regex = { workspace = true }
|
||||
secrecy = { workspace = true }
|
||||
sha2 = { workspace = true }
|
||||
hex = { workspace = true }
|
||||
tracing = { workspace = true }
|
||||
thiserror = { workspace = true }
|
||||
walkdir = "2"
|
||||
futures-util = "0.3"
|
||||
@@ -1,16 +0,0 @@
|
||||
//! Error type for the dynamic-execution logic.
|
||||
|
||||
/// Anything that can go wrong provisioning and testing a soft-PLC. The compliance
|
||||
/// agent maps this into its own `AgentError` at the call boundary.
|
||||
#[derive(thiserror::Error, Debug)]
|
||||
pub enum ExecError {
|
||||
/// An HTTP request (to OpenPLC) failed.
|
||||
#[error("HTTP error: {0}")]
|
||||
Http(#[from] reqwest::Error),
|
||||
/// A local IO / process error (e.g. invoking `docker`).
|
||||
#[error("IO error: {0}")]
|
||||
Io(#[from] std::io::Error),
|
||||
/// Any other failure, with a message.
|
||||
#[error("{0}")]
|
||||
Other(String),
|
||||
}
|
||||
@@ -1,32 +0,0 @@
|
||||
//! Finding fingerprint helper (a SHA-256 over the salient parts), shared by the
|
||||
//! probe modules for stable dedup keys. Mirrors the agent's `dedup` helper.
|
||||
|
||||
use sha2::{Digest, Sha256};
|
||||
|
||||
/// A stable fingerprint over the given parts (order-sensitive, separated so
|
||||
/// `["ab","c"]` and `["a","bc"]` differ).
|
||||
pub fn compute_fingerprint(parts: &[&str]) -> String {
|
||||
let mut hasher = Sha256::new();
|
||||
for part in parts {
|
||||
hasher.update(part.as_bytes());
|
||||
hasher.update(b"|");
|
||||
}
|
||||
hex::encode(hasher.finalize())
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn deterministic_and_hex() {
|
||||
let a = compute_fingerprint(&["repo", "rule", "1"]);
|
||||
assert_eq!(a, compute_fingerprint(&["repo", "rule", "1"]));
|
||||
assert_eq!(a.len(), 64);
|
||||
assert!(a.chars().all(|c| c.is_ascii_hexdigit()));
|
||||
assert_ne!(
|
||||
compute_fingerprint(&["ab", "c"]),
|
||||
compute_fingerprint(&["a", "bc"])
|
||||
);
|
||||
}
|
||||
}
|
||||
@@ -1,18 +0,0 @@
|
||||
//! Shared dynamic-execution logic for Werkbank.
|
||||
//!
|
||||
//! The soft-PLC provisioning + industrial-protocol probing that turns a control-
|
||||
//! logic artifact into findings: provision an ephemeral OpenPLC, load the program,
|
||||
//! start it, probe it over Modbus/OPC-UA/EtherNet-IP, DAST its web endpoint, tear
|
||||
//! it down. Extracted from the compliance agent (#183) so both the agent (in
|
||||
//! process) and the Werkbank runner (WB-04) run identical logic.
|
||||
//!
|
||||
//! - [`ics`] — read-only industrial-protocol probing.
|
||||
//! - [`plc`] — ephemeral soft-PLC provisioning + the provision-and-test loop.
|
||||
|
||||
pub mod error;
|
||||
pub mod ics;
|
||||
pub mod plc;
|
||||
|
||||
mod fingerprint;
|
||||
|
||||
pub use error::ExecError;
|
||||
Reference in New Issue
Block a user