From 6fbca2fedc1dcc9dcff26f82451e015769ff3911 Mon Sep 17 00:00:00 2001 From: Alexey Aristov Date: Thu, 28 Aug 2025 13:57:44 +0200 Subject: [PATCH] store small blobs inline Signed-off-by: Alexey Aristov --- src/blob.rs | 110 +++++++++++++++++++++++++++++++++++++++++++++--- src/handlers.rs | 15 +++++-- 2 files changed, 116 insertions(+), 9 deletions(-) diff --git a/src/blob.rs b/src/blob.rs index 573c8861ba..6e01a76e54 100644 --- a/src/blob.rs +++ b/src/blob.rs @@ -1,8 +1,12 @@ -use actix_web::web::Payload; +use actix_web::dev::ServiceRequest; +use actix_web::http::header::ContentLength; +use actix_web::web::{Header, Payload}; +use aws_sdk_s3::primitives::ByteStream; use aws_sdk_s3::types::{CompletedMultipartUpload, CompletedPart}; use blake3::Hasher; use bytes::BytesMut; use futures_util::StreamExt; +use size::Size; use tracing::*; use crate::s3::S3Client; @@ -16,17 +20,112 @@ use crate::handlers::{ApiError, HandlerResult}; pub struct Blob { pub s3_key: String, pub size: u64, + pub inline: Option>, } -// upload and deduplicate blob -#[instrument(level = "debug", skip_all, fields(s3_bucket, s3_key, upload))] -pub async fn upload(s3: &S3Client, pool: &Pool, mut payload: Payload) -> Result { +const MULTIPART_THRESHOLD: usize = 4; // mb +const INLINE_THRESHHOLD: usize = 100; // kb + +fn random_key() -> String { + ksuid::Ksuid::generate().to_base62() +} + +#[instrument(level = "debug", skip_all, fields(s3_bucket))] +pub async fn upload( + s3: &S3Client, + pool: &Pool, + request: &mut ServiceRequest, + payload: Payload, +) -> Result { let span = Span::current(); let s3_bucket = &CONFIG.s3_bucket; - let s3_key = ksuid::Ksuid::generate().to_base62(); span.record("s3_bucket", &s3_bucket); + + if let Ok(length) = request.extract::>().await + && length.0 < Size::from_megabytes(MULTIPART_THRESHOLD).bytes() as usize + { + upload_regular( + s3, + pool, + &s3_bucket, + length.0 < Size::from_kilobytes(INLINE_THRESHHOLD).bytes() as usize, + payload, + ) + .await + } else { + upload_multipart(s3, pool, &s3_bucket, payload).await + } +} + +#[instrument(level = "debug", skip_all, fields(s3_key))] +async fn upload_regular( + s3: &S3Client, + pool: &Pool, + s3_bucket: &str, + require_inline: bool, + payload: Payload, +) -> Result { + let span = Span::current(); + + let payload = payload + .to_bytes_limited(Size::from_megabytes(MULTIPART_THRESHOLD).bytes() as usize) + .await + .map_err(|_| actix_web::error::ErrorPayloadTooLarge("payload too large"))??; + + let size = payload.len() as u64; + let inline = if require_inline { + Some(payload.to_vec()) + } else { + None + }; + + let mut hash = Hasher::new(); + hash.update(&payload); + + let hash = hash.finalize().to_hex().to_string(); + + let s3_key = if let Some(s3_key_found) = postgres::find_blob_by_hash(&pool, &hash).await? { + span.record("s3_key", &s3_key_found); + debug!(s3_key_found, "blob deduplicated"); + s3_key_found + } else { + let s3_key = random_key(); + span.record("s3_key", &s3_key); + + s3.put_object() + .bucket(s3_bucket) + .key(&s3_key) + .body(ByteStream::from(payload)) + .send() + .await?; + + postgres::insert_blob(&pool, &s3_key, &hash).await?; + + debug!("blob created"); + + s3_key + }; + + Ok(Blob { + s3_key, + size, + inline, + }) +} + +#[instrument(level = "debug", skip_all, fields(upload, s3_key))] +async fn upload_multipart( + s3: &S3Client, + pool: &Pool, + s3_bucket: &str, + mut payload: Payload, +) -> Result { + let span = Span::current(); + + let s3_key = random_key(); + span.record("s3_key", &s3_key); let create_multipart = s3 @@ -139,5 +238,6 @@ pub async fn upload(s3: &S3Client, pool: &Pool, mut payload: Payload) -> Result< Ok(Blob { s3_key, size: total_uploaded as u64, + inline: None, }) } diff --git a/src/handlers.rs b/src/handlers.rs index cbb00af263..b7285b7112 100644 --- a/src/handlers.rs +++ b/src/handlers.rs @@ -129,7 +129,7 @@ pub async fn put(request: HttpRequest, payload: Payload) -> HandlerResult HandlerResult HandlerResult>().unwrap().to_owned(); let s3 = request.app_data::>().unwrap().to_owned(); - let uploaded = upload(&s3, &pool, payload).await?; + let uploaded = upload(&s3, &pool, &mut request, payload).await?; let parts = postgres::find_parts::(&pool, path.workspace, &path.key).await?; @@ -199,7 +206,7 @@ pub async fn post(request: HttpRequest, payload: Payload) -> HandlerResult