diff --git a/crates/ecstore/src/error.rs b/crates/ecstore/src/error.rs index 4de5c1595..43a524766 100644 --- a/crates/ecstore/src/error.rs +++ b/crates/ecstore/src/error.rs @@ -156,11 +156,12 @@ pub enum StorageError { #[error("Object exists on :{0} as directory {1}")] ObjectExistsAsDirectory(String, String), - // #[error("Storage resources are insufficient for the read operation")] - // InsufficientReadQuorum, + #[error("Storage resources are insufficient for the read operation: {0}/{1}")] + InsufficientReadQuorum(String, String), + + #[error("Storage resources are insufficient for the write operation: {0}/{1}")] + InsufficientWriteQuorum(String, String), - // #[error("Storage resources are insufficient for the write operation")] - // InsufficientWriteQuorum, #[error("Decommission not started")] DecommissionNotStarted, #[error("Decommission already running")] @@ -413,6 +414,8 @@ impl Clone for StorageError { StorageError::TooManyOpenFiles => StorageError::TooManyOpenFiles, StorageError::NoHealRequired => StorageError::NoHealRequired, StorageError::Lock(e) => StorageError::Lock(e.clone()), + StorageError::InsufficientReadQuorum(a, b) => StorageError::InsufficientReadQuorum(a.clone(), b.clone()), + StorageError::InsufficientWriteQuorum(a, b) => StorageError::InsufficientWriteQuorum(a.clone(), b.clone()), } } } @@ -476,6 +479,8 @@ impl StorageError { StorageError::TooManyOpenFiles => 0x36, StorageError::NoHealRequired => 0x37, StorageError::Lock(_) => 0x38, + StorageError::InsufficientReadQuorum(_, _) => 0x39, + StorageError::InsufficientWriteQuorum(_, _) => 0x3A, } } @@ -541,6 +546,8 @@ impl StorageError { 0x36 => Some(StorageError::TooManyOpenFiles), 0x37 => Some(StorageError::NoHealRequired), 0x38 => Some(StorageError::Lock(rustfs_lock::LockError::internal("Generic lock error".to_string()))), + 0x39 => Some(StorageError::InsufficientReadQuorum(Default::default(), Default::default())), + 0x3A => Some(StorageError::InsufficientWriteQuorum(Default::default(), Default::default())), _ => None, } } @@ -753,6 +760,17 @@ pub fn to_object_err(err: Error, params: Vec<&str>) -> Error { StorageError::PrefixAccessDenied(bucket, object) } + StorageError::ErasureReadQuorum => { + let bucket = params.first().cloned().unwrap_or_default().to_owned(); + let object = params.get(1).cloned().map(decode_dir_object).unwrap_or_default(); + StorageError::InsufficientReadQuorum(bucket, object) + } + StorageError::ErasureWriteQuorum => { + let bucket = params.first().cloned().unwrap_or_default().to_owned(); + let object = params.get(1).cloned().map(decode_dir_object).unwrap_or_default(); + StorageError::InsufficientWriteQuorum(bucket, object) + } + _ => err, } } diff --git a/crates/ecstore/src/set_disk.rs b/crates/ecstore/src/set_disk.rs index 6878f2320..58dc5ae0c 100644 --- a/crates/ecstore/src/set_disk.rs +++ b/crates/ecstore/src/set_disk.rs @@ -17,6 +17,8 @@ use crate::bitrot::{create_bitrot_reader, create_bitrot_writer}; use crate::bucket::lifecycle::lifecycle::TRANSITION_COMPLETE; +use crate::bucket::versioning::VersioningApi; +use crate::bucket::versioning_sys::BucketVersioningSys; use crate::client::{object_api_utils::extract_etag, transition_api::ReaderImpl}; use crate::disk::STORAGE_FORMAT_FILE; use crate::disk::error_reduce::{OBJECT_OP_IGNORED_ERRS, reduce_read_quorum_errs, reduce_write_quorum_errs}; @@ -2027,6 +2029,24 @@ impl SetDisks { Ok((fi, parts_metadata, op_online_disks)) } + async fn get_object_info_and_quorum(&self, bucket: &str, object: &str, opts: &ObjectOptions) -> Result<(ObjectInfo, usize)> { + let (fi, _, _) = self.get_object_fileinfo(bucket, object, opts, false).await?; + + let write_quorum = fi.write_quorum(self.default_write_quorum()); + + let oi = ObjectInfo::from_file_info(&fi, bucket, object, opts.versioned || opts.version_suspended); + // TODO: replicatio + + if fi.deleted { + if opts.version_id.is_none() || opts.delete_marker { + return Err(to_object_err(StorageError::FileNotFound, vec![bucket, object])); + } else { + return Err(to_object_err(StorageError::MethodNotAllowed, vec![bucket, object])); + } + } + + Ok((oi, write_quorum)) + } #[allow(clippy::too_many_arguments)] #[tracing::instrument( @@ -3769,39 +3789,47 @@ impl StorageAPI for SetDisks { } // Per-object guards to keep until function end - let mut _guards: Vec> = Vec::with_capacity(objects.len()); + let mut _guards: HashMap = HashMap::new(); // Acquire locks for all objects first; mark errors for failures for (i, dobj) in objects.iter().enumerate() { - match self - .namespace_lock - .lock_guard(&dobj.object_name, &self.locker_owner, Duration::from_secs(5), Duration::from_secs(10)) - .await? - { - Some(g) => _guards.push(Some(g)), - None => { - del_errs[i] = Some(Error::other("can not get lock. please retry")); - _guards.push(None); + if !_guards.contains_key(&dobj.object_name) { + match self + .namespace_lock + .lock_guard(&dobj.object_name, &self.locker_owner, Duration::from_secs(5), Duration::from_secs(10)) + .await? + { + Some(g) => { + _guards.insert(dobj.object_name.clone(), g); + } + None => { + del_errs[i] = Some(Error::other("can not get lock. please retry")); + } } } } // let mut del_fvers = Vec::with_capacity(objects.len()); + let ver_cfg = BucketVersioningSys::get(bucket).await.unwrap_or_default(); + let mut vers_map: HashMap<&String, FileInfoVersions> = HashMap::new(); for (i, dobj) in objects.iter().enumerate() { let mut vr = FileInfo { name: dobj.object_name.clone(), version_id: dobj.version_id, + idx: i, ..Default::default() }; - // 删除 - del_objects[i].object_name.clone_from(&vr.name); - del_objects[i].version_id = vr.version_id.map(|v| v.to_string()); + vr.set_tier_free_version_id(&Uuid::new_v4().to_string()); - if del_objects[i].version_id.is_none() { - let (suspended, versioned) = (opts.version_suspended, opts.versioned); + // 删除 + // del_objects[i].object_name.clone_from(&vr.name); + // del_objects[i].version_id = vr.version_id.map(|v| v.to_string()); + + if dobj.version_id.is_none() { + let (suspended, versioned) = (ver_cfg.suspended(), ver_cfg.prefix_enabled(dobj.object_name.as_str())); if suspended || versioned { vr.mod_time = Some(OffsetDateTime::now_utc()); vr.deleted = true; @@ -3842,15 +3870,22 @@ impl StorageAPI for SetDisks { } // Only add to vers_map if we hold the lock - if _guards[i].is_some() { + if _guards.contains_key(&dobj.object_name) { vers_map.insert(&dobj.object_name, v); } } let mut vers = Vec::with_capacity(vers_map.len()); - for (_, ver) in vers_map { - vers.push(ver); + for (_, mut fi_vers) in vers_map { + fi_vers.versions.sort_by(|a, b| a.deleted.cmp(&b.deleted)); + fi_vers.versions.reverse(); + + if let Some(index) = fi_vers.versions.iter().position(|fi| fi.deleted) { + fi_vers.versions.truncate(index + 1); + } + + vers.push(fi_vers); } let disks = self.disks.read().await; @@ -3906,6 +3941,61 @@ impl StorageAPI for SetDisks { return Ok(ObjectInfo::default()); } + let (oi, write_quorum) = match self.get_object_info_and_quorum(bucket, object, &opts).await { + Ok((oi, wq)) => (oi, wq), + Err(e) => { + return Err(to_object_err(e, vec![bucket, object])); + } + }; + + let mark_delete = oi.version_id.is_some(); + + let mut delete_marker = opts.versioned; + + let mod_time = if let Some(mt) = opts.mod_time { + mt + } else { + OffsetDateTime::now_utc() + }; + + let find_vid = Uuid::new_v4(); + + if mark_delete && (opts.versioned || opts.version_suspended) { + if !delete_marker { + delete_marker = opts.version_suspended && opts.version_id.is_none(); + } + + let mut fi = FileInfo { + name: object.to_string(), + deleted: delete_marker, + mark_deleted: mark_delete, + mod_time: Some(mod_time), + ..Default::default() // TODO: replication + }; + + fi.set_tier_free_version_id(&find_vid.to_string()); + + if opts.skip_free_version { + fi.set_skip_tier_free_version(); + } + + fi.version_id = if let Some(vid) = opts.version_id { + Some(Uuid::parse_str(vid.as_str())?) + } else if opts.versioned { + Some(Uuid::new_v4()) + } else { + None + }; + + self.delete_object_version(bucket, object, &fi, opts.delete_marker) + .await + .map_err(|e| to_object_err(e, vec![bucket, object]))?; + + return Ok(ObjectInfo::from_file_info(&fi, bucket, object, opts.versioned || opts.version_suspended)); + } + + let version_id = opts.version_id.as_ref().and_then(|v| Uuid::parse_str(v).ok()); + // Create a single object deletion request let mut vr = FileInfo { name: object.to_string(), diff --git a/crates/ecstore/src/store_api.rs b/crates/ecstore/src/store_api.rs index 7373734fa..5f9a2fe31 100644 --- a/crates/ecstore/src/store_api.rs +++ b/crates/ecstore/src/store_api.rs @@ -310,6 +310,8 @@ pub struct ObjectOptions { pub replication_request: bool, pub delete_marker: bool, + pub skip_free_version: bool, + pub transition: TransitionOptions, pub expiration: ExpirationOptions, pub lifecycle_audit_event: LcAuditEvent, diff --git a/crates/filemeta/src/filemeta.rs b/crates/filemeta/src/filemeta.rs index c8126b79d..cde30da54 100644 --- a/crates/filemeta/src/filemeta.rs +++ b/crates/filemeta/src/filemeta.rs @@ -496,39 +496,32 @@ impl FileMeta { } pub fn add_version_filemata(&mut self, ver: FileMetaVersion) -> Result<()> { - let mod_time = ver.get_mod_time().unwrap().nanosecond(); if !ver.valid() { return Err(Error::other("attempted to add invalid version")); } - let encoded = ver.marshal_msg()?; - if self.versions.len() + 1 > 100 { + if self.versions.len() + 1 >= 100 { return Err(Error::other( "You've exceeded the limit on the number of versions you can create on this object", )); } - self.versions.push(FileMetaShallowVersion { - header: FileMetaVersionHeader { - mod_time: Some(OffsetDateTime::from_unix_timestamp(-1)?), - ..Default::default() - }, - ..Default::default() - }); + let mod_time = ver.get_mod_time(); + let encoded = ver.marshal_msg()?; + let new_version = FileMetaShallowVersion { + header: ver.header(), + meta: encoded, + }; - let len = self.versions.len(); - for (i, existing) in self.versions.iter().enumerate() { - if existing.header.mod_time.unwrap().nanosecond() <= mod_time { - let vers = self.versions[i..len - 1].to_vec(); - self.versions[i + 1..].clone_from_slice(vers.as_slice()); - self.versions[i] = FileMetaShallowVersion { - header: ver.header(), - meta: encoded, - }; - return Ok(()); - } - } - Err(Error::other("addVersion: Internal error, unable to add version")) + // Find the insertion position: insert before the first element with mod_time >= new mod_time + // This maintains descending order by mod_time (newest first) + let insert_pos = self + .versions + .iter() + .position(|existing| existing.header.mod_time <= mod_time) + .unwrap_or(self.versions.len()); + self.versions.insert(insert_pos, new_version); + Ok(()) } // delete_version deletes version, returns data_dir @@ -554,7 +547,15 @@ impl FileMeta { match ver.header.version_type { VersionType::Invalid | VersionType::Legacy => return Err(Error::other("invalid file meta version")), - VersionType::Delete => return Ok(None), + VersionType::Delete => { + self.versions.remove(i); + if fi.deleted && fi.version_id.is_none() { + self.add_version_filemata(ventry)?; + return Ok(None); + } + + return Ok(None); + } VersionType::Object => { let v = self.get_idx(i)?; @@ -600,6 +601,7 @@ impl FileMeta { if fi.deleted { self.add_version_filemata(ventry)?; + return Ok(None); } Err(Error::FileVersionNotFound) @@ -961,7 +963,8 @@ impl FileMetaVersion { pub fn get_version_id(&self) -> Option { match self.version_type { - VersionType::Object | VersionType::Delete => self.object.as_ref().map(|v| v.version_id).unwrap_or_default(), + VersionType::Object => self.object.as_ref().map(|v| v.version_id).unwrap_or_default(), + VersionType::Delete => self.delete_marker.as_ref().map(|v| v.version_id).unwrap_or_default(), _ => None, } } diff --git a/rustfs/src/storage/ecfs.rs b/rustfs/src/storage/ecfs.rs index 7ad743b81..a6031c349 100644 --- a/rustfs/src/storage/ecfs.rs +++ b/rustfs/src/storage/ecfs.rs @@ -1318,7 +1318,7 @@ impl S3 for FS { let objects: Vec = object_infos .objects .iter() - .filter(|v| !v.name.is_empty()) + .filter(|v| !v.name.is_empty() && !v.delete_marker) .map(|v| { ObjectVersion { key: Some(v.name.to_owned()), @@ -1340,6 +1340,19 @@ impl S3 for FS { .map(|v| CommonPrefix { prefix: Some(v) }) .collect(); + let delete_markers = object_infos + .objects + .iter() + .filter(|o| o.delete_marker) + .map(|o| DeleteMarkerEntry { + key: Some(o.name.clone()), + version_id: o.version_id.map(|v| v.to_string()), + is_latest: Some(o.is_latest), + last_modified: o.mod_time.map(Timestamp::from), + ..Default::default() + }) + .collect::>(); + let output = ListObjectVersionsOutput { // is_truncated: Some(object_infos.is_truncated), max_keys: Some(key_count), @@ -1348,6 +1361,7 @@ impl S3 for FS { prefix: Some(prefix), common_prefixes: Some(common_prefixes), versions: Some(objects), + delete_markers: Some(delete_markers), ..Default::default() }; diff --git a/rustfs/src/storage/options.rs b/rustfs/src/storage/options.rs index 6fec66f22..a3405de7f 100644 --- a/rustfs/src/storage/options.rs +++ b/rustfs/src/storage/options.rs @@ -722,8 +722,7 @@ mod tests { assert_eq!( metadata.get("content-type"), Some(&expected_content_type.to_string()), - "Failed for filename: {}", - filename + "Failed for filename: {filename}" ); } }