diff --git a/compliance-agent/src/api/handlers/mod.rs b/compliance-agent/src/api/handlers/mod.rs index 9478e63..c69fd18 100644 --- a/compliance-agent/src/api/handlers/mod.rs +++ b/compliance-agent/src/api/handlers/mod.rs @@ -14,6 +14,7 @@ pub mod pentest_handlers; pub use pentest_handlers as pentest; pub mod sbom; pub mod scans; +pub mod werkbank_jobs; // Re-export all handler functions so routes.rs can use `handlers::function_name` pub use dto::*; diff --git a/compliance-agent/src/api/handlers/werkbank_jobs.rs b/compliance-agent/src/api/handlers/werkbank_jobs.rs new file mode 100644 index 0000000..41692c7 --- /dev/null +++ b/compliance-agent/src/api/handlers/werkbank_jobs.rs @@ -0,0 +1,193 @@ +//! 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, Request}; +use axum::http::{header, StatusCode}; +use axum::middleware::Next; +use axum::response::{IntoResponse, Response}; +use axum::Json; +use mongodb::bson::doc; +use secrecy::ExposeSecret; +use std::time::Duration; + +use compliance_core::models::werkbank::{ + CompleteRequest, CompleteResponse, HeartbeatRequest, JobResult, LeaseRequest, +}; + +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, +) -> Result { + 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, +) -> Result { + 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, +) -> Result, 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 })) +} + +/// 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 { + agent.db_pool.for_tenant_id(tenant).await.map_err(internal) +} + +/// Map any internal error to a 500. +fn internal(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")); + } +} diff --git a/compliance-agent/src/api/server.rs b/compliance-agent/src/api/server.rs index c60cf71..0f97477 100644 --- a/compliance-agent/src/api/server.rs +++ b/compliance-agent/src/api/server.rs @@ -4,7 +4,7 @@ use axum::extract::{DefaultBodyLimit, Request}; use axum::http::HeaderValue; use axum::middleware::Next; use axum::response::Response; -use axum::routing::{delete, get}; +use axum::routing::{delete, get, post}; use axum::{middleware, Extension, Router}; use tokio::sync::RwLock; use tower_http::cors::CorsLayer; @@ -72,8 +72,35 @@ pub async fn start_api_server(agent: ComplianceAgent, port: u16) -> Result<(), A 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), + ) + .layer(middleware::from_fn( + handlers::werkbank_jobs::require_runner_token, + )) + } else { + Router::new() + }; + let mut app = routes::build_router() .merge(admin_router) + .merge(werkbank_router) // Allow large artifact uploads (PLC .projectarchive, firmware images, // mobile packages) — axum's default request-body limit is only 2 MiB. .layer(DefaultBodyLimit::max(512 * 1024 * 1024)) diff --git a/compliance-agent/src/config.rs b/compliance-agent/src/config.rs index fc7d2d2..2878f3d 100644 --- a/compliance-agent/src/config.rs +++ b/compliance-agent/src/config.rs @@ -65,6 +65,7 @@ pub fn load_config() -> Result { admin_api_token: env_secret_opt("ADMIN_API_TOKEN"), tenant_registry_url: env_var_opt("TENANT_REGISTRY_URL"), plc_runtime: load_plc_runtime_config(), + werkbank_runner_token: env_secret_opt("WERKBANK_RUNNER_TOKEN"), }) } diff --git a/compliance-agent/src/pentest/cleanup.rs b/compliance-agent/src/pentest/cleanup.rs index c58c6e9..0c3fe8b 100644 --- a/compliance-agent/src/pentest/cleanup.rs +++ b/compliance-agent/src/pentest/cleanup.rs @@ -343,6 +343,7 @@ mod tests { admin_api_token: None, tenant_registry_url: None, plc_runtime: compliance_core::PlcRuntimeConfig::default(), + werkbank_runner_token: None, } } diff --git a/compliance-agent/tests/common/mod.rs b/compliance-agent/tests/common/mod.rs index fce6fa2..57ccd5b 100644 --- a/compliance-agent/tests/common/mod.rs +++ b/compliance-agent/tests/common/mod.rs @@ -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); diff --git a/compliance-agent/tests/werkbank_api.rs b/compliance-agent/tests/werkbank_api.rs new file mode 100644 index 0000000..487ba8c --- /dev/null +++ b/compliance-agent/tests/werkbank_api.rs @@ -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 { + 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::().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; +} diff --git a/compliance-core/src/auth.rs b/compliance-core/src/auth.rs index 849503f..e3bcc59 100644 --- a/compliance-core/src/auth.rs +++ b/compliance-core/src/auth.rs @@ -64,11 +64,11 @@ struct Claims { const PUBLIC_ENDPOINTS: &[&str] = &["/api/v1/health"]; /// Path prefixes that bypass JWT validation. The admin sub-router -/// (`/api/v1/admin/*`) has its own static-bearer middleware and must -/// not be routed through the customer-JWT path — a Keycloak token -/// always carries a single tenant_id and would semantically conflict -/// with cross-tenant admin operations. -const PUBLIC_PREFIXES: &[&str] = &["/api/v1/admin/"]; +/// (`/api/v1/admin/*`) and the Werkbank runner API (`/api/v1/werkbank/*`) +/// have their own static-bearer middleware and must not be routed through the +/// customer-JWT path — a Keycloak token always carries a single tenant_id and +/// would semantically conflict with these cross-tenant / machine operations. +const PUBLIC_PREFIXES: &[&str] = &["/api/v1/admin/", "/api/v1/werkbank/"]; /// Middleware that validates Bearer JWT tokens against Keycloak's JWKS /// and attaches a `TenantContext` extension on success. diff --git a/compliance-core/src/config.rs b/compliance-core/src/config.rs index 0686533..1389826 100644 --- a/compliance-core/src/config.rs +++ b/compliance-core/src/config.rs @@ -53,6 +53,11 @@ pub struct AgentConfig { /// default: it needs Docker access in the agent's runtime, which is a /// deployment opt-in. 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, } /// Configuration for the ephemeral soft-PLC "provision-and-test" path (#183). diff --git a/compliance-core/src/models/mod.rs b/compliance-core/src/models/mod.rs index 9ef2323..fb9bd58 100644 --- a/compliance-core/src/models/mod.rs +++ b/compliance-core/src/models/mod.rs @@ -49,6 +49,7 @@ pub use repository::ScanTrigger; pub use sbom::{SbomEntry, VulnRef}; pub use scan::{ScanPhase, ScanRun, ScanRunStatus, ScanType}; pub use werkbank::{ - DastCollect, Executor, HeartbeatAck, InputRef, Job, JobCollect, JobRecord, JobResult, - JobRuntime, JobStatus, JobType, LeasedJob, + CompleteRequest, CompleteResponse, DastCollect, Executor, HeartbeatAck, HeartbeatRequest, + InputRef, Job, JobCollect, JobRecord, JobResult, JobRuntime, JobStatus, JobType, LeaseRequest, + LeasedJob, }; diff --git a/compliance-core/src/models/werkbank.rs b/compliance-core/src/models/werkbank.rs index 2be2771..babea2a 100644 --- a/compliance-core/src/models/werkbank.rs +++ b/compliance-core/src/models/werkbank.rs @@ -349,6 +349,59 @@ pub struct HeartbeatAck { 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, + /// 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)] #[allow(clippy::expect_used, clippy::unwrap_used)] mod tests {