feat(job-registry): wire /api/admin/jobs/*

This commit is contained in:
Edouard Vanbelle
2026-07-27 23:23:23 +02:00
parent 72b99f5b9a
commit dfedde54a4
12 changed files with 372 additions and 79 deletions
+25 -12
View File
@@ -21,7 +21,7 @@ use chrono::Utc;
use tokio::task::JoinHandle;
use super::registry::{JobEntry, JobRegistry};
use super::types::{ErrCause, JobOutcome};
use super::types::{ErrCause, JobOutcome, JobRunArgs};
/// Public handle to the running supervisor.
///
@@ -93,8 +93,9 @@ async fn run(registry: Arc<JobRegistry>) {
// Fire and forget from the supervisor's perspective — we
// don't care about the outcome, `dispatch` records it on the
// entry and emits the log line itself.
let _ = dispatch(&name, entry).await;
// entry and emits the log line itself. Periodic ticks never
// force — that's an admin-trigger-only affordance.
let _ = dispatch(&name, entry, &JobRunArgs::default()).await;
}
}
@@ -112,7 +113,15 @@ async fn run(registry: Arc<JobRegistry>) {
///
/// Non-panicking; every failure path resolves to a `JobOutcome::Err`
/// with a `cause` log field.
pub(super) async fn dispatch(name: &str, entry: Arc<JobEntry>) -> JobOutcome {
///
/// `args` is passed through to `JobHandler::run`. The supervisor's
/// periodic ticks pass `JobRunArgs::default()`; the admin trigger
/// endpoint forwards parsed query params such as `?force=true`.
pub(super) async fn dispatch(
name: &str,
entry: Arc<JobEntry>,
args: &JobRunArgs,
) -> JobOutcome {
// Try to acquire the single-permit gate. `try_acquire` is
// non-blocking — if held, we know the previous run is still
// executing and skip this tick.
@@ -157,9 +166,11 @@ pub(super) async fn dispatch(name: &str, entry: Arc<JobEntry>) -> JobOutcome {
let start_instant = Instant::now();
// Spawn so panics land as `JoinError::is_panic()` instead of
// unwinding into the supervisor loop.
// unwinding into the supervisor loop. Args cloned into the spawn
// scope so the borrow doesn't outlive the caller.
let handler = entry.handler.clone();
let join = tokio::spawn(async move { handler.run().await });
let args_owned = args.clone();
let join = tokio::spawn(async move { handler.run(&args_owned).await });
let (outcome, cause) = match entry.timeout {
Some(dur) => match tokio::time::timeout(dur, join).await {
@@ -305,7 +316,7 @@ mod tests {
fn name(&self) -> &str {
&self.name
}
async fn run(&self) -> JobOutcome {
async fn run(&self, _args: &JobRunArgs) -> JobOutcome {
self.calls.fetch_add(1, Ordering::SeqCst);
if !self.sleep.is_zero() {
tokio::time::sleep(self.sleep).await;
@@ -321,7 +332,7 @@ mod tests {
fn name(&self) -> &str {
"panicker"
}
async fn run(&self) -> JobOutcome {
async fn run(&self, _args: &JobRunArgs) -> JobOutcome {
panic!("intentional test panic");
}
}
@@ -331,7 +342,7 @@ mod tests {
// Directly exercise translate_join with a spawned panic — the
// supervisor loop's dispatch path uses this same helper.
let handler = Arc::new(PanickingHandler);
let join = tokio::spawn(async move { handler.run().await });
let join = tokio::spawn(async move { handler.run(&JobRunArgs::default()).await });
let (outcome, cause) = translate_join(join.await);
assert!(!outcome.is_ok());
assert_eq!(cause, Some(ErrCause::Panicked));
@@ -361,13 +372,15 @@ mod tests {
// Kick off dispatch 1 in the background — it holds the permit
// for ~200 ms.
let entry_bg = entry.clone();
let bg = tokio::spawn(async move { dispatch("overrun", entry_bg).await });
let bg = tokio::spawn(async move {
dispatch("overrun", entry_bg, &JobRunArgs::default()).await
});
// Give dispatch 1 time to grab the permit.
tokio::time::sleep(Duration::from_millis(50)).await;
// Dispatch 2 should observe the permit taken and skip.
dispatch("overrun", entry.clone()).await;
dispatch("overrun", entry.clone(), &JobRunArgs::default()).await;
// Only dispatch 1's handler should have actually run so far.
assert_eq!(calls.load(Ordering::SeqCst), 1);
@@ -396,7 +409,7 @@ mod tests {
.unwrap();
let entry = registry.get("slow").await.unwrap();
dispatch("slow", entry.clone()).await;
dispatch("slow", entry.clone(), &JobRunArgs::default()).await;
// The timeout fired; last_outcome must be Err.
let state = entry.state.lock().unwrap();
+16 -7
View File
@@ -3,12 +3,13 @@
//! Everything a native service needs to write to plug into the periodic
//! scheduler is on this page. See `docs/plan/job-registry.md` Part 1
//! for the design rationale and migration criterion (the "operator
//! trigger" question — if an operator would never `POST /trigger-job`
//! for this loop, it doesn't belong here; keep it as a core worker).
//! trigger" question — if an operator would never
//! `POST /api/admin/jobs/{name}/trigger` for this loop, it doesn't
//! belong here; keep it as a core worker).
use async_trait::async_trait;
use super::types::JobOutcome;
use super::types::{JobOutcome, JobRunArgs};
/// Implemented by every service that wants to run on a fixed interval
/// through the periodic scheduler.
@@ -31,7 +32,7 @@ use super::types::JobOutcome;
///
/// Return a stable, unique snake_case identifier. Log lines
/// (`job = %name`), admin listing, admin trigger URLs
/// (`POST /api/admin/internal/trigger-job/{name}`) and env vars
/// (`POST /api/admin/jobs/{name}/trigger`) and env vars
/// (`OXICLOUD_JOB_<NAME>_INTERVAL_HOURS`) all key on this. Renaming
/// after release is a breaking change to operator scripts and log
/// dashboards.
@@ -64,7 +65,15 @@ pub trait JobHandler: Send + Sync {
fn name(&self) -> &str;
/// One execution. Called at the registered interval and (optionally)
/// on admin trigger. See trait-level docs for guidance on when to
/// return Ok vs Err.
async fn run(&self) -> JobOutcome;
/// on admin trigger.
///
/// `args` carries per-dispatch parameters (`force: bool` today).
/// Periodic ticks pass [`JobRunArgs::default()`]; admin triggers
/// forward query params such as `?force=true`. Handlers that don't
/// understand a given arg silently ignore it — the arg exists to
/// give per-job acceleration semantics without spreading per-job
/// knowledge into every caller.
///
/// See trait-level docs for guidance on when to return Ok vs Err.
async fn run(&self, args: &JobRunArgs) -> JobOutcome;
}
+2 -2
View File
@@ -29,5 +29,5 @@ mod types;
pub use engine::SchedulerEngine;
pub use handler::JobHandler;
pub use registry::{JobEntry, JobRegistry, RegisterError};
pub use types::{ErrCause, JobOutcome};
pub use registry::{JobEntry, JobRegistry, JobSummary, RegisterError};
pub use types::{ErrCause, JobOutcome, JobRunArgs};
+72 -7
View File
@@ -16,10 +16,11 @@ use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant};
use chrono::{DateTime, Utc};
use serde::Serialize;
use tokio::sync::{RwLock, Semaphore};
use super::handler::JobHandler;
use super::types::JobOutcome;
use super::types::{JobOutcome, JobRunArgs};
/// A registered job plus its runtime state. Held as `Arc<JobEntry>`
/// inside the registry so the engine can hold a snapshot across an
@@ -158,6 +159,32 @@ impl JobRegistry {
guard.iter().map(|(k, v)| (k.clone(), v.clone())).collect()
}
/// Serialisable snapshot for `GET /api/admin/jobs`. Each entry
/// captures the operator-visible state: interval (null for on-
/// demand), next scheduled dispatch (null for on-demand), when
/// the last run started, and its outcome.
pub async fn snapshot(&self) -> Vec<JobSummary> {
let entries = self.snapshot_all().await;
entries
.into_iter()
.map(|(name, entry)| {
let state = entry.state.lock().expect("JobState mutex poisoned");
let (last_run_at, last_outcome) = match &state.last_outcome {
Some((at, outcome)) => (Some(*at), Some(outcome.clone())),
None => (None, None),
};
JobSummary {
name,
interval_ms: entry.interval.map(|d| d.as_millis() as u64),
next_run_at: state.next_run_at,
last_run_at,
last_outcome,
running: state.current_run_start.is_some(),
}
})
.collect()
}
/// Count of registered jobs — used for the startup log line.
pub async fn len(&self) -> usize {
self.entries.read().await.len()
@@ -171,7 +198,7 @@ impl JobRegistry {
/// Manual dispatch — the single entry point for running a
/// registered job outside the scheduler's tick loop. Called by:
///
/// - The admin endpoint `POST /api/admin/internal/trigger-job/{name}`.
/// - The admin endpoint `POST /api/admin/jobs/{name}/trigger`.
/// - Any service that wants a scheduler-uniform dispatch of a
/// peer job (uniform log line, exclusivity, panic containment,
/// timeout enforcement).
@@ -185,9 +212,17 @@ impl JobRegistry {
///
/// Works for BOTH scheduled and on-demand jobs — for on-demand
/// jobs this is the only way they ever run.
pub async fn trigger(self: &Arc<Self>, name: &str) -> Option<JobOutcome> {
///
/// `args` is forwarded to `JobHandler::run`. Admin trigger routes
/// use `JobRunArgs { force: query.force }`; programmatic callers
/// that just want a plain run pass `JobRunArgs::default()`.
pub async fn trigger(
self: &Arc<Self>,
name: &str,
args: &JobRunArgs,
) -> Option<JobOutcome> {
let entry = self.get(name).await?;
Some(super::engine::dispatch(name, entry).await)
Some(super::engine::dispatch(name, entry, args).await)
}
}
@@ -203,6 +238,29 @@ pub enum RegisterError {
DuplicateName(String),
}
/// Per-job row in the `GET /api/admin/jobs` response.
///
/// - `interval_ms` — periodic cadence; `null` for on-demand jobs.
/// - `next_run_at` — next scheduled dispatch; `null` for on-demand.
/// - `last_run_at` / `last_outcome` — most recent completed run;
/// `null` until the first run finishes.
/// - `running` — true iff the in-flight permit is currently held
/// (either the supervisor tick is in progress or an admin trigger
/// raced in).
#[derive(Debug, Clone, Serialize)]
pub struct JobSummary {
pub name: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub interval_ms: Option<u64>,
#[serde(skip_serializing_if = "Option::is_none")]
pub next_run_at: Option<DateTime<Utc>>,
#[serde(skip_serializing_if = "Option::is_none")]
pub last_run_at: Option<DateTime<Utc>>,
#[serde(skip_serializing_if = "Option::is_none")]
pub last_outcome: Option<JobOutcome>,
pub running: bool,
}
#[cfg(test)]
mod tests {
use super::*;
@@ -217,7 +275,7 @@ mod tests {
fn name(&self) -> &str {
&self.name
}
async fn run(&self) -> JobOutcome {
async fn run(&self, _args: &JobRunArgs) -> JobOutcome {
JobOutcome::ok(0)
}
}
@@ -299,13 +357,20 @@ mod tests {
let reg = Arc::new(JobRegistry::new());
reg.register(handler("gc"), None, None).await.unwrap();
let outcome = reg.trigger("gc").await.expect("job exists");
let outcome = reg
.trigger("gc", &JobRunArgs::default())
.await
.expect("job exists");
assert!(outcome.is_ok());
}
#[tokio::test]
async fn trigger_returns_none_for_unknown_job() {
let reg = Arc::new(JobRegistry::new());
assert!(reg.trigger("nope").await.is_none());
assert!(
reg.trigger("nope", &JobRunArgs::default())
.await
.is_none()
);
}
}
+21
View File
@@ -9,6 +9,27 @@ use std::fmt;
use serde::{Deserialize, Serialize};
/// Per-dispatch parameters passed from the caller (scheduler tick or
/// admin trigger) into [`JobHandler::run`](super::handler::JobHandler::run).
///
/// Deliberately a struct — not a bare `bool` — so we don't churn every
/// handler signature the next time a job needs another knob. Grows by
/// addition; renaming a field is a breaking change to admin scripts
/// that pass query params, so treat like SQL columns.
///
/// **Handlers that don't understand a given arg silently ignore it.**
/// No error path just because a caller set an unused flag — that would
/// leak per-job semantics into callers who don't need to know.
///
/// Semantics of `force`, per job:
/// - `dedup_gc` — skip the orphan grace window (grace = 0).
/// - `grant_cleanup` — grace = 0.
/// - Others (trash_cleanup, storage_reconcile, …) — ignored.
#[derive(Debug, Clone, Default)]
pub struct JobRunArgs {
pub force: bool,
}
/// Uniform outcome the supervisor logs and stores for every job dispatch.
///
/// Two variants, deliberately. Distinguishing *why* a job failed
+23 -7
View File
@@ -3124,7 +3124,7 @@ impl DedupPort for DedupService {
/// Registered name for the dedup GC job. Stable identifier used in
/// log lines, `admin.background_runs.job_name` (when Part 2 lands),
/// and admin URLs (`POST /api/admin/internal/trigger-job/dedup_gc`).
/// and admin URLs (`POST /api/admin/jobs/dedup_gc/trigger`).
pub const DEDUP_GC_JOB_NAME: &str = "dedup_gc";
#[async_trait::async_trait]
@@ -3136,7 +3136,7 @@ impl crate::infrastructure::scheduler::JobHandler for DedupService {
/// Runs one `garbage_collect` sweep — the same reclamation that
/// `TrashCleanupService` invokes inline as its tail step, exposed
/// through the scheduler so operators can trigger it uniformly via
/// `POST /api/admin/internal/trigger-job/dedup_gc`.
/// `POST /api/admin/jobs/dedup_gc/trigger`.
///
/// Registered with `interval = None` (on-demand only): the periodic
/// tick belongs to trash cleanup, whose sweep already runs GC as
@@ -3148,12 +3148,28 @@ impl crate::infrastructure::scheduler::JobHandler for DedupService {
/// `count` reports blobs reclaimed; `extra.bytes_reclaimed` reports
/// the freed disk. GC returning `(0, 0)` is normal — it means trash
/// cleanup already reaped everything.
async fn run(&self) -> crate::infrastructure::scheduler::JobOutcome {
///
/// `args.force = true` skips the orphan grace window
/// (`garbage_collect_force` — grace_secs = 0), matching the legacy
/// `POST /admin/internal/trigger-gc?force=true` semantics. Unsafe
/// under concurrent uploads: only reachable through the admin
/// endpoint and only intentionally used by tests + operator
/// diagnostic sessions.
async fn run(
&self,
args: &crate::infrastructure::scheduler::JobRunArgs,
) -> crate::infrastructure::scheduler::JobOutcome {
use crate::infrastructure::scheduler::JobOutcome;
match self.garbage_collect().await {
Ok((items, bytes)) => {
JobOutcome::ok_with(items, serde_json::json!({ "bytes_reclaimed": bytes }))
}
let result = if args.force {
self.garbage_collect_force().await
} else {
self.garbage_collect().await
};
match result {
Ok((items, bytes)) => JobOutcome::ok_with(
items,
serde_json::json!({ "bytes_reclaimed": bytes, "forced": args.force }),
),
Err(e) => JobOutcome::Err(format!("dedup GC failed: {e}")),
}
}
@@ -23,7 +23,7 @@ use tracing::{error, info};
use crate::application::ports::authorization_ports::AuthorizationEngine;
use crate::common::errors::DomainError;
use crate::infrastructure::scheduler::{JobHandler, JobOutcome};
use crate::infrastructure::scheduler::{JobHandler, JobOutcome, JobRunArgs};
use crate::infrastructure::services::pg_acl_engine::PgAclEngine;
use async_trait::async_trait;
@@ -112,19 +112,26 @@ impl JobHandler for GrantCleanupService {
GRANT_CLEANUP_JOB_NAME
}
/// Runs one purge with the configured grace window. `count` on the
/// returned `JobOutcome::Ok` is the number of `role_grants` rows
/// physically deleted; `extra.grace_days` records which grace was
/// applied so admin listings can see it without a second lookup.
/// Runs one purge. `count` on the returned `JobOutcome::Ok` is
/// the number of `role_grants` rows physically deleted;
/// `extra.grace_days` records which grace was applied so admin
/// listings can see it without a second lookup.
///
/// Admin `?force=true` (grace = 0) does NOT come through here —
/// that path calls `purge(Some(0))` directly on the shared
/// `Arc<GrantCleanupService>` from the handler.
async fn run(&self) -> JobOutcome {
match self.purge(None).await {
Ok(count) => {
JobOutcome::ok_with(count, serde_json::json!({ "grace_days": self.grace_days }))
}
/// `args.force = true` collapses the grace window to zero for
/// this run only — matches the legacy
/// `POST /admin/internal/trigger-grant-cleanup?force=true` shape.
/// The configured `self.grace_days` is not mutated.
async fn run(&self, args: &JobRunArgs) -> JobOutcome {
let grace_override = if args.force { Some(0) } else { None };
let effective_grace = grace_override.unwrap_or(self.grace_days);
match self.purge(grace_override).await {
Ok(count) => JobOutcome::ok_with(
count,
serde_json::json!({
"grace_days": effective_grace,
"forced": args.force,
}),
),
Err(e) => JobOutcome::Err(format!("grant cleanup failed: {e}")),
}
}
@@ -6,7 +6,7 @@ use tracing::{debug, error, info, instrument};
use crate::common::errors::Result;
use crate::domain::repositories::trash_repository::TrashRepository;
use crate::infrastructure::repositories::pg::trash_db_repository::TrashDbRepository;
use crate::infrastructure::scheduler::{JobHandler, JobOutcome};
use crate::infrastructure::scheduler::{JobHandler, JobOutcome, JobRunArgs};
use crate::infrastructure::services::dedup_service::DedupService;
use async_trait::async_trait;
@@ -165,7 +165,11 @@ impl JobHandler for TrashCleanupService {
///
/// Failure of the trash sweep itself → `Err`. GC failure alone is
/// non-fatal and stays logged only.
async fn run(&self) -> JobOutcome {
///
/// `args.force` is ignored — trash cleanup has no acceleration
/// concept (retention windows are per-item metadata, not a runtime
/// knob).
async fn run(&self, _args: &JobRunArgs) -> JobOutcome {
match self.run_once().await {
Ok(stats) => {
let removed = stats.files_purged + stats.folders_purged;