From dfc633dceece7c948c32cbe152d0273acc1ab216 Mon Sep 17 00:00:00 2001 From: weisd Date: Sun, 29 Dec 2024 17:15:43 +0800 Subject: [PATCH] fix: #186 multipart upload version --- ecstore/src/global.rs | 13 +++----- ecstore/src/set_disk.rs | 64 +++++++++++++++++++++++++------------ ecstore/src/store.rs | 2 +- ecstore/src/utils/crypto.rs | 14 ++++++++ rustfs/src/storage/ecfs.rs | 2 ++ 5 files changed, 65 insertions(+), 30 deletions(-) diff --git a/ecstore/src/global.rs b/ecstore/src/global.rs index 65012804e..8735249c9 100644 --- a/ecstore/src/global.rs +++ b/ecstore/src/global.rs @@ -36,7 +36,7 @@ lazy_static! { pub static ref GLOBAL_BackgroundHealState: Arc = AllHealState::new(false); pub static ref GLOBAL_ALlHealState: Arc = AllHealState::new(false); pub static ref GLOBAL_MRFState: Arc = Arc::new(MRFState::new()); - static ref globalDeploymentIDPtr: RwLock = RwLock::new(Uuid::nil()); + static ref globalDeploymentIDPtr: OnceLock = OnceLock::new(); } pub fn global_rustfs_port() -> u16 { @@ -51,14 +51,11 @@ pub fn set_global_rustfs_port(value: u16) { GLOBAL_RUSTFS_PORT.set(value).expect("set_global_rustfs_port fail"); } -pub async fn set_global_deployment_id(id: Uuid) { - let mut id_ptr = globalDeploymentIDPtr.write().await; - *id_ptr = id +pub fn set_global_deployment_id(id: Uuid) { + globalDeploymentIDPtr.set(id).unwrap(); } -pub async fn get_global_deployment_id() -> Uuid { - let id_ptr = globalDeploymentIDPtr.read().await; - - *id_ptr +pub fn get_global_deployment_id() -> Option { + globalDeploymentIDPtr.get().map(|v| v.to_string()) } pub fn set_global_endpoints(eps: Vec) { diff --git a/ecstore/src/set_disk.rs b/ecstore/src/set_disk.rs index 73c646249..c9088a7a6 100644 --- a/ecstore/src/set_disk.rs +++ b/ecstore/src/set_disk.rs @@ -606,19 +606,20 @@ impl SetDisks { } fn get_upload_id_dir(bucket: &str, object: &str, upload_id: &str) -> String { - // warn!("get_upload_id_dir upload_id {:?}", upload_id); - let upload_uuid = match base64_decode(upload_id.as_bytes()) { - Ok(res) => { - let decoded_str = String::from_utf8(res).expect("Failed to convert decoded bytes to a UTF-8 string"); - let parts: Vec<&str> = decoded_str.splitn(2, '.').collect(); - if parts.len() == 2 { - parts[1].to_string() - } else { - upload_id.to_string() - } - } - Err(_) => upload_id.to_string(), - }; + warn!("get_upload_id_dir upload_id {:?}", upload_id); + + let upload_uuid = base64_decode(upload_id.as_bytes()) + .and_then(|v| { + String::from_utf8(v).map_or(Ok(upload_id.to_owned()), |v| { + let parts: Vec<_> = v.splitn(2, '.').collect(); + if parts.len() == 2 { + Ok(parts[1].to_string()) + } else { + Ok(upload_id.to_string()) + } + }) + }) + .unwrap_or_default(); format!("{}/{}", Self::get_multipart_sha_dir(bucket, object), upload_uuid) } @@ -4263,12 +4264,10 @@ impl StorageAPI for SetDisks { } }; - let deployment_id = { get_global_deployment_id().await }; - uploads.push(MultipartInfo { bucket: bucket.to_owned(), object: object.to_owned(), - upload_id: base64_encode(format!("{}.{}", deployment_id, upload_id).as_bytes()), + upload_id: base64_encode(format!("{}.{}", get_global_deployment_id().unwrap_or_default(), upload_id).as_bytes()), initiated: Some(start_time), ..Default::default() }); @@ -4325,6 +4324,8 @@ impl StorageAPI for SetDisks { }) } async fn new_multipart_upload(&self, bucket: &str, object: &str, opts: &ObjectOptions) -> Result { + warn!("new_multipart_upload opt {:?}", opts); + let disks = self.disks.read().await; let disks = disks.clone(); @@ -4371,25 +4372,46 @@ impl StorageAPI for SetDisks { let mut fi = FileInfo::new([bucket, object].join("/").as_str(), data_drives, parity_drives); + fi.version_id = if let Some(vid) = &opts.version_id { + Some(Uuid::parse_str(vid)?) + } else { + None + }; + + if opts.versioned && opts.version_id.is_none() { + fi.version_id = Some(Uuid::new_v4()); + } + fi.data_dir = Some(Uuid::new_v4()); fi.fresh = true; let parts_metadata = vec![fi.clone(); disks.len()]; + if !user_defined.contains_key("content-type") { + // TODO: get content-type + } + + if let Some(sc) = user_defined.get(xhttp::AMZ_STORAGE_CLASS) { + if sc == storageclass::STANDARD { + let _ = user_defined.remove(xhttp::AMZ_STORAGE_CLASS); + } + } + let (shuffle_disks, mut parts_metadatas) = Self::shuffle_disks_and_parts_metadata(&disks, &parts_metadata, &fi); - let now: OffsetDateTime = OffsetDateTime::now_utc(); + let mod_time = opts.mod_time.unwrap_or(OffsetDateTime::now_utc()); + for fi in parts_metadatas.iter_mut() { fi.metadata = Some(user_defined.clone()); - fi.mod_time = Some(now); + fi.mod_time = Some(mod_time); fi.fresh = true; } - fi.mod_time = Some(now); + // fi.mod_time = Some(now); - let upload_uuid = format!("{}x{}", Uuid::new_v4(), fi.mod_time.unwrap().unix_timestamp()); + let upload_uuid = format!("{}x{}", Uuid::new_v4(), mod_time.unix_timestamp_nanos()); - let upload_id = base64_encode(format!("{}.{}", "globalDeploymentID", upload_uuid).as_bytes()); + let upload_id = base64_encode(format!("{}.{}", get_global_deployment_id().unwrap_or_default(), upload_uuid).as_bytes()); let upload_path = Self::get_upload_id_dir(bucket, object, upload_uuid.as_str()); diff --git a/ecstore/src/store.rs b/ecstore/src/store.rs index 3a1038457..bd27a35ff 100644 --- a/ecstore/src/store.rs +++ b/ecstore/src/store.rs @@ -224,7 +224,7 @@ impl ECStore { set_object_layer(ec.clone()).await; if let Some(dep_id) = deployment_id { - set_global_deployment_id(dep_id).await; + set_global_deployment_id(dep_id); } Ok(ec) diff --git a/ecstore/src/utils/crypto.rs b/ecstore/src/utils/crypto.rs index 38032182d..08ede89c7 100644 --- a/ecstore/src/utils/crypto.rs +++ b/ecstore/src/utils/crypto.rs @@ -23,3 +23,17 @@ pub fn hex(data: impl AsRef<[u8]>) -> String { // h.update(data).unwrap(); // h.finish().unwrap() // } + +#[test] +fn test_base64() { + let src = "c0194290-d911-45cb-8e12-79ec563f46a8x1735460504394878000"; + + let s = base64_encode(src.as_bytes()); + + println!("{}", &s); + + let de = base64_decode(s.clone().as_bytes()).unwrap(); + let decoded_str = String::from_utf8(de).unwrap(); + + assert_eq!(decoded_str, src) +} diff --git a/rustfs/src/storage/ecfs.rs b/rustfs/src/storage/ecfs.rs index 54724646c..51b4182c6 100644 --- a/rustfs/src/storage/ecfs.rs +++ b/rustfs/src/storage/ecfs.rs @@ -308,6 +308,8 @@ impl S3 for FS { async fn get_object(&self, req: S3Request) -> S3Result> { // mc get 3 + warn!("get_object input {:?}, vid {:?}", &req.input, req.input.version_id); + let GetObjectInput { bucket, key, version_id, .. } = req.input;