feat(werkbank): make the loop runnable — enqueue + artifact serve/fetch (WB-05b) (#209)
CI / Check (push) Has been skipped
CI / Detect Changes (push) Successful in 3s
CI / Deploy Agent (push) Successful in 3m43s
CI / Deploy Dashboard (push) Has been skipped
CI / Deploy Docs (push) Has been skipped
CI / Deploy MCP (push) Has been skipped

This commit was merged in pull request #209.
This commit is contained in:
2026-07-17 14:34:00 +00:00
parent 70a4ee55ab
commit 94c0d51a11
6 changed files with 208 additions and 6 deletions
+5
View File
@@ -47,6 +47,11 @@ 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
@@ -10,18 +10,20 @@
//! 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::extract::{Extension, Path, Request};
use axum::http::{header, StatusCode};
use axum::middleware::Next;
use axum::response::{IntoResponse, Response};
use axum::Json;
use mongodb::bson::doc;
use mongodb::bson::{doc, oid::ObjectId};
use secrecy::ExposeSecret;
use serde::{Deserialize, Serialize};
use std::time::Duration;
use compliance_core::models::werkbank::{
CompleteRequest, CompleteResponse, HeartbeatRequest, JobResult, LeaseRequest,
CompleteRequest, CompleteResponse, HeartbeatRequest, InputRef, Job, JobResult, LeaseRequest,
};
use compliance_core::models::ArtifactKind;
use super::dto::AgentExt;
use crate::database::Database;
@@ -122,6 +124,100 @@ 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,6 +91,14 @@ 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,
))
+24
View File
@@ -32,6 +32,30 @@ 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`).
mod blob;
pub(crate) mod blob;
use std::collections::HashMap;
use std::path::{Path, PathBuf};
+71 -2
View File
@@ -11,7 +11,7 @@ mod common;
use std::sync::Arc;
use axum::routing::post;
use axum::routing::{get, post};
use axum::{middleware, Extension, Router};
use compliance_agent::agent::ComplianceAgent;
@@ -19,7 +19,9 @@ 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 compliance_core::models::{
Artifact, Finding, OnboardedTarget, PlcFormat, ScanType, Severity, TargetType,
};
use common::{dev_config, TEST_RUNNER_TOKEN};
@@ -58,6 +60,14 @@ 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)));
@@ -177,6 +187,65 @@ 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 };