Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
c1e0a80f90 |
@@ -34,19 +34,6 @@ SCAN_SCHEDULE=0 0 */6 * * *
|
|||||||
CVE_MONITOR_SCHEDULE=0 0 0 * * *
|
CVE_MONITOR_SCHEDULE=0 0 0 * * *
|
||||||
GIT_CLONE_BASE_PATH=/tmp/compliance-scanner/repos
|
GIT_CLONE_BASE_PATH=/tmp/compliance-scanner/repos
|
||||||
|
|
||||||
# Dynamic PLC testing — ephemeral soft-PLC provisioning (#183). Off unless
|
|
||||||
# enabled; requires the agent container to have Docker access (socket mount).
|
|
||||||
# When on, a PLC/SPS target with control logic but no reachable device gets its
|
|
||||||
# logic instantiated on a throwaway OpenPLC, probed, then torn down.
|
|
||||||
PLC_RUNTIME_ENABLED=0
|
|
||||||
PLC_RUNTIME_IMAGE=registry.meghsakha.com/openplc:latest
|
|
||||||
PLC_RUNTIME_NETWORK=certifai
|
|
||||||
PLC_RUNTIME_MEMORY=512m
|
|
||||||
PLC_RUNTIME_CPUS=0.5
|
|
||||||
PLC_RUNTIME_MAX_LIFETIME_SECS=180
|
|
||||||
PLC_RUNTIME_OPENPLC_USER=openplc
|
|
||||||
PLC_RUNTIME_OPENPLC_PASSWORD=openplc
|
|
||||||
|
|
||||||
# Dashboard
|
# Dashboard
|
||||||
DASHBOARD_PORT=8080
|
DASHBOARD_PORT=8080
|
||||||
AGENT_API_URL=http://localhost:3001
|
AGENT_API_URL=http://localhost:3001
|
||||||
|
|||||||
Generated
-1
@@ -723,7 +723,6 @@ dependencies = [
|
|||||||
"sha2",
|
"sha2",
|
||||||
"thiserror 2.0.18",
|
"thiserror 2.0.18",
|
||||||
"tokio",
|
"tokio",
|
||||||
"toml",
|
|
||||||
"tracing",
|
"tracing",
|
||||||
"tracing-opentelemetry",
|
"tracing-opentelemetry",
|
||||||
"tracing-subscriber",
|
"tracing-subscriber",
|
||||||
|
|||||||
+1
-1
@@ -23,7 +23,7 @@ tracing = "0.1"
|
|||||||
tracing-subscriber = { version = "0.3", features = ["env-filter"] }
|
tracing-subscriber = { version = "0.3", features = ["env-filter"] }
|
||||||
chrono = { version = "0.4", features = ["serde"] }
|
chrono = { version = "0.4", features = ["serde"] }
|
||||||
mongodb = { version = "3", features = ["rustls-tls", "compat-3-0-0"] }
|
mongodb = { version = "3", features = ["rustls-tls", "compat-3-0-0"] }
|
||||||
reqwest = { version = "0.12", features = ["json", "rustls-tls", "multipart", "cookies"], default-features = false }
|
reqwest = { version = "0.12", features = ["json", "rustls-tls", "multipart"], default-features = false }
|
||||||
thiserror = "2"
|
thiserror = "2"
|
||||||
sha2 = "0.10"
|
sha2 = "0.10"
|
||||||
hex = "0.4"
|
hex = "0.4"
|
||||||
|
|||||||
@@ -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))
|
||||||
|
|||||||
@@ -1,4 +1,3 @@
|
|||||||
use compliance_core::config::PlcRuntimeConfig;
|
|
||||||
use compliance_core::AgentConfig;
|
use compliance_core::AgentConfig;
|
||||||
use secrecy::SecretString;
|
use secrecy::SecretString;
|
||||||
|
|
||||||
@@ -64,29 +63,5 @@ pub fn load_config() -> Result<AgentConfig, AgentError> {
|
|||||||
pentest_imap_password: env_secret_opt("PENTEST_IMAP_PASSWORD"),
|
pentest_imap_password: env_secret_opt("PENTEST_IMAP_PASSWORD"),
|
||||||
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(),
|
|
||||||
werkbank_runner_token: env_secret_opt("WERKBANK_RUNNER_TOKEN"),
|
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Build the ephemeral soft-PLC provisioning config from the environment,
|
|
||||||
/// falling back to [`PlcRuntimeConfig::default`] for any unset knob. Disabled
|
|
||||||
/// unless `PLC_RUNTIME_ENABLED` is truthy — it requires Docker access.
|
|
||||||
fn load_plc_runtime_config() -> PlcRuntimeConfig {
|
|
||||||
let d = PlcRuntimeConfig::default();
|
|
||||||
PlcRuntimeConfig {
|
|
||||||
enabled: env_var_opt("PLC_RUNTIME_ENABLED")
|
|
||||||
.map(|v| v == "1" || v.eq_ignore_ascii_case("true"))
|
|
||||||
.unwrap_or(d.enabled),
|
|
||||||
image: env_var_opt("PLC_RUNTIME_IMAGE").unwrap_or(d.image),
|
|
||||||
network: env_var_opt("PLC_RUNTIME_NETWORK").unwrap_or(d.network),
|
|
||||||
memory: env_var_opt("PLC_RUNTIME_MEMORY").unwrap_or(d.memory),
|
|
||||||
cpus: env_var_opt("PLC_RUNTIME_CPUS").unwrap_or(d.cpus),
|
|
||||||
max_lifetime_secs: env_var_opt("PLC_RUNTIME_MAX_LIFETIME_SECS")
|
|
||||||
.and_then(|v| v.parse().ok())
|
|
||||||
.unwrap_or(d.max_lifetime_secs),
|
|
||||||
openplc_user: env_var_opt("PLC_RUNTIME_OPENPLC_USER").unwrap_or(d.openplc_user),
|
|
||||||
openplc_password: env_secret_opt("PLC_RUNTIME_OPENPLC_PASSWORD")
|
|
||||||
.unwrap_or(d.openplc_password),
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|||||||
@@ -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;
|
|
||||||
|
|||||||
@@ -342,8 +342,6 @@ mod tests {
|
|||||||
pentest_imap_password: None,
|
pentest_imap_password: None,
|
||||||
admin_api_token: None,
|
admin_api_token: None,
|
||||||
tenant_registry_url: None,
|
tenant_registry_url: None,
|
||||||
plc_runtime: compliance_core::PlcRuntimeConfig::default(),
|
|
||||||
werkbank_runner_token: None,
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -103,39 +103,6 @@ async fn modbus_findings(host: &str, port: u16, repo_id: &str, budget: Duration)
|
|||||||
);
|
);
|
||||||
findings.push(f);
|
findings.push(f);
|
||||||
}
|
}
|
||||||
|
|
||||||
// Exposed process points: coils / holding registers that a read enumerated
|
|
||||||
// and that, over unauthenticated Modbus/TCP, are also writable. This is the
|
|
||||||
// concrete attack surface behind the exposure — the live variables an
|
|
||||||
// attacker can overwrite. (Read-only to detect: we never write.)
|
|
||||||
let coils = probe.coils_readable.unwrap_or(0);
|
|
||||||
let registers = probe.holding_registers_readable.unwrap_or(0);
|
|
||||||
if coils > 0 || registers > 0 {
|
|
||||||
let fp = dedup::compute_fingerprint(&[repo_id, "ics-modbus-exposed-points", &target]);
|
|
||||||
let mut f = Finding::new(
|
|
||||||
repo_id.to_string(),
|
|
||||||
fp,
|
|
||||||
"ics-probe".to_string(),
|
|
||||||
ScanType::IcsProbe,
|
|
||||||
"Writable process points exposed over unauthenticated Modbus/TCP".to_string(),
|
|
||||||
format!(
|
|
||||||
"Reading the device at {target} enumerated {coils} coil(s) and {registers} \
|
|
||||||
holding register(s). Coils and holding registers are read/write process points \
|
|
||||||
in Modbus, so any host that can reach this port can not only read but overwrite \
|
|
||||||
live process state (force coils, change setpoints) without authentication."
|
|
||||||
),
|
|
||||||
Severity::High,
|
|
||||||
);
|
|
||||||
f.rule_id = Some("ics-modbus-exposed-points".to_string());
|
|
||||||
f.cwe = Some("CWE-306".to_string());
|
|
||||||
f.remediation = Some(
|
|
||||||
"Segment the Modbus/TCP port to a trusted control network; where the device \
|
|
||||||
supports it use Modbus/TLS or an authenticating protocol gateway; restrict which \
|
|
||||||
function codes and register ranges are reachable from outside the control zone."
|
|
||||||
.to_string(),
|
|
||||||
);
|
|
||||||
findings.push(f);
|
|
||||||
}
|
|
||||||
findings
|
findings
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -21,13 +21,6 @@ pub struct ModbusProbe {
|
|||||||
pub speaks_modbus: bool,
|
pub speaks_modbus: bool,
|
||||||
/// Device identity, if disclosed via Read Device Identification (FC 43 / 14).
|
/// Device identity, if disclosed via Read Device Identification (FC 43 / 14).
|
||||||
pub device: Option<DeviceId>,
|
pub device: Option<DeviceId>,
|
||||||
/// Coils returned by a Read Coils of the first block, if that address range
|
|
||||||
/// exists. Coils are read/write process bits, so an exposed block is an
|
|
||||||
/// unauthenticated write surface on the live process.
|
|
||||||
pub coils_readable: Option<u16>,
|
|
||||||
/// Holding registers returned by a Read Holding Registers of the first block,
|
|
||||||
/// if that range exists. Holding registers are read/write process words.
|
|
||||||
pub holding_registers_readable: Option<u16>,
|
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Vendor / product / revision from Read Device Identification.
|
/// Vendor / product / revision from Read Device Identification.
|
||||||
@@ -38,13 +31,8 @@ pub struct DeviceId {
|
|||||||
pub revision: Option<String>,
|
pub revision: Option<String>,
|
||||||
}
|
}
|
||||||
|
|
||||||
/// How many coils / holding registers to request when enumerating the exposed
|
/// Probe a Modbus/TCP endpoint. Read-only: issues a Read Holding Registers and a
|
||||||
/// process surface. Read-only: a normal reply means the block exists and is,
|
/// Read Device Identification request; never writes to the device.
|
||||||
/// over unauthenticated Modbus/TCP, also writable.
|
|
||||||
const ENUM_QTY: u16 = 16;
|
|
||||||
|
|
||||||
/// Probe a Modbus/TCP endpoint. Read-only: issues Read Holding Registers, Read
|
|
||||||
/// Coils, and Read Device Identification requests; never writes to the device.
|
|
||||||
pub async fn probe(host: &str, port: u16, budget: Duration) -> ModbusProbe {
|
pub async fn probe(host: &str, port: u16, budget: Duration) -> ModbusProbe {
|
||||||
let mut out = ModbusProbe::default();
|
let mut out = ModbusProbe::default();
|
||||||
let Ok(Ok(mut stream)) = timeout(budget, TcpStream::connect((host, port))).await else {
|
let Ok(Ok(mut stream)) = timeout(budget, TcpStream::connect((host, port))).await else {
|
||||||
@@ -52,28 +40,13 @@ pub async fn probe(host: &str, port: u16, budget: Duration) -> ModbusProbe {
|
|||||||
};
|
};
|
||||||
out.reachable = true;
|
out.reachable = true;
|
||||||
|
|
||||||
// Read Holding Registers (FC 0x03), unit 1, addr 0 — a benign read that also
|
// Read Holding Registers (FC 0x03), unit 1, addr 0, qty 1 — a benign read.
|
||||||
// enumerates the exposed register block.
|
let rhr = [0x03u8, 0x00, 0x00, 0x00, 0x01];
|
||||||
let rhr = [0x03u8, 0x00, 0x00, (ENUM_QTY >> 8) as u8, ENUM_QTY as u8];
|
|
||||||
if let Some(resp) = txn(&mut stream, 1, &rhr, budget).await {
|
if let Some(resp) = txn(&mut stream, 1, &rhr, budget).await {
|
||||||
// A normal reply (0x03) or an exception (0x83) both prove it speaks Modbus.
|
// A normal reply (0x03) or an exception (0x83) both prove it speaks Modbus.
|
||||||
if matches!(resp.first(), Some(0x03) | Some(0x83)) {
|
if matches!(resp.first(), Some(0x03) | Some(0x83)) {
|
||||||
out.speaks_modbus = true;
|
out.speaks_modbus = true;
|
||||||
}
|
}
|
||||||
if resp.first() == Some(&0x03) {
|
|
||||||
out.holding_registers_readable = Some(register_count_from_reply(&resp));
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// Read Coils (FC 0x01), addr 0 — enumerates the exposed coil (bit) block.
|
|
||||||
let rc = [0x01u8, 0x00, 0x00, (ENUM_QTY >> 8) as u8, ENUM_QTY as u8];
|
|
||||||
if let Some(resp) = txn(&mut stream, 1, &rc, budget).await {
|
|
||||||
if matches!(resp.first(), Some(0x01) | Some(0x81)) {
|
|
||||||
out.speaks_modbus = true;
|
|
||||||
}
|
|
||||||
if resp.first() == Some(&0x01) {
|
|
||||||
out.coils_readable = Some(coil_count_from_reply(&resp));
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// Read Device Identification (FC 0x2B / MEI 0x0E), basic (0x01), object 0.
|
// Read Device Identification (FC 0x2B / MEI 0x0E), basic (0x01), object 0.
|
||||||
@@ -87,17 +60,6 @@ pub async fn probe(host: &str, port: u16, budget: Duration) -> ModbusProbe {
|
|||||||
out
|
out
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Coils reported by a Read Coils reply `[0x01, byte_count, data…]` (8 per byte).
|
|
||||||
fn coil_count_from_reply(pdu: &[u8]) -> u16 {
|
|
||||||
pdu.get(1).map(|&b| u16::from(b) * 8).unwrap_or(0)
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Registers reported by a Read Holding Registers reply `[0x03, byte_count,
|
|
||||||
/// data…]` (2 bytes per register).
|
|
||||||
fn register_count_from_reply(pdu: &[u8]) -> u16 {
|
|
||||||
pdu.get(1).map(|&b| u16::from(b) / 2).unwrap_or(0)
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Send one Modbus PDU and return the response PDU (function code + data), or
|
/// Send one Modbus PDU and return the response PDU (function code + data), or
|
||||||
/// `None` on timeout / malformed reply.
|
/// `None` on timeout / malformed reply.
|
||||||
async fn txn(stream: &mut TcpStream, unit: u8, pdu: &[u8], budget: Duration) -> Option<Vec<u8>> {
|
async fn txn(stream: &mut TcpStream, unit: u8, pdu: &[u8], budget: Duration) -> Option<Vec<u8>> {
|
||||||
@@ -192,8 +154,7 @@ mod tests {
|
|||||||
break;
|
break;
|
||||||
}
|
}
|
||||||
let reply_pdu: Vec<u8> = match pdu.first() {
|
let reply_pdu: Vec<u8> = match pdu.first() {
|
||||||
Some(0x03) => vec![0x03, 0x02, 0x00, 0x00], // 1 register (byte_count 2)
|
Some(0x03) => vec![0x03, 0x02, 0x00, 0x00], // 1 register = 0
|
||||||
Some(0x01) => vec![0x01, 0x02, 0xFF, 0xFF], // 16 coils (byte_count 2)
|
|
||||||
Some(0x2B) if with_device => vec![
|
Some(0x2B) if with_device => vec![
|
||||||
0x2B, 0x0E, 0x01, 0x81, 0x00, 0x00, 0x02, // 2 objects
|
0x2B, 0x0E, 0x01, 0x81, 0x00, 0x00, 0x02, // 2 objects
|
||||||
0x00, 0x04, b'A', b'C', b'M', b'E', // vendor
|
0x00, 0x04, b'A', b'C', b'M', b'E', // vendor
|
||||||
@@ -224,24 +185,6 @@ mod tests {
|
|||||||
assert_eq!(dev.product.as_deref(), Some("PLC"));
|
assert_eq!(dev.product.as_deref(), Some("PLC"));
|
||||||
}
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
|
||||||
async fn probe_enumerates_exposed_process_points() {
|
|
||||||
let addr = mock_server(false).await;
|
|
||||||
let p = probe(&addr.ip().to_string(), addr.port(), Duration::from_secs(2)).await;
|
|
||||||
assert!(p.speaks_modbus);
|
|
||||||
// The mock returns a 2-byte holding-register block (1 register) and a
|
|
||||||
// 2-byte coil block (16 coils).
|
|
||||||
assert_eq!(p.holding_registers_readable, Some(1));
|
|
||||||
assert_eq!(p.coils_readable, Some(16));
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn reply_counts_decode_byte_counts() {
|
|
||||||
assert_eq!(register_count_from_reply(&[0x03, 0x08]), 4); // 8 bytes → 4 regs
|
|
||||||
assert_eq!(coil_count_from_reply(&[0x01, 0x03]), 24); // 3 bytes → 24 coils
|
|
||||||
assert_eq!(register_count_from_reply(&[0x03]), 0); // malformed → 0
|
|
||||||
}
|
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn probe_reports_unreachable_for_a_closed_port() {
|
async fn probe_reports_unreachable_for_a_closed_port() {
|
||||||
// 127.0.0.1:1 is (almost certainly) closed.
|
// 127.0.0.1:1 is (almost certainly) closed.
|
||||||
|
|||||||
@@ -404,21 +404,6 @@ impl PipelineOrchestrator {
|
|||||||
let ics = plan.has(ScanType::IcsProbe);
|
let ics = plan.has(ScanType::IcsProbe);
|
||||||
if plc {
|
if plc {
|
||||||
new_count += self.run_plc_scan(target, &target_id, scan_run_id).await?;
|
new_count += self.run_plc_scan(target, &target_id, scan_run_id).await?;
|
||||||
// Provision-and-test (#183): with the control logic but no reachable
|
|
||||||
// device, instantiate it on an ephemeral soft-PLC and probe that
|
|
||||||
// instead of the customer's OT network. Opt-in (needs Docker) and only
|
|
||||||
// when there is no live URL to probe directly. Never fails the scan.
|
|
||||||
if self.config.plc_runtime.enabled && target.live_url().is_none() {
|
|
||||||
match self
|
|
||||||
.run_provisioned_plc_test(target, &target_id, scan_run_id)
|
|
||||||
.await
|
|
||||||
{
|
|
||||||
Ok(n) => new_count += n,
|
|
||||||
Err(e) => {
|
|
||||||
tracing::warn!(target_id = %target_id, error = %e, "provision-and-test failed")
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
if ics {
|
if ics {
|
||||||
new_count += self.run_ics_probe(target, &target_id, scan_run_id).await?;
|
new_count += self.run_ics_probe(target, &target_id, scan_run_id).await?;
|
||||||
@@ -555,98 +540,6 @@ impl PipelineOrchestrator {
|
|||||||
Ok(new_count)
|
Ok(new_count)
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Provision-and-test (#183): instantiate the target's control logic on an
|
|
||||||
/// ephemeral soft-PLC (OpenPLC), start it, probe the provisioned Modbus
|
|
||||||
/// endpoint, and tear the instance down. Used when a PLC/SPS target has the
|
|
||||||
/// control logic but no reachable live device to probe directly. Guarded by
|
|
||||||
/// `plc_runtime.enabled` (needs Docker); persists the same [`ScanType::IcsProbe`]
|
|
||||||
/// findings as a live probe.
|
|
||||||
async fn run_provisioned_plc_test(
|
|
||||||
&self,
|
|
||||||
target: &OnboardedTarget,
|
|
||||||
target_id: &str,
|
|
||||||
scan_run_id: &str,
|
|
||||||
) -> Result<u32, AgentError> {
|
|
||||||
self.update_phase(scan_run_id, "plc_provision").await;
|
|
||||||
|
|
||||||
// Locate a loadable control-logic program among the PLC-source artifacts
|
|
||||||
// (same selection as the static PLC scan: dedicated PLC projects plus code
|
|
||||||
// artifacts holding PLCopen XML / ST exports).
|
|
||||||
let ctx = crate::ingest::IngestContext::from_config(&self.config, target_id);
|
|
||||||
let ingest_set = crate::ingest::ingest_all(target, &ctx)?;
|
|
||||||
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())?;
|
|
||||||
crate::pipeline::plc::runtime::extract_program(&path)
|
|
||||||
});
|
|
||||||
let Some(program) = program else {
|
|
||||||
tracing::info!(
|
|
||||||
target_id,
|
|
||||||
"provision-and-test: no loadable control-logic program"
|
|
||||||
);
|
|
||||||
return Ok(0);
|
|
||||||
};
|
|
||||||
|
|
||||||
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,
|
|
||||||
&program,
|
|
||||||
target_id,
|
|
||||||
)
|
|
||||||
.await?;
|
|
||||||
tracing::info!(
|
|
||||||
target_id,
|
|
||||||
found = outcome.findings.len(),
|
|
||||||
dast = outcome.dast.is_some(),
|
|
||||||
"provision-and-test complete"
|
|
||||||
);
|
|
||||||
|
|
||||||
let mut new_count = 0u32;
|
|
||||||
for mut finding in outcome.findings {
|
|
||||||
finding.scan_run_id = Some(scan_run_id.to_string());
|
|
||||||
if self
|
|
||||||
.db
|
|
||||||
.findings()
|
|
||||||
.find_one(doc! { "fingerprint": &finding.fingerprint })
|
|
||||||
.await?
|
|
||||||
.is_none()
|
|
||||||
{
|
|
||||||
self.db.findings().insert_one(&finding).await?;
|
|
||||||
new_count += 1;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// Persist the DAST scan of the provisioned web endpoint, linked to this
|
|
||||||
// scan run (mirrors `maybe_trigger_dast`).
|
|
||||||
if let Some(dast) = outcome.dast {
|
|
||||||
let mut scan_run = dast.scan_run;
|
|
||||||
scan_run.sast_scan_run_id = Some(scan_run_id.to_string());
|
|
||||||
if let Err(e) = self.db.dast_scan_runs().insert_one(&scan_run).await {
|
|
||||||
tracing::warn!(target_id, error = %e, "failed to store provisioned DAST scan run");
|
|
||||||
}
|
|
||||||
for finding in &dast.findings {
|
|
||||||
if let Err(e) = self.db.dast_findings().insert_one(finding).await {
|
|
||||||
tracing::warn!(target_id, error = %e, "failed to store provisioned DAST finding");
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
Ok(new_count)
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Probe a running PLC/SPS device over industrial protocols (Modbus/TCP, …)
|
/// Probe a running PLC/SPS device over industrial protocols (Modbus/TCP, …)
|
||||||
/// and persist findings for exposed / unauthenticated control access. The
|
/// and persist findings for exposed / unauthenticated control access. The
|
||||||
/// probe is read-only; it targets the Modbus port of the target's live URL.
|
/// probe is read-only; it targets the Modbus port of the target's live URL.
|
||||||
|
|||||||
@@ -9,7 +9,6 @@ pub mod lexer;
|
|||||||
pub mod parser;
|
pub mod parser;
|
||||||
pub mod plcopen;
|
pub mod plcopen;
|
||||||
pub mod rules;
|
pub mod rules;
|
||||||
pub mod runtime;
|
|
||||||
pub mod sbom;
|
pub mod sbom;
|
||||||
|
|
||||||
use std::path::Path;
|
use std::path::Path;
|
||||||
|
|||||||
@@ -1,423 +0,0 @@
|
|||||||
//! Dynamic PLC testing via an ephemeral soft-PLC (#183).
|
|
||||||
//!
|
|
||||||
//! When a PLC/SPS target ships control logic but no reachable live device, the
|
|
||||||
//! agent instantiates that logic itself instead of trying to reach the customer's
|
|
||||||
//! OT network: it provisions a throwaway soft-PLC (OpenPLC) container in-cluster,
|
|
||||||
//! loads the program, starts the runtime, probes it over industrial protocols,
|
|
||||||
//! then tears the instance down. No customer network access, sandboxed, and
|
|
||||||
//! reproducible — destructive tests become safe because the target is ours.
|
|
||||||
//!
|
|
||||||
//! - [`provision`] owns the container lifecycle (sub-task 1 + 5).
|
|
||||||
//! - [`openplc`] loads the program into the running instance (sub-task 2).
|
|
||||||
//! - [`provision_and_test`] composes them with a hard deadline and guaranteed
|
|
||||||
//! teardown, and runs the ICS probe against the provisioned endpoint.
|
|
||||||
|
|
||||||
pub mod openplc;
|
|
||||||
pub mod provision;
|
|
||||||
|
|
||||||
use std::path::Path;
|
|
||||||
use std::time::Duration;
|
|
||||||
|
|
||||||
use secrecy::ExposeSecret;
|
|
||||||
|
|
||||||
use compliance_core::models::dast::{DastFinding, DastScanRun, DastTarget, DastTargetType};
|
|
||||||
use compliance_core::models::Finding;
|
|
||||||
use compliance_core::PlcRuntimeConfig;
|
|
||||||
|
|
||||||
use crate::error::AgentError;
|
|
||||||
|
|
||||||
pub use provision::{DockerSoftPlc, ProvisionedRuntime, SoftPlc};
|
|
||||||
|
|
||||||
/// The result of a DAST scan against a provisioned web endpoint.
|
|
||||||
#[derive(Debug)]
|
|
||||||
pub struct DastRunResult {
|
|
||||||
/// The scan-run record (linked to the onboarded target).
|
|
||||||
pub scan_run: DastScanRun,
|
|
||||||
/// The DAST findings.
|
|
||||||
pub findings: Vec<DastFinding>,
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Everything a provision-and-test run produced: the ICS-probe findings plus, if
|
|
||||||
/// it ran, the DAST scan of the provisioned web endpoint. The caller persists
|
|
||||||
/// both — keeping this a plain data return means the whole run is portable to a
|
|
||||||
/// remote execution backend that just hands the results back.
|
|
||||||
#[derive(Debug, Default)]
|
|
||||||
pub struct ProvisionOutcome {
|
|
||||||
/// ICS-probe findings from the provisioned Modbus endpoint.
|
|
||||||
pub findings: Vec<Finding>,
|
|
||||||
/// DAST scan of the provisioned web endpoint, if it ran.
|
|
||||||
pub dast: Option<DastRunResult>,
|
|
||||||
}
|
|
||||||
|
|
||||||
/// A control-logic program ready to load into a soft-PLC: the source text plus a
|
|
||||||
/// cosmetic file name (OpenPLC re-stores it under its own name).
|
|
||||||
#[derive(Debug, Clone)]
|
|
||||||
pub struct PlcProgram {
|
|
||||||
/// The original file name (for the upload form; OpenPLC renames on storage).
|
|
||||||
pub file_name: String,
|
|
||||||
/// The program source — Structured Text or PLCopen XML.
|
|
||||||
pub source: String,
|
|
||||||
}
|
|
||||||
|
|
||||||
/// 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, AgentError> {
|
|
||||||
reqwest::Client::builder()
|
|
||||||
.cookie_store(true)
|
|
||||||
.timeout(Duration::from_secs(30))
|
|
||||||
.build()
|
|
||||||
.map_err(AgentError::Http)
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Pick the control-logic program to run from an ingested PLC source tree.
|
|
||||||
///
|
|
||||||
/// OpenPLC runs one program, so we choose the best single candidate: a complete
|
|
||||||
/// Structured Text program (one carrying a `CONFIGURATION` block) is ideal;
|
|
||||||
/// failing that the largest ST file; failing that a PLCopen XML export. Returns
|
|
||||||
/// `None` when the tree holds no loadable control logic.
|
|
||||||
pub fn extract_program(root: &Path) -> Option<PlcProgram> {
|
|
||||||
let mut st: Vec<(String, String)> = Vec::new();
|
|
||||||
let mut xml: Vec<(String, String)> = Vec::new();
|
|
||||||
for entry in walkdir::WalkDir::new(root)
|
|
||||||
.into_iter()
|
|
||||||
.filter_map(Result::ok)
|
|
||||||
{
|
|
||||||
if !entry.file_type().is_file() {
|
|
||||||
continue;
|
|
||||||
}
|
|
||||||
let path = entry.path();
|
|
||||||
let ext = path
|
|
||||||
.extension()
|
|
||||||
.and_then(|e| e.to_str())
|
|
||||||
.unwrap_or("")
|
|
||||||
.to_ascii_lowercase();
|
|
||||||
let is_st = matches!(ext.as_str(), "st" | "iecst" | "scl" | "exp" | "il");
|
|
||||||
let is_xml = matches!(ext.as_str(), "xml" | "plcopen" | "project");
|
|
||||||
if !is_st && !is_xml {
|
|
||||||
continue;
|
|
||||||
}
|
|
||||||
let Ok(content) = std::fs::read_to_string(path) else {
|
|
||||||
continue;
|
|
||||||
};
|
|
||||||
let name = path
|
|
||||||
.file_name()
|
|
||||||
.and_then(|n| n.to_str())
|
|
||||||
.unwrap_or("program")
|
|
||||||
.to_string();
|
|
||||||
if is_st {
|
|
||||||
st.push((name, content));
|
|
||||||
} else if looks_like_plcopen(&content) {
|
|
||||||
xml.push((name, content));
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
if let Some((name, source)) = st.iter().find(|(_, c)| has_configuration(c)) {
|
|
||||||
return Some(PlcProgram {
|
|
||||||
file_name: name.clone(),
|
|
||||||
source: source.clone(),
|
|
||||||
});
|
|
||||||
}
|
|
||||||
if let Some((name, source)) = st.iter().max_by_key(|(_, c)| c.len()) {
|
|
||||||
return Some(PlcProgram {
|
|
||||||
file_name: name.clone(),
|
|
||||||
source: source.clone(),
|
|
||||||
});
|
|
||||||
}
|
|
||||||
xml.into_iter()
|
|
||||||
.max_by_key(|(_, c)| c.len())
|
|
||||||
.map(|(file_name, source)| PlcProgram { file_name, source })
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Whether an ST source is a complete, runnable program (has a `CONFIGURATION`).
|
|
||||||
fn has_configuration(source: &str) -> bool {
|
|
||||||
source.to_ascii_uppercase().contains("CONFIGURATION")
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Whether an XML file looks like a PLCopen project export.
|
|
||||||
fn looks_like_plcopen(source: &str) -> bool {
|
|
||||||
let lower = source.to_ascii_lowercase();
|
|
||||||
lower.contains("<project") || lower.contains("plcopen")
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Provision an ephemeral soft-PLC, load `program`, start it, probe it over
|
|
||||||
/// industrial protocols, and tear it down. Returns the ICS-probe findings.
|
|
||||||
///
|
|
||||||
/// Teardown is guaranteed: the load/probe work runs under a hard deadline
|
|
||||||
/// (`max_lifetime_secs`) and the instance is removed afterwards on every path —
|
|
||||||
/// success, error, or deadline expiry.
|
|
||||||
pub async fn provision_and_test<P: SoftPlc>(
|
|
||||||
provisioner: &P,
|
|
||||||
http: &reqwest::Client,
|
|
||||||
cfg: &PlcRuntimeConfig,
|
|
||||||
program: &PlcProgram,
|
|
||||||
target_id: &str,
|
|
||||||
) -> Result<ProvisionOutcome, AgentError> {
|
|
||||||
let handle = provisioner.provision(target_id).await?;
|
|
||||||
tracing::info!(
|
|
||||||
target_id,
|
|
||||||
instance = %handle.name,
|
|
||||||
modbus = %handle.modbus_endpoint,
|
|
||||||
"provisioned ephemeral soft-PLC"
|
|
||||||
);
|
|
||||||
|
|
||||||
let deadline = Duration::from_secs(cfg.max_lifetime_secs);
|
|
||||||
let result = tokio::time::timeout(
|
|
||||||
deadline,
|
|
||||||
run_dynamic_test(http, cfg, program, target_id, &handle),
|
|
||||||
)
|
|
||||||
.await;
|
|
||||||
|
|
||||||
// Guaranteed teardown — runs on success, error, and deadline expiry. The
|
|
||||||
// inner future is panic-free (the workspace lint bans unwrap/expect), so no
|
|
||||||
// unwind can skip this; a container leaked by an agent *crash* is swept by
|
|
||||||
// the next run's stale reaper.
|
|
||||||
provisioner.teardown(&handle).await;
|
|
||||||
|
|
||||||
match result {
|
|
||||||
Ok(inner) => inner,
|
|
||||||
Err(_) => {
|
|
||||||
tracing::warn!(
|
|
||||||
target_id,
|
|
||||||
instance = %handle.name,
|
|
||||||
"provision-and-test hit the lifetime deadline; torn down"
|
|
||||||
);
|
|
||||||
Ok(ProvisionOutcome::default())
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/// The load → start → probe → DAST body, run under the caller's deadline.
|
|
||||||
async fn run_dynamic_test(
|
|
||||||
http: &reqwest::Client,
|
|
||||||
cfg: &PlcRuntimeConfig,
|
|
||||||
program: &PlcProgram,
|
|
||||||
target_id: &str,
|
|
||||||
handle: &ProvisionedRuntime,
|
|
||||||
) -> 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?;
|
|
||||||
|
|
||||||
let compile_budget = Duration::from_secs((cfg.max_lifetime_secs / 2).clamp(20, 120));
|
|
||||||
openplc::load_and_start(
|
|
||||||
http,
|
|
||||||
&handle.webvisu_url,
|
|
||||||
&cfg.openplc_user,
|
|
||||||
cfg.openplc_password.expose_secret(),
|
|
||||||
program,
|
|
||||||
compile_budget,
|
|
||||||
)
|
|
||||||
.await?;
|
|
||||||
|
|
||||||
// Give the runtime a moment to open the Modbus/TCP server before probing.
|
|
||||||
tokio::time::sleep(Duration::from_secs(3)).await;
|
|
||||||
|
|
||||||
let probe_budget = Duration::from_secs(5);
|
|
||||||
let findings =
|
|
||||||
crate::pipeline::ics::probe_target(&handle.modbus_endpoint, target_id, probe_budget).await;
|
|
||||||
tracing::info!(
|
|
||||||
target_id,
|
|
||||||
instance = %handle.name,
|
|
||||||
found = findings.len(),
|
|
||||||
"provision-and-test probe complete"
|
|
||||||
);
|
|
||||||
|
|
||||||
// DAST the provisioned web endpoint (independently bounded so it can't eat
|
|
||||||
// the whole lifetime). On the OpenPLC substrate this is OpenPLC's own web UI,
|
|
||||||
// not a customer HMI — the CODESYS-runtime follow-up raises the fidelity —
|
|
||||||
// but it proves the deploy→run→probe→DAST loop end to end.
|
|
||||||
let dast_budget = Duration::from_secs((cfg.max_lifetime_secs / 2).clamp(20, 120));
|
|
||||||
let dast = match tokio::time::timeout(dast_budget, run_webvisu_dast(handle, target_id)).await {
|
|
||||||
Ok(d) => d,
|
|
||||||
Err(_) => {
|
|
||||||
tracing::warn!(target_id, instance = %handle.name, "provision-and-test DAST timed out");
|
|
||||||
None
|
|
||||||
}
|
|
||||||
};
|
|
||||||
|
|
||||||
Ok(ProvisionOutcome { findings, dast })
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Run a bounded DAST scan against the provisioned web endpoint and tag the
|
|
||||||
/// results with our target id. Best-effort — a DAST failure never fails the run.
|
|
||||||
async fn run_webvisu_dast(handle: &ProvisionedRuntime, target_id: &str) -> Option<DastRunResult> {
|
|
||||||
let mut dt = DastTarget::new(
|
|
||||||
"provisioned-webvisu".to_string(),
|
|
||||||
handle.webvisu_url.clone(),
|
|
||||||
DastTargetType::WebApp,
|
|
||||||
);
|
|
||||||
dt.repo_id = Some(target_id.to_string());
|
|
||||||
dt.max_crawl_depth = 2; // shallow — the instance is ephemeral
|
|
||||||
|
|
||||||
let orchestrator = compliance_dast::DastOrchestrator::new(100);
|
|
||||||
match orchestrator.run_scan(&dt, Vec::new()).await {
|
|
||||||
Ok((mut scan_run, mut findings)) => {
|
|
||||||
scan_run.target_id = target_id.to_string();
|
|
||||||
for f in &mut findings {
|
|
||||||
f.target_id = target_id.to_string();
|
|
||||||
}
|
|
||||||
tracing::info!(
|
|
||||||
target_id,
|
|
||||||
instance = %handle.name,
|
|
||||||
dast_findings = findings.len(),
|
|
||||||
"provision-and-test DAST complete"
|
|
||||||
);
|
|
||||||
Some(DastRunResult { scan_run, findings })
|
|
||||||
}
|
|
||||||
Err(e) => {
|
|
||||||
tracing::warn!(target_id, instance = %handle.name, error = %e, "provision-and-test DAST failed");
|
|
||||||
None
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
#[cfg(test)]
|
|
||||||
#[allow(clippy::expect_used, clippy::unwrap_used)]
|
|
||||||
mod tests {
|
|
||||||
use super::*;
|
|
||||||
use std::sync::atomic::{AtomicUsize, Ordering};
|
|
||||||
use std::sync::Arc;
|
|
||||||
|
|
||||||
/// A scratch dir removed on drop.
|
|
||||||
struct Scratch(std::path::PathBuf);
|
|
||||||
impl Scratch {
|
|
||||||
fn new() -> Self {
|
|
||||||
let p = std::env::temp_dir().join(format!("cs-plc-rt-{}", uuid::Uuid::new_v4()));
|
|
||||||
std::fs::create_dir_all(&p).expect("mkdir");
|
|
||||||
Self(p)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
impl Drop for Scratch {
|
|
||||||
fn drop(&mut self) {
|
|
||||||
let _ = std::fs::remove_dir_all(&self.0);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn extract_prefers_a_complete_st_program() {
|
|
||||||
let s = Scratch::new();
|
|
||||||
std::fs::write(s.0.join("fragment.st"), "PROGRAM P\nEND_PROGRAM\n").expect("w");
|
|
||||||
std::fs::write(
|
|
||||||
s.0.join("full.st"),
|
|
||||||
"PROGRAM Main\nEND_PROGRAM\nCONFIGURATION Config0\n RESOURCE R\nEND_CONFIGURATION\n",
|
|
||||||
)
|
|
||||||
.expect("w");
|
|
||||||
let prog = extract_program(&s.0).expect("program");
|
|
||||||
assert_eq!(prog.file_name, "full.st");
|
|
||||||
assert!(prog.source.contains("CONFIGURATION"));
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn extract_falls_back_to_largest_st_then_plcopen() {
|
|
||||||
let s = Scratch::new();
|
|
||||||
std::fs::write(s.0.join("small.st"), "PROGRAM A\nEND_PROGRAM\n").expect("w");
|
|
||||||
std::fs::write(
|
|
||||||
s.0.join("big.st"),
|
|
||||||
"PROGRAM B\nVAR x : INT; y : INT; z : INT; END_VAR\nEND_PROGRAM\n",
|
|
||||||
)
|
|
||||||
.expect("w");
|
|
||||||
let prog = extract_program(&s.0).expect("program");
|
|
||||||
assert_eq!(
|
|
||||||
prog.file_name, "big.st",
|
|
||||||
"largest ST wins when none complete"
|
|
||||||
);
|
|
||||||
|
|
||||||
// Only a PLCopen XML present.
|
|
||||||
let s2 = Scratch::new();
|
|
||||||
std::fs::write(
|
|
||||||
s2.0.join("proj.xml"),
|
|
||||||
"<?xml version='1.0'?><project xmlns='http://www.plcopen.org/xml/tc6_0201'><pou/></project>",
|
|
||||||
)
|
|
||||||
.expect("w");
|
|
||||||
let prog2 = extract_program(&s2.0).expect("program");
|
|
||||||
assert_eq!(prog2.file_name, "proj.xml");
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn extract_returns_none_without_control_logic() {
|
|
||||||
let s = Scratch::new();
|
|
||||||
std::fs::write(s.0.join("readme.md"), "# not a plc program").expect("w");
|
|
||||||
std::fs::write(s.0.join("data.xml"), "<config><db/></config>").expect("w");
|
|
||||||
assert!(extract_program(&s.0).is_none());
|
|
||||||
}
|
|
||||||
|
|
||||||
/// A fake provisioner recording provision/teardown calls, for lifecycle tests.
|
|
||||||
struct FakeSoftPlc {
|
|
||||||
provisions: Arc<AtomicUsize>,
|
|
||||||
teardowns: Arc<AtomicUsize>,
|
|
||||||
fail_provision: bool,
|
|
||||||
}
|
|
||||||
|
|
||||||
impl SoftPlc for FakeSoftPlc {
|
|
||||||
async fn provision(&self, _target_id: &str) -> Result<ProvisionedRuntime, AgentError> {
|
|
||||||
self.provisions.fetch_add(1, Ordering::SeqCst);
|
|
||||||
if self.fail_provision {
|
|
||||||
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.
|
|
||||||
Ok(ProvisionedRuntime {
|
|
||||||
name: "fake-plc".into(),
|
|
||||||
modbus_endpoint: "fake-plc:502".into(),
|
|
||||||
webvisu_url: "http://fake-plc.invalid:8080".into(),
|
|
||||||
})
|
|
||||||
}
|
|
||||||
async fn teardown(&self, _handle: &ProvisionedRuntime) {
|
|
||||||
self.teardowns.fetch_add(1, Ordering::SeqCst);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
fn short_cfg() -> PlcRuntimeConfig {
|
|
||||||
PlcRuntimeConfig {
|
|
||||||
enabled: true,
|
|
||||||
max_lifetime_secs: 1, // keep the deadline path fast
|
|
||||||
..PlcRuntimeConfig::default()
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
#[tokio::test]
|
|
||||||
async fn teardown_runs_even_when_the_test_never_completes() {
|
|
||||||
let provisions = Arc::new(AtomicUsize::new(0));
|
|
||||||
let teardowns = Arc::new(AtomicUsize::new(0));
|
|
||||||
let fake = FakeSoftPlc {
|
|
||||||
provisions: provisions.clone(),
|
|
||||||
teardowns: teardowns.clone(),
|
|
||||||
fail_provision: false,
|
|
||||||
};
|
|
||||||
let http = http_client().expect("client");
|
|
||||||
let prog = PlcProgram {
|
|
||||||
file_name: "p.st".into(),
|
|
||||||
source: "PROGRAM P\nEND_PROGRAM\n".into(),
|
|
||||||
};
|
|
||||||
let out = provision_and_test(&fake, &http, &short_cfg(), &prog, "t1")
|
|
||||||
.await
|
|
||||||
.expect("ok on deadline");
|
|
||||||
assert!(out.findings.is_empty(), "deadline path yields no findings");
|
|
||||||
assert!(out.dast.is_none(), "deadline path runs no DAST");
|
|
||||||
assert_eq!(provisions.load(Ordering::SeqCst), 1);
|
|
||||||
assert_eq!(teardowns.load(Ordering::SeqCst), 1, "teardown must run");
|
|
||||||
}
|
|
||||||
|
|
||||||
#[tokio::test]
|
|
||||||
async fn provision_failure_propagates_and_skips_teardown() {
|
|
||||||
let provisions = Arc::new(AtomicUsize::new(0));
|
|
||||||
let teardowns = Arc::new(AtomicUsize::new(0));
|
|
||||||
let fake = FakeSoftPlc {
|
|
||||||
provisions: provisions.clone(),
|
|
||||||
teardowns: teardowns.clone(),
|
|
||||||
fail_provision: true,
|
|
||||||
};
|
|
||||||
let http = http_client().expect("client");
|
|
||||||
let prog = PlcProgram {
|
|
||||||
file_name: "p.st".into(),
|
|
||||||
source: String::new(),
|
|
||||||
};
|
|
||||||
let err = provision_and_test(&fake, &http, &short_cfg(), &prog, "t1").await;
|
|
||||||
assert!(err.is_err(), "provision failure propagates");
|
|
||||||
assert_eq!(provisions.load(Ordering::SeqCst), 1);
|
|
||||||
assert_eq!(
|
|
||||||
teardowns.load(Ordering::SeqCst),
|
|
||||||
0,
|
|
||||||
"nothing to tear down when provisioning failed"
|
|
||||||
);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -1,254 +0,0 @@
|
|||||||
//! Loading a control-logic program into a provisioned OpenPLC (#183, sub-task 2).
|
|
||||||
//!
|
|
||||||
//! Drives the OpenPLC v3 web UI over HTTP to turn a static control-logic artifact
|
|
||||||
//! into a *running* PLC: log in, upload the program, save it, compile it (MatIEC),
|
|
||||||
//! and start the runtime — at which point OpenPLC opens its Modbus/TCP server on
|
|
||||||
//! 502 and the ICS probe has something to talk to. The endpoint sequence mirrors
|
|
||||||
//! the OpenPLC web UI: `POST /login` → `POST /upload-program` (which hands back a
|
|
||||||
//! server-assigned `prog_file`) → `POST /upload-program-action` →
|
|
||||||
//! `GET /compile-program?file=<prog_file>` → `GET /start_plc`.
|
|
||||||
|
|
||||||
use std::time::Duration;
|
|
||||||
|
|
||||||
use crate::error::AgentError;
|
|
||||||
|
|
||||||
use super::PlcProgram;
|
|
||||||
|
|
||||||
/// Default OpenPLC program name/description recorded in its UI.
|
|
||||||
const PROG_NAME: &str = "certifai-provisioned";
|
|
||||||
const PROG_DESCR: &str = "Uploaded by the Certifai provision-and-test scan";
|
|
||||||
|
|
||||||
/// Poll interval while waiting for readiness / compilation.
|
|
||||||
const POLL_INTERVAL: Duration = Duration::from_secs(2);
|
|
||||||
|
|
||||||
/// Wait until the OpenPLC web UI answers (any non-5xx reply to `/login`), or the
|
|
||||||
/// budget elapses. A freshly-started container needs a few seconds to boot.
|
|
||||||
pub async fn wait_ready(
|
|
||||||
http: &reqwest::Client,
|
|
||||||
base_url: &str,
|
|
||||||
budget: Duration,
|
|
||||||
) -> Result<(), AgentError> {
|
|
||||||
let login = format!("{base_url}/login");
|
|
||||||
let outcome = tokio::time::timeout(budget, async {
|
|
||||||
loop {
|
|
||||||
if let Ok(resp) = http.get(&login).send().await {
|
|
||||||
if !resp.status().is_server_error() {
|
|
||||||
return;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
tokio::time::sleep(POLL_INTERVAL).await;
|
|
||||||
}
|
|
||||||
})
|
|
||||||
.await;
|
|
||||||
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
|
|
||||||
/// OpenPLC Modbus/TCP server is listening on 502.
|
|
||||||
pub async fn load_and_start(
|
|
||||||
http: &reqwest::Client,
|
|
||||||
base_url: &str,
|
|
||||||
user: &str,
|
|
||||||
password: &str,
|
|
||||||
program: &PlcProgram,
|
|
||||||
compile_budget: Duration,
|
|
||||||
) -> 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?;
|
|
||||||
compile(http, base_url, &prog_file, compile_budget).await?;
|
|
||||||
start(http, base_url).await?;
|
|
||||||
Ok(())
|
|
||||||
}
|
|
||||||
|
|
||||||
/// `POST /login` — establishes the session cookie (the client must have a cookie
|
|
||||||
/// store; see the provision-and-test entry point).
|
|
||||||
async fn login(
|
|
||||||
http: &reqwest::Client,
|
|
||||||
base_url: &str,
|
|
||||||
user: &str,
|
|
||||||
password: &str,
|
|
||||||
) -> 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(AgentError::Other(format!(
|
|
||||||
"OpenPLC login failed: HTTP {}",
|
|
||||||
resp.status()
|
|
||||||
)));
|
|
||||||
}
|
|
||||||
Ok(())
|
|
||||||
}
|
|
||||||
|
|
||||||
/// `POST /upload-program` (multipart `file`) — OpenPLC stores the program under a
|
|
||||||
/// server-assigned name and returns it in a hidden `prog_file` form field, which
|
|
||||||
/// we parse out for the follow-up save/compile steps.
|
|
||||||
async fn upload_program(
|
|
||||||
http: &reqwest::Client,
|
|
||||||
base_url: &str,
|
|
||||||
program: &PlcProgram,
|
|
||||||
) -> Result<String, AgentError> {
|
|
||||||
let part = reqwest::multipart::Part::text(program.source.clone())
|
|
||||||
.file_name(program.file_name.clone())
|
|
||||||
.mime_str("application/octet-stream")?;
|
|
||||||
let form = reqwest::multipart::Form::new().part("file", part);
|
|
||||||
let resp = http
|
|
||||||
.post(format!("{base_url}/upload-program"))
|
|
||||||
.multipart(form)
|
|
||||||
.send()
|
|
||||||
.await?;
|
|
||||||
let html = resp.text().await?;
|
|
||||||
parse_prog_file(&html).ok_or_else(|| {
|
|
||||||
AgentError::Other("OpenPLC upload did not return a prog_file handle".to_string())
|
|
||||||
})
|
|
||||||
}
|
|
||||||
|
|
||||||
/// `POST /upload-program-action` — records the uploaded program in OpenPLC's
|
|
||||||
/// program list. `epoch_time` must be close to the server's clock (OpenPLC
|
|
||||||
/// rejects stale timestamps), so we send the current time.
|
|
||||||
async fn save_program(
|
|
||||||
http: &reqwest::Client,
|
|
||||||
base_url: &str,
|
|
||||||
prog_file: &str,
|
|
||||||
) -> Result<(), AgentError> {
|
|
||||||
let epoch = std::time::SystemTime::now()
|
|
||||||
.duration_since(std::time::UNIX_EPOCH)
|
|
||||||
.map(|d| d.as_secs())
|
|
||||||
.unwrap_or(0)
|
|
||||||
.to_string();
|
|
||||||
let resp = http
|
|
||||||
.post(format!("{base_url}/upload-program-action"))
|
|
||||||
.form(&[
|
|
||||||
("prog_name", PROG_NAME),
|
|
||||||
("prog_descr", PROG_DESCR),
|
|
||||||
("prog_file", prog_file),
|
|
||||||
("epoch_time", &epoch),
|
|
||||||
])
|
|
||||||
.send()
|
|
||||||
.await?;
|
|
||||||
if resp.status().is_server_error() {
|
|
||||||
return Err(AgentError::Other(format!(
|
|
||||||
"OpenPLC save-program failed: HTTP {}",
|
|
||||||
resp.status()
|
|
||||||
)));
|
|
||||||
}
|
|
||||||
Ok(())
|
|
||||||
}
|
|
||||||
|
|
||||||
/// `GET /compile-program?file=<prog_file>` then poll `/compilation-logs` until
|
|
||||||
/// MatIEC reports it finished (or the budget elapses). Errors if compilation
|
|
||||||
/// finishes with errors — a program that won't compile can't be started.
|
|
||||||
async fn compile(
|
|
||||||
http: &reqwest::Client,
|
|
||||||
base_url: &str,
|
|
||||||
prog_file: &str,
|
|
||||||
budget: Duration,
|
|
||||||
) -> Result<(), AgentError> {
|
|
||||||
http.get(format!("{base_url}/compile-program"))
|
|
||||||
.query(&[("file", prog_file)])
|
|
||||||
.send()
|
|
||||||
.await?;
|
|
||||||
|
|
||||||
let logs_url = format!("{base_url}/compilation-logs");
|
|
||||||
let outcome = tokio::time::timeout(budget, async {
|
|
||||||
loop {
|
|
||||||
if let Ok(resp) = http.get(&logs_url).send().await {
|
|
||||||
if let Ok(text) = resp.text().await {
|
|
||||||
if compilation_finished(&text) {
|
|
||||||
return !compilation_failed(&text);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
tokio::time::sleep(POLL_INTERVAL).await;
|
|
||||||
}
|
|
||||||
})
|
|
||||||
.await;
|
|
||||||
match outcome {
|
|
||||||
Ok(true) => Ok(()),
|
|
||||||
Ok(false) => Err(AgentError::Other(
|
|
||||||
"OpenPLC compilation finished with errors".to_string(),
|
|
||||||
)),
|
|
||||||
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<(), AgentError> {
|
|
||||||
let resp = http.get(format!("{base_url}/start_plc")).send().await?;
|
|
||||||
if resp.status().is_server_error() {
|
|
||||||
return Err(AgentError::Other(format!(
|
|
||||||
"OpenPLC start_plc failed: HTTP {}",
|
|
||||||
resp.status()
|
|
||||||
)));
|
|
||||||
}
|
|
||||||
Ok(())
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Extract the server-assigned `prog_file` from the `/upload-program` response,
|
|
||||||
/// which embeds it in a hidden input. Attribute order varies, so accept both
|
|
||||||
/// `value=… name='prog_file'` and `name='prog_file' … value=…`.
|
|
||||||
fn parse_prog_file(html: &str) -> Option<String> {
|
|
||||||
// The OpenPLC template renders `value='<name>.st' id='prog_file'
|
|
||||||
// name='prog_file'`. Match the value bound to that input, either order.
|
|
||||||
let value_then_name =
|
|
||||||
regex::Regex::new(r#"(?is)value=['"]([^'"]+)['"][^>]*name=['"]prog_file['"]"#).ok()?;
|
|
||||||
if let Some(c) = value_then_name.captures(html) {
|
|
||||||
return c.get(1).map(|m| m.as_str().to_string());
|
|
||||||
}
|
|
||||||
let name_then_value =
|
|
||||||
regex::Regex::new(r#"(?is)name=['"]prog_file['"][^>]*value=['"]([^'"]+)['"]"#).ok()?;
|
|
||||||
name_then_value
|
|
||||||
.captures(html)
|
|
||||||
.and_then(|c| c.get(1))
|
|
||||||
.map(|m| m.as_str().to_string())
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Whether the MatIEC compilation log shows the build has finished (either way).
|
|
||||||
fn compilation_finished(log: &str) -> bool {
|
|
||||||
log.contains("Compilation finished")
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Whether a finished compilation ended in failure.
|
|
||||||
fn compilation_failed(log: &str) -> bool {
|
|
||||||
log.contains("Compilation finished with errors")
|
|
||||||
}
|
|
||||||
|
|
||||||
#[cfg(test)]
|
|
||||||
#[allow(clippy::expect_used, clippy::unwrap_used)]
|
|
||||||
mod tests {
|
|
||||||
use super::*;
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn parses_prog_file_value_then_name() {
|
|
||||||
let html = "<form><input type='hidden' value='483927.st' id='prog_file' \
|
|
||||||
name='prog_file'/></form>";
|
|
||||||
assert_eq!(parse_prog_file(html), Some("483927.st".to_string()));
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn parses_prog_file_name_then_value() {
|
|
||||||
let html = r#"<input name="prog_file" id="prog_file" value="12.st" />"#;
|
|
||||||
assert_eq!(parse_prog_file(html), Some("12.st".to_string()));
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn parse_prog_file_none_when_absent() {
|
|
||||||
assert_eq!(parse_prog_file("<html>no form here</html>"), None);
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn compilation_predicates() {
|
|
||||||
assert!(!compilation_finished("Compiling..."));
|
|
||||||
assert!(compilation_finished(
|
|
||||||
"...\nCompilation finished successfully!\n"
|
|
||||||
));
|
|
||||||
assert!(compilation_finished("Compilation finished with errors!"));
|
|
||||||
assert!(compilation_failed("Compilation finished with errors!"));
|
|
||||||
assert!(!compilation_failed("Compilation finished successfully!"));
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -1,306 +0,0 @@
|
|||||||
//! Ephemeral soft-PLC container lifecycle (#183, sub-task 1 + 5).
|
|
||||||
//!
|
|
||||||
//! Provisions a throwaway OpenPLC container per scan, isolated on the agent's own
|
|
||||||
//! Docker network with hard resource caps and **no host port exposure**, then
|
|
||||||
//! guarantees teardown. The container is reachable in-cluster only, by its name
|
|
||||||
//! (the shared user-defined network's embedded DNS resolves it); it is never
|
|
||||||
//! published to the host.
|
|
||||||
//!
|
|
||||||
//! The `docker` argv is produced by pure functions so provisioning is unit-tested
|
|
||||||
//! without a Docker daemon — only the thin [`run_docker`] wrapper touches the OS.
|
|
||||||
//! It requires the agent's runtime to have Docker access (a socket mount), which
|
|
||||||
//! is why the whole path is gated behind [`PlcRuntimeConfig::enabled`].
|
|
||||||
|
|
||||||
use std::time::{SystemTime, UNIX_EPOCH};
|
|
||||||
|
|
||||||
use compliance_core::PlcRuntimeConfig;
|
|
||||||
|
|
||||||
use crate::error::AgentError;
|
|
||||||
|
|
||||||
/// The Modbus/TCP port an OpenPLC instance opens once a program is running.
|
|
||||||
const MODBUS_PORT: u16 = 502;
|
|
||||||
/// The OpenPLC web-UI / WebVisu port.
|
|
||||||
const WEBVISU_PORT: u16 = 8080;
|
|
||||||
|
|
||||||
/// Label key marking a container as an ephemeral PLC runtime we own.
|
|
||||||
const OWNER_LABEL_KEY: &str = "certifai.ephemeral";
|
|
||||||
/// Label value for our ephemeral PLC runtimes.
|
|
||||||
const OWNER_LABEL_VALUE: &str = "plc-runtime";
|
|
||||||
|
|
||||||
/// A running ephemeral soft-PLC instance. Reachable in-cluster by `name`.
|
|
||||||
#[derive(Debug, Clone)]
|
|
||||||
pub struct ProvisionedRuntime {
|
|
||||||
/// The container name — also its in-network DNS alias.
|
|
||||||
pub name: String,
|
|
||||||
/// `name:502` — the Modbus/TCP endpoint the ICS probe targets.
|
|
||||||
pub modbus_endpoint: String,
|
|
||||||
/// `http://name:8080` — the WebVisu / OpenPLC web UI.
|
|
||||||
pub webvisu_url: String,
|
|
||||||
}
|
|
||||||
|
|
||||||
/// A source of ephemeral soft-PLC instances. Abstracted so the provision-and-test
|
|
||||||
/// orchestration is unit-testable with a fake that never touches Docker.
|
|
||||||
pub trait SoftPlc {
|
|
||||||
/// Start a fresh instance for a target and return its handle.
|
|
||||||
fn provision(
|
|
||||||
&self,
|
|
||||||
target_id: &str,
|
|
||||||
) -> 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)
|
|
||||||
-> impl std::future::Future<Output = ()> + Send;
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Provisions OpenPLC instances by shelling out to the Docker CLI.
|
|
||||||
pub struct DockerSoftPlc {
|
|
||||||
cfg: PlcRuntimeConfig,
|
|
||||||
}
|
|
||||||
|
|
||||||
impl DockerSoftPlc {
|
|
||||||
/// Build a provisioner from the PLC-runtime config.
|
|
||||||
pub fn new(cfg: PlcRuntimeConfig) -> Self {
|
|
||||||
Self { cfg }
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
impl SoftPlc for DockerSoftPlc {
|
|
||||||
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.
|
|
||||||
reap_stale(&self.cfg, now_epoch()).await;
|
|
||||||
|
|
||||||
let name = instance_name(target_id, now_epoch(), &random_suffix());
|
|
||||||
let args = run_args(&self.cfg, &name, target_id);
|
|
||||||
let out = run_docker(&args).await?;
|
|
||||||
if !out.status.success() {
|
|
||||||
return Err(AgentError::Other(format!(
|
|
||||||
"docker run for soft-PLC {name} failed: {}",
|
|
||||||
String::from_utf8_lossy(&out.stderr).trim()
|
|
||||||
)));
|
|
||||||
}
|
|
||||||
Ok(ProvisionedRuntime {
|
|
||||||
modbus_endpoint: format!("{name}:{MODBUS_PORT}"),
|
|
||||||
webvisu_url: format!("http://{name}:{WEBVISU_PORT}"),
|
|
||||||
name,
|
|
||||||
})
|
|
||||||
}
|
|
||||||
|
|
||||||
async fn teardown(&self, handle: &ProvisionedRuntime) {
|
|
||||||
match run_docker(&rm_args(&handle.name)).await {
|
|
||||||
Ok(out) if out.status.success() => {
|
|
||||||
tracing::info!(instance = %handle.name, "soft-PLC instance torn down");
|
|
||||||
}
|
|
||||||
Ok(out) => tracing::warn!(
|
|
||||||
instance = %handle.name,
|
|
||||||
"soft-PLC teardown non-zero exit: {}",
|
|
||||||
String::from_utf8_lossy(&out.stderr).trim()
|
|
||||||
),
|
|
||||||
Err(e) => {
|
|
||||||
tracing::warn!(instance = %handle.name, error = %e, "soft-PLC teardown failed")
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Seconds since the Unix epoch (0 if the clock is before 1970, which never
|
|
||||||
/// happens in practice).
|
|
||||||
fn now_epoch() -> u64 {
|
|
||||||
SystemTime::now()
|
|
||||||
.duration_since(UNIX_EPOCH)
|
|
||||||
.map(|d| d.as_secs())
|
|
||||||
.unwrap_or(0)
|
|
||||||
}
|
|
||||||
|
|
||||||
/// A short random, docker-name-safe suffix.
|
|
||||||
fn random_suffix() -> String {
|
|
||||||
uuid::Uuid::new_v4().simple().to_string()
|
|
||||||
}
|
|
||||||
|
|
||||||
/// A unique, docker-safe container name that encodes the creation epoch (for the
|
|
||||||
/// stale reaper) and the target it belongs to. Shape:
|
|
||||||
/// `certifai-plc-<epoch>-<target12>-<rand6>`.
|
|
||||||
fn instance_name(target_id: &str, epoch: u64, rand: &str) -> String {
|
|
||||||
let short: String = target_id
|
|
||||||
.chars()
|
|
||||||
.filter(char::is_ascii_alphanumeric)
|
|
||||||
.take(12)
|
|
||||||
.collect();
|
|
||||||
let rand: String = rand
|
|
||||||
.chars()
|
|
||||||
.filter(char::is_ascii_alphanumeric)
|
|
||||||
.take(6)
|
|
||||||
.collect();
|
|
||||||
format!("certifai-plc-{epoch}-{short}-{rand}")
|
|
||||||
}
|
|
||||||
|
|
||||||
/// The creation epoch encoded in an instance name, if it is one of ours.
|
|
||||||
fn parse_epoch(name: &str) -> Option<u64> {
|
|
||||||
name.strip_prefix("certifai-plc-")?
|
|
||||||
.split('-')
|
|
||||||
.next()?
|
|
||||||
.parse()
|
|
||||||
.ok()
|
|
||||||
}
|
|
||||||
|
|
||||||
/// The `docker run` argv for an ephemeral soft-PLC: detached, joined to the
|
|
||||||
/// agent's network, resource-capped, hardened, labelled for reaping, and — by
|
|
||||||
/// omitting any `-p` — never published to the host.
|
|
||||||
fn run_args(cfg: &PlcRuntimeConfig, name: &str, target_id: &str) -> Vec<String> {
|
|
||||||
vec![
|
|
||||||
"run".into(),
|
|
||||||
"-d".into(),
|
|
||||||
"--name".into(),
|
|
||||||
name.into(),
|
|
||||||
"--network".into(),
|
|
||||||
cfg.network.clone(),
|
|
||||||
"--memory".into(),
|
|
||||||
cfg.memory.clone(),
|
|
||||||
"--cpus".into(),
|
|
||||||
cfg.cpus.clone(),
|
|
||||||
"--pids-limit".into(),
|
|
||||||
"512".into(),
|
|
||||||
"--security-opt".into(),
|
|
||||||
"no-new-privileges".into(),
|
|
||||||
"--stop-timeout".into(),
|
|
||||||
"5".into(),
|
|
||||||
"--label".into(),
|
|
||||||
format!("{OWNER_LABEL_KEY}={OWNER_LABEL_VALUE}"),
|
|
||||||
"--label".into(),
|
|
||||||
format!("certifai.target={target_id}"),
|
|
||||||
cfg.image.clone(),
|
|
||||||
]
|
|
||||||
}
|
|
||||||
|
|
||||||
/// The `docker rm -f` argv that stops and removes an instance.
|
|
||||||
fn rm_args(name: &str) -> Vec<String> {
|
|
||||||
vec!["rm".into(), "-f".into(), name.into()]
|
|
||||||
}
|
|
||||||
|
|
||||||
/// The `docker ps` argv listing the names of every ephemeral PLC container we own.
|
|
||||||
fn reap_list_args() -> Vec<String> {
|
|
||||||
vec![
|
|
||||||
"ps".into(),
|
|
||||||
"-a".into(),
|
|
||||||
"--filter".into(),
|
|
||||||
format!("label={OWNER_LABEL_KEY}={OWNER_LABEL_VALUE}"),
|
|
||||||
"--format".into(),
|
|
||||||
"{{.Names}}".into(),
|
|
||||||
]
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Remove any ephemeral PLC container older than twice the configured max
|
|
||||||
/// lifetime — i.e. one a crashed run leaked. The generous threshold guarantees a
|
|
||||||
/// container from a *live* run (still within its own deadline) is never swept.
|
|
||||||
/// Best-effort: any Docker error (e.g. no daemon) is ignored.
|
|
||||||
async fn reap_stale(cfg: &PlcRuntimeConfig, now: u64) {
|
|
||||||
let cutoff = cfg.max_lifetime_secs.saturating_mul(2);
|
|
||||||
let Ok(out) = run_docker(&reap_list_args()).await else {
|
|
||||||
return;
|
|
||||||
};
|
|
||||||
if !out.status.success() {
|
|
||||||
return;
|
|
||||||
}
|
|
||||||
let names = String::from_utf8_lossy(&out.stdout);
|
|
||||||
for name in names.lines().map(str::trim).filter(|n| !n.is_empty()) {
|
|
||||||
let Some(epoch) = parse_epoch(name) else {
|
|
||||||
continue;
|
|
||||||
};
|
|
||||||
if now.saturating_sub(epoch) > cutoff {
|
|
||||||
tracing::warn!(instance = %name, "reaping stale soft-PLC instance");
|
|
||||||
let _ = run_docker(&rm_args(name)).await;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Run a `docker` subcommand, capturing its output.
|
|
||||||
async fn run_docker(args: &[String]) -> Result<std::process::Output, AgentError> {
|
|
||||||
tokio::process::Command::new("docker")
|
|
||||||
.args(args)
|
|
||||||
.output()
|
|
||||||
.await
|
|
||||||
.map_err(AgentError::Io)
|
|
||||||
}
|
|
||||||
|
|
||||||
#[cfg(test)]
|
|
||||||
#[allow(clippy::expect_used, clippy::unwrap_used)]
|
|
||||||
mod tests {
|
|
||||||
use super::*;
|
|
||||||
|
|
||||||
fn cfg() -> PlcRuntimeConfig {
|
|
||||||
PlcRuntimeConfig {
|
|
||||||
enabled: true,
|
|
||||||
image: "registry.example.com/openplc:latest".into(),
|
|
||||||
network: "certifai".into(),
|
|
||||||
memory: "512m".into(),
|
|
||||||
cpus: "0.5".into(),
|
|
||||||
max_lifetime_secs: 180,
|
|
||||||
..PlcRuntimeConfig::default()
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn instance_name_is_unique_docker_safe_and_reaper_parseable() {
|
|
||||||
let a = instance_name("64f0aabbccddeeff00112233", 1_700_000_000, "abcdef123456");
|
|
||||||
assert_eq!(a, "certifai-plc-1700000000-64f0aabbccdd-abcdef");
|
|
||||||
assert_eq!(parse_epoch(&a), Some(1_700_000_000));
|
|
||||||
// Docker names: only [A-Za-z0-9_.-].
|
|
||||||
assert!(a
|
|
||||||
.chars()
|
|
||||||
.all(|c| c.is_ascii_alphanumeric() || matches!(c, '_' | '.' | '-')));
|
|
||||||
// A different random suffix yields a different name for the same target.
|
|
||||||
let b = instance_name("64f0aabbccddeeff00112233", 1_700_000_000, "zzzzzz999999");
|
|
||||||
assert_ne!(a, b);
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn parse_epoch_rejects_foreign_names() {
|
|
||||||
assert_eq!(parse_epoch("some-other-container"), None);
|
|
||||||
assert_eq!(parse_epoch("certifai-plc-notanumber-x"), None);
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn run_args_cap_resources_harden_label_and_never_publish_a_port() {
|
|
||||||
let args = run_args(&cfg(), "certifai-plc-1-t-r", "target-123");
|
|
||||||
// No host port publishing.
|
|
||||||
assert!(!args.iter().any(|a| a == "-p" || a == "--publish"));
|
|
||||||
// Detached.
|
|
||||||
assert!(args.contains(&"-d".to_string()));
|
|
||||||
// Joined to the agent's own network.
|
|
||||||
let net = args.iter().position(|a| a == "--network").expect("network");
|
|
||||||
assert_eq!(args[net + 1], "certifai");
|
|
||||||
// Resource caps.
|
|
||||||
let mem = args.iter().position(|a| a == "--memory").expect("memory");
|
|
||||||
assert_eq!(args[mem + 1], "512m");
|
|
||||||
let cpu = args.iter().position(|a| a == "--cpus").expect("cpus");
|
|
||||||
assert_eq!(args[cpu + 1], "0.5");
|
|
||||||
assert!(args.iter().any(|a| a == "--pids-limit"));
|
|
||||||
// Hardening.
|
|
||||||
let so = args
|
|
||||||
.iter()
|
|
||||||
.position(|a| a == "--security-opt")
|
|
||||||
.expect("secopt");
|
|
||||||
assert_eq!(args[so + 1], "no-new-privileges");
|
|
||||||
// Ownership + target labels for reaping / attribution.
|
|
||||||
assert!(args.contains(&"certifai.ephemeral=plc-runtime".to_string()));
|
|
||||||
assert!(args.contains(&"certifai.target=target-123".to_string()));
|
|
||||||
// Image is last.
|
|
||||||
assert_eq!(
|
|
||||||
args.last().map(String::as_str),
|
|
||||||
Some("registry.example.com/openplc:latest")
|
|
||||||
);
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn rm_args_force_remove() {
|
|
||||||
assert_eq!(rm_args("x"), vec!["rm", "-f", "x"]);
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn reap_list_filters_by_owner_label() {
|
|
||||||
let args = reap_list_args();
|
|
||||||
assert!(args.contains(&"label=certifai.ephemeral=plc-runtime".to_string()));
|
|
||||||
assert!(args.contains(&"{{.Names}}".to_string()));
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -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,44 @@ 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,
|
||||||
|
};
|
||||||
|
|
||||||
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();
|
|
||||||
}
|
|
||||||
@@ -50,7 +50,3 @@ axum = { version = "0.8", optional = true }
|
|||||||
jsonwebtoken = { version = "9", optional = true }
|
jsonwebtoken = { version = "9", optional = true }
|
||||||
reqwest = { workspace = true, optional = true }
|
reqwest = { workspace = true, optional = true }
|
||||||
tokio = { workspace = true, optional = true }
|
tokio = { workspace = true, optional = true }
|
||||||
|
|
||||||
[dev-dependencies]
|
|
||||||
# Parse the declarative TOML job specs in the Werkbank contract tests.
|
|
||||||
toml = "0.8"
|
|
||||||
|
|||||||
@@ -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.
|
||||||
|
|||||||
@@ -49,62 +49,6 @@ pub struct AgentConfig {
|
|||||||
/// of tenants to iterate. When `None` or unreachable, scheduler
|
/// of tenants to iterate. When `None` or unreachable, scheduler
|
||||||
/// falls back to `SCHEDULER_TENANT_IDS` env (M7.2-C).
|
/// falls back to `SCHEDULER_TENANT_IDS` env (M7.2-C).
|
||||||
pub tenant_registry_url: Option<String>,
|
pub tenant_registry_url: Option<String>,
|
||||||
/// Ephemeral soft-PLC provisioning for dynamic PLC testing (#183). Off by
|
|
||||||
/// 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).
|
|
||||||
///
|
|
||||||
/// When a PLC/SPS target ships control logic but no reachable live device, the
|
|
||||||
/// agent can instantiate that logic itself: spin up a throwaway soft-PLC
|
|
||||||
/// (OpenPLC) container in-cluster, load the program, start the runtime, probe it
|
|
||||||
/// over industrial protocols, then tear it down. This struct carries the knobs
|
|
||||||
/// for that container's lifecycle and the OpenPLC web-UI credentials used to
|
|
||||||
/// upload the program.
|
|
||||||
#[derive(Clone, Debug)]
|
|
||||||
pub struct PlcRuntimeConfig {
|
|
||||||
/// Master switch. Provision-and-test does nothing unless this is set — it
|
|
||||||
/// shells out to `docker`, which requires the agent container to have Docker
|
|
||||||
/// access (socket mount), an explicit deployment decision.
|
|
||||||
pub enabled: bool,
|
|
||||||
/// Container image for the ephemeral soft-PLC (OpenPLC).
|
|
||||||
pub image: String,
|
|
||||||
/// Docker network the instance joins. Must be the agent's own network so it
|
|
||||||
/// is reachable in-cluster by container name and never published to the host.
|
|
||||||
pub network: String,
|
|
||||||
/// Memory cap passed to `docker run --memory` (e.g. `512m`).
|
|
||||||
pub memory: String,
|
|
||||||
/// CPU cap passed to `docker run --cpus` (e.g. `0.5`).
|
|
||||||
pub cpus: String,
|
|
||||||
/// Hard ceiling on a provisioned instance's lifetime. Teardown is guaranteed
|
|
||||||
/// no later than this even if a load/probe step hangs.
|
|
||||||
pub max_lifetime_secs: u64,
|
|
||||||
/// OpenPLC web-UI username for the program upload (image default `openplc`).
|
|
||||||
pub openplc_user: String,
|
|
||||||
/// OpenPLC web-UI password (image default `openplc`).
|
|
||||||
pub openplc_password: SecretString,
|
|
||||||
}
|
|
||||||
|
|
||||||
impl Default for PlcRuntimeConfig {
|
|
||||||
fn default() -> Self {
|
|
||||||
Self {
|
|
||||||
enabled: false,
|
|
||||||
image: "registry.meghsakha.com/openplc:latest".to_string(),
|
|
||||||
network: "certifai".to_string(),
|
|
||||||
memory: "512m".to_string(),
|
|
||||||
cpus: "0.5".to_string(),
|
|
||||||
max_lifetime_secs: 180,
|
|
||||||
openplc_user: "openplc".to_string(),
|
|
||||||
openplc_password: SecretString::from("openplc".to_string()),
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
#[derive(Clone, Debug, Serialize, Deserialize)]
|
#[derive(Clone, Debug, Serialize, Deserialize)]
|
||||||
|
|||||||
@@ -13,6 +13,6 @@ pub mod auth;
|
|||||||
#[cfg(feature = "axum")]
|
#[cfg(feature = "axum")]
|
||||||
pub mod tenant_ctx;
|
pub mod tenant_ctx;
|
||||||
|
|
||||||
pub use config::{AgentConfig, DashboardConfig, PlcRuntimeConfig};
|
pub use config::{AgentConfig, DashboardConfig};
|
||||||
pub use error::CoreError;
|
pub use error::CoreError;
|
||||||
pub use tenant::{OrgRole, TenantContext, TenantStatus};
|
pub use tenant::{OrgRole, TenantContext, TenantStatus};
|
||||||
|
|||||||
@@ -15,7 +15,6 @@ pub mod repository;
|
|||||||
pub mod sbom;
|
pub mod sbom;
|
||||||
pub mod scan;
|
pub mod scan;
|
||||||
pub(crate) mod serde_helpers;
|
pub(crate) mod serde_helpers;
|
||||||
pub mod werkbank;
|
|
||||||
|
|
||||||
pub use auth::AuthInfo;
|
pub use auth::AuthInfo;
|
||||||
pub use chat::{ChatMessage, ChatRequest, ChatResponse, SourceReference};
|
pub use chat::{ChatMessage, ChatRequest, ChatResponse, SourceReference};
|
||||||
@@ -48,8 +47,3 @@ pub use pentest::{
|
|||||||
pub use repository::ScanTrigger;
|
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::{
|
|
||||||
CompleteRequest, CompleteResponse, DastCollect, Executor, HeartbeatAck, HeartbeatRequest,
|
|
||||||
InputRef, Job, JobCollect, JobRecord, JobResult, JobRuntime, JobStatus, JobType, LeaseRequest,
|
|
||||||
LeasedJob,
|
|
||||||
};
|
|
||||||
|
|||||||
@@ -1,514 +0,0 @@
|
|||||||
//! The Werkbank job/result contract (WB-01).
|
|
||||||
//!
|
|
||||||
//! The shared, dependency-free vocabulary the control plane and the Werkbank
|
|
||||||
//! execution runner agree on: what a [`Job`] is, which [`Executor`] runs it, how
|
|
||||||
//! it moves through the queue ([`JobStatus`]), and what a [`JobResult`] carries
|
|
||||||
//! back. Jobs are declarative — TOML on disk, JSON on the wire — and results
|
|
||||||
//! reuse the existing scanner result types ([`Finding`], [`DastFinding`],
|
|
||||||
//! [`SbomEntry`]) so the runner produces exactly what the control plane persists.
|
|
||||||
//!
|
|
||||||
//! This module is intentionally free of the `mongodb`/`axum` features so the
|
|
||||||
//! runner can depend on `compliance-core` without pulling the server stack.
|
|
||||||
|
|
||||||
use std::collections::BTreeMap;
|
|
||||||
|
|
||||||
use chrono::{DateTime, Utc};
|
|
||||||
use serde::{Deserialize, Serialize};
|
|
||||||
|
|
||||||
use super::dast::DastFinding;
|
|
||||||
use super::finding::Finding;
|
|
||||||
use super::sbom::SbomEntry;
|
|
||||||
|
|
||||||
/// The kind of dynamic-execution job.
|
|
||||||
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
|
|
||||||
#[serde(rename_all = "kebab-case")]
|
|
||||||
pub enum JobType {
|
|
||||||
/// Instantiate control logic on an ephemeral soft-PLC and probe it.
|
|
||||||
PlcProvision,
|
|
||||||
/// Boot a firmware image under QEMU and run dynamic checks.
|
|
||||||
QemuBoot,
|
|
||||||
/// Crawl and dynamically test a running web endpoint.
|
|
||||||
Dast,
|
|
||||||
/// Run an active penetration test against a running target.
|
|
||||||
Pentest,
|
|
||||||
}
|
|
||||||
|
|
||||||
/// How a runner executes a job — the CI-runner-style classification. A runner
|
|
||||||
/// advertises exactly one; a job requires one.
|
|
||||||
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
|
|
||||||
#[serde(rename_all = "lowercase")]
|
|
||||||
pub enum Executor {
|
|
||||||
/// A subprocess on the runner host (dev / trusted single-node).
|
|
||||||
Shell,
|
|
||||||
/// One or more containers on the runner's Docker (default; QEMU runs here).
|
|
||||||
Docker,
|
|
||||||
/// A Pod/Job in a Kubernetes cluster (scale-out / multi-tenant).
|
|
||||||
K8s,
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Lifecycle state of a job in the queue.
|
|
||||||
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
|
|
||||||
#[serde(rename_all = "lowercase")]
|
|
||||||
pub enum JobStatus {
|
|
||||||
/// Waiting to be leased.
|
|
||||||
Queued,
|
|
||||||
/// Leased by a runner but not yet started.
|
|
||||||
Leased,
|
|
||||||
/// Executing on a runner.
|
|
||||||
Running,
|
|
||||||
/// Completed successfully.
|
|
||||||
Succeeded,
|
|
||||||
/// Completed with an error.
|
|
||||||
Failed,
|
|
||||||
/// The lease/lifetime deadline elapsed before completion.
|
|
||||||
Expired,
|
|
||||||
/// Cancelled by the control plane.
|
|
||||||
Cancelled,
|
|
||||||
}
|
|
||||||
|
|
||||||
impl JobStatus {
|
|
||||||
/// Whether the job has reached a terminal state (no further transitions).
|
|
||||||
pub fn is_terminal(self) -> bool {
|
|
||||||
matches!(
|
|
||||||
self,
|
|
||||||
JobStatus::Succeeded | JobStatus::Failed | JobStatus::Expired | JobStatus::Cancelled
|
|
||||||
)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/// A reference to an input artifact. Resolved by the runner from a source it can
|
|
||||||
/// reach; the blob itself never flows through the control plane (so an on-prem
|
|
||||||
/// runner keeps customer data local). Exactly one of `blob`/`url` should be set.
|
|
||||||
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
|
|
||||||
pub struct InputRef {
|
|
||||||
/// Content-addressed blob (e.g. `sha256:…`) the runner fetches from its store.
|
|
||||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
|
||||||
pub blob: Option<String>,
|
|
||||||
/// A URL the runner can reach (git repo, internal artifact store, …).
|
|
||||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
|
||||||
pub url: Option<String>,
|
|
||||||
}
|
|
||||||
|
|
||||||
impl InputRef {
|
|
||||||
/// A content-addressed blob reference.
|
|
||||||
pub fn blob(id: impl Into<String>) -> Self {
|
|
||||||
Self {
|
|
||||||
blob: Some(id.into()),
|
|
||||||
url: None,
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Sandbox runtime knobs. Fields are executor/job-type specific and all optional;
|
|
||||||
/// `extra` carries anything not modelled explicitly.
|
|
||||||
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
|
|
||||||
pub struct JobRuntime {
|
|
||||||
/// Container image (Docker executor).
|
|
||||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
|
||||||
pub image: Option<String>,
|
|
||||||
/// Memory cap (e.g. `512m`).
|
|
||||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
|
||||||
pub memory: Option<String>,
|
|
||||||
/// CPU cap (e.g. `0.5`).
|
|
||||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
|
||||||
pub cpus: Option<String>,
|
|
||||||
/// Network to join (e.g. `isolated`).
|
|
||||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
|
||||||
pub network: Option<String>,
|
|
||||||
/// QEMU machine type (qemu-boot).
|
|
||||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
|
||||||
pub machine: Option<String>,
|
|
||||||
/// QEMU target architecture (qemu-boot).
|
|
||||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
|
||||||
pub arch: Option<String>,
|
|
||||||
/// Executor-specific extras not modelled above.
|
|
||||||
#[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
|
|
||||||
pub extra: BTreeMap<String, String>,
|
|
||||||
}
|
|
||||||
|
|
||||||
/// DAST collection settings for jobs that scan a web endpoint.
|
|
||||||
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
|
|
||||||
pub struct DastCollect {
|
|
||||||
/// Maximum crawl depth (kept shallow for ephemeral instances).
|
|
||||||
pub max_crawl_depth: u32,
|
|
||||||
}
|
|
||||||
|
|
||||||
/// What to collect from a run.
|
|
||||||
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
|
|
||||||
pub struct JobCollect {
|
|
||||||
/// Run the industrial-protocol probe (Modbus/OPC-UA/EtherNet-IP).
|
|
||||||
#[serde(default)]
|
|
||||||
pub ics_probe: bool,
|
|
||||||
/// Run DAST against the provisioned/booted web endpoint.
|
|
||||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
|
||||||
pub dast: Option<DastCollect>,
|
|
||||||
/// Run an active pentest.
|
|
||||||
#[serde(default)]
|
|
||||||
pub pentest: bool,
|
|
||||||
/// Collect an SBOM.
|
|
||||||
#[serde(default)]
|
|
||||||
pub sbom: bool,
|
|
||||||
}
|
|
||||||
|
|
||||||
/// A declarative dynamic-execution job the control plane enqueues and a Werkbank
|
|
||||||
/// runner leases and executes.
|
|
||||||
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
|
|
||||||
pub struct Job {
|
|
||||||
/// Unique job id (assigned by the control plane on enqueue).
|
|
||||||
pub id: String,
|
|
||||||
/// What kind of job this is.
|
|
||||||
#[serde(rename = "type")]
|
|
||||||
pub job_type: JobType,
|
|
||||||
/// Owning tenant.
|
|
||||||
pub tenant: String,
|
|
||||||
/// The onboarded target this job tests.
|
|
||||||
pub target_id: String,
|
|
||||||
/// The executor a runner must provide to run this job.
|
|
||||||
pub executor: Executor,
|
|
||||||
/// Runner capabilities this job requires (e.g. `arch=amd64`, `kvm=true`).
|
|
||||||
#[serde(default, skip_serializing_if = "Vec::is_empty")]
|
|
||||||
pub labels: Vec<String>,
|
|
||||||
/// Hard lifetime deadline for the whole job.
|
|
||||||
pub timeout_secs: u64,
|
|
||||||
/// Named input artifacts (e.g. `program`, `firmware`), by reference.
|
|
||||||
#[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
|
|
||||||
pub inputs: BTreeMap<String, InputRef>,
|
|
||||||
/// Sandbox runtime knobs.
|
|
||||||
#[serde(default)]
|
|
||||||
pub runtime: JobRuntime,
|
|
||||||
/// What to collect from the run.
|
|
||||||
#[serde(default)]
|
|
||||||
pub collect: JobCollect,
|
|
||||||
}
|
|
||||||
|
|
||||||
impl Job {
|
|
||||||
/// A `plc-provision` job: instantiate the control logic named `program` on an
|
|
||||||
/// ephemeral soft-PLC (Docker executor) and collect the ICS probe + DAST.
|
|
||||||
pub fn plc_provision(
|
|
||||||
id: impl Into<String>,
|
|
||||||
tenant: impl Into<String>,
|
|
||||||
target_id: impl Into<String>,
|
|
||||||
program: InputRef,
|
|
||||||
timeout_secs: u64,
|
|
||||||
) -> Self {
|
|
||||||
let mut inputs = BTreeMap::new();
|
|
||||||
inputs.insert("program".to_string(), program);
|
|
||||||
Self {
|
|
||||||
id: id.into(),
|
|
||||||
job_type: JobType::PlcProvision,
|
|
||||||
tenant: tenant.into(),
|
|
||||||
target_id: target_id.into(),
|
|
||||||
executor: Executor::Docker,
|
|
||||||
labels: Vec::new(),
|
|
||||||
timeout_secs,
|
|
||||||
inputs,
|
|
||||||
runtime: JobRuntime::default(),
|
|
||||||
collect: JobCollect {
|
|
||||||
ics_probe: true,
|
|
||||||
dast: Some(DastCollect { max_crawl_depth: 2 }),
|
|
||||||
pentest: false,
|
|
||||||
sbom: false,
|
|
||||||
},
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/// The outcome of running a job, posted back to the control plane. Findings and
|
|
||||||
/// SBOM reuse the shared scanner types, so the control plane persists them
|
|
||||||
/// unchanged. Submission is idempotent — keyed by [`JobResult::job_id`].
|
|
||||||
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
|
|
||||||
pub struct JobResult {
|
|
||||||
/// The job this result is for.
|
|
||||||
pub job_id: String,
|
|
||||||
/// Terminal status of the job.
|
|
||||||
pub status: Option<JobStatus>,
|
|
||||||
/// General scanner findings (e.g. ICS-probe findings).
|
|
||||||
#[serde(default, skip_serializing_if = "Vec::is_empty")]
|
|
||||||
pub findings: Vec<Finding>,
|
|
||||||
/// DAST findings from a web-endpoint scan.
|
|
||||||
#[serde(default, skip_serializing_if = "Vec::is_empty")]
|
|
||||||
pub dast_findings: Vec<DastFinding>,
|
|
||||||
/// SBOM components collected from the run.
|
|
||||||
#[serde(default, skip_serializing_if = "Vec::is_empty")]
|
|
||||||
pub sbom: Vec<SbomEntry>,
|
|
||||||
/// Error message when the job failed.
|
|
||||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
|
||||||
pub error: Option<String>,
|
|
||||||
/// Captured execution log (truncated by the runner).
|
|
||||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
|
||||||
pub logs: Option<String>,
|
|
||||||
/// When execution started on the runner.
|
|
||||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
|
||||||
pub started_at: Option<DateTime<Utc>>,
|
|
||||||
/// When execution finished.
|
|
||||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
|
||||||
pub finished_at: Option<DateTime<Utc>>,
|
|
||||||
}
|
|
||||||
|
|
||||||
impl JobResult {
|
|
||||||
/// A successful result for a job.
|
|
||||||
pub fn succeeded(job_id: impl Into<String>) -> Self {
|
|
||||||
Self {
|
|
||||||
job_id: job_id.into(),
|
|
||||||
status: Some(JobStatus::Succeeded),
|
|
||||||
..Default::default()
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/// A failed result carrying an error message.
|
|
||||||
pub fn failed(job_id: impl Into<String>, error: impl Into<String>) -> Self {
|
|
||||||
Self {
|
|
||||||
job_id: job_id.into(),
|
|
||||||
status: Some(JobStatus::Failed),
|
|
||||||
error: Some(error.into()),
|
|
||||||
..Default::default()
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/// A queued job as persisted by the control plane (WB-02): the [`Job`] contract
|
|
||||||
/// plus the queue bookkeeping — status, lease ownership, attempt count, and the
|
|
||||||
/// eventual result. The runner never sees this record; on lease it receives a
|
|
||||||
/// [`LeasedJob`] (the job plus a token it presents to heartbeat/complete).
|
|
||||||
///
|
|
||||||
/// Timestamps persist as native BSON dates so the queue's range queries (lease
|
|
||||||
/// FIFO by `created_at`, visibility-timeout sweep by `lease_expires_at`) compare
|
|
||||||
/// correctly.
|
|
||||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
|
||||||
pub struct JobRecord {
|
|
||||||
/// The job to run.
|
|
||||||
pub job: Job,
|
|
||||||
/// Current queue state.
|
|
||||||
pub status: JobStatus,
|
|
||||||
/// The lease token held by the current runner (proves lease ownership).
|
|
||||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
|
||||||
pub lease_token: Option<String>,
|
|
||||||
/// Id of the runner holding the lease.
|
|
||||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
|
||||||
pub leased_by: Option<String>,
|
|
||||||
/// When the current lease expires — the visibility timeout after which a
|
|
||||||
/// crashed runner's job is swept back to `queued`.
|
|
||||||
#[serde(default, with = "super::serde_helpers::opt_bson_datetime")]
|
|
||||||
pub lease_expires_at: Option<DateTime<Utc>>,
|
|
||||||
/// Last heartbeat from the runner.
|
|
||||||
#[serde(default, with = "super::serde_helpers::opt_bson_datetime")]
|
|
||||||
pub heartbeat_at: Option<DateTime<Utc>>,
|
|
||||||
/// How many times the job has been leased (incremented on each lease).
|
|
||||||
#[serde(default)]
|
|
||||||
pub attempts: u32,
|
|
||||||
/// Set when the control plane requests cancellation; the runner sees it on
|
|
||||||
/// its next heartbeat and aborts.
|
|
||||||
#[serde(default)]
|
|
||||||
pub cancel_requested: bool,
|
|
||||||
/// The result, once the job reaches a terminal state.
|
|
||||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
|
||||||
pub result: Option<JobResult>,
|
|
||||||
/// When the job was enqueued.
|
|
||||||
#[serde(with = "super::serde_helpers::bson_datetime")]
|
|
||||||
pub created_at: DateTime<Utc>,
|
|
||||||
/// Last modification.
|
|
||||||
#[serde(with = "super::serde_helpers::bson_datetime")]
|
|
||||||
pub updated_at: DateTime<Utc>,
|
|
||||||
}
|
|
||||||
|
|
||||||
impl JobRecord {
|
|
||||||
/// A freshly-enqueued (`queued`) record for a job.
|
|
||||||
pub fn queued(job: Job, now: DateTime<Utc>) -> Self {
|
|
||||||
Self {
|
|
||||||
job,
|
|
||||||
status: JobStatus::Queued,
|
|
||||||
lease_token: None,
|
|
||||||
leased_by: None,
|
|
||||||
lease_expires_at: None,
|
|
||||||
heartbeat_at: None,
|
|
||||||
attempts: 0,
|
|
||||||
cancel_requested: false,
|
|
||||||
result: None,
|
|
||||||
created_at: now,
|
|
||||||
updated_at: now,
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/// A job handed to a runner on lease: what to run plus the token the runner must
|
|
||||||
/// present to heartbeat and complete it.
|
|
||||||
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
|
|
||||||
pub struct LeasedJob {
|
|
||||||
/// The job to execute.
|
|
||||||
pub job: Job,
|
|
||||||
/// The lease token proving ownership (opaque to the runner).
|
|
||||||
pub lease_token: String,
|
|
||||||
}
|
|
||||||
|
|
||||||
/// The runner's view of a heartbeat: whether the control plane has asked the job
|
|
||||||
/// to stop. `None` from the queue means the lease was lost (token mismatch or the
|
|
||||||
/// job already terminal) and the runner should abandon the work.
|
|
||||||
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
|
|
||||||
pub struct HeartbeatAck {
|
|
||||||
/// The control plane requested cancellation — the runner should tear down.
|
|
||||||
pub cancelled: bool,
|
|
||||||
}
|
|
||||||
|
|
||||||
// --- 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 {
|
|
||||||
use super::*;
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn job_round_trips_through_json() {
|
|
||||||
let job = Job::plc_provision("job_1", "acme", "64f0aa", InputRef::blob("sha256:abc"), 180);
|
|
||||||
let json = serde_json::to_string(&job).expect("serialize");
|
|
||||||
let back: Job = serde_json::from_str(&json).expect("deserialize");
|
|
||||||
assert_eq!(job, back);
|
|
||||||
// Enum wire forms are the kebab/lowercase the contract documents.
|
|
||||||
assert!(json.contains("\"type\":\"plc-provision\""));
|
|
||||||
assert!(json.contains("\"executor\":\"docker\""));
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn parses_the_design_doc_plc_provision_toml() {
|
|
||||||
// The exact shape from docs/DESIGN.md §5 (wrapped in a [job] table).
|
|
||||||
#[derive(Deserialize)]
|
|
||||||
struct JobFile {
|
|
||||||
job: Job,
|
|
||||||
}
|
|
||||||
let src = r#"
|
|
||||||
[job]
|
|
||||||
id = "job_01H"
|
|
||||||
type = "plc-provision"
|
|
||||||
tenant = "acme"
|
|
||||||
target_id = "64f0"
|
|
||||||
executor = "docker"
|
|
||||||
labels = ["arch=amd64"]
|
|
||||||
timeout_secs = 180
|
|
||||||
|
|
||||||
[job.inputs]
|
|
||||||
program = { blob = "sha256:deadbeef" }
|
|
||||||
|
|
||||||
[job.runtime]
|
|
||||||
image = "openplc:latest"
|
|
||||||
memory = "512m"
|
|
||||||
cpus = "0.5"
|
|
||||||
network = "isolated"
|
|
||||||
|
|
||||||
[job.collect]
|
|
||||||
ics_probe = true
|
|
||||||
dast = { max_crawl_depth = 2 }
|
|
||||||
"#;
|
|
||||||
let file: JobFile = toml::from_str(src).expect("parse job toml");
|
|
||||||
let job = file.job;
|
|
||||||
assert_eq!(job.job_type, JobType::PlcProvision);
|
|
||||||
assert_eq!(job.executor, Executor::Docker);
|
|
||||||
assert_eq!(job.labels, vec!["arch=amd64".to_string()]);
|
|
||||||
assert_eq!(
|
|
||||||
job.inputs.get("program").and_then(|i| i.blob.as_deref()),
|
|
||||||
Some("sha256:deadbeef")
|
|
||||||
);
|
|
||||||
assert_eq!(job.runtime.image.as_deref(), Some("openplc:latest"));
|
|
||||||
assert!(job.collect.ics_probe);
|
|
||||||
assert_eq!(job.collect.dast.map(|d| d.max_crawl_depth), Some(2));
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn qemu_boot_runtime_fields_parse() {
|
|
||||||
#[derive(Deserialize)]
|
|
||||||
struct JobFile {
|
|
||||||
job: Job,
|
|
||||||
}
|
|
||||||
let src = r#"
|
|
||||||
[job]
|
|
||||||
id = "j2"
|
|
||||||
type = "qemu-boot"
|
|
||||||
tenant = "acme"
|
|
||||||
target_id = "t"
|
|
||||||
executor = "docker"
|
|
||||||
labels = ["kvm=true"]
|
|
||||||
timeout_secs = 600
|
|
||||||
[job.inputs]
|
|
||||||
firmware = { blob = "sha256:cafe" }
|
|
||||||
[job.runtime]
|
|
||||||
machine = "virt"
|
|
||||||
arch = "arm"
|
|
||||||
memory = "1g"
|
|
||||||
"#;
|
|
||||||
let file: JobFile = toml::from_str(src).expect("parse");
|
|
||||||
assert_eq!(file.job.job_type, JobType::QemuBoot);
|
|
||||||
assert_eq!(file.job.runtime.arch.as_deref(), Some("arm"));
|
|
||||||
assert_eq!(
|
|
||||||
file.job
|
|
||||||
.inputs
|
|
||||||
.get("firmware")
|
|
||||||
.and_then(|i| i.blob.as_deref()),
|
|
||||||
Some("sha256:cafe")
|
|
||||||
);
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn status_terminality() {
|
|
||||||
assert!(JobStatus::Succeeded.is_terminal());
|
|
||||||
assert!(JobStatus::Expired.is_terminal());
|
|
||||||
assert!(!JobStatus::Queued.is_terminal());
|
|
||||||
assert!(!JobStatus::Running.is_terminal());
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn result_constructors() {
|
|
||||||
assert_eq!(JobResult::succeeded("j").status, Some(JobStatus::Succeeded));
|
|
||||||
let f = JobResult::failed("j", "boom");
|
|
||||||
assert_eq!(f.status, Some(JobStatus::Failed));
|
|
||||||
assert_eq!(f.error.as_deref(), Some("boom"));
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -44,7 +44,6 @@ export default withMermaid(defineConfig({
|
|||||||
items: [
|
items: [
|
||||||
{ text: 'Glossary', link: '/reference/glossary' },
|
{ text: 'Glossary', link: '/reference/glossary' },
|
||||||
{ text: 'Tools & Scanners', link: '/reference/tools' },
|
{ text: 'Tools & Scanners', link: '/reference/tools' },
|
||||||
{ text: 'PLC Runtime Landscape', link: '/reference/plc-runtimes' },
|
|
||||||
],
|
],
|
||||||
},
|
},
|
||||||
],
|
],
|
||||||
|
|||||||
@@ -1,97 +0,0 @@
|
|||||||
# PLC Runtime Landscape & Support
|
|
||||||
|
|
||||||
A soft PLC is a **SoC + Linux + a software runtime + an IEC 61131-3 control app**
|
|
||||||
(see [PLC / SPS Projects](/guide/plc)).
|
|
||||||
The **runtime** is what defines the device — it provides the IEC engine, the
|
|
||||||
Modbus / OPC UA / EtherNet/IP servers, and the WebVisu. This page tracks the
|
|
||||||
runtime ecosystems Certifai may encounter.
|
|
||||||
|
|
||||||
We do **not** aim to support every runtime up front. Certifai supports the
|
|
||||||
**CODESYS family** today; everything else is a **watch-list** — when a customer
|
|
||||||
shows up using one, we add the parser/support for it then. The dynamic OT probe
|
|
||||||
(Modbus / OPC UA / EtherNet/IP) is **vendor-agnostic** and works regardless of
|
|
||||||
the runtime.
|
|
||||||
|
|
||||||
## Support status
|
|
||||||
|
|
||||||
| Status | Meaning |
|
|
||||||
| --- | --- |
|
|
||||||
| ✅ **Supported** | Static analysis works today (control-logic SAST + library/runtime SBOM + CVE). |
|
|
||||||
| 🟡 **Covered via CODESYS** | A rebranded CODESYS runtime — our CODESYS parsing applies (may need minor per-vendor tweaks). |
|
|
||||||
| 🔭 **Watch-list** | Own project format — we add a format parser when a customer needs it. The dynamic OT probe already applies. |
|
|
||||||
| 🧪 **Test-bench** | A free runtime we use to *reconstruct and dynamically test* a device (see epic: provision-and-test). |
|
|
||||||
|
|
||||||
## 1. CODESYS and rebranded CODESYS (the largest slice)
|
|
||||||
|
|
||||||
Much of the market licenses the CODESYS runtime and rebrands the IDE. If a
|
|
||||||
customer "doesn't use CODESYS", they often do — under another name.
|
|
||||||
|
|
||||||
| Product / vendor | Based on | Status |
|
|
||||||
| --- | --- | --- |
|
|
||||||
| **CODESYS** (3S-Smart Software Solutions) | CODESYS | ✅ Supported |
|
|
||||||
| Schneider **EcoStruxure Machine Expert** (ex-SoMachine) | CODESYS | 🟡 Covered via CODESYS |
|
|
||||||
| **WAGO** e!COCKPIT / PFC controllers | CODESYS | 🟡 Covered via CODESYS |
|
|
||||||
| **ABB** AC500 / Automation Builder | CODESYS | 🟡 Covered via CODESYS |
|
|
||||||
| **Bosch Rexroth** ctrlX / IndraLogic | CODESYS | 🟡 Covered via CODESYS |
|
|
||||||
| **Eaton** XSoft-CODESYS, **KEBA** KeStudio, Berghof, Kontron, Festo (CPX-E), IFM, Turck, … | CODESYS | 🟡 Covered via CODESYS |
|
|
||||||
|
|
||||||
## 2. Other embeddable IEC 61131-3 runtime toolkits
|
|
||||||
|
|
||||||
Same model as CODESYS (an OEM licenses a runtime + IDE and bakes it into a
|
|
||||||
device), but with **different project formats and libraries**.
|
|
||||||
|
|
||||||
| Toolkit | Vendor | Status |
|
|
||||||
| --- | --- | --- |
|
|
||||||
| **ProConOS / MULTIPROG** | Phoenix Contact / KW-Software | 🔭 Watch-list |
|
|
||||||
| **ISaGRAF** (also does IEC 61499) | Rockwell | 🔭 Watch-list |
|
|
||||||
| **straton** | COPA-DATA | 🔭 Watch-list |
|
|
||||||
| **logi.CAD** | logi.cals | 🔭 Watch-list |
|
|
||||||
|
|
||||||
## 3. Fully proprietary ecosystems (own runtime + IDE + protocols)
|
|
||||||
|
|
||||||
Static analysis here needs a **per-vendor project parser**; the **dynamic OT
|
|
||||||
probe still works** (they speak Modbus / OPC UA / EtherNet/IP, plus vendor
|
|
||||||
protocols like S7comm / CIP).
|
|
||||||
|
|
||||||
| Ecosystem | Vendor | Notes | Status |
|
|
||||||
| --- | --- | --- | --- |
|
|
||||||
| **TIA Portal / STEP 7** (S7-1200/1500), S7-1500 **Software Controller**, **Virtual PLC** | Siemens | Largest install base; the soft/virtual variants are Linux/container | 🔭 Watch-list |
|
|
||||||
| **Studio 5000** (ControlLogix / CompactLogix) | Rockwell / Allen-Bradley | Strong in North America | 🔭 Watch-list |
|
|
||||||
| **TwinCAT 3** | Beckhoff | Genuine PC-based control on Windows / TwinCAT-BSD; IEC 61131-3 **+ C++ + Simulink** | 🔭 Watch-list |
|
|
||||||
| **Automation Studio** | B&R (ABB) | Own Automation Runtime | 🔭 Watch-list |
|
|
||||||
| **GX Works** (MELSEC) | Mitsubishi | | 🔭 Watch-list |
|
|
||||||
| **Sysmac Studio** (NX / NJ) | Omron | | 🔭 Watch-list |
|
|
||||||
| **Proficy Machine Edition** (PACSystems) | Emerson / GE | | 🔭 Watch-list |
|
|
||||||
|
|
||||||
## 4. Linux-native / containerized soft-PLC (the direction of travel)
|
|
||||||
|
|
||||||
| Product | Vendor | Notes | Status |
|
|
||||||
| --- | --- | --- | --- |
|
|
||||||
| **PLCnext** | Phoenix Contact | Open, Linux-based; native runtime is eCLR (not CODESYS), but can also run CODESYS as an app | 🔭 Watch-list |
|
|
||||||
| **ctrlX** | Bosch Rexroth | Ubuntu-core, app-store model (CODESYS runtime inside) | 🟡 Covered via CODESYS |
|
|
||||||
| **Virtual PLC** / **CODESYS Virtual Control** | Siemens / CODESYS | Containerized PLCs (Docker / K8s) | 🟡 / 🔭 |
|
|
||||||
|
|
||||||
## 5. Open-source runtimes (free — our test-bench substrates)
|
|
||||||
|
|
||||||
Used to **reconstruct and dynamically test** a customer device without touching
|
|
||||||
their network (provision-and-test).
|
|
||||||
|
|
||||||
| Runtime | Standard | Notes | Status |
|
|
||||||
| --- | --- | --- | --- |
|
|
||||||
| **OpenPLC** | IEC 61131-3 | Modbus-centric, education/small automation; uses MatIEC | 🧪 Test-bench (current) |
|
|
||||||
| **Beremiz + MatIEC** | IEC 61131-3 | Fuller open-source IDE; compiles ST/IL → C. Natural fidelity step-up from OpenPLC | 🧪 Test-bench (candidate) |
|
|
||||||
| **Eclipse 4diac (FORTE)** | IEC **61499** | Distributed, event-driven — a *different paradigm* from 61131-3's scan cycle | 🔭 Watch-list |
|
|
||||||
| **ProView** | — | Open-source process control + SCADA | 🔭 Watch-list |
|
|
||||||
|
|
||||||
## How we add support for a new runtime
|
|
||||||
|
|
||||||
- **Static (SAST / SBOM):** needs a parser for that runtime's **project format**
|
|
||||||
(and its library/package convention). This is the per-vendor work.
|
|
||||||
- **Dynamic (ICS probe / DAST):** already **vendor-agnostic** — it targets the
|
|
||||||
device's OT ports and WebVisu, not the runtime's file format. So a brand-new
|
|
||||||
ecosystem still gets dynamic coverage on day one.
|
|
||||||
|
|
||||||
::: tip Rule of thumb
|
|
||||||
Confirm whether a "non-CODESYS" controller is actually a **rebranded CODESYS**
|
|
||||||
runtime (Section 1) before assuming new work — most of the long tail is.
|
|
||||||
:::
|
|
||||||
Reference in New Issue
Block a user