feat(storage-migration): move storage mig. to recoverable job
This commit is contained in:
@@ -1,273 +0,0 @@
|
||||
//! `MigrationBlobBackend` — decorator that enables zero-downtime migration
|
||||
//! between blob storage backends.
|
||||
//!
|
||||
//! During a migration the decorator writes to the **target** backend and reads
|
||||
//! from **target-first-then-source** (dual-read). A background job
|
||||
//! (see `migration_job.rs`) copies remaining blobs in the background.
|
||||
|
||||
use std::future::Future;
|
||||
use std::path::{Path, PathBuf};
|
||||
use std::pin::Pin;
|
||||
use std::sync::Arc;
|
||||
|
||||
use bytes::Bytes;
|
||||
use chrono::{DateTime, Utc};
|
||||
use serde::Serialize;
|
||||
use tokio::sync::RwLock;
|
||||
|
||||
use crate::application::ports::blob_storage_ports::{
|
||||
BlobStorageBackend, BlobStream, StorageHealthStatus,
|
||||
};
|
||||
use crate::common::errors::DomainError;
|
||||
|
||||
// ── Migration state ────────────────────────────────────────────────
|
||||
|
||||
/// Progress of an ongoing (or completed) backend migration.
|
||||
#[derive(Debug, Clone, Serialize)]
|
||||
pub struct MigrationState {
|
||||
pub status: MigrationStatus,
|
||||
pub total_blobs: u64,
|
||||
pub migrated_blobs: u64,
|
||||
pub migrated_bytes: u64,
|
||||
pub failed_blobs: Vec<String>,
|
||||
pub started_at: Option<DateTime<Utc>>,
|
||||
pub completed_at: Option<DateTime<Utc>>,
|
||||
}
|
||||
|
||||
impl Default for MigrationState {
|
||||
fn default() -> Self {
|
||||
Self {
|
||||
status: MigrationStatus::Idle,
|
||||
total_blobs: 0,
|
||||
migrated_blobs: 0,
|
||||
migrated_bytes: 0,
|
||||
failed_blobs: Vec::new(),
|
||||
started_at: None,
|
||||
completed_at: None,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Status of the migration job.
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
|
||||
#[serde(rename_all = "snake_case")]
|
||||
pub enum MigrationStatus {
|
||||
Idle,
|
||||
Running,
|
||||
Paused,
|
||||
Completed,
|
||||
Failed,
|
||||
}
|
||||
|
||||
// ── MigrationBlobBackend ───────────────────────────────────────────
|
||||
|
||||
/// A `BlobStorageBackend` decorator that proxies requests to a *source*
|
||||
/// (old) and *target* (new) backend, enabling live migration.
|
||||
pub struct MigrationBlobBackend {
|
||||
source: Arc<dyn BlobStorageBackend>,
|
||||
target: Arc<dyn BlobStorageBackend>,
|
||||
state: Arc<RwLock<MigrationState>>,
|
||||
}
|
||||
|
||||
impl MigrationBlobBackend {
|
||||
pub fn new(
|
||||
source: Arc<dyn BlobStorageBackend>,
|
||||
target: Arc<dyn BlobStorageBackend>,
|
||||
state: Arc<RwLock<MigrationState>>,
|
||||
) -> Self {
|
||||
Self {
|
||||
source,
|
||||
target,
|
||||
state,
|
||||
}
|
||||
}
|
||||
|
||||
pub fn state(&self) -> &Arc<RwLock<MigrationState>> {
|
||||
&self.state
|
||||
}
|
||||
|
||||
pub fn source(&self) -> &Arc<dyn BlobStorageBackend> {
|
||||
&self.source
|
||||
}
|
||||
|
||||
pub fn target(&self) -> &Arc<dyn BlobStorageBackend> {
|
||||
&self.target
|
||||
}
|
||||
}
|
||||
|
||||
/// Boxed future alias (same as in the trait module).
|
||||
type BoxFut<'a, T> = Pin<Box<dyn Future<Output = T> + Send + 'a>>;
|
||||
|
||||
impl BlobStorageBackend for MigrationBlobBackend {
|
||||
fn initialize(&self) -> BoxFut<'_, Result<(), DomainError>> {
|
||||
Box::pin(async move {
|
||||
self.target.initialize().await?;
|
||||
// Source is already initialised; call anyway for idempotency.
|
||||
self.source.initialize().await?;
|
||||
Ok(())
|
||||
})
|
||||
}
|
||||
|
||||
/// Writes go to **target** only.
|
||||
fn put_blob(&self, hash: &str, source_path: &Path) -> BoxFut<'_, Result<u64, DomainError>> {
|
||||
let hash = hash.to_string();
|
||||
let path = source_path.to_path_buf();
|
||||
Box::pin(async move { self.target.put_blob(&hash, &path).await })
|
||||
}
|
||||
|
||||
/// Writes bytes to **target** only.
|
||||
fn put_blob_from_bytes(&self, hash: &str, data: Bytes) -> BoxFut<'_, Result<u64, DomainError>> {
|
||||
let hash = hash.to_string();
|
||||
Box::pin(async move { self.target.put_blob_from_bytes(&hash, data).await })
|
||||
}
|
||||
|
||||
/// Unsynced writes go to **target** only (same as the synced variant).
|
||||
fn put_blob_from_bytes_unsynced(
|
||||
&self,
|
||||
hash: &str,
|
||||
data: Bytes,
|
||||
) -> BoxFut<'_, Result<u64, DomainError>> {
|
||||
let hash = hash.to_string();
|
||||
Box::pin(async move { self.target.put_blob_from_bytes_unsynced(&hash, data).await })
|
||||
}
|
||||
|
||||
/// Durability sweep goes to **target**, where unsynced writes land.
|
||||
fn sync_blobs(&self, hashes: &[String]) -> BoxFut<'_, Result<(), DomainError>> {
|
||||
self.target.sync_blobs(hashes)
|
||||
}
|
||||
|
||||
/// Read from target first; fall back to source.
|
||||
fn get_blob_stream(&self, hash: &str) -> BoxFut<'_, Result<BlobStream, DomainError>> {
|
||||
let hash = hash.to_string();
|
||||
Box::pin(async move {
|
||||
match self.target.get_blob_stream(&hash).await {
|
||||
Ok(stream) => Ok(stream),
|
||||
Err(_) => self.source.get_blob_stream(&hash).await,
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
fn get_blob_range_stream(
|
||||
&self,
|
||||
hash: &str,
|
||||
start: u64,
|
||||
end: Option<u64>,
|
||||
) -> BoxFut<'_, Result<BlobStream, DomainError>> {
|
||||
let hash = hash.to_string();
|
||||
Box::pin(async move {
|
||||
match self.target.get_blob_range_stream(&hash, start, end).await {
|
||||
Ok(stream) => Ok(stream),
|
||||
Err(_) => self.source.get_blob_range_stream(&hash, start, end).await,
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
/// Delete from **both** backends (best-effort on source).
|
||||
fn delete_blob(&self, hash: &str) -> BoxFut<'_, Result<(), DomainError>> {
|
||||
let hash = hash.to_string();
|
||||
Box::pin(async move {
|
||||
self.target.delete_blob(&hash).await?;
|
||||
// Best-effort on source — ignore errors (blob may already be gone).
|
||||
let _ = self.source.delete_blob(&hash).await;
|
||||
Ok(())
|
||||
})
|
||||
}
|
||||
|
||||
/// Exists in either backend.
|
||||
fn blob_exists(&self, hash: &str) -> BoxFut<'_, Result<bool, DomainError>> {
|
||||
let hash = hash.to_string();
|
||||
Box::pin(async move {
|
||||
if self.target.blob_exists(&hash).await? {
|
||||
return Ok(true);
|
||||
}
|
||||
self.source.blob_exists(&hash).await
|
||||
})
|
||||
}
|
||||
|
||||
fn blob_size(&self, hash: &str) -> BoxFut<'_, Result<u64, DomainError>> {
|
||||
let hash = hash.to_string();
|
||||
Box::pin(async move {
|
||||
match self.target.blob_size(&hash).await {
|
||||
Ok(sz) => Ok(sz),
|
||||
Err(_) => self.source.blob_size(&hash).await,
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
fn health_check(&self) -> BoxFut<'_, Result<StorageHealthStatus, DomainError>> {
|
||||
Box::pin(async move {
|
||||
let target_health = self.target.health_check().await?;
|
||||
let source_health = self.source.health_check().await?;
|
||||
Ok(StorageHealthStatus {
|
||||
connected: target_health.connected && source_health.connected,
|
||||
backend_type: format!(
|
||||
"migration({} → {})",
|
||||
source_health.backend_type, target_health.backend_type
|
||||
),
|
||||
message: format!(
|
||||
"Source: {} | Target: {}",
|
||||
source_health.message, target_health.message
|
||||
),
|
||||
available_bytes: target_health.available_bytes,
|
||||
})
|
||||
})
|
||||
}
|
||||
|
||||
fn backend_type(&self) -> &'static str {
|
||||
"migration"
|
||||
}
|
||||
|
||||
/// Reads are served target-first (see `get_blob_stream`), so adopt the
|
||||
/// target's read-ahead.
|
||||
fn read_prefetch(&self) -> usize {
|
||||
self.target.read_prefetch()
|
||||
}
|
||||
|
||||
fn local_blob_path(&self, hash: &str) -> Option<PathBuf> {
|
||||
// Prefer target, fall back to source.
|
||||
self.target
|
||||
.local_blob_path(hash)
|
||||
.or_else(|| self.source.local_blob_path(hash))
|
||||
}
|
||||
|
||||
/// Enumeration during migration is intentionally REFUSED. Both
|
||||
/// source and target legitimately hold bytes concurrently
|
||||
/// mid-migration: a blob copied to target but not yet deleted
|
||||
/// from source would be reported "twice"; a blob in-flight from
|
||||
/// source to target could be flagged as orphan on whichever
|
||||
/// side the consistency scan doesn't walk. There's no single
|
||||
/// authoritative "what's on the backend" answer while a
|
||||
/// migration is running.
|
||||
///
|
||||
/// Operators wanting to run `backend_consistency` during a
|
||||
/// migration should either wait for the migration to complete
|
||||
/// (target becomes authoritative) or cancel it. The
|
||||
/// `operation_not_supported` error is surfaced by the tenant as
|
||||
/// a single run-level `backend_unenumerable` finding — no
|
||||
/// per-blob probes attempted.
|
||||
fn list_blob_hashes(
|
||||
&self,
|
||||
_cursor: Option<String>,
|
||||
_limit: usize,
|
||||
) -> Pin<
|
||||
Box<
|
||||
dyn std::future::Future<
|
||||
Output = Result<
|
||||
crate::application::ports::blob_storage_ports::BlobListPage,
|
||||
DomainError,
|
||||
>,
|
||||
> + Send
|
||||
+ '_,
|
||||
>,
|
||||
> {
|
||||
Box::pin(async {
|
||||
Err(DomainError::operation_not_supported(
|
||||
"list_blob_hashes",
|
||||
"backend_consistency cannot enumerate while a storage \
|
||||
migration is in progress — source and target hold bytes \
|
||||
concurrently; wait for migration completion or cancel it \
|
||||
before running the scan",
|
||||
))
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -1,240 +0,0 @@
|
||||
//! Background migration job — copies blobs from a source backend to a target
|
||||
//! backend with configurable concurrency and progress tracking.
|
||||
|
||||
use std::sync::Arc;
|
||||
|
||||
use futures::StreamExt;
|
||||
use serde::Serialize;
|
||||
use sqlx::PgPool;
|
||||
use tokio::sync::RwLock;
|
||||
|
||||
use crate::application::ports::blob_storage_ports::BlobStorageBackend;
|
||||
use crate::common::errors::DomainError;
|
||||
use crate::infrastructure::services::migration_blob_backend::{MigrationState, MigrationStatus};
|
||||
|
||||
/// Run the migration: stream all blob hashes from `storage.blobs` and copy
|
||||
/// each one from `source` to `target`.
|
||||
///
|
||||
/// * The job respects `Paused` / `Failed` status in `state` — it will stop
|
||||
/// streaming when the status is no longer `Running`.
|
||||
/// * Errors on individual blobs are logged and collected in `failed_blobs`
|
||||
/// but do **not** abort the full run.
|
||||
/// * `concurrency` controls `buffer_unordered` parallelism (default: 4).
|
||||
pub async fn run_migration(
|
||||
source: Arc<dyn BlobStorageBackend>,
|
||||
target: Arc<dyn BlobStorageBackend>,
|
||||
pool: Arc<PgPool>,
|
||||
state: Arc<RwLock<MigrationState>>,
|
||||
concurrency: usize,
|
||||
) -> Result<(), DomainError> {
|
||||
// Count total blobs for progress tracking.
|
||||
let total: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM storage.blobs")
|
||||
.fetch_one(pool.as_ref())
|
||||
.await
|
||||
.unwrap_or(0);
|
||||
|
||||
{
|
||||
let mut s = state.write().await;
|
||||
s.status = MigrationStatus::Running;
|
||||
s.total_blobs = total as u64;
|
||||
s.migrated_blobs = 0;
|
||||
s.migrated_bytes = 0;
|
||||
s.failed_blobs.clear();
|
||||
s.started_at = Some(chrono::Utc::now());
|
||||
s.completed_at = None;
|
||||
}
|
||||
|
||||
// Stream all hashes+sizes with a cursor.
|
||||
let mut rows =
|
||||
sqlx::query_as::<_, (String, i64)>("SELECT hash, size FROM storage.blobs ORDER BY hash")
|
||||
.fetch(pool.as_ref());
|
||||
|
||||
// Collect all hashes first to avoid holding the cursor across awaits.
|
||||
let mut work: Vec<(String, i64)> = Vec::with_capacity(total as usize);
|
||||
while let Some(row) = rows.next().await {
|
||||
match row {
|
||||
Ok(r) => work.push(r),
|
||||
Err(e) => {
|
||||
tracing::warn!("Error fetching blob row during migration: {}", e);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Process in parallel chunks.
|
||||
let results = futures::stream::iter(work.into_iter().map(|(hash, size)| {
|
||||
let src = source.clone();
|
||||
let tgt = target.clone();
|
||||
let st = state.clone();
|
||||
async move {
|
||||
// Check if we should keep running.
|
||||
{
|
||||
let s = st.read().await;
|
||||
if s.status != MigrationStatus::Running {
|
||||
return;
|
||||
}
|
||||
}
|
||||
|
||||
// Skip if already in target.
|
||||
match tgt.blob_exists(&hash).await {
|
||||
Ok(true) => {
|
||||
let mut s = st.write().await;
|
||||
s.migrated_blobs += 1;
|
||||
s.migrated_bytes += size as u64;
|
||||
return;
|
||||
}
|
||||
Ok(false) => {}
|
||||
Err(e) => {
|
||||
tracing::warn!("blob_exists check failed for {}: {}", hash, e);
|
||||
}
|
||||
}
|
||||
|
||||
// Copy: stream from source → temp file → put into target.
|
||||
if let Err(e) = copy_blob(&src, &tgt, &hash).await {
|
||||
tracing::warn!("Failed to migrate blob {}: {}", hash, e);
|
||||
let mut s = st.write().await;
|
||||
s.failed_blobs.push(hash);
|
||||
return;
|
||||
}
|
||||
|
||||
let mut s = st.write().await;
|
||||
s.migrated_blobs += 1;
|
||||
s.migrated_bytes += size as u64;
|
||||
}
|
||||
}))
|
||||
.buffer_unordered(concurrency)
|
||||
.collect::<Vec<()>>()
|
||||
.await;
|
||||
|
||||
drop(results);
|
||||
|
||||
// Finalize state.
|
||||
let mut s = state.write().await;
|
||||
if s.status == MigrationStatus::Running {
|
||||
if s.failed_blobs.is_empty() {
|
||||
s.status = MigrationStatus::Completed;
|
||||
} else {
|
||||
s.status = MigrationStatus::Failed;
|
||||
}
|
||||
s.completed_at = Some(chrono::Utc::now());
|
||||
}
|
||||
|
||||
tracing::info!(
|
||||
"Migration finished: {}/{} blobs, {} failures",
|
||||
s.migrated_blobs,
|
||||
s.total_blobs,
|
||||
s.failed_blobs.len()
|
||||
);
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Copy a single blob: stream from source → spool to temp file → put_blob into target.
|
||||
async fn copy_blob(
|
||||
source: &Arc<dyn BlobStorageBackend>,
|
||||
target: &Arc<dyn BlobStorageBackend>,
|
||||
hash: &str,
|
||||
) -> Result<(), DomainError> {
|
||||
use tokio::io::AsyncWriteExt;
|
||||
|
||||
// Create a temp file to spool content.
|
||||
let tmp_dir = std::env::temp_dir().join("oxicloud-migration");
|
||||
tokio::fs::create_dir_all(&tmp_dir).await.map_err(|e| {
|
||||
DomainError::internal_error("Migration", format!("Failed to create temp dir: {}", e))
|
||||
})?;
|
||||
|
||||
let tmp_path = tmp_dir.join(format!("{}.tmp", hash));
|
||||
|
||||
// Stream from source.
|
||||
let stream = source.get_blob_stream(hash).await?;
|
||||
|
||||
// Write to temp file.
|
||||
let mut file = tokio::fs::File::create(&tmp_path).await.map_err(|e| {
|
||||
DomainError::internal_error("Migration", format!("Failed to create temp file: {}", e))
|
||||
})?;
|
||||
|
||||
let mut stream = std::pin::pin!(stream);
|
||||
while let Some(chunk) = stream.next().await {
|
||||
let bytes = chunk.map_err(|e| {
|
||||
DomainError::internal_error("Migration", format!("Stream error: {}", e))
|
||||
})?;
|
||||
file.write_all(&bytes)
|
||||
.await
|
||||
.map_err(|e| DomainError::internal_error("Migration", format!("Write error: {}", e)))?;
|
||||
}
|
||||
file.flush()
|
||||
.await
|
||||
.map_err(|e| DomainError::internal_error("Migration", format!("Flush error: {}", e)))?;
|
||||
drop(file);
|
||||
|
||||
// Put into target.
|
||||
target.put_blob(hash, &tmp_path).await?;
|
||||
|
||||
// Clean up temp file.
|
||||
let _ = tokio::fs::remove_file(&tmp_path).await;
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Verify migration integrity by comparing blob counts and sampling random hashes.
|
||||
pub async fn verify_migration(
|
||||
target: Arc<dyn BlobStorageBackend>,
|
||||
pool: Arc<PgPool>,
|
||||
sample_size: usize,
|
||||
) -> Result<VerificationResult, DomainError> {
|
||||
// 1. Count blobs in PG.
|
||||
let pg_count: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM storage.blobs")
|
||||
.fetch_one(pool.as_ref())
|
||||
.await
|
||||
.unwrap_or(0);
|
||||
|
||||
// 2. Verify sample of blobs exist in target.
|
||||
let sample_rows: Vec<(String, i64)> =
|
||||
sqlx::query_as("SELECT hash, size FROM storage.blobs ORDER BY random() LIMIT $1")
|
||||
.bind(sample_size as i64)
|
||||
.fetch_all(pool.as_ref())
|
||||
.await
|
||||
.map_err(|e| {
|
||||
DomainError::internal_error("Migration", format!("Sample query failed: {}", e))
|
||||
})?;
|
||||
|
||||
let mut missing = Vec::new();
|
||||
let mut size_mismatches = Vec::new();
|
||||
|
||||
for (hash, expected_size) in &sample_rows {
|
||||
match target.blob_exists(hash).await {
|
||||
Ok(false) => missing.push(hash.clone()),
|
||||
Err(e) => {
|
||||
tracing::warn!("blob_exists failed for {}: {}", hash, e);
|
||||
missing.push(hash.clone());
|
||||
}
|
||||
Ok(true) => {
|
||||
// Verify size matches.
|
||||
if let Ok(actual_size) = target.blob_size(hash).await
|
||||
&& actual_size != *expected_size as u64
|
||||
{
|
||||
size_mismatches.push(hash.clone());
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
let passed = missing.is_empty() && size_mismatches.is_empty();
|
||||
|
||||
Ok(VerificationResult {
|
||||
pg_blob_count: pg_count as u64,
|
||||
sample_checked: sample_rows.len() as u64,
|
||||
missing_in_target: missing,
|
||||
size_mismatches,
|
||||
passed,
|
||||
})
|
||||
}
|
||||
|
||||
/// Result of a post-migration integrity check.
|
||||
#[derive(Debug, Clone, Serialize, serde::Deserialize)]
|
||||
pub struct VerificationResult {
|
||||
pub pg_blob_count: u64,
|
||||
pub sample_checked: u64,
|
||||
pub missing_in_target: Vec<String>,
|
||||
pub size_mismatches: Vec<String>,
|
||||
pub passed: bool,
|
||||
}
|
||||
@@ -25,8 +25,6 @@ pub mod local_blob_backend;
|
||||
pub mod local_fs_mount_provider;
|
||||
pub mod login_lockout_service;
|
||||
pub mod media_metadata_service;
|
||||
pub mod migration_blob_backend;
|
||||
pub mod migration_job;
|
||||
pub mod mock_email_sender;
|
||||
pub mod mount_provider_factory;
|
||||
pub mod nextcloud_chunked_upload_service;
|
||||
@@ -46,6 +44,7 @@ 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 thumbnail_service;
|
||||
#[cfg(test)]
|
||||
mod thumbnail_service_test;
|
||||
|
||||
@@ -0,0 +1,486 @@
|
||||
//! Storage-backend migration as a recoverable-run tenant (Part 2 engine).
|
||||
//!
|
||||
//! Iterates `storage.blobs` and copies each byte payload from the SOURCE
|
||||
//! backend (whatever the app booted with) to the TARGET backend
|
||||
//! (whatever the current admin storage-settings config describes).
|
||||
//! Both legacy whole-file blobs AND CDC chunk blobs are covered by the
|
||||
//! single walk — they share `storage.blobs` as their physical registry
|
||||
//! (see memory `project_cdc_dual_storage_registries`).
|
||||
//! `storage.chunk_manifests` is pure PG state, holds no backend bytes,
|
||||
//! and needs no migration.
|
||||
//!
|
||||
//! Retires the in-memory `Arc<RwLock<MigrationState>>` + one-shot
|
||||
//! `tokio::spawn` in `migration_job.rs`. The recoverable engine
|
||||
//! provides cursor persistence, cooperative cancel, boot-time crash
|
||||
//! recovery, and the uniform `/api/admin/jobs/*` admin surface.
|
||||
//!
|
||||
//! ### Restart survival
|
||||
//!
|
||||
//! Cursor + per-blob failure findings are persisted after every batch.
|
||||
//! On restart the boot-time sweep flips any abandoned `Running` row to
|
||||
//! `Paused`; a subsequent admin trigger resumes from the persisted
|
||||
//! cursor via `run_or_resume`. At most one batch of already-copied
|
||||
//! blobs replays, and the `target.blob_exists` short-circuit makes
|
||||
//! even that replay effectively free.
|
||||
//!
|
||||
//! ### Design notes
|
||||
//!
|
||||
//! * **Cursor.** UTF-8 hex of the last-processed blob hash (64 chars).
|
||||
//! Natural lex order matches `ORDER BY hash ASC`. Same encoding
|
||||
//! `blobs_consistency` uses.
|
||||
//! * **Target resolution.** Rebuilt at the START of every fresh or
|
||||
//! resumed run via `StorageSettingsService::build_effective_backend`.
|
||||
//! Held for the duration of the run; a mid-run settings change is
|
||||
//! ignored until the next run. On resume the admin may have paused
|
||||
//! *specifically* to fix a broken target config, so we re-derive
|
||||
//! rather than pin.
|
||||
//! * **Per-blob failures don't fail the run.** Each failure records a
|
||||
//! `migration_failed` finding (severity `data_loss` — the bytes
|
||||
//! didn't cross) and the walk continues. A run that completes with
|
||||
//! zero findings is proof the target has every blob.
|
||||
//! * **Skip already-present blobs.** `target.blob_exists(hash)` before
|
||||
//! the copy — makes cheap re-runs safe and lets a paused run resume
|
||||
//! without redoing bytes.
|
||||
//! * **No per-batch concurrency knob.** The old code buffered N copies
|
||||
//! in parallel. Sequential is easier to reason about with cooperative
|
||||
//! cancel + cursor discipline; the batch loop is I/O-bound anyway.
|
||||
//! Add concurrency later if a real throughput need appears.
|
||||
|
||||
use std::path::Path;
|
||||
use std::sync::Arc;
|
||||
|
||||
use async_trait::async_trait;
|
||||
use futures::StreamExt;
|
||||
use sqlx::PgPool;
|
||||
|
||||
use crate::application::ports::blob_storage_ports::BlobStorageBackend;
|
||||
use crate::application::services::storage_settings_service::StorageSettingsService;
|
||||
use crate::common::errors::DomainError;
|
||||
use crate::infrastructure::scheduler::{
|
||||
JobRegistry, JobRunArgs, JobStore, JobStoreProvider, RecoverableJobHandler, RunOutcome,
|
||||
RunStatus, record_or_log,
|
||||
};
|
||||
|
||||
pub const STORAGE_MIGRATION_JOB_NAME: &str = "storage_migration";
|
||||
|
||||
/// Rows per batch. Copies are I/O-bound (source read + target write);
|
||||
/// larger batches amortise fewer SQL round-trips but the checkpoint
|
||||
/// / cancel-poll cadence lengthens. 100 balances the two — one
|
||||
/// checkpoint per ~hundred blobs is fine, and the cancel-poll comes
|
||||
/// every 100 rows too. Match `blobs_consistency` for consistency.
|
||||
const BATCH_SIZE: i64 = 100;
|
||||
|
||||
pub struct StorageMigrationService {
|
||||
pool: Arc<PgPool>,
|
||||
source: Arc<dyn BlobStorageBackend>,
|
||||
storage_settings: Arc<StorageSettingsService>,
|
||||
}
|
||||
|
||||
impl StorageMigrationService {
|
||||
pub fn new(
|
||||
pool: Arc<PgPool>,
|
||||
source: Arc<dyn BlobStorageBackend>,
|
||||
storage_settings: Arc<StorageSettingsService>,
|
||||
) -> Self {
|
||||
Self {
|
||||
pool,
|
||||
source,
|
||||
storage_settings,
|
||||
}
|
||||
}
|
||||
|
||||
/// Chainable self-registration — mirrors the `*_consistency`
|
||||
/// tenants. 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 StorageMigrationService {
|
||||
fn name(&self) -> &str {
|
||||
STORAGE_MIGRATION_JOB_NAME
|
||||
}
|
||||
|
||||
/// Definitive count — one row per blob. `SELECT COUNT(*) FROM
|
||||
/// storage.blobs` on a modern PG is a sub-second index-only scan
|
||||
/// even at millions of 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::migration",
|
||||
event = "storage_migration.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 {
|
||||
// No-op guard — refuse when the effective (target) config
|
||||
// points at the same physical storage as the source (boot
|
||||
// config). Without this, a misclick on an S3 deployment
|
||||
// issues one HEAD per blob for zero useful work — cheap on
|
||||
// local, expensive on remote. Same-type-different-location
|
||||
// migrations (local dir change, S3 bucket change) pass this
|
||||
// check and proceed normally.
|
||||
match self.storage_settings.is_source_target_identical().await {
|
||||
Ok(true) => {
|
||||
tracing::warn!(
|
||||
target: "audit",
|
||||
event = "storage_migration.refused_noop",
|
||||
run_id = %store.run_id(),
|
||||
"storage_migration refused: source and target point at the same storage"
|
||||
);
|
||||
return RunOutcome::Failed {
|
||||
message:
|
||||
"target equals source; change storage settings before triggering a migration"
|
||||
.to_string(),
|
||||
};
|
||||
}
|
||||
Ok(false) => {}
|
||||
Err(e) => {
|
||||
return RunOutcome::Failed {
|
||||
message: format!("identity check: {e}"),
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
// Resolve target at run start.
|
||||
let target = match self.storage_settings.build_effective_backend().await {
|
||||
Ok(t) => t,
|
||||
Err(e) => {
|
||||
return RunOutcome::Failed {
|
||||
message: format!("resolve target backend: {e}"),
|
||||
};
|
||||
}
|
||||
};
|
||||
if let Err(e) = target.initialize().await {
|
||||
return RunOutcome::Failed {
|
||||
message: format!("target backend init: {e}"),
|
||||
};
|
||||
}
|
||||
|
||||
let source_kind = self.source.backend_type();
|
||||
let target_kind = target.backend_type();
|
||||
tracing::info!(
|
||||
target: "audit",
|
||||
event = "storage_migration.run_started",
|
||||
run_id = %store.run_id(),
|
||||
source = source_kind,
|
||||
target = target_kind,
|
||||
resuming = resume_cursor.is_some(),
|
||||
"storage_migration starting {source_kind} → {target_kind}"
|
||||
);
|
||||
|
||||
// Cursor = the last-visited blob hash, UTF-8-encoded. On resume
|
||||
// walk `WHERE hash > $cursor`. `None` / empty = start from the
|
||||
// smallest hash. Same shape `blobs_consistency` uses.
|
||||
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) => {
|
||||
return RunOutcome::Failed {
|
||||
message: format!("invalid cursor: not valid UTF-8: {e}"),
|
||||
};
|
||||
}
|
||||
},
|
||||
};
|
||||
|
||||
let mut copied_count = 0u64;
|
||||
let mut skipped_count = 0u64;
|
||||
let mut failed_count = 0u64;
|
||||
let mut source_missing_count = 0u64;
|
||||
|
||||
loop {
|
||||
// Cooperative cancel poll between batches.
|
||||
match store.status().await {
|
||||
Ok(RunStatus::CancelRequested) => {
|
||||
tracing::info!(
|
||||
target: "oxicloud::migration",
|
||||
event = "storage_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"
|
||||
);
|
||||
return RunOutcome::Paused {
|
||||
cursor: cursor
|
||||
.as_ref()
|
||||
.map(|s| s.as_bytes().to_vec())
|
||||
.unwrap_or_default(),
|
||||
};
|
||||
}
|
||||
Ok(_) => {}
|
||||
Err(e) => {
|
||||
return RunOutcome::Failed {
|
||||
message: format!("status poll: {e}"),
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
// Fetch the next batch. `hash > $1` keyset pagination on
|
||||
// the PK; index-only scan.
|
||||
let rows: Vec<(String, i64)> = match sqlx::query_as(
|
||||
r#"
|
||||
SELECT hash, size
|
||||
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) => {
|
||||
return RunOutcome::Failed {
|
||||
message: format!("batch fetch: {e}"),
|
||||
};
|
||||
}
|
||||
};
|
||||
|
||||
if rows.is_empty() {
|
||||
tracing::info!(
|
||||
target: "oxicloud::migration",
|
||||
event = "storage_migration.completed",
|
||||
run_id = %store.run_id(),
|
||||
copied = copied_count,
|
||||
skipped = skipped_count,
|
||||
failed = failed_count,
|
||||
source_missing = source_missing_count,
|
||||
"storage_migration completed"
|
||||
);
|
||||
return RunOutcome::Completed;
|
||||
}
|
||||
|
||||
for (hash, size) in &rows {
|
||||
// Probe SOURCE first — without this a run would
|
||||
// silently "succeed" against a source that's missing
|
||||
// blobs the DB expects, and the audit intent of the
|
||||
// walk is lost (relevant on any post-migration state
|
||||
// where target may already have every blob). A
|
||||
// missing-on-source blob is a real data-loss
|
||||
// condition; record it and move on — we never
|
||||
// "copy" from nothing.
|
||||
match self.source.blob_exists(hash).await {
|
||||
Ok(true) => {}
|
||||
Ok(false) => {
|
||||
source_missing_count += 1;
|
||||
tracing::warn!(
|
||||
target: "oxicloud::migration",
|
||||
event = "storage_migration.source_missing",
|
||||
run_id = %store.run_id(),
|
||||
hash = %hash,
|
||||
source = source_kind,
|
||||
"blob absent from source; recording data-loss finding, no copy"
|
||||
);
|
||||
record_or_log(
|
||||
store,
|
||||
STORAGE_MIGRATION_JOB_NAME,
|
||||
"source_missing",
|
||||
"data_loss",
|
||||
None,
|
||||
serde_json::json!({
|
||||
"hash": hash,
|
||||
"size": size,
|
||||
"source": source_kind,
|
||||
"target": target_kind,
|
||||
}),
|
||||
)
|
||||
.await;
|
||||
continue;
|
||||
}
|
||||
Err(e) => {
|
||||
// Transient probe failure on source is NOT a
|
||||
// finding — treat like a network blip.
|
||||
// Skipping this row on this run; a re-run
|
||||
// will re-probe. If the failure is
|
||||
// persistent, `blobs_consistency` catches
|
||||
// it.
|
||||
tracing::warn!(
|
||||
target: "oxicloud::migration",
|
||||
event = "storage_migration.source_probe_error",
|
||||
run_id = %store.run_id(),
|
||||
hash = %hash,
|
||||
error = %e,
|
||||
"source blob_exists probe failed; skipping this row"
|
||||
);
|
||||
continue;
|
||||
}
|
||||
}
|
||||
|
||||
// Skip when the target already has it — supports
|
||||
// idempotent resume and cheap re-runs against a
|
||||
// partially-migrated target.
|
||||
match target.blob_exists(hash).await {
|
||||
Ok(true) => {
|
||||
skipped_count += 1;
|
||||
continue;
|
||||
}
|
||||
Ok(false) => {}
|
||||
Err(e) => {
|
||||
tracing::warn!(
|
||||
target: "oxicloud::migration",
|
||||
event = "storage_migration.blob_exists_error",
|
||||
run_id = %store.run_id(),
|
||||
hash = %hash,
|
||||
error = %e,
|
||||
"blob_exists probe on target failed; attempting copy anyway"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
match copy_blob(self.source.as_ref(), target.as_ref(), hash).await {
|
||||
Ok(()) => {
|
||||
copied_count += 1;
|
||||
}
|
||||
Err(e) => {
|
||||
failed_count += 1;
|
||||
tracing::warn!(
|
||||
target: "oxicloud::migration",
|
||||
event = "storage_migration.blob_failed",
|
||||
run_id = %store.run_id(),
|
||||
hash = %hash,
|
||||
error = %e,
|
||||
"failed to migrate blob; recording finding, continuing"
|
||||
);
|
||||
// resource_id stays None — blob hash isn't a
|
||||
// UUID. Real identifier lives in `detail.hash`
|
||||
// where the admin UI reads it.
|
||||
record_or_log(
|
||||
store,
|
||||
STORAGE_MIGRATION_JOB_NAME,
|
||||
"migration_failed",
|
||||
"data_loss",
|
||||
None,
|
||||
serde_json::json!({
|
||||
"hash": hash,
|
||||
"size": size,
|
||||
"source": source_kind,
|
||||
"target": target_kind,
|
||||
"error": e.to_string(),
|
||||
}),
|
||||
)
|
||||
.await;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Advance cursor + checkpoint. `delta_count` counts WORK
|
||||
// ATTEMPTED (copied + skipped + failed), not successful
|
||||
// copies alone — otherwise the progress bar stalls whenever
|
||||
// a batch is dominated by already-present blobs, which is
|
||||
// exactly the case on a resume.
|
||||
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 {
|
||||
return RunOutcome::Failed {
|
||||
message: format!("checkpoint: {e}"),
|
||||
};
|
||||
}
|
||||
|
||||
if (rows.len() as i64) < BATCH_SIZE {
|
||||
tracing::info!(
|
||||
target: "oxicloud::migration",
|
||||
event = "storage_migration.completed",
|
||||
run_id = %store.run_id(),
|
||||
copied = copied_count,
|
||||
skipped = skipped_count,
|
||||
failed = failed_count,
|
||||
source_missing = source_missing_count,
|
||||
"storage_migration completed"
|
||||
);
|
||||
return RunOutcome::Completed;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Copy one blob: stream source bytes to a temp file, then hand the
|
||||
/// path to `target.put_blob`. The spool-through-disk shape matches
|
||||
/// what the old `migration_job::copy_blob` did — some backends'
|
||||
/// `put_blob` want a path they can `rename(2)` or multi-part upload
|
||||
/// from, not an in-memory buffer. The temp file lives in
|
||||
/// `std::env::temp_dir()/oxicloud-migration/{hash}.tmp` and is
|
||||
/// removed on success (best-effort on the failure paths — the OS
|
||||
/// cleans up on reboot).
|
||||
async fn copy_blob(
|
||||
source: &dyn BlobStorageBackend,
|
||||
target: &dyn BlobStorageBackend,
|
||||
hash: &str,
|
||||
) -> Result<(), DomainError> {
|
||||
let tmp_dir = std::env::temp_dir().join("oxicloud-migration");
|
||||
tokio::fs::create_dir_all(&tmp_dir).await.map_err(|e| {
|
||||
DomainError::internal_error(
|
||||
"StorageMigration",
|
||||
format!("create temp dir {}: {e}", tmp_dir.display()),
|
||||
)
|
||||
})?;
|
||||
let tmp_path = tmp_dir.join(format!("{hash}.tmp"));
|
||||
|
||||
if let Err(e) = write_source_to_tmp(source, hash, &tmp_path).await {
|
||||
let _ = tokio::fs::remove_file(&tmp_path).await;
|
||||
return Err(e);
|
||||
}
|
||||
|
||||
let put_result = target.put_blob(hash, &tmp_path).await;
|
||||
let _ = tokio::fs::remove_file(&tmp_path).await;
|
||||
put_result.map(|_bytes_written| ())
|
||||
}
|
||||
|
||||
async fn write_source_to_tmp(
|
||||
source: &dyn BlobStorageBackend,
|
||||
hash: &str,
|
||||
tmp_path: &Path,
|
||||
) -> Result<(), DomainError> {
|
||||
use tokio::io::AsyncWriteExt;
|
||||
|
||||
let stream = source.get_blob_stream(hash).await?;
|
||||
let mut file = tokio::fs::File::create(tmp_path).await.map_err(|e| {
|
||||
DomainError::internal_error(
|
||||
"StorageMigration",
|
||||
format!("create temp file {}: {e}", tmp_path.display()),
|
||||
)
|
||||
})?;
|
||||
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}"))
|
||||
})?;
|
||||
file.write_all(&bytes).await.map_err(|e| {
|
||||
DomainError::internal_error("StorageMigration", format!("temp file write: {e}"))
|
||||
})?;
|
||||
}
|
||||
file.flush()
|
||||
.await
|
||||
.map_err(|e| DomainError::internal_error("StorageMigration", format!("temp flush: {e}")))?;
|
||||
Ok(())
|
||||
}
|
||||
Reference in New Issue
Block a user