feat(storage migration): make logs more verbose
This commit is contained in:
@@ -62,7 +62,7 @@ use crate::infrastructure::scheduler::{
|
|||||||
JobRegistry, JobRunArgs, JobStore, JobStoreProvider, RecoverableJobHandler, RunOutcome,
|
JobRegistry, JobRunArgs, JobStore, JobStoreProvider, RecoverableJobHandler, RunOutcome,
|
||||||
RunStatus, record_or_log,
|
RunStatus, record_or_log,
|
||||||
};
|
};
|
||||||
use crate::infrastructure::services::encrypted_blob_backend::EncryptedBlobBackend;
|
use crate::infrastructure::services::encrypted_blob_backend::{EncryptedBlobBackend, HeadCheck};
|
||||||
use crate::infrastructure::services::entry_backend::{
|
use crate::infrastructure::services::entry_backend::{
|
||||||
build_entry_backend_typed, persist_active_backend_name, persist_migration_readonly,
|
build_entry_backend_typed, persist_active_backend_name, persist_migration_readonly,
|
||||||
};
|
};
|
||||||
@@ -666,7 +666,8 @@ impl RecoverableJobHandler for BackendMigrationService {
|
|||||||
// both correct (checks the header, not just
|
// both correct (checks the header, not just
|
||||||
// existence) AND fast (skips the ~99% of resume
|
// existence) AND fast (skips the ~99% of resume
|
||||||
// blobs already at head format).
|
// blobs already at head format).
|
||||||
if target.is_at_head_format(hash).await.unwrap_or(false) {
|
let head_check = target.head_check(hash).await;
|
||||||
|
if matches!(head_check, HeadCheck::Match) {
|
||||||
skipped_count += 1;
|
skipped_count += 1;
|
||||||
tracing::debug!(
|
tracing::debug!(
|
||||||
target: "oxicloud::migration",
|
target: "oxicloud::migration",
|
||||||
@@ -682,6 +683,33 @@ impl RecoverableJobHandler for BackendMigrationService {
|
|||||||
match copy_blob(self.source.as_ref(), target.clone(), hash).await {
|
match copy_blob(self.source.as_ref(), target.clone(), hash).await {
|
||||||
Ok(()) => {
|
Ok(()) => {
|
||||||
copied_count += 1;
|
copied_count += 1;
|
||||||
|
// Log the concrete action: overwrite (blob
|
||||||
|
// existed with wrong format — the K1.2 repair
|
||||||
|
// case) vs fresh write (blob absent). Info
|
||||||
|
// level for both so a single log tail shows
|
||||||
|
// operators exactly what happened per blob.
|
||||||
|
match head_check {
|
||||||
|
HeadCheck::Mismatch(prev) => tracing::info!(
|
||||||
|
target: "oxicloud::migration",
|
||||||
|
event = "backend_migration.blob_overwritten",
|
||||||
|
run_id = %store.run_id(),
|
||||||
|
hash = %hash,
|
||||||
|
previous_format = %prev,
|
||||||
|
new_format = %target.head_format(),
|
||||||
|
"🔄 target blob existed with different format — overwritten"
|
||||||
|
),
|
||||||
|
HeadCheck::Absent => tracing::info!(
|
||||||
|
target: "oxicloud::migration",
|
||||||
|
event = "backend_migration.blob_written",
|
||||||
|
run_id = %store.run_id(),
|
||||||
|
hash = %hash,
|
||||||
|
new_format = %target.head_format(),
|
||||||
|
"✍️ fresh blob written to target"
|
||||||
|
),
|
||||||
|
// Unreachable in practice — we already
|
||||||
|
// early-`continue`d on Match above.
|
||||||
|
HeadCheck::Match => {}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
Err(e) => {
|
Err(e) => {
|
||||||
failed_count += 1;
|
failed_count += 1;
|
||||||
|
|||||||
@@ -237,46 +237,49 @@ impl EncryptedBlobBackend {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Smart-skip probe used by `backend_migration` (and potentially
|
/// Richer variant of [`Self::is_at_head_format`] — returns not
|
||||||
/// `backend_rotate` if it ever gains a fast-path). Reads the
|
/// just "matches head?" but also *why* it doesn't match, so
|
||||||
/// first [`HEADER_SIZE`] bytes of the on-disk blob and returns
|
/// callers (currently `backend_migration`) can log the concrete
|
||||||
/// `true` iff:
|
/// action they're about to take: **skip** (match), **overwrite**
|
||||||
///
|
/// (blob exists but with wrong format/key), or **fresh write**
|
||||||
/// * the blob exists, AND
|
/// (blob absent). Same one-round-trip cost as the boolean version.
|
||||||
/// * its first bytes carry the v1 header (`OXCPT | 0001`), AND
|
|
||||||
/// * the header's `key_fp` matches this wrapper's head key
|
|
||||||
/// (encrypted-v1 → `head_key_fp`; plaintext-v1 → all-zero).
|
|
||||||
///
|
|
||||||
/// Any error (blob missing, short read, IO failure) returns
|
|
||||||
/// `Ok(false)` — the caller re-writes, which is always safe.
|
|
||||||
/// Only propagates errors that the caller genuinely cannot
|
|
||||||
/// distinguish from "blob absent" and needs to see.
|
|
||||||
///
|
///
|
||||||
/// Backend-agnostic: reads through the trait's
|
/// Backend-agnostic: reads through the trait's
|
||||||
/// `get_blob_range_stream` on the *inner* backend (bypasses this
|
/// `get_blob_range_stream` on the *inner* backend (bypasses this
|
||||||
/// wrapper's decrypt so we see the raw on-disk header bytes).
|
/// wrapper's decrypt so we see the raw on-disk header bytes).
|
||||||
/// Local pays one `pread` syscall; S3 pays one HEAD/GET with
|
/// Local pays one `pread` syscall; S3 pays one HEAD/GET with
|
||||||
/// `Range: bytes=0-14`; Azure the same.
|
/// `Range: bytes=0-14`; Azure the same.
|
||||||
pub async fn is_at_head_format(&self, hash: &str) -> Result<bool, DomainError> {
|
pub async fn head_check(&self, hash: &str) -> HeadCheck {
|
||||||
// Skip the probe entirely for backends that don't stream —
|
|
||||||
// `get_blob_range_stream` would fail on a legit-missing blob
|
|
||||||
// with `NotFound` which we want to translate to `Ok(false)`.
|
|
||||||
let stream = match self
|
let stream = match self
|
||||||
.inner
|
.inner
|
||||||
.get_blob_range_stream(hash, 0, Some(HEADER_SIZE as u64))
|
.get_blob_range_stream(hash, 0, Some(HEADER_SIZE as u64))
|
||||||
.await
|
.await
|
||||||
{
|
{
|
||||||
Ok(s) => s,
|
Ok(s) => s,
|
||||||
Err(_) => return Ok(false),
|
Err(_) => return HeadCheck::Absent,
|
||||||
};
|
};
|
||||||
let raw = match collect_stream(stream).await {
|
let raw = match collect_stream(stream).await {
|
||||||
Ok(b) => b,
|
Ok(b) => b,
|
||||||
Err(_) => return Ok(false),
|
Err(_) => return HeadCheck::Absent,
|
||||||
};
|
};
|
||||||
if raw.len() < HEADER_SIZE {
|
if raw.is_empty() {
|
||||||
return Ok(false);
|
return HeadCheck::Absent;
|
||||||
}
|
}
|
||||||
Ok(BlobFormat::classify(&raw) == self.head_format())
|
// Bytes present but short-header — treat as "exists, needs
|
||||||
|
// rewrite" (malformed blob, classify falls back to Legacy).
|
||||||
|
let current = BlobFormat::classify(&raw);
|
||||||
|
if current == self.head_format() {
|
||||||
|
HeadCheck::Match
|
||||||
|
} else {
|
||||||
|
HeadCheck::Mismatch(current)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Boolean convenience wrapper over [`Self::head_check`] for
|
||||||
|
/// callers that only need "matches head?" — retained for API
|
||||||
|
/// symmetry with the original design.
|
||||||
|
pub async fn is_at_head_format(&self, hash: &str) -> Result<bool, DomainError> {
|
||||||
|
Ok(matches!(self.head_check(hash).await, HeadCheck::Match))
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Fetch, classify, and decrypt a blob in one round-trip. Used by
|
/// Fetch, classify, and decrypt a blob in one round-trip. Used by
|
||||||
@@ -366,6 +369,28 @@ impl BlobFormat {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Result of [`EncryptedBlobBackend::head_check`] — describes the
|
||||||
|
/// exact state of a target blob relative to the wrapper's current
|
||||||
|
/// head format, so callers can log/act with precision.
|
||||||
|
///
|
||||||
|
/// [`HeadCheck::Absent`] and [`HeadCheck::Mismatch`] both require a
|
||||||
|
/// write; the distinction is purely observability — operators seeing
|
||||||
|
/// `Mismatch` in migration logs know the K1.2 residue is being
|
||||||
|
/// repaired, while `Absent` is a plain first-time copy.
|
||||||
|
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||||
|
pub enum HeadCheck {
|
||||||
|
/// Blob exists AND its header bytes match the wrapper's
|
||||||
|
/// head format exactly — skip the write.
|
||||||
|
Match,
|
||||||
|
/// Blob exists but header differs (legacy shape, different
|
||||||
|
/// `key_fp`, plaintext vs encrypted, or malformed shorter than
|
||||||
|
/// [`HEADER_SIZE`]). The current on-disk shape is reported so
|
||||||
|
/// the caller can name it in the log.
|
||||||
|
Mismatch(BlobFormat),
|
||||||
|
/// No bytes at that hash — fresh write, not an overwrite.
|
||||||
|
Absent,
|
||||||
|
}
|
||||||
|
|
||||||
impl std::fmt::Display for BlobFormat {
|
impl std::fmt::Display for BlobFormat {
|
||||||
/// Human-friendly format for audit logs + finding details.
|
/// Human-friendly format for audit logs + finding details.
|
||||||
/// Renders `key_fp` as SSH-style colon-hex (e.g.
|
/// Renders `key_fp` as SSH-style colon-hex (e.g.
|
||||||
|
|||||||
Reference in New Issue
Block a user