feat:#165 copyobject v1

This commit is contained in:
weisd
2024-12-30 13:59:14 +08:00
parent db5b247d58
commit 47f365d385
11 changed files with 671 additions and 161 deletions
+198 -75
View File
@@ -8,8 +8,6 @@ use std::{
use crate::config::error::is_not_found;
use crate::global::GLOBAL_MRFState;
use crate::heal::heal_ops::{HealEntryFn, HealSequence};
use crate::heal::mrf::PartialOperation;
use crate::{
bitrot::{bitrot_verify, close_bitrot_writers, new_bitrot_filereader, new_bitrot_filewriter, BitrotFileWriter},
cache_value::metacache_set::{list_path_raw, ListPathRawOptions},
@@ -54,11 +52,16 @@ use crate::{
},
xhttp,
};
use crate::{disk::STORAGE_FORMAT_FILE, heal::mrf::PartialOperation};
use crate::{file_meta::file_info_from_raw, heal::data_usage_cache::DataUsageCache};
use crate::{
heal::data_scanner::{globalHealConfig, HEAL_DELETE_DANGLING},
store_api::ListObjectVersionsInfo,
};
use crate::{
heal::heal_ops::{HealEntryFn, HealSequence},
utils::path::path_join_buf,
};
use bytesize::ByteSize;
use chrono::Utc;
use futures::future::join_all;
@@ -76,7 +79,7 @@ use rand::{
{seq::SliceRandom, Rng},
};
use reader::reader::EtagReader;
use s3s::dto::StreamingBlob;
use s3s::{dto::StreamingBlob, Body};
use sha2::{Digest, Sha256};
use std::hash::Hash;
use std::time::SystemTime;
@@ -574,10 +577,10 @@ impl SetDisks {
bucket: &str,
prefix: &str,
files: &[FileInfo],
// write_quorum: usize,
) -> Vec<Option<Error>> {
write_quorum: usize,
) -> Result<()> {
let mut futures = Vec::with_capacity(disks.len());
let mut errors = Vec::with_capacity(disks.len());
let mut errs = Vec::with_capacity(disks.len());
for (i, disk) in disks.iter().enumerate() {
let mut file_info = files[i].clone();
@@ -595,14 +598,42 @@ impl SetDisks {
for result in results {
match result {
Ok(_) => {
errors.push(None);
errs.push(None);
}
Err(e) => {
errors.push(Some(e));
errs.push(Some(e));
}
}
}
errors
if let Some(err) = reduce_write_quorum_errs(&errs, object_op_ignored_errs().as_ref(), write_quorum) {
// TODO: 并发
for (i, err) in errs.iter().enumerate() {
if err.is_some() {
continue;
}
if let Some(disk) = disks[i].as_ref() {
let _ = disk
.delete(
bucket,
&path_join_buf(&[prefix, STORAGE_FORMAT_FILE]),
DeleteOptions {
recursive: true,
..Default::default()
},
)
.await
.map_err(|e| {
warn!("write meta revert err {:?}", e);
e
});
}
}
return Err(err);
}
Ok(())
}
fn get_upload_id_dir(bucket: &str, object: &str, upload_id: &str) -> String {
@@ -1747,12 +1778,17 @@ impl SetDisks {
}
#[allow(clippy::too_many_arguments)]
#[tracing::instrument(
level = "debug",
skip( writer,disks,fi,files),
fields(start_time=?time::OffsetDateTime::now_utc())
)]
async fn get_object_with_fileinfo(
// &self,
bucket: &str,
object: &str,
offset: i64,
length: i64,
offset: usize,
length: usize,
writer: &mut DuplexStream,
fi: FileInfo,
files: Vec<FileInfo>,
@@ -1762,10 +1798,10 @@ impl SetDisks {
) -> Result<()> {
let (disks, files) = Self::shuffle_disks_and_parts_metadata_by_index(disks, &files, &fi);
let total_size = fi.size as i64;
let total_size = fi.size;
let length = {
if length < 0 {
if length == 0 {
total_size - offset
} else {
length
@@ -1791,13 +1827,13 @@ impl SetDisks {
let (last_part_index, _) = fi.to_part_offset(end_offset)?;
// debug!(
// "get_object_with_fileinfo end offset:{}, part_index:{},part_offset:{}",
// "get_object_with_fileinfo end offset:{}, last_part_index:{},part_offset:{}",
// end_offset, last_part_index, 0
// );
let erasure = Erasure::new(fi.erasure.data_blocks, fi.erasure.parity_blocks, fi.erasure.block_size);
let mut total_readed: i64 = 0;
let mut total_readed = 0;
for i in part_index..=last_part_index {
if total_readed == length {
break;
@@ -1805,12 +1841,12 @@ impl SetDisks {
let part_number = fi.parts[i].number;
let part_size = fi.parts[i].size;
let mut part_length = part_size - part_offset as usize;
if part_length > (length - total_readed) as usize {
part_length = (length - total_readed) as usize
let mut part_length = part_size - part_offset;
if part_length > length - total_readed {
part_length = length - total_readed
}
let till_offset = erasure.shard_file_offset(part_offset.try_into().unwrap(), part_length, part_size);
let till_offset = erasure.shard_file_offset(part_offset, part_length, part_size);
let mut readers = Vec::with_capacity(disks.len());
for (idx, disk_op) in disks.iter().enumerate() {
// debug!("read part_path {}", &part_path);
@@ -1844,13 +1880,12 @@ impl SetDisks {
// "read part {} part_offset {},part_length {},part_size {} ",
// part_number, part_offset, part_length, part_size
// );
let (written, mut err) = erasure
.decode(writer, readers, part_offset as usize, part_length, part_size)
.await;
let (written, mut err) = erasure.decode(writer, readers, part_offset, part_length, part_size).await;
if let Some(e) = err.as_ref() {
if written == part_length {
match e.downcast_ref::<DiskError>() {
Some(DiskError::FileNotFound) | Some(DiskError::FileCorrupt) => {
error!("erasure.decode err 111 {:?}", &e);
GLOBAL_MRFState
.add_partial(PartialOperation {
bucket: bucket.to_string(),
@@ -1870,11 +1905,12 @@ impl SetDisks {
}
}
if let Some(err) = err {
error!("erasure.decode err {} {:?}", written, &err);
return Err(err);
}
// debug!("ec decode {} writed size {}", part_number, n);
total_readed += part_length as i64;
total_readed += part_length;
part_offset = 0;
}
@@ -3495,19 +3531,26 @@ impl SetDisks {
#[async_trait::async_trait]
impl ObjectIO for SetDisks {
#[tracing::instrument(level = "debug", skip(self))]
async fn get_object_reader(
&self,
bucket: &str,
object: &str,
_rs: HTTPRangeSpec,
_h: HeaderMap,
range: Option<HTTPRangeSpec>,
h: HeaderMap,
opts: &ObjectOptions,
) -> Result<GetObjectReader> {
let (fi, files, disks) = self.get_object_fileinfo(bucket, object, opts, true).await?;
let (fi, files, disks) = self
.get_object_fileinfo(bucket, object, opts, true)
.await
.map_err(|err| to_object_err(err, vec![bucket, object]))?;
let object_info = fi.to_object_info(bucket, object, opts.versioned || opts.version_suspended);
if object_info.delete_marker {
return Err(Error::new(DiskError::FileNotFound));
if opts.version_id.is_none() {
return Err(to_object_err(Error::new(DiskError::FileNotFound), vec![bucket, object]));
}
return Err(to_object_err(Error::new(StorageError::MethodNotAllowed), vec![bucket, object]));
}
// if object_info.size == 0 {
@@ -3519,17 +3562,28 @@ impl ObjectIO for SetDisks {
// });
// }
let rs = HTTPRangeSpec::from_object_info(&object_info, opts.part_number);
let (offset, length) = rs.get_offset_length(object_info.size.try_into().unwrap())?;
if object_info.size == 0 {
if let Some(rs) = range {
let _ = rs.get_offset_length(object_info.size)?;
}
// debug!("get_object_reader offset:{}, length:{}", offset, length);
let reader = GetObjectReader {
stream: StreamingBlob::from(Body::from(Vec::new())),
object_info,
};
return Ok(reader);
}
// TODO: remote
let (rd, mut wd) = tokio::io::duplex(fi.erasure.block_size);
// let disks = self.disks.read().await;
let (reader, offset, length) =
GetObjectReader::new(StreamingBlob::wrap(tokio_util::io::ReaderStream::new(rd)), range, &object_info, opts, &h)?;
// let disks = disks.clone();
let bucket = String::from(bucket);
let object = String::from(object);
let bucket = bucket.to_owned();
let object = object.to_owned();
let set_index = self.set_index;
let pool_index = self.pool_index;
tokio::spawn(async move {
@@ -3542,14 +3596,6 @@ impl ObjectIO for SetDisks {
};
});
let read_stream = tokio_util::io::ReaderStream::new(rd);
// let rd: Box<dyn AsyncRead> = Box::new(rd);
let reader = GetObjectReader {
stream: StreamingBlob::wrap(read_stream),
object_info,
};
Ok(reader)
}
@@ -3585,7 +3631,7 @@ impl ObjectIO for SetDisks {
_ns = Some(ns_lock);
}
let mut user_defined = opts.user_defined.clone();
let mut user_defined = opts.user_defined.clone().unwrap_or_default();
let sc_parity_drives = {
if let Some(sc) = GLOBAL_StorageClass.get() {
@@ -3789,6 +3835,110 @@ impl StorageAPI for SetDisks {
unimplemented!()
}
async fn copy_object(
&self,
src_bucket: &str,
src_object: &str,
_dst_bucket: &str,
_dst_object: &str,
src_info: &mut ObjectInfo,
src_opts: &ObjectOptions,
dst_opts: &ObjectOptions,
) -> Result<ObjectInfo> {
// FIXME: TODO:
if !src_info.metadata_only {
return Err(Error::new(StorageError::NotImplemented));
}
let disks = self.get_disks_internal().await;
let (mut metas, errs) = {
if let Some(vid) = &src_opts.version_id {
Self::read_all_fileinfo(&disks, "", src_bucket, src_object, vid, true, false).await
} else {
Self::read_all_xl(&disks, src_bucket, src_object, true, false).await
}
};
let (read_quorum, write_quorum) = match Self::object_quorum_from_meta(&metas, &errs, self.default_parity_count) {
Ok((r, w)) => (r as usize, w as usize),
Err(mut err) => {
if ErasureError::ErasureReadQuorum.is(&err)
&& !src_bucket.starts_with(RUSTFS_META_BUCKET)
&& self
.delete_if_dang_ling(src_bucket, src_object, &metas, &errs, &HashMap::new(), src_opts.clone())
.await
.is_ok()
{
if src_opts.version_id.is_some() {
err = Error::new(DiskError::FileVersionNotFound)
} else {
err = Error::new(DiskError::FileNotFound)
}
}
return Err(to_object_err(err, vec![src_bucket, src_object]));
}
};
let (online_disks, mod_time, etag) = Self::list_online_disks(&disks, &metas, &errs, read_quorum);
let mut fi = Self::pick_valid_fileinfo(&metas, mod_time, etag, read_quorum)
.map_err(|e| to_object_err(e, vec![src_bucket, src_object]))?;
if fi.deleted {
if src_opts.version_id.is_none() {
return Err(to_object_err(Error::new(DiskError::FileNotFound), vec![src_bucket, src_object]));
}
return Err(to_object_err(Error::new(StorageError::MethodNotAllowed), vec![src_bucket, src_object]));
}
let version_id = {
if src_info.version_only {
if let Some(vid) = &dst_opts.version_id {
Some(Uuid::parse_str(vid)?)
} else {
Some(Uuid::new_v4())
}
} else {
src_info.version_id
}
};
let inline_data = fi.inline_data();
fi.metadata = src_info.user_defined.clone();
if let Some(ud) = src_info.user_defined.as_mut() {
if let Some(etag) = &src_info.etag {
ud.insert("etag".to_owned(), etag.clone());
}
}
let mod_time = OffsetDateTime::now_utc();
for fi in metas.iter_mut() {
if fi.is_valid() {
fi.metadata = src_info.user_defined.clone();
fi.mod_time = Some(mod_time);
fi.version_id = version_id;
fi.versioned = src_opts.versioned || src_opts.version_suspended;
if !fi.inline_data() {
fi.data = None;
}
if inline_data {
fi.set_inline_data();
}
}
}
Self::write_unique_file_info(&online_disks, "", src_bucket, src_object, &metas, write_quorum)
.await
.map_err(|e| to_object_err(e, vec![src_bucket, src_object]))?;
Ok(fi.to_object_info(src_bucket, src_object, src_opts.versioned || src_opts.version_suspended))
}
async fn delete_objects(
&self,
bucket: &str,
@@ -4333,7 +4483,7 @@ impl StorageAPI for SetDisks {
let disks = disks.clone();
let mut user_defined = opts.user_defined.clone();
let mut user_defined = opts.user_defined.clone().unwrap_or_default();
if let Some(ref etag) = opts.preserve_etag {
user_defined.insert("etag".to_owned(), etag.clone());
@@ -4371,8 +4521,6 @@ impl StorageAPI for SetDisks {
write_quorum += 1
}
let _ = write_quorum;
let mut fi = FileInfo::new([bucket, object].join("/").as_str(), data_drives, parity_drives);
fi.version_id = if let Some(vid) = &opts.version_id {
@@ -4418,42 +4566,17 @@ impl StorageAPI for SetDisks {
let upload_path = Self::get_upload_id_dir(bucket, object, upload_uuid.as_str());
let errs = Self::write_unique_file_info(
Self::write_unique_file_info(
&shuffle_disks,
bucket,
RUSTFS_META_MULTIPART_BUCKET,
upload_path.as_str(),
&parts_metadatas,
write_quorum,
)
.await;
.await
.map_err(|e| to_object_err(e, vec![bucket, object]))?;
if let Some(err) = reduce_write_quorum_errs(&errs, object_op_ignored_errs().as_ref(), write_quorum) {
// TODO: 并发
for (i, err) in errs.iter().enumerate() {
if err.is_some() {
continue;
}
if let Some(disk) = shuffle_disks[i].as_ref() {
let _ = disk
.delete(
RUSTFS_META_MULTIPART_BUCKET,
upload_path.as_str(),
DeleteOptions {
recursive: true,
..Default::default()
},
)
.await
.map_err(|e| {
warn!("write meta revert err {:?}", e);
e
});
}
}
return Err(err);
}
// evalDisks
Ok(MultipartUploadResult { upload_id })
@@ -4612,7 +4735,7 @@ impl StorageAPI for SetDisks {
// etag
let etag = {
if let Some(etag) = opts.user_defined.get("etag") {
if let Some(Some(etag)) = opts.user_defined.as_ref().map(|v| v.get("etag")) {
etag.clone()
} else {
get_complete_multipart_md5(&uploaded_parts)