mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-27 15:37:02 +00:00
Merge branch 'main' of github.com:rustfs/s3-rustfs into feature/observability-metrics
# Conflicts: # ecstore/src/file_meta.rs # ecstore/src/set_disk.rs # ecstore/src/store_api.rs
This commit is contained in:
+17
-18
@@ -71,7 +71,7 @@ impl FileMeta {
|
|||||||
Ok(xl)
|
Ok(xl)
|
||||||
}
|
}
|
||||||
|
|
||||||
// check_xl2_v1 读xl文件头,返回后续内容,版本信息
|
// check_xl2_v1 读 xl 文件头,返回后续内容,版本信息
|
||||||
// checkXL2V1
|
// checkXL2V1
|
||||||
#[tracing::instrument]
|
#[tracing::instrument]
|
||||||
pub fn check_xl2_v1(buf: &[u8]) -> Result<(&[u8], u16, u16)> {
|
pub fn check_xl2_v1(buf: &[u8]) -> Result<(&[u8], u16, u16)> {
|
||||||
@@ -92,11 +92,11 @@ impl FileMeta {
|
|||||||
Ok((&buf[8..], major, minor))
|
Ok((&buf[8..], major, minor))
|
||||||
}
|
}
|
||||||
|
|
||||||
// 固定u32
|
// 固定 u32
|
||||||
pub fn read_bytes_header(buf: &[u8]) -> Result<(u32, &[u8])> {
|
pub fn read_bytes_header(buf: &[u8]) -> Result<(u32, &[u8])> {
|
||||||
let (mut size_buf, _) = buf.split_at(5);
|
let (mut size_buf, _) = buf.split_at(5);
|
||||||
|
|
||||||
// 取meta数据,buf = crc + data
|
// 取 meta 数据,buf = crc + data
|
||||||
let bin_len = rmp::decode::read_bin_len(&mut size_buf)?;
|
let bin_len = rmp::decode::read_bin_len(&mut size_buf)?;
|
||||||
|
|
||||||
Ok((bin_len, &buf[5..]))
|
Ok((bin_len, &buf[5..]))
|
||||||
@@ -110,7 +110,7 @@ impl FileMeta {
|
|||||||
|
|
||||||
let (mut size_buf, buf) = buf.split_at(5);
|
let (mut size_buf, buf) = buf.split_at(5);
|
||||||
|
|
||||||
// 取meta数据,buf = crc + data
|
// 取 meta 数据,buf = crc + data
|
||||||
let bin_len = rmp::decode::read_bin_len(&mut size_buf)?;
|
let bin_len = rmp::decode::read_bin_len(&mut size_buf)?;
|
||||||
|
|
||||||
let (meta, buf) = buf.split_at(bin_len as usize);
|
let (meta, buf) = buf.split_at(bin_len as usize);
|
||||||
@@ -130,7 +130,7 @@ impl FileMeta {
|
|||||||
self.data.validate()?;
|
self.data.validate()?;
|
||||||
}
|
}
|
||||||
|
|
||||||
// 解析meta
|
// 解析 meta
|
||||||
if !meta.is_empty() {
|
if !meta.is_empty() {
|
||||||
let (versions_len, _, meta_ver, meta) = Self::decode_xl_headers(meta)?;
|
let (versions_len, _, meta_ver, meta) = Self::decode_xl_headers(meta)?;
|
||||||
|
|
||||||
@@ -168,7 +168,7 @@ impl FileMeta {
|
|||||||
Ok(i)
|
Ok(i)
|
||||||
}
|
}
|
||||||
|
|
||||||
// decode_xl_headers 解析 meta 头,返回 (versions数量,xl_header_version, xl_meta_version, 已读数据长度)
|
// decode_xl_headers 解析 meta 头,返回 (versions 数量,xl_header_version, xl_meta_version, 已读数据长度)
|
||||||
#[tracing::instrument]
|
#[tracing::instrument]
|
||||||
fn decode_xl_headers(buf: &[u8]) -> Result<(usize, u8, u8, &[u8])> {
|
fn decode_xl_headers(buf: &[u8]) -> Result<(usize, u8, u8, &[u8])> {
|
||||||
let mut cur = Cursor::new(buf);
|
let mut cur = Cursor::new(buf);
|
||||||
@@ -280,7 +280,7 @@ impl FileMeta {
|
|||||||
rmp::encode::write_bin(&mut wr, &ver.meta)?;
|
rmp::encode::write_bin(&mut wr, &ver.meta)?;
|
||||||
}
|
}
|
||||||
|
|
||||||
// 更新bin长度
|
// 更新 bin 长度
|
||||||
let data_len = wr.len() - offset;
|
let data_len = wr.len() - offset;
|
||||||
byteorder::BigEndian::write_u32(&mut wr[offset - 4..offset], data_len as u32);
|
byteorder::BigEndian::write_u32(&mut wr[offset - 4..offset], data_len as u32);
|
||||||
|
|
||||||
@@ -368,7 +368,7 @@ impl FileMeta {
|
|||||||
Err(Error::new(DiskError::FileVersionNotFound))
|
Err(Error::new(DiskError::FileVersionNotFound))
|
||||||
}
|
}
|
||||||
|
|
||||||
// shard_data_dir_count 查询 vid下data_dir的数量
|
// shard_data_dir_count 查询 vid 下 data_dir 的数量
|
||||||
#[tracing::instrument(level = "debug", skip_all)]
|
#[tracing::instrument(level = "debug", skip_all)]
|
||||||
pub fn shard_data_dir_count(&self, vid: &Option<Uuid>, data_dir: &Option<Uuid>) -> usize {
|
pub fn shard_data_dir_count(&self, vid: &Option<Uuid>, data_dir: &Option<Uuid>) -> usize {
|
||||||
self.versions
|
self.versions
|
||||||
@@ -494,7 +494,7 @@ impl FileMeta {
|
|||||||
Err(Error::msg("add_version failed"))
|
Err(Error::msg("add_version failed"))
|
||||||
}
|
}
|
||||||
|
|
||||||
// delete_version 删除版本,返回data_dir
|
// delete_version 删除版本,返回 data_dir
|
||||||
pub fn delete_version(&mut self, fi: &FileInfo) -> Result<Option<Uuid>> {
|
pub fn delete_version(&mut self, fi: &FileInfo) -> Result<Option<Uuid>> {
|
||||||
let mut ventry = FileMetaVersion::default();
|
let mut ventry = FileMetaVersion::default();
|
||||||
if fi.deleted {
|
if fi.deleted {
|
||||||
@@ -710,7 +710,7 @@ impl FileMetaVersion {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// decode_data_dir_from_meta 从 meta中读取data_dir TODO: 直接从meta buf中只解析出data_dir, msg.skip
|
// decode_data_dir_from_meta 从 meta 中读取 data_dir TODO: 直接从 meta buf 中只解析出 data_dir, msg.skip
|
||||||
pub fn decode_data_dir_from_meta(buf: &[u8]) -> Result<Option<Uuid>> {
|
pub fn decode_data_dir_from_meta(buf: &[u8]) -> Result<Option<Uuid>> {
|
||||||
let mut ver = Self::default();
|
let mut ver = Self::default();
|
||||||
ver.unmarshal_msg(buf)?;
|
ver.unmarshal_msg(buf)?;
|
||||||
@@ -733,7 +733,7 @@ impl FileMetaVersion {
|
|||||||
|
|
||||||
// println!("unmarshal_msg fields name len() {}", &str_len);
|
// println!("unmarshal_msg fields name len() {}", &str_len);
|
||||||
|
|
||||||
// !!! Vec::with_capacity(str_len) 失败,vec!正常
|
// !!!Vec::with_capacity(str_len) 失败,vec! 正常
|
||||||
let mut field_buff = vec![0u8; str_len as usize];
|
let mut field_buff = vec![0u8; str_len as usize];
|
||||||
|
|
||||||
cur.read_exact(&mut field_buff)?;
|
cur.read_exact(&mut field_buff)?;
|
||||||
@@ -1143,7 +1143,7 @@ impl From<FileMetaVersion> for FileMetaVersionHeader {
|
|||||||
}
|
}
|
||||||
|
|
||||||
#[derive(Serialize, Deserialize, Debug, Clone, Default, PartialEq)]
|
#[derive(Serialize, Deserialize, Debug, Clone, Default, PartialEq)]
|
||||||
// 因为自定义message_pack,所以一定要保证字段顺序
|
// 因为自定义 message_pack,所以一定要保证字段顺序
|
||||||
pub struct MetaObject {
|
pub struct MetaObject {
|
||||||
pub version_id: Option<Uuid>, // Version ID
|
pub version_id: Option<Uuid>, // Version ID
|
||||||
pub data_dir: Option<Uuid>, // Data dir ID
|
pub data_dir: Option<Uuid>, // Data dir ID
|
||||||
@@ -1182,7 +1182,7 @@ impl MetaObject {
|
|||||||
|
|
||||||
// println!("unmarshal_msg fields name len() {}", &str_len);
|
// println!("unmarshal_msg fields name len() {}", &str_len);
|
||||||
|
|
||||||
// !!! Vec::with_capacity(str_len) 失败,vec!正常
|
// !!!Vec::with_capacity(str_len) 失败,vec! 正常
|
||||||
let mut field_buff = vec![0u8; str_len as usize];
|
let mut field_buff = vec![0u8; str_len as usize];
|
||||||
|
|
||||||
cur.read_exact(&mut field_buff)?;
|
cur.read_exact(&mut field_buff)?;
|
||||||
@@ -1413,7 +1413,7 @@ impl MetaObject {
|
|||||||
|
|
||||||
Ok(cur.position())
|
Ok(cur.position())
|
||||||
}
|
}
|
||||||
// marshal_msg 自定义 messagepack 命名与go一致
|
// marshal_msg 自定义 messagepack 命名与 go 一致
|
||||||
pub fn marshal_msg(&self) -> Result<Vec<u8>> {
|
pub fn marshal_msg(&self) -> Result<Vec<u8>> {
|
||||||
let mut len: u32 = 18;
|
let mut len: u32 = 18;
|
||||||
let mut mask: u32 = 0;
|
let mut mask: u32 = 0;
|
||||||
@@ -1682,7 +1682,7 @@ impl MetaDeleteMarker {
|
|||||||
|
|
||||||
let str_len = rmp::decode::read_str_len(&mut cur)?;
|
let str_len = rmp::decode::read_str_len(&mut cur)?;
|
||||||
|
|
||||||
// !!! Vec::with_capacity(str_len) 失败,vec!正常
|
// !!!Vec::with_capacity(str_len) 失败,vec! 正常
|
||||||
let mut field_buff = vec![0u8; str_len as usize];
|
let mut field_buff = vec![0u8; str_len as usize];
|
||||||
|
|
||||||
cur.read_exact(&mut field_buff)?;
|
cur.read_exact(&mut field_buff)?;
|
||||||
@@ -2175,7 +2175,6 @@ pub async fn read_xl_meta_no_data<R: AsyncRead + Unpin>(reader: &mut R, size: us
|
|||||||
}
|
}
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
mod test {
|
mod test {
|
||||||
|
|
||||||
use super::*;
|
use super::*;
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
@@ -2257,7 +2256,7 @@ mod test {
|
|||||||
|
|
||||||
// println!("obj2 {:?}", &obj2);
|
// println!("obj2 {:?}", &obj2);
|
||||||
|
|
||||||
// 时间截不一致- -
|
// 时间截不一致 - -
|
||||||
assert_eq!(obj, obj2);
|
assert_eq!(obj, obj2);
|
||||||
assert_eq!(obj.get_version_id(), obj2.get_version_id());
|
assert_eq!(obj.get_version_id(), obj2.get_version_id());
|
||||||
assert_eq!(obj.write_version, obj2.write_version);
|
assert_eq!(obj.write_version, obj2.write_version);
|
||||||
@@ -2276,7 +2275,7 @@ mod test {
|
|||||||
let mut obj2 = FileMetaVersionHeader::default();
|
let mut obj2 = FileMetaVersionHeader::default();
|
||||||
obj2.unmarshal_msg(&encoded).unwrap();
|
obj2.unmarshal_msg(&encoded).unwrap();
|
||||||
|
|
||||||
// 时间截不一致- -
|
// 时间截不一致 - -
|
||||||
assert_eq!(obj, obj2);
|
assert_eq!(obj, obj2);
|
||||||
assert_eq!(obj.version_id, obj2.version_id);
|
assert_eq!(obj.version_id, obj2.version_id);
|
||||||
assert_eq!(obj.version_id, vid);
|
assert_eq!(obj.version_id, vid);
|
||||||
|
|||||||
+38
-34
@@ -1643,7 +1643,7 @@ impl SetDisks {
|
|||||||
..Default::default()
|
..Default::default()
|
||||||
},
|
},
|
||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
} else {
|
} else {
|
||||||
Err(Error::new(DiskError::DiskNotFound))
|
Err(Error::new(DiskError::DiskNotFound))
|
||||||
}
|
}
|
||||||
@@ -2232,7 +2232,7 @@ impl SetDisks {
|
|||||||
object,
|
object,
|
||||||
opts.scan_mode,
|
opts.scan_mode,
|
||||||
)
|
)
|
||||||
.await?;
|
.await?;
|
||||||
|
|
||||||
// info!(
|
// info!(
|
||||||
// "disks_with_all_parts: got available_disks: {:?}, data_errs_by_disk: {:?}, data_errs_by_part: {:?}, lastest_meta: {:?}",
|
// "disks_with_all_parts: got available_disks: {:?}, data_errs_by_disk: {:?}, data_errs_by_part: {:?}, lastest_meta: {:?}",
|
||||||
@@ -2378,7 +2378,7 @@ impl SetDisks {
|
|||||||
|
|
||||||
if !lastest_meta.deleted && lastest_meta.erasure.distribution.len() != available_disks.len() {
|
if !lastest_meta.deleted && lastest_meta.erasure.distribution.len() != available_disks.len() {
|
||||||
let err_str = format!("unexpected file distribution ({:?}) from available disks ({:?}), looks like backend disks have been manually modified refusing to heal {}/{}({})",
|
let err_str = format!("unexpected file distribution ({:?}) from available disks ({:?}), looks like backend disks have been manually modified refusing to heal {}/{}({})",
|
||||||
lastest_meta.erasure.distribution, available_disks, bucket, object, version_id);
|
lastest_meta.erasure.distribution, available_disks, bucket, object, version_id);
|
||||||
warn!(err_str);
|
warn!(err_str);
|
||||||
let err = Error::from_string(err_str);
|
let err = Error::from_string(err_str);
|
||||||
return Ok((
|
return Ok((
|
||||||
@@ -2391,7 +2391,7 @@ impl SetDisks {
|
|||||||
let latest_disks = Self::shuffle_disks(&available_disks, &lastest_meta.erasure.distribution);
|
let latest_disks = Self::shuffle_disks(&available_disks, &lastest_meta.erasure.distribution);
|
||||||
if !lastest_meta.deleted && lastest_meta.erasure.distribution.len() != outdate_disks.len() {
|
if !lastest_meta.deleted && lastest_meta.erasure.distribution.len() != outdate_disks.len() {
|
||||||
let err_str = format!("unexpected file distribution ({:?}) from outdated disks ({:?}), looks like backend disks have been manually modified refusing to heal {}/{}({})",
|
let err_str = format!("unexpected file distribution ({:?}) from outdated disks ({:?}), looks like backend disks have been manually modified refusing to heal {}/{}({})",
|
||||||
lastest_meta.erasure.distribution, outdate_disks, bucket, object, version_id);
|
lastest_meta.erasure.distribution, outdate_disks, bucket, object, version_id);
|
||||||
warn!(err_str);
|
warn!(err_str);
|
||||||
let err = Error::from_string(err_str);
|
let err = Error::from_string(err_str);
|
||||||
return Ok((
|
return Ok((
|
||||||
@@ -2403,7 +2403,7 @@ impl SetDisks {
|
|||||||
|
|
||||||
if !lastest_meta.deleted && lastest_meta.erasure.distribution.len() != parts_metadata.len() {
|
if !lastest_meta.deleted && lastest_meta.erasure.distribution.len() != parts_metadata.len() {
|
||||||
let err_str = format!("unexpected file distribution ({:?}) from metadata entries ({:?}), looks like backend disks have been manually modified refusing to heal {}/{}({})",
|
let err_str = format!("unexpected file distribution ({:?}) from metadata entries ({:?}), looks like backend disks have been manually modified refusing to heal {}/{}({})",
|
||||||
lastest_meta.erasure.distribution, parts_metadata.len(), bucket, object, version_id);
|
lastest_meta.erasure.distribution, parts_metadata.len(), bucket, object, version_id);
|
||||||
warn!(err_str);
|
warn!(err_str);
|
||||||
let err = Error::from_string(err_str);
|
let err = Error::from_string(err_str);
|
||||||
return Ok((
|
return Ok((
|
||||||
@@ -2509,7 +2509,7 @@ impl SetDisks {
|
|||||||
DEFAULT_BITROT_ALGO,
|
DEFAULT_BITROT_ALGO,
|
||||||
erasure.shard_size(erasure.block_size),
|
erasure.shard_size(erasure.block_size),
|
||||||
)
|
)
|
||||||
.await?;
|
.await?;
|
||||||
|
|
||||||
writers.push(Some(writer));
|
writers.push(Some(writer));
|
||||||
} else {
|
} else {
|
||||||
@@ -2601,7 +2601,7 @@ impl SetDisks {
|
|||||||
..Default::default()
|
..Default::default()
|
||||||
},
|
},
|
||||||
)
|
)
|
||||||
.await?;
|
.await?;
|
||||||
}
|
}
|
||||||
|
|
||||||
for (i, v) in result.before.drives.iter().enumerate() {
|
for (i, v) in result.before.drives.iter().enumerate() {
|
||||||
@@ -3359,10 +3359,10 @@ impl SetDisks {
|
|||||||
}
|
}
|
||||||
if bucket == RUSTFS_META_BUCKET
|
if bucket == RUSTFS_META_BUCKET
|
||||||
&& (Pattern::new("buckets/*/.metacache/*")
|
&& (Pattern::new("buckets/*/.metacache/*")
|
||||||
.map(|p| p.matches(&entry.name))
|
.map(|p| p.matches(&entry.name))
|
||||||
.unwrap_or(false)
|
.unwrap_or(false)
|
||||||
|| Pattern::new("tmp/.trash/*").map(|p| p.matches(&entry.name)).unwrap_or(false)
|
|| Pattern::new("tmp/.trash/*").map(|p| p.matches(&entry.name)).unwrap_or(false)
|
||||||
|| Pattern::new("multipart/*").map(|p| p.matches(&entry.name)).unwrap_or(false))
|
|| Pattern::new("multipart/*").map(|p| p.matches(&entry.name)).unwrap_or(false))
|
||||||
{
|
{
|
||||||
defer.await;
|
defer.await;
|
||||||
return;
|
return;
|
||||||
@@ -3550,7 +3550,7 @@ impl SetDisks {
|
|||||||
..Default::default()
|
..Default::default()
|
||||||
},
|
},
|
||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
{
|
{
|
||||||
ret_err = Some(err);
|
ret_err = Some(err);
|
||||||
}
|
}
|
||||||
@@ -3601,7 +3601,7 @@ impl SetDisks {
|
|||||||
..Default::default()
|
..Default::default()
|
||||||
},
|
},
|
||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
} else {
|
} else {
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
@@ -3687,7 +3687,7 @@ impl ObjectIO for SetDisks {
|
|||||||
set_index,
|
set_index,
|
||||||
pool_index,
|
pool_index,
|
||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
{
|
{
|
||||||
error!("get_object_with_fileinfo err {:?}", e);
|
error!("get_object_with_fileinfo err {:?}", e);
|
||||||
};
|
};
|
||||||
@@ -3814,7 +3814,7 @@ impl ObjectIO for SetDisks {
|
|||||||
DEFAULT_BITROT_ALGO,
|
DEFAULT_BITROT_ALGO,
|
||||||
erasure.shard_size(erasure.block_size),
|
erasure.shard_size(erasure.block_size),
|
||||||
)
|
)
|
||||||
.await?;
|
.await?;
|
||||||
|
|
||||||
writers.push(Some(writer));
|
writers.push(Some(writer));
|
||||||
} else {
|
} else {
|
||||||
@@ -3880,7 +3880,7 @@ impl ObjectIO for SetDisks {
|
|||||||
object,
|
object,
|
||||||
write_quorum,
|
write_quorum,
|
||||||
)
|
)
|
||||||
.await?;
|
.await?;
|
||||||
|
|
||||||
if let Some(old_dir) = op_old_dir {
|
if let Some(old_dir) = op_old_dir {
|
||||||
self.commit_rename_data_dir(&shuffle_disks, bucket, object, &old_dir.to_string(), write_quorum)
|
self.commit_rename_data_dir(&shuffle_disks, bucket, object, &old_dir.to_string(), write_quorum)
|
||||||
@@ -3984,9 +3984,9 @@ impl StorageAPI for SetDisks {
|
|||||||
if ErasureError::ErasureReadQuorum.is(&err)
|
if ErasureError::ErasureReadQuorum.is(&err)
|
||||||
&& !src_bucket.starts_with(RUSTFS_META_BUCKET)
|
&& !src_bucket.starts_with(RUSTFS_META_BUCKET)
|
||||||
&& self
|
&& self
|
||||||
.delete_if_dang_ling(src_bucket, src_object, &metas, &errs, &HashMap::new(), src_opts.clone())
|
.delete_if_dang_ling(src_bucket, src_object, &metas, &errs, &HashMap::new(), src_opts.clone())
|
||||||
.await
|
.await
|
||||||
.is_ok()
|
.is_ok()
|
||||||
{
|
{
|
||||||
if src_opts.version_id.is_some() {
|
if src_opts.version_id.is_some() {
|
||||||
err = Error::new(DiskError::FileVersionNotFound)
|
err = Error::new(DiskError::FileVersionNotFound)
|
||||||
@@ -4266,7 +4266,7 @@ impl StorageAPI for SetDisks {
|
|||||||
false,
|
false,
|
||||||
false,
|
false,
|
||||||
)
|
)
|
||||||
.await?
|
.await?
|
||||||
} else {
|
} else {
|
||||||
Self::read_all_xl(&disks, bucket, object, false, false).await
|
Self::read_all_xl(&disks, bucket, object, false, false).await
|
||||||
}
|
}
|
||||||
@@ -4278,9 +4278,9 @@ impl StorageAPI for SetDisks {
|
|||||||
if ErasureError::ErasureReadQuorum.is(&err)
|
if ErasureError::ErasureReadQuorum.is(&err)
|
||||||
&& !bucket.starts_with(RUSTFS_META_BUCKET)
|
&& !bucket.starts_with(RUSTFS_META_BUCKET)
|
||||||
&& self
|
&& self
|
||||||
.delete_if_dang_ling(bucket, object, &metas, &errs, &HashMap::new(), opts.clone())
|
.delete_if_dang_ling(bucket, object, &metas, &errs, &HashMap::new(), opts.clone())
|
||||||
.await
|
.await
|
||||||
.is_ok()
|
.is_ok()
|
||||||
{
|
{
|
||||||
if opts.version_id.is_some() {
|
if opts.version_id.is_some() {
|
||||||
err = Error::new(DiskError::FileVersionNotFound)
|
err = Error::new(DiskError::FileVersionNotFound)
|
||||||
@@ -4446,7 +4446,7 @@ impl StorageAPI for SetDisks {
|
|||||||
DEFAULT_BITROT_ALGO,
|
DEFAULT_BITROT_ALGO,
|
||||||
shared_size,
|
shared_size,
|
||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
{
|
{
|
||||||
Ok(writer) => Ok(Some(writer)),
|
Ok(writer) => Ok(Some(writer)),
|
||||||
Err(e) => Err(e),
|
Err(e) => Err(e),
|
||||||
@@ -4502,7 +4502,7 @@ impl StorageAPI for SetDisks {
|
|||||||
fi_buff,
|
fi_buff,
|
||||||
write_quorum,
|
write_quorum,
|
||||||
)
|
)
|
||||||
.await?;
|
.await?;
|
||||||
|
|
||||||
let ret: PartInfo = PartInfo {
|
let ret: PartInfo = PartInfo {
|
||||||
etag: Some(etag.clone()),
|
etag: Some(etag.clone()),
|
||||||
@@ -4758,8 +4758,8 @@ impl StorageAPI for SetDisks {
|
|||||||
&parts_metadatas,
|
&parts_metadatas,
|
||||||
write_quorum,
|
write_quorum,
|
||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
.map_err(|e| to_object_err(e, vec![bucket, object]))?;
|
.map_err(|e| to_object_err(e, vec![bucket, object]))?;
|
||||||
|
|
||||||
// evalDisks
|
// evalDisks
|
||||||
|
|
||||||
@@ -5007,7 +5007,7 @@ impl StorageAPI for SetDisks {
|
|||||||
object,
|
object,
|
||||||
write_quorum,
|
write_quorum,
|
||||||
)
|
)
|
||||||
.await?;
|
.await?;
|
||||||
|
|
||||||
for (i, op_disk) in online_disks.iter().enumerate() {
|
for (i, op_disk) in online_disks.iter().enumerate() {
|
||||||
if let Some(disk) = op_disk {
|
if let Some(disk) = op_disk {
|
||||||
@@ -5428,7 +5428,7 @@ async fn disks_with_all_parts(
|
|||||||
checksum_info.hash,
|
checksum_info.hash,
|
||||||
meta.erasure.shard_size(meta.erasure.block_size),
|
meta.erasure.shard_size(meta.erasure.block_size),
|
||||||
)
|
)
|
||||||
.await)
|
.await)
|
||||||
.err();
|
.err();
|
||||||
|
|
||||||
if let Some(vec) = data_errs_by_part.get_mut(&0) {
|
if let Some(vec) = data_errs_by_part.get_mut(&0) {
|
||||||
@@ -5890,7 +5890,7 @@ mod tests {
|
|||||||
erasure: ErasureInfo {
|
erasure: ErasureInfo {
|
||||||
data_blocks: 4,
|
data_blocks: 4,
|
||||||
parity_blocks: 2,
|
parity_blocks: 2,
|
||||||
index: 1, // Must be > 0 for is_valid() to return true
|
index: 1, // Must be > 0 for is_valid() to return true
|
||||||
distribution: vec![1, 2, 3, 4, 5, 6], // Must match data_blocks + parity_blocks
|
distribution: vec![1, 2, 3, 4, 5, 6], // Must match data_blocks + parity_blocks
|
||||||
..Default::default()
|
..Default::default()
|
||||||
},
|
},
|
||||||
@@ -5902,7 +5902,7 @@ mod tests {
|
|||||||
erasure: ErasureInfo {
|
erasure: ErasureInfo {
|
||||||
data_blocks: 6,
|
data_blocks: 6,
|
||||||
parity_blocks: 3,
|
parity_blocks: 3,
|
||||||
index: 1, // Must be > 0 for is_valid() to return true
|
index: 1, // Must be > 0 for is_valid() to return true
|
||||||
distribution: vec![1, 2, 3, 4, 5, 6, 7, 8, 9], // Must match data_blocks + parity_blocks
|
distribution: vec![1, 2, 3, 4, 5, 6, 7, 8, 9], // Must match data_blocks + parity_blocks
|
||||||
..Default::default()
|
..Default::default()
|
||||||
},
|
},
|
||||||
@@ -5914,7 +5914,7 @@ mod tests {
|
|||||||
erasure: ErasureInfo {
|
erasure: ErasureInfo {
|
||||||
data_blocks: 2,
|
data_blocks: 2,
|
||||||
parity_blocks: 1,
|
parity_blocks: 1,
|
||||||
index: 1, // Must be > 0 for is_valid() to return true
|
index: 1, // Must be > 0 for is_valid() to return true
|
||||||
distribution: vec![1, 2, 3], // Must match data_blocks + parity_blocks
|
distribution: vec![1, 2, 3], // Must match data_blocks + parity_blocks
|
||||||
..Default::default()
|
..Default::default()
|
||||||
},
|
},
|
||||||
@@ -6019,7 +6019,11 @@ mod tests {
|
|||||||
#[test]
|
#[test]
|
||||||
fn test_join_errs() {
|
fn test_join_errs() {
|
||||||
// Test joining error messages
|
// Test joining error messages
|
||||||
let errs = vec![None, Some(Error::from_string("error1")), Some(Error::from_string("error2"))];
|
let errs = vec![
|
||||||
|
None,
|
||||||
|
Some(Error::from_string("error1")),
|
||||||
|
Some(Error::from_string("error2")),
|
||||||
|
];
|
||||||
let joined = join_errs(&errs);
|
let joined = join_errs(&errs);
|
||||||
assert!(joined.contains("<nil>"));
|
assert!(joined.contains("<nil>"));
|
||||||
assert!(joined.contains("error1"));
|
assert!(joined.contains("error1"));
|
||||||
@@ -6094,4 +6098,4 @@ mod tests {
|
|||||||
assert_eq!(result2.len(), 3);
|
assert_eq!(result2.len(), 3);
|
||||||
assert!(result2.iter().all(|d| d.is_none()));
|
assert!(result2.iter().all(|d| d.is_none()));
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -229,7 +229,7 @@ impl FileInfo {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// to_part_offset 取offset 所在的part index, 返回part index, offset
|
// `to_part_offset` takes the `part index` where the `offset` is located, and returns `part index`, `offset`
|
||||||
pub fn to_part_offset(&self, offset: usize) -> Result<(usize, usize)> {
|
pub fn to_part_offset(&self, offset: usize) -> Result<(usize, usize)> {
|
||||||
if offset == 0 {
|
if offset == 0 {
|
||||||
return Ok((0, 0));
|
return Ok((0, 0));
|
||||||
@@ -356,7 +356,7 @@ impl ErasureInfo {
|
|||||||
let last_shard_size = last_block_size.div_ceil(self.data_blocks);
|
let last_shard_size = last_block_size.div_ceil(self.data_blocks);
|
||||||
num_shards * self.shard_size(self.block_size) + last_shard_size
|
num_shards * self.shard_size(self.block_size) + last_shard_size
|
||||||
|
|
||||||
// // 因为写入的时候ec需要补全,所以最后一个长度应该也是一样的
|
// // 因为写入的时候 ec 需要补全,所以最后一个长度应该也是一样的
|
||||||
// if last_block_size != 0 {
|
// if last_block_size != 0 {
|
||||||
// num_shards += 1
|
// num_shards += 1
|
||||||
// }
|
// }
|
||||||
@@ -1250,7 +1250,7 @@ mod tests {
|
|||||||
assert_eq!(object_info.etag, Some("test-etag".to_string()));
|
assert_eq!(object_info.etag, Some("test-etag".to_string()));
|
||||||
}
|
}
|
||||||
|
|
||||||
// to_part_offset 取offset 所在的part index, 返回part index, offset
|
// to_part_offset 取 offset 所在的 part index, 返回 part index, offset
|
||||||
#[test]
|
#[test]
|
||||||
fn test_file_info_to_part_offset() {
|
fn test_file_info_to_part_offset() {
|
||||||
let mut file_info = FileInfo::new("test", 4, 2);
|
let mut file_info = FileInfo::new("test", 4, 2);
|
||||||
|
|||||||
Reference in New Issue
Block a user