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
32 changed files with 47 additions and 3624 deletions
-5
View File
@@ -47,11 +47,6 @@ PLC_RUNTIME_MAX_LIFETIME_SECS=180
PLC_RUNTIME_OPENPLC_USER=openplc
PLC_RUNTIME_OPENPLC_PASSWORD=openplc
# Werkbank runner API (/api/v1/werkbank/jobs/*, /api/v1/werkbank/artifacts/*).
# When set, mounts the runner-facing queue + artifact endpoints behind this
# bearer token; runners present the same token. Unset = endpoints not mounted.
WERKBANK_RUNNER_TOKEN=
# Dashboard
DASHBOARD_PORT=8080
AGENT_API_URL=http://localhost:3001
+2 -4
View File
@@ -107,8 +107,6 @@ jobs:
run: cargo clippy -p compliance-dashboard --features web --no-default-features -- -D warnings
- name: Clippy (mcp)
run: cargo clippy -p compliance-mcp -- -D warnings
- name: Clippy (werkbank-exec)
run: cargo clippy -p werkbank-exec -- -D warnings
# Security audit
- name: Security Audit
@@ -117,8 +115,8 @@ jobs:
RUSTC_WRAPPER: ""
# Tests (reuses compilation artifacts from clippy)
- name: Tests (core + agent + werkbank-exec)
run: cargo test -p compliance-core -p compliance-agent -p werkbank-exec --lib
- name: Tests (core + agent)
run: cargo test -p compliance-core -p compliance-agent --lib
- name: Tests (dashboard server)
run: cargo test -p compliance-dashboard --features server --no-default-features
- name: Tests (dashboard web)
Generated
-20
View File
@@ -699,7 +699,6 @@ dependencies = [
"urlencoding",
"uuid",
"walkdir",
"werkbank-exec",
"zip",
]
@@ -6715,25 +6714,6 @@ dependencies = [
"rustls-pki-types",
]
[[package]]
name = "werkbank-exec"
version = "0.1.0"
dependencies = [
"compliance-core",
"compliance-dast",
"futures-util",
"hex",
"regex",
"reqwest",
"secrecy",
"sha2",
"thiserror 2.0.18",
"tokio",
"tracing",
"uuid",
"walkdir",
]
[[package]]
name = "which"
version = "6.0.3"
-1
View File
@@ -7,7 +7,6 @@ members = [
"compliance-dast",
"compliance-mcp",
"compliance-smoke",
"werkbank-exec",
]
resolver = "2"
-3
View File
@@ -10,9 +10,6 @@ workspace = true
compliance-core = { workspace = true, features = ["mongodb", "telemetry", "axum"] }
compliance-graph = { path = "../compliance-graph" }
compliance-dast = { path = "../compliance-dast" }
# Shared dynamic-execution logic (soft-PLC provisioning + ICS probing), also
# used by the Werkbank runner.
werkbank-exec = { path = "../werkbank-exec" }
# Native firmware build/target detection for bare-metal & RTOS artifacts.
# Same-company IP, used directly (not via CLI) so the whole tramiton suite is
# available to the onboarding classifier. NOTE: CI must be able to fetch this
@@ -10,20 +10,18 @@
//! so a job run by a remote runner lands the same findings an in-process run
//! would (WB-05, the control-plane cut-over).
use axum::extract::{Extension, Path, Request};
use axum::extract::{Extension, Request};
use axum::http::{header, StatusCode};
use axum::middleware::Next;
use axum::response::{IntoResponse, Response};
use axum::Json;
use mongodb::bson::{doc, oid::ObjectId};
use mongodb::bson::doc;
use secrecy::ExposeSecret;
use serde::{Deserialize, Serialize};
use std::time::Duration;
use compliance_core::models::werkbank::{
CompleteRequest, CompleteResponse, HeartbeatRequest, InputRef, Job, JobResult, LeaseRequest,
CompleteRequest, CompleteResponse, HeartbeatRequest, JobResult, LeaseRequest,
};
use compliance_core::models::ArtifactKind;
use super::dto::AgentExt;
use crate::database::Database;
@@ -124,100 +122,6 @@ pub async fn complete(
Ok(Json(CompleteResponse { recorded }))
}
/// `GET /api/v1/werkbank/artifacts/{hash}` — serve a content-addressed blob (the
/// program a runner needs to load). The hash is validated against traversal by
/// [`crate::ingest::blob::read_blob`]; a runner fetches this for a job's `blob`
/// input.
#[tracing::instrument(skip_all, fields(hash = %hash))]
pub async fn serve_artifact(
Extension(agent): AgentExt,
Path(hash): Path<String>,
) -> Result<Response, StatusCode> {
let base = std::path::Path::new(&agent.config.artifact_store_base_path);
match crate::ingest::blob::read_blob(base, &hash) {
Ok(bytes) => {
Ok(([(header::CONTENT_TYPE, "application/octet-stream")], bytes).into_response())
}
Err(_) => Err(StatusCode::NOT_FOUND),
}
}
/// Enqueue a `plc-provision` job for a target: extract its control-logic program,
/// stash it as a content-addressed blob (which the runner fetches via
/// [`serve_artifact`]), and queue the job. This is the control-plane "enqueue"
/// half of the loop — a runner then leases it, provisions, and posts results.
#[derive(Debug, Deserialize)]
pub struct EnqueueRequest {
/// The tenant whose queue to enqueue into.
pub tenant: String,
/// The onboarded target to test.
pub target_id: String,
}
/// The enqueued job's id.
#[derive(Debug, Serialize)]
pub struct EnqueueResponse {
/// The new job id.
pub job_id: String,
/// Whether this call inserted it (false = already queued).
pub enqueued: bool,
}
#[tracing::instrument(skip_all, fields(tenant = %req.tenant, target = %req.target_id))]
pub async fn enqueue(
Extension(agent): AgentExt,
Json(req): Json<EnqueueRequest>,
) -> Result<Json<EnqueueResponse>, StatusCode> {
let db = tenant_db(&agent, &req.tenant).await?;
let oid = ObjectId::parse_str(&req.target_id).map_err(|_| StatusCode::BAD_REQUEST)?;
let target = db
.onboarded_targets()
.find_one(doc! { "_id": oid })
.await
.map_err(internal)?
.ok_or(StatusCode::NOT_FOUND)?;
// Extract the control-logic program from the target's PLC-source artifacts
// (same selection as the in-process PLC scan).
let ctx = crate::ingest::IngestContext::from_config(&agent.config, &req.target_id);
let ingest_set = crate::ingest::ingest_all(&target, &ctx).map_err(internal)?;
let program = target
.artifacts
.iter()
.filter(|a| {
matches!(
a.kind,
ArtifactKind::PlcProject | ArtifactKind::GitRepo | ArtifactKind::SourceArchive
)
})
.find_map(|a| {
let path = ingest_set
.get(&a.id)
.and_then(|ia| ia.working_path.clone())?;
werkbank_exec::plc::extract_program(&path)
})
.ok_or(StatusCode::UNPROCESSABLE_ENTITY)?;
// Stash the program source so the runner can fetch it by hash.
let base = std::path::Path::new(&agent.config.artifact_store_base_path);
let hash =
crate::ingest::blob::store_bytes(base, program.source.as_bytes()).map_err(internal)?;
let job_id = format!("job_{}", uuid::Uuid::new_v4().simple());
let job = Job::plc_provision(
&job_id,
&req.tenant,
&req.target_id,
InputRef::blob(hash),
agent.config.plc_runtime.max_lifetime_secs,
);
let enqueued = JobQueue::new(&db)
.enqueue(job, chrono::Utc::now())
.await
.map_err(internal)?;
Ok(Json(EnqueueResponse { job_id, enqueued }))
}
/// Persist a job result's findings against its target: general findings
/// (dedup'd by fingerprint) and DAST findings. Best-effort — a persistence hiccup
/// is logged, not surfaced to the runner (its result is already recorded).
-8
View File
@@ -91,14 +91,6 @@ pub async fn start_api_server(agent: ComplianceAgent, port: u16) -> Result<(), A
"/api/v1/werkbank/jobs/complete",
post(handlers::werkbank_jobs::complete),
)
.route(
"/api/v1/werkbank/jobs/enqueue",
post(handlers::werkbank_jobs::enqueue),
)
.route(
"/api/v1/werkbank/artifacts/{hash}",
get(handlers::werkbank_jobs::serve_artifact),
)
.layer(middleware::from_fn(
handlers::werkbank_jobs::require_runner_token,
))
-10
View File
@@ -1,10 +0,0 @@
//! Controls corpus providers.
//!
//! Implementations of [`compliance_core::traits::ControlsProvider`] that supply
//! the control corpus the mapping engine assesses findings against. Currently:
//! [`OscalControlsProvider`], which pulls breakpilot-compliance's OSCAL catalog
//! and snapshots it locally.
mod oscal_provider;
pub use oscal_provider::OscalControlsProvider;
@@ -1,233 +0,0 @@
//! Pull + snapshot [`ControlsProvider`] backed by breakpilot-compliance's OSCAL
//! catalog export.
//!
//! Fetches `GET {base}/api/compliance/v1/oscal/catalog?framework=<fw>`, snapshots
//! the exact bytes to disk (so scans are deterministic and keep working offline /
//! on-prem), and maps the catalog into the corpus controls the mapping engine
//! consumes. The producer owns the catalog; we own the assessment — this is the
//! ingest half of the loop.
use std::path::PathBuf;
use secrecy::{ExposeSecret, SecretString};
use compliance_core::error::CoreError;
use compliance_core::models::onboarding::ComplianceFramework;
use compliance_core::models::oscal::OscalDocument;
use compliance_core::traits::{Control, ControlQuery, ControlsProvider};
/// A [`ControlsProvider`] that pulls the OSCAL catalog from breakpilot-compliance
/// and snapshots it locally for deterministic / offline reuse.
pub struct OscalControlsProvider {
http: reqwest::Client,
base_url: String,
token: Option<SecretString>,
snapshot_dir: PathBuf,
}
impl OscalControlsProvider {
/// Create a provider. `base_url` is the breakpilot-compliance root (e.g.
/// `http://backend-compliance:8002`); `snapshot_dir` is where catalog
/// snapshots are written so a later scan can reuse them without the network.
pub fn new(
http: reqwest::Client,
base_url: impl Into<String>,
token: Option<SecretString>,
snapshot_dir: impl Into<PathBuf>,
) -> Self {
Self {
http,
base_url: base_url.into(),
token,
snapshot_dir: snapshot_dir.into(),
}
}
fn catalog_url(&self, framework: ComplianceFramework) -> String {
format!(
"{}/api/compliance/v1/oscal/catalog?framework={framework}",
self.base_url.trim_end_matches('/')
)
}
fn snapshot_path(&self, framework: ComplianceFramework) -> PathBuf {
self.snapshot_dir
.join(format!("oscal-catalog-{framework}.json"))
}
/// Fetch the raw catalog bytes for a framework over HTTP.
async fn fetch_raw(&self, framework: ComplianceFramework) -> Result<Vec<u8>, CoreError> {
let mut req = self.http.get(self.catalog_url(framework));
if let Some(token) = &self.token {
req = req.bearer_auth(token.expose_secret());
}
let resp = req
.send()
.await
.map_err(|e| CoreError::Http(e.to_string()))?;
if !resp.status().is_success() {
return Err(CoreError::Http(format!(
"catalog fetch for {framework} returned HTTP {}",
resp.status()
)));
}
resp.bytes()
.await
.map(|b| b.to_vec())
.map_err(|e| CoreError::Http(e.to_string()))
}
/// Write a catalog snapshot atomically (temp file + rename).
async fn write_snapshot(
&self,
framework: ComplianceFramework,
raw: &[u8],
) -> Result<(), CoreError> {
tokio::fs::create_dir_all(&self.snapshot_dir).await?;
let path = self.snapshot_path(framework);
let tmp = path.with_extension("json.tmp");
tokio::fs::write(&tmp, raw).await?;
tokio::fs::rename(&tmp, &path).await?;
Ok(())
}
/// Read a previously written snapshot, if one exists.
async fn read_snapshot(
&self,
framework: ComplianceFramework,
) -> Result<Option<OscalDocument>, CoreError> {
match tokio::fs::read(self.snapshot_path(framework)).await {
Ok(raw) => Ok(Some(serde_json::from_slice(&raw)?)),
Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(None),
Err(e) => Err(e.into()),
}
}
/// Load the catalog for a framework: fetch fresh + snapshot the exact bytes;
/// on network failure, fall back to the last snapshot so scans still run.
pub async fn load(&self, framework: ComplianceFramework) -> Result<OscalDocument, CoreError> {
match self.fetch_raw(framework).await {
Ok(raw) => {
let doc: OscalDocument = serde_json::from_slice(&raw)?;
if let Err(e) = self.write_snapshot(framework, &raw).await {
tracing::warn!(%framework, error = %e, "failed to write OSCAL snapshot");
}
Ok(doc)
}
Err(fetch_err) => match self.read_snapshot(framework).await? {
Some(doc) => {
tracing::warn!(
%framework, error = %fetch_err,
"OSCAL catalog fetch failed; falling back to snapshot"
);
Ok(doc)
}
None => Err(fetch_err),
},
}
}
}
/// Order controls whose title/text mention the query context first (stable), then
/// truncate to the requested limit. Naive relevance — refined when the assessment
/// layer lands.
fn rank_and_truncate(mut controls: Vec<Control>, context: &str, limit: usize) -> Vec<Control> {
if !context.is_empty() {
let needle = context.to_lowercase();
controls.sort_by_key(|c| {
let hit =
c.title.to_lowercase().contains(&needle) || c.text.to_lowercase().contains(&needle);
u8::from(!hit)
});
}
controls.truncate(limit);
controls
}
impl ControlsProvider for OscalControlsProvider {
fn name(&self) -> &str {
"breakpilot-oscal"
}
async fn controls(&self, query: &ControlQuery<'_>) -> Result<Vec<Control>, CoreError> {
let mut out: Vec<Control> = Vec::new();
for &framework in query.frameworks {
match self.load(framework).await {
Ok(doc) => out.extend(doc.to_controls()),
Err(e) => {
tracing::warn!(%framework, error = %e, "skipping framework: catalog unavailable")
}
}
}
Ok(rank_and_truncate(out, query.context, query.limit))
}
}
#[cfg(test)]
#[allow(clippy::unwrap_used)]
mod tests {
use super::*;
const MINI_CATALOG: &str = r#"{"catalog":{"uuid":"u","metadata":{"title":"T",
"version":"1.0.0","oscal-version":"1.1.2","props":[{"name":"framework","value":"cra"}]},
"groups":[{"id":"g","title":"G","controls":[{"id":"cra-ai-1","title":"MFA",
"props":[],"parts":[{"name":"statement","prose":"require mfa"}]}]}]}}"#;
fn provider(dir: &std::path::Path) -> OscalControlsProvider {
OscalControlsProvider::new(reqwest::Client::new(), "http://unused/", None, dir)
}
#[test]
fn builds_catalog_url_and_snapshot_path() {
let p = provider(std::path::Path::new("/snap"));
assert_eq!(
p.catalog_url(ComplianceFramework::Cra),
"http://unused/api/compliance/v1/oscal/catalog?framework=cra"
);
assert_eq!(
p.snapshot_path(ComplianceFramework::Cra),
std::path::Path::new("/snap/oscal-catalog-cra.json")
);
}
#[test]
fn ranks_context_hits_first_then_truncates() {
let mk = |id: &str, title: &str| Control {
id: id.into(),
framework: ComplianceFramework::Cra,
title: title.into(),
text: String::new(),
source: None,
};
let controls = vec![
mk("a", "logging policy"),
mk("b", "multi-factor auth"),
mk("c", "backup"),
];
let ranked = rank_and_truncate(controls, "auth", 2);
assert_eq!(ranked.len(), 2);
assert_eq!(ranked[0].id, "b"); // the "auth" hit floats to the top
}
#[tokio::test]
async fn snapshot_round_trip_and_offline_fallback() {
let dir = std::env::temp_dir().join(format!("oscal-test-{}", uuid::Uuid::new_v4()));
let p = provider(&dir);
assert!(p
.read_snapshot(ComplianceFramework::Cra)
.await
.unwrap()
.is_none());
p.write_snapshot(ComplianceFramework::Cra, MINI_CATALOG.as_bytes())
.await
.unwrap();
let doc = p
.read_snapshot(ComplianceFramework::Cra)
.await
.unwrap()
.unwrap();
assert_eq!(doc.to_controls().len(), 1);
assert_eq!(doc.framework(), Some(ComplianceFramework::Cra));
let _ = std::fs::remove_dir_all(&dir);
}
}
-3
View File
@@ -27,9 +27,6 @@ pub enum AgentError {
#[error("Configuration error: {0}")]
Config(String),
#[error("Dynamic-execution error: {0}")]
Exec(#[from] werkbank_exec::ExecError),
#[error("{0}")]
Other(String),
}
-24
View File
@@ -32,30 +32,6 @@ pub fn hash_file(path: &Path) -> Result<(String, u64), AgentError> {
Ok((hex::encode(hasher.finalize()), total))
}
/// Store raw bytes in the content-addressed blob store under `base`, returning
/// the SHA-256 digest. Used to stash a small derived artifact (e.g. the extracted
/// PLC program source) so a Werkbank runner can fetch it by hash. Idempotent.
pub fn store_bytes(base: &Path, bytes: &[u8]) -> Result<String, AgentError> {
let sha = hex::encode(Sha256::digest(bytes));
let dir = base.join("blobs").join(&sha[0..2]);
fs::create_dir_all(&dir)?;
let dest = dir.join(&sha);
if !dest.exists() {
fs::write(&dest, bytes)?;
}
Ok(sha)
}
/// Read a blob's bytes by its SHA-256 digest. Rejects a non-hex/wrong-length hash
/// so a request can't traverse outside the blob store.
pub fn read_blob(base: &Path, sha: &str) -> Result<Vec<u8>, AgentError> {
if sha.len() != 64 || !sha.bytes().all(|b| b.is_ascii_hexdigit()) {
return Err(AgentError::Other(format!("invalid content hash '{sha}'")));
}
let path = base.join("blobs").join(&sha[0..2]).join(sha);
Ok(fs::read(path)?)
}
/// Copy `src` into the content-addressed blob store under `base`, returning the
/// stored path. Idempotent: an already-present blob is not rewritten.
pub fn store_file(base: &Path, src: &Path, sha: &str) -> Result<PathBuf, AgentError> {
+1 -1
View File
@@ -6,7 +6,7 @@
//! is also the reconciliation key against sibling products (a firmware sha256
//! matches tramiton's `Artifact.sha256`).
pub(crate) mod blob;
mod blob;
use std::collections::HashMap;
use std::path::{Path, PathBuf};
-1
View File
@@ -4,7 +4,6 @@ pub mod agent;
pub mod api;
pub mod classify;
pub mod config;
pub mod controls;
pub mod database;
pub mod error;
pub mod ingest;
@@ -14,7 +14,7 @@ use std::time::Duration;
use compliance_core::models::{Finding, ScanType, Severity};
use crate::fingerprint as dedup;
use crate::pipeline::dedup;
/// Well-known deep-probe ports (each independent of any WebVisu HTTP port).
const MODBUS_PORT: u16 = 502;
+1
View File
@@ -5,6 +5,7 @@ pub mod firmware_sbom;
pub mod git;
pub mod gitleaks;
mod graph_build;
pub mod ics;
mod issue_creation;
pub mod lint;
pub mod orchestrator;
@@ -587,7 +587,7 @@ impl PipelineOrchestrator {
let path = ingest_set
.get(&a.id)
.and_then(|ia| ia.working_path.clone())?;
werkbank_exec::plc::extract_program(&path)
crate::pipeline::plc::runtime::extract_program(&path)
});
let Some(program) = program else {
tracing::info!(
@@ -597,9 +597,10 @@ impl PipelineOrchestrator {
return Ok(0);
};
let http = werkbank_exec::plc::http_client()?;
let provisioner = werkbank_exec::plc::DockerSoftPlc::new(self.config.plc_runtime.clone());
let outcome = werkbank_exec::plc::provision_and_test(
let http = crate::pipeline::plc::runtime::http_client()?;
let provisioner =
crate::pipeline::plc::runtime::DockerSoftPlc::new(self.config.plc_runtime.clone());
let outcome = crate::pipeline::plc::runtime::provision_and_test(
&provisioner,
&http,
&self.config.plc_runtime,
@@ -662,7 +663,7 @@ impl PipelineOrchestrator {
};
// Short per-request budget so an unreachable device doesn't stall the scan.
let budget = std::time::Duration::from_secs(5);
let findings = werkbank_exec::ics::probe_target(&endpoint, target_id, budget).await;
let findings = crate::pipeline::ics::probe_target(&endpoint, target_id, budget).await;
tracing::info!(
target_id,
endpoint = %endpoint,
+1
View File
@@ -9,6 +9,7 @@ pub mod lexer;
pub mod parser;
pub mod plcopen;
pub mod rules;
pub mod runtime;
pub mod sbom;
use std::path::Path;
@@ -24,7 +24,7 @@ use compliance_core::models::dast::{DastFinding, DastScanRun, DastTarget, DastTa
use compliance_core::models::Finding;
use compliance_core::PlcRuntimeConfig;
use crate::error::ExecError;
use crate::error::AgentError;
pub use provision::{DockerSoftPlc, ProvisionedRuntime, SoftPlc};
@@ -61,12 +61,12 @@ pub struct PlcProgram {
/// A cookie-aware HTTP client for the OpenPLC web UI. A fresh client per scan
/// isolates the OpenPLC session (its Flask login cookie) from every other scan.
pub fn http_client() -> Result<reqwest::Client, ExecError> {
pub fn http_client() -> Result<reqwest::Client, AgentError> {
reqwest::Client::builder()
.cookie_store(true)
.timeout(Duration::from_secs(30))
.build()
.map_err(ExecError::Http)
.map_err(AgentError::Http)
}
/// Pick the control-logic program to run from an ingested PLC source tree.
@@ -151,7 +151,7 @@ pub async fn provision_and_test<P: SoftPlc>(
cfg: &PlcRuntimeConfig,
program: &PlcProgram,
target_id: &str,
) -> Result<ProvisionOutcome, ExecError> {
) -> Result<ProvisionOutcome, AgentError> {
let handle = provisioner.provision(target_id).await?;
tracing::info!(
target_id,
@@ -193,7 +193,7 @@ async fn run_dynamic_test(
program: &PlcProgram,
target_id: &str,
handle: &ProvisionedRuntime,
) -> Result<ProvisionOutcome, ExecError> {
) -> Result<ProvisionOutcome, AgentError> {
let ready_budget = Duration::from_secs((cfg.max_lifetime_secs / 3).clamp(10, 60));
openplc::wait_ready(http, &handle.webvisu_url, ready_budget).await?;
@@ -212,7 +212,8 @@ async fn run_dynamic_test(
tokio::time::sleep(Duration::from_secs(3)).await;
let probe_budget = Duration::from_secs(5);
let findings = crate::ics::probe_target(&handle.modbus_endpoint, target_id, probe_budget).await;
let findings =
crate::pipeline::ics::probe_target(&handle.modbus_endpoint, target_id, probe_budget).await;
tracing::info!(
target_id,
instance = %handle.name,
@@ -347,10 +348,10 @@ mod tests {
}
impl SoftPlc for FakeSoftPlc {
async fn provision(&self, _target_id: &str) -> Result<ProvisionedRuntime, ExecError> {
async fn provision(&self, _target_id: &str) -> Result<ProvisionedRuntime, AgentError> {
self.provisions.fetch_add(1, Ordering::SeqCst);
if self.fail_provision {
return Err(ExecError::Other("provision failed".into()));
return Err(AgentError::Other("provision failed".into()));
}
// Unreachable address so run_dynamic_test blocks on readiness until the
// deadline fires — exercising the teardown-on-deadline path.
@@ -10,7 +10,7 @@
use std::time::Duration;
use crate::error::ExecError;
use crate::error::AgentError;
use super::PlcProgram;
@@ -27,7 +27,7 @@ pub async fn wait_ready(
http: &reqwest::Client,
base_url: &str,
budget: Duration,
) -> Result<(), ExecError> {
) -> Result<(), AgentError> {
let login = format!("{base_url}/login");
let outcome = tokio::time::timeout(budget, async {
loop {
@@ -40,7 +40,7 @@ pub async fn wait_ready(
}
})
.await;
outcome.map_err(|_| ExecError::Other(format!("OpenPLC at {base_url} did not become ready")))
outcome.map_err(|_| AgentError::Other(format!("OpenPLC at {base_url} did not become ready")))
}
/// Log in, upload the program, compile it, and start the runtime. On success the
@@ -52,7 +52,7 @@ pub async fn load_and_start(
password: &str,
program: &PlcProgram,
compile_budget: Duration,
) -> Result<(), ExecError> {
) -> Result<(), AgentError> {
login(http, base_url, user, password).await?;
let prog_file = upload_program(http, base_url, program).await?;
save_program(http, base_url, &prog_file).await?;
@@ -68,14 +68,14 @@ async fn login(
base_url: &str,
user: &str,
password: &str,
) -> Result<(), ExecError> {
) -> Result<(), AgentError> {
let resp = http
.post(format!("{base_url}/login"))
.form(&[("username", user), ("password", password)])
.send()
.await?;
if resp.status().is_server_error() {
return Err(ExecError::Other(format!(
return Err(AgentError::Other(format!(
"OpenPLC login failed: HTTP {}",
resp.status()
)));
@@ -90,7 +90,7 @@ async fn upload_program(
http: &reqwest::Client,
base_url: &str,
program: &PlcProgram,
) -> Result<String, ExecError> {
) -> Result<String, AgentError> {
let part = reqwest::multipart::Part::text(program.source.clone())
.file_name(program.file_name.clone())
.mime_str("application/octet-stream")?;
@@ -102,7 +102,7 @@ async fn upload_program(
.await?;
let html = resp.text().await?;
parse_prog_file(&html).ok_or_else(|| {
ExecError::Other("OpenPLC upload did not return a prog_file handle".to_string())
AgentError::Other("OpenPLC upload did not return a prog_file handle".to_string())
})
}
@@ -113,7 +113,7 @@ async fn save_program(
http: &reqwest::Client,
base_url: &str,
prog_file: &str,
) -> Result<(), ExecError> {
) -> Result<(), AgentError> {
let epoch = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_secs())
@@ -130,7 +130,7 @@ async fn save_program(
.send()
.await?;
if resp.status().is_server_error() {
return Err(ExecError::Other(format!(
return Err(AgentError::Other(format!(
"OpenPLC save-program failed: HTTP {}",
resp.status()
)));
@@ -146,7 +146,7 @@ async fn compile(
base_url: &str,
prog_file: &str,
budget: Duration,
) -> Result<(), ExecError> {
) -> Result<(), AgentError> {
http.get(format!("{base_url}/compile-program"))
.query(&[("file", prog_file)])
.send()
@@ -168,20 +168,20 @@ async fn compile(
.await;
match outcome {
Ok(true) => Ok(()),
Ok(false) => Err(ExecError::Other(
Ok(false) => Err(AgentError::Other(
"OpenPLC compilation finished with errors".to_string(),
)),
Err(_) => Err(ExecError::Other(
Err(_) => Err(AgentError::Other(
"OpenPLC compilation did not finish in time".to_string(),
)),
}
}
/// `GET /start_plc` — starts the runtime, opening Modbus/TCP on 502.
async fn start(http: &reqwest::Client, base_url: &str) -> Result<(), ExecError> {
async fn start(http: &reqwest::Client, base_url: &str) -> Result<(), AgentError> {
let resp = http.get(format!("{base_url}/start_plc")).send().await?;
if resp.status().is_server_error() {
return Err(ExecError::Other(format!(
return Err(AgentError::Other(format!(
"OpenPLC start_plc failed: HTTP {}",
resp.status()
)));
@@ -15,7 +15,7 @@ use std::time::{SystemTime, UNIX_EPOCH};
use compliance_core::PlcRuntimeConfig;
use crate::error::ExecError;
use crate::error::AgentError;
/// The Modbus/TCP port an OpenPLC instance opens once a program is running.
const MODBUS_PORT: u16 = 502;
@@ -45,7 +45,7 @@ pub trait SoftPlc {
fn provision(
&self,
target_id: &str,
) -> impl std::future::Future<Output = Result<ProvisionedRuntime, ExecError>> + Send;
) -> impl std::future::Future<Output = Result<ProvisionedRuntime, AgentError>> + Send;
/// Tear an instance down. Best-effort and idempotent — never fails the scan.
fn teardown(&self, handle: &ProvisionedRuntime)
@@ -65,7 +65,7 @@ impl DockerSoftPlc {
}
impl SoftPlc for DockerSoftPlc {
async fn provision(&self, target_id: &str) -> Result<ProvisionedRuntime, ExecError> {
async fn provision(&self, target_id: &str) -> Result<ProvisionedRuntime, AgentError> {
// Best-effort sweep of any container leaked by a crashed earlier run
// before we add another. Only removes instances past their max lifetime,
// so it can never disturb a concurrent run.
@@ -75,7 +75,7 @@ impl SoftPlc for DockerSoftPlc {
let args = run_args(&self.cfg, &name, target_id);
let out = run_docker(&args).await?;
if !out.status.success() {
return Err(ExecError::Other(format!(
return Err(AgentError::Other(format!(
"docker run for soft-PLC {name} failed: {}",
String::from_utf8_lossy(&out.stderr).trim()
)));
@@ -215,12 +215,12 @@ async fn reap_stale(cfg: &PlcRuntimeConfig, now: u64) {
}
/// Run a `docker` subcommand, capturing its output.
async fn run_docker(args: &[String]) -> Result<std::process::Output, ExecError> {
async fn run_docker(args: &[String]) -> Result<std::process::Output, AgentError> {
tokio::process::Command::new("docker")
.args(args)
.output()
.await
.map_err(ExecError::Io)
.map_err(AgentError::Io)
}
#[cfg(test)]
+2 -71
View File
@@ -11,7 +11,7 @@ mod common;
use std::sync::Arc;
use axum::routing::{get, post};
use axum::routing::post;
use axum::{middleware, Extension, Router};
use compliance_agent::agent::ComplianceAgent;
@@ -19,9 +19,7 @@ use compliance_agent::api::handlers::werkbank_jobs;
use compliance_agent::database::DatabasePool;
use compliance_agent::werkbank::JobQueue;
use compliance_core::models::werkbank::{InputRef, Job, JobResult, JobStatus, LeasedJob};
use compliance_core::models::{
Artifact, Finding, OnboardedTarget, PlcFormat, ScanType, Severity, TargetType,
};
use compliance_core::models::{Finding, ScanType, Severity};
use common::{dev_config, TEST_RUNNER_TOKEN};
@@ -60,14 +58,6 @@ async fn start() -> Option<Harness> {
"/api/v1/werkbank/jobs/complete",
post(werkbank_jobs::complete),
)
.route(
"/api/v1/werkbank/jobs/enqueue",
post(werkbank_jobs::enqueue),
)
.route(
"/api/v1/werkbank/artifacts/{hash}",
get(werkbank_jobs::serve_artifact),
)
.layer(middleware::from_fn(werkbank_jobs::require_runner_token))
.layer(Extension(Arc::new(agent)));
@@ -187,65 +177,6 @@ async fn lease_complete_persists_findings_against_the_target() {
h.cleanup().await;
}
#[tokio::test]
async fn enqueue_extracts_program_stores_a_blob_and_serves_it() {
let Some(h) = start().await else { return };
let db = h.pool.for_tenant_id(TENANT).await.unwrap();
// A PlcSps target with a single complete ST program uploaded.
let dir = std::env::temp_dir().join(format!("wbq-prog-{}", uuid::Uuid::new_v4()));
std::fs::create_dir_all(&dir).unwrap();
let st = dir.join("main.st");
std::fs::write(
&st,
"PROGRAM Main\nEND_PROGRAM\nCONFIGURATION C\n RESOURCE R\nEND_CONFIGURATION\n",
)
.unwrap();
let mut target = OnboardedTarget::new("plc".into(), TargetType::PlcSps);
let mut art = Artifact::plc_project("main.st", PlcFormat::StructuredText);
art.stored_path = Some(st.to_string_lossy().to_string());
target.artifacts.push(art);
let ins = db.onboarded_targets().insert_one(&target).await.unwrap();
let target_id = ins.inserted_id.as_object_id().unwrap().to_hex();
// Enqueue → a plc-provision job whose program is a content-addressed blob.
let resp = h
.post(
"/api/v1/werkbank/jobs/enqueue",
Some(TEST_RUNNER_TOKEN),
serde_json::json!({ "tenant": TENANT, "target_id": target_id }),
)
.send()
.await
.unwrap();
assert_eq!(resp.status(), 200, "enqueue should succeed");
let body: serde_json::Value = resp.json().await.unwrap();
let job_id = body["job_id"].as_str().unwrap().to_string();
let rec = JobQueue::new(&db).get(&job_id).await.unwrap().unwrap();
let hash = rec
.job
.inputs
.get("program")
.and_then(|i| i.blob.clone())
.expect("program blob");
// Serve the blob back and confirm it's the program source (what the runner
// would fetch).
let served = h
.client
.get(format!("{}/api/v1/werkbank/artifacts/{hash}", h.base_url))
.bearer_auth(TEST_RUNNER_TOKEN)
.send()
.await
.unwrap();
assert_eq!(served.status(), 200);
assert!(served.text().await.unwrap().contains("CONFIGURATION"));
h.cleanup().await;
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn empty_queue_leases_nothing() {
let Some(h) = start().await else { return };
-2
View File
@@ -10,7 +10,6 @@ pub mod mcp;
pub mod mcp_token;
pub mod notification;
pub mod onboarding;
pub mod oscal;
pub mod pentest;
pub mod repository;
pub mod sbom;
@@ -40,7 +39,6 @@ pub use onboarding::{
GitArtifactConfig, IssueTrackerConfig, OnboardedTarget, PlcArtifactConfig, PlcFormat,
TargetScanConfig, TargetType, TargetTypeCandidate, WebArtifactConfig,
};
pub use oscal::OscalDocument;
pub use pentest::{
AttackChainNode, AttackNodeStatus, AuthMode, CodeContextHint, Environment, IdentityProvider,
PentestAuthConfig, PentestConfig, PentestEvent, PentestMessage, PentestSession, PentestStats,
-250
View File
@@ -1,250 +0,0 @@
//! OSCAL 1.1 catalog types + mapping into the controls corpus.
//!
//! Deserialises the OSCAL catalog served by breakpilot-compliance
//! (`GET /api/compliance/v1/oscal/catalog`) and maps its controls into the
//! framework-agnostic [`crate::traits::Control`] that the mapping engine consumes.
//! Only the fields we use are modelled; unknown OSCAL fields are ignored so the
//! producer can add detail without breaking us.
//!
//! Scope boundary: this is the *catalog* (domain content). Assessment objectives
//! and scanner routing live in our assessment layer, not here — see
//! [`crate::traits::ControlsProvider`].
use serde::Deserialize;
use crate::models::onboarding::ComplianceFramework;
use crate::traits::Control as CorpusControl;
/// A parsed OSCAL catalog document (`{"catalog": {...}}`).
#[derive(Debug, Clone, Deserialize)]
pub struct OscalDocument {
pub catalog: Catalog,
}
/// An OSCAL catalog: metadata + a tree of control groups.
#[derive(Debug, Clone, Deserialize)]
pub struct Catalog {
pub uuid: String,
pub metadata: Metadata,
#[serde(default)]
pub groups: Vec<Group>,
#[serde(rename = "back-matter", default)]
pub back_matter: Option<BackMatter>,
}
/// Catalog metadata (title/version + provenance props).
#[derive(Debug, Clone, Deserialize)]
pub struct Metadata {
pub title: String,
pub version: String,
#[serde(rename = "oscal-version")]
pub oscal_version: String,
#[serde(default)]
pub props: Vec<Prop>,
}
/// A name/value property, optionally namespaced.
#[derive(Debug, Clone, Deserialize)]
pub struct Prop {
pub name: String,
pub value: String,
#[serde(default)]
pub ns: Option<String>,
}
/// A control group (may nest sub-groups and controls).
#[derive(Debug, Clone, Deserialize)]
pub struct Group {
#[serde(default)]
pub id: String,
#[serde(default)]
pub title: String,
#[serde(default)]
pub controls: Vec<Control>,
#[serde(default)]
pub groups: Vec<Group>,
}
/// An OSCAL control (may nest enhancement controls).
#[derive(Debug, Clone, Deserialize)]
pub struct Control {
pub id: String,
#[serde(default)]
pub title: String,
#[serde(default)]
pub props: Vec<Prop>,
#[serde(default)]
pub parts: Vec<Part>,
#[serde(default)]
pub links: Vec<Link>,
#[serde(default)]
pub controls: Vec<Control>,
}
/// A control part (e.g. the `statement`), may nest sub-parts.
#[derive(Debug, Clone, Deserialize)]
pub struct Part {
#[serde(default)]
pub name: String,
#[serde(default)]
pub prose: Option<String>,
#[serde(default)]
pub parts: Vec<Part>,
}
/// A link, e.g. a `reference` to a back-matter resource.
#[derive(Debug, Clone, Deserialize)]
pub struct Link {
pub href: String,
#[serde(default)]
pub rel: Option<String>,
}
/// Back-matter holding referenced resources (e.g. the CRA measures).
#[derive(Debug, Clone, Deserialize)]
pub struct BackMatter {
#[serde(default)]
pub resources: Vec<Resource>,
}
/// A back-matter resource referenced by control links.
#[derive(Debug, Clone, Deserialize)]
pub struct Resource {
pub uuid: String,
#[serde(default)]
pub title: Option<String>,
#[serde(default)]
pub description: Option<String>,
}
impl Metadata {
/// First prop value with the given name.
pub fn prop(&self, name: &str) -> Option<&str> {
self.props
.iter()
.find(|p| p.name == name)
.map(|p| p.value.as_str())
}
}
impl Control {
/// First prop value with the given name.
pub fn prop(&self, name: &str) -> Option<&str> {
self.props
.iter()
.find(|p| p.name == name)
.map(|p| p.value.as_str())
}
/// The control's `statement` prose, if present.
pub fn statement(&self) -> Option<&str> {
self.parts
.iter()
.find(|p| p.name == "statement")
.and_then(|p| p.prose.as_deref())
}
}
impl OscalDocument {
/// The framework this catalog declares (`metadata.props[name="framework"]`).
pub fn framework(&self) -> Option<ComplianceFramework> {
framework_from_str(self.catalog.metadata.prop("framework")?)
}
/// The catalog `content-hash` prop — consumers pin this to snapshot/detect drift.
pub fn content_hash(&self) -> Option<&str> {
self.catalog.metadata.prop("content-hash")
}
/// Flatten the catalog into the corpus controls the mapping engine consumes.
pub fn to_controls(&self) -> Vec<CorpusControl> {
let framework = self.framework().unwrap_or(ComplianceFramework::Cra);
let source_label = self.catalog.metadata.title.as_str();
let mut out = Vec::new();
for group in &self.catalog.groups {
collect_group(group, framework, source_label, &mut out);
}
out
}
}
/// Map an OSCAL framework token (e.g. `"cra"`) to [`ComplianceFramework`] via its
/// serde snake_case representation.
fn framework_from_str(raw: &str) -> Option<ComplianceFramework> {
serde_json::from_value(serde_json::Value::String(raw.to_string())).ok()
}
fn collect_group(
group: &Group,
framework: ComplianceFramework,
source_label: &str,
out: &mut Vec<CorpusControl>,
) {
for control in &group.controls {
collect_control(control, framework, source_label, out);
}
for sub in &group.groups {
collect_group(sub, framework, source_label, out);
}
}
fn collect_control(
control: &Control,
framework: ComplianceFramework,
source_label: &str,
out: &mut Vec<CorpusControl>,
) {
let source = match control.prop("annex-anchor") {
Some(anchor) => Some(format!("{source_label} · {anchor}")),
None => Some(source_label.to_string()),
};
out.push(CorpusControl {
id: control.id.clone(),
framework,
title: control.title.clone(),
text: control.statement().unwrap_or_default().to_string(),
source,
});
for enhancement in &control.controls {
collect_control(enhancement, framework, source_label, out);
}
}
#[cfg(test)]
#[allow(clippy::unwrap_used)]
mod tests {
use super::*;
const CATALOG: &str = include_str!("../../tests/data/cra_catalog.json");
fn parse() -> OscalDocument {
serde_json::from_str(CATALOG).unwrap()
}
#[test]
fn parses_full_catalog() {
let doc = parse();
assert_eq!(doc.catalog.metadata.oscal_version, "1.1.2");
assert!(!doc.catalog.groups.is_empty());
assert!(doc.catalog.back_matter.is_some());
}
#[test]
fn maps_all_controls_to_corpus() {
let doc = parse();
let controls = doc.to_controls();
assert_eq!(controls.len(), 40);
assert_eq!(doc.framework(), Some(ComplianceFramework::Cra));
let c8 = controls.iter().find(|c| c.id == "cra-ai-8").unwrap();
assert_eq!(c8.framework, ComplianceFramework::Cra);
assert!(!c8.title.is_empty());
assert!(!c8.text.is_empty(), "statement prose should map into text");
assert!(c8.source.as_deref().unwrap_or_default().contains("Annex I"));
}
#[test]
fn exposes_content_hash_for_snapshotting() {
assert_eq!(parse().content_hash().map(str::len), Some(64));
}
}
File diff suppressed because it is too large Load Diff
-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;