refactor: Reimplement bucket replication system with enhanced architecture (#590)

* feat:refactor replication

* use aws sdk for replication client

* refactor/replication

* merge main

* fix lifecycle test
This commit is contained in:
weisd
2025-09-26 14:27:53 +08:00
committed by GitHub
parent 9b029d18b2
commit 90f21a9102
91 changed files with 10532 additions and 4917 deletions
+163 -115
View File
@@ -34,7 +34,9 @@ use crate::global::{
use crate::notification_sys::get_global_notification_sys;
use crate::pools::PoolMeta;
use crate::rebalance::RebalanceMeta;
use crate::store_api::{ListMultipartsInfo, ListObjectVersionsInfo, ListPartsInfo, MultipartInfo, ObjectIO};
use crate::store_api::{
ListMultipartsInfo, ListObjectVersionsInfo, ListPartsInfo, MultipartInfo, ObjectIO, ObjectInfoOrErr, WalkOptions,
};
use crate::store_init::{check_disk_fatal_errs, ec_drives_no_config};
use crate::{
bucket::{lifecycle::bucket_lifecycle_ops::TransitionState, metadata::BucketMetadata},
@@ -68,8 +70,9 @@ use std::time::SystemTime;
use std::{collections::HashMap, sync::Arc, time::Duration};
use time::OffsetDateTime;
use tokio::select;
use tokio::sync::{RwLock, broadcast};
use tokio::sync::RwLock;
use tokio::time::sleep;
use tokio_util::sync::CancellationToken;
use tracing::{debug, info};
use tracing::{error, warn};
use uuid::Uuid;
@@ -109,7 +112,7 @@ pub struct ECStore {
impl ECStore {
#[allow(clippy::new_ret_no_self)]
#[tracing::instrument(level = "debug", skip(endpoint_pools))]
pub async fn new(address: SocketAddr, endpoint_pools: EndpointServerPools) -> Result<Arc<Self>> {
pub async fn new(address: SocketAddr, endpoint_pools: EndpointServerPools, ctx: CancellationToken) -> Result<Arc<Self>> {
// let layouts = DisksLayout::from_volumes(endpoints.as_slice())?;
let mut deployment_id = None;
@@ -251,7 +254,7 @@ impl ECStore {
let wait_sec = 5;
let mut exit_count = 0;
loop {
if let Err(err) = ec.init().await {
if let Err(err) = ec.init(ctx.clone()).await {
error!("init err: {}", err);
error!("retry after {} second", wait_sec);
sleep(Duration::from_secs(wait_sec)).await;
@@ -273,7 +276,7 @@ impl ECStore {
Ok(ec)
}
pub async fn init(self: &Arc<Self>) -> Result<()> {
pub async fn init(self: &Arc<Self>, rx: CancellationToken) -> Result<()> {
GLOBAL_BOOT_TIME.get_or_init(|| async { SystemTime::now() }).await;
if self.load_rebalance_meta().await.is_ok() {
@@ -317,18 +320,16 @@ impl ECStore {
if !pool_indices.is_empty() {
let idx = pool_indices[0];
if endpoints.as_ref()[idx].endpoints.as_ref()[0].is_local {
let (_tx, rx) = broadcast::channel(1);
let store = self.clone();
tokio::spawn(async move {
// wait 3 minutes for cluster init
tokio::time::sleep(Duration::from_secs(60 * 3)).await;
if let Err(err) = store.decommission(rx.resubscribe(), pool_indices.clone()).await {
if let Err(err) = store.decommission(rx.clone(), pool_indices.clone()).await {
if err == StorageError::DecommissionAlreadyRunning {
for i in pool_indices.iter() {
store.do_decommission_in_routine(rx.resubscribe(), *i).await;
store.do_decommission_in_routine(rx.clone(), *i).await;
}
return;
}
@@ -700,9 +701,13 @@ impl ECStore {
opts: &ObjectOptions,
) -> Result<(PoolObjInfo, Vec<PoolErr>)> {
let mut futures = Vec::new();
for pool in self.pools.iter() {
futures.push(pool.get_object_info(bucket, object, opts));
let mut pool_opts = opts.clone();
if !pool_opts.metadata_chg {
pool_opts.version_id = None;
}
futures.push(async move { pool.get_object_info(bucket, object, &pool_opts).await });
}
let results = join_all(futures).await;
@@ -1351,6 +1356,17 @@ impl StorageAPI for ECStore {
.await
}
async fn walk(
self: Arc<Self>,
rx: CancellationToken,
bucket: &str,
prefix: &str,
result: tokio::sync::mpsc::Sender<ObjectInfoOrErr>,
opts: WalkOptions,
) -> Result<()> {
self.walk_internal(rx, bucket, prefix, result, opts).await
}
#[tracing::instrument(skip(self))]
async fn get_object_info(&self, bucket: &str, object: &str, opts: &ObjectOptions) -> Result<ObjectInfo> {
check_object_args(bucket, object)?;
@@ -1450,9 +1466,12 @@ impl StorageAPI for ECStore {
let object = encode_dir_object(object);
let object = object.as_str();
let mut gopts = opts.clone();
gopts.no_lock = true;
// 查询在哪个 pool
let (mut pinfo, errs) = self
.get_pool_info_existing_with_opts(bucket, object, &opts)
.get_pool_info_existing_with_opts(bucket, object, &gopts)
.await
.map_err(|e| {
if is_err_read_quorum(&e) {
@@ -1513,7 +1532,7 @@ impl StorageAPI for ECStore {
bucket: &str,
objects: Vec<ObjectToDelete>,
opts: ObjectOptions,
) -> Result<(Vec<DeletedObject>, Vec<Option<Error>>)> {
) -> (Vec<DeletedObject>, Vec<Option<Error>>) {
// encode object name
let objects: Vec<ObjectToDelete> = objects
.iter()
@@ -1534,131 +1553,160 @@ impl StorageAPI for ECStore {
// TODO: nslock
let mut futures = Vec::with_capacity(objects.len());
let mut futures = Vec::with_capacity(self.pools.len());
for obj in objects.iter() {
futures.push(async move {
self.internal_get_pool_info_existing_with_opts(
bucket,
&obj.object_name,
&ObjectOptions {
no_lock: true,
..Default::default()
},
)
.await
});
for pool in self.pools.iter() {
futures.push(pool.delete_objects(bucket, objects.clone(), opts.clone()));
}
let results = join_all(futures).await;
// let mut jhs = Vec::new();
// let semaphore = Arc::new(Semaphore::new(num_cpus::get()));
// let pools = Arc::new(self.pools.clone());
for idx in 0..del_objects.len() {
for (dels, errs) in results.iter() {
if errs[idx].is_none() && dels[idx].found {
del_errs[idx] = None;
del_objects[idx] = dels[idx].clone();
break;
}
if del_errs[idx].is_none() {
del_errs[idx] = errs[idx].clone();
del_objects[idx] = dels[idx].clone();
}
}
}
del_objects.iter_mut().for_each(|v| {
v.object_name = decode_dir_object(&v.object_name);
});
(del_objects, del_errs)
// let mut futures = Vec::with_capacity(objects.len());
// for obj in objects.iter() {
// let (semaphore, pools, bucket, object_name, opt) = (
// semaphore.clone(),
// pools.clone(),
// bucket.to_string(),
// obj.object_name.to_string(),
// ObjectOptions::default(),
// );
// let jh = tokio::spawn(async move {
// let _permit = semaphore.acquire().await.unwrap();
// self.internal_get_pool_info_existing_with_opts(pools.as_ref(), &bucket, &object_name, &opt)
// .await
// futures.push(async move {
// self.internal_get_pool_info_existing_with_opts(
// bucket,
// &obj.object_name,
// &ObjectOptions {
// no_lock: true,
// ..Default::default()
// },
// )
// .await
// });
// jhs.push(jh);
// }
// let mut results = Vec::new();
// for jh in jhs {
// results.push(jh.await.unwrap());
// }
// 记录 pool Index 对应的 objects pool_idx -> objects idx
let mut pool_obj_idx_map = HashMap::new();
let mut orig_index_map = HashMap::new();
// let results = join_all(futures).await;
for (i, res) in results.into_iter().enumerate() {
match res {
Ok((pinfo, _)) => {
if let Some(obj) = objects.get(i) {
if pinfo.object_info.delete_marker && obj.version_id.is_none() {
del_objects[i] = DeletedObject {
delete_marker: pinfo.object_info.delete_marker,
delete_marker_version_id: pinfo.object_info.version_id.map(|v| v.to_string()),
object_name: decode_dir_object(&pinfo.object_info.name),
delete_marker_mtime: pinfo.object_info.mod_time,
..Default::default()
};
continue;
}
// // let mut jhs = Vec::new();
// // let semaphore = Arc::new(Semaphore::new(num_cpus::get()));
// // let pools = Arc::new(self.pools.clone());
if !pool_obj_idx_map.contains_key(&pinfo.index) {
pool_obj_idx_map.insert(pinfo.index, vec![obj.clone()]);
} else if let Some(val) = pool_obj_idx_map.get_mut(&pinfo.index) {
val.push(obj.clone());
}
// // for obj in objects.iter() {
// // let (semaphore, pools, bucket, object_name, opt) = (
// // semaphore.clone(),
// // pools.clone(),
// // bucket.to_string(),
// // obj.object_name.to_string(),
// // ObjectOptions::default(),
// // );
if !orig_index_map.contains_key(&pinfo.index) {
orig_index_map.insert(pinfo.index, vec![i]);
} else if let Some(val) = orig_index_map.get_mut(&pinfo.index) {
val.push(i);
}
}
}
Err(e) => {
if !is_err_object_not_found(&e) && is_err_version_not_found(&e) {
del_errs[i] = Some(e)
}
// // let jh = tokio::spawn(async move {
// // let _permit = semaphore.acquire().await.unwrap();
// // self.internal_get_pool_info_existing_with_opts(pools.as_ref(), &bucket, &object_name, &opt)
// // .await
// // });
// // jhs.push(jh);
// // }
// // let mut results = Vec::new();
// // for jh in jhs {
// // results.push(jh.await.unwrap());
// // }
if let Some(obj) = objects.get(i) {
del_objects[i] = DeletedObject {
object_name: decode_dir_object(&obj.object_name),
version_id: obj.version_id.map(|v| v.to_string()),
..Default::default()
}
}
}
}
}
// // 记录 pool Index 对应的 objects pool_idx -> objects idx
// let mut pool_obj_idx_map = HashMap::new();
// let mut orig_index_map = HashMap::new();
if !pool_obj_idx_map.is_empty() {
for (i, sets) in self.pools.iter().enumerate() {
// 取 pool idx 对应的 objects index
if let Some(objs) = pool_obj_idx_map.get(&i) {
// 取对应 obj,理论上不会 none
// let objs: Vec<ObjectToDelete> = obj_idxs.iter().filter_map(|&idx| objects.get(idx).cloned()).collect();
// for (i, res) in results.into_iter().enumerate() {
// match res {
// Ok((pinfo, _)) => {
// if let Some(obj) = objects.get(i) {
// if pinfo.object_info.delete_marker && obj.version_id.is_none() {
// del_objects[i] = DeletedObject {
// delete_marker: pinfo.object_info.delete_marker,
// delete_marker_version_id: pinfo.object_info.version_id.map(|v| v.to_string()),
// object_name: decode_dir_object(&pinfo.object_info.name),
// delete_marker_mtime: pinfo.object_info.mod_time,
// ..Default::default()
// };
// continue;
// }
if objs.is_empty() {
continue;
}
// if !pool_obj_idx_map.contains_key(&pinfo.index) {
// pool_obj_idx_map.insert(pinfo.index, vec![obj.clone()]);
// } else if let Some(val) = pool_obj_idx_map.get_mut(&pinfo.index) {
// val.push(obj.clone());
// }
let (pdel_objs, perrs) = sets.delete_objects(bucket, objs.clone(), opts.clone()).await?;
// if !orig_index_map.contains_key(&pinfo.index) {
// orig_index_map.insert(pinfo.index, vec![i]);
// } else if let Some(val) = orig_index_map.get_mut(&pinfo.index) {
// val.push(i);
// }
// }
// }
// Err(e) => {
// if !is_err_object_not_found(&e) && is_err_version_not_found(&e) {
// del_errs[i] = Some(e)
// }
// 同时存入不可能为 none
let org_indexes = orig_index_map.get(&i).unwrap();
// if let Some(obj) = objects.get(i) {
// del_objects[i] = DeletedObject {
// object_name: decode_dir_object(&obj.object_name),
// version_id: obj.version_id.map(|v| v.to_string()),
// ..Default::default()
// }
// }
// }
// }
// }
// perrs 的顺序理论上跟 obj_idxs 顺序一致
for (i, err) in perrs.into_iter().enumerate() {
let obj_idx = org_indexes[i];
// if !pool_obj_idx_map.is_empty() {
// for (i, sets) in self.pools.iter().enumerate() {
// // 取 pool idx 对应的 objects index
// if let Some(objs) = pool_obj_idx_map.get(&i) {
// // 取对应 obj,理论上不会 none
// // let objs: Vec<ObjectToDelete> = obj_idxs.iter().filter_map(|&idx| objects.get(idx).cloned()).collect();
if err.is_some() {
del_errs[obj_idx] = err;
}
// if objs.is_empty() {
// continue;
// }
let mut dobj = pdel_objs.get(i).unwrap().clone();
dobj.object_name = decode_dir_object(&dobj.object_name);
// let (pdel_objs, perrs) = sets.delete_objects(bucket, objs.clone(), opts.clone()).await?;
del_objects[obj_idx] = dobj;
}
}
}
}
// // 同时存入不可能为 none
// let org_indexes = orig_index_map.get(&i).unwrap();
Ok((del_objects, del_errs))
// // perrs 的顺序理论上跟 obj_idxs 顺序一致
// for (i, err) in perrs.into_iter().enumerate() {
// let obj_idx = org_indexes[i];
// if err.is_some() {
// del_errs[obj_idx] = err;
// }
// let mut dobj = pdel_objs.get(i).unwrap().clone();
// dobj.object_name = decode_dir_object(&dobj.object_name);
// del_objects[obj_idx] = dobj;
// }
// }
// }
// }
// Ok((del_objects, del_errs))
}
#[tracing::instrument(skip(self))]