feat(job-registry): wire grant cleaner
This commit is contained in:
+19
-3
@@ -1487,9 +1487,25 @@ impl AppServiceFactory {
|
|||||||
core.config.features.grant_cleanup.interval_hours,
|
core.config.features.grant_cleanup.interval_hours,
|
||||||
),
|
),
|
||||||
);
|
);
|
||||||
// First tick fires immediately inside start_cleanup_job —
|
// Registered with the periodic-job scheduler
|
||||||
// matches the trash/storage-usage daemon shape.
|
// (`docs/plan/job-registry.md` Part 1); the retired
|
||||||
svc.clone().start_cleanup_job().await;
|
// `start_cleanup_job` used to spawn its own interval loop.
|
||||||
|
// Admin `?force=true` trigger still calls `svc.purge(Some(0))`
|
||||||
|
// directly — grace override doesn't fit the JobHandler shape.
|
||||||
|
let interval = svc.interval();
|
||||||
|
if let Err(e) = core
|
||||||
|
.job_registry
|
||||||
|
.register(svc.clone(), Some(interval), None)
|
||||||
|
.await
|
||||||
|
{
|
||||||
|
tracing::error!("Failed to register grant_cleanup job: {e}");
|
||||||
|
} else {
|
||||||
|
tracing::info!(
|
||||||
|
"Grant cleanup registered with scheduler (every {}h, grace = {}d)",
|
||||||
|
interval.as_secs() / 3600,
|
||||||
|
core.config.features.grant_cleanup.grace_days,
|
||||||
|
);
|
||||||
|
}
|
||||||
Some(svc)
|
Some(svc)
|
||||||
} else {
|
} else {
|
||||||
tracing::info!(
|
tracing::info!(
|
||||||
|
|||||||
@@ -1,36 +1,39 @@
|
|||||||
//! Background daemon that purges expired `storage.role_grants` rows.
|
//! Service that purges expired `storage.role_grants` rows.
|
||||||
//!
|
//!
|
||||||
//! The AuthZ engine already filters expired grants out of every
|
//! The AuthZ engine already filters expired grants out of every
|
||||||
//! permission check at read time (`expires_at IS NULL OR
|
//! permission check at read time (`expires_at IS NULL OR
|
||||||
//! expires_at > NOW()` on every `check` / `list_grants_*` path in
|
//! expires_at > NOW()` on every `check` / `list_grants_*` path in
|
||||||
//! `PgAclEngine`), so expired rows never leak permission. They just
|
//! `PgAclEngine`), so expired rows never leak permission. They just
|
||||||
//! accumulate. This daemon garbage-collects them once per
|
//! accumulate. This service garbage-collects them, with a grace window
|
||||||
//! [`GrantCleanupService::interval_hours`], with a grace window past
|
//! past `expires_at` that preserves the audit / support answer to
|
||||||
//! `expires_at` that preserves the audit / support answer to "what
|
//! "what happened to my access?" for a few weeks.
|
||||||
//! happened to my access?" for a few weeks.
|
|
||||||
//!
|
//!
|
||||||
//! Shape mirrors [`TrashCleanupService`] verbatim (fire-and-forget
|
//! **Scheduling.** Registered with the periodic-job scheduler
|
||||||
//! `tokio::spawn`, `tokio::time::interval`, first-tick-immediate). The
|
//! (`docs/plan/job-registry.md` Part 1). The retired `start_cleanup_job`
|
||||||
//! authoritative pattern for background daemons in this codebase; see
|
//! used to spawn its own `tokio::interval` loop; the scheduler now
|
||||||
//! the plan doc `docs/plan/` (deferred future work: fold all daemons
|
//! dispatches [`GrantCleanupService::purge`] on the configured cadence
|
||||||
//! into a central `JobRegistry` that plugins can also register into).
|
//! and handles panic containment + exclusivity + admin trigger routing.
|
||||||
//!
|
//! Admin trigger with `?force=true` still bypasses the registered
|
||||||
//! [`TrashCleanupService`]: crate::infrastructure::services::trash_cleanup_service::TrashCleanupService
|
//! job and calls `purge(Some(0))` directly so the grace override reaches
|
||||||
|
//! the underlying SQL.
|
||||||
|
|
||||||
use std::sync::Arc;
|
use std::sync::Arc;
|
||||||
use std::time::{Duration, Instant};
|
use std::time::{Duration, Instant};
|
||||||
use tokio::time;
|
|
||||||
use tracing::{error, info};
|
use tracing::{error, info};
|
||||||
|
|
||||||
use crate::application::ports::authorization_ports::AuthorizationEngine;
|
use crate::application::ports::authorization_ports::AuthorizationEngine;
|
||||||
|
use crate::common::errors::DomainError;
|
||||||
|
use crate::infrastructure::scheduler::{JobHandler, JobOutcome};
|
||||||
use crate::infrastructure::services::pg_acl_engine::PgAclEngine;
|
use crate::infrastructure::services::pg_acl_engine::PgAclEngine;
|
||||||
|
use async_trait::async_trait;
|
||||||
|
|
||||||
/// Daemon that periodically deletes expired grants.
|
pub const GRANT_CLEANUP_JOB_NAME: &str = "grant_cleanup";
|
||||||
|
|
||||||
|
/// Service that deletes expired grants.
|
||||||
///
|
///
|
||||||
/// Owns an `Arc<PgAclEngine>` (not a `dyn AuthorizationEngine`) to avoid
|
/// Owns an `Arc<PgAclEngine>` (not a `dyn AuthorizationEngine`) to avoid
|
||||||
/// the wrapper allocation on every SQL call — the daemon is the sole
|
/// the wrapper allocation on every SQL call — the caller set is small
|
||||||
/// caller of `purge_expired_grants` outside of the admin trigger
|
/// (scheduler tick + admin trigger endpoint), both statically dispatched.
|
||||||
/// endpoint, both statically dispatched.
|
|
||||||
pub struct GrantCleanupService {
|
pub struct GrantCleanupService {
|
||||||
authz: Arc<PgAclEngine>,
|
authz: Arc<PgAclEngine>,
|
||||||
grace_days: u32,
|
grace_days: u32,
|
||||||
@@ -48,50 +51,35 @@ impl GrantCleanupService {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Grace period the daemon uses on its scheduled ticks. Exposed
|
/// Grace period the service uses on its scheduled runs. Exposed
|
||||||
/// for the admin trigger's default-response field.
|
/// for the admin trigger's default-response field.
|
||||||
pub fn grace_days(&self) -> u32 {
|
pub fn grace_days(&self) -> u32 {
|
||||||
self.grace_days
|
self.grace_days
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Fire-and-forget the periodic purge. Never joins; killed
|
/// Cadence exposed as `Duration` so DI passes a sanitised value
|
||||||
/// implicitly at `tokio::runtime::shutdown`.
|
/// (post-`.max(1)`) to `JobRegistry::register`.
|
||||||
pub async fn start_cleanup_job(self: Arc<Self>) {
|
pub fn interval(&self) -> Duration {
|
||||||
let interval_hours = self.interval_hours;
|
Duration::from_secs(self.interval_hours * 3600)
|
||||||
let grace_days = self.grace_days;
|
|
||||||
info!(
|
|
||||||
"Starting grant-cleanup daemon: every {}h, grace = {}d",
|
|
||||||
interval_hours, grace_days
|
|
||||||
);
|
|
||||||
|
|
||||||
tokio::spawn(async move {
|
|
||||||
let mut interval = time::interval(Duration::from_secs(interval_hours * 60 * 60));
|
|
||||||
// First tick fires immediately — matches TrashCleanupService.
|
|
||||||
// Any accumulated backlog at boot gets flushed straight away.
|
|
||||||
loop {
|
|
||||||
interval.tick().await;
|
|
||||||
self.run_once().await;
|
|
||||||
}
|
|
||||||
});
|
|
||||||
}
|
}
|
||||||
|
|
||||||
/// One scheduled pass. Also called by the admin trigger endpoint
|
/// Run one purge pass.
|
||||||
/// (via a shared `Arc<GrantCleanupService>` on `AppState`).
|
|
||||||
///
|
///
|
||||||
/// `grace_override`:
|
/// `grace_override`:
|
||||||
/// - `None` → use the configured grace (`self.grace_days`).
|
/// - `None` → use the configured grace (`self.grace_days`).
|
||||||
/// - `Some(n)` → override with `n`. The admin `?force=true` trigger
|
/// - `Some(n)` → override with `n`. The admin `?force=true` trigger
|
||||||
/// passes `Some(0)` so Hurl regressions can hit expired grants
|
/// passes `Some(0)` so Hurl regressions can hit expired grants
|
||||||
/// without waiting the configured grace out.
|
/// without waiting the configured grace out.
|
||||||
pub async fn purge(&self, grace_override: Option<u32>) -> u64 {
|
///
|
||||||
|
/// Returns `Ok(count)` on success, `Err(_)` on DB error. Audit-log
|
||||||
|
/// lines fire on both paths (success + failure) — bulk deletion of
|
||||||
|
/// authorization rows is security-relevant enough to log even a
|
||||||
|
/// zero-count run, and failures MUST reach the audit channel.
|
||||||
|
pub async fn purge(&self, grace_override: Option<u32>) -> Result<u64, DomainError> {
|
||||||
let grace = grace_override.unwrap_or(self.grace_days);
|
let grace = grace_override.unwrap_or(self.grace_days);
|
||||||
let start = Instant::now();
|
let start = Instant::now();
|
||||||
match self.authz.purge_expired_grants(grace).await {
|
match self.authz.purge_expired_grants(grace).await {
|
||||||
Ok(count) => {
|
Ok(count) => {
|
||||||
// Audit-channel logging: bulk deletion of authorization
|
|
||||||
// rows is security-relevant enough to keep it in the
|
|
||||||
// audit stream even when the count is zero (proves the
|
|
||||||
// daemon is reachable).
|
|
||||||
info!(
|
info!(
|
||||||
target: "audit",
|
target: "audit",
|
||||||
event = "grant_cleanup.purged",
|
event = "grant_cleanup.purged",
|
||||||
@@ -102,7 +90,7 @@ impl GrantCleanupService {
|
|||||||
count,
|
count,
|
||||||
grace,
|
grace,
|
||||||
);
|
);
|
||||||
count
|
Ok(count)
|
||||||
}
|
}
|
||||||
Err(e) => {
|
Err(e) => {
|
||||||
error!(
|
error!(
|
||||||
@@ -112,13 +100,33 @@ impl GrantCleanupService {
|
|||||||
error = %e,
|
error = %e,
|
||||||
"Grant cleanup failed"
|
"Grant cleanup failed"
|
||||||
);
|
);
|
||||||
0
|
Err(e)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
|
||||||
/// Convenience for the scheduled loop.
|
#[async_trait]
|
||||||
async fn run_once(&self) {
|
impl JobHandler for GrantCleanupService {
|
||||||
let _ = self.purge(None).await;
|
fn name(&self) -> &str {
|
||||||
|
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.
|
||||||
|
///
|
||||||
|
/// 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 }),
|
||||||
|
),
|
||||||
|
Err(e) => JobOutcome::Err(format!("grant cleanup failed: {e}")),
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -2280,7 +2280,12 @@ pub async fn internal_trigger_grant_cleanup(
|
|||||||
// only — the daemon's configured grace is untouched. Mirrors the
|
// only — the daemon's configured grace is untouched. Mirrors the
|
||||||
// `trigger-gc?force=true` shape.
|
// `trigger-gc?force=true` shape.
|
||||||
let grace_override = if query.force { Some(0) } else { None };
|
let grace_override = if query.force { Some(0) } else { None };
|
||||||
let grants_deleted = svc.purge(grace_override).await;
|
let grants_deleted = match svc.purge(grace_override).await {
|
||||||
|
Ok(n) => n,
|
||||||
|
Err(e) => {
|
||||||
|
return AppError::internal_error(format!("grant cleanup failed: {e}")).into_response();
|
||||||
|
}
|
||||||
|
};
|
||||||
let grace_days = grace_override.unwrap_or_else(|| svc.grace_days());
|
let grace_days = grace_override.unwrap_or_else(|| svc.grace_days());
|
||||||
(
|
(
|
||||||
StatusCode::OK,
|
StatusCode::OK,
|
||||||
|
|||||||
Reference in New Issue
Block a user