perf(upload): one durability barrier per file instead of two fsyncs per chunk
The CDC chunk path paid `sync_all` + parent-dir fsync + one ref-count upsert round-trip PER ~256 KB chunk: a new 1 GB file ≈ 8 000 fsyncs + 4 000 sequential INSERTs on the upload critical path (seconds on SSD, tens of seconds on HDD). - New `put_blob_from_bytes_unsynced` + `sync_blobs` on BlobStorageBackend. Defaults delegate to the durable variants, so remote backends (S3/Azure — durable on PUT) are untouched. LocalBlobBackend writes chunks without fsync and `sync_blobs` then fsyncs all files concurrently (kernel coalesces the writeback) plus one fsync per DISTINCT shard directory — at most 256 — instead of one per chunk. EncryptedBlobBackend delegates both so encryption over the local backend keeps the optimization. - store_chunks: uploads use the unsynced variant; one `sync_blobs` barrier runs before anything references the chunks, and the per-chunk ref-count upserts collapse into a single `INSERT ... SELECT unnest(...) ON CONFLICT` statement. Durability contract is unchanged: every chunk is on stable storage before the manifest row that references it commits. A crash mid-upload leaves chunk files without DB rows — the same orphan class the per-chunk scheme already produced, just a wider window. https://claude.ai/code/session_01Dp3oWon5GBMVn4j3QXZdgx
This commit is contained in:
@@ -60,6 +60,36 @@ pub trait BlobStorageBackend: Send + Sync + 'static {
|
||||
/// 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 a durability barrier**.
|
||||
///
|
||||
/// The CDC chunk path writes thousands of small chunks per file; paying
|
||||
/// two fsyncs per chunk (file + parent dir) put ~8 000 fsyncs on the
|
||||
/// critical path of a 1 GB upload. Callers using this MUST issue one
|
||||
/// [`Self::sync_blobs`] barrier over the written hashes before
|
||||
/// persisting any record that references them (the chunk manifest).
|
||||
///
|
||||
/// Default: delegates to [`Self::put_blob_from_bytes`] — correct for
|
||||
/// remote backends (S3, Azure) where a successful PUT is already
|
||||
/// durable and the "unsynced" notion does not exist.
|
||||
fn put_blob_from_bytes_unsynced(
|
||||
&self,
|
||||
hash: &str,
|
||||
data: Bytes,
|
||||
) -> BoxFut<'_, Result<u64, DomainError>> {
|
||||
self.put_blob_from_bytes(hash, data)
|
||||
}
|
||||
|
||||
/// Durability barrier for blobs previously written with
|
||||
/// [`Self::put_blob_from_bytes_unsynced`]: when this resolves, the
|
||||
/// listed blobs and their directory entries survive a power loss.
|
||||
///
|
||||
/// Default: no-op — matches the default `put_blob_from_bytes_unsynced`,
|
||||
/// which is already durable on completion.
|
||||
fn sync_blobs(&self, hashes: &[String]) -> BoxFut<'_, Result<(), DomainError>> {
|
||||
let _ = hashes;
|
||||
Box::pin(async { Ok(()) })
|
||||
}
|
||||
|
||||
/// Stream the full blob content in chunks.
|
||||
fn get_blob_stream(&self, hash: &str) -> BoxFut<'_, Result<BlobStream, DomainError>>;
|
||||
|
||||
|
||||
@@ -535,15 +535,24 @@ impl DedupService {
|
||||
.filter(|c| !existing_hashes.contains(&c.hash))
|
||||
.map(|c| (c.hash.clone(), c.offset as u64, c.length))
|
||||
.collect();
|
||||
let new_hashes: Vec<String> = new_ops.iter().map(|(h, _, _)| h.clone()).collect();
|
||||
let new_sizes: Vec<i64> = new_ops.iter().map(|(_, _, l)| *l as i64).collect();
|
||||
|
||||
let source = Arc::new(std::fs::File::open(source_path).map_err(|e| {
|
||||
DomainError::internal_error("Dedup", format!("Failed to open source file: {}", e))
|
||||
})?);
|
||||
|
||||
// Chunks are written WITHOUT a per-chunk durability barrier — the
|
||||
// former two fsyncs per chunk put ~8 000 fsyncs on the critical path
|
||||
// of a 1 GB upload. One `sync_blobs` barrier below makes every chunk
|
||||
// durable before the manifest (the only record referencing them) is
|
||||
// committed, so the durability contract observed by the caller is
|
||||
// unchanged. A crash before the barrier leaves orphan chunk files
|
||||
// with no DB row — the same orphan class the per-chunk scheme had,
|
||||
// just a wider window.
|
||||
let results: Vec<Result<(), DomainError>> = stream::iter(new_ops)
|
||||
.map(|(hash, offset, length)| {
|
||||
let source = source.clone();
|
||||
let pool = pool.clone();
|
||||
let backend = backend.clone();
|
||||
async move {
|
||||
// Positioned read of just this chunk (≤ CDC_MAX_CHUNK) off
|
||||
@@ -563,27 +572,8 @@ impl DedupService {
|
||||
})?;
|
||||
|
||||
backend
|
||||
.put_blob_from_bytes(&hash, Bytes::from(bytes))
|
||||
.put_blob_from_bytes_unsynced(&hash, Bytes::from(bytes))
|
||||
.await?;
|
||||
// ON CONFLICT covers a concurrent uploader inserting the
|
||||
// same brand-new chunk between the existence check above and
|
||||
// this INSERT.
|
||||
sqlx::query(
|
||||
"INSERT INTO storage.blobs (hash, size, ref_count)
|
||||
VALUES ($1, $2, 1)
|
||||
ON CONFLICT (hash) DO UPDATE
|
||||
SET ref_count = storage.blobs.ref_count + 1",
|
||||
)
|
||||
.bind(&hash)
|
||||
.bind(length as i64)
|
||||
.execute(pool.as_ref())
|
||||
.await
|
||||
.map_err(|e| {
|
||||
DomainError::internal_error(
|
||||
"Dedup",
|
||||
format!("Failed to upsert chunk: {}", e),
|
||||
)
|
||||
})?;
|
||||
Ok(())
|
||||
}
|
||||
})
|
||||
@@ -595,6 +585,31 @@ impl DedupService {
|
||||
result?;
|
||||
}
|
||||
|
||||
if !new_hashes.is_empty() {
|
||||
// Durability barrier: every new chunk (and its dirent) is on
|
||||
// stable storage before any DB row references it.
|
||||
backend.sync_blobs(&new_hashes).await?;
|
||||
|
||||
// One batched upsert for all new chunks (was one round-trip per
|
||||
// chunk, interleaved with the uploads). ON CONFLICT covers a
|
||||
// concurrent uploader inserting the same brand-new chunk between
|
||||
// the existence check above and this INSERT.
|
||||
sqlx::query(
|
||||
"INSERT INTO storage.blobs (hash, size, ref_count)
|
||||
SELECT t.hash, t.size, 1
|
||||
FROM unnest($1::text[], $2::bigint[]) AS t(hash, size)
|
||||
ON CONFLICT (hash) DO UPDATE
|
||||
SET ref_count = storage.blobs.ref_count + 1",
|
||||
)
|
||||
.bind(&new_hashes)
|
||||
.bind(&new_sizes)
|
||||
.execute(pool.as_ref())
|
||||
.await
|
||||
.map_err(|e| {
|
||||
DomainError::internal_error("Dedup", format!("Failed to upsert chunks: {}", e))
|
||||
})?;
|
||||
}
|
||||
|
||||
// chunk_hashes/chunk_sizes keep the full per-occurrence CDC sequence —
|
||||
// the manifest needs every chunk, in order, to reassemble the file.
|
||||
let chunk_hashes: Vec<String> = chunks.iter().map(|c| c.hash.clone()).collect();
|
||||
|
||||
@@ -129,6 +129,41 @@ impl BlobStorageBackend for EncryptedBlobBackend {
|
||||
})
|
||||
}
|
||||
|
||||
fn put_blob_from_bytes_unsynced(
|
||||
&self,
|
||||
hash: &str,
|
||||
data: Bytes,
|
||||
) -> Pin<Box<dyn std::future::Future<Output = Result<u64, DomainError>> + Send + '_>> {
|
||||
let inner = self.inner.clone();
|
||||
let hash = hash.to_string();
|
||||
let cipher = self.cipher.clone();
|
||||
Box::pin(async move {
|
||||
// Encrypt in memory: nonce || ciphertext (includes GCM tag)
|
||||
let nonce = Aes256Gcm::generate_nonce(&mut OsRng);
|
||||
let ciphertext = cipher.encrypt(&nonce, data.as_ref()).map_err(|e| {
|
||||
DomainError::internal_error("Encryption", format!("encrypt failed: {e}"))
|
||||
})?;
|
||||
|
||||
let mut encrypted = Vec::with_capacity(NONCE_SIZE + ciphertext.len());
|
||||
encrypted.extend_from_slice(nonce.as_slice());
|
||||
encrypted.extend_from_slice(&ciphertext);
|
||||
|
||||
// Delegate the relaxed-durability variant so encryption over
|
||||
// the local backend keeps the batched `sync_blobs` barrier.
|
||||
inner
|
||||
.put_blob_from_bytes_unsynced(&hash, Bytes::from(encrypted))
|
||||
.await
|
||||
})
|
||||
}
|
||||
|
||||
fn sync_blobs(
|
||||
&self,
|
||||
hashes: &[String],
|
||||
) -> Pin<Box<dyn std::future::Future<Output = Result<(), DomainError>> + Send + '_>> {
|
||||
// Pure passthrough — encryption does not change blob addressing.
|
||||
self.inner.sync_blobs(hashes)
|
||||
}
|
||||
|
||||
fn get_blob_stream(
|
||||
&self,
|
||||
hash: &str,
|
||||
|
||||
@@ -36,11 +36,16 @@ async fn fsync_parent_dir(child_path: &Path) {
|
||||
let Some(parent) = child_path.parent() else {
|
||||
return;
|
||||
};
|
||||
let parent = parent.to_owned();
|
||||
fsync_dir(parent).await;
|
||||
}
|
||||
|
||||
/// Fsync a directory (best-effort — see [`fsync_parent_dir`]).
|
||||
async fn fsync_dir(dir: &Path) {
|
||||
let dir_owned = dir.to_owned();
|
||||
// std::fs (synchronous) opens directories reliably on Linux/macOS;
|
||||
// do it on the blocking pool so we don't park the tokio worker.
|
||||
let result = tokio::task::spawn_blocking(move || -> std::io::Result<()> {
|
||||
let dir = std::fs::File::open(&parent)?;
|
||||
let dir = std::fs::File::open(&dir_owned)?;
|
||||
dir.sync_all()
|
||||
})
|
||||
.await;
|
||||
@@ -49,20 +54,35 @@ async fn fsync_parent_dir(child_path: &Path) {
|
||||
Ok(Err(e)) => {
|
||||
tracing::warn!(
|
||||
error = %e,
|
||||
path = %child_path.display(),
|
||||
"Blob parent-dir fsync failed (rename durability not guaranteed)"
|
||||
path = %dir.display(),
|
||||
"Blob dir fsync failed (rename durability not guaranteed)"
|
||||
);
|
||||
}
|
||||
Err(e) => {
|
||||
tracing::warn!(
|
||||
error = %e,
|
||||
path = %child_path.display(),
|
||||
"Blob parent-dir fsync task join failed"
|
||||
path = %dir.display(),
|
||||
"Blob dir fsync task join failed"
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Open + fsync a blob file so its content survives a power loss.
|
||||
async fn fsync_blob_file(path: &Path) -> Result<(), DomainError> {
|
||||
let file = fs::File::open(path).await.map_err(|e| {
|
||||
DomainError::internal_error("Blob", format!("Failed to open blob for fsync: {}", e))
|
||||
})?;
|
||||
file.sync_all().await.map_err(|e| {
|
||||
DomainError::internal_error("Blob", format!("Failed to fsync blob file: {}", e))
|
||||
})?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Concurrent fsyncs in flight during a [`BlobStorageBackend::sync_blobs`]
|
||||
/// barrier — lets the kernel coalesce writeback across many small chunks.
|
||||
const SYNC_BLOBS_CONCURRENCY: usize = 16;
|
||||
|
||||
/// Chunk size for streaming file reads (256 KB).
|
||||
const STREAM_CHUNK_SIZE: usize = 256 * 1024;
|
||||
|
||||
@@ -117,6 +137,28 @@ impl LocalBlobBackend {
|
||||
pub fn blob_root(&self) -> &Path {
|
||||
&self.blob_root
|
||||
}
|
||||
|
||||
/// Write `data` to `blob_path` unless it already exists (idempotent
|
||||
/// dedup-skip). Returns `true` when a new file was written. No
|
||||
/// durability barrier here — callers choose between the per-file
|
||||
/// fsync (`put_blob_from_bytes`) and the batched `sync_blobs` barrier
|
||||
/// (`put_blob_from_bytes_unsynced`).
|
||||
async fn write_blob_bytes(&self, blob_path: &Path, data: &Bytes) -> Result<bool, DomainError> {
|
||||
if fs::try_exists(blob_path).await.unwrap_or(false) {
|
||||
return Ok(false);
|
||||
}
|
||||
|
||||
// `fs::write` is `create + write_all + close` — but the close on
|
||||
// tokio::fs::File does NOT fsync, so durability is layered on
|
||||
// explicitly by the caller.
|
||||
let mut file = fs::File::create(blob_path).await.map_err(|e| {
|
||||
DomainError::internal_error("Blob", format!("Failed to create blob file: {}", e))
|
||||
})?;
|
||||
file.write_all(data).await.map_err(|e| {
|
||||
DomainError::internal_error("Blob", format!("Failed to write blob from bytes: {}", e))
|
||||
})?;
|
||||
Ok(true)
|
||||
}
|
||||
}
|
||||
|
||||
impl BlobStorageBackend for LocalBlobBackend {
|
||||
@@ -220,36 +262,65 @@ impl BlobStorageBackend for LocalBlobBackend {
|
||||
let hash = hash.to_owned();
|
||||
Box::pin(async move {
|
||||
let blob_path = self.blob_path(&hash);
|
||||
let size = data.len() as u64;
|
||||
if self.write_blob_bytes(&blob_path, &data).await? {
|
||||
// Same durability story as `put_blob`: the blob file is
|
||||
// fsync'd before the parent directory is, so both the
|
||||
// content and the dirent creation survive a power loss in
|
||||
// the same step.
|
||||
fsync_blob_file(&blob_path).await?;
|
||||
fsync_parent_dir(&blob_path).await;
|
||||
}
|
||||
Ok(data.len() as u64)
|
||||
})
|
||||
}
|
||||
|
||||
// Idempotent: if blob already exists, skip
|
||||
if fs::try_exists(&blob_path).await.unwrap_or(false) {
|
||||
return Ok(size);
|
||||
fn put_blob_from_bytes_unsynced(
|
||||
&self,
|
||||
hash: &str,
|
||||
data: Bytes,
|
||||
) -> Pin<Box<dyn std::future::Future<Output = Result<u64, DomainError>> + Send + '_>> {
|
||||
let hash = hash.to_owned();
|
||||
Box::pin(async move {
|
||||
let blob_path = self.blob_path(&hash);
|
||||
// No fsync here — the CDC chunk path issues one `sync_blobs`
|
||||
// barrier over all written chunks before the manifest row that
|
||||
// references them is committed.
|
||||
self.write_blob_bytes(&blob_path, &data).await?;
|
||||
Ok(data.len() as u64)
|
||||
})
|
||||
}
|
||||
|
||||
fn sync_blobs(
|
||||
&self,
|
||||
hashes: &[String],
|
||||
) -> Pin<Box<dyn std::future::Future<Output = Result<(), DomainError>> + Send + '_>> {
|
||||
let paths: Vec<PathBuf> = hashes.iter().map(|h| self.blob_path(h)).collect();
|
||||
// One fsync per DISTINCT parent directory — the hash-sharded layout
|
||||
// spreads chunks over at most 256 dirs, so this replaces the former
|
||||
// one-dir-fsync-per-chunk. Best-effort, same as `put_blob`.
|
||||
let parents: std::collections::HashSet<PathBuf> = paths
|
||||
.iter()
|
||||
.filter_map(|p| p.parent().map(Path::to_owned))
|
||||
.collect();
|
||||
Box::pin(async move {
|
||||
// Fsync every blob file concurrently — the kernel coalesces
|
||||
// batched writeback far better than an fsync after every write
|
||||
// on the upload critical path.
|
||||
use futures::stream::{self, StreamExt};
|
||||
let results: Vec<Result<(), DomainError>> = stream::iter(paths)
|
||||
.map(|path| async move { fsync_blob_file(&path).await })
|
||||
.buffer_unordered(SYNC_BLOBS_CONCURRENCY)
|
||||
.collect()
|
||||
.await;
|
||||
for result in results {
|
||||
result?;
|
||||
}
|
||||
|
||||
// Write directly to blob path. `fs::write` is `create +
|
||||
// write_all + close` — but the close on tokio::fs::File
|
||||
// does NOT fsync, so we open explicitly to keep the
|
||||
// `sync_all` call site obvious. Same durability story as
|
||||
// `put_blob`: the blob file is fsync'd before the parent
|
||||
// directory is, so both the content and the dirent
|
||||
// creation survive a power loss in the same step.
|
||||
let mut file = fs::File::create(&blob_path).await.map_err(|e| {
|
||||
DomainError::internal_error("Blob", format!("Failed to create blob file: {}", e))
|
||||
})?;
|
||||
file.write_all(&data).await.map_err(|e| {
|
||||
DomainError::internal_error(
|
||||
"Blob",
|
||||
format!("Failed to write blob from bytes: {}", e),
|
||||
)
|
||||
})?;
|
||||
file.sync_all().await.map_err(|e| {
|
||||
DomainError::internal_error("Blob", format!("Failed to fsync blob file: {}", e))
|
||||
})?;
|
||||
drop(file);
|
||||
fsync_parent_dir(&blob_path).await;
|
||||
for parent in &parents {
|
||||
fsync_dir(parent).await;
|
||||
}
|
||||
|
||||
Ok(size)
|
||||
Ok(())
|
||||
})
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user