diff --git a/docs/plan/job-registry.md b/docs/plan/job-registry.md index 0c78127b..291349a4 100644 --- a/docs/plan/job-registry.md +++ b/docs/plan/job-registry.md @@ -865,8 +865,14 @@ backend slice. Order of appearance: 3. `consistency_batch` + more tenants — done (Slices 5–6: drives + folders + files, plus batch). 4. Frontend `/admin/jobs` page — takes the completed backend surface as-is; no backend changes required by the UI landing. -5. **Deferred, post-UI:** progress estimation (`fraction`, `kind` on - `RunSummary`). See memory `project_job_progress_estimation`. +5. Progress estimation on `RunSummary.progress` (`fraction`, `kind`, + `scanned`, `total`) — **done (Slice 9)**. Tenants that CAN count + their subject override `RecoverableJobHandler::count_total()`; + `run_or_resume` seeds `params.total_rows` + `params.progress_kind` + on fresh runs; `row_to_summary` derives the `progress` block at + serialisation time. UI renders a bar; `kind = "approximate"` runs + get a striped fill so operators recognise proxy-derived + estimates. See memory `project_job_progress_estimation`. ### Notifications & alerting diff --git a/frontend/src/lib/api/types.ts b/frontend/src/lib/api/types.ts index dd3a7fb9..c23efc17 100644 --- a/frontend/src/lib/api/types.ts +++ b/frontend/src/lib/api/types.ts @@ -578,6 +578,31 @@ export interface RunSummary { params: Record; cursor_hex?: string; error_message?: string; + /** Populated when the tenant reported a countable subject at run + * start (`RecoverableJobHandler::count_total`). Absent when the + * tenant can't count — the UI hides the progress bar and falls + * back to raw `scanned_count`. */ + progress?: RunProgress; +} + +/** + * Confidence level of a `RunProgress` fraction. Wire lowercase per + * the `#[serde(rename_all = "lowercase")]` on the Rust enum. + * + * - `count` — `scanned_count / total_rows` where `total_rows` came + * from a definitive `COUNT(*)` on the subject table. + * - `approximate` — proxy-derived total (e.g. `storage_consistency` + * using DB blob count as a stand-in for backend object count). + * Fraction can legitimately exceed 1.0 at run end — the deviation + * quantifies the drift the check is looking for. + */ +export type ProgressKind = 'count' | 'approximate'; + +export interface RunProgress { + fraction: number; + kind: ProgressKind; + scanned: number; + total: number; } /** diff --git a/frontend/src/lib/components/AdminJobsPanel.svelte b/frontend/src/lib/components/AdminJobsPanel.svelte index 1f089a16..9086d07c 100644 --- a/frontend/src/lib/components/AdminJobsPanel.svelte +++ b/frontend/src/lib/components/AdminJobsPanel.svelte @@ -483,7 +483,9 @@ {t('admin.jobs.col_started_at', 'Started')} {t('admin.jobs.col_status', 'Status')} {t('admin.jobs.col_duration', 'Duration')} - {t('admin.jobs.col_scanned', 'Scanned')} + + {t('admin.jobs.col_progress', 'Progress')} + {t('admin.jobs.col_findings', 'Findings')} {t('admin.jobs.col_error', 'Error')} @@ -517,8 +519,52 @@ {runDurationLabel(run)} - - {scanned ?? '—'} + + {#if run.progress} + {@const barPct = Math.min( + 100, + Math.max(0, run.progress.fraction * 100) + )} + {@const pctLabel = (run.progress.fraction * 100).toFixed(1) + '%'} +
+
+
+
+ + {run.progress.scanned}/{run.progress.total} + +
+ {:else} + {scanned ?? '—'} + {/if} {findingCount ?? 0} @@ -941,4 +987,45 @@ opacity: 0.5; cursor: not-allowed; } + + .jobs-panel__col-progress { + min-width: 12rem; + } + + .jobs-panel__progress { + display: flex; + align-items: center; + gap: 0.5rem; + } + + .jobs-panel__progress-bar { + flex: 1; + height: 0.5rem; + background: var(--color-bg-subtle); + border-radius: 999px; + overflow: hidden; + } + + .jobs-panel__progress-fill { + height: 100%; + background: var(--color-accent); + transition: width 0.25s ease-out; + } + + .jobs-panel__progress-bar--approx .jobs-panel__progress-fill { + background: repeating-linear-gradient( + 45deg, + var(--color-accent), + var(--color-accent) 6px, + var(--color-accent-hover) 6px, + var(--color-accent-hover) 12px + ); + } + + .jobs-panel__progress-label { + font-variant-numeric: tabular-nums; + font-size: 0.8rem; + color: var(--color-text-muted); + white-space: nowrap; + } diff --git a/frontend/static/locales/en.json b/frontend/static/locales/en.json index a2963f72..384534ee 100644 --- a/frontend/static/locales/en.json +++ b/frontend/static/locales/en.json @@ -1209,7 +1209,10 @@ "col_status": "Status", "col_duration": "Duration", "col_scanned": "Scanned", + "col_progress": "Progress", "col_findings": "Findings", + "progress_exact_tooltip": "{{pct}} ({{scanned}} / {{total}})", + "progress_approx_tooltip": "{{pct}} ({{scanned}} / {{total}} — approximate, backend proxy)", "col_error": "Error", "col_kind": "Kind", "col_severity": "Severity", diff --git a/src/infrastructure/scheduler/mod.rs b/src/infrastructure/scheduler/mod.rs index f48f40de..16d9963a 100644 --- a/src/infrastructure/scheduler/mod.rs +++ b/src/infrastructure/scheduler/mod.rs @@ -33,8 +33,9 @@ pub use engine::SchedulerEngine; pub use handler::JobHandler; pub use pg_job_store::{PgJobStore, PgJobStoreProvider}; pub use recoverable::{ - Finding, JobStore, JobStoreProvider, OpenedRun, RecoverableAdapter, RecoverableJobHandler, - RunOutcome, RunStatus, RunSummary, record_or_log, run_or_resume, + Finding, JobStore, JobStoreProvider, OpenedRun, ProgressKind, RecoverableAdapter, + RecoverableJobHandler, RunOutcome, RunProgress, RunStatus, RunSummary, derive_progress, + record_or_log, run_or_resume, }; pub use registry::{JobEntry, JobRegistry, JobSummary, RegisterError}; pub use types::{ErrCause, JobOutcome, JobRunArgs}; diff --git a/src/infrastructure/scheduler/pg_job_store.rs b/src/infrastructure/scheduler/pg_job_store.rs index 6e31ae26..2542545d 100644 --- a/src/infrastructure/scheduler/pg_job_store.rs +++ b/src/infrastructure/scheduler/pg_job_store.rs @@ -16,7 +16,10 @@ use uuid::Uuid; use crate::common::errors::DomainError; -use super::recoverable::{Finding, JobStore, JobStoreProvider, OpenedRun, RunStatus, RunSummary}; +use super::recoverable::{ + Finding, JobStore, JobStoreProvider, OpenedRun, ProgressKind, RunStatus, RunSummary, + derive_progress, +}; // ─── PgJobStore — bound to one run ────────────────────────────────────────── @@ -161,6 +164,40 @@ impl JobStore for PgJobStore { Ok(()) } + async fn seed_progress_params( + &self, + total: u64, + kind: ProgressKind, + ) -> Result<(), DomainError> { + // Stamp `params.total_rows` + `params.progress_kind` in one + // UPDATE. Two `jsonb_set` calls compose left-to-right so both + // keys land atomically. `bigint` cast handles the (theoretical) + // > 2^31 subject-row case. + let total_i64 = total as i64; + sqlx::query( + r#" + UPDATE jobs.recoverable_runs + SET params = jsonb_set( + jsonb_set( + params, + '{total_rows}', + to_jsonb($2::bigint) + ), + '{progress_kind}', + to_jsonb($3::text) + ) + WHERE id = $1 + "#, + ) + .bind(self.run_id) + .bind(total_i64) + .bind(kind.as_str()) + .execute(self.pool.as_ref()) + .await + .map_err(|e| map_sqlx_err("seed_progress_params", e))?; + Ok(()) + } + async fn mark_completed(&self) -> Result<(), DomainError> { sqlx::query( r#" @@ -454,6 +491,23 @@ fn row_to_summary(row: RunSummaryRow) -> Result { let status = RunStatus::parse(&status_str).ok_or_else(|| { DomainError::internal_error("JobStore", format!("unknown status: {status_str}")) })?; + + // Derive progress from stats.scanned_count + params.total_rows + + // params.progress_kind. All three are optional — if the tenant + // didn't seed a total (no count_total override) the block is None + // and the UI hides the bar. `derive_progress` also guards against + // total = 0 (empty-subject run — bar would be meaningless). + let scanned = stats + .get("scanned_count") + .and_then(|v| v.as_u64()) + .unwrap_or(0); + let total = params.get("total_rows").and_then(|v| v.as_u64()); + let kind = params + .get("progress_kind") + .and_then(|v| v.as_str()) + .and_then(ProgressKind::parse); + let progress = derive_progress(scanned, total, kind); + Ok(RunSummary { id, job_name, @@ -465,6 +519,7 @@ fn row_to_summary(row: RunSummaryRow) -> Result { params, cursor_hex: cursor.map(hex::encode), error_message, + progress, }) } diff --git a/src/infrastructure/scheduler/recoverable.rs b/src/infrastructure/scheduler/recoverable.rs index ffb39bd0..b1b3d70f 100644 --- a/src/infrastructure/scheduler/recoverable.rs +++ b/src/infrastructure/scheduler/recoverable.rs @@ -179,6 +179,34 @@ pub trait RecoverableJobHandler: Send + Sync { args: &JobRunArgs, resume_cursor: Option>, ) -> RunOutcome; + + /// **Optional** — override to enable progress estimation on the + /// admin UI. Called ONCE at fresh-run start by [`run_or_resume`]; + /// the returned count is stashed in `params.total_rows` and paired + /// with `stats.scanned_count` at serialisation time to produce a + /// `RunProgress` fraction on `RunSummary`. + /// + /// Return `None` (the default) when the tenant cannot count its + /// subject — an external crawler, a streaming source, or any + /// unbounded workload. The UI then hides the bar and falls back + /// to raw `scanned_count`. + /// + /// **Not called on resume.** A Paused run keeps the `total_rows` + /// stamped at its original start — mid-scan re-counts would make + /// the fraction jump around every time the operator resumed. + async fn count_total(&self) -> Option { + None + } + + /// Confidence level of the count returned by [`count_total`]. + /// Default is [`ProgressKind::Count`] — assume the count is + /// authoritative unless the tenant overrides. Tenants whose + /// `count_total` is a proxy (backend enumeration counting DB + /// blobs instead of backend objects) return + /// [`ProgressKind::Approximate`]. + fn progress_kind(&self) -> ProgressKind { + ProgressKind::Count + } } /// Bound-to-a-run handle. The handler polls status + writes @@ -209,6 +237,15 @@ pub trait JobStore: Send + Sync { /// batches — the run's heartbeat. async fn checkpoint(&self, cursor: Vec, delta_count: u64) -> Result<(), DomainError>; + /// **Engine-only.** Called by [`run_or_resume`] on a Fresh run + /// after the tenant's [`RecoverableJobHandler::count_total`] + /// reports a countable subject. Stamps `params.total_rows` + + /// `params.progress_kind` on the row so subsequent `RunSummary` + /// projections can derive `progress` without asking the tenant + /// again. Handler code MUST NOT call this. + async fn seed_progress_params(&self, total: u64, kind: ProgressKind) + -> Result<(), DomainError>; + /// Persist one finding to `jobs.run_findings` and bump /// `stats.finding_count` on the parent run. Consistency handlers /// call this in place of the transitional @@ -323,6 +360,91 @@ pub trait JobStoreProvider: Send + Sync { ) -> Result, DomainError>; } +/// How a `RunProgress` fraction was derived. Lets the UI communicate +/// confidence to the operator — a `count`-derived 47% is authoritative, +/// an `approximate`-derived 47% is a proxy (e.g. `storage_consistency` +/// using DB blob count as a stand-in for backend object count). +/// +/// A future `cursor` variant will cover UUID-cursor-position-derived +/// fractions (`cursor_position / 2^128`) — useful when `COUNT(*)` on +/// the subject table is too expensive to run at start. Not implemented +/// yet; all shipped tenants override [`RecoverableJobHandler::count_total`]. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)] +#[serde(rename_all = "lowercase")] +pub enum ProgressKind { + /// `scanned_count / total_rows` where `total_rows` came from a + /// definitive `COUNT(*)` on the tenant's subject table. + Count, + /// `scanned_count / total_rows` where `total_rows` is a proxy + /// (e.g. DB blob count for a backend enumeration). The fraction + /// deviating from 1.0 at run end IS informative — it quantifies + /// the drift the check is looking for. + Approximate, +} + +impl ProgressKind { + pub fn as_str(self) -> &'static str { + match self { + ProgressKind::Count => "count", + ProgressKind::Approximate => "approximate", + } + } + + pub fn parse(s: &str) -> Option { + match s { + "count" => Some(ProgressKind::Count), + "approximate" => Some(ProgressKind::Approximate), + _ => None, + } + } +} + +/// Progress estimate on a recoverable run. Populated on `RunSummary` +/// only when the tenant's [`RecoverableJobHandler::count_total`] +/// returned `Some(n)` at run start — a tenant that cannot count its +/// subject (external crawler, streaming source) leaves this `None` and +/// the UI hides the progress bar. +/// +/// `fraction` CAN exceed 1.0 at the end of an +/// [`ProgressKind::Approximate`] run — the deviation IS the finding. +/// The UI should clamp for the bar width but surface the raw fraction +/// in the tooltip. +#[derive(Debug, Clone, Copy, PartialEq, Serialize)] +pub struct RunProgress { + pub fraction: f32, + pub kind: ProgressKind, + /// Included so the UI can render "347 / 1200" alongside the bar + /// without recomputing from `stats.scanned_count`. + pub scanned: u64, + pub total: u64, +} + +/// Build a `RunProgress` from the persisted scanned / total / kind. +/// `None` when `total` is absent (tenant didn't count) OR zero (avoid +/// dividing by zero and rendering a bar for an empty-subject run). +pub fn derive_progress( + scanned: u64, + total: Option, + kind: Option, +) -> Option { + let total = total?; + if total == 0 { + return None; + } + let kind = kind.unwrap_or(ProgressKind::Count); + // We deliberately DON'T clamp — an approximate-kind run can + // legitimately exceed 1.0 (backend has orphans), and that + // deviation is informative signal. The UI clamps for bar width + // but shows raw fraction in the tooltip. + let fraction = scanned as f32 / total as f32; + Some(RunProgress { + fraction, + kind, + scanned, + total, + }) +} + /// Serialisable snapshot of one `jobs.run_findings` row, returned by /// `GET /api/admin/jobs/{name}/runs/{id}/findings`. Consumers key off /// `kind` to know the shape of `detail`. @@ -362,6 +484,12 @@ pub struct RunSummary { pub cursor_hex: Option, #[serde(skip_serializing_if = "Option::is_none")] pub error_message: Option, + /// Present when the tenant reported a countable subject at run + /// start (see [`RecoverableJobHandler::count_total`]). `None` + /// tells the UI "hide the progress bar, show scanned_count as a + /// raw number instead." + #[serde(skip_serializing_if = "Option::is_none")] + pub progress: Option, } /// Result of [`JobStoreProvider::open_or_start`]. @@ -398,7 +526,7 @@ pub async fn run_or_resume( Ok(o) => o, Err(e) => return JobOutcome::err(format!("open_or_start failed: {e}")), }; - let (store, resume_cursor) = match opened { + let (store, resume_cursor, is_fresh) = match opened { OpenedRun::AlreadyActive { run_id, status } => { return JobOutcome::ok_with( 0, @@ -409,11 +537,30 @@ pub async fn run_or_resume( }), ); } - OpenedRun::Fresh { store } => (store, None), - OpenedRun::Resumed { store, cursor } => (store, Some(cursor)), + OpenedRun::Fresh { store } => (store, None, true), + OpenedRun::Resumed { store, cursor } => (store, Some(cursor), false), }; let run_id = store.run_id(); + // Seed progress params on a Fresh run only — a resumed run keeps + // the total_rows stamped when it originally started, otherwise + // the fraction would jump every time the operator resumed. A + // failed count is not fatal; the progress block just stays None + // on the summary (UI falls back to raw scanned_count). + if is_fresh && let Some(total) = job.count_total().await { + let kind = job.progress_kind(); + if let Err(e) = store.seed_progress_params(total, kind).await { + tracing::warn!( + target: "oxicloud::scheduler", + event = "recoverable.seed_progress_failed", + job = job.name(), + run_id = %run_id, + error = %e, + "failed to seed progress params; run continues without a bar" + ); + } + } + // Dispatch. Terminal writes to `jobs.recoverable_runs` happen // here (NOT in the handler) so the row always ends in a state // that matches what the handler returned. @@ -581,6 +728,8 @@ mod tests { scanned_count: u64, error_message: Option, findings: Vec, + progress_total: Option, + progress_kind: Option, } #[async_trait] @@ -619,6 +768,16 @@ mod tests { }); Ok(()) } + async fn seed_progress_params( + &self, + total: u64, + kind: ProgressKind, + ) -> Result<(), DomainError> { + let mut s = self.state.lock().unwrap(); + s.progress_total = Some(total); + s.progress_kind = Some(kind); + Ok(()) + } async fn mark_completed(&self) -> Result<(), DomainError> { self.state.lock().unwrap().status = RunStatus::Completed; Ok(()) @@ -668,6 +827,8 @@ mod tests { scanned_count: 0, error_message: None, findings: Vec::new(), + progress_total: None, + progress_kind: None, }), }); let id = store.run_id; @@ -724,6 +885,8 @@ mod tests { scanned_count: 0, error_message: None, findings: Vec::new(), + progress_total: None, + progress_kind: None, }), }); stores.push(store.clone()); @@ -760,6 +923,11 @@ mod tests { .take(limit as usize) .map(|s| { let state = s.state.lock().unwrap(); + let progress = derive_progress( + state.scanned_count, + state.progress_total, + state.progress_kind, + ); RunSummary { id: s.run_id, job_name: job_name.to_string(), @@ -771,6 +939,7 @@ mod tests { params: serde_json::json!({}), cursor_hex: state.cursor.as_ref().map(hex::encode), error_message: state.error_message.clone(), + progress, } }) .collect(); @@ -782,6 +951,11 @@ mod tests { let now = Utc::now(); Ok(stores.iter().find(|s| s.run_id == run_id).map(|s| { let state = s.state.lock().unwrap(); + let progress = derive_progress( + state.scanned_count, + state.progress_total, + state.progress_kind, + ); RunSummary { id: s.run_id, job_name: "mem".to_string(), @@ -793,6 +967,7 @@ mod tests { params: serde_json::json!({}), cursor_hex: state.cursor.as_ref().map(hex::encode), error_message: state.error_message.clone(), + progress, } })) } diff --git a/src/infrastructure/services/drives_consistency_service.rs b/src/infrastructure/services/drives_consistency_service.rs index 29516841..4781f766 100644 --- a/src/infrastructure/services/drives_consistency_service.rs +++ b/src/infrastructure/services/drives_consistency_service.rs @@ -69,6 +69,28 @@ impl RecoverableJobHandler for DrivesConsistencyCheck { DRIVES_CONSISTENCY_JOB_NAME } + /// Definitive count — one row per drive, table is tiny (dozens per + /// install), COUNT(*) is trivially fast. Enables progress bar on + /// the admin UI. + async fn count_total(&self) -> Option { + let row: Result<(i64,), sqlx::Error> = + sqlx::query_as("SELECT COUNT(*) FROM storage.drives") + .fetch_one(self.pool.as_ref()) + .await; + match row { + Ok((n,)) => Some(n.max(0) as u64), + Err(e) => { + tracing::debug!( + target: "oxicloud::consistency", + event = "drives_consistency.count_total_failed", + error = %e, + "count_total failed — run will not surface a progress bar" + ); + None + } + } + } + async fn run_resumable( &self, store: &dyn JobStore, diff --git a/src/infrastructure/services/files_consistency_service.rs b/src/infrastructure/services/files_consistency_service.rs index cc587b56..541a8a39 100644 --- a/src/infrastructure/services/files_consistency_service.rs +++ b/src/infrastructure/services/files_consistency_service.rs @@ -113,6 +113,30 @@ impl RecoverableJobHandler for FilesConsistencyCheck { FILES_CONSISTENCY_JOB_NAME } + /// Definitive count — one row per file. This is the largest table + /// of the trio (millions on big installs); COUNT(*) is still an + /// index-only scan but can take ~seconds. The tradeoff is worth + /// it — an operator staring at a running files_consistency scan + /// wants a bar, and a seconds-scale one-off at run start is + /// invisible compared to the multi-minute scan that follows. + async fn count_total(&self) -> Option { + let row: Result<(i64,), sqlx::Error> = sqlx::query_as("SELECT COUNT(*) FROM storage.files") + .fetch_one(self.pool.as_ref()) + .await; + match row { + Ok((n,)) => Some(n.max(0) as u64), + Err(e) => { + tracing::debug!( + target: "oxicloud::consistency", + event = "files_consistency.count_total_failed", + error = %e, + "count_total failed — run will not surface a progress bar" + ); + None + } + } + } + async fn run_resumable( &self, store: &dyn JobStore, diff --git a/src/infrastructure/services/folders_consistency_service.rs b/src/infrastructure/services/folders_consistency_service.rs index 8494d73d..c3519197 100644 --- a/src/infrastructure/services/folders_consistency_service.rs +++ b/src/infrastructure/services/folders_consistency_service.rs @@ -117,6 +117,29 @@ impl RecoverableJobHandler for FoldersConsistencyCheck { FOLDERS_CONSISTENCY_JOB_NAME } + /// Definitive count — one row per folder. Larger table than drives + /// but the COUNT(*) is still index-only on PG. On multi-million-row + /// deployments this is ~100ms at run start; acceptable given the + /// progress bar is only rendered when the operator is watching. + async fn count_total(&self) -> Option { + let row: Result<(i64,), sqlx::Error> = + sqlx::query_as("SELECT COUNT(*) FROM storage.folders") + .fetch_one(self.pool.as_ref()) + .await; + match row { + Ok((n,)) => Some(n.max(0) as u64), + Err(e) => { + tracing::debug!( + target: "oxicloud::consistency", + event = "folders_consistency.count_total_failed", + error = %e, + "count_total failed — run will not surface a progress bar" + ); + None + } + } + } + async fn run_resumable( &self, store: &dyn JobStore,