mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-24 13:16:28 +00:00
fix:rm object versions (#385)
This commit is contained in:
@@ -156,11 +156,12 @@ pub enum StorageError {
|
|||||||
#[error("Object exists on :{0} as directory {1}")]
|
#[error("Object exists on :{0} as directory {1}")]
|
||||||
ObjectExistsAsDirectory(String, String),
|
ObjectExistsAsDirectory(String, String),
|
||||||
|
|
||||||
// #[error("Storage resources are insufficient for the read operation")]
|
#[error("Storage resources are insufficient for the read operation: {0}/{1}")]
|
||||||
// InsufficientReadQuorum,
|
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")]
|
#[error("Decommission not started")]
|
||||||
DecommissionNotStarted,
|
DecommissionNotStarted,
|
||||||
#[error("Decommission already running")]
|
#[error("Decommission already running")]
|
||||||
@@ -413,6 +414,8 @@ impl Clone for StorageError {
|
|||||||
StorageError::TooManyOpenFiles => StorageError::TooManyOpenFiles,
|
StorageError::TooManyOpenFiles => StorageError::TooManyOpenFiles,
|
||||||
StorageError::NoHealRequired => StorageError::NoHealRequired,
|
StorageError::NoHealRequired => StorageError::NoHealRequired,
|
||||||
StorageError::Lock(e) => StorageError::Lock(e.clone()),
|
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::TooManyOpenFiles => 0x36,
|
||||||
StorageError::NoHealRequired => 0x37,
|
StorageError::NoHealRequired => 0x37,
|
||||||
StorageError::Lock(_) => 0x38,
|
StorageError::Lock(_) => 0x38,
|
||||||
|
StorageError::InsufficientReadQuorum(_, _) => 0x39,
|
||||||
|
StorageError::InsufficientWriteQuorum(_, _) => 0x3A,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -541,6 +546,8 @@ impl StorageError {
|
|||||||
0x36 => Some(StorageError::TooManyOpenFiles),
|
0x36 => Some(StorageError::TooManyOpenFiles),
|
||||||
0x37 => Some(StorageError::NoHealRequired),
|
0x37 => Some(StorageError::NoHealRequired),
|
||||||
0x38 => Some(StorageError::Lock(rustfs_lock::LockError::internal("Generic lock error".to_string()))),
|
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,
|
_ => None,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -753,6 +760,17 @@ pub fn to_object_err(err: Error, params: Vec<&str>) -> Error {
|
|||||||
StorageError::PrefixAccessDenied(bucket, object)
|
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,
|
_ => err,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
+108
-18
@@ -17,6 +17,8 @@
|
|||||||
|
|
||||||
use crate::bitrot::{create_bitrot_reader, create_bitrot_writer};
|
use crate::bitrot::{create_bitrot_reader, create_bitrot_writer};
|
||||||
use crate::bucket::lifecycle::lifecycle::TRANSITION_COMPLETE;
|
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::client::{object_api_utils::extract_etag, transition_api::ReaderImpl};
|
||||||
use crate::disk::STORAGE_FORMAT_FILE;
|
use crate::disk::STORAGE_FORMAT_FILE;
|
||||||
use crate::disk::error_reduce::{OBJECT_OP_IGNORED_ERRS, reduce_read_quorum_errs, reduce_write_quorum_errs};
|
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))
|
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)]
|
#[allow(clippy::too_many_arguments)]
|
||||||
#[tracing::instrument(
|
#[tracing::instrument(
|
||||||
@@ -3769,39 +3789,47 @@ impl StorageAPI for SetDisks {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// Per-object guards to keep until function end
|
// Per-object guards to keep until function end
|
||||||
let mut _guards: Vec<Option<rustfs_lock::LockGuard>> = Vec::with_capacity(objects.len());
|
let mut _guards: HashMap<String, rustfs_lock::LockGuard> = HashMap::new();
|
||||||
// Acquire locks for all objects first; mark errors for failures
|
// Acquire locks for all objects first; mark errors for failures
|
||||||
for (i, dobj) in objects.iter().enumerate() {
|
for (i, dobj) in objects.iter().enumerate() {
|
||||||
match self
|
if !_guards.contains_key(&dobj.object_name) {
|
||||||
.namespace_lock
|
match self
|
||||||
.lock_guard(&dobj.object_name, &self.locker_owner, Duration::from_secs(5), Duration::from_secs(10))
|
.namespace_lock
|
||||||
.await?
|
.lock_guard(&dobj.object_name, &self.locker_owner, Duration::from_secs(5), Duration::from_secs(10))
|
||||||
{
|
.await?
|
||||||
Some(g) => _guards.push(Some(g)),
|
{
|
||||||
None => {
|
Some(g) => {
|
||||||
del_errs[i] = Some(Error::other("can not get lock. please retry"));
|
_guards.insert(dobj.object_name.clone(), g);
|
||||||
_guards.push(None);
|
}
|
||||||
|
None => {
|
||||||
|
del_errs[i] = Some(Error::other("can not get lock. please retry"));
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// let mut del_fvers = Vec::with_capacity(objects.len());
|
// 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();
|
let mut vers_map: HashMap<&String, FileInfoVersions> = HashMap::new();
|
||||||
|
|
||||||
for (i, dobj) in objects.iter().enumerate() {
|
for (i, dobj) in objects.iter().enumerate() {
|
||||||
let mut vr = FileInfo {
|
let mut vr = FileInfo {
|
||||||
name: dobj.object_name.clone(),
|
name: dobj.object_name.clone(),
|
||||||
version_id: dobj.version_id,
|
version_id: dobj.version_id,
|
||||||
|
idx: i,
|
||||||
..Default::default()
|
..Default::default()
|
||||||
};
|
};
|
||||||
|
|
||||||
// 删除
|
vr.set_tier_free_version_id(&Uuid::new_v4().to_string());
|
||||||
del_objects[i].object_name.clone_from(&vr.name);
|
|
||||||
del_objects[i].version_id = vr.version_id.map(|v| v.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 {
|
if suspended || versioned {
|
||||||
vr.mod_time = Some(OffsetDateTime::now_utc());
|
vr.mod_time = Some(OffsetDateTime::now_utc());
|
||||||
vr.deleted = true;
|
vr.deleted = true;
|
||||||
@@ -3842,15 +3870,22 @@ impl StorageAPI for SetDisks {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// Only add to vers_map if we hold the lock
|
// 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);
|
vers_map.insert(&dobj.object_name, v);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
let mut vers = Vec::with_capacity(vers_map.len());
|
let mut vers = Vec::with_capacity(vers_map.len());
|
||||||
|
|
||||||
for (_, ver) in vers_map {
|
for (_, mut fi_vers) in vers_map {
|
||||||
vers.push(ver);
|
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;
|
let disks = self.disks.read().await;
|
||||||
@@ -3906,6 +3941,61 @@ impl StorageAPI for SetDisks {
|
|||||||
return Ok(ObjectInfo::default());
|
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
|
// Create a single object deletion request
|
||||||
let mut vr = FileInfo {
|
let mut vr = FileInfo {
|
||||||
name: object.to_string(),
|
name: object.to_string(),
|
||||||
|
|||||||
@@ -310,6 +310,8 @@ pub struct ObjectOptions {
|
|||||||
pub replication_request: bool,
|
pub replication_request: bool,
|
||||||
pub delete_marker: bool,
|
pub delete_marker: bool,
|
||||||
|
|
||||||
|
pub skip_free_version: bool,
|
||||||
|
|
||||||
pub transition: TransitionOptions,
|
pub transition: TransitionOptions,
|
||||||
pub expiration: ExpirationOptions,
|
pub expiration: ExpirationOptions,
|
||||||
pub lifecycle_audit_event: LcAuditEvent,
|
pub lifecycle_audit_event: LcAuditEvent,
|
||||||
|
|||||||
@@ -496,39 +496,32 @@ impl FileMeta {
|
|||||||
}
|
}
|
||||||
|
|
||||||
pub fn add_version_filemata(&mut self, ver: FileMetaVersion) -> Result<()> {
|
pub fn add_version_filemata(&mut self, ver: FileMetaVersion) -> Result<()> {
|
||||||
let mod_time = ver.get_mod_time().unwrap().nanosecond();
|
|
||||||
if !ver.valid() {
|
if !ver.valid() {
|
||||||
return Err(Error::other("attempted to add invalid version"));
|
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(
|
return Err(Error::other(
|
||||||
"You've exceeded the limit on the number of versions you can create on this object",
|
"You've exceeded the limit on the number of versions you can create on this object",
|
||||||
));
|
));
|
||||||
}
|
}
|
||||||
|
|
||||||
self.versions.push(FileMetaShallowVersion {
|
let mod_time = ver.get_mod_time();
|
||||||
header: FileMetaVersionHeader {
|
let encoded = ver.marshal_msg()?;
|
||||||
mod_time: Some(OffsetDateTime::from_unix_timestamp(-1)?),
|
let new_version = FileMetaShallowVersion {
|
||||||
..Default::default()
|
header: ver.header(),
|
||||||
},
|
meta: encoded,
|
||||||
..Default::default()
|
};
|
||||||
});
|
|
||||||
|
|
||||||
let len = self.versions.len();
|
// Find the insertion position: insert before the first element with mod_time >= new mod_time
|
||||||
for (i, existing) in self.versions.iter().enumerate() {
|
// This maintains descending order by mod_time (newest first)
|
||||||
if existing.header.mod_time.unwrap().nanosecond() <= mod_time {
|
let insert_pos = self
|
||||||
let vers = self.versions[i..len - 1].to_vec();
|
.versions
|
||||||
self.versions[i + 1..].clone_from_slice(vers.as_slice());
|
.iter()
|
||||||
self.versions[i] = FileMetaShallowVersion {
|
.position(|existing| existing.header.mod_time <= mod_time)
|
||||||
header: ver.header(),
|
.unwrap_or(self.versions.len());
|
||||||
meta: encoded,
|
self.versions.insert(insert_pos, new_version);
|
||||||
};
|
Ok(())
|
||||||
return Ok(());
|
|
||||||
}
|
|
||||||
}
|
|
||||||
Err(Error::other("addVersion: Internal error, unable to add version"))
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// delete_version deletes version, returns data_dir
|
// delete_version deletes version, returns data_dir
|
||||||
@@ -554,7 +547,15 @@ impl FileMeta {
|
|||||||
|
|
||||||
match ver.header.version_type {
|
match ver.header.version_type {
|
||||||
VersionType::Invalid | VersionType::Legacy => return Err(Error::other("invalid file meta version")),
|
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 => {
|
VersionType::Object => {
|
||||||
let v = self.get_idx(i)?;
|
let v = self.get_idx(i)?;
|
||||||
|
|
||||||
@@ -600,6 +601,7 @@ impl FileMeta {
|
|||||||
|
|
||||||
if fi.deleted {
|
if fi.deleted {
|
||||||
self.add_version_filemata(ventry)?;
|
self.add_version_filemata(ventry)?;
|
||||||
|
return Ok(None);
|
||||||
}
|
}
|
||||||
|
|
||||||
Err(Error::FileVersionNotFound)
|
Err(Error::FileVersionNotFound)
|
||||||
@@ -961,7 +963,8 @@ impl FileMetaVersion {
|
|||||||
|
|
||||||
pub fn get_version_id(&self) -> Option<Uuid> {
|
pub fn get_version_id(&self) -> Option<Uuid> {
|
||||||
match self.version_type {
|
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,
|
_ => None,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1318,7 +1318,7 @@ impl S3 for FS {
|
|||||||
let objects: Vec<ObjectVersion> = object_infos
|
let objects: Vec<ObjectVersion> = object_infos
|
||||||
.objects
|
.objects
|
||||||
.iter()
|
.iter()
|
||||||
.filter(|v| !v.name.is_empty())
|
.filter(|v| !v.name.is_empty() && !v.delete_marker)
|
||||||
.map(|v| {
|
.map(|v| {
|
||||||
ObjectVersion {
|
ObjectVersion {
|
||||||
key: Some(v.name.to_owned()),
|
key: Some(v.name.to_owned()),
|
||||||
@@ -1340,6 +1340,19 @@ impl S3 for FS {
|
|||||||
.map(|v| CommonPrefix { prefix: Some(v) })
|
.map(|v| CommonPrefix { prefix: Some(v) })
|
||||||
.collect();
|
.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::<Vec<_>>();
|
||||||
|
|
||||||
let output = ListObjectVersionsOutput {
|
let output = ListObjectVersionsOutput {
|
||||||
// is_truncated: Some(object_infos.is_truncated),
|
// is_truncated: Some(object_infos.is_truncated),
|
||||||
max_keys: Some(key_count),
|
max_keys: Some(key_count),
|
||||||
@@ -1348,6 +1361,7 @@ impl S3 for FS {
|
|||||||
prefix: Some(prefix),
|
prefix: Some(prefix),
|
||||||
common_prefixes: Some(common_prefixes),
|
common_prefixes: Some(common_prefixes),
|
||||||
versions: Some(objects),
|
versions: Some(objects),
|
||||||
|
delete_markers: Some(delete_markers),
|
||||||
..Default::default()
|
..Default::default()
|
||||||
};
|
};
|
||||||
|
|
||||||
|
|||||||
@@ -722,8 +722,7 @@ mod tests {
|
|||||||
assert_eq!(
|
assert_eq!(
|
||||||
metadata.get("content-type"),
|
metadata.get("content-type"),
|
||||||
Some(&expected_content_type.to_string()),
|
Some(&expected_content_type.to_string()),
|
||||||
"Failed for filename: {}",
|
"Failed for filename: {filename}"
|
||||||
filename
|
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user