Compare commits
1
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
e5f4b562c3 |
@@ -47,11 +47,6 @@ PLC_RUNTIME_MAX_LIFETIME_SECS=180
|
|||||||
PLC_RUNTIME_OPENPLC_USER=openplc
|
PLC_RUNTIME_OPENPLC_USER=openplc
|
||||||
PLC_RUNTIME_OPENPLC_PASSWORD=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
|
||||||
DASHBOARD_PORT=8080
|
DASHBOARD_PORT=8080
|
||||||
AGENT_API_URL=http://localhost:3001
|
AGENT_API_URL=http://localhost:3001
|
||||||
|
|||||||
@@ -14,7 +14,6 @@ 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::*;
|
||||||
|
|||||||
@@ -1,289 +0,0 @@
|
|||||||
//! 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, Path, 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 secrecy::ExposeSecret;
|
|
||||||
use serde::{Deserialize, Serialize};
|
|
||||||
use std::time::Duration;
|
|
||||||
|
|
||||||
use compliance_core::models::werkbank::{
|
|
||||||
CompleteRequest, CompleteResponse, HeartbeatRequest, InputRef, Job, JobResult, LeaseRequest,
|
|
||||||
};
|
|
||||||
use compliance_core::models::ArtifactKind;
|
|
||||||
|
|
||||||
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 }))
|
|
||||||
}
|
|
||||||
|
|
||||||
/// `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).
|
|
||||||
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"));
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -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, post};
|
use axum::routing::{delete, get};
|
||||||
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,43 +72,8 @@ 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),
|
|
||||||
)
|
|
||||||
.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,
|
|
||||||
))
|
|
||||||
} 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))
|
||||||
|
|||||||
@@ -65,7 +65,6 @@ 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"),
|
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -32,30 +32,6 @@ pub fn hash_file(path: &Path) -> Result<(String, u64), AgentError> {
|
|||||||
Ok((hex::encode(hasher.finalize()), total))
|
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
|
/// Copy `src` into the content-addressed blob store under `base`, returning the
|
||||||
/// stored path. Idempotent: an already-present blob is not rewritten.
|
/// stored path. Idempotent: an already-present blob is not rewritten.
|
||||||
pub fn store_file(base: &Path, src: &Path, sha: &str) -> Result<PathBuf, AgentError> {
|
pub fn store_file(base: &Path, src: &Path, sha: &str) -> Result<PathBuf, AgentError> {
|
||||||
|
|||||||
@@ -6,7 +6,7 @@
|
|||||||
//! is also the reconciliation key against sibling products (a firmware sha256
|
//! is also the reconciliation key against sibling products (a firmware sha256
|
||||||
//! matches tramiton's `Artifact.sha256`).
|
//! matches tramiton's `Artifact.sha256`).
|
||||||
|
|
||||||
pub(crate) mod blob;
|
mod blob;
|
||||||
|
|
||||||
use std::collections::HashMap;
|
use std::collections::HashMap;
|
||||||
use std::path::{Path, PathBuf};
|
use std::path::{Path, PathBuf};
|
||||||
|
|||||||
@@ -343,7 +343,6 @@ 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,
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -2,10 +2,6 @@
|
|||||||
//
|
//
|
||||||
// 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;
|
||||||
|
|
||||||
@@ -15,54 +11,6 @@ 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,
|
||||||
@@ -85,7 +33,45 @@ impl TestServer {
|
|||||||
.await
|
.await
|
||||||
.expect("Failed to build DatabasePool");
|
.expect("Failed to build DatabasePool");
|
||||||
|
|
||||||
let config = dev_config(mongodb_uri.clone(), db_name.clone());
|
let config = AgentConfig {
|
||||||
|
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);
|
||||||
|
|
||||||
|
|||||||
@@ -1,291 +0,0 @@
|
|||||||
//! 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::{get, 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::{
|
|
||||||
Artifact, Finding, OnboardedTarget, PlcFormat, ScanType, Severity, TargetType,
|
|
||||||
};
|
|
||||||
|
|
||||||
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),
|
|
||||||
)
|
|
||||||
.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)));
|
|
||||||
|
|
||||||
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 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 };
|
|
||||||
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;
|
|
||||||
}
|
|
||||||
@@ -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/*`) and the Werkbank runner API (`/api/v1/werkbank/*`)
|
/// (`/api/v1/admin/*`) has its own static-bearer middleware and must
|
||||||
/// have their own static-bearer middleware and must not be routed through the
|
/// not be routed through the customer-JWT path — a Keycloak token
|
||||||
/// customer-JWT path — a Keycloak token always carries a single tenant_id and
|
/// always carries a single tenant_id and would semantically conflict
|
||||||
/// would semantically conflict with these cross-tenant / machine operations.
|
/// with cross-tenant admin operations.
|
||||||
const PUBLIC_PREFIXES: &[&str] = &["/api/v1/admin/", "/api/v1/werkbank/"];
|
const PUBLIC_PREFIXES: &[&str] = &["/api/v1/admin/"];
|
||||||
|
|
||||||
/// 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.
|
||||||
|
|||||||
@@ -53,11 +53,6 @@ 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).
|
||||||
|
|||||||
@@ -49,7 +49,6 @@ 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::{
|
||||||
CompleteRequest, CompleteResponse, DastCollect, Executor, HeartbeatAck, HeartbeatRequest,
|
DastCollect, Executor, HeartbeatAck, InputRef, Job, JobCollect, JobRecord, JobResult,
|
||||||
InputRef, Job, JobCollect, JobRecord, JobResult, JobRuntime, JobStatus, JobType, LeaseRequest,
|
JobRuntime, JobStatus, JobType, LeasedJob,
|
||||||
LeasedJob,
|
|
||||||
};
|
};
|
||||||
|
|||||||
@@ -349,59 +349,6 @@ 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 {
|
||||||
|
|||||||
Reference in New Issue
Block a user