From 908b8f4d4b8e583ce3e764c74a65d250b8adf056 Mon Sep 17 00:00:00 2001 From: Claude Date: Wed, 10 Jun 2026 09:52:47 +0000 Subject: [PATCH] perf(upload): one durability barrier per file instead of two fsyncs per chunk MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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 --- src/application/ports/blob_storage_ports.rs | 30 ++++ src/infrastructure/services/dedup_service.rs | 57 +++++--- .../services/encrypted_blob_backend.rs | 35 +++++ .../services/local_blob_backend.rs | 135 +++++++++++++----- 4 files changed, 204 insertions(+), 53 deletions(-) diff --git a/src/application/ports/blob_storage_ports.rs b/src/application/ports/blob_storage_ports.rs index 391fe164..f7f7092c 100644 --- a/src/application/ports/blob_storage_ports.rs +++ b/src/application/ports/blob_storage_ports.rs @@ -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>; + /// 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> { + 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>; diff --git a/src/infrastructure/services/dedup_service.rs b/src/infrastructure/services/dedup_service.rs index a916d60d..53e09834 100644 --- a/src/infrastructure/services/dedup_service.rs +++ b/src/infrastructure/services/dedup_service.rs @@ -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 = new_ops.iter().map(|(h, _, _)| h.clone()).collect(); + let new_sizes: Vec = 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> = 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 = chunks.iter().map(|c| c.hash.clone()).collect(); diff --git a/src/infrastructure/services/encrypted_blob_backend.rs b/src/infrastructure/services/encrypted_blob_backend.rs index 47a3997a..f2620c2e 100644 --- a/src/infrastructure/services/encrypted_blob_backend.rs +++ b/src/infrastructure/services/encrypted_blob_backend.rs @@ -129,6 +129,41 @@ impl BlobStorageBackend for EncryptedBlobBackend { }) } + fn put_blob_from_bytes_unsynced( + &self, + hash: &str, + data: Bytes, + ) -> Pin> + 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> + Send + '_>> { + // Pure passthrough — encryption does not change blob addressing. + self.inner.sync_blobs(hashes) + } + fn get_blob_stream( &self, hash: &str, diff --git a/src/infrastructure/services/local_blob_backend.rs b/src/infrastructure/services/local_blob_backend.rs index 926bbbc8..7bf57c85 100644 --- a/src/infrastructure/services/local_blob_backend.rs +++ b/src/infrastructure/services/local_blob_backend.rs @@ -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 { + 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> + 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> + Send + '_>> { + let paths: Vec = 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 = 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> = 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(()) }) }