perf(thumbnails): hold a decode permit before reading source blobs
Request-path thumbnail generation (REST get_thumbnail and NC preview) read the full source blob into RAM and decoded it with no concurrency bound — a first-view gallery of K images stacked K full-size buffers and K parallel decodes on the blocking pool. The background hook had the inverse ordering problem: it read the blob eagerly and only then queued on the decode semaphore, so N concurrent uploads held N originals in memory while waiting. All three paths now acquire the decode semaphore first and read the blob under the permit, capping peak RAM at permits x image size: - new ThumbnailService::get_thumbnail_from_blob defers the blob read into the moka init closure (cache and disk hits never touch the blob); both handlers use it and no longer pre-read. - get_thumbnail_from_bytes (direct-bytes variant) now also takes a permit before decoding; shared generate_and_persist core keeps the two entrypoints duplicate-free. - generate_all_sizes_background_from_blob (renamed from _from_bytes) checks blob existence and disk-dedup state, then acquires the permit, then reads; the redundant pre-read spawn wrapper in ThumbnailRefreshHook is gone. https://claude.ai/code/session_01QxwJDHqQhbMkHK333QtMme
This commit is contained in:
@@ -258,7 +258,10 @@ impl ThumbnailService {
|
|||||||
/// Get a thumbnail from raw image bytes, generating it if needed.
|
/// Get a thumbnail from raw image bytes, generating it if needed.
|
||||||
///
|
///
|
||||||
/// This is the storage-model-safe entrypoint for CDC/manifest-backed
|
/// This is the storage-model-safe entrypoint for CDC/manifest-backed
|
||||||
/// blobs where no single local source file exists on disk.
|
/// blobs where no single local source file exists on disk. Prefer
|
||||||
|
/// [`Self::get_thumbnail_from_blob`] on request paths — it defers the
|
||||||
|
/// full blob read until a decode permit is held, so a stampede of
|
||||||
|
/// cache misses cannot stack one source image per request in RAM.
|
||||||
pub async fn get_thumbnail_from_bytes(
|
pub async fn get_thumbnail_from_bytes(
|
||||||
&self,
|
&self,
|
||||||
file_id: &str,
|
file_id: &str,
|
||||||
@@ -287,30 +290,12 @@ impl ThumbnailService {
|
|||||||
return Bytes::from(data);
|
return Bytes::from(data);
|
||||||
}
|
}
|
||||||
|
|
||||||
tracing::info!("🎨 Generating thumbnail: {} {:?}", file_id_owned, size);
|
let Ok(_permit) = self.decode_semaphore.acquire().await else {
|
||||||
match Self::generate_thumbnail_from_data(
|
tracing::warn!("Decode semaphore closed, skipping {}", file_id_owned);
|
||||||
original_data,
|
return Bytes::new();
|
||||||
size,
|
};
|
||||||
self.generation_timeout,
|
self.generate_and_persist(&file_id_owned, &thumb_path, size, original_data)
|
||||||
)
|
|
||||||
.await
|
.await
|
||||||
{
|
|
||||||
Ok(bytes) => {
|
|
||||||
if let Some(parent) = thumb_path.parent() {
|
|
||||||
let _ = fs::create_dir_all(parent).await;
|
|
||||||
}
|
|
||||||
let _ = fs::write(&thumb_path, &bytes).await;
|
|
||||||
bytes
|
|
||||||
}
|
|
||||||
Err(e) => {
|
|
||||||
tracing::warn!(
|
|
||||||
"Thumbnail generation failed for {} {:?}: {e}",
|
|
||||||
file_id_owned,
|
|
||||||
size
|
|
||||||
);
|
|
||||||
Bytes::new()
|
|
||||||
}
|
|
||||||
}
|
|
||||||
})
|
})
|
||||||
.await;
|
.await;
|
||||||
|
|
||||||
@@ -325,6 +310,106 @@ impl ThumbnailService {
|
|||||||
Ok(bytes)
|
Ok(bytes)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Get a thumbnail for a content-addressed blob, generating it if needed.
|
||||||
|
///
|
||||||
|
/// Request-path entrypoint: on a memory+disk cache miss the source blob
|
||||||
|
/// is read **after** a decode permit is acquired, so peak RAM under a
|
||||||
|
/// thumbnail stampede is `permits × image size` instead of
|
||||||
|
/// `in-flight requests × image size`. moka's per-key init additionally
|
||||||
|
/// collapses concurrent requests for the same thumbnail into one read.
|
||||||
|
pub async fn get_thumbnail_from_blob(
|
||||||
|
&self,
|
||||||
|
file_id: &str,
|
||||||
|
blob_hash: &str,
|
||||||
|
size: ThumbnailSize,
|
||||||
|
dedup: Arc<DedupService>,
|
||||||
|
) -> Result<Bytes, ThumbnailError> {
|
||||||
|
let cache_key = ThumbnailCacheKey {
|
||||||
|
file_id: file_id.to_string(),
|
||||||
|
size,
|
||||||
|
};
|
||||||
|
|
||||||
|
let thumb_path = self.get_thumbnail_path(blob_hash, size);
|
||||||
|
let file_id_owned = file_id.to_string();
|
||||||
|
let blob_hash_owned = blob_hash.to_string();
|
||||||
|
|
||||||
|
let entry = self
|
||||||
|
.cache
|
||||||
|
.entry(cache_key)
|
||||||
|
.or_insert_with(async move {
|
||||||
|
if let Ok(data) = fs::read(&thumb_path).await {
|
||||||
|
tracing::debug!(
|
||||||
|
"💾 Thumbnail loaded from disk: {} {:?}",
|
||||||
|
file_id_owned,
|
||||||
|
size
|
||||||
|
);
|
||||||
|
return Bytes::from(data);
|
||||||
|
}
|
||||||
|
|
||||||
|
let Ok(_permit) = self.decode_semaphore.acquire().await else {
|
||||||
|
tracing::warn!("Decode semaphore closed, skipping {}", file_id_owned);
|
||||||
|
return Bytes::new();
|
||||||
|
};
|
||||||
|
let original_data = match dedup.read_blob_bytes(&blob_hash_owned).await {
|
||||||
|
Ok(bytes) => bytes,
|
||||||
|
Err(e) => {
|
||||||
|
tracing::warn!(
|
||||||
|
"Failed to read blob for thumbnail {} {:?}: {e}",
|
||||||
|
file_id_owned,
|
||||||
|
size
|
||||||
|
);
|
||||||
|
return Bytes::new();
|
||||||
|
}
|
||||||
|
};
|
||||||
|
self.generate_and_persist(&file_id_owned, &thumb_path, size, original_data)
|
||||||
|
.await
|
||||||
|
})
|
||||||
|
.await;
|
||||||
|
|
||||||
|
let bytes = entry.into_value();
|
||||||
|
if bytes.is_empty() {
|
||||||
|
return Err(ThumbnailError::ImageError(
|
||||||
|
"Thumbnail generation failed".to_string(),
|
||||||
|
));
|
||||||
|
}
|
||||||
|
|
||||||
|
tracing::debug!("🔥 Thumbnail served: {} {:?}", file_id, size);
|
||||||
|
Ok(bytes)
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Decode `original_data` into one thumbnail size, persist it to its
|
||||||
|
/// blob-keyed disk path, and return the encoded bytes — empty `Bytes`
|
||||||
|
/// on failure (moka's zero-weight negative-entry convention).
|
||||||
|
///
|
||||||
|
/// Callers must hold a `decode_semaphore` permit.
|
||||||
|
async fn generate_and_persist(
|
||||||
|
&self,
|
||||||
|
file_id: &str,
|
||||||
|
thumb_path: &Path,
|
||||||
|
size: ThumbnailSize,
|
||||||
|
original_data: Bytes,
|
||||||
|
) -> Bytes {
|
||||||
|
tracing::info!("🎨 Generating thumbnail: {} {:?}", file_id, size);
|
||||||
|
match Self::generate_thumbnail_from_data(original_data, size, self.generation_timeout).await
|
||||||
|
{
|
||||||
|
Ok(bytes) => {
|
||||||
|
if let Some(parent) = thumb_path.parent() {
|
||||||
|
let _ = fs::create_dir_all(parent).await;
|
||||||
|
}
|
||||||
|
let _ = fs::write(&thumb_path, &bytes).await;
|
||||||
|
bytes
|
||||||
|
}
|
||||||
|
Err(e) => {
|
||||||
|
tracing::warn!(
|
||||||
|
"Thumbnail generation failed for {} {:?}: {e}",
|
||||||
|
file_id,
|
||||||
|
size
|
||||||
|
);
|
||||||
|
Bytes::new()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
/// Try to serve a thumbnail from cache only (memory → disk).
|
/// Try to serve a thumbnail from cache only (memory → disk).
|
||||||
///
|
///
|
||||||
/// Unlike `get_thumbnail`, this does **not** generate a new thumbnail.
|
/// Unlike `get_thumbnail`, this does **not** generate a new thumbnail.
|
||||||
@@ -823,15 +908,16 @@ impl ThumbnailService {
|
|||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Generate all thumbnail sizes in the background from raw image bytes.
|
/// Generate all thumbnail sizes in the background for a content-addressed
|
||||||
|
/// blob (CDC/manifest-safe — no physical source file required).
|
||||||
///
|
///
|
||||||
/// This is compatible with CDC/manifest-backed blobs because it does not
|
/// The source blob is read **after** the decode permit is acquired, so N
|
||||||
/// require a single physical source file on disk.
|
/// concurrent uploads queue as N small tasks, not N full images in RAM:
|
||||||
pub fn generate_all_sizes_background_from_bytes(
|
/// peak memory is `permits × image size` regardless of upload concurrency.
|
||||||
|
pub fn generate_all_sizes_background_from_blob(
|
||||||
self: Arc<Self>,
|
self: Arc<Self>,
|
||||||
file_id: String,
|
file_id: String,
|
||||||
blob_hash: String,
|
blob_hash: String,
|
||||||
original_data: Bytes,
|
|
||||||
dedup: Arc<DedupService>,
|
dedup: Arc<DedupService>,
|
||||||
) {
|
) {
|
||||||
tokio::spawn(async move {
|
tokio::spawn(async move {
|
||||||
@@ -889,6 +975,20 @@ impl ThumbnailService {
|
|||||||
}
|
}
|
||||||
};
|
};
|
||||||
|
|
||||||
|
// Read the source only now that a permit bounds how many of
|
||||||
|
// these full-image buffers can exist at once.
|
||||||
|
let original_data = match dedup.read_blob_bytes(&blob_hash).await {
|
||||||
|
Ok(bytes) => bytes,
|
||||||
|
Err(e) => {
|
||||||
|
tracing::warn!(
|
||||||
|
"Failed to read blob for thumbnail generation {}: {}",
|
||||||
|
file_id,
|
||||||
|
e
|
||||||
|
);
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
let results = tokio::task::spawn_blocking(move || {
|
let results = tokio::task::spawn_blocking(move || {
|
||||||
Self::render_all_thumbnails_from_data(original_data.as_ref())
|
Self::render_all_thumbnails_from_data(original_data.as_ref())
|
||||||
})
|
})
|
||||||
@@ -1013,11 +1113,12 @@ impl crate::application::ports::file_lifecycle::FileLifecycleHook for ThumbnailR
|
|||||||
if !is_new_blob || !ThumbnailService::is_supported_image(content_type) {
|
if !is_new_blob || !ThumbnailService::is_supported_image(content_type) {
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
Self::spawn_thumbnail_generation(
|
self.thumbnail
|
||||||
self.thumbnail.clone(),
|
.clone()
|
||||||
self.dedup.clone(),
|
.generate_all_sizes_background_from_blob(
|
||||||
file_id.to_string(),
|
file_id.to_string(),
|
||||||
blob_hash.to_string(),
|
blob_hash.to_string(),
|
||||||
|
self.dedup.clone(),
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -1047,7 +1148,7 @@ impl crate::application::ports::file_lifecycle::FileLifecycleHook for ThumbnailR
|
|||||||
e
|
e
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
Self::spawn_thumbnail_generation(thumbnail, dedup, file_id, blob_hash);
|
thumbnail.generate_all_sizes_background_from_blob(file_id, blob_hash, dedup);
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -1066,30 +1167,6 @@ impl crate::application::ports::file_lifecycle::FileLifecycleHook for ThumbnailR
|
|||||||
// to avoid a circular Arc: DedupService→BlobLifecycleService→ThumbnailRefreshHook→DedupService.
|
// to avoid a circular Arc: DedupService→BlobLifecycleService→ThumbnailRefreshHook→DedupService.
|
||||||
// ThumbnailService does not hold DedupService so no cycle exists.
|
// ThumbnailService does not hold DedupService so no cycle exists.
|
||||||
|
|
||||||
impl ThumbnailRefreshHook {
|
|
||||||
fn spawn_thumbnail_generation(
|
|
||||||
ts: Arc<ThumbnailService>,
|
|
||||||
ds: Arc<DedupService>,
|
|
||||||
file_id: String,
|
|
||||||
hash: String,
|
|
||||||
) {
|
|
||||||
tokio::spawn(async move {
|
|
||||||
match ds.read_blob_bytes(&hash).await {
|
|
||||||
Ok(bytes) => {
|
|
||||||
ts.generate_all_sizes_background_from_bytes(file_id, hash, bytes, ds.clone());
|
|
||||||
}
|
|
||||||
Err(e) => {
|
|
||||||
tracing::warn!(
|
|
||||||
"Failed to read blob for thumbnail generation {}: {}",
|
|
||||||
file_id,
|
|
||||||
e
|
|
||||||
);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
});
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// ─── BlobLifecycleHook ───────────────────────────────────────────────────────
|
// ─── BlobLifecycleHook ───────────────────────────────────────────────────────
|
||||||
|
|
||||||
impl crate::application::ports::blob_lifecycle::BlobLifecycleHook for ThumbnailService {
|
impl crate::application::ports::blob_lifecycle::BlobLifecycleHook for ThumbnailService {
|
||||||
|
|||||||
@@ -442,19 +442,13 @@ impl FileHandler {
|
|||||||
.into_response();
|
.into_response();
|
||||||
}
|
}
|
||||||
|
|
||||||
let original_bytes = match state.core.dedup_service.read_blob_bytes(&blob_hash).await {
|
|
||||||
Ok(bytes) => bytes,
|
|
||||||
Err(err) => {
|
|
||||||
return AppError::internal_error(format!(
|
|
||||||
"Failed to load source image for thumbnail generation: {}",
|
|
||||||
err
|
|
||||||
))
|
|
||||||
.into_response();
|
|
||||||
}
|
|
||||||
};
|
|
||||||
|
|
||||||
match thumbnail_service
|
match thumbnail_service
|
||||||
.get_thumbnail_from_bytes(&id, &blob_hash, thumb_size.into(), original_bytes)
|
.get_thumbnail_from_blob(
|
||||||
|
&id,
|
||||||
|
&blob_hash,
|
||||||
|
thumb_size.into(),
|
||||||
|
state.core.dedup_service.clone(),
|
||||||
|
)
|
||||||
.await
|
.await
|
||||||
{
|
{
|
||||||
Ok(data) => Response::builder()
|
Ok(data) => Response::builder()
|
||||||
|
|||||||
@@ -157,26 +157,18 @@ pub async fn handle_preview(
|
|||||||
.unwrap();
|
.unwrap();
|
||||||
}
|
}
|
||||||
|
|
||||||
let original_bytes = match state.core.dedup_service.read_blob_bytes(&blob_hash).await {
|
// Generate/get thumbnail — the blob is read inside the service once a
|
||||||
Ok(bytes) => bytes,
|
// decode permit is held, so preview stampedes cannot stack source
|
||||||
Err(err) => {
|
// images in RAM.
|
||||||
tracing::error!(
|
|
||||||
"Failed to load source image for preview {}: {}",
|
|
||||||
object_id,
|
|
||||||
err
|
|
||||||
);
|
|
||||||
return Response::builder()
|
|
||||||
.status(StatusCode::INTERNAL_SERVER_ERROR)
|
|
||||||
.body(Body::from("Failed to load preview source"))
|
|
||||||
.unwrap();
|
|
||||||
}
|
|
||||||
};
|
|
||||||
|
|
||||||
// Generate/get thumbnail
|
|
||||||
match state
|
match state
|
||||||
.core
|
.core
|
||||||
.thumbnail_service
|
.thumbnail_service
|
||||||
.get_thumbnail_from_bytes(&object_id, &blob_hash, thumb_size.into(), original_bytes)
|
.get_thumbnail_from_blob(
|
||||||
|
&object_id,
|
||||||
|
&blob_hash,
|
||||||
|
thumb_size.into(),
|
||||||
|
state.core.dedup_service.clone(),
|
||||||
|
)
|
||||||
.await
|
.await
|
||||||
{
|
{
|
||||||
Ok(data) => {
|
Ok(data) => {
|
||||||
|
|||||||
Reference in New Issue
Block a user