feat(werkbank): runner queue endpoints + result persistence (WB-05) (#207)
CI / Check (push) Has been skipped
CI / Detect Changes (push) Successful in 4s
CI / Deploy Agent (push) Successful in 3m44s
CI / Deploy Dashboard (push) Successful in 2m37s
CI / Deploy Docs (push) Has been skipped
CI / Deploy MCP (push) Successful in 1m49s
CI / Check (push) Has been skipped
CI / Detect Changes (push) Successful in 4s
CI / Deploy Agent (push) Successful in 3m44s
CI / Deploy Dashboard (push) Successful in 2m37s
CI / Deploy Docs (push) Has been skipped
CI / Deploy MCP (push) Successful in 1m49s
This commit was merged in pull request #207.
This commit is contained in:
@@ -2,6 +2,10 @@
|
||||
//
|
||||
// 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.
|
||||
//
|
||||
// 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;
|
||||
|
||||
@@ -11,6 +15,54 @@ use compliance_agent::database::DatabasePool;
|
||||
use compliance_core::AgentConfig;
|
||||
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.
|
||||
pub struct TestServer {
|
||||
pub base_url: String,
|
||||
@@ -33,45 +85,7 @@ impl TestServer {
|
||||
.await
|
||||
.expect("Failed to build DatabasePool");
|
||||
|
||||
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 config = dev_config(mongodb_uri.clone(), db_name.clone());
|
||||
|
||||
let agent = ComplianceAgent::new(config, db_pool);
|
||||
|
||||
|
||||
@@ -0,0 +1,222 @@
|
||||
//! Integration tests for the Werkbank runner endpoints (WB-05).
|
||||
//!
|
||||
//! Drives the real HTTP handlers (lease/heartbeat/complete) against a live Mongo:
|
||||
//! a runner leases a seeded job, completes it, and the result's findings are
|
||||
//! persisted against the job's target. Also checks the bearer-token gate. Skips
|
||||
//! cleanly when no Mongo is reachable.
|
||||
|
||||
#![allow(clippy::expect_used, clippy::unwrap_used)]
|
||||
|
||||
mod common;
|
||||
|
||||
use std::sync::Arc;
|
||||
|
||||
use axum::routing::post;
|
||||
use axum::{middleware, Extension, Router};
|
||||
|
||||
use compliance_agent::agent::ComplianceAgent;
|
||||
use compliance_agent::api::handlers::werkbank_jobs;
|
||||
use compliance_agent::database::DatabasePool;
|
||||
use compliance_agent::werkbank::JobQueue;
|
||||
use compliance_core::models::werkbank::{InputRef, Job, JobResult, JobStatus, LeasedJob};
|
||||
use compliance_core::models::{Finding, ScanType, Severity};
|
||||
|
||||
use common::{dev_config, TEST_RUNNER_TOKEN};
|
||||
|
||||
const TENANT: &str = "dev";
|
||||
|
||||
/// A running werkbank API on a random port, or `None` if no Mongo.
|
||||
struct Harness {
|
||||
base_url: String,
|
||||
client: reqwest::Client,
|
||||
pool: DatabasePool,
|
||||
db_name: String,
|
||||
}
|
||||
|
||||
async fn start() -> Option<Harness> {
|
||||
let uri = std::env::var("TEST_MONGODB_URI")
|
||||
.unwrap_or_else(|_| "mongodb://root:example@localhost:27017/?authSource=admin".into());
|
||||
let db_name = format!("wba_{}", &uuid::Uuid::new_v4().simple().to_string()[..12]);
|
||||
let pool = match DatabasePool::connect(&uri, &db_name).await {
|
||||
Ok(p) => p,
|
||||
Err(_) => {
|
||||
eprintln!("SKIP werkbank_api: no MongoDB reachable at {uri}");
|
||||
return None;
|
||||
}
|
||||
};
|
||||
// Touch the tenant DB so indexes are ensured before the queue is used.
|
||||
pool.for_tenant_id(TENANT).await.expect("tenant db");
|
||||
|
||||
let agent = ComplianceAgent::new(dev_config(uri, db_name.clone()), pool.clone());
|
||||
let app = Router::new()
|
||||
.route("/api/v1/werkbank/jobs/lease", post(werkbank_jobs::lease))
|
||||
.route(
|
||||
"/api/v1/werkbank/jobs/heartbeat",
|
||||
post(werkbank_jobs::heartbeat),
|
||||
)
|
||||
.route(
|
||||
"/api/v1/werkbank/jobs/complete",
|
||||
post(werkbank_jobs::complete),
|
||||
)
|
||||
.layer(middleware::from_fn(werkbank_jobs::require_runner_token))
|
||||
.layer(Extension(Arc::new(agent)));
|
||||
|
||||
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
|
||||
let port = listener.local_addr().unwrap().port();
|
||||
tokio::spawn(async move {
|
||||
axum::serve(listener, app).await.ok();
|
||||
});
|
||||
|
||||
Some(Harness {
|
||||
base_url: format!("http://127.0.0.1:{port}"),
|
||||
client: reqwest::Client::new(),
|
||||
pool,
|
||||
db_name,
|
||||
})
|
||||
}
|
||||
|
||||
impl Harness {
|
||||
fn post(
|
||||
&self,
|
||||
path: &str,
|
||||
token: Option<&str>,
|
||||
body: serde_json::Value,
|
||||
) -> reqwest::RequestBuilder {
|
||||
let mut r = self
|
||||
.client
|
||||
.post(format!("{}{path}", self.base_url))
|
||||
.json(&body);
|
||||
if let Some(t) = token {
|
||||
r = r.bearer_auth(t);
|
||||
}
|
||||
r
|
||||
}
|
||||
async fn cleanup(&self) {
|
||||
let _ = self
|
||||
.pool
|
||||
.client()
|
||||
.database(&format!("{}_{TENANT}", self.db_name))
|
||||
.drop()
|
||||
.await;
|
||||
}
|
||||
}
|
||||
|
||||
fn finding_for(target: &str, fp: &str) -> Finding {
|
||||
let mut f = Finding::new(
|
||||
target.to_string(),
|
||||
fp.to_string(),
|
||||
"ics-probe".to_string(),
|
||||
ScanType::IcsProbe,
|
||||
"Modbus exposed".to_string(),
|
||||
"unauthenticated".to_string(),
|
||||
Severity::Critical,
|
||||
);
|
||||
f.rule_id = Some("ics-modbus-exposed".to_string());
|
||||
f
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn lease_complete_persists_findings_against_the_target() {
|
||||
let Some(h) = start().await else { return };
|
||||
let db = h.pool.for_tenant_id(TENANT).await.unwrap();
|
||||
let queue = JobQueue::new(&db);
|
||||
|
||||
// Seed a queued job.
|
||||
let job = Job::plc_provision("job-1", TENANT, "target-1", InputRef::blob("sha256:x"), 180);
|
||||
assert!(queue.enqueue(job, chrono::Utc::now()).await.unwrap());
|
||||
|
||||
// Lease it over HTTP.
|
||||
let resp = h
|
||||
.post(
|
||||
"/api/v1/werkbank/jobs/lease",
|
||||
Some(TEST_RUNNER_TOKEN),
|
||||
serde_json::json!({
|
||||
"tenant": TENANT, "runner_id": "r1", "executor": "docker",
|
||||
"labels": [], "lease_ttl_secs": 60
|
||||
}),
|
||||
)
|
||||
.send()
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(resp.status(), 200, "lease should return a job");
|
||||
let leased: LeasedJob = resp.json().await.unwrap();
|
||||
assert_eq!(leased.job.id, "job-1");
|
||||
|
||||
// Complete it with a finding.
|
||||
let mut result = JobResult::succeeded("job-1");
|
||||
result.findings = vec![finding_for("target-1", "fp-abc")];
|
||||
let resp = h
|
||||
.post(
|
||||
"/api/v1/werkbank/jobs/complete",
|
||||
Some(TEST_RUNNER_TOKEN),
|
||||
serde_json::json!({
|
||||
"tenant": TENANT, "job_id": "job-1",
|
||||
"lease_token": leased.lease_token, "result": result
|
||||
}),
|
||||
)
|
||||
.send()
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(resp.status(), 200);
|
||||
assert!(resp.json::<serde_json::Value>().await.unwrap()["recorded"]
|
||||
.as_bool()
|
||||
.unwrap());
|
||||
|
||||
// The job is now succeeded, and the finding was persisted to the target.
|
||||
assert_eq!(
|
||||
queue.get("job-1").await.unwrap().unwrap().status,
|
||||
JobStatus::Succeeded
|
||||
);
|
||||
let stored = db
|
||||
.findings()
|
||||
.find_one(mongodb::bson::doc! { "fingerprint": "fp-abc" })
|
||||
.await
|
||||
.unwrap();
|
||||
assert!(stored.is_some(), "finding should be persisted");
|
||||
|
||||
h.cleanup().await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn empty_queue_leases_nothing() {
|
||||
let Some(h) = start().await else { return };
|
||||
let resp = h
|
||||
.post(
|
||||
"/api/v1/werkbank/jobs/lease",
|
||||
Some(TEST_RUNNER_TOKEN),
|
||||
serde_json::json!({
|
||||
"tenant": TENANT, "runner_id": "r1", "executor": "docker",
|
||||
"labels": [], "lease_ttl_secs": 60
|
||||
}),
|
||||
)
|
||||
.send()
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(resp.status(), 204, "no job → 204");
|
||||
h.cleanup().await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn runner_endpoints_require_the_bearer_token() {
|
||||
let Some(h) = start().await else { return };
|
||||
let body = serde_json::json!({
|
||||
"tenant": TENANT, "runner_id": "r1", "executor": "docker",
|
||||
"labels": [], "lease_ttl_secs": 60
|
||||
});
|
||||
|
||||
let no_token = h
|
||||
.post("/api/v1/werkbank/jobs/lease", None, body.clone())
|
||||
.send()
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(no_token.status(), 401, "missing token → 401");
|
||||
|
||||
let bad_token = h
|
||||
.post("/api/v1/werkbank/jobs/lease", Some("wrong"), body)
|
||||
.send()
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(bad_token.status(), 401, "wrong token → 401");
|
||||
|
||||
h.cleanup().await;
|
||||
}
|
||||
Reference in New Issue
Block a user