Compare commits
3
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
896a06e8f6 | ||
|
|
69bce2f07c | ||
|
|
1ae6025286 |
@@ -107,8 +107,6 @@ jobs:
|
|||||||
run: cargo clippy -p compliance-dashboard --features web --no-default-features -- -D warnings
|
run: cargo clippy -p compliance-dashboard --features web --no-default-features -- -D warnings
|
||||||
- name: Clippy (mcp)
|
- name: Clippy (mcp)
|
||||||
run: cargo clippy -p compliance-mcp -- -D warnings
|
run: cargo clippy -p compliance-mcp -- -D warnings
|
||||||
- name: Clippy (werkbank-exec)
|
|
||||||
run: cargo clippy -p werkbank-exec -- -D warnings
|
|
||||||
|
|
||||||
# Security audit
|
# Security audit
|
||||||
- name: Security Audit
|
- name: Security Audit
|
||||||
@@ -117,8 +115,8 @@ jobs:
|
|||||||
RUSTC_WRAPPER: ""
|
RUSTC_WRAPPER: ""
|
||||||
|
|
||||||
# Tests (reuses compilation artifacts from clippy)
|
# Tests (reuses compilation artifacts from clippy)
|
||||||
- name: Tests (core + agent + werkbank-exec)
|
- name: Tests (core + agent)
|
||||||
run: cargo test -p compliance-core -p compliance-agent -p werkbank-exec --lib
|
run: cargo test -p compliance-core -p compliance-agent --lib
|
||||||
- name: Tests (dashboard server)
|
- name: Tests (dashboard server)
|
||||||
run: cargo test -p compliance-dashboard --features server --no-default-features
|
run: cargo test -p compliance-dashboard --features server --no-default-features
|
||||||
- name: Tests (dashboard web)
|
- name: Tests (dashboard web)
|
||||||
|
|||||||
Generated
-21
@@ -699,7 +699,6 @@ dependencies = [
|
|||||||
"urlencoding",
|
"urlencoding",
|
||||||
"uuid",
|
"uuid",
|
||||||
"walkdir",
|
"walkdir",
|
||||||
"werkbank-exec",
|
|
||||||
"zip",
|
"zip",
|
||||||
]
|
]
|
||||||
|
|
||||||
@@ -724,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",
|
||||||
@@ -6715,25 +6713,6 @@ dependencies = [
|
|||||||
"rustls-pki-types",
|
"rustls-pki-types",
|
||||||
]
|
]
|
||||||
|
|
||||||
[[package]]
|
|
||||||
name = "werkbank-exec"
|
|
||||||
version = "0.1.0"
|
|
||||||
dependencies = [
|
|
||||||
"compliance-core",
|
|
||||||
"compliance-dast",
|
|
||||||
"futures-util",
|
|
||||||
"hex",
|
|
||||||
"regex",
|
|
||||||
"reqwest",
|
|
||||||
"secrecy",
|
|
||||||
"sha2",
|
|
||||||
"thiserror 2.0.18",
|
|
||||||
"tokio",
|
|
||||||
"tracing",
|
|
||||||
"uuid",
|
|
||||||
"walkdir",
|
|
||||||
]
|
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "which"
|
name = "which"
|
||||||
version = "6.0.3"
|
version = "6.0.3"
|
||||||
|
|||||||
@@ -7,7 +7,6 @@ members = [
|
|||||||
"compliance-dast",
|
"compliance-dast",
|
||||||
"compliance-mcp",
|
"compliance-mcp",
|
||||||
"compliance-smoke",
|
"compliance-smoke",
|
||||||
"werkbank-exec",
|
|
||||||
]
|
]
|
||||||
resolver = "2"
|
resolver = "2"
|
||||||
|
|
||||||
|
|||||||
@@ -10,9 +10,6 @@ workspace = true
|
|||||||
compliance-core = { workspace = true, features = ["mongodb", "telemetry", "axum"] }
|
compliance-core = { workspace = true, features = ["mongodb", "telemetry", "axum"] }
|
||||||
compliance-graph = { path = "../compliance-graph" }
|
compliance-graph = { path = "../compliance-graph" }
|
||||||
compliance-dast = { path = "../compliance-dast" }
|
compliance-dast = { path = "../compliance-dast" }
|
||||||
# Shared dynamic-execution logic (soft-PLC provisioning + ICS probing), also
|
|
||||||
# used by the Werkbank runner.
|
|
||||||
werkbank-exec = { path = "../werkbank-exec" }
|
|
||||||
# Native firmware build/target detection for bare-metal & RTOS artifacts.
|
# Native firmware build/target detection for bare-metal & RTOS artifacts.
|
||||||
# Same-company IP, used directly (not via CLI) so the whole tramiton suite is
|
# Same-company IP, used directly (not via CLI) so the whole tramiton suite is
|
||||||
# available to the onboarding classifier. NOTE: CI must be able to fetch this
|
# available to the onboarding classifier. NOTE: CI must be able to fetch this
|
||||||
|
|||||||
@@ -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)
|
||||||
|
|||||||
@@ -27,9 +27,6 @@ pub enum AgentError {
|
|||||||
#[error("Configuration error: {0}")]
|
#[error("Configuration error: {0}")]
|
||||||
Config(String),
|
Config(String),
|
||||||
|
|
||||||
#[error("Dynamic-execution error: {0}")]
|
|
||||||
Exec(#[from] werkbank_exec::ExecError),
|
|
||||||
|
|
||||||
#[error("{0}")]
|
#[error("{0}")]
|
||||||
Other(String),
|
Other(String),
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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;
|
|
||||||
|
|||||||
@@ -14,7 +14,7 @@ use std::time::Duration;
|
|||||||
|
|
||||||
use compliance_core::models::{Finding, ScanType, Severity};
|
use compliance_core::models::{Finding, ScanType, Severity};
|
||||||
|
|
||||||
use crate::fingerprint as dedup;
|
use crate::pipeline::dedup;
|
||||||
|
|
||||||
/// Well-known deep-probe ports (each independent of any WebVisu HTTP port).
|
/// Well-known deep-probe ports (each independent of any WebVisu HTTP port).
|
||||||
const MODBUS_PORT: u16 = 502;
|
const MODBUS_PORT: u16 = 502;
|
||||||
@@ -5,6 +5,7 @@ pub mod firmware_sbom;
|
|||||||
pub mod git;
|
pub mod git;
|
||||||
pub mod gitleaks;
|
pub mod gitleaks;
|
||||||
mod graph_build;
|
mod graph_build;
|
||||||
|
pub mod ics;
|
||||||
mod issue_creation;
|
mod issue_creation;
|
||||||
pub mod lint;
|
pub mod lint;
|
||||||
pub mod orchestrator;
|
pub mod orchestrator;
|
||||||
|
|||||||
@@ -587,7 +587,7 @@ impl PipelineOrchestrator {
|
|||||||
let path = ingest_set
|
let path = ingest_set
|
||||||
.get(&a.id)
|
.get(&a.id)
|
||||||
.and_then(|ia| ia.working_path.clone())?;
|
.and_then(|ia| ia.working_path.clone())?;
|
||||||
werkbank_exec::plc::extract_program(&path)
|
crate::pipeline::plc::runtime::extract_program(&path)
|
||||||
});
|
});
|
||||||
let Some(program) = program else {
|
let Some(program) = program else {
|
||||||
tracing::info!(
|
tracing::info!(
|
||||||
@@ -597,9 +597,10 @@ impl PipelineOrchestrator {
|
|||||||
return Ok(0);
|
return Ok(0);
|
||||||
};
|
};
|
||||||
|
|
||||||
let http = werkbank_exec::plc::http_client()?;
|
let http = crate::pipeline::plc::runtime::http_client()?;
|
||||||
let provisioner = werkbank_exec::plc::DockerSoftPlc::new(self.config.plc_runtime.clone());
|
let provisioner =
|
||||||
let outcome = werkbank_exec::plc::provision_and_test(
|
crate::pipeline::plc::runtime::DockerSoftPlc::new(self.config.plc_runtime.clone());
|
||||||
|
let outcome = crate::pipeline::plc::runtime::provision_and_test(
|
||||||
&provisioner,
|
&provisioner,
|
||||||
&http,
|
&http,
|
||||||
&self.config.plc_runtime,
|
&self.config.plc_runtime,
|
||||||
@@ -662,7 +663,7 @@ impl PipelineOrchestrator {
|
|||||||
};
|
};
|
||||||
// Short per-request budget so an unreachable device doesn't stall the scan.
|
// Short per-request budget so an unreachable device doesn't stall the scan.
|
||||||
let budget = std::time::Duration::from_secs(5);
|
let budget = std::time::Duration::from_secs(5);
|
||||||
let findings = werkbank_exec::ics::probe_target(&endpoint, target_id, budget).await;
|
let findings = crate::pipeline::ics::probe_target(&endpoint, target_id, budget).await;
|
||||||
tracing::info!(
|
tracing::info!(
|
||||||
target_id,
|
target_id,
|
||||||
endpoint = %endpoint,
|
endpoint = %endpoint,
|
||||||
|
|||||||
@@ -9,6 +9,7 @@ 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;
|
||||||
|
|||||||
@@ -24,7 +24,7 @@ use compliance_core::models::dast::{DastFinding, DastScanRun, DastTarget, DastTa
|
|||||||
use compliance_core::models::Finding;
|
use compliance_core::models::Finding;
|
||||||
use compliance_core::PlcRuntimeConfig;
|
use compliance_core::PlcRuntimeConfig;
|
||||||
|
|
||||||
use crate::error::ExecError;
|
use crate::error::AgentError;
|
||||||
|
|
||||||
pub use provision::{DockerSoftPlc, ProvisionedRuntime, SoftPlc};
|
pub use provision::{DockerSoftPlc, ProvisionedRuntime, SoftPlc};
|
||||||
|
|
||||||
@@ -61,12 +61,12 @@ pub struct PlcProgram {
|
|||||||
|
|
||||||
/// A cookie-aware HTTP client for the OpenPLC web UI. A fresh client per scan
|
/// 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.
|
/// isolates the OpenPLC session (its Flask login cookie) from every other scan.
|
||||||
pub fn http_client() -> Result<reqwest::Client, ExecError> {
|
pub fn http_client() -> Result<reqwest::Client, AgentError> {
|
||||||
reqwest::Client::builder()
|
reqwest::Client::builder()
|
||||||
.cookie_store(true)
|
.cookie_store(true)
|
||||||
.timeout(Duration::from_secs(30))
|
.timeout(Duration::from_secs(30))
|
||||||
.build()
|
.build()
|
||||||
.map_err(ExecError::Http)
|
.map_err(AgentError::Http)
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Pick the control-logic program to run from an ingested PLC source tree.
|
/// Pick the control-logic program to run from an ingested PLC source tree.
|
||||||
@@ -151,7 +151,7 @@ pub async fn provision_and_test<P: SoftPlc>(
|
|||||||
cfg: &PlcRuntimeConfig,
|
cfg: &PlcRuntimeConfig,
|
||||||
program: &PlcProgram,
|
program: &PlcProgram,
|
||||||
target_id: &str,
|
target_id: &str,
|
||||||
) -> Result<ProvisionOutcome, ExecError> {
|
) -> Result<ProvisionOutcome, AgentError> {
|
||||||
let handle = provisioner.provision(target_id).await?;
|
let handle = provisioner.provision(target_id).await?;
|
||||||
tracing::info!(
|
tracing::info!(
|
||||||
target_id,
|
target_id,
|
||||||
@@ -193,7 +193,7 @@ async fn run_dynamic_test(
|
|||||||
program: &PlcProgram,
|
program: &PlcProgram,
|
||||||
target_id: &str,
|
target_id: &str,
|
||||||
handle: &ProvisionedRuntime,
|
handle: &ProvisionedRuntime,
|
||||||
) -> Result<ProvisionOutcome, ExecError> {
|
) -> Result<ProvisionOutcome, AgentError> {
|
||||||
let ready_budget = Duration::from_secs((cfg.max_lifetime_secs / 3).clamp(10, 60));
|
let ready_budget = Duration::from_secs((cfg.max_lifetime_secs / 3).clamp(10, 60));
|
||||||
openplc::wait_ready(http, &handle.webvisu_url, ready_budget).await?;
|
openplc::wait_ready(http, &handle.webvisu_url, ready_budget).await?;
|
||||||
|
|
||||||
@@ -212,7 +212,8 @@ async fn run_dynamic_test(
|
|||||||
tokio::time::sleep(Duration::from_secs(3)).await;
|
tokio::time::sleep(Duration::from_secs(3)).await;
|
||||||
|
|
||||||
let probe_budget = Duration::from_secs(5);
|
let probe_budget = Duration::from_secs(5);
|
||||||
let findings = crate::ics::probe_target(&handle.modbus_endpoint, target_id, probe_budget).await;
|
let findings =
|
||||||
|
crate::pipeline::ics::probe_target(&handle.modbus_endpoint, target_id, probe_budget).await;
|
||||||
tracing::info!(
|
tracing::info!(
|
||||||
target_id,
|
target_id,
|
||||||
instance = %handle.name,
|
instance = %handle.name,
|
||||||
@@ -347,10 +348,10 @@ mod tests {
|
|||||||
}
|
}
|
||||||
|
|
||||||
impl SoftPlc for FakeSoftPlc {
|
impl SoftPlc for FakeSoftPlc {
|
||||||
async fn provision(&self, _target_id: &str) -> Result<ProvisionedRuntime, ExecError> {
|
async fn provision(&self, _target_id: &str) -> Result<ProvisionedRuntime, AgentError> {
|
||||||
self.provisions.fetch_add(1, Ordering::SeqCst);
|
self.provisions.fetch_add(1, Ordering::SeqCst);
|
||||||
if self.fail_provision {
|
if self.fail_provision {
|
||||||
return Err(ExecError::Other("provision failed".into()));
|
return Err(AgentError::Other("provision failed".into()));
|
||||||
}
|
}
|
||||||
// Unreachable address so run_dynamic_test blocks on readiness until the
|
// Unreachable address so run_dynamic_test blocks on readiness until the
|
||||||
// deadline fires — exercising the teardown-on-deadline path.
|
// deadline fires — exercising the teardown-on-deadline path.
|
||||||
+15
-15
@@ -10,7 +10,7 @@
|
|||||||
|
|
||||||
use std::time::Duration;
|
use std::time::Duration;
|
||||||
|
|
||||||
use crate::error::ExecError;
|
use crate::error::AgentError;
|
||||||
|
|
||||||
use super::PlcProgram;
|
use super::PlcProgram;
|
||||||
|
|
||||||
@@ -27,7 +27,7 @@ pub async fn wait_ready(
|
|||||||
http: &reqwest::Client,
|
http: &reqwest::Client,
|
||||||
base_url: &str,
|
base_url: &str,
|
||||||
budget: Duration,
|
budget: Duration,
|
||||||
) -> Result<(), ExecError> {
|
) -> Result<(), AgentError> {
|
||||||
let login = format!("{base_url}/login");
|
let login = format!("{base_url}/login");
|
||||||
let outcome = tokio::time::timeout(budget, async {
|
let outcome = tokio::time::timeout(budget, async {
|
||||||
loop {
|
loop {
|
||||||
@@ -40,7 +40,7 @@ pub async fn wait_ready(
|
|||||||
}
|
}
|
||||||
})
|
})
|
||||||
.await;
|
.await;
|
||||||
outcome.map_err(|_| ExecError::Other(format!("OpenPLC at {base_url} did not become ready")))
|
outcome.map_err(|_| AgentError::Other(format!("OpenPLC at {base_url} did not become ready")))
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Log in, upload the program, compile it, and start the runtime. On success the
|
/// Log in, upload the program, compile it, and start the runtime. On success the
|
||||||
@@ -52,7 +52,7 @@ pub async fn load_and_start(
|
|||||||
password: &str,
|
password: &str,
|
||||||
program: &PlcProgram,
|
program: &PlcProgram,
|
||||||
compile_budget: Duration,
|
compile_budget: Duration,
|
||||||
) -> Result<(), ExecError> {
|
) -> Result<(), AgentError> {
|
||||||
login(http, base_url, user, password).await?;
|
login(http, base_url, user, password).await?;
|
||||||
let prog_file = upload_program(http, base_url, program).await?;
|
let prog_file = upload_program(http, base_url, program).await?;
|
||||||
save_program(http, base_url, &prog_file).await?;
|
save_program(http, base_url, &prog_file).await?;
|
||||||
@@ -68,14 +68,14 @@ async fn login(
|
|||||||
base_url: &str,
|
base_url: &str,
|
||||||
user: &str,
|
user: &str,
|
||||||
password: &str,
|
password: &str,
|
||||||
) -> Result<(), ExecError> {
|
) -> Result<(), AgentError> {
|
||||||
let resp = http
|
let resp = http
|
||||||
.post(format!("{base_url}/login"))
|
.post(format!("{base_url}/login"))
|
||||||
.form(&[("username", user), ("password", password)])
|
.form(&[("username", user), ("password", password)])
|
||||||
.send()
|
.send()
|
||||||
.await?;
|
.await?;
|
||||||
if resp.status().is_server_error() {
|
if resp.status().is_server_error() {
|
||||||
return Err(ExecError::Other(format!(
|
return Err(AgentError::Other(format!(
|
||||||
"OpenPLC login failed: HTTP {}",
|
"OpenPLC login failed: HTTP {}",
|
||||||
resp.status()
|
resp.status()
|
||||||
)));
|
)));
|
||||||
@@ -90,7 +90,7 @@ async fn upload_program(
|
|||||||
http: &reqwest::Client,
|
http: &reqwest::Client,
|
||||||
base_url: &str,
|
base_url: &str,
|
||||||
program: &PlcProgram,
|
program: &PlcProgram,
|
||||||
) -> Result<String, ExecError> {
|
) -> Result<String, AgentError> {
|
||||||
let part = reqwest::multipart::Part::text(program.source.clone())
|
let part = reqwest::multipart::Part::text(program.source.clone())
|
||||||
.file_name(program.file_name.clone())
|
.file_name(program.file_name.clone())
|
||||||
.mime_str("application/octet-stream")?;
|
.mime_str("application/octet-stream")?;
|
||||||
@@ -102,7 +102,7 @@ async fn upload_program(
|
|||||||
.await?;
|
.await?;
|
||||||
let html = resp.text().await?;
|
let html = resp.text().await?;
|
||||||
parse_prog_file(&html).ok_or_else(|| {
|
parse_prog_file(&html).ok_or_else(|| {
|
||||||
ExecError::Other("OpenPLC upload did not return a prog_file handle".to_string())
|
AgentError::Other("OpenPLC upload did not return a prog_file handle".to_string())
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -113,7 +113,7 @@ async fn save_program(
|
|||||||
http: &reqwest::Client,
|
http: &reqwest::Client,
|
||||||
base_url: &str,
|
base_url: &str,
|
||||||
prog_file: &str,
|
prog_file: &str,
|
||||||
) -> Result<(), ExecError> {
|
) -> Result<(), AgentError> {
|
||||||
let epoch = std::time::SystemTime::now()
|
let epoch = std::time::SystemTime::now()
|
||||||
.duration_since(std::time::UNIX_EPOCH)
|
.duration_since(std::time::UNIX_EPOCH)
|
||||||
.map(|d| d.as_secs())
|
.map(|d| d.as_secs())
|
||||||
@@ -130,7 +130,7 @@ async fn save_program(
|
|||||||
.send()
|
.send()
|
||||||
.await?;
|
.await?;
|
||||||
if resp.status().is_server_error() {
|
if resp.status().is_server_error() {
|
||||||
return Err(ExecError::Other(format!(
|
return Err(AgentError::Other(format!(
|
||||||
"OpenPLC save-program failed: HTTP {}",
|
"OpenPLC save-program failed: HTTP {}",
|
||||||
resp.status()
|
resp.status()
|
||||||
)));
|
)));
|
||||||
@@ -146,7 +146,7 @@ async fn compile(
|
|||||||
base_url: &str,
|
base_url: &str,
|
||||||
prog_file: &str,
|
prog_file: &str,
|
||||||
budget: Duration,
|
budget: Duration,
|
||||||
) -> Result<(), ExecError> {
|
) -> Result<(), AgentError> {
|
||||||
http.get(format!("{base_url}/compile-program"))
|
http.get(format!("{base_url}/compile-program"))
|
||||||
.query(&[("file", prog_file)])
|
.query(&[("file", prog_file)])
|
||||||
.send()
|
.send()
|
||||||
@@ -168,20 +168,20 @@ async fn compile(
|
|||||||
.await;
|
.await;
|
||||||
match outcome {
|
match outcome {
|
||||||
Ok(true) => Ok(()),
|
Ok(true) => Ok(()),
|
||||||
Ok(false) => Err(ExecError::Other(
|
Ok(false) => Err(AgentError::Other(
|
||||||
"OpenPLC compilation finished with errors".to_string(),
|
"OpenPLC compilation finished with errors".to_string(),
|
||||||
)),
|
)),
|
||||||
Err(_) => Err(ExecError::Other(
|
Err(_) => Err(AgentError::Other(
|
||||||
"OpenPLC compilation did not finish in time".to_string(),
|
"OpenPLC compilation did not finish in time".to_string(),
|
||||||
)),
|
)),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// `GET /start_plc` — starts the runtime, opening Modbus/TCP on 502.
|
/// `GET /start_plc` — starts the runtime, opening Modbus/TCP on 502.
|
||||||
async fn start(http: &reqwest::Client, base_url: &str) -> Result<(), ExecError> {
|
async fn start(http: &reqwest::Client, base_url: &str) -> Result<(), AgentError> {
|
||||||
let resp = http.get(format!("{base_url}/start_plc")).send().await?;
|
let resp = http.get(format!("{base_url}/start_plc")).send().await?;
|
||||||
if resp.status().is_server_error() {
|
if resp.status().is_server_error() {
|
||||||
return Err(ExecError::Other(format!(
|
return Err(AgentError::Other(format!(
|
||||||
"OpenPLC start_plc failed: HTTP {}",
|
"OpenPLC start_plc failed: HTTP {}",
|
||||||
resp.status()
|
resp.status()
|
||||||
)));
|
)));
|
||||||
+6
-6
@@ -15,7 +15,7 @@ use std::time::{SystemTime, UNIX_EPOCH};
|
|||||||
|
|
||||||
use compliance_core::PlcRuntimeConfig;
|
use compliance_core::PlcRuntimeConfig;
|
||||||
|
|
||||||
use crate::error::ExecError;
|
use crate::error::AgentError;
|
||||||
|
|
||||||
/// The Modbus/TCP port an OpenPLC instance opens once a program is running.
|
/// The Modbus/TCP port an OpenPLC instance opens once a program is running.
|
||||||
const MODBUS_PORT: u16 = 502;
|
const MODBUS_PORT: u16 = 502;
|
||||||
@@ -45,7 +45,7 @@ pub trait SoftPlc {
|
|||||||
fn provision(
|
fn provision(
|
||||||
&self,
|
&self,
|
||||||
target_id: &str,
|
target_id: &str,
|
||||||
) -> impl std::future::Future<Output = Result<ProvisionedRuntime, ExecError>> + Send;
|
) -> impl std::future::Future<Output = Result<ProvisionedRuntime, AgentError>> + Send;
|
||||||
|
|
||||||
/// Tear an instance down. Best-effort and idempotent — never fails the scan.
|
/// Tear an instance down. Best-effort and idempotent — never fails the scan.
|
||||||
fn teardown(&self, handle: &ProvisionedRuntime)
|
fn teardown(&self, handle: &ProvisionedRuntime)
|
||||||
@@ -65,7 +65,7 @@ impl DockerSoftPlc {
|
|||||||
}
|
}
|
||||||
|
|
||||||
impl SoftPlc for DockerSoftPlc {
|
impl SoftPlc for DockerSoftPlc {
|
||||||
async fn provision(&self, target_id: &str) -> Result<ProvisionedRuntime, ExecError> {
|
async fn provision(&self, target_id: &str) -> Result<ProvisionedRuntime, AgentError> {
|
||||||
// Best-effort sweep of any container leaked by a crashed earlier run
|
// Best-effort sweep of any container leaked by a crashed earlier run
|
||||||
// before we add another. Only removes instances past their max lifetime,
|
// before we add another. Only removes instances past their max lifetime,
|
||||||
// so it can never disturb a concurrent run.
|
// so it can never disturb a concurrent run.
|
||||||
@@ -75,7 +75,7 @@ impl SoftPlc for DockerSoftPlc {
|
|||||||
let args = run_args(&self.cfg, &name, target_id);
|
let args = run_args(&self.cfg, &name, target_id);
|
||||||
let out = run_docker(&args).await?;
|
let out = run_docker(&args).await?;
|
||||||
if !out.status.success() {
|
if !out.status.success() {
|
||||||
return Err(ExecError::Other(format!(
|
return Err(AgentError::Other(format!(
|
||||||
"docker run for soft-PLC {name} failed: {}",
|
"docker run for soft-PLC {name} failed: {}",
|
||||||
String::from_utf8_lossy(&out.stderr).trim()
|
String::from_utf8_lossy(&out.stderr).trim()
|
||||||
)));
|
)));
|
||||||
@@ -215,12 +215,12 @@ async fn reap_stale(cfg: &PlcRuntimeConfig, now: u64) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
/// Run a `docker` subcommand, capturing its output.
|
/// Run a `docker` subcommand, capturing its output.
|
||||||
async fn run_docker(args: &[String]) -> Result<std::process::Output, ExecError> {
|
async fn run_docker(args: &[String]) -> Result<std::process::Output, AgentError> {
|
||||||
tokio::process::Command::new("docker")
|
tokio::process::Command::new("docker")
|
||||||
.args(args)
|
.args(args)
|
||||||
.output()
|
.output()
|
||||||
.await
|
.await
|
||||||
.map_err(ExecError::Io)
|
.map_err(AgentError::Io)
|
||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
@@ -1,10 +0,0 @@
|
|||||||
//! Werkbank control-plane: the dynamic-execution job queue.
|
|
||||||
//!
|
|
||||||
//! The control plane enqueues declarative [`Job`](compliance_core::models::werkbank::Job)s
|
|
||||||
//! and Werkbank runners lease, run, and complete them. [`queue::JobQueue`] is the
|
|
||||||
//! Mongo-backed queue behind that flow (WB-02); the runner-facing HTTP transport
|
|
||||||
//! and the runner itself land in later stories.
|
|
||||||
|
|
||||||
pub mod queue;
|
|
||||||
|
|
||||||
pub use queue::{JobQueue, SweepOutcome};
|
|
||||||
@@ -1,309 +0,0 @@
|
|||||||
//! The Mongo-backed Werkbank job queue (WB-02).
|
|
||||||
//!
|
|
||||||
//! A pull queue: the control plane [`enqueue`](JobQueue::enqueue)s jobs; a runner
|
|
||||||
//! [`lease`](JobQueue::lease)s the oldest queued job it can run (matched by
|
|
||||||
//! executor + labels), [`heartbeat`](JobQueue::heartbeat)s while it works, and
|
|
||||||
//! [`complete`](JobQueue::complete)s it. Leases carry a visibility timeout: if a
|
|
||||||
//! runner dies mid-job its heartbeats stop, the lease expires, and
|
|
||||||
//! [`sweep_expired`](JobQueue::sweep_expired) returns the job to `queued` (or
|
|
||||||
//! `expired` once it has been retried too many times).
|
|
||||||
//!
|
|
||||||
//! All state transitions are single atomic Mongo updates guarded by the lease
|
|
||||||
//! token, so two runners can never both own a job. Every operation takes an
|
|
||||||
//! explicit `now` so the queue's time-dependent behaviour is deterministically
|
|
||||||
//! testable.
|
|
||||||
|
|
||||||
use std::time::Duration;
|
|
||||||
|
|
||||||
use chrono::{DateTime, Utc};
|
|
||||||
use mongodb::bson::{doc, Bson, DateTime as BsonDateTime};
|
|
||||||
use mongodb::error::{ErrorKind, WriteFailure};
|
|
||||||
use mongodb::options::ReturnDocument;
|
|
||||||
use mongodb::Collection;
|
|
||||||
|
|
||||||
use compliance_core::models::werkbank::{
|
|
||||||
Executor, HeartbeatAck, Job, JobRecord, JobResult, JobStatus, LeasedJob,
|
|
||||||
};
|
|
||||||
|
|
||||||
use crate::database::Database;
|
|
||||||
use crate::error::AgentError;
|
|
||||||
|
|
||||||
/// The non-terminal states a job can be swept or cancelled from.
|
|
||||||
const ACTIVE_STATES: [&str; 2] = ["leased", "running"];
|
|
||||||
/// Every terminal state (no further transitions).
|
|
||||||
const TERMINAL_STATES: [&str; 4] = ["succeeded", "failed", "expired", "cancelled"];
|
|
||||||
|
|
||||||
/// What a visibility-timeout sweep did.
|
|
||||||
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
|
|
||||||
pub struct SweepOutcome {
|
|
||||||
/// Expired-lease jobs returned to `queued` for another runner.
|
|
||||||
pub requeued: u64,
|
|
||||||
/// Jobs that had exhausted their attempts and were marked `expired`.
|
|
||||||
pub expired: u64,
|
|
||||||
}
|
|
||||||
|
|
||||||
/// The Mongo-backed job queue.
|
|
||||||
pub struct JobQueue {
|
|
||||||
coll: Collection<JobRecord>,
|
|
||||||
}
|
|
||||||
|
|
||||||
impl JobQueue {
|
|
||||||
/// Build a queue over a tenant database's `werkbank_jobs` collection.
|
|
||||||
pub fn new(db: &Database) -> Self {
|
|
||||||
Self {
|
|
||||||
coll: db.werkbank_jobs(),
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Enqueue a job. Idempotent by job id: a job that is already present is a
|
|
||||||
/// no-op. Returns `true` if this call inserted it, `false` if it existed.
|
|
||||||
pub async fn enqueue(&self, job: Job, now: DateTime<Utc>) -> Result<bool, AgentError> {
|
|
||||||
let record = JobRecord::queued(job, now);
|
|
||||||
match self.coll.insert_one(&record).await {
|
|
||||||
Ok(_) => Ok(true),
|
|
||||||
Err(e) if is_duplicate_key(&e) => Ok(false),
|
|
||||||
Err(e) => Err(e.into()),
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Atomically lease the oldest `queued` job this runner can run — matched by
|
|
||||||
/// executor and by labels (every label the job requires must be one the
|
|
||||||
/// runner advertises). Returns the job plus a lease token, or `None` if
|
|
||||||
/// nothing is runnable.
|
|
||||||
pub async fn lease(
|
|
||||||
&self,
|
|
||||||
runner_id: &str,
|
|
||||||
executor: Executor,
|
|
||||||
runner_labels: &[String],
|
|
||||||
lease_ttl: Duration,
|
|
||||||
now: DateTime<Utc>,
|
|
||||||
) -> Result<Option<LeasedJob>, AgentError> {
|
|
||||||
let token = uuid::Uuid::new_v4().to_string();
|
|
||||||
let expires = bson_dt(now + ttl(lease_ttl));
|
|
||||||
let executor_bson = mongodb::bson::to_bson(&executor).unwrap_or(Bson::Null);
|
|
||||||
|
|
||||||
let filter = doc! {
|
|
||||||
"status": "queued",
|
|
||||||
"cancel_requested": { "$ne": true },
|
|
||||||
"job.executor": executor_bson,
|
|
||||||
// Every label the job requires must be in the runner's set — i.e. the
|
|
||||||
// job has no label that is not offered by the runner. Absent/empty
|
|
||||||
// job labels match any runner.
|
|
||||||
"job.labels": { "$not": { "$elemMatch": { "$nin": runner_labels.to_vec() } } },
|
|
||||||
};
|
|
||||||
let update = doc! {
|
|
||||||
"$set": {
|
|
||||||
"status": "leased",
|
|
||||||
"lease_token": &token,
|
|
||||||
"leased_by": runner_id,
|
|
||||||
"lease_expires_at": expires,
|
|
||||||
"heartbeat_at": bson_dt(now),
|
|
||||||
"updated_at": bson_dt(now),
|
|
||||||
},
|
|
||||||
"$inc": { "attempts": 1 },
|
|
||||||
};
|
|
||||||
|
|
||||||
let record = self
|
|
||||||
.coll
|
|
||||||
.find_one_and_update(filter, update)
|
|
||||||
.sort(doc! { "created_at": 1 }) // FIFO
|
|
||||||
.return_document(ReturnDocument::After)
|
|
||||||
.await?;
|
|
||||||
Ok(record.map(|r| LeasedJob {
|
|
||||||
job: r.job,
|
|
||||||
lease_token: token,
|
|
||||||
}))
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Extend a lease and report whether the job has been asked to cancel.
|
|
||||||
/// Transitions the job to `running` on the first heartbeat. Returns `None`
|
|
||||||
/// when the lease is no longer valid (token mismatch, or the job is already
|
|
||||||
/// terminal) — the runner should then abandon the work.
|
|
||||||
pub async fn heartbeat(
|
|
||||||
&self,
|
|
||||||
job_id: &str,
|
|
||||||
lease_token: &str,
|
|
||||||
lease_ttl: Duration,
|
|
||||||
now: DateTime<Utc>,
|
|
||||||
) -> Result<Option<HeartbeatAck>, AgentError> {
|
|
||||||
let filter = doc! {
|
|
||||||
"job.id": job_id,
|
|
||||||
"lease_token": lease_token,
|
|
||||||
"status": { "$in": ACTIVE_STATES.to_vec() },
|
|
||||||
};
|
|
||||||
let update = doc! {
|
|
||||||
"$set": {
|
|
||||||
"status": "running",
|
|
||||||
"lease_expires_at": bson_dt(now + ttl(lease_ttl)),
|
|
||||||
"heartbeat_at": bson_dt(now),
|
|
||||||
"updated_at": bson_dt(now),
|
|
||||||
},
|
|
||||||
};
|
|
||||||
let record = self
|
|
||||||
.coll
|
|
||||||
.find_one_and_update(filter, update)
|
|
||||||
.return_document(ReturnDocument::After)
|
|
||||||
.await?;
|
|
||||||
Ok(record.map(|r| HeartbeatAck {
|
|
||||||
cancelled: r.cancel_requested,
|
|
||||||
}))
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Record a job's terminal result. Guarded by the lease token and only from
|
|
||||||
/// an active (`leased`/`running`) state, so it is idempotent — a duplicate or
|
|
||||||
/// late submission after the job already finished matches nothing. Returns
|
|
||||||
/// `true` if this call recorded the result.
|
|
||||||
pub async fn complete(
|
|
||||||
&self,
|
|
||||||
job_id: &str,
|
|
||||||
lease_token: &str,
|
|
||||||
result: &JobResult,
|
|
||||||
now: DateTime<Utc>,
|
|
||||||
) -> Result<bool, AgentError> {
|
|
||||||
let status = result.status.unwrap_or(JobStatus::Failed);
|
|
||||||
let status_bson = mongodb::bson::to_bson(&status).unwrap_or(Bson::String("failed".into()));
|
|
||||||
let result_bson =
|
|
||||||
mongodb::bson::to_bson(result).map_err(|e| AgentError::Other(e.to_string()))?;
|
|
||||||
|
|
||||||
let filter = doc! {
|
|
||||||
"job.id": job_id,
|
|
||||||
"lease_token": lease_token,
|
|
||||||
"status": { "$in": ACTIVE_STATES.to_vec() },
|
|
||||||
};
|
|
||||||
let update = doc! {
|
|
||||||
"$set": {
|
|
||||||
"status": status_bson,
|
|
||||||
"result": result_bson,
|
|
||||||
"lease_token": Bson::Null,
|
|
||||||
"lease_expires_at": Bson::Null,
|
|
||||||
"updated_at": bson_dt(now),
|
|
||||||
},
|
|
||||||
};
|
|
||||||
let res = self.coll.update_one(filter, update).await?;
|
|
||||||
Ok(res.modified_count == 1)
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Request cancellation of a job. A still-`queued` job is cancelled outright;
|
|
||||||
/// an in-flight one is flagged so the runner sees it on its next heartbeat and
|
|
||||||
/// tears down. Returns `true` if a non-terminal job matched.
|
|
||||||
pub async fn cancel(&self, job_id: &str, now: DateTime<Utc>) -> Result<bool, AgentError> {
|
|
||||||
let filter = doc! {
|
|
||||||
"job.id": job_id,
|
|
||||||
"status": { "$nin": TERMINAL_STATES.to_vec() },
|
|
||||||
};
|
|
||||||
// Pipeline update: flag cancellation, and if still queued flip straight to
|
|
||||||
// cancelled (nothing is running it).
|
|
||||||
let pipeline = vec![doc! {
|
|
||||||
"$set": {
|
|
||||||
"cancel_requested": true,
|
|
||||||
"status": {
|
|
||||||
"$cond": [ { "$eq": ["$status", "queued"] }, "cancelled", "$status" ]
|
|
||||||
},
|
|
||||||
"updated_at": bson_dt(now),
|
|
||||||
}
|
|
||||||
}];
|
|
||||||
let res = self.coll.update_one(filter, pipeline).await?;
|
|
||||||
Ok(res.matched_count == 1)
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Sweep leases whose visibility timeout has elapsed: return them to `queued`
|
|
||||||
/// for another runner, or mark them `expired` once they have been leased
|
|
||||||
/// `max_attempts` times. This is what makes a crashed runner's job recover.
|
|
||||||
pub async fn sweep_expired(
|
|
||||||
&self,
|
|
||||||
now: DateTime<Utc>,
|
|
||||||
max_attempts: u32,
|
|
||||||
// (kept explicit rather than a const so callers can tune retry policy)
|
|
||||||
) -> Result<SweepOutcome, AgentError> {
|
|
||||||
let now_bson = bson_dt(now);
|
|
||||||
let max = i64::from(max_attempts);
|
|
||||||
|
|
||||||
let requeue = self
|
|
||||||
.coll
|
|
||||||
.update_many(
|
|
||||||
doc! {
|
|
||||||
"status": { "$in": ACTIVE_STATES.to_vec() },
|
|
||||||
"lease_expires_at": { "$lt": &now_bson },
|
|
||||||
"attempts": { "$lt": max },
|
|
||||||
},
|
|
||||||
doc! { "$set": {
|
|
||||||
"status": "queued",
|
|
||||||
"lease_token": Bson::Null,
|
|
||||||
"leased_by": Bson::Null,
|
|
||||||
"lease_expires_at": Bson::Null,
|
|
||||||
"updated_at": &now_bson,
|
|
||||||
} },
|
|
||||||
)
|
|
||||||
.await?;
|
|
||||||
|
|
||||||
let expire = self
|
|
||||||
.coll
|
|
||||||
.update_many(
|
|
||||||
doc! {
|
|
||||||
"status": { "$in": ACTIVE_STATES.to_vec() },
|
|
||||||
"lease_expires_at": { "$lt": &now_bson },
|
|
||||||
"attempts": { "$gte": max },
|
|
||||||
},
|
|
||||||
doc! { "$set": {
|
|
||||||
"status": "expired",
|
|
||||||
"lease_token": Bson::Null,
|
|
||||||
"lease_expires_at": Bson::Null,
|
|
||||||
"updated_at": &now_bson,
|
|
||||||
} },
|
|
||||||
)
|
|
||||||
.await?;
|
|
||||||
|
|
||||||
Ok(SweepOutcome {
|
|
||||||
requeued: requeue.modified_count,
|
|
||||||
expired: expire.modified_count,
|
|
||||||
})
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Fetch a job record by job id (inspection / control-plane reads).
|
|
||||||
pub async fn get(&self, job_id: &str) -> Result<Option<JobRecord>, AgentError> {
|
|
||||||
Ok(self.coll.find_one(doc! { "job.id": job_id }).await?)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/// A `chrono::Duration` for a lease TTL, saturating rather than panicking on an
|
|
||||||
/// absurd input (`chrono::Duration::seconds` panics past its internal bound).
|
|
||||||
fn ttl(d: Duration) -> chrono::Duration {
|
|
||||||
let secs = i64::try_from(d.as_secs()).unwrap_or(i64::MAX);
|
|
||||||
chrono::Duration::try_seconds(secs).unwrap_or(chrono::Duration::MAX)
|
|
||||||
}
|
|
||||||
|
|
||||||
/// A chrono instant as a BSON date (so Mongo stores/compares it as a real date).
|
|
||||||
fn bson_dt(dt: DateTime<Utc>) -> BsonDateTime {
|
|
||||||
BsonDateTime::from_chrono(dt)
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Whether a Mongo error is a duplicate-key (E11000) violation — a job with this
|
|
||||||
/// id is already enqueued.
|
|
||||||
fn is_duplicate_key(e: &mongodb::error::Error) -> bool {
|
|
||||||
match &*e.kind {
|
|
||||||
ErrorKind::Write(WriteFailure::WriteError(we)) => we.code == 11000,
|
|
||||||
_ => false,
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
#[cfg(test)]
|
|
||||||
mod tests {
|
|
||||||
use super::*;
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn ttl_saturates_and_converts() {
|
|
||||||
assert_eq!(ttl(Duration::from_secs(30)), chrono::Duration::seconds(30));
|
|
||||||
// An absurd TTL saturates instead of panicking.
|
|
||||||
assert_eq!(ttl(Duration::from_secs(u64::MAX)), chrono::Duration::MAX);
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn state_constants_are_disjoint() {
|
|
||||||
for s in ACTIVE_STATES {
|
|
||||||
assert!(
|
|
||||||
!TERMINAL_STATES.contains(&s),
|
|
||||||
"{s} cannot be both active and terminal"
|
|
||||||
);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -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"
|
|
||||||
|
|||||||
@@ -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,7 +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::{
|
|
||||||
DastCollect, Executor, HeartbeatAck, InputRef, Job, JobCollect, JobRecord, JobResult,
|
|
||||||
JobRuntime, JobStatus, JobType, LeasedJob,
|
|
||||||
};
|
|
||||||
|
|||||||
@@ -1,461 +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,
|
|
||||||
}
|
|
||||||
|
|
||||||
#[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"));
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -1,23 +0,0 @@
|
|||||||
[package]
|
|
||||||
name = "werkbank-exec"
|
|
||||||
version = "0.1.0"
|
|
||||||
edition = "2021"
|
|
||||||
description = "Shared dynamic-execution logic: soft-PLC provisioning + industrial-protocol probing, used by the compliance agent and the Werkbank runner."
|
|
||||||
|
|
||||||
[lints]
|
|
||||||
workspace = true
|
|
||||||
|
|
||||||
[dependencies]
|
|
||||||
compliance-core = { workspace = true }
|
|
||||||
compliance-dast = { path = "../compliance-dast" }
|
|
||||||
tokio = { workspace = true }
|
|
||||||
reqwest = { workspace = true }
|
|
||||||
uuid = { workspace = true }
|
|
||||||
regex = { workspace = true }
|
|
||||||
secrecy = { workspace = true }
|
|
||||||
sha2 = { workspace = true }
|
|
||||||
hex = { workspace = true }
|
|
||||||
tracing = { workspace = true }
|
|
||||||
thiserror = { workspace = true }
|
|
||||||
walkdir = "2"
|
|
||||||
futures-util = "0.3"
|
|
||||||
@@ -1,16 +0,0 @@
|
|||||||
//! Error type for the dynamic-execution logic.
|
|
||||||
|
|
||||||
/// Anything that can go wrong provisioning and testing a soft-PLC. The compliance
|
|
||||||
/// agent maps this into its own `AgentError` at the call boundary.
|
|
||||||
#[derive(thiserror::Error, Debug)]
|
|
||||||
pub enum ExecError {
|
|
||||||
/// An HTTP request (to OpenPLC) failed.
|
|
||||||
#[error("HTTP error: {0}")]
|
|
||||||
Http(#[from] reqwest::Error),
|
|
||||||
/// A local IO / process error (e.g. invoking `docker`).
|
|
||||||
#[error("IO error: {0}")]
|
|
||||||
Io(#[from] std::io::Error),
|
|
||||||
/// Any other failure, with a message.
|
|
||||||
#[error("{0}")]
|
|
||||||
Other(String),
|
|
||||||
}
|
|
||||||
@@ -1,32 +0,0 @@
|
|||||||
//! Finding fingerprint helper (a SHA-256 over the salient parts), shared by the
|
|
||||||
//! probe modules for stable dedup keys. Mirrors the agent's `dedup` helper.
|
|
||||||
|
|
||||||
use sha2::{Digest, Sha256};
|
|
||||||
|
|
||||||
/// A stable fingerprint over the given parts (order-sensitive, separated so
|
|
||||||
/// `["ab","c"]` and `["a","bc"]` differ).
|
|
||||||
pub fn compute_fingerprint(parts: &[&str]) -> String {
|
|
||||||
let mut hasher = Sha256::new();
|
|
||||||
for part in parts {
|
|
||||||
hasher.update(part.as_bytes());
|
|
||||||
hasher.update(b"|");
|
|
||||||
}
|
|
||||||
hex::encode(hasher.finalize())
|
|
||||||
}
|
|
||||||
|
|
||||||
#[cfg(test)]
|
|
||||||
mod tests {
|
|
||||||
use super::*;
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn deterministic_and_hex() {
|
|
||||||
let a = compute_fingerprint(&["repo", "rule", "1"]);
|
|
||||||
assert_eq!(a, compute_fingerprint(&["repo", "rule", "1"]));
|
|
||||||
assert_eq!(a.len(), 64);
|
|
||||||
assert!(a.chars().all(|c| c.is_ascii_hexdigit()));
|
|
||||||
assert_ne!(
|
|
||||||
compute_fingerprint(&["ab", "c"]),
|
|
||||||
compute_fingerprint(&["a", "bc"])
|
|
||||||
);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -1,18 +0,0 @@
|
|||||||
//! Shared dynamic-execution logic for Werkbank.
|
|
||||||
//!
|
|
||||||
//! The soft-PLC provisioning + industrial-protocol probing that turns a control-
|
|
||||||
//! logic artifact into findings: provision an ephemeral OpenPLC, load the program,
|
|
||||||
//! start it, probe it over Modbus/OPC-UA/EtherNet-IP, DAST its web endpoint, tear
|
|
||||||
//! it down. Extracted from the compliance agent (#183) so both the agent (in
|
|
||||||
//! process) and the Werkbank runner (WB-04) run identical logic.
|
|
||||||
//!
|
|
||||||
//! - [`ics`] — read-only industrial-protocol probing.
|
|
||||||
//! - [`plc`] — ephemeral soft-PLC provisioning + the provision-and-test loop.
|
|
||||||
|
|
||||||
pub mod error;
|
|
||||||
pub mod ics;
|
|
||||||
pub mod plc;
|
|
||||||
|
|
||||||
mod fingerprint;
|
|
||||||
|
|
||||||
pub use error::ExecError;
|
|
||||||
Reference in New Issue
Block a user