1ea3826660
The admin panel had no repair toggle wired to anything but a hardcoded
name list naming the two refcount tenants, so `thumb_derived_import` and
`thumb_attached_import` could not be run in repair mode from the UI at
all despite supporting it. And nothing in the job list said what any
given job does or whether clicking Run on production writes anything.
Three defaulted methods on `JobHandler` and `RecoverableJobHandler`:
fn description(&self) -> &'static str
fn mutates(&self) -> Mutates // Never | Always | OnRepairOnly
fn repair_description(&self) -> Option<&'static str>
`RecoverableAdapter` forwards them — the registry only holds
`dyn JobHandler`, so a tenant's metadata is invisible otherwise, and
falling back to the defaults would report every recoverable job as
read-only, including the ones that delete files.
Three values rather than a boolean because a job can be read-only by
default and destructive under `?repair=true`; a boolean answers wrongly
for one of its two modes, and `false` on something that unlinks files is
the dangerous direction to be wrong in. `repair_description` returning
`Option` collapses "does it repair" and "what does repair do" into one
method: presence gates the toggle, content is the confirmation text —
which the frontend cannot invent, since correcting a counter and
deleting sidecars are not the same warning.
`OnRepairOnly` with no `repair_description` is rejected at registration:
it claims to mutate only under a flag it does not support.
All 17 registered jobs declare all three. The panel now renders the
description under each name, badges read-only jobs, confirms before a
plain run of a mutating one, and offers the repair variant off the
backend flag instead of the name list.
Descriptions are English in the trait, next to the behaviour: one in
`locales/*.json` rots invisibly the moment a job changes, and a
translator cannot know what `manifests_consistency` reconciles. i18n can
layer on later keyed by job name with these as the fallback.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
561 lines
22 KiB
Rust
561 lines
22 KiB
Rust
//! Storage-format rotation as a recoverable-run tenant (K3 of
|
|
//! `docs/plan/storage-key-rotation.md`).
|
|
//!
|
|
//! Iterates `storage.blobs` for a target entry, decides per blob
|
|
//! whether the on-disk format matches what the entry's head pair
|
|
//! would write, and rewrites in place when it doesn't. Covers four
|
|
//! transitions with a single equality check:
|
|
//!
|
|
//! * Legacy blob (no `OXCPT` magic) → rewrite as v1 with the head
|
|
//! pair's format.
|
|
//! * v1 encrypted, decrypted under a pair-index other than head →
|
|
//! rewrite (key rotation).
|
|
//! * v1 plaintext with head=`aes:K` → rewrite (encrypt-in-place).
|
|
//! * v1 encrypted with head=`none:` → rewrite (decrypt-in-place).
|
|
//!
|
|
//! ### No readonly, no cutover
|
|
//!
|
|
//! `backend_rotate` is per-blob idempotent — repeat rewrites are
|
|
//! byte-safe (content-addressability holds; the wrapper always
|
|
//! produces the head format). Concurrent user writes coexist: they
|
|
//! land as head-format themselves, so when the walk reaches that
|
|
//! hash the classifier reports "already at head format" and the
|
|
//! decision tree collapses to `skip`. No app-wide read-only gate is
|
|
//! ever engaged — a critical improvement over `backend_migration`,
|
|
//! whose target-different-from-source cutover forces one.
|
|
//!
|
|
//! ### Restart survival
|
|
//!
|
|
//! Cursor + per-blob failure findings are persisted after every
|
|
//! batch. On restart, boot flips any abandoned `Running` row to
|
|
//! `Paused`; an admin trigger resumes from the checkpointed cursor.
|
|
//! The last checkpoint window (~100 blobs) re-processes; each of
|
|
//! those blobs is now head-format from the previous run's rewrite,
|
|
//! so the walk short-circuits without re-writing. Effectively free.
|
|
//!
|
|
//! ### Design notes
|
|
//!
|
|
//! * **Cursor** — UTF-8 hex of the last-processed blob hash (64
|
|
//! chars). Same encoding as `backend_migration` and
|
|
//! `blobs_consistency`.
|
|
//! * **Target lookup** — the entry NAME is stashed in `params` at
|
|
//! Fresh-open time and re-read on Resume. The wrapper for that
|
|
//! entry is rebuilt at the top of every run via
|
|
//! `build_entry_backend_typed`; mid-run config changes are
|
|
//! ignored until the next run (mirrors `backend_migration`).
|
|
//! * **Per-blob failures don't fail the run** — each failure records
|
|
//! a `rotation_failed` finding (severity `data_loss` — the bytes
|
|
//! didn't get rewritten) and the walk continues. A run that
|
|
//! completes with zero findings is proof every blob is at head
|
|
//! format.
|
|
//! * **`?deep=true` is unused** — rotation has no slow variant.
|
|
//! Parameter accepted for uniformity with other tenants; ignored.
|
|
|
|
use std::path::PathBuf;
|
|
use std::sync::Arc;
|
|
|
|
use async_trait::async_trait;
|
|
use bytes::Bytes;
|
|
use sqlx::PgPool;
|
|
|
|
use crate::application::ports::blob_storage_ports::BlobStorageBackend;
|
|
use crate::common::config::NamedStorageEntry;
|
|
use crate::common::migration_progress::MigrationProgress;
|
|
use crate::infrastructure::scheduler::{
|
|
JobRegistry, JobRunArgs, JobStore, JobStoreProvider, Mutates, RecoverableJobHandler,
|
|
RunOutcome, RunStatus, record_or_log,
|
|
};
|
|
use crate::infrastructure::services::encrypted_blob_backend::BlobFormat;
|
|
use crate::infrastructure::services::entry_backend::build_entry_backend_typed;
|
|
|
|
pub const BACKEND_ROTATE_JOB_NAME: &str = "backend_rotate";
|
|
|
|
/// The `params` JSONB key under which the run's target entry name is
|
|
/// stashed at Fresh-open time via `JobStore::set_string_param`.
|
|
/// Kept identical to `backend_migration`'s TARGET_NAME_PARAM so
|
|
/// operators grepping run rows see the same convention across both
|
|
/// storage-touching tenants.
|
|
pub const TARGET_NAME_PARAM: &str = "target_name";
|
|
|
|
/// Rows per batch. Matches `backend_migration` / `blobs_consistency`
|
|
/// so the checkpoint + cancel-poll cadence is uniform across tenants.
|
|
const BATCH_SIZE: i64 = 100;
|
|
|
|
pub struct BackendRotateService {
|
|
pool: Arc<PgPool>,
|
|
/// Immutable per-deploy snapshot; used to look up the target
|
|
/// entry by name at run start. Matches `AppConfig.storage_entries`.
|
|
storage_entries: Vec<NamedStorageEntry>,
|
|
/// Ambient `AppConfig.storage_path` used as the `root_dir`
|
|
/// fallback for a Local target entry that doesn't declare its
|
|
/// own `_ROOT_DIR`. Same fallback rule as boot
|
|
/// (`build_entry_backend`).
|
|
storage_path_fallback: PathBuf,
|
|
/// Shared in-memory progress snapshot for the server-status
|
|
/// header middleware. `Some(_)` while a rotation is
|
|
/// running/paused, `None` otherwise. Distinct from
|
|
/// `AppState.migration_progress` so the header can broadcast
|
|
/// migration + rotation states independently.
|
|
rotation_progress: Arc<std::sync::RwLock<Option<MigrationProgress>>>,
|
|
}
|
|
|
|
impl BackendRotateService {
|
|
pub fn new(
|
|
pool: Arc<PgPool>,
|
|
storage_entries: Vec<NamedStorageEntry>,
|
|
storage_path_fallback: PathBuf,
|
|
rotation_progress: Arc<std::sync::RwLock<Option<MigrationProgress>>>,
|
|
) -> Self {
|
|
Self {
|
|
pool,
|
|
storage_entries,
|
|
storage_path_fallback,
|
|
rotation_progress,
|
|
}
|
|
}
|
|
|
|
/// Chainable self-registration — mirrors the `*_consistency`
|
|
/// tenants and `backend_migration`. On-demand only (no periodic
|
|
/// tick).
|
|
pub async fn register_recoverable_job(
|
|
self: Arc<Self>,
|
|
registry: &JobRegistry,
|
|
provider: &Arc<dyn JobStoreProvider>,
|
|
) -> Arc<Self> {
|
|
registry
|
|
.register_recoverable_job(self.clone(), provider.clone(), None)
|
|
.await;
|
|
self
|
|
}
|
|
}
|
|
|
|
#[async_trait]
|
|
impl RecoverableJobHandler for BackendRotateService {
|
|
fn name(&self) -> &str {
|
|
BACKEND_ROTATE_JOB_NAME
|
|
}
|
|
|
|
fn description(&self) -> &'static str {
|
|
"Brings every blob's on-disk format in line with the storage \
|
|
entry's current head key: encrypts plaintext, re-encrypts under a \
|
|
rotated key, decrypts when the head is 'none', and upgrades \
|
|
legacy blobs to v1. Blobs already in the right format are skipped, \
|
|
so re-running after a key change is cheap."
|
|
}
|
|
|
|
/// Rewrites blobs **in place**. Unlike a migration this has no additive
|
|
/// fallback — the previous ciphertext is gone once a blob is rewritten.
|
|
fn mutates(&self) -> Mutates {
|
|
Mutates::Always
|
|
}
|
|
|
|
/// Definitive count — one row per blob. Same query as
|
|
/// `backend_migration::count_total`; the two walk the same rows.
|
|
async fn count_total(&self) -> Option<u64> {
|
|
let row: Result<(i64,), sqlx::Error> = sqlx::query_as("SELECT COUNT(*) FROM storage.blobs")
|
|
.fetch_one(self.pool.as_ref())
|
|
.await;
|
|
match row {
|
|
Ok((n,)) => Some(n.max(0) as u64),
|
|
Err(e) => {
|
|
tracing::debug!(
|
|
target: "oxicloud::rotate",
|
|
event = "backend_rotate.count_total_failed",
|
|
error = %e,
|
|
"count_total failed — run will not surface a progress bar"
|
|
);
|
|
None
|
|
}
|
|
}
|
|
}
|
|
|
|
async fn run_resumable(
|
|
&self,
|
|
store: &dyn JobStore,
|
|
args: &JobRunArgs,
|
|
resume_cursor: Option<Vec<u8>>,
|
|
) -> RunOutcome {
|
|
// Resolve target entry name — same shape as `backend_migration`.
|
|
let is_fresh = resume_cursor.is_none();
|
|
let target_name = if is_fresh {
|
|
let Some(name) = args.storage.clone() else {
|
|
return RunOutcome::Failed {
|
|
message: "backend_rotate requires `target_name` on a fresh run — trigger via \
|
|
POST /api/admin/storage/entries/{name}/rotate"
|
|
.to_string(),
|
|
};
|
|
};
|
|
if let Err(e) = store.set_string_param(TARGET_NAME_PARAM, &name).await {
|
|
return RunOutcome::Failed {
|
|
message: format!("failed to persist target_name to params: {e}"),
|
|
};
|
|
}
|
|
name
|
|
} else {
|
|
match store.get_string_param(TARGET_NAME_PARAM).await {
|
|
Ok(Some(name)) => name,
|
|
Ok(None) => {
|
|
return RunOutcome::Failed {
|
|
message: format!(
|
|
"resumed run has no {TARGET_NAME_PARAM} in params — cancel + trigger \
|
|
fresh."
|
|
),
|
|
};
|
|
}
|
|
Err(e) => {
|
|
return RunOutcome::Failed {
|
|
message: format!("read {TARGET_NAME_PARAM} from params: {e}"),
|
|
};
|
|
}
|
|
}
|
|
};
|
|
|
|
// Look up the target entry.
|
|
let target_entry = match self.storage_entries.iter().find(|e| e.name == target_name) {
|
|
Some(e) => e,
|
|
None => {
|
|
let available = if self.storage_entries.is_empty() {
|
|
"(none)".to_string()
|
|
} else {
|
|
self.storage_entries
|
|
.iter()
|
|
.map(|e| e.name.as_str())
|
|
.collect::<Vec<_>>()
|
|
.join(", ")
|
|
};
|
|
return RunOutcome::Failed {
|
|
message: format!(
|
|
"target entry `{target_name}` not declared in `OXICLOUD_STORAGE_ENTRIES` — \
|
|
available: {available}."
|
|
),
|
|
};
|
|
}
|
|
};
|
|
|
|
// Build the wrapper for this entry — typed so we can call
|
|
// `read_and_classify` + `head_format` directly.
|
|
let wrapper = build_entry_backend_typed(target_entry, &self.storage_path_fallback);
|
|
if let Err(e) = wrapper.initialize().await {
|
|
return RunOutcome::Failed {
|
|
message: format!("target entry `{target_name}` failed to initialize: {e}"),
|
|
};
|
|
}
|
|
let head_format = wrapper.head_format();
|
|
|
|
tracing::info!(
|
|
target: "audit",
|
|
event = "backend_rotate.run_started",
|
|
run_id = %store.run_id(),
|
|
target_name = %target_name,
|
|
// `%` (Display) → SSH-style `encrypted-v1 key_fp=83:96:...`
|
|
// instead of the raw `[131, 150, 255, ...]` byte-array
|
|
// shape Debug produces. Matches how `xxd` renders the
|
|
// header bytes on disk.
|
|
head_format = %head_format,
|
|
resuming = !is_fresh,
|
|
"backend_rotate started on `{target_name}` (head_format = {head_format})"
|
|
);
|
|
|
|
// Seed the progress snapshot. Total = count_total's estimate;
|
|
// if that failed we still surface the header without a
|
|
// denominator so the banner shows "rotation in progress" at
|
|
// minimum.
|
|
let total = self.count_total().await.unwrap_or(0);
|
|
{
|
|
let mut guard = self
|
|
.rotation_progress
|
|
.write()
|
|
.unwrap_or_else(std::sync::PoisonError::into_inner);
|
|
*guard = Some(MigrationProgress::new(target_name.clone(), total));
|
|
}
|
|
|
|
let mut cursor: Option<String> = match resume_cursor {
|
|
None => None,
|
|
Some(bytes) if bytes.is_empty() => None,
|
|
Some(bytes) => match String::from_utf8(bytes) {
|
|
Ok(s) => Some(s),
|
|
Err(e) => {
|
|
self.clear_progress();
|
|
return RunOutcome::Failed {
|
|
message: format!("invalid cursor: not valid UTF-8: {e}"),
|
|
};
|
|
}
|
|
},
|
|
};
|
|
|
|
let mut rewritten_count = 0u64;
|
|
let mut skipped_count = 0u64;
|
|
let mut failed_count = 0u64;
|
|
|
|
loop {
|
|
// Cooperative cancel poll between batches.
|
|
match store.status().await {
|
|
Ok(RunStatus::CancelRequested) => {
|
|
self.clear_progress();
|
|
tracing::info!(
|
|
target: "oxicloud::rotate",
|
|
event = "backend_rotate.cancelled",
|
|
run_id = %store.run_id(),
|
|
rewritten = rewritten_count,
|
|
skipped = skipped_count,
|
|
failed = failed_count,
|
|
"backend_rotate cancelled cooperatively, pausing"
|
|
);
|
|
return RunOutcome::Paused {
|
|
cursor: cursor
|
|
.as_ref()
|
|
.map(|s| s.as_bytes().to_vec())
|
|
.unwrap_or_default(),
|
|
};
|
|
}
|
|
Ok(_) => {}
|
|
Err(e) => {
|
|
self.clear_progress();
|
|
return RunOutcome::Failed {
|
|
message: format!("status poll: {e}"),
|
|
};
|
|
}
|
|
}
|
|
|
|
// Fetch the next batch. Same keyset pagination shape as
|
|
// `backend_migration` — `hash > $1` on the PK, index-only.
|
|
let rows: Vec<(String,)> = match sqlx::query_as(
|
|
r#"
|
|
SELECT hash
|
|
FROM storage.blobs
|
|
WHERE ($1::text IS NULL OR hash > $1)
|
|
ORDER BY hash
|
|
LIMIT $2
|
|
"#,
|
|
)
|
|
.bind(cursor.as_deref())
|
|
.bind(BATCH_SIZE)
|
|
.fetch_all(self.pool.as_ref())
|
|
.await
|
|
{
|
|
Ok(r) => r,
|
|
Err(e) => {
|
|
self.clear_progress();
|
|
return RunOutcome::Failed {
|
|
message: format!("batch fetch: {e}"),
|
|
};
|
|
}
|
|
};
|
|
|
|
if rows.is_empty() {
|
|
return self
|
|
.finish_completed(
|
|
store,
|
|
&target_name,
|
|
head_format,
|
|
rewritten_count,
|
|
skipped_count,
|
|
failed_count,
|
|
)
|
|
.await;
|
|
}
|
|
|
|
for (hash,) in &rows {
|
|
// Read + classify in one round-trip. Failure here is
|
|
// a real read failure (e.g. blob missing on disk),
|
|
// recorded as a finding.
|
|
let (plaintext, current_format) = match wrapper.read_and_classify(hash).await {
|
|
Ok(pair) => pair,
|
|
Err(e) => {
|
|
failed_count += 1;
|
|
tracing::warn!(
|
|
target: "oxicloud::rotate",
|
|
event = "backend_rotate.read_failed",
|
|
run_id = %store.run_id(),
|
|
hash = %hash,
|
|
error = %e,
|
|
"failed to read blob for classification; recording finding"
|
|
);
|
|
record_or_log(
|
|
store,
|
|
BACKEND_ROTATE_JOB_NAME,
|
|
"rotation_failed",
|
|
"data_loss",
|
|
None,
|
|
serde_json::json!({
|
|
"hash": hash,
|
|
"phase": "read",
|
|
"error": e.to_string(),
|
|
}),
|
|
)
|
|
.await;
|
|
continue;
|
|
}
|
|
};
|
|
|
|
// The whole decision tree collapses to one equality
|
|
// check thanks to `BlobFormat`'s `PartialEq`. Six
|
|
// cases in the plan → one branch here.
|
|
if current_format == head_format {
|
|
skipped_count += 1;
|
|
continue;
|
|
}
|
|
|
|
// Rewrite via the atomic-replace write path.
|
|
//
|
|
// **NOT** `put_blob_from_bytes`: that variant is
|
|
// idempotent-skip (`O_CREAT|O_EXCL` on
|
|
// `LocalBlobBackend`) — correct for uploads (same
|
|
// plaintext ↔ any ciphertext at hash decrypts back)
|
|
// but a silent no-op for us. Rotate NEEDS the on-disk
|
|
// bytes to change (legacy → v1 header, old key → new
|
|
// key, plaintext ↔ encrypted). Ed hit this on
|
|
// 2026-08-02: rotation reported success in 9s but
|
|
// every blob on disk still had the legacy shape.
|
|
// `put_blob_from_bytes_replace` writes to a tempfile
|
|
// + atomic `rename(2)`s over the existing object key.
|
|
if let Err(e) = wrapper
|
|
.put_blob_from_bytes_replace(hash, Bytes::from(plaintext.to_vec()))
|
|
.await
|
|
{
|
|
failed_count += 1;
|
|
tracing::warn!(
|
|
target: "oxicloud::rotate",
|
|
event = "backend_rotate.write_failed",
|
|
run_id = %store.run_id(),
|
|
hash = %hash,
|
|
error = %e,
|
|
"failed to rewrite blob; recording finding"
|
|
);
|
|
record_or_log(
|
|
store,
|
|
BACKEND_ROTATE_JOB_NAME,
|
|
"rotation_failed",
|
|
"data_loss",
|
|
None,
|
|
serde_json::json!({
|
|
"hash": hash,
|
|
"phase": "write",
|
|
"from": format!("{current_format}"),
|
|
"to": format!("{head_format}"),
|
|
"error": e.to_string(),
|
|
}),
|
|
)
|
|
.await;
|
|
continue;
|
|
}
|
|
rewritten_count += 1;
|
|
}
|
|
|
|
// Advance cursor + checkpoint. `delta_count` = work
|
|
// attempted this batch, so the progress bar advances even
|
|
// when a batch is dominated by skips (steady-state
|
|
// re-run) or failures.
|
|
let last_hash = rows.last().map(|(h,)| h.clone()).expect("non-empty rows");
|
|
cursor = Some(last_hash.clone());
|
|
let batch_len = rows.len() as u64;
|
|
if let Err(e) = store.checkpoint(last_hash.into_bytes(), batch_len).await {
|
|
self.clear_progress();
|
|
return RunOutcome::Failed {
|
|
message: format!("checkpoint: {e}"),
|
|
};
|
|
}
|
|
{
|
|
let mut guard = self
|
|
.rotation_progress
|
|
.write()
|
|
.unwrap_or_else(std::sync::PoisonError::into_inner);
|
|
if let Some(progress) = guard.as_mut() {
|
|
progress.bump(batch_len);
|
|
}
|
|
}
|
|
|
|
if (rows.len() as i64) < BATCH_SIZE {
|
|
return self
|
|
.finish_completed(
|
|
store,
|
|
&target_name,
|
|
head_format,
|
|
rewritten_count,
|
|
skipped_count,
|
|
failed_count,
|
|
)
|
|
.await;
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
impl BackendRotateService {
|
|
/// Terminal successful path — clear the header snapshot and log a
|
|
/// final audit line. Unlike `backend_migration::finish_completed`
|
|
/// there's no cutover / hot-swap step: rotation writes in place
|
|
/// on the entry that's already there.
|
|
///
|
|
/// `head_format` — the target format at run completion. Persisted
|
|
/// into the run row's `stats` as `head_format` (Display) +
|
|
/// `head_key_fp` (raw hex) so operators have a durable record of
|
|
/// "at time T, all blobs on entry E were normalised to fingerprint
|
|
/// F". Combined with `failed = 0`, that's the signal to remove
|
|
/// obsolete keys from `.env` — any key NOT matching `head_key_fp`
|
|
/// no longer decrypts any live blob and can be safely dropped.
|
|
async fn finish_completed(
|
|
&self,
|
|
store: &dyn JobStore,
|
|
target_name: &str,
|
|
head_format: BlobFormat,
|
|
rewritten: u64,
|
|
skipped: u64,
|
|
failed: u64,
|
|
) -> RunOutcome {
|
|
self.clear_progress();
|
|
|
|
// Render two fingerprint shapes:
|
|
// * `head_format` — Display impl, e.g.
|
|
// `encrypted-v1 key_fp=15:f3:8f:80:2c:ae:2c:50` — human
|
|
// friendly for audit logs + admin UI.
|
|
// * `head_key_fp` — bare 16-hex string, matches what an
|
|
// operator gets from `openssl dgst -sha256 <keyfile> | head -c 16`
|
|
// so post-hoc verification against the raw key material
|
|
// is trivial.
|
|
let head_display = format!("{head_format}");
|
|
let head_key_fp_hex = match head_format {
|
|
BlobFormat::EncryptedV1 { key_fp } => hex::encode(key_fp),
|
|
BlobFormat::PlaintextV1 => String::new(), // all-zero, uninformative
|
|
BlobFormat::Legacy => String::new(), // never emitted at head
|
|
};
|
|
|
|
tracing::info!(
|
|
target: "audit",
|
|
event = "backend_rotate.run_completed",
|
|
run_id = %store.run_id(),
|
|
target_name = %target_name,
|
|
rewritten = rewritten,
|
|
skipped = skipped,
|
|
failed = failed,
|
|
head_format = %head_display,
|
|
"backend_rotate completed on `{target_name}` — {rewritten} rewritten, {skipped} skipped, {failed} failed; head = {head_display}"
|
|
);
|
|
|
|
// Surface the per-run summary counters as extras merged into
|
|
// the run row's `stats` JSONB. Frontend renders whatever keys
|
|
// are present, so no wire-format bumping is needed — the
|
|
// admin UI's run drawer just picks these up alongside the
|
|
// engine-owned `finding_count` + `scanned_count`.
|
|
//
|
|
// `head_key_fp` empty string when head is not an encrypted
|
|
// pair (plaintext-v1) — frontend can render "all in clear"
|
|
// vs "all under fp <X>" based on that discriminator.
|
|
RunOutcome::completed_with(serde_json::json!({
|
|
"rewritten": rewritten,
|
|
"skipped": skipped,
|
|
"failed": failed,
|
|
"head_format": head_display,
|
|
"head_key_fp": head_key_fp_hex,
|
|
}))
|
|
}
|
|
|
|
fn clear_progress(&self) {
|
|
let mut guard = self
|
|
.rotation_progress
|
|
.write()
|
|
.unwrap_or_else(std::sync::PoisonError::into_inner);
|
|
*guard = None;
|
|
}
|
|
}
|