Compare commits

..
Author SHA1 Message Date
Sharang ParnerkarandClaude Fable 5 420af3af9e feat(werkbank): runner queue endpoints + result persistence (WB-05)
CI / Deploy Agent (pull_request) Has been skipped
CI / Deploy Dashboard (pull_request) Has been skipped
CI / Deploy Docs (pull_request) Has been skipped
CI / Check (pull_request) Successful in 5m35s
CI / Detect Changes (pull_request) Has been skipped
CI / Deploy MCP (pull_request) Has been skipped
The control-plane side of the pull API (implements sharang/werkbank#6), so a
Werkbank runner can reach a real queue end-to-end:

- POST /api/v1/werkbank/jobs/{lease,heartbeat,complete} — thin handlers over the
  JobQueue (WB-02), tenant-scoped from the request's `tenant` (db_pool
  .for_tenant_id). lease→204 when empty; heartbeat→409 on a lost lease.
- Machine auth: a static WERKBANK_RUNNER_TOKEN bearer (require_runner_token),
  mounted only when the token is set — like the admin API, and NOT a Keycloak JWT
  (a runner acts across tenants). /api/v1/werkbank/* is added to the JWT
  PUBLIC_PREFIXES so it routes to the runner-token gate, not the customer-JWT one.
- On completion, the runner's findings + DAST findings are persisted against the
  job's target (dedup'd by fingerprint), so a job run by a remote runner lands
  the same findings an in-process run would.
- Shared transport types (LeaseRequest/HeartbeatRequest/CompleteRequest/
  CompleteResponse) live in compliance-core so the runner (client) and control
  plane (server) agree on shapes.

Tests: 3 HTTP integration tests against a live Mongo (lease→complete→persist,
empty-queue 204, and the bearer-token gate) + a token-compare unit test. Skips
cleanly with no Mongo. clippy + fmt clean.

Follow-up: wiring the scan pipeline to enqueue plc-provision jobs needs an
artifact-fetch path for the runner (so it can pull the program blob); tracked
with the on-prem work.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-17 12:28:02 +02:00
31 changed files with 606 additions and 202 deletions
+2 -4
View File
@@ -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
-20
View File
@@ -699,7 +699,6 @@ dependencies = [
"urlencoding", "urlencoding",
"uuid", "uuid",
"walkdir", "walkdir",
"werkbank-exec",
"zip", "zip",
] ]
@@ -6715,25 +6714,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"
-1
View File
@@ -7,7 +7,6 @@ members = [
"compliance-dast", "compliance-dast",
"compliance-mcp", "compliance-mcp",
"compliance-smoke", "compliance-smoke",
"werkbank-exec",
] ]
resolver = "2" resolver = "2"
-3
View File
@@ -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
+1
View File
@@ -14,6 +14,7 @@ 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::*;
@@ -0,0 +1,193 @@
//! 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"));
}
}
+28 -1
View File
@@ -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}; use axum::routing::{delete, get, post};
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,8 +72,35 @@ 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
View File
@@ -65,6 +65,7 @@ pub fn load_config() -> Result<AgentConfig, AgentError> {
admin_api_token: env_secret_opt("ADMIN_API_TOKEN"), admin_api_token: env_secret_opt("ADMIN_API_TOKEN"),
tenant_registry_url: env_var_opt("TENANT_REGISTRY_URL"), tenant_registry_url: env_var_opt("TENANT_REGISTRY_URL"),
plc_runtime: load_plc_runtime_config(), plc_runtime: load_plc_runtime_config(),
werkbank_runner_token: env_secret_opt("WERKBANK_RUNNER_TOKEN"),
}) })
} }
-3
View File
@@ -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),
} }
+1
View File
@@ -343,6 +343,7 @@ mod tests {
admin_api_token: None, admin_api_token: None,
tenant_registry_url: None, tenant_registry_url: None,
plc_runtime: compliance_core::PlcRuntimeConfig::default(), plc_runtime: compliance_core::PlcRuntimeConfig::default(),
werkbank_runner_token: None,
} }
} }
@@ -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;
+1
View File
@@ -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,
+1
View File
@@ -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.
@@ -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()
))); )));
@@ -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)]
+53 -39
View File
@@ -2,6 +2,10 @@
// //
// 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;
@@ -11,6 +15,54 @@ 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,
@@ -33,45 +85,7 @@ impl TestServer {
.await .await
.expect("Failed to build DatabasePool"); .expect("Failed to build DatabasePool");
let config = AgentConfig { let config = dev_config(mongodb_uri.clone(), db_name.clone());
mongodb_uri: mongodb_uri.clone(),
mongodb_database: db_name.clone(),
litellm_url: std::env::var("TEST_LITELLM_URL")
.unwrap_or_else(|_| "http://localhost:4000".into()),
litellm_api_key: SecretString::from(String::new()),
litellm_model: "gpt-4o".into(),
litellm_embed_model: "text-embedding-3-small".into(),
agent_port: 0, // not used — we bind ourselves
scan_schedule: String::new(),
cve_monitor_schedule: String::new(),
git_clone_base_path: "/tmp/compliance-scanner-tests/repos".into(),
artifact_store_base_path: "/tmp/compliance-scanner-tests/artifacts".into(),
ssh_key_path: "/tmp/compliance-scanner-tests/ssh/id_ed25519".into(),
github_token: None,
github_webhook_secret: None,
gitlab_url: None,
gitlab_token: None,
gitlab_webhook_secret: None,
jira_url: None,
jira_email: None,
jira_api_token: None,
jira_project_key: None,
searxng_url: None,
nvd_api_key: None,
keycloak_url: None,
keycloak_realm: None,
keycloak_admin_username: None,
keycloak_admin_password: None,
pentest_verification_email: None,
pentest_imap_host: None,
pentest_imap_port: None,
pentest_imap_tls: false,
pentest_imap_username: None,
pentest_imap_password: None,
admin_api_token: None,
tenant_registry_url: None,
plc_runtime: compliance_core::PlcRuntimeConfig::default(),
};
let agent = ComplianceAgent::new(config, db_pool); let agent = ComplianceAgent::new(config, db_pool);
+222
View File
@@ -0,0 +1,222 @@
//! 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;
}
+5 -5
View File
@@ -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/*`) has its own static-bearer middleware and must /// (`/api/v1/admin/*`) and the Werkbank runner API (`/api/v1/werkbank/*`)
/// not be routed through the customer-JWT path — a Keycloak token /// have their own static-bearer middleware and must not be routed through the
/// always carries a single tenant_id and would semantically conflict /// customer-JWT path — a Keycloak token always carries a single tenant_id and
/// with cross-tenant admin operations. /// would semantically conflict with these cross-tenant / machine operations.
const PUBLIC_PREFIXES: &[&str] = &["/api/v1/admin/"]; const PUBLIC_PREFIXES: &[&str] = &["/api/v1/admin/", "/api/v1/werkbank/"];
/// 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.
+5
View File
@@ -53,6 +53,11 @@ pub struct AgentConfig {
/// default: it needs Docker access in the agent's runtime, which is a /// default: it needs Docker access in the agent's runtime, which is a
/// deployment opt-in. /// deployment opt-in.
pub plc_runtime: PlcRuntimeConfig, pub plc_runtime: PlcRuntimeConfig,
/// Static bearer for the Werkbank runner endpoints
/// (`/api/v1/werkbank/jobs/*`). Machine auth for runners leasing/completing
/// jobs — NOT a Keycloak JWT, since a runner acts across tenants. When
/// `None`, those endpoints are not mounted at all.
pub werkbank_runner_token: Option<SecretString>,
} }
/// Configuration for the ephemeral soft-PLC "provision-and-test" path (#183). /// Configuration for the ephemeral soft-PLC "provision-and-test" path (#183).
+3 -2
View File
@@ -49,6 +49,7 @@ pub use repository::ScanTrigger;
pub use sbom::{SbomEntry, VulnRef}; pub use sbom::{SbomEntry, VulnRef};
pub use scan::{ScanPhase, ScanRun, ScanRunStatus, ScanType}; pub use scan::{ScanPhase, ScanRun, ScanRunStatus, ScanType};
pub use werkbank::{ pub use werkbank::{
DastCollect, Executor, HeartbeatAck, InputRef, Job, JobCollect, JobRecord, JobResult, CompleteRequest, CompleteResponse, DastCollect, Executor, HeartbeatAck, HeartbeatRequest,
JobRuntime, JobStatus, JobType, LeasedJob, InputRef, Job, JobCollect, JobRecord, JobResult, JobRuntime, JobStatus, JobType, LeaseRequest,
LeasedJob,
}; };
+53
View File
@@ -349,6 +349,59 @@ pub struct HeartbeatAck {
pub cancelled: bool, pub cancelled: bool,
} }
// --- Runner ↔ control-plane transport (the pull API wire types) ---------------
// Shared so the runner (client) and the control plane (server) agree on shapes.
/// Runner → control plane: lease the oldest runnable job for this runner.
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct LeaseRequest {
/// The tenant queue to lease from.
pub tenant: String,
/// The runner id (advertised for attribution).
pub runner_id: String,
/// The executor this runner provides.
pub executor: Executor,
/// The capability labels this runner advertises.
#[serde(default)]
pub labels: Vec<String>,
/// Requested lease lifetime (the visibility timeout), in seconds.
pub lease_ttl_secs: u64,
}
/// Runner → control plane: prove lease ownership and extend it.
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct HeartbeatRequest {
/// The tenant queue.
pub tenant: String,
/// The job being worked.
pub job_id: String,
/// The lease token from the [`LeasedJob`].
pub lease_token: String,
/// Lease lifetime to extend to, in seconds.
pub lease_ttl_secs: u64,
}
/// Runner → control plane: record a job's terminal result.
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct CompleteRequest {
/// The tenant queue.
pub tenant: String,
/// The job being completed.
pub job_id: String,
/// The lease token proving ownership.
pub lease_token: String,
/// The result to record.
pub result: JobResult,
}
/// Control plane → runner: whether the completion was recorded (false if the
/// lease was already lost — token mismatch or the job had become terminal).
#[derive(Debug, Clone, Copy, Serialize, Deserialize)]
pub struct CompleteResponse {
/// Whether the result was recorded.
pub recorded: bool,
}
#[cfg(test)] #[cfg(test)]
#[allow(clippy::expect_used, clippy::unwrap_used)] #[allow(clippy::expect_used, clippy::unwrap_used)]
mod tests { mod tests {
-23
View File
@@ -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"
-16
View File
@@ -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),
}
-32
View File
@@ -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"])
);
}
}
-18
View File
@@ -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;