refactor(backend): normalize naming convention to backend rather storage
no ambiguity with the backend rather storage
This commit is contained in:
@@ -259,7 +259,7 @@ pub struct StorageEntrySummaryDto {
|
||||
/// - See which pairs are configured + their SSH-style
|
||||
/// fingerprints without inspecting `.env`.
|
||||
/// - Cross-reference the head pair against the `head_key_fp`
|
||||
/// from the last `storage_rotate` completion — if they
|
||||
/// from the last `backend_rotate` completion — if they
|
||||
/// match AND `failed = 0`, every on-disk blob is under the
|
||||
/// head, and non-head pairs are safe to remove.
|
||||
#[serde(default)]
|
||||
|
||||
@@ -136,12 +136,12 @@ pub trait BlobStorageBackend: Send + Sync + 'static {
|
||||
/// that need the on-disk BYTES to change even when the CONTENT
|
||||
/// hash doesn't:
|
||||
///
|
||||
/// * `storage_rotate` — rewrites every blob under the head pair's
|
||||
/// * `backend_rotate` — rewrites every blob under the head pair's
|
||||
/// format (legacy → v1 header, old key → new key, plaintext ↔
|
||||
/// encrypted). If the target's `put_blob_from_bytes` silently
|
||||
/// skipped, rotation would report success while leaving the old
|
||||
/// format on disk.
|
||||
/// * `storage_migration` — same story when a target already has a
|
||||
/// * `backend_migration` — same story when a target already has a
|
||||
/// blob at that hash from an earlier state (Ed hit this on
|
||||
/// 2026-08-01 in the S3 → local migration test).
|
||||
///
|
||||
|
||||
@@ -510,7 +510,7 @@ impl KeyPair {
|
||||
/// SSH-style colon-hex fingerprint of the key material — 8 bytes
|
||||
/// of SHA-256 truncation rendered as `xx:yy:zz:...`. Same
|
||||
/// truncation as the v1 header's `<key_fp>` field and the
|
||||
/// `head_key_fp` reported by `storage_rotate` on completion, so
|
||||
/// `head_key_fp` reported by `backend_rotate` on completion, so
|
||||
/// operators can cross-reference the boot log against a rotate
|
||||
/// report or the CLI's `oxicloud --fingerprint <base64key>`
|
||||
/// output without any format conversion.
|
||||
@@ -725,7 +725,7 @@ pub fn parse_encryption_pair_list(entry_name: &str, raw: &str) -> Result<Vec<Key
|
||||
///
|
||||
/// Used by the `oxicloud --fingerprint <base64>` CLI subcommand so
|
||||
/// admins can identify which key in their `.env` corresponds to the
|
||||
/// `head_key_fp` a `storage_rotate` run reported on completion —
|
||||
/// `head_key_fp` a `backend_rotate` run reported on completion —
|
||||
/// see `docs/plan/storage-key-rotation.md`.
|
||||
///
|
||||
/// Errors on non-base64 input or on decoded length ≠ 32 bytes (the
|
||||
@@ -911,7 +911,7 @@ impl NamedStorageEntry {
|
||||
|
||||
/// The whole pair list, or an empty slice when the entry is
|
||||
/// unencrypted. Used by K2's read path to walk pairs and by
|
||||
/// `storage_rotate` to enumerate legacy pairs. Callers that only
|
||||
/// `backend_rotate` to enumerate legacy pairs. Callers that only
|
||||
/// need the write pair should prefer [`Self::head_key_material`].
|
||||
pub fn encryption_pairs(&self) -> &[KeyPair] {
|
||||
self.encryption.as_deref().unwrap_or(&[])
|
||||
@@ -3931,7 +3931,7 @@ mod tests {
|
||||
// K3.7: display fp switched from 12-char raw hex to
|
||||
// SSH-style 8-byte colon-hex (16 hex + 7 colons = 23 chars)
|
||||
// so operators can cross-reference against the v1 header's
|
||||
// `<key_fp>` field + `storage_rotate`'s `head_key_fp`
|
||||
// `<key_fp>` field + `backend_rotate`'s `head_key_fp`
|
||||
// output + the `oxicloud --fingerprint` CLI.
|
||||
let pairs =
|
||||
parse_encryption_pair_list("t", &format!("aes-256-gcm:{K1_B64},none:")).unwrap();
|
||||
|
||||
+9
-9
@@ -2245,7 +2245,7 @@ impl AppServiceFactory {
|
||||
dyn crate::infrastructure::scheduler::JobStoreProvider,
|
||||
> = app_state.core.job_store_provider.clone();
|
||||
let _ = Arc::new(
|
||||
crate::infrastructure::services::storage_migration_service::StorageMigrationService::new(
|
||||
crate::infrastructure::services::backend_migration_service::BackendMigrationService::new(
|
||||
app_state
|
||||
.maintenance_pool
|
||||
.clone()
|
||||
@@ -2262,13 +2262,13 @@ impl AppServiceFactory {
|
||||
.register_recoverable_job(&app_state.core.job_registry, &job_store_provider_dyn)
|
||||
.await;
|
||||
|
||||
// K3: `storage_rotate` recoverable-job tenant. Same
|
||||
// pattern as `storage_migration` but without the
|
||||
// K3: `backend_rotate` recoverable-job tenant. Same
|
||||
// pattern as `backend_migration` but without the
|
||||
// cutover/readonly plumbing — rotation writes in place on
|
||||
// whichever entry the trigger endpoint names. Target name
|
||||
// comes from `params.target_name` per run.
|
||||
let _ = Arc::new(
|
||||
crate::infrastructure::services::storage_rotate_service::StorageRotateService::new(
|
||||
crate::infrastructure::services::backend_rotate_service::BackendRotateService::new(
|
||||
app_state
|
||||
.maintenance_pool
|
||||
.clone()
|
||||
@@ -2477,7 +2477,7 @@ impl AppServiceFactory {
|
||||
// Migration-readonly boot-clear rule. See
|
||||
// `docs/plan/storage-multi-entry.md` §"Read-only mode".
|
||||
//
|
||||
// If the flag was set true at boot AND no storage_migration
|
||||
// If the flag was set true at boot AND no backend_migration
|
||||
// run is currently non-terminal AND active_backend_name
|
||||
// matches the entry the app actually booted onto — that means
|
||||
// the cutover completed on a prior boot (the run reached
|
||||
@@ -2495,11 +2495,11 @@ impl AppServiceFactory {
|
||||
.migration_readonly
|
||||
.load(std::sync::atomic::Ordering::Relaxed)
|
||||
{
|
||||
use crate::infrastructure::services::storage_migration_service::STORAGE_MIGRATION_JOB_NAME;
|
||||
use crate::infrastructure::services::backend_migration_service::BACKEND_MIGRATION_JOB_NAME;
|
||||
let has_in_flight = match app_state
|
||||
.core
|
||||
.job_store_provider
|
||||
.list_runs(STORAGE_MIGRATION_JOB_NAME, 5)
|
||||
.list_runs(BACKEND_MIGRATION_JOB_NAME, 5)
|
||||
.await
|
||||
{
|
||||
Ok(runs) => runs.iter().any(|r| {
|
||||
@@ -2515,7 +2515,7 @@ impl AppServiceFactory {
|
||||
target: "oxicloud::scheduler",
|
||||
event = "storage.migration_readonly.clear_check_failed",
|
||||
error = %e,
|
||||
"failed to list storage_migration runs during readonly-clear check; \
|
||||
"failed to list backend_migration runs during readonly-clear check; \
|
||||
leaving migration_readonly flag as-is"
|
||||
);
|
||||
// Play it safe: assume in-flight to avoid clearing prematurely.
|
||||
@@ -2816,7 +2816,7 @@ pub struct AppState {
|
||||
/// polling. See `MigrationProgress` for the field shape.
|
||||
pub migration_progress: Arc<std::sync::RwLock<Option<crate::common::migration_progress::MigrationProgress>>>,
|
||||
/// Live progress snapshot for the storage-rotate handler
|
||||
/// (`storage_rotate` — K3 of the storage-key-rotation plan).
|
||||
/// (`backend_rotate` — K3 of the storage-key-rotation plan).
|
||||
/// `Some(_)` while a rotation is running; `None` otherwise.
|
||||
/// Held separately from `migration_progress` so the
|
||||
/// server-status header can broadcast the two states
|
||||
|
||||
@@ -124,7 +124,7 @@ pub enum RunOutcome {
|
||||
/// `extra_stats` is merged into the run row's `stats` JSONB
|
||||
/// alongside the engine-owned `scanned_count` + `finding_count`
|
||||
/// / `severity_counts`. Handlers use it to surface per-run
|
||||
/// summary counters (e.g. `storage_rotate` reports
|
||||
/// summary counters (e.g. `backend_rotate` reports
|
||||
/// `{"rewritten": N, "skipped": M, "failed": K}`) — the outcome
|
||||
/// message in `JobOutcome.extra` and every downstream reader
|
||||
/// of `RunSummary.stats` see the merged fields.
|
||||
@@ -304,7 +304,7 @@ pub trait JobStore: Send + Sync {
|
||||
|
||||
/// Set an arbitrary string field on `params` (JSONB). Used by
|
||||
/// handlers on a Fresh run to persist per-run configuration that
|
||||
/// must survive a mid-run restart — e.g. `storage_migration`
|
||||
/// must survive a mid-run restart — e.g. `backend_migration`
|
||||
/// stamping `params.target_name` at run start so a resume can
|
||||
/// pick up the same target without the admin re-specifying it.
|
||||
///
|
||||
|
||||
@@ -37,7 +37,7 @@ use serde::{Deserialize, Serialize};
|
||||
///
|
||||
/// Semantics of `storage`, per job (added for the multi-entry storage
|
||||
/// design — see `docs/plan/storage-multi-entry.md`):
|
||||
/// - `storage_migration` — the NAME of the target storage entry to
|
||||
/// - `backend_migration` — the NAME of the target storage entry to
|
||||
/// copy blobs INTO. Required on a Fresh run (handler refuses
|
||||
/// without it); ignored on a Resumed run (target read from the
|
||||
/// persisted `params.target_name`).
|
||||
|
||||
@@ -173,8 +173,8 @@ impl BlobStorageBackend for AzureBlobBackend {
|
||||
})
|
||||
}
|
||||
|
||||
/// Atomic overwrite path used by `storage_rotate` and
|
||||
/// `storage_migration` when re-writing an already-present blob
|
||||
/// Atomic overwrite path used by `backend_rotate` and
|
||||
/// `backend_migration` when re-writing an already-present blob
|
||||
/// under a new head key/format. Trait default delegates to
|
||||
/// `put_blob_from_bytes` which `get_properties`-probes and
|
||||
/// silently skips — exactly wrong for the rotate/migrate use
|
||||
|
||||
+31
-31
@@ -67,7 +67,7 @@ use crate::infrastructure::services::entry_backend::{
|
||||
build_entry_backend_typed, persist_active_backend_name, persist_migration_readonly,
|
||||
};
|
||||
|
||||
pub const STORAGE_MIGRATION_JOB_NAME: &str = "storage_migration";
|
||||
pub const BACKEND_MIGRATION_JOB_NAME: &str = "backend_migration";
|
||||
|
||||
/// The `params` JSONB key under which the run's target entry name is
|
||||
/// stashed at Fresh-open time via `JobStore::set_string_param`.
|
||||
@@ -93,7 +93,7 @@ pub const SOURCE_NAME_PARAM: &str = "source_name";
|
||||
/// every 100 rows too. Match `blobs_consistency` for consistency.
|
||||
const BATCH_SIZE: i64 = 100;
|
||||
|
||||
pub struct StorageMigrationService {
|
||||
pub struct BackendMigrationService {
|
||||
pool: Arc<PgPool>,
|
||||
/// Backend the running app is bound to at handler-construction
|
||||
/// time. Refers to the hot-swap wrapper when multi-entry is
|
||||
@@ -140,7 +140,7 @@ pub struct StorageMigrationService {
|
||||
Arc<std::sync::RwLock<Option<crate::common::migration_progress::MigrationProgress>>>,
|
||||
}
|
||||
|
||||
impl StorageMigrationService {
|
||||
impl BackendMigrationService {
|
||||
#[allow(clippy::too_many_arguments)]
|
||||
pub fn new(
|
||||
pool: Arc<PgPool>,
|
||||
@@ -183,9 +183,9 @@ impl StorageMigrationService {
|
||||
}
|
||||
|
||||
#[async_trait]
|
||||
impl RecoverableJobHandler for StorageMigrationService {
|
||||
impl RecoverableJobHandler for BackendMigrationService {
|
||||
fn name(&self) -> &str {
|
||||
STORAGE_MIGRATION_JOB_NAME
|
||||
BACKEND_MIGRATION_JOB_NAME
|
||||
}
|
||||
|
||||
/// Definitive count — one row per blob. `SELECT COUNT(*) FROM
|
||||
@@ -200,7 +200,7 @@ impl RecoverableJobHandler for StorageMigrationService {
|
||||
Err(e) => {
|
||||
tracing::debug!(
|
||||
target: "oxicloud::migration",
|
||||
event = "storage_migration.count_total_failed",
|
||||
event = "backend_migration.count_total_failed",
|
||||
error = %e,
|
||||
"count_total failed — run will not surface a progress bar"
|
||||
);
|
||||
@@ -234,7 +234,7 @@ impl RecoverableJobHandler for StorageMigrationService {
|
||||
let Some(name) = args.storage.clone() else {
|
||||
return RunOutcome::Failed {
|
||||
message:
|
||||
"storage_migration requires `target_name` on a fresh run — trigger via \
|
||||
"backend_migration requires `target_name` on a fresh run — trigger via \
|
||||
POST /api/admin/storage/migration/start with `{\"target_name\": \"<entry>\"}`."
|
||||
.to_string(),
|
||||
};
|
||||
@@ -305,7 +305,7 @@ impl RecoverableJobHandler for StorageMigrationService {
|
||||
.clone();
|
||||
tracing::warn!(
|
||||
target: "oxicloud::migration",
|
||||
event = "storage_migration.legacy_paused_row_source_defaulted",
|
||||
event = "backend_migration.legacy_paused_row_source_defaulted",
|
||||
run_id = %store.run_id(),
|
||||
fallback_source = %fallback,
|
||||
"resumed run has no source_name in params (pre-K3.8 row) — defaulting \
|
||||
@@ -330,11 +330,11 @@ impl RecoverableJobHandler for StorageMigrationService {
|
||||
if target_name == active_backend_name {
|
||||
tracing::warn!(
|
||||
target: "audit",
|
||||
event = "storage_migration.refused_noop",
|
||||
event = "backend_migration.refused_noop",
|
||||
run_id = %store.run_id(),
|
||||
target_name = %target_name,
|
||||
active = %active_backend_name,
|
||||
"storage_migration refused: target equals the currently-active entry"
|
||||
"backend_migration refused: target equals the currently-active entry"
|
||||
);
|
||||
return RunOutcome::Failed {
|
||||
message: format!(
|
||||
@@ -397,12 +397,12 @@ impl RecoverableJobHandler for StorageMigrationService {
|
||||
};
|
||||
tracing::warn!(
|
||||
target: "audit",
|
||||
event = "storage_migration.refused_same_physical_storage",
|
||||
event = "backend_migration.refused_same_physical_storage",
|
||||
run_id = %store.run_id(),
|
||||
target_name = %target_name,
|
||||
source_name = %active_backend_name,
|
||||
encryption_differs = key_differs,
|
||||
"storage_migration refused: named target differs from source but physical storage matches"
|
||||
"backend_migration refused: named target differs from source but physical storage matches"
|
||||
);
|
||||
return RunOutcome::Failed {
|
||||
message: format!(
|
||||
@@ -484,14 +484,14 @@ impl RecoverableJobHandler for StorageMigrationService {
|
||||
let target_kind = target.backend_type();
|
||||
tracing::info!(
|
||||
target: "audit",
|
||||
event = "storage_migration.run_started",
|
||||
event = "backend_migration.run_started",
|
||||
run_id = %store.run_id(),
|
||||
source_name = %active_backend_name,
|
||||
target_name = %target_name,
|
||||
source_kind = source_kind,
|
||||
target_kind = target_kind,
|
||||
resuming = !is_fresh,
|
||||
"storage_migration starting {active_backend_name} ({source_kind}) → {target_name} ({target_kind})"
|
||||
"backend_migration starting {active_backend_name} ({source_kind}) → {target_name} ({target_kind})"
|
||||
);
|
||||
|
||||
// Cursor = the last-visited blob hash, UTF-8-encoded. On resume
|
||||
@@ -527,13 +527,13 @@ impl RecoverableJobHandler for StorageMigrationService {
|
||||
Ok(RunStatus::CancelRequested) => {
|
||||
tracing::info!(
|
||||
target: "oxicloud::migration",
|
||||
event = "storage_migration.cancelled",
|
||||
event = "backend_migration.cancelled",
|
||||
run_id = %store.run_id(),
|
||||
copied = copied_count,
|
||||
skipped = skipped_count,
|
||||
failed = failed_count,
|
||||
source_missing = source_missing_count,
|
||||
"storage_migration cancelled cooperatively, pausing"
|
||||
"backend_migration cancelled cooperatively, pausing"
|
||||
);
|
||||
return RunOutcome::Paused {
|
||||
cursor: cursor
|
||||
@@ -604,7 +604,7 @@ impl RecoverableJobHandler for StorageMigrationService {
|
||||
source_missing_count += 1;
|
||||
tracing::warn!(
|
||||
target: "oxicloud::migration",
|
||||
event = "storage_migration.source_missing",
|
||||
event = "backend_migration.source_missing",
|
||||
run_id = %store.run_id(),
|
||||
hash = %hash,
|
||||
source = source_kind,
|
||||
@@ -612,7 +612,7 @@ impl RecoverableJobHandler for StorageMigrationService {
|
||||
);
|
||||
record_or_log(
|
||||
store,
|
||||
STORAGE_MIGRATION_JOB_NAME,
|
||||
BACKEND_MIGRATION_JOB_NAME,
|
||||
"source_missing",
|
||||
"data_loss",
|
||||
None,
|
||||
@@ -635,7 +635,7 @@ impl RecoverableJobHandler for StorageMigrationService {
|
||||
// it.
|
||||
tracing::warn!(
|
||||
target: "oxicloud::migration",
|
||||
event = "storage_migration.source_probe_error",
|
||||
event = "backend_migration.source_probe_error",
|
||||
run_id = %store.run_id(),
|
||||
hash = %hash,
|
||||
error = %e,
|
||||
@@ -670,7 +670,7 @@ impl RecoverableJobHandler for StorageMigrationService {
|
||||
skipped_count += 1;
|
||||
tracing::debug!(
|
||||
target: "oxicloud::migration",
|
||||
event = "storage_migration.blob_skipped_head_match",
|
||||
event = "backend_migration.blob_skipped_head_match",
|
||||
run_id = %store.run_id(),
|
||||
hash = %hash,
|
||||
head_format = %target.head_format(),
|
||||
@@ -687,7 +687,7 @@ impl RecoverableJobHandler for StorageMigrationService {
|
||||
failed_count += 1;
|
||||
tracing::warn!(
|
||||
target: "oxicloud::migration",
|
||||
event = "storage_migration.blob_failed",
|
||||
event = "backend_migration.blob_failed",
|
||||
run_id = %store.run_id(),
|
||||
hash = %hash,
|
||||
error = %e,
|
||||
@@ -698,7 +698,7 @@ impl RecoverableJobHandler for StorageMigrationService {
|
||||
// where the admin UI reads it.
|
||||
record_or_log(
|
||||
store,
|
||||
STORAGE_MIGRATION_JOB_NAME,
|
||||
BACKEND_MIGRATION_JOB_NAME,
|
||||
"migration_failed",
|
||||
"data_loss",
|
||||
None,
|
||||
@@ -760,7 +760,7 @@ impl RecoverableJobHandler for StorageMigrationService {
|
||||
}
|
||||
}
|
||||
|
||||
impl StorageMigrationService {
|
||||
impl BackendMigrationService {
|
||||
/// Terminal successful path — reached from both Completed sites
|
||||
/// in the batch loop (empty-first-batch and short-batch).
|
||||
///
|
||||
@@ -831,7 +831,7 @@ impl StorageMigrationService {
|
||||
if !readonly_persisted {
|
||||
tracing::warn!(
|
||||
target: "oxicloud::migration",
|
||||
event = "storage_migration.readonly_clear_persist_failed",
|
||||
event = "backend_migration.readonly_clear_persist_failed",
|
||||
run_id = %store.run_id(),
|
||||
"cleared migration_readonly in memory (writes allowed against source) but the \
|
||||
DB persist failed. Boot-clear rule will fix on next restart."
|
||||
@@ -839,7 +839,7 @@ impl StorageMigrationService {
|
||||
}
|
||||
tracing::info!(
|
||||
target: "audit",
|
||||
event = "storage_migration.aborted",
|
||||
event = "backend_migration.aborted",
|
||||
reason = "blobs_failed",
|
||||
run_id = %store.run_id(),
|
||||
active_backend_name = previous_active,
|
||||
@@ -848,7 +848,7 @@ impl StorageMigrationService {
|
||||
skipped = skipped,
|
||||
failed = failed,
|
||||
source_missing = source_missing,
|
||||
"🛑 storage_migration aborted — {failed} blob(s) failed, active backend left at \
|
||||
"🛑 backend_migration aborted — {failed} blob(s) failed, active backend left at \
|
||||
`{previous_active}`, readonly cleared. Inspect findings and retry, or accept \
|
||||
the partial migration via `oxicloud --select-storage {target_name}`."
|
||||
);
|
||||
@@ -907,7 +907,7 @@ impl StorageMigrationService {
|
||||
if !readonly_persisted {
|
||||
tracing::warn!(
|
||||
target: "oxicloud::migration",
|
||||
event = "storage_migration.readonly_clear_persist_failed",
|
||||
event = "backend_migration.readonly_clear_persist_failed",
|
||||
run_id = %store.run_id(),
|
||||
"cleared migration_readonly in memory (writes allowed) but the DB persist \
|
||||
failed. If the server crashes before next boot, boot will re-seed the flag \
|
||||
@@ -917,7 +917,7 @@ impl StorageMigrationService {
|
||||
|
||||
tracing::info!(
|
||||
target: "audit",
|
||||
event = "storage_migration.completed",
|
||||
event = "backend_migration.completed",
|
||||
run_id = %store.run_id(),
|
||||
active_backend_name = target_name,
|
||||
previous_active = previous_active,
|
||||
@@ -925,11 +925,11 @@ impl StorageMigrationService {
|
||||
skipped = skipped,
|
||||
failed = failed,
|
||||
source_missing = source_missing,
|
||||
"✅ storage_migration completed — hot-swapped runtime backend to `{target_name}`, \
|
||||
"✅ backend_migration completed — hot-swapped runtime backend to `{target_name}`, \
|
||||
writes resumed. No restart required."
|
||||
);
|
||||
// Per-run summary counters merged into `stats` for the admin
|
||||
// UI drawer. Same shape as `storage_rotate`'s extras + one
|
||||
// UI drawer. Same shape as `backend_rotate`'s extras + one
|
||||
// extra `source_missing` counter unique to migration.
|
||||
RunOutcome::completed_with(serde_json::json!({
|
||||
"copied": copied,
|
||||
@@ -1022,7 +1022,7 @@ async fn collect_stream_bytes(
|
||||
let mut stream = std::pin::pin!(stream);
|
||||
while let Some(chunk) = stream.next().await {
|
||||
let bytes = chunk.map_err(|e| {
|
||||
DomainError::internal_error("StorageMigration", format!("source stream read: {e}"))
|
||||
DomainError::internal_error("BackendMigration", format!("source stream read: {e}"))
|
||||
})?;
|
||||
buf.extend_from_slice(&bytes);
|
||||
}
|
||||
+29
-29
@@ -15,13 +15,13 @@
|
||||
//!
|
||||
//! ### No readonly, no cutover
|
||||
//!
|
||||
//! `storage_rotate` is per-blob idempotent — repeat rewrites are
|
||||
//! `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 `storage_migration`,
|
||||
//! ever engaged — a critical improvement over `backend_migration`,
|
||||
//! whose target-different-from-source cutover forces one.
|
||||
//!
|
||||
//! ### Restart survival
|
||||
@@ -36,13 +36,13 @@
|
||||
//! ### Design notes
|
||||
//!
|
||||
//! * **Cursor** — UTF-8 hex of the last-processed blob hash (64
|
||||
//! chars). Same encoding as `storage_migration` and
|
||||
//! 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 `storage_migration`).
|
||||
//! 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
|
||||
@@ -68,20 +68,20 @@ use crate::infrastructure::scheduler::{
|
||||
use crate::infrastructure::services::encrypted_blob_backend::BlobFormat;
|
||||
use crate::infrastructure::services::entry_backend::build_entry_backend_typed;
|
||||
|
||||
pub const STORAGE_ROTATE_JOB_NAME: &str = "storage_rotate";
|
||||
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 `storage_migration`'s TARGET_NAME_PARAM so
|
||||
/// 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 `storage_migration` / `blobs_consistency`
|
||||
/// 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 StorageRotateService {
|
||||
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`.
|
||||
@@ -99,7 +99,7 @@ pub struct StorageRotateService {
|
||||
rotation_progress: Arc<std::sync::RwLock<Option<MigrationProgress>>>,
|
||||
}
|
||||
|
||||
impl StorageRotateService {
|
||||
impl BackendRotateService {
|
||||
pub fn new(
|
||||
pool: Arc<PgPool>,
|
||||
storage_entries: Vec<NamedStorageEntry>,
|
||||
@@ -115,7 +115,7 @@ impl StorageRotateService {
|
||||
}
|
||||
|
||||
/// Chainable self-registration — mirrors the `*_consistency`
|
||||
/// tenants and `storage_migration`. On-demand only (no periodic
|
||||
/// tenants and `backend_migration`. On-demand only (no periodic
|
||||
/// tick).
|
||||
pub async fn register_recoverable_job(
|
||||
self: Arc<Self>,
|
||||
@@ -130,13 +130,13 @@ impl StorageRotateService {
|
||||
}
|
||||
|
||||
#[async_trait]
|
||||
impl RecoverableJobHandler for StorageRotateService {
|
||||
impl RecoverableJobHandler for BackendRotateService {
|
||||
fn name(&self) -> &str {
|
||||
STORAGE_ROTATE_JOB_NAME
|
||||
BACKEND_ROTATE_JOB_NAME
|
||||
}
|
||||
|
||||
/// Definitive count — one row per blob. Same query as
|
||||
/// `storage_migration::count_total`; the two walk the same rows.
|
||||
/// `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())
|
||||
@@ -146,7 +146,7 @@ impl RecoverableJobHandler for StorageRotateService {
|
||||
Err(e) => {
|
||||
tracing::debug!(
|
||||
target: "oxicloud::rotate",
|
||||
event = "storage_rotate.count_total_failed",
|
||||
event = "backend_rotate.count_total_failed",
|
||||
error = %e,
|
||||
"count_total failed — run will not surface a progress bar"
|
||||
);
|
||||
@@ -161,12 +161,12 @@ impl RecoverableJobHandler for StorageRotateService {
|
||||
args: &JobRunArgs,
|
||||
resume_cursor: Option<Vec<u8>>,
|
||||
) -> RunOutcome {
|
||||
// Resolve target entry name — same shape as `storage_migration`.
|
||||
// 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: "storage_rotate requires `target_name` on a fresh run — trigger via \
|
||||
message: "backend_rotate requires `target_name` on a fresh run — trigger via \
|
||||
POST /api/admin/storage/entries/{name}/rotate"
|
||||
.to_string(),
|
||||
};
|
||||
@@ -230,7 +230,7 @@ impl RecoverableJobHandler for StorageRotateService {
|
||||
|
||||
tracing::info!(
|
||||
target: "audit",
|
||||
event = "storage_rotate.run_started",
|
||||
event = "backend_rotate.run_started",
|
||||
run_id = %store.run_id(),
|
||||
target_name = %target_name,
|
||||
// `%` (Display) → SSH-style `encrypted-v1 key_fp=83:96:...`
|
||||
@@ -239,7 +239,7 @@ impl RecoverableJobHandler for StorageRotateService {
|
||||
// header bytes on disk.
|
||||
head_format = %head_format,
|
||||
resuming = !is_fresh,
|
||||
"storage_rotate started on `{target_name}` (head_format = {head_format})"
|
||||
"backend_rotate started on `{target_name}` (head_format = {head_format})"
|
||||
);
|
||||
|
||||
// Seed the progress snapshot. Total = count_total's estimate;
|
||||
@@ -280,12 +280,12 @@ impl RecoverableJobHandler for StorageRotateService {
|
||||
self.clear_progress();
|
||||
tracing::info!(
|
||||
target: "oxicloud::rotate",
|
||||
event = "storage_rotate.cancelled",
|
||||
event = "backend_rotate.cancelled",
|
||||
run_id = %store.run_id(),
|
||||
rewritten = rewritten_count,
|
||||
skipped = skipped_count,
|
||||
failed = failed_count,
|
||||
"storage_rotate cancelled cooperatively, pausing"
|
||||
"backend_rotate cancelled cooperatively, pausing"
|
||||
);
|
||||
return RunOutcome::Paused {
|
||||
cursor: cursor
|
||||
@@ -304,7 +304,7 @@ impl RecoverableJobHandler for StorageRotateService {
|
||||
}
|
||||
|
||||
// Fetch the next batch. Same keyset pagination shape as
|
||||
// `storage_migration` — `hash > $1` on the PK, index-only.
|
||||
// `backend_migration` — `hash > $1` on the PK, index-only.
|
||||
let rows: Vec<(String,)> = match sqlx::query_as(
|
||||
r#"
|
||||
SELECT hash
|
||||
@@ -351,7 +351,7 @@ impl RecoverableJobHandler for StorageRotateService {
|
||||
failed_count += 1;
|
||||
tracing::warn!(
|
||||
target: "oxicloud::rotate",
|
||||
event = "storage_rotate.read_failed",
|
||||
event = "backend_rotate.read_failed",
|
||||
run_id = %store.run_id(),
|
||||
hash = %hash,
|
||||
error = %e,
|
||||
@@ -359,7 +359,7 @@ impl RecoverableJobHandler for StorageRotateService {
|
||||
);
|
||||
record_or_log(
|
||||
store,
|
||||
STORAGE_ROTATE_JOB_NAME,
|
||||
BACKEND_ROTATE_JOB_NAME,
|
||||
"rotation_failed",
|
||||
"data_loss",
|
||||
None,
|
||||
@@ -402,7 +402,7 @@ impl RecoverableJobHandler for StorageRotateService {
|
||||
failed_count += 1;
|
||||
tracing::warn!(
|
||||
target: "oxicloud::rotate",
|
||||
event = "storage_rotate.write_failed",
|
||||
event = "backend_rotate.write_failed",
|
||||
run_id = %store.run_id(),
|
||||
hash = %hash,
|
||||
error = %e,
|
||||
@@ -410,7 +410,7 @@ impl RecoverableJobHandler for StorageRotateService {
|
||||
);
|
||||
record_or_log(
|
||||
store,
|
||||
STORAGE_ROTATE_JOB_NAME,
|
||||
BACKEND_ROTATE_JOB_NAME,
|
||||
"rotation_failed",
|
||||
"data_loss",
|
||||
None,
|
||||
@@ -467,9 +467,9 @@ impl RecoverableJobHandler for StorageRotateService {
|
||||
}
|
||||
}
|
||||
|
||||
impl StorageRotateService {
|
||||
impl BackendRotateService {
|
||||
/// Terminal successful path — clear the header snapshot and log a
|
||||
/// final audit line. Unlike `storage_migration::finish_completed`
|
||||
/// 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.
|
||||
///
|
||||
@@ -508,14 +508,14 @@ impl StorageRotateService {
|
||||
|
||||
tracing::info!(
|
||||
target: "audit",
|
||||
event = "storage_rotate.run_completed",
|
||||
event = "backend_rotate.run_completed",
|
||||
run_id = %store.run_id(),
|
||||
target_name = %target_name,
|
||||
rewritten = rewritten,
|
||||
skipped = skipped,
|
||||
failed = failed,
|
||||
head_format = %head_display,
|
||||
"storage_rotate completed on `{target_name}` — {rewritten} rewritten, {skipped} skipped, {failed} failed; head = {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
|
||||
@@ -75,7 +75,7 @@ pub const BLOBS_CONSISTENCY_JOB_NAME: &str = "blobs_consistency";
|
||||
|
||||
/// `params` JSONB key under which the entry name being probed is
|
||||
/// stashed on a Fresh run (matches `TARGET_NAME_PARAM` on
|
||||
/// `storage_migration`). Resumed runs re-read it so a paused audit
|
||||
/// `backend_migration`). Resumed runs re-read it so a paused audit
|
||||
/// survives restart without the admin re-specifying the target.
|
||||
pub const PROBED_STORAGE_PARAM: &str = "probed_storage";
|
||||
|
||||
@@ -189,7 +189,7 @@ impl RecoverableJobHandler for BlobsConsistencyCheck {
|
||||
resume_cursor: Option<Vec<u8>>,
|
||||
) -> RunOutcome {
|
||||
// Resolve the backend to probe. Two paths, mirroring the
|
||||
// Fresh/Resumed split the storage_migration handler uses:
|
||||
// Fresh/Resumed split the backend_migration handler uses:
|
||||
//
|
||||
// * Fresh + args.storage=Some — probe that named entry
|
||||
// instead of the live backend. Stamp probed_storage in
|
||||
|
||||
@@ -142,7 +142,7 @@ pub struct EncryptedBlobBackend {
|
||||
/// * `read_dispatch` legacy-fallback path — iterates in order
|
||||
/// (oldest → newest) to try every real-cipher pair when a
|
||||
/// legacy blob's head-key decrypt fails.
|
||||
/// * K3 `storage_rotate` — needs to walk pair indices.
|
||||
/// * K3 `backend_rotate` — needs to walk pair indices.
|
||||
pairs: Vec<KeyPair>,
|
||||
/// `<key_fp>` → per-pair cipher, for O(1) read dispatch on v1
|
||||
/// blobs. Excludes any `none:` pair (nothing to build). Cloned
|
||||
@@ -218,7 +218,7 @@ impl EncryptedBlobBackend {
|
||||
key
|
||||
}
|
||||
|
||||
/// The format `storage_rotate` should normalise every blob TO —
|
||||
/// The format `backend_rotate` should normalise every blob TO —
|
||||
/// derived from the wrapper's head pair. When
|
||||
/// `head_cipher.is_some()` we're writing encrypted-v1 with the
|
||||
/// head pair's `key_fp`; when it's `None` we're writing
|
||||
@@ -237,8 +237,8 @@ impl EncryptedBlobBackend {
|
||||
}
|
||||
}
|
||||
|
||||
/// Smart-skip probe used by `storage_migration` (and potentially
|
||||
/// `storage_rotate` if it ever gains a fast-path). Reads the
|
||||
/// Smart-skip probe used by `backend_migration` (and potentially
|
||||
/// `backend_rotate` if it ever gains a fast-path). Reads the
|
||||
/// first [`HEADER_SIZE`] bytes of the on-disk blob and returns
|
||||
/// `true` iff:
|
||||
///
|
||||
@@ -280,7 +280,7 @@ impl EncryptedBlobBackend {
|
||||
}
|
||||
|
||||
/// Fetch, classify, and decrypt a blob in one round-trip. Used by
|
||||
/// K3's `storage_rotate` per-blob step: it needs both the
|
||||
/// K3's `backend_rotate` per-blob step: it needs both the
|
||||
/// plaintext (to re-encrypt under the head pair) AND the current
|
||||
/// on-disk format (to decide whether a rewrite is needed at all).
|
||||
///
|
||||
@@ -317,7 +317,7 @@ impl EncryptedBlobBackend {
|
||||
}
|
||||
|
||||
/// Classification of a raw blob's on-disk format. Exposed for K3's
|
||||
/// `storage_rotate` decision tree; not used on the hot request path.
|
||||
/// `backend_rotate` decision tree; not used on the hot request path.
|
||||
///
|
||||
/// PartialEq is derived so `current == head_format` collapses the
|
||||
/// plan's six-case decision tree into a single equality check:
|
||||
@@ -541,7 +541,7 @@ fn read_dispatch(
|
||||
hash = %expected_hash,
|
||||
size = encrypted.len(),
|
||||
"🩹 legacy plaintext blob served via BLAKE3 rescue — no configured key \
|
||||
decrypted it, but content hash matched. Run storage_rotate to re-write \
|
||||
decrypted it, but content hash matched. Run backend_rotate to re-write \
|
||||
under the current head."
|
||||
);
|
||||
return Ok(Bytes::from(encrypted));
|
||||
@@ -737,7 +737,7 @@ impl BlobStorageBackend for EncryptedBlobBackend {
|
||||
|
||||
/// Frame the plaintext with the head pair's format (encrypted-v1
|
||||
/// or plaintext-v1), then delegate the atomic replace to the
|
||||
/// inner backend. Used by `storage_rotate` to actually change the
|
||||
/// inner backend. Used by `backend_rotate` to actually change the
|
||||
/// on-disk bytes — `put_blob_from_bytes` would silently no-op on
|
||||
/// `LocalBlobBackend` when the object key already exists.
|
||||
fn put_blob_from_bytes_replace(
|
||||
@@ -1373,7 +1373,7 @@ mod tests {
|
||||
// ─────────────────────────────────────────────────────────────
|
||||
// K3 tests — BlobFormat classifier + head_format + read_and_classify.
|
||||
//
|
||||
// These pin the format-inspection contract that `storage_rotate`
|
||||
// These pin the format-inspection contract that `backend_rotate`
|
||||
// depends on. The rotate job's per-blob decision tree collapses
|
||||
// to `current != head_format ? rewrite : skip`, so any drift in
|
||||
// either helper would silently change rotation semantics.
|
||||
|
||||
@@ -201,7 +201,7 @@ pub async fn resolve_active_entry<'a>(
|
||||
///
|
||||
/// Same construction path as `build_entry_backend`; the trait-object
|
||||
/// version delegates through this. Preferred for job handlers
|
||||
/// (`storage_rotate`) that need typed access. The trait-object
|
||||
/// (`backend_rotate`) that need typed access. The trait-object
|
||||
/// version stays for the DI hot-path where the caller only needs
|
||||
/// the generic `BlobStorageBackend` contract.
|
||||
pub fn build_entry_backend_typed(
|
||||
|
||||
@@ -470,7 +470,7 @@ impl BlobStorageBackend for LocalBlobBackend {
|
||||
/// then `rename(2)` over the target. `write_blob_bytes`'s
|
||||
/// `O_CREAT|O_EXCL` idempotent-skip (the right choice for uploads)
|
||||
/// silently no-ops when the target already exists — wrong for
|
||||
/// callers like `storage_rotate` that need the bytes to change.
|
||||
/// callers like `backend_rotate` that need the bytes to change.
|
||||
/// See the trait doc for the full picture.
|
||||
///
|
||||
/// Tempfile lives beside the target under the same shard directory
|
||||
|
||||
@@ -1,6 +1,8 @@
|
||||
pub mod audio_metadata_service;
|
||||
pub mod azure_blob_backend;
|
||||
pub mod backend_consistency_service;
|
||||
pub mod backend_migration_service;
|
||||
pub mod backend_rotate_service;
|
||||
pub mod blobs_consistency_service;
|
||||
pub mod cached_blob_backend;
|
||||
pub mod chunked_upload_service;
|
||||
@@ -45,8 +47,6 @@ pub mod s3_blob_backend;
|
||||
pub mod search_index;
|
||||
pub mod share_unlock_cookie;
|
||||
pub mod smtp_email_sender;
|
||||
pub mod storage_migration_service;
|
||||
pub mod storage_rotate_service;
|
||||
pub mod swappable_blob_backend;
|
||||
pub mod thumbnail_service;
|
||||
#[cfg(test)]
|
||||
|
||||
@@ -237,8 +237,8 @@ impl BlobStorageBackend for S3BlobBackend {
|
||||
})
|
||||
}
|
||||
|
||||
/// Atomic overwrite path used by `storage_rotate` and
|
||||
/// `storage_migration` when re-writing an already-present blob
|
||||
/// Atomic overwrite path used by `backend_rotate` and
|
||||
/// `backend_migration` when re-writing an already-present blob
|
||||
/// under a new head key/format. Trait default delegates to
|
||||
/// `put_blob_from_bytes` which HEAD-probes and silently skips —
|
||||
/// exactly wrong for the rotate/migrate use case (the whole
|
||||
|
||||
@@ -75,9 +75,9 @@ pub fn admin_routes() -> Router<Arc<AppState>> {
|
||||
.route("/settings/storage", put(save_storage_settings))
|
||||
.route("/settings/storage/test", post(test_storage_connection))
|
||||
// Storage migration — thin shims over the recoverable-run
|
||||
// engine (job_name = "storage_migration"). Retained under
|
||||
// engine (job_name = "backend_migration"). Retained under
|
||||
// /storage/migration/* until the admin UI is rewired to
|
||||
// /api/admin/jobs/storage_migration/*; both paths route to
|
||||
// /api/admin/jobs/backend_migration/*; both paths route to
|
||||
// the same underlying JobRegistry dispatch. The old /complete
|
||||
// endpoint is retired — a finished run is a Completed row,
|
||||
// there's nothing to acknowledge.
|
||||
@@ -91,7 +91,7 @@ pub fn admin_routes() -> Router<Arc<AppState>> {
|
||||
// new-key). No readonly mode; safe under normal traffic.
|
||||
.route(
|
||||
"/storage/entries/{name}/rotate",
|
||||
post(trigger_storage_rotate),
|
||||
post(trigger_backend_rotate),
|
||||
)
|
||||
// NOTE: /storage/migration/verify retired in slice 7 (see the
|
||||
// comment near where `verify_migration` used to live). Use
|
||||
@@ -383,7 +383,7 @@ async fn test_storage_connection(
|
||||
/// GET /api/admin/storage/migration — current migration progress.
|
||||
///
|
||||
/// Shim over the recoverable-run engine: reads the latest
|
||||
/// `storage_migration` run from `jobs.recoverable_runs` (via the
|
||||
/// `backend_migration` run from `jobs.recoverable_runs` (via the
|
||||
/// `JobStoreProvider`) and projects it into the legacy
|
||||
/// `MigrationStateDto` shape the admin storage tab expects. When no
|
||||
/// run has ever been triggered the response is an empty "idle" DTO —
|
||||
@@ -403,11 +403,11 @@ async fn test_storage_connection(
|
||||
pub async fn get_migration_status(
|
||||
State(state): State<Arc<AppState>>,
|
||||
) -> Result<impl IntoResponse, AppError> {
|
||||
use crate::infrastructure::services::storage_migration_service::STORAGE_MIGRATION_JOB_NAME;
|
||||
use crate::infrastructure::services::backend_migration_service::BACKEND_MIGRATION_JOB_NAME;
|
||||
|
||||
let provider = state.core.job_store_provider.clone();
|
||||
let latest = provider
|
||||
.list_runs(STORAGE_MIGRATION_JOB_NAME, 1)
|
||||
.list_runs(BACKEND_MIGRATION_JOB_NAME, 1)
|
||||
.await
|
||||
.map_err(AppError::from)?
|
||||
.into_iter()
|
||||
@@ -440,7 +440,7 @@ pub async fn get_migration_status(
|
||||
|
||||
/// POST /api/admin/storage/migration/start — begin background migration.
|
||||
///
|
||||
/// Shim that forwards to `JobRegistry::trigger("storage_migration",
|
||||
/// Shim that forwards to `JobRegistry::trigger("backend_migration",
|
||||
/// ...)`. `run_or_resume` (the RecoverableAdapter's inner dispatch)
|
||||
/// resumes a Paused run or starts a fresh one — one endpoint covers
|
||||
/// both. Exclusivity is enforced at the DB layer (the partial unique
|
||||
@@ -500,7 +500,7 @@ pub async fn start_migration(
|
||||
dto.target_name
|
||||
)));
|
||||
}
|
||||
trigger_storage_migration(state, Some(dto.target_name)).await
|
||||
trigger_backend_migration(state, Some(dto.target_name)).await
|
||||
}
|
||||
|
||||
/// POST /api/admin/storage/migration/pause — pause a running migration.
|
||||
@@ -524,18 +524,18 @@ pub async fn start_migration(
|
||||
pub async fn pause_migration(
|
||||
State(state): State<Arc<AppState>>,
|
||||
) -> Result<impl IntoResponse, AppError> {
|
||||
use crate::infrastructure::services::storage_migration_service::STORAGE_MIGRATION_JOB_NAME;
|
||||
use crate::infrastructure::services::backend_migration_service::BACKEND_MIGRATION_JOB_NAME;
|
||||
|
||||
tracing::info!(
|
||||
target: "audit",
|
||||
event = "storage_migration.pause_requested",
|
||||
"👮🏻♂️ Admin requested storage_migration pause"
|
||||
event = "backend_migration.pause_requested",
|
||||
"👮🏻♂️ Admin requested backend_migration pause"
|
||||
);
|
||||
|
||||
let flipped = state
|
||||
.core
|
||||
.job_store_provider
|
||||
.request_cancel(STORAGE_MIGRATION_JOB_NAME)
|
||||
.request_cancel(BACKEND_MIGRATION_JOB_NAME)
|
||||
.await
|
||||
.map_err(AppError::from)?;
|
||||
|
||||
@@ -576,7 +576,7 @@ pub async fn resume_migration(
|
||||
// it from `params.target_name` stamped on the original Fresh
|
||||
// open. Refuses gracefully via RunOutcome::Failed if there is
|
||||
// no Paused row to resume.
|
||||
trigger_storage_migration(state, None).await
|
||||
trigger_backend_migration(state, None).await
|
||||
}
|
||||
|
||||
// verify_migration endpoint retired (slice 7 of
|
||||
@@ -596,18 +596,18 @@ pub async fn resume_migration(
|
||||
/// desync `current_run_start` from the actually-running task). The
|
||||
/// admin UI polls `GET /storage/migration` for progress; the trigger
|
||||
/// itself is fire-and-forget.
|
||||
async fn trigger_storage_migration(
|
||||
async fn trigger_backend_migration(
|
||||
state: Arc<AppState>,
|
||||
target_name: Option<String>,
|
||||
) -> Result<axum::response::Response, AppError> {
|
||||
use crate::infrastructure::scheduler::JobRunArgs;
|
||||
use crate::infrastructure::services::storage_migration_service::STORAGE_MIGRATION_JOB_NAME;
|
||||
use crate::infrastructure::services::backend_migration_service::BACKEND_MIGRATION_JOB_NAME;
|
||||
|
||||
tracing::info!(
|
||||
target: "audit",
|
||||
event = "storage_migration.trigger_requested",
|
||||
event = "backend_migration.trigger_requested",
|
||||
target_name = target_name.as_deref().unwrap_or("<resume>"),
|
||||
"👮🏻♂️ Admin triggered storage_migration"
|
||||
"👮🏻♂️ Admin triggered backend_migration"
|
||||
);
|
||||
|
||||
let registry = state.core.job_registry.clone();
|
||||
@@ -616,7 +616,7 @@ async fn trigger_storage_migration(
|
||||
..JobRunArgs::default()
|
||||
};
|
||||
tokio::spawn(async move {
|
||||
registry.trigger(STORAGE_MIGRATION_JOB_NAME, &args).await;
|
||||
registry.trigger(BACKEND_MIGRATION_JOB_NAME, &args).await;
|
||||
});
|
||||
|
||||
Ok((
|
||||
@@ -630,7 +630,7 @@ async fn trigger_storage_migration(
|
||||
}
|
||||
|
||||
/// POST /api/admin/storage/entries/{name}/rotate — trigger the
|
||||
/// `storage_rotate` recoverable job on a specific entry.
|
||||
/// `backend_rotate` recoverable job on a specific entry.
|
||||
///
|
||||
/// Normalises every blob on `<name>` to the entry's head-pair
|
||||
/// format: legacy → v1, plaintext ↔ encrypted, old-key → new-key.
|
||||
@@ -658,13 +658,13 @@ async fn trigger_storage_migration(
|
||||
security(("bearerAuth" = [])),
|
||||
tag = "admin"
|
||||
)]
|
||||
pub async fn trigger_storage_rotate(
|
||||
pub async fn trigger_backend_rotate(
|
||||
State(state): State<Arc<AppState>>,
|
||||
axum::extract::Path(name): axum::extract::Path<String>,
|
||||
) -> Result<impl IntoResponse, AppError> {
|
||||
use crate::infrastructure::scheduler::JobRunArgs;
|
||||
use crate::infrastructure::services::storage_migration_service::STORAGE_MIGRATION_JOB_NAME;
|
||||
use crate::infrastructure::services::storage_rotate_service::STORAGE_ROTATE_JOB_NAME;
|
||||
use crate::infrastructure::services::backend_migration_service::BACKEND_MIGRATION_JOB_NAME;
|
||||
use crate::infrastructure::services::backend_rotate_service::BACKEND_ROTATE_JOB_NAME;
|
||||
|
||||
// Synchronous entry-existence check — a bad name would fail the
|
||||
// run anyway, but returning 400 here spares the operator an
|
||||
@@ -699,7 +699,7 @@ pub async fn trigger_storage_rotate(
|
||||
.clone();
|
||||
if name != active {
|
||||
return Err(AppError::bad_request(format!(
|
||||
"storage_rotate refuses non-active entry `{name}` — the DB blob registry \
|
||||
"backend_rotate refuses non-active entry `{name}` — the DB blob registry \
|
||||
describes the active entry (`{active}`), so walking it against a stale \
|
||||
target produces spurious `rotation_failed` findings. Activate `{name}` \
|
||||
first via `Migrate & activate`, then rotate."
|
||||
@@ -713,7 +713,7 @@ pub async fn trigger_storage_rotate(
|
||||
// hash. Cheap check — `list_runs` limit 1 with the status
|
||||
// filter is an index scan.
|
||||
let provider = state.core.job_store_provider.clone();
|
||||
for job_name in [STORAGE_ROTATE_JOB_NAME, STORAGE_MIGRATION_JOB_NAME] {
|
||||
for job_name in [BACKEND_ROTATE_JOB_NAME, BACKEND_MIGRATION_JOB_NAME] {
|
||||
let in_flight = provider
|
||||
.list_runs(job_name, 5)
|
||||
.await
|
||||
@@ -729,7 +729,7 @@ pub async fn trigger_storage_rotate(
|
||||
});
|
||||
if in_flight {
|
||||
return Err(AppError::bad_request(format!(
|
||||
"cannot start storage_rotate on `{name}` — `{job_name}` is already Running / \
|
||||
"cannot start backend_rotate on `{name}` — `{job_name}` is already Running / \
|
||||
Paused / CancelRequested. Wait for it to finish (or cancel via \
|
||||
`POST /api/admin/jobs/{job_name}/cancel`)."
|
||||
)));
|
||||
@@ -738,9 +738,9 @@ pub async fn trigger_storage_rotate(
|
||||
|
||||
tracing::info!(
|
||||
target: "audit",
|
||||
event = "storage_rotate.trigger_requested",
|
||||
event = "backend_rotate.trigger_requested",
|
||||
target_name = %name,
|
||||
"👮🏻♂️ Admin triggered storage_rotate on `{name}`"
|
||||
"👮🏻♂️ Admin triggered backend_rotate on `{name}`"
|
||||
);
|
||||
|
||||
let registry = state.core.job_registry.clone();
|
||||
@@ -749,13 +749,13 @@ pub async fn trigger_storage_rotate(
|
||||
..JobRunArgs::default()
|
||||
};
|
||||
tokio::spawn(async move {
|
||||
registry.trigger(STORAGE_ROTATE_JOB_NAME, &args).await;
|
||||
registry.trigger(BACKEND_ROTATE_JOB_NAME, &args).await;
|
||||
});
|
||||
|
||||
Ok((
|
||||
StatusCode::ACCEPTED,
|
||||
Json(serde_json::json!({
|
||||
"message": format!("Rotation dispatched on `{name}` — poll GET /api/admin/jobs/{STORAGE_ROTATE_JOB_NAME} for progress"),
|
||||
"message": format!("Rotation dispatched on `{name}` — poll GET /api/admin/jobs/{BACKEND_ROTATE_JOB_NAME} for progress"),
|
||||
"detached": true,
|
||||
})),
|
||||
)
|
||||
@@ -2348,7 +2348,7 @@ pub struct TriggerJobQuery {
|
||||
pub deep: bool,
|
||||
/// Optional named storage entry to scope the run against — used by
|
||||
/// tenants that respect `JobRunArgs.storage` (currently
|
||||
/// `storage_migration` for its target; `blobs_consistency` /
|
||||
/// `backend_migration` for its target; `blobs_consistency` /
|
||||
/// `backend_consistency` will pick this up in slice 7 to probe a
|
||||
/// non-active entry). Ignored by tenants that don't declare a
|
||||
/// semantic for it. Unknown-name validation is per-tenant — the
|
||||
@@ -2405,7 +2405,7 @@ pub async fn trigger_job(
|
||||
storage: query.storage.clone(),
|
||||
};
|
||||
|
||||
// Jobs that can run for hours (storage_migration, future
|
||||
// Jobs that can run for hours (backend_migration, future
|
||||
// reextract_*) are detached: `tokio::spawn` the trigger so the
|
||||
// HTTP request returns immediately. Without this, browser HTTP
|
||||
// timeouts drop the request future mid-await → the SemaphorePermit
|
||||
@@ -2461,7 +2461,7 @@ pub async fn trigger_job(
|
||||
/// there's a second long-running tenant that justifies the plumbing.
|
||||
/// See the comment in `trigger_job` for why detach matters.
|
||||
fn is_detached_job(name: &str) -> bool {
|
||||
matches!(name, "storage_migration")
|
||||
matches!(name, "backend_migration")
|
||||
}
|
||||
|
||||
/// `POST /api/admin/jobs/{name}/cancel` — cooperative cancel of the
|
||||
|
||||
+3
-3
@@ -172,7 +172,7 @@ fn main() -> Result<(), Box<dyn std::error::Error>> {
|
||||
// One-shot helper: compute the SSH-style colon-hex
|
||||
// fingerprint of a base64-encoded AES-256 key and
|
||||
// print to stdout. Same truncation used by the v1
|
||||
// header's `<key_fp>` field + the `storage_rotate`
|
||||
// header's `<key_fp>` field + the `backend_rotate`
|
||||
// completion summary — so an admin can:
|
||||
// 1. Look at the `head_key_fp` reported by the
|
||||
// last rotate run.
|
||||
@@ -297,7 +297,7 @@ fn print_help() {
|
||||
println!(" oxicloud --fingerprint <base64key|-> One-shot helper — print the SSH-style");
|
||||
println!(" fingerprint of a base64 AES-256 key.");
|
||||
println!(" Same shape used by the v1 blob header");
|
||||
println!(" + `storage_rotate` completion summary.");
|
||||
println!(" + `backend_rotate` completion summary.");
|
||||
println!(" Read stdin with `-` to keep keys out");
|
||||
println!(" of shell history.");
|
||||
println!();
|
||||
@@ -327,7 +327,7 @@ fn print_help() {
|
||||
println!(" --fingerprint <base64key | ->");
|
||||
println!(" Compute the SSH-style colon-hex fingerprint (16-hex, 8-byte");
|
||||
println!(" truncation of sha256) of a base64-encoded AES-256 key. Matches the");
|
||||
println!(" `head_key_fp` field the `storage_rotate` job reports on completion,");
|
||||
println!(" `head_key_fp` field the `backend_rotate` job reports on completion,");
|
||||
println!(" and the raw <key_fp> field embedded in every v1 blob header. Used");
|
||||
println!(" to identify which key in `OXICLOUD_STORAGE_<N>_ENCRYPTION_KEY`");
|
||||
println!(" corresponds to the current on-disk head — safe to drop any key");
|
||||
|
||||
Reference in New Issue
Block a user