From 32b043573803453ccf6d32031422f6bf40664ce5 Mon Sep 17 00:00:00 2001 From: Sharang Parnerkar <30073382+mighty840@users.noreply.github.com> Date: Fri, 17 Jul 2026 16:26:26 +0200 Subject: [PATCH] =?UTF-8?q?feat(werkbank):=20make=20the=20loop=20runnable?= =?UTF-8?q?=20=E2=80=94=20enqueue=20+=20artifact=20serve/fetch=20(WB-05b)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Closes the gap between "endpoints exist" and "a runner can actually run a job": - POST /api/v1/werkbank/jobs/enqueue {tenant, target_id} — the control-plane "enqueue" half: extract the target's control-logic program, stash the source as a content-addressed blob, and queue a plc-provision job referencing it by hash. - GET /api/v1/werkbank/artifacts/{hash} — serve a blob so the runner can fetch the program (traversal-safe: the hash is validated). Both runner-token gated. - blob.rs: store_bytes / read_blob helpers; ingest::blob is now pub(crate). - .env.example documents WERKBANK_RUNNER_TOKEN. With the runner-side auth + blob fetch (werkbank repo), the loop runs end to end: enqueue -> lease -> Docker executor fetches the program by hash, provisions, probes, DASTs -> completes -> findings persisted against the target. Tests: a new integration test enqueues from a PlcSps target, confirms the job carries a program blob, and serves it back. All 4 werkbank_api tests pass; clippy + fmt clean. Co-Authored-By: Claude Fable 5 --- .env.example | 5 + .../src/api/handlers/werkbank_jobs.rs | 102 +++++++++++++++++- compliance-agent/src/api/server.rs | 8 ++ compliance-agent/src/ingest/blob.rs | 24 +++++ compliance-agent/src/ingest/mod.rs | 2 +- compliance-agent/tests/werkbank_api.rs | 73 ++++++++++++- 6 files changed, 208 insertions(+), 6 deletions(-) diff --git a/.env.example b/.env.example index 88b8b83..974bfbb 100644 --- a/.env.example +++ b/.env.example @@ -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 diff --git a/compliance-agent/src/api/handlers/werkbank_jobs.rs b/compliance-agent/src/api/handlers/werkbank_jobs.rs index 41692c7..212d668 100644 --- a/compliance-agent/src/api/handlers/werkbank_jobs.rs +++ b/compliance-agent/src/api/handlers/werkbank_jobs.rs @@ -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, +) -> Result { + 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, +) -> Result, 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). diff --git a/compliance-agent/src/api/server.rs b/compliance-agent/src/api/server.rs index 0f97477..50d8c3d 100644 --- a/compliance-agent/src/api/server.rs +++ b/compliance-agent/src/api/server.rs @@ -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, )) diff --git a/compliance-agent/src/ingest/blob.rs b/compliance-agent/src/ingest/blob.rs index f279434..fb90972 100644 --- a/compliance-agent/src/ingest/blob.rs +++ b/compliance-agent/src/ingest/blob.rs @@ -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 { + 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, 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 { diff --git a/compliance-agent/src/ingest/mod.rs b/compliance-agent/src/ingest/mod.rs index c1b2009..338f544 100644 --- a/compliance-agent/src/ingest/mod.rs +++ b/compliance-agent/src/ingest/mod.rs @@ -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}; diff --git a/compliance-agent/tests/werkbank_api.rs b/compliance-agent/tests/werkbank_api.rs index 487ba8c..b382530 100644 --- a/compliance-agent/tests/werkbank_api.rs +++ b/compliance-agent/tests/werkbank_api.rs @@ -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 { "/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 }; -- 2.54.0