Files
Oxicloud/src/application/ports/blob_storage_ports.rs
T
Claude 5b8b740233 perf(blob): small read-ahead for the local backend (read_prefetch 1 -> 2)
The local backend inherited the trait's conservative read_prefetch() = 1
(strictly sequential chunk reassembly), while S3/Azure already use 8. The
prior rationale was that concurrent opens over scattered content-addressed
chunk files turn one sequential read into competing random I/O ("slower
cold"). Benchmarked that assumption with examples/bench_blob_prefetch
(sweeps the buffered(N) depth over a real LocalBlobBackend under disk-bound
vs network-bound consumers and warm vs cold page cache).

Result on SSD-class storage (median MB/s vs N=1):
  warm  disk-bound   N=2 +11.8%   N=8 +3.9%   N=16 -4.4%
  cold  disk-bound   N=2  +7.2%   (no cold regression on SSD)
  network-bound (throttled)  ~0% at any N — the socket, not the disk, caps it

So N=8 is wrong for local (leaves gain on the table, risks HDD seek thrash)
and the network-bound win the analysis assumed doesn't materialize: buffered()
here overlaps the per-chunk File::open (cheap on local disk), not the data
read. N=2 captures most of the disk-bound gain — which covers localhost/LAN
downloads AND the internal blob reads that drain as fast as the disk delivers
(thumbnail render, transcode, ZIP export, content extraction) — at the lowest
fan-out. Env-tunable via OXICLOUD_LOCAL_READ_PREFETCH (set 1 on seek-bound
HDDs to restore the old behaviour; raise on fast NVMe). Signature unchanged,
so all ~16 LocalBlobBackend::new call sites are untouched.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01JG5yYZ9s868mJwqT2Qz7ez
2026-06-22 08:10:46 +00:00

149 lines
6.8 KiB
Rust
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
//! Blob Storage Backend Port — abstracts raw byte I/O for content-addressable storage.
//!
//! This trait decouples `DedupService` from any specific storage medium.
//! Implementations include:
//! - `LocalBlobBackend` — local filesystem (default)
//! - `S3BlobBackend` — any S3-compatible service (AWS, Backblaze B2, MinIO, R2…)
//!
//! `DedupService` owns an `Arc<dyn BlobStorageBackend>` and delegates all
//! byte-level I/O through this trait, keeping BLAKE3 hashing, ref-counting
//! and PostgreSQL index logic in `DedupService` itself.
use bytes::Bytes;
use futures::Stream;
use serde::Serialize;
use std::future::Future;
use std::path::{Path, PathBuf};
use std::pin::Pin;
use crate::domain::errors::DomainError;
/// Boxed future alias used by [`BlobStorageBackend`] to keep the trait dyn-compatible.
type BoxFut<'a, T> = Pin<Box<dyn Future<Output = T> + Send + 'a>>;
/// Pinned boxed byte stream — the return type for blob reads.
pub type BlobStream = Pin<Box<dyn Stream<Item = Result<Bytes, std::io::Error>> + Send>>;
/// Health-check result returned by [`BlobStorageBackend::health_check`].
#[derive(Debug, Clone, Serialize)]
pub struct StorageHealthStatus {
/// Whether the backend is reachable and functional.
pub connected: bool,
/// Human-readable backend identifier (e.g. `"local"`, `"s3"`).
pub backend_type: String,
/// Descriptive status message.
pub message: String,
/// Available space in bytes, if the backend can report it.
pub available_bytes: Option<u64>,
}
/// Minimal trait for blob byte I/O — decoupled from dedup logic.
///
/// Every method operates on a *hash key* that uniquely identifies a blob.
/// The backend is responsible for mapping the hash to its own addressing
/// scheme (filesystem path, S3 key, etc.).
///
/// Returns boxed futures so the trait is dyn-compatible (`Arc<dyn BlobStorageBackend>`).
pub trait BlobStorageBackend: Send + Sync + 'static {
/// Perform any one-time setup (create directories, verify bucket, etc.).
fn initialize(&self) -> BoxFut<'_, Result<(), DomainError>>;
/// Store a blob from a local temporary file.
///
/// Must be **idempotent**: if the blob already exists the call succeeds
/// without overwriting. Returns the number of bytes stored.
fn put_blob(&self, hash: &str, source_path: &Path) -> BoxFut<'_, Result<u64, DomainError>>;
/// Store a blob from in-memory bytes (used by CDC chunk storage).
///
/// Must be **idempotent**: if the blob already exists the call succeeds
/// without overwriting. Returns the number of bytes stored.
fn put_blob_from_bytes(&self, hash: &str, data: Bytes) -> BoxFut<'_, Result<u64, DomainError>>;
/// Store a blob from in-memory bytes **without forcing durability**.
///
/// Same idempotency contract as [`Self::put_blob_from_bytes`], but the
/// bytes may still sit in volatile caches (e.g. the OS page cache) when
/// the future resolves. Durability is only guaranteed after a subsequent
/// [`Self::sync_blobs`] covering this hash returns `Ok`. Callers MUST NOT
/// record a durable reference to the blob (e.g. a PostgreSQL row) before
/// that sync completes.
///
/// Default: delegates to `put_blob_from_bytes` (immediately durable),
/// pairing with the no-op `sync_blobs` default so backends that don't
/// opt in keep today's per-write durability semantics.
fn put_blob_from_bytes_unsynced(
&self,
hash: &str,
data: Bytes,
) -> BoxFut<'_, Result<u64, DomainError>> {
self.put_blob_from_bytes(hash, data)
}
/// Make previously written blobs durable in one batched operation.
///
/// Durability barrier for blobs written via `put_blob_from_bytes_unsynced`:
/// when this returns `Ok`, every listed blob is crash-safe. Local
/// filesystem backends fsync each listed blob file plus each distinct
/// parent directory once — one sweep per upload instead of two fsyncs
/// per chunk. Remote object stores are durable on PUT, so the default
/// is a no-op.
fn sync_blobs(&self, _hashes: &[String]) -> BoxFut<'_, Result<(), DomainError>> {
Box::pin(async { Ok(()) })
}
/// Stream the full blob content in chunks.
fn get_blob_stream(&self, hash: &str) -> BoxFut<'_, Result<BlobStream, DomainError>>;
/// Stream the byte range `[start, end)` of the blob (for HTTP Range
/// requests / video seek). `end` is **exclusive**; `None` means "to the
/// end of the blob". Callers translating inclusive HTTP Range headers
/// must pass `last_byte + 1`.
fn get_blob_range_stream(
&self,
hash: &str,
start: u64,
end: Option<u64>,
) -> BoxFut<'_, Result<BlobStream, DomainError>>;
/// Delete a blob by hash. Must be **idempotent** (no error if already gone).
fn delete_blob(&self, hash: &str) -> BoxFut<'_, Result<(), DomainError>>;
/// Check if a blob exists in the backend.
fn blob_exists(&self, hash: &str) -> BoxFut<'_, Result<bool, DomainError>>;
/// Get blob size in bytes without downloading content.
fn blob_size(&self, hash: &str) -> BoxFut<'_, Result<u64, DomainError>>;
/// Verify connectivity and permissions (used by the admin "Test Connection" button).
fn health_check(&self) -> BoxFut<'_, Result<StorageHealthStatus, DomainError>>;
/// Return the backend type name for display (e.g. `"local"`, `"s3"`).
fn backend_type(&self) -> &'static str;
/// Return the local filesystem path for a blob, if available.
///
/// Only meaningful for local-filesystem backends. Remote backends
/// return `None`; callers that need a local file must stream + spool.
fn local_blob_path(&self, hash: &str) -> Option<PathBuf>;
/// How many chunk fetches the CDC reader may run concurrently when
/// reassembling a file (`read_blob_stream`'s `buffered(N)` read-ahead).
///
/// The trait default is a conservative **1** (strictly sequential) — the
/// safe fallback for any backend that doesn't know its own I/O profile.
/// Concrete backends override it:
/// • Local disk → a small benchmarked depth (default `2`, env-tunable):
/// overlapping the next chunk's `File::open` with the current chunk's
/// drain measured +7–12% on disk-bound reads with no cold regression on
/// SSDs; deeper queues buy little (the data read, not the open, is then
/// the cost) and risk competing random I/O over scattered content-
/// addressed chunk files on seek-bound HDDs. See `LocalBlobBackend`.
/// • Remote (S3/Azure) → `8`: there the dominant cost is per-chunk request
/// latency, and overlapping fetches hides it (≈ N× faster reassembly).
/// Wrapping backends delegate to the backend that actually serves the bytes.
fn read_prefetch(&self) -> usize {
1
}
}