From 67e2ae59cacd2819403d289570419112c41c7b76 Mon Sep 17 00:00:00 2001 From: loonghao Date: Thu, 28 May 2026 02:51:45 +0800 Subject: [PATCH] feat: expose farm observability endpoints --- crates/farm-controller/src/main.rs | 98 ++++++++- crates/farm-core/src/models.rs | 71 ++++++ crates/farm-core/src/scheduler.rs | 341 ++++++++++++++++++++++++++--- docs/architecture.md | 16 ++ 4 files changed, 490 insertions(+), 36 deletions(-) diff --git a/crates/farm-controller/src/main.rs b/crates/farm-controller/src/main.rs index 1bcbbff..bb81dee 100644 --- a/crates/farm-controller/src/main.rs +++ b/crates/farm-controller/src/main.rs @@ -9,10 +9,11 @@ use axum::routing::{get, post}; use axum::{Json, Router}; use clap::{Parser, ValueEnum}; use farm_core::{ - DashboardSnapshot, FarmError, FarmLogEntry, FarmStats, InMemoryScheduler, Job, JobId, - JobPriorityUpdate, JobSubmit, ResourceLimitDefinition, ResourceLimitSnapshot, SchedulerConfig, - SqliteScheduler, Task, TaskComplete, TaskId, TaskLease, TaskLeaseRenewal, TaskStarted, - WorkerId, WorkerInfo, WorkerLogBatch, WorkerRegister, + AuditEvent, DashboardSnapshot, FarmError, FarmLogEntry, FarmMetrics, FarmStats, + HealthComponent, HealthReport, HealthStatus, InMemoryScheduler, Job, JobId, JobPriorityUpdate, + JobSubmit, ResourceLimitDefinition, ResourceLimitSnapshot, SchedulerConfig, SqliteScheduler, + Task, TaskComplete, TaskId, TaskLease, TaskLeaseRenewal, TaskStarted, WorkerId, WorkerInfo, + WorkerLogBatch, WorkerRegister, }; use serde_json::json; use tower_http::cors::CorsLayer; @@ -85,6 +86,60 @@ impl AppScheduler { } } + fn metrics_snapshot(&self) -> Result { + match self { + Self::Memory(scheduler) => scheduler.metrics_snapshot(), + Self::Sqlite(scheduler) => scheduler.metrics_snapshot(), + } + } + + fn list_audit_events(&self) -> Result, FarmError> { + match self { + Self::Memory(scheduler) => scheduler.list_audit_events(), + Self::Sqlite(scheduler) => scheduler.list_audit_events(), + } + } + + fn backend_name(&self) -> &'static str { + match self { + Self::Memory(_) => "memory", + Self::Sqlite(_) => "sqlite", + } + } + + fn health_report(&self) -> HealthReport { + match self.metrics_snapshot() { + Ok(_) => HealthReport { + status: HealthStatus::Ready, + controller: HealthComponent { + status: HealthStatus::Ready, + backend: None, + message: Some("controller ready".to_string()), + }, + scheduler: HealthComponent { + status: HealthStatus::Ready, + backend: Some(self.backend_name().to_string()), + message: Some("scheduler ready".to_string()), + }, + degraded: Vec::new(), + }, + Err(error) => HealthReport { + status: HealthStatus::Degraded, + controller: HealthComponent { + status: HealthStatus::Ready, + backend: None, + message: Some("controller ready".to_string()), + }, + scheduler: HealthComponent { + status: HealthStatus::Degraded, + backend: Some(self.backend_name().to_string()), + message: Some(error.to_string()), + }, + degraded: vec!["scheduler".to_string()], + }, + } + } + fn list_worker_logs(&self, worker_id: WorkerId) -> Result, FarmError> { match self { Self::Memory(scheduler) => scheduler.list_worker_logs(worker_id), @@ -260,8 +315,12 @@ async fn main() -> anyhow::Result<()> { fn app(scheduler: AppScheduler) -> Router { Router::new() .route("/healthz", get(healthz)) + .route("/readyz", get(healthz)) + .route("/v1/health", get(healthz)) .route("/v1/dashboard", get(get_dashboard)) .route("/v1/logs", get(list_logs)) + .route("/v1/audit", get(list_audit_events)) + .route("/v1/metrics", get(get_metrics)) .route("/v1/stats", get(get_stats)) .route("/v1/limits", get(list_limits).post(define_limit)) .route("/v1/jobs", get(list_jobs).post(submit_job)) @@ -296,8 +355,8 @@ fn app(scheduler: AppScheduler) -> Router { .with_state(scheduler) } -async fn healthz() -> Json { - Json(json!({ "status": "ok" })) +async fn healthz(State(scheduler): State) -> Json { + Json(scheduler.health_report()) } async fn submit_job( @@ -359,6 +418,10 @@ async fn get_stats(State(scheduler): State) -> Result) -> Result, ApiError> { + Ok(Json(scheduler.metrics_snapshot()?)) +} + async fn get_dashboard( State(scheduler): State, ) -> Result, ApiError> { @@ -371,6 +434,12 @@ async fn list_logs( Ok(Json(scheduler.list_logs()?)) } +async fn list_audit_events( + State(scheduler): State, +) -> Result>, ApiError> { + Ok(Json(scheduler.list_audit_events()?)) +} + async fn define_limit( State(scheduler): State, Json(definition): Json, @@ -547,3 +616,20 @@ fn content_type_for(name: &str) -> &'static str { _ => "application/octet-stream", } } + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn health_report_includes_controller_and_scheduler_status() { + let scheduler = AppScheduler::Memory(InMemoryScheduler::default()); + let report = scheduler.health_report(); + + assert_eq!(report.status, HealthStatus::Ready); + assert_eq!(report.controller.status, HealthStatus::Ready); + assert_eq!(report.scheduler.status, HealthStatus::Ready); + assert_eq!(report.scheduler.backend.as_deref(), Some("memory")); + assert!(report.degraded.is_empty()); + } +} diff --git a/crates/farm-core/src/models.rs b/crates/farm-core/src/models.rs index fdef237..f7fc1c1 100644 --- a/crates/farm-core/src/models.rs +++ b/crates/farm-core/src/models.rs @@ -462,6 +462,44 @@ pub struct FarmStats { pub worker_slots_available: u32, } +#[derive(Debug, Clone, Default, Serialize, Deserialize)] +pub struct FarmMetrics { + pub queue_depth: usize, + pub running_tasks: usize, + pub failed_tasks: usize, + pub lease_renewals: u64, + pub scheduler_errors: u64, + pub workers_online: usize, + pub workers_offline: usize, + pub stats: FarmStats, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct HealthReport { + pub status: HealthStatus, + pub controller: HealthComponent, + pub scheduler: HealthComponent, + #[serde(default)] + pub degraded: Vec, +} + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +#[serde(rename_all = "snake_case")] +pub enum HealthStatus { + #[serde(rename = "ok")] + Ready, + Degraded, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct HealthComponent { + pub status: HealthStatus, + #[serde(default)] + pub backend: Option, + #[serde(default)] + pub message: Option, +} + #[derive(Debug, Clone, Serialize, Deserialize)] pub struct Task { pub id: TaskId, @@ -640,6 +678,39 @@ pub struct FarmLogEntry { pub worker_id: Option, } +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct AuditEvent { + pub id: Uuid, + pub timestamp: DateTime, + pub actor: String, + pub action: String, + pub target_type: String, + #[serde(default)] + pub target_id: Option, + pub outcome: AuditOutcome, + #[serde(default)] + pub message: Option, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct AuditEventInput { + pub actor: String, + pub action: String, + pub target_type: String, + #[serde(default)] + pub target_id: Option, + pub outcome: AuditOutcome, + #[serde(default)] + pub message: Option, +} + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +#[serde(rename_all = "snake_case")] +pub enum AuditOutcome { + Success, + Failure, +} + #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] #[serde(rename_all = "snake_case")] pub enum LogLevel { diff --git a/crates/farm-core/src/scheduler.rs b/crates/farm-core/src/scheduler.rs index 59d7037..1c65492 100644 --- a/crates/farm-core/src/scheduler.rs +++ b/crates/farm-core/src/scheduler.rs @@ -9,11 +9,12 @@ use thiserror::Error; use uuid::Uuid; use crate::models::{ - DashboardSnapshot, FarmLogEntry, FarmStats, Job, JobId, JobState, JobSubmit, LogLevel, - LogSource, OpenJdAmountRequirement, OpenJdAttributeRequirement, ResourceLimitDefinition, - ResourceLimitSnapshot, Task, TaskArtifact, TaskAttempt, TaskAttemptState, TaskComplete, TaskId, - TaskLease, TaskLeaseInfo, TaskLeaseRenewal, TaskStarted, TaskState, TaskSubmit, WorkerCapacity, - WorkerId, WorkerInfo, WorkerLogBatch, WorkerRegister, WorkerState, + AuditEvent, AuditEventInput, AuditOutcome, DashboardSnapshot, FarmLogEntry, FarmMetrics, + FarmStats, Job, JobId, JobState, JobSubmit, LogLevel, LogSource, OpenJdAmountRequirement, + OpenJdAttributeRequirement, ResourceLimitDefinition, ResourceLimitSnapshot, Task, TaskArtifact, + TaskAttempt, TaskAttemptState, TaskComplete, TaskId, TaskLease, TaskLeaseInfo, + TaskLeaseRenewal, TaskStarted, TaskState, TaskSubmit, WorkerCapacity, WorkerId, WorkerInfo, + WorkerLogBatch, WorkerRegister, WorkerState, }; use crate::openjd::openjd_to_tasks; @@ -68,10 +69,22 @@ struct SchedulerState { jobs: HashMap, workers: HashMap, logs: Vec, + #[serde(default)] + audit_events: Vec, + #[serde(default)] + observability: ObservabilityCounters, + #[serde(default)] limits: HashMap, } const MAX_LOG_ENTRIES: usize = 500; +const MAX_AUDIT_EVENTS: usize = 500; + +#[derive(Debug, Clone, Default, Serialize, Deserialize)] +struct ObservabilityCounters { + lease_renewals: u64, + scheduler_errors: u64, +} impl InMemoryScheduler { pub fn with_config(config: SchedulerConfig) -> Self { @@ -163,6 +176,14 @@ impl InMemoryScheduler { stream: None, }, ); + audit_success( + &mut state, + "api", + "job.submit", + "job", + Some(job.id.to_string()), + Some(job.name.clone()), + ); Ok(job) } @@ -209,6 +230,21 @@ impl InMemoryScheduler { Ok(state.logs.clone()) } + pub fn metrics_snapshot(&self) -> Result { + let state = self.inner.lock().map_err(|_| FarmError::LockPoisoned)?; + Ok(metrics_snapshot(&state)) + } + + pub fn list_audit_events(&self) -> Result, FarmError> { + let state = self.inner.lock().map_err(|_| FarmError::LockPoisoned)?; + Ok(state.audit_events.clone()) + } + + pub fn record_audit_event(&self, input: AuditEventInput) -> Result { + let mut state = self.inner.lock().map_err(|_| FarmError::LockPoisoned)?; + Ok(push_audit(&mut state, input)) + } + pub fn list_worker_logs(&self, worker_id: WorkerId) -> Result, FarmError> { let state = self.inner.lock().map_err(|_| FarmError::LockPoisoned)?; if !state.workers.contains_key(&worker_id) { @@ -248,6 +284,14 @@ impl InMemoryScheduler { stream: None, }, ); + audit_success( + &mut state, + "worker", + "worker.register", + "worker", + Some(worker.id.to_string()), + Some(worker.name.clone()), + ); Ok(worker) } @@ -331,6 +375,14 @@ impl InMemoryScheduler { state .limits .insert(definition.name.clone(), definition.clone()); + audit_success( + &mut state, + "api", + "limit.define", + "limit", + Some(definition.name.clone()), + Some(format!("max_count={}", definition.max_count)), + ); Ok(limit_snapshot(&state, &definition)) } @@ -347,7 +399,16 @@ impl InMemoryScheduler { .ok_or(FarmError::WorkerNotFound(worker_id))?; worker.state = WorkerState::Online; worker.last_seen_at = Utc::now(); - Ok(worker.clone()) + let worker = worker.clone(); + audit_success( + &mut state, + format!("worker:{}", worker.id), + "worker.heartbeat", + "worker", + Some(worker.id.to_string()), + None, + ); + Ok(worker) } pub fn pause_job(&self, job_id: JobId) -> Result { @@ -369,7 +430,16 @@ impl InMemoryScheduler { ))); } } - Ok(job.clone()) + let job = job.clone(); + audit_success( + &mut state, + "api", + "job.pause", + "job", + Some(job.id.to_string()), + None, + ); + Ok(job) } pub fn resume_job(&self, job_id: JobId) -> Result { @@ -388,7 +458,16 @@ impl InMemoryScheduler { ))); } } - Ok(job.clone()) + let job = job.clone(); + audit_success( + &mut state, + "api", + "job.resume", + "job", + Some(job.id.to_string()), + None, + ); + Ok(job) } pub fn cancel_job(&self, job_id: JobId) -> Result { @@ -434,7 +513,16 @@ impl InMemoryScheduler { } } } - Ok(job.clone()) + let job = job.clone(); + audit_success( + &mut state, + "api", + "job.cancel", + "job", + Some(job.id.to_string()), + None, + ); + Ok(job) } pub fn update_job_priority(&self, job_id: JobId, priority: i32) -> Result { @@ -445,7 +533,16 @@ impl InMemoryScheduler { .ok_or(FarmError::JobNotFound(job_id))?; job.priority = priority; job.updated_at = Utc::now(); - Ok(job.clone()) + let job = job.clone(); + audit_success( + &mut state, + "api", + "job.priority", + "job", + Some(job.id.to_string()), + Some(format!("priority={priority}")), + ); + Ok(job) } pub fn lease_task(&self, worker_id: WorkerId) -> Result, FarmError> { @@ -547,6 +644,14 @@ impl InMemoryScheduler { stream: None, }, ); + audit_success( + &mut state, + format!("worker:{worker_id}"), + "task.lease", + "task", + Some(task_id.to_string()), + Some(task_name), + ); Ok(Some(TaskLease { task: leased_task, @@ -597,6 +702,14 @@ impl InMemoryScheduler { } let task = task.clone(); update_job_state(job); + audit_success( + &mut state, + "api", + "task.cancel", + "task", + Some(task.id.to_string()), + Some(task.name.clone()), + ); Ok(task) } @@ -630,6 +743,14 @@ impl InMemoryScheduler { } let task = task.clone(); update_job_state(job); + audit_success( + &mut state, + "api", + "task.requeue", + "task", + Some(task.id.to_string()), + Some(task.name.clone()), + ); Ok(task) } @@ -670,6 +791,14 @@ impl InMemoryScheduler { stream: None, }, ); + audit_success( + &mut state, + format!("worker:{}", started.worker_id), + "task.start", + "task", + Some(task.id.to_string()), + Some(task.name.clone()), + ); Ok(task) } @@ -680,22 +809,34 @@ impl InMemoryScheduler { ) -> Result { let mut state = self.inner.lock().map_err(|_| FarmError::LockPoisoned)?; let (job_id, task_index) = find_task_location(&state, task_id)?; - let job = state - .jobs - .get_mut(&job_id) - .ok_or(FarmError::JobNotFound(job_id))?; - let task = job - .tasks - .get_mut(task_index) - .ok_or(FarmError::TaskNotFound(task_id))?; - validate_lease(task, renewal.worker_id, &renewal.lease_token)?; + let task = { + let job = state + .jobs + .get_mut(&job_id) + .ok_or(FarmError::JobNotFound(job_id))?; + let task = job + .tasks + .get_mut(task_index) + .ok_or(FarmError::TaskNotFound(task_id))?; + validate_lease(task, renewal.worker_id, &renewal.lease_token)?; - let now = Utc::now(); - if let Some(lease) = &mut task.lease { - lease.expires_at = now + self.lease_ttl(); - } - task.updated_at = now; - Ok(task.clone()) + let now = Utc::now(); + if let Some(lease) = &mut task.lease { + lease.expires_at = now + self.lease_ttl(); + } + task.updated_at = now; + task.clone() + }; + state.observability.lease_renewals += 1; + audit_success( + &mut state, + format!("worker:{}", renewal.worker_id), + "task.lease_renew", + "task", + Some(task.id.to_string()), + Some(task.name.clone()), + ); + Ok(task) } pub fn complete_task( @@ -736,7 +877,16 @@ impl InMemoryScheduler { }, ); task.lease = None; - return Ok(task.clone()); + let task = task.clone(); + audit_success( + &mut state, + format!("worker:{}", completion.worker_id), + "task.complete_cancelled", + "task", + Some(task.id.to_string()), + Some(format!("exit_code={}", completion.exit_code)), + ); + return Ok(task); } task.completed_at = Some(completed_at); @@ -802,6 +952,14 @@ impl InMemoryScheduler { stream: None, }, ); + audit_success( + &mut state, + format!("worker:{}", completion.worker_id), + "task.complete", + "task", + Some(task.id.to_string()), + Some(format!("exit_code={}", completion.exit_code)), + ); Ok(task) } @@ -882,6 +1040,18 @@ impl SqliteScheduler { self.inner.list_logs() } + pub fn metrics_snapshot(&self) -> Result { + self.inner.metrics_snapshot() + } + + pub fn list_audit_events(&self) -> Result, FarmError> { + self.inner.list_audit_events() + } + + pub fn record_audit_event(&self, input: AuditEventInput) -> Result { + self.write_durably(|inner| inner.record_audit_event(input)) + } + pub fn list_worker_logs(&self, worker_id: WorkerId) -> Result, FarmError> { self.inner.list_worker_logs(worker_id) } @@ -1446,6 +1616,49 @@ fn push_log(state: &mut SchedulerState, log: NewLog) -> FarmLogEntry { entry } +fn push_audit(state: &mut SchedulerState, input: AuditEventInput) -> AuditEvent { + if input.outcome == AuditOutcome::Failure { + state.observability.scheduler_errors += 1; + } + let event = AuditEvent { + id: Uuid::new_v4(), + timestamp: Utc::now(), + actor: input.actor, + action: input.action, + target_type: input.target_type, + target_id: input.target_id, + outcome: input.outcome, + message: input.message, + }; + state.audit_events.push(event.clone()); + if state.audit_events.len() > MAX_AUDIT_EVENTS { + let overflow = state.audit_events.len() - MAX_AUDIT_EVENTS; + state.audit_events.drain(0..overflow); + } + event +} + +fn audit_success( + state: &mut SchedulerState, + actor: impl Into, + action: impl Into, + target_type: impl Into, + target_id: Option, + message: Option, +) { + push_audit( + state, + AuditEventInput { + actor: actor.into(), + action: action.into(), + target_type: target_type.into(), + target_id, + outcome: AuditOutcome::Success, + message, + }, + ); +} + fn find_task_location( state: &SchedulerState, task_id: TaskId, @@ -1549,13 +1762,27 @@ fn compute_stats(state: &SchedulerState) -> FarmStats { stats } +fn metrics_snapshot(state: &SchedulerState) -> FarmMetrics { + let stats = compute_stats(state); + FarmMetrics { + queue_depth: stats.tasks_pending, + running_tasks: stats.tasks_running + stats.tasks_leased, + failed_tasks: stats.tasks_failed, + lease_renewals: state.observability.lease_renewals, + scheduler_errors: state.observability.scheduler_errors, + workers_online: stats.workers_online, + workers_offline: stats.workers_offline, + stats, + } +} + #[cfg(test)] mod tests { use crate::models::{ - ArtifactKind, CommandSpec, LogLevel, OpenJdAmountRequirement, OpenJdAttributeRequirement, - ResourceLimitDefinition, TaskArtifact, TaskAttemptState, TaskLease, TaskLeaseRenewal, - TaskRequirements, TaskStarted, WorkerCapacity, WorkerId, WorkerInfo, WorkerLogBatch, - WorkerLogInput, + ArtifactKind, AuditEventInput, AuditOutcome, CommandSpec, LogLevel, + OpenJdAmountRequirement, OpenJdAttributeRequirement, ResourceLimitDefinition, TaskArtifact, + TaskAttemptState, TaskLease, TaskLeaseRenewal, TaskRequirements, TaskStarted, + WorkerCapacity, WorkerId, WorkerInfo, WorkerLogBatch, WorkerLogInput, }; use super::*; @@ -1669,6 +1896,60 @@ mod tests { assert_eq!(snapshot.workers[0].name, "render-node-01"); } + #[test] + fn metrics_snapshot_reports_operational_counters() { + let scheduler = InMemoryScheduler::default(); + let worker = register_worker(&scheduler); + let lease = lease_single_task(&scheduler, worker.id); + scheduler + .renew_task_lease( + lease.task.id, + TaskLeaseRenewal { + worker_id: worker.id, + lease_token: lease.lease_token.clone(), + }, + ) + .expect("lease should renew"); + scheduler + .complete_task( + lease.task.id, + TaskComplete { + worker_id: worker.id, + lease_token: lease.lease_token, + exit_code: 1, + stdout_tail: None, + stderr_tail: Some("failed".to_string()), + artifacts: Vec::new(), + }, + ) + .expect("task should fail"); + scheduler + .record_audit_event(AuditEventInput { + actor: "test".to_string(), + action: "scheduler.error".to_string(), + target_type: "scheduler".to_string(), + target_id: None, + outcome: AuditOutcome::Failure, + message: Some("synthetic failure".to_string()), + }) + .expect("audit should record"); + + let metrics = scheduler.metrics_snapshot().expect("metrics should load"); + assert_eq!(metrics.queue_depth, 0); + assert_eq!(metrics.running_tasks, 0); + assert_eq!(metrics.failed_tasks, 1); + assert_eq!(metrics.lease_renewals, 1); + assert_eq!(metrics.scheduler_errors, 1); + assert_eq!(metrics.workers_online, 1); + assert_eq!(metrics.workers_offline, 0); + + let audit = scheduler.list_audit_events().expect("audit should load"); + assert!(audit.iter().any(|event| event.action == "task.lease_renew")); + assert!(audit + .iter() + .any(|event| event.outcome == AuditOutcome::Failure)); + } + #[test] fn renews_task_lease_before_it_expires() { let scheduler = InMemoryScheduler::with_config(SchedulerConfig { diff --git a/docs/architecture.md b/docs/architecture.md index c0d48d8..aa393ec 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -20,6 +20,9 @@ Core endpoints: - `POST /v1/jobs`: submit a direct task job or an OpenJD-backed job. - `GET /v1/jobs/{job_id}`: inspect state and task output tails. +- `GET /healthz`, `GET /readyz`, `GET /v1/health`: inspect controller and scheduler readiness. +- `GET /v1/metrics`: read queue, worker, lease-renewal, and scheduler-error metrics. +- `GET /v1/audit`: read recent state-changing control-plane events. - `POST /v1/workers/register`: register a worker. - `POST /v1/workers/{worker_id}/heartbeat`: keep worker online. - `POST /v1/workers/{worker_id}/lease`: get the next runnable task. @@ -64,6 +67,19 @@ runnable tasks until usage is released by completion, cancellation, requeue, or lease expiry. Limit snapshots are exposed in the dashboard payload and the dedicated `/v1/limits` API. +## Observability + +Health endpoints return a machine-readable `HealthReport` with controller +readiness, scheduler backend name, scheduler readiness, and degraded component +names. Small farms can poll `/healthz`; deployment systems that separate +startup and readiness can use `/readyz` with the same payload. + +`/v1/metrics` returns a `FarmMetrics` payload with queue depth, running tasks, +failed tasks, worker online/offline counts, lease renewal count, scheduler error +count, and the full queue `FarmStats` snapshot. `/v1/audit` returns recent +state-changing events with actor, action, target type/id, timestamp, outcome, +and message so operators can review queue mutations and worker transitions. + ## OpenJD path The controller accepts an OpenJD template bundle under `openjd.template_yaml`. It uses the official OpenJD Rust crates instead of a hand-written parser: