diff --git a/server/src/blob.rs b/server/src/blob.rs index a78ec0c2ff..772377dd8e 100644 --- a/server/src/blob.rs +++ b/server/src/blob.rs @@ -11,6 +11,7 @@ use tracing::*; use crate::handlers::ApiError; use crate::postgres::DbError; +use crate::recovery; use crate::s3::S3Client; use crate::{ config::CONFIG, @@ -20,6 +21,7 @@ use crate::{ #[derive(Debug)] pub struct Blob { pub s3_key: String, + pub hash: String, pub length: usize, pub inline: Option, pub parts_count: Option, @@ -66,7 +68,7 @@ where let buffer = buffer.freeze(); - let hash = hash.update(&buffer).finalize().to_hex(); + let hash = hash.update(&buffer).finalize().to_hex().to_string(); let length = buffer.len(); let inline = Some(buffer.clone()); @@ -108,6 +110,7 @@ where }; Blob { + hash, s3_key, length, inline, @@ -150,6 +153,7 @@ where }; Blob { + hash, s3_key, length: upload.length, inline: None, @@ -158,6 +162,10 @@ where } }; + if !blob.deduplicated { + recovery::set_blob(&s3, &blob.s3_key, &blob.hash).await?; + } + Ok(blob) } diff --git a/server/src/handlers.rs b/server/src/handlers.rs index 285f7e5d58..06ff317d20 100644 --- a/server/src/handlers.rs +++ b/server/src/handlers.rs @@ -17,13 +17,13 @@ use size::Size; use tracing::*; use uuid::Uuid; -use crate::merge::MergeStrategy; use crate::s3::S3Client; use crate::{blob, merge}; use crate::{ config::CONFIG, postgres::{self, Pool}, }; +use crate::{merge::MergeStrategy, recovery}; #[derive(Deserialize, Debug)] pub struct ObjectPath { @@ -212,6 +212,9 @@ pub async fn put(request: HttpRequest, payload: Payload) -> HandlerResult HandlerResult>(); + recovery::set_object(&s3, path.workspace, &part_data.key, obj_parts).await?; + postgres::append_part( &pool, path.workspace, diff --git a/server/src/main.rs b/server/src/main.rs index 370dd8b64c..29d0335696 100644 --- a/server/src/main.rs +++ b/server/src/main.rs @@ -18,6 +18,7 @@ mod blob; mod config; mod handlers; mod merge; +mod recovery; mod patch; mod postgres; mod s3; diff --git a/server/src/recovery.rs b/server/src/recovery.rs new file mode 100644 index 0000000000..0f0c807b32 --- /dev/null +++ b/server/src/recovery.rs @@ -0,0 +1,43 @@ +use bytes::Bytes; + +use crate::config::CONFIG; +use crate::{handlers::PartData, s3::S3Client}; + +pub async fn set_object( + s3: &S3Client, + workspace: uuid::Uuid, + key: &str, + parts: Vec<&PartData>, +) -> anyhow::Result<()> { + let s3_bucket = &CONFIG.s3_bucket; + + let key = format!("blob/{}/{}", workspace, key); + let body = Bytes::from(serde_json::to_string(&parts)?); + + s3.put_object() + .bucket(s3_bucket) + .key(key) + .body(body.into()) + .content_type("application/json") + .send() + .await?; + + Ok(()) +} + +pub async fn set_blob(s3: &S3Client, key: &str, hash: &str) -> anyhow::Result<()> { + let s3_bucket = &CONFIG.s3_bucket; + + let key = format!("hash/{}", key); + let body = Bytes::from(hash.to_string()); + + s3.put_object() + .bucket(s3_bucket) + .key(key) + .body(body.into()) + .content_type("text/plain") + .send() + .await?; + + Ok(()) +}