mirror of
https://github.com/hcengineering/platform.git
synced 2026-09-12 12:47:45 +02:00
+105
-5
@@ -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<Vec<u8>>,
|
||||
}
|
||||
|
||||
// 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<Blob, ApiError> {
|
||||
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<Blob, ApiError> {
|
||||
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::<Header<ContentLength>>().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<Blob, ApiError> {
|
||||
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<Blob, ApiError> {
|
||||
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,
|
||||
})
|
||||
}
|
||||
|
||||
+11
-4
@@ -129,7 +129,7 @@ pub async fn put(request: HttpRequest, payload: Payload) -> HandlerResult<HttpRe
|
||||
}
|
||||
}
|
||||
|
||||
let uploaded = upload(&s3, &pool, payload).await?;
|
||||
let uploaded = upload(&s3, &pool, &mut request, payload).await?;
|
||||
|
||||
let part_data = PartData {
|
||||
workspace: path.workspace,
|
||||
@@ -142,7 +142,14 @@ pub async fn put(request: HttpRequest, payload: Payload) -> HandlerResult<HttpRe
|
||||
meta: Some(meta.into_iter().collect()),
|
||||
};
|
||||
|
||||
postgres::set_part(&pool, path.workspace, &part_data.key, None, &part_data).await?;
|
||||
postgres::set_part(
|
||||
&pool,
|
||||
path.workspace,
|
||||
&part_data.key,
|
||||
uploaded.inline,
|
||||
&part_data,
|
||||
)
|
||||
.await?;
|
||||
|
||||
let mut response = HttpResponse::Created();
|
||||
response.insert_header((header::CONTENT_LOCATION, part_data.key));
|
||||
@@ -168,7 +175,7 @@ pub async fn post(request: HttpRequest, payload: Payload) -> HandlerResult<HttpR
|
||||
let pool = request.app_data::<Data<Pool>>().unwrap().to_owned();
|
||||
let s3 = request.app_data::<Data<S3Client>>().unwrap().to_owned();
|
||||
|
||||
let uploaded = upload(&s3, &pool, payload).await?;
|
||||
let uploaded = upload(&s3, &pool, &mut request, payload).await?;
|
||||
|
||||
let parts = postgres::find_parts::<PartData>(&pool, path.workspace, &path.key).await?;
|
||||
|
||||
@@ -199,7 +206,7 @@ pub async fn post(request: HttpRequest, payload: Payload) -> HandlerResult<HttpR
|
||||
path.workspace,
|
||||
&part_data.key,
|
||||
part_data.part,
|
||||
None,
|
||||
uploaded.inline,
|
||||
&part_data,
|
||||
)
|
||||
.await?;
|
||||
|
||||
Reference in New Issue
Block a user