mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-22 20:36:38 +00:00
Compare commits
4 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| e14d87b70e | |||
| da90d02c15 | |||
| a930152d5a | |||
| 4ad27860e3 |
@@ -3425,6 +3425,44 @@ mod tests {
|
||||
assert!(mutexes.contains_key("second"));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn update_all_targets_publishes_disable_proxy_on_target_client() {
|
||||
// The read-proxy selector (replication_proxy::get_proxy_targets) skips
|
||||
// targets whose TargetClient carries disable_proxy — the persisted
|
||||
// per-target opt-out must survive client publication.
|
||||
let sys = BucketTargetSys::default();
|
||||
let target = |arn: &str, disable_proxy: bool| BucketTarget {
|
||||
arn: arn.to_string(),
|
||||
endpoint: "192.168.1.10:9000".to_string(),
|
||||
target_bucket: "target-bucket".to_string(),
|
||||
region: "us-east-1".to_string(),
|
||||
disable_proxy,
|
||||
credentials: Some(Credentials {
|
||||
access_key: "access".to_string(),
|
||||
secret_key: "secret".to_string(),
|
||||
session_token: None,
|
||||
expiration: None,
|
||||
}),
|
||||
..Default::default()
|
||||
};
|
||||
let targets = BucketTargets {
|
||||
targets: vec![target("arn:proxied", false), target("arn:opted-out", true)],
|
||||
};
|
||||
|
||||
sys.update_all_targets("bucket", Some(&targets)).await;
|
||||
|
||||
let proxied = sys
|
||||
.get_remote_target_client("bucket", "arn:proxied")
|
||||
.await
|
||||
.expect("client should be published");
|
||||
assert!(!proxied.disable_proxy);
|
||||
let opted_out = sys
|
||||
.get_remote_target_client("bucket", "arn:opted-out")
|
||||
.await
|
||||
.expect("client should be published");
|
||||
assert!(opted_out.disable_proxy, "disable_proxy must reach the published TargetClient");
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
|
||||
async fn target_updates_serialize_client_build_through_publication_per_bucket() {
|
||||
let sys = Arc::new(BucketTargetSys::default());
|
||||
|
||||
@@ -865,16 +865,22 @@ fn should_replace_pool_status_for_status_refresh(
|
||||
!has_active_worker && persisted.last_update > current.last_update
|
||||
}
|
||||
|
||||
fn merge_pool_status_refresh(current: &mut PoolMeta, persisted: PoolMeta, active_workers: &[bool]) {
|
||||
/// Merges a persisted pool metadata snapshot into `current` monotonically:
|
||||
/// a pool entry is replaced only when no active worker covers it and the
|
||||
/// snapshot is strictly newer, so delayed snapshots never roll back local
|
||||
/// queued/terminal progressions. Returns whether any entry was replaced or
|
||||
/// appended.
|
||||
pub(crate) fn merge_pool_status_refresh(current: &mut PoolMeta, persisted: PoolMeta, active_workers: &[bool]) -> bool {
|
||||
if persisted.pools.is_empty() {
|
||||
return;
|
||||
return false;
|
||||
}
|
||||
|
||||
if current.pools.is_empty() {
|
||||
*current = persisted;
|
||||
return;
|
||||
return true;
|
||||
}
|
||||
|
||||
let mut merged_newer = false;
|
||||
for (idx, persisted_pool) in persisted.pools.into_iter().enumerate() {
|
||||
if persisted_pool.id != idx {
|
||||
continue;
|
||||
@@ -884,11 +890,14 @@ fn merge_pool_status_refresh(current: &mut PoolMeta, persisted: PoolMeta, active
|
||||
if idx < current.pools.len() {
|
||||
if should_replace_pool_status_for_status_refresh(current.pools.get(idx), &persisted_pool, has_active_worker) {
|
||||
current.pools[idx] = persisted_pool;
|
||||
merged_newer = true;
|
||||
}
|
||||
} else if idx == current.pools.len() && !has_active_worker {
|
||||
current.pools.push(persisted_pool);
|
||||
merged_newer = true;
|
||||
}
|
||||
}
|
||||
merged_newer
|
||||
}
|
||||
|
||||
fn resolve_start_decommission_pool_meta_reload_result(result: Result<()>) -> Result<()> {
|
||||
@@ -5707,6 +5716,55 @@ mod pools_tests {
|
||||
assert_eq!(info.bytes_done, 1_024);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_merge_pool_status_refresh_fails_closed_on_missing_persisted_pools() {
|
||||
let newer = OffsetDateTime::from_unix_timestamp(2_000).expect("test timestamp should be valid");
|
||||
let mut current = PoolMeta {
|
||||
pools: vec![decommission_test_pool_status(
|
||||
0,
|
||||
Some(PoolDecommissionInfo {
|
||||
complete: true,
|
||||
..Default::default()
|
||||
}),
|
||||
)],
|
||||
..Default::default()
|
||||
};
|
||||
current.pools[0].last_update = newer;
|
||||
|
||||
assert!(
|
||||
!merge_pool_status_refresh(&mut current, PoolMeta::default(), &[false]),
|
||||
"an empty persisted snapshot must fail closed instead of replacing local state"
|
||||
);
|
||||
|
||||
let info = current.pools[0]
|
||||
.decommission
|
||||
.as_ref()
|
||||
.expect("local decommission info should survive a missing snapshot");
|
||||
assert!(info.complete);
|
||||
assert_eq!(current.pools[0].last_update, newer);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_merge_pool_status_refresh_ignores_mislabeled_pool_entries() {
|
||||
let older = OffsetDateTime::from_unix_timestamp(1_000).expect("test timestamp should be valid");
|
||||
let mut current = PoolMeta {
|
||||
pools: vec![decommission_test_pool_status(0, None)],
|
||||
..Default::default()
|
||||
};
|
||||
let mut persisted = PoolMeta {
|
||||
pools: vec![decommission_test_pool_status(0, Some(PoolDecommissionInfo::default()))],
|
||||
..Default::default()
|
||||
};
|
||||
persisted.pools[0].id = 7;
|
||||
persisted.pools[0].last_update = older;
|
||||
|
||||
assert!(
|
||||
!merge_pool_status_refresh(&mut current, persisted, &[false]),
|
||||
"a pool entry whose id does not match its index must be ignored"
|
||||
);
|
||||
assert!(current.pools[0].decommission.is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_dedup_indices_removes_duplicates_preserving_order() {
|
||||
assert_eq!(dedup_indices(&[0, 2, 1, 2, 3, 0]), vec![0, 2, 1, 3]);
|
||||
|
||||
@@ -297,10 +297,16 @@ impl ECStore {
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use crate::bucket::metadata_sys;
|
||||
use crate::core::pools::{PoolDecommissionInfo, PoolStatus};
|
||||
use crate::disk::{DiskOption, format::FormatV3, new_disk};
|
||||
use crate::layout::endpoints::{Endpoints, PoolEndpoints};
|
||||
use crate::disk::{DeleteOptions, DiskOption, format::FormatV3, new_disk};
|
||||
use crate::layout::endpoints::{EndpointServerPools, Endpoints, PoolEndpoints};
|
||||
use crate::runtime::instance::InstanceContext;
|
||||
use crate::storage_api_contracts::bucket::{BucketOperations, MakeBucketOptions};
|
||||
use crate::storage_api_contracts::object::{ObjectIO as _, ObjectOperations};
|
||||
use crate::store::init_format::{load_format_erasure, save_format_file};
|
||||
use crate::store::init_local_disks_with_instance_ctx;
|
||||
use tokio_util::sync::CancellationToken;
|
||||
|
||||
async fn minimal_heal_pool(pool_idx: usize) -> Arc<Sets> {
|
||||
let format = FormatV3::new(1, 1);
|
||||
@@ -347,6 +353,51 @@ mod tests {
|
||||
}
|
||||
}
|
||||
|
||||
async fn multi_pool_heal_store() -> (tempfile::TempDir, Arc<ECStore>, CancellationToken) {
|
||||
let temp_dir = tempfile::tempdir().expect("multi-pool heal test directory should be created");
|
||||
let mut pool_endpoints = Vec::new();
|
||||
for pool_index in 0..2 {
|
||||
let mut endpoints = Vec::new();
|
||||
for disk_index in 0..4 {
|
||||
let disk_path = temp_dir.path().join(format!("pool{pool_index}-disk{disk_index}"));
|
||||
tokio::fs::create_dir_all(&disk_path)
|
||||
.await
|
||||
.expect("multi-pool heal test disk should be created");
|
||||
let mut endpoint = Endpoint::try_from(disk_path.to_str().expect("disk path should be utf8"))
|
||||
.expect("test endpoint should parse");
|
||||
endpoint.set_pool_index(pool_index);
|
||||
endpoint.set_set_index(0);
|
||||
endpoint.set_disk_index(disk_index);
|
||||
endpoints.push(endpoint);
|
||||
}
|
||||
pool_endpoints.push(PoolEndpoints {
|
||||
legacy: false,
|
||||
set_count: 1,
|
||||
drives_per_set: 4,
|
||||
endpoints: Endpoints::from(endpoints),
|
||||
cmd_line: format!("heal-owner-pool-{pool_index}"),
|
||||
platform: "test".to_string(),
|
||||
});
|
||||
}
|
||||
|
||||
let endpoint_pools = EndpointServerPools::from(pool_endpoints);
|
||||
let instance_ctx = Arc::new(InstanceContext::new());
|
||||
init_local_disks_with_instance_ctx(&instance_ctx, endpoint_pools.clone())
|
||||
.await
|
||||
.expect("multi-pool local disks should initialize");
|
||||
let shutdown = CancellationToken::new();
|
||||
let store = ECStore::new_with_instance_ctx(
|
||||
"127.0.0.1:0".parse().expect("test address should parse"),
|
||||
endpoint_pools,
|
||||
shutdown.clone(),
|
||||
instance_ctx,
|
||||
)
|
||||
.await
|
||||
.expect("multi-pool test store should initialize");
|
||||
metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
|
||||
(temp_dir, store, shutdown)
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn heal_object_pool_scope_selects_only_requested_pool() {
|
||||
let store = minimal_heal_store().await;
|
||||
@@ -506,6 +557,229 @@ mod tests {
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial_test::serial]
|
||||
async fn unscoped_heal_object_suspended_owner_semantics() {
|
||||
let (_temp_dir, store, shutdown) = multi_pool_heal_store().await;
|
||||
let bucket = format!("heal-owner-{}", Uuid::new_v4().simple());
|
||||
let active_object = "active-owner";
|
||||
let suspended_only_object = "suspended-only";
|
||||
let duplicate_object = "duplicate-owner";
|
||||
let marker_object = "marker-owner";
|
||||
let quorum_object = "quorum-owner";
|
||||
store
|
||||
.make_bucket(&bucket, &MakeBucketOptions::default())
|
||||
.await
|
||||
.expect("bucket should be created in all pools");
|
||||
|
||||
let mut active_reader = PutObjReader::from_vec(b"active owner".to_vec());
|
||||
store.pools[0]
|
||||
.put_object(&bucket, active_object, &mut active_reader, &ObjectOptions::default())
|
||||
.await
|
||||
.expect("active owner object should be written");
|
||||
let active_disks = store.pools[0].disk_set[0].disks.read().await.clone();
|
||||
let missing_active_disk = active_disks[0].clone().expect("active disk should be online");
|
||||
missing_active_disk
|
||||
.delete(
|
||||
&bucket,
|
||||
active_object,
|
||||
DeleteOptions {
|
||||
recursive: true,
|
||||
immediate: true,
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await
|
||||
.expect("active owner shard should be removed for repair");
|
||||
assert!(
|
||||
missing_active_disk.read_xl(&bucket, active_object, false).await.is_err(),
|
||||
"the active owner fixture must start with one missing metadata copy"
|
||||
);
|
||||
|
||||
let mut suspended_reader = PutObjReader::from_vec(b"suspended owner".to_vec());
|
||||
store.pools[1]
|
||||
.put_object(&bucket, suspended_only_object, &mut suspended_reader, &ObjectOptions::default())
|
||||
.await
|
||||
.expect("suspended owner object should be written");
|
||||
for (pool_index, mod_time) in [1_i64, 2_i64].into_iter().enumerate() {
|
||||
let mut duplicate_reader = PutObjReader::from_vec(format!("duplicate-pool-{pool_index}").into_bytes());
|
||||
store.pools[pool_index]
|
||||
.put_object(
|
||||
&bucket,
|
||||
duplicate_object,
|
||||
&mut duplicate_reader,
|
||||
&ObjectOptions {
|
||||
mod_time: Some(OffsetDateTime::UNIX_EPOCH + time::Duration::seconds(mod_time)),
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await
|
||||
.expect("duplicate owner object should be written");
|
||||
}
|
||||
let duplicate_missing_disk = store.pools[0].disk_set[0].disks.read().await[0]
|
||||
.clone()
|
||||
.expect("duplicate active owner disk should be online");
|
||||
duplicate_missing_disk
|
||||
.delete(
|
||||
&bucket,
|
||||
duplicate_object,
|
||||
DeleteOptions {
|
||||
recursive: true,
|
||||
immediate: true,
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await
|
||||
.expect("duplicate active owner shard should be removed for repair");
|
||||
let history_version = Uuid::new_v4();
|
||||
let mut history_reader = PutObjReader::from_vec(b"marker history".to_vec());
|
||||
store.pools[0]
|
||||
.put_object(
|
||||
&bucket,
|
||||
marker_object,
|
||||
&mut history_reader,
|
||||
&ObjectOptions {
|
||||
versioned: true,
|
||||
version_id: Some(history_version.to_string()),
|
||||
mod_time: Some(OffsetDateTime::UNIX_EPOCH + time::Duration::seconds(1)),
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await
|
||||
.expect("versioned marker history should be written");
|
||||
store.pools[0]
|
||||
.delete_object(
|
||||
&bucket,
|
||||
marker_object,
|
||||
ObjectOptions {
|
||||
versioned: true,
|
||||
mod_time: Some(OffsetDateTime::UNIX_EPOCH + time::Duration::seconds(2)),
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await
|
||||
.expect("delete marker should be written");
|
||||
let mut quorum_reader = PutObjReader::from_vec(b"quorum boundary".to_vec());
|
||||
store.pools[0]
|
||||
.put_object(&bucket, quorum_object, &mut quorum_reader, &ObjectOptions::default())
|
||||
.await
|
||||
.expect("quorum boundary object should be written");
|
||||
{
|
||||
let mut pool_meta = store.pool_meta.write().await;
|
||||
let mut next = PoolMeta::new(&store.pools, &pool_meta);
|
||||
next.pools[1].decommission = Some(PoolDecommissionInfo {
|
||||
start_time: Some(OffsetDateTime::UNIX_EPOCH),
|
||||
..Default::default()
|
||||
});
|
||||
*pool_meta = next;
|
||||
}
|
||||
|
||||
let (_, duplicate_owner) = store
|
||||
.get_latest_object_info_with_idx(&bucket, duplicate_object, &ObjectOptions::default())
|
||||
.await
|
||||
.expect("duplicate owner should resolve");
|
||||
assert_eq!(duplicate_owner, 1, "latest duplicate must win when all pools are eligible");
|
||||
let (_, active_duplicate_owner) = store
|
||||
.get_latest_object_info_with_idx(
|
||||
&bucket,
|
||||
duplicate_object,
|
||||
&ObjectOptions {
|
||||
skip_decommissioned: true,
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await
|
||||
.expect("active duplicate owner should resolve");
|
||||
assert_eq!(
|
||||
active_duplicate_owner, 0,
|
||||
"suspended duplicate must be excluded from active owner selection"
|
||||
);
|
||||
let (duplicate_result, duplicate_err) = store
|
||||
.handle_heal_object(&bucket, duplicate_object, "", &HealOpts::default())
|
||||
.await
|
||||
.expect("duplicate owner heal should complete through the production path");
|
||||
assert_eq!(duplicate_result.object, duplicate_object);
|
||||
assert!(duplicate_err.is_none(), "active duplicate should be repaired: {duplicate_err:?}");
|
||||
assert!(
|
||||
duplicate_missing_disk.read_xl(&bucket, duplicate_object, false).await.is_ok(),
|
||||
"production heal must repair the active duplicate owner rather than the suspended owner"
|
||||
);
|
||||
let (marker_info, marker_owner) = store
|
||||
.get_latest_object_info_with_idx(
|
||||
&bucket,
|
||||
marker_object,
|
||||
&ObjectOptions {
|
||||
skip_decommissioned: true,
|
||||
versioned: true,
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await
|
||||
.expect("latest delete marker should resolve");
|
||||
assert_eq!(marker_owner, 0);
|
||||
assert!(marker_info.delete_marker, "latest version must preserve delete-marker semantics");
|
||||
|
||||
let (active_result, active_err) = store
|
||||
.handle_heal_object(&bucket, active_object, "", &HealOpts::default())
|
||||
.await
|
||||
.expect("unscoped active-owner heal should complete");
|
||||
assert_eq!(active_result.object, active_object);
|
||||
assert!(active_err.is_none(), "active owner must be selected even with a suspended pool");
|
||||
assert!(
|
||||
missing_active_disk.read_xl(&bucket, active_object, false).await.is_ok(),
|
||||
"active owner heal must write the missing disk metadata: result={active_result:?}, err={active_err:?}"
|
||||
);
|
||||
assert!(
|
||||
store.pools[1]
|
||||
.get_object_info(&bucket, active_object, &ObjectOptions::default())
|
||||
.await
|
||||
.is_err(),
|
||||
"the suspended pool must not be written for an active-owner object"
|
||||
);
|
||||
|
||||
let (suspended_result, suspended_err) = store
|
||||
.handle_heal_object(&bucket, suspended_only_object, "", &HealOpts::default())
|
||||
.await
|
||||
.expect("unscoped suspended-only heal should return a terminal result");
|
||||
assert!(suspended_result.object.is_empty());
|
||||
assert!(matches!(suspended_err, Some(Error::FileNotFound)));
|
||||
assert!(
|
||||
store.pools[1]
|
||||
.get_object_info(&bucket, suspended_only_object, &ObjectOptions::default())
|
||||
.await
|
||||
.is_ok(),
|
||||
"suspended-only data must remain untouched when unscoped heal reports absent"
|
||||
);
|
||||
|
||||
let (_, explicit_err) = store
|
||||
.handle_heal_object(
|
||||
&bucket,
|
||||
suspended_only_object,
|
||||
"",
|
||||
&HealOpts {
|
||||
pool: Some(1),
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await
|
||||
.expect("explicit suspended-owner heal should return a mapped error");
|
||||
assert!(matches!(explicit_err, Some(Error::SlowDown)));
|
||||
|
||||
let original_quorum_disks = store.pools[0].disk_set[0].disks.read().await.clone();
|
||||
let surviving_quorum_disk = original_quorum_disks[3].clone();
|
||||
*store.pools[0].disk_set[0].disks.write().await = vec![None, None, None, surviving_quorum_disk];
|
||||
let (_, quorum_err) = store
|
||||
.handle_heal_object(&bucket, quorum_object, "", &HealOpts::default())
|
||||
.await
|
||||
.expect("quorum boundary heal should return a mapped result");
|
||||
*store.pools[0].disk_set[0].disks.write().await = original_quorum_disks;
|
||||
assert!(
|
||||
matches!(quorum_err, Some(Error::ErasureReadQuorum)),
|
||||
"quorum-boundary heal must preserve quorum error, got {quorum_err:?}"
|
||||
);
|
||||
shutdown.cancel();
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn handle_heal_format_continues_after_a_pool_error() {
|
||||
let canonical_format = FormatV3::new(1, 3);
|
||||
|
||||
@@ -14,10 +14,15 @@
|
||||
|
||||
use super::*;
|
||||
use crate::config::storageclass;
|
||||
use crate::core::pools::merge_pool_status_refresh;
|
||||
use crate::layout::pool_space::{ServerPoolsAvailableSpace, build_server_pools_available_space};
|
||||
use crate::runtime::sources as runtime_sources;
|
||||
use crate::storage_api_contracts::{admin::StorageAdminApi, namespace::NamespaceLocking as _, object::ObjectOperations as _};
|
||||
pub(in crate::store) mod support;
|
||||
|
||||
const LOG_COMPONENT_ECSTORE: &str = "ecstore";
|
||||
const LOG_SUBSYSTEM_POOLS: &str = "pools";
|
||||
const EVENT_POOL_META_RELOAD: &str = "pool_meta_reload";
|
||||
use support::{
|
||||
LatestObjectInfoCandidate, PoolErr, PoolObjInfo, RebalanceDeletePoolResult, pool_lookup_not_found_error,
|
||||
rebalance_disk_set_lookup_error, resolve_latest_object_info_candidates, resolve_rebalance_delete_from_all_pools_result,
|
||||
@@ -684,17 +689,49 @@ impl ECStore {
|
||||
)
|
||||
}
|
||||
|
||||
pub async fn reload_pool_meta(&self) -> Result<()> {
|
||||
let mut meta = PoolMeta::default();
|
||||
/// Peer reload entry: refreshes in-memory pool metadata from the shared
|
||||
/// persisted snapshot. Returns whether newer state was actually merged so
|
||||
/// callers only trigger missing-worker recovery after a real state change;
|
||||
/// delayed snapshots are merged monotonically and never blind-assigned.
|
||||
pub async fn reload_pool_meta(&self) -> Result<bool> {
|
||||
let mut reloaded = PoolMeta::default();
|
||||
resolve_store_rebalance_pool_meta_reload_result(
|
||||
meta.load(self.pools[0].clone(), self.pools.clone()).await,
|
||||
reloaded.load(self.pools[0].clone(), self.pools.clone()).await,
|
||||
"reload_pool_meta",
|
||||
)?;
|
||||
|
||||
// Lock order: release the decommission_cancelers guard before taking
|
||||
// the pool_meta write guard; neither is held across the disk read.
|
||||
let active_workers = {
|
||||
let cancelers = self.decommission_cancelers.read().await;
|
||||
cancelers.iter().map(Option::is_some).collect::<Vec<_>>()
|
||||
};
|
||||
|
||||
let incoming_has_pools = !reloaded.pools.is_empty();
|
||||
let mut pool_meta = self.pool_meta.write().await;
|
||||
*pool_meta = meta;
|
||||
// *self.pool_meta.write().expect("operation should succeed") = meta;
|
||||
Ok(())
|
||||
let merged_newer = merge_pool_status_refresh(&mut pool_meta, reloaded, &active_workers);
|
||||
|
||||
if !merged_newer && !incoming_has_pools {
|
||||
warn!(
|
||||
event = EVENT_POOL_META_RELOAD,
|
||||
component = LOG_COMPONENT_ECSTORE,
|
||||
subsystem = LOG_SUBSYSTEM_POOLS,
|
||||
result = "ignored",
|
||||
reason = "missing_metadata",
|
||||
"Peer pool meta reload ignored because persisted metadata is missing"
|
||||
);
|
||||
} else if !merged_newer {
|
||||
debug!(
|
||||
event = EVENT_POOL_META_RELOAD,
|
||||
component = LOG_COMPONENT_ECSTORE,
|
||||
subsystem = LOG_SUBSYSTEM_POOLS,
|
||||
result = "ignored",
|
||||
reason = "stale_snapshot",
|
||||
"Peer pool meta reload ignored as a stale snapshot"
|
||||
);
|
||||
}
|
||||
|
||||
Ok(merged_newer)
|
||||
}
|
||||
|
||||
/// Disk information deduplication function
|
||||
@@ -860,6 +897,7 @@ fn lifecycle_delete_all_test_failure(phase: crate::object_api::LifecycleDeleteAl
|
||||
mod tests {
|
||||
use super::*;
|
||||
use crate::config::storageclass::{CLASS_RRS, CLASS_STANDARD, lookup_config_for_pools_without_env};
|
||||
use crate::core::pools::{POOL_META_VERSION, PoolDecommissionInfo, PoolStatus};
|
||||
use crate::disk::error::DiskError;
|
||||
use crate::layout::endpoint::Endpoint;
|
||||
use crate::layout::endpoints::{EndpointServerPools, Endpoints, PoolEndpoints};
|
||||
@@ -870,6 +908,7 @@ mod tests {
|
||||
use rustfs_config::server_config::KVS;
|
||||
use rustfs_filemeta::FileInfo;
|
||||
use std::sync::Arc;
|
||||
use time::{Duration as TimeDuration, OffsetDateTime};
|
||||
use tokio_util::sync::CancellationToken;
|
||||
|
||||
async fn setup_multi_pool_test_store(
|
||||
@@ -1712,4 +1751,280 @@ mod tests {
|
||||
.contains("failed to resolve rebalance disk set: pool index 2, set index 7, pool count 3")
|
||||
);
|
||||
}
|
||||
|
||||
fn reload_test_pool_status(decommission: Option<PoolDecommissionInfo>, last_update: time::OffsetDateTime) -> PoolStatus {
|
||||
PoolStatus {
|
||||
id: 0,
|
||||
cmd_line: "pool-0".to_string(),
|
||||
last_update,
|
||||
decommission,
|
||||
}
|
||||
}
|
||||
|
||||
fn reload_test_pool_meta(pool: PoolStatus) -> PoolMeta {
|
||||
PoolMeta {
|
||||
version: POOL_META_VERSION,
|
||||
pools: vec![pool],
|
||||
dont_save: false,
|
||||
}
|
||||
}
|
||||
|
||||
async fn persist_reload_snapshot(store: &ECStore, snapshot: &PoolMeta) {
|
||||
snapshot
|
||||
.save(store.pools.clone())
|
||||
.await
|
||||
.expect("pool meta snapshot should persist to every pool");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial_test::serial]
|
||||
async fn peer_pool_meta_reload_does_not_rollback_newer_local_states() {
|
||||
let (_temp_dir, store, shutdown) = setup_multi_pool_test_store("pool-meta-reload-stale", &[2]).await;
|
||||
|
||||
let stale_time = OffsetDateTime::now_utc();
|
||||
let newer_time = stale_time + TimeDuration::seconds(30);
|
||||
let progressed_states: [(&str, PoolDecommissionInfo); 4] = [
|
||||
(
|
||||
"queued",
|
||||
PoolDecommissionInfo {
|
||||
queued: true,
|
||||
start_time: Some(stale_time),
|
||||
..Default::default()
|
||||
},
|
||||
),
|
||||
(
|
||||
"canceled",
|
||||
PoolDecommissionInfo {
|
||||
canceled: true,
|
||||
start_time: Some(stale_time),
|
||||
..Default::default()
|
||||
},
|
||||
),
|
||||
(
|
||||
"failed",
|
||||
PoolDecommissionInfo {
|
||||
failed: true,
|
||||
start_time: Some(stale_time),
|
||||
..Default::default()
|
||||
},
|
||||
),
|
||||
(
|
||||
"complete",
|
||||
PoolDecommissionInfo {
|
||||
complete: true,
|
||||
start_time: Some(stale_time),
|
||||
..Default::default()
|
||||
},
|
||||
),
|
||||
];
|
||||
|
||||
for (state_label, local_state) in progressed_states {
|
||||
{
|
||||
let mut pool_meta = store.pool_meta.write().await;
|
||||
*pool_meta = reload_test_pool_meta(reload_test_pool_status(Some(local_state.clone()), newer_time));
|
||||
}
|
||||
|
||||
// A delayed peer message carries a snapshot that predates the local progression.
|
||||
let stale_snapshot = reload_test_pool_meta(reload_test_pool_status(
|
||||
Some(PoolDecommissionInfo {
|
||||
start_time: Some(stale_time),
|
||||
..Default::default()
|
||||
}),
|
||||
stale_time,
|
||||
));
|
||||
persist_reload_snapshot(&store, &stale_snapshot).await;
|
||||
|
||||
let merged_newer = store.reload_pool_meta().await.expect("stale reload should succeed");
|
||||
assert!(
|
||||
!merged_newer,
|
||||
"a delayed reload must not report merged newer state for the {state_label} progression"
|
||||
);
|
||||
|
||||
let pool_meta = store.pool_meta.read().await;
|
||||
let info = pool_meta.pools[0]
|
||||
.decommission
|
||||
.as_ref()
|
||||
.expect("local decommission state should survive a stale reload");
|
||||
assert_eq!(info.queued, local_state.queued, "{state_label} queued flag must not roll back");
|
||||
assert_eq!(info.canceled, local_state.canceled, "{state_label} canceled flag must not roll back");
|
||||
assert_eq!(info.failed, local_state.failed, "{state_label} failed flag must not roll back");
|
||||
assert_eq!(info.complete, local_state.complete, "{state_label} complete flag must not roll back");
|
||||
assert_eq!(
|
||||
pool_meta.pools[0].last_update, newer_time,
|
||||
"{state_label} progress timestamp must be kept"
|
||||
);
|
||||
}
|
||||
|
||||
shutdown.cancel();
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial_test::serial]
|
||||
async fn peer_pool_meta_reload_merges_newer_state_and_is_idempotent_on_duplicate_delivery() {
|
||||
let (_temp_dir, store, shutdown) = setup_multi_pool_test_store("pool-meta-reload-duplicate", &[2]).await;
|
||||
|
||||
let older_time = OffsetDateTime::now_utc();
|
||||
let newer_time = older_time + TimeDuration::seconds(30);
|
||||
|
||||
{
|
||||
let mut pool_meta = store.pool_meta.write().await;
|
||||
*pool_meta = reload_test_pool_meta(reload_test_pool_status(
|
||||
Some(PoolDecommissionInfo {
|
||||
items_decommissioned: 1,
|
||||
..Default::default()
|
||||
}),
|
||||
older_time,
|
||||
));
|
||||
}
|
||||
let newer_snapshot = reload_test_pool_meta(reload_test_pool_status(
|
||||
Some(PoolDecommissionInfo {
|
||||
complete: true,
|
||||
items_decommissioned: 10,
|
||||
..Default::default()
|
||||
}),
|
||||
newer_time,
|
||||
));
|
||||
persist_reload_snapshot(&store, &newer_snapshot).await;
|
||||
|
||||
let merged_newer = store.reload_pool_meta().await.expect("first reload should succeed");
|
||||
assert!(merged_newer, "a strictly newer persisted snapshot must merge");
|
||||
|
||||
{
|
||||
let pool_meta = store.pool_meta.read().await;
|
||||
let info = pool_meta.pools[0].decommission.as_ref().expect("merged decommission state");
|
||||
assert!(info.complete);
|
||||
assert_eq!(info.items_decommissioned, 10);
|
||||
assert_eq!(pool_meta.pools[0].last_update, newer_time);
|
||||
}
|
||||
|
||||
// Redelivering the same generation must be a no-op.
|
||||
let duplicate_merged = store.reload_pool_meta().await.expect("duplicate reload should succeed");
|
||||
assert!(!duplicate_merged, "a duplicate delivery must not re-apply merged state");
|
||||
|
||||
{
|
||||
let pool_meta = store.pool_meta.read().await;
|
||||
let info = pool_meta.pools[0].decommission.as_ref().expect("merged decommission state");
|
||||
assert!(info.complete);
|
||||
assert_eq!(info.items_decommissioned, 10);
|
||||
assert_eq!(pool_meta.pools[0].last_update, newer_time);
|
||||
}
|
||||
|
||||
shutdown.cancel();
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial_test::serial]
|
||||
async fn peer_pool_meta_reload_keeps_active_worker_progress_over_newer_snapshot() {
|
||||
let (_temp_dir, store, shutdown) = setup_multi_pool_test_store("pool-meta-reload-worker", &[2]).await;
|
||||
*store.decommission_cancelers.write().await = vec![Some(CancellationToken::new())];
|
||||
|
||||
let worker_time = OffsetDateTime::now_utc();
|
||||
let newer_time = worker_time + TimeDuration::seconds(30);
|
||||
|
||||
{
|
||||
let mut pool_meta = store.pool_meta.write().await;
|
||||
*pool_meta = reload_test_pool_meta(reload_test_pool_status(
|
||||
Some(PoolDecommissionInfo {
|
||||
start_time: Some(worker_time),
|
||||
items_decommissioned: 10,
|
||||
bytes_done: 1_024,
|
||||
..Default::default()
|
||||
}),
|
||||
worker_time,
|
||||
));
|
||||
}
|
||||
// Even a strictly newer terminal snapshot must not override a live worker.
|
||||
let newer_terminal_snapshot = reload_test_pool_meta(reload_test_pool_status(
|
||||
Some(PoolDecommissionInfo {
|
||||
complete: true,
|
||||
..Default::default()
|
||||
}),
|
||||
newer_time,
|
||||
));
|
||||
persist_reload_snapshot(&store, &newer_terminal_snapshot).await;
|
||||
|
||||
let merged_newer = store
|
||||
.reload_pool_meta()
|
||||
.await
|
||||
.expect("reload under an active worker should succeed");
|
||||
assert!(!merged_newer, "an active local worker must block snapshot replacement");
|
||||
|
||||
let pool_meta = store.pool_meta.read().await;
|
||||
let info = pool_meta.pools[0]
|
||||
.decommission
|
||||
.as_ref()
|
||||
.expect("worker progress should remain");
|
||||
assert!(!info.complete);
|
||||
assert_eq!(info.items_decommissioned, 10);
|
||||
assert_eq!(info.bytes_done, 1_024);
|
||||
assert_eq!(pool_meta.pools[0].last_update, worker_time);
|
||||
|
||||
shutdown.cancel();
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial_test::serial]
|
||||
async fn peer_pool_meta_reload_fails_closed_when_persisted_metadata_is_missing() {
|
||||
let (temp_dir, store, shutdown) = setup_multi_pool_test_store("pool-meta-reload-missing", &[2]).await;
|
||||
|
||||
let kept_time = OffsetDateTime::now_utc();
|
||||
{
|
||||
let mut pool_meta = store.pool_meta.write().await;
|
||||
*pool_meta = reload_test_pool_meta(reload_test_pool_status(
|
||||
Some(PoolDecommissionInfo {
|
||||
complete: true,
|
||||
..Default::default()
|
||||
}),
|
||||
kept_time,
|
||||
));
|
||||
}
|
||||
// Persist first so the test controls exactly what exists on disk.
|
||||
persist_reload_snapshot(
|
||||
&store,
|
||||
&reload_test_pool_meta(reload_test_pool_status(
|
||||
Some(PoolDecommissionInfo {
|
||||
complete: true,
|
||||
..Default::default()
|
||||
}),
|
||||
kept_time,
|
||||
)),
|
||||
)
|
||||
.await;
|
||||
|
||||
let mut deleted_any = false;
|
||||
for disk_index in 0..2 {
|
||||
let pool_bin_dir = temp_dir
|
||||
.path()
|
||||
.join(format!("pool0-disk{disk_index}"))
|
||||
.join(crate::disk::RUSTFS_META_BUCKET)
|
||||
.join(crate::core::pools::POOL_META_NAME);
|
||||
if pool_bin_dir.exists() {
|
||||
tokio::fs::remove_dir_all(&pool_bin_dir)
|
||||
.await
|
||||
.expect("persisted pool metadata object dir should be removable");
|
||||
deleted_any = true;
|
||||
}
|
||||
}
|
||||
// The meta-bucket layout may nest objects per pool; fall back to removing
|
||||
// every pool.bin object directory below the temp root.
|
||||
if !deleted_any {
|
||||
panic!("no pool.bin found under {:?}", temp_dir.path());
|
||||
}
|
||||
|
||||
let merged_newer = store
|
||||
.reload_pool_meta()
|
||||
.await
|
||||
.expect("reload with missing metadata should fail closed, not error");
|
||||
assert!(!merged_newer, "missing persisted metadata must not count as merged state");
|
||||
|
||||
let pool_meta = store.pool_meta.read().await;
|
||||
let info = pool_meta.pools[0]
|
||||
.decommission
|
||||
.as_ref()
|
||||
.expect("missing persisted metadata must not default local state away");
|
||||
assert!(info.complete);
|
||||
assert_eq!(pool_meta.pools[0].last_update, kept_time);
|
||||
|
||||
shutdown.cancel();
|
||||
}
|
||||
}
|
||||
|
||||
@@ -60,7 +60,9 @@ pub const REPLICATION_READ_ONLY_HISTORICAL_FIELDS: &[&str] = &[
|
||||
"Destination.ReplicationTime",
|
||||
];
|
||||
|
||||
pub const REMOTE_TARGET_CAPABILITY_CONTRACT_VERSION: u32 = 1;
|
||||
// v2: disableProxy moved from unsupported to writable (per-target read-proxy
|
||||
// opt-out is accepted by set-remote-target and the `proxy` update op).
|
||||
pub const REMOTE_TARGET_CAPABILITY_CONTRACT_VERSION: u32 = 2;
|
||||
|
||||
pub const REMOTE_TARGET_WRITABLE_FIELDS: &[&str] = &[
|
||||
"sourcebucket",
|
||||
@@ -83,9 +85,12 @@ pub const REMOTE_TARGET_WRITABLE_FIELDS: &[&str] = &[
|
||||
// madmin default of 60s); the per-target health-check interval is not
|
||||
// yet applied — the heartbeat keeps its global env-configured interval.
|
||||
"healthCheckDuration",
|
||||
// Per-target read-proxy opt-out, consumed by the proxy-target selector
|
||||
// (contract v2; previously only importable via MinIO bucket-targets.json).
|
||||
"disableProxy",
|
||||
];
|
||||
|
||||
pub const REMOTE_TARGET_UNSUPPORTED_FIELDS: &[&str] = &["disableProxy", "edge", "edgeSyncBeforeExpiry"];
|
||||
pub const REMOTE_TARGET_UNSUPPORTED_FIELDS: &[&str] = &["edge", "edgeSyncBeforeExpiry"];
|
||||
|
||||
#[derive(Debug, Clone, Serialize, Deserialize, Default)]
|
||||
pub struct ObjectOpts {
|
||||
|
||||
@@ -73,6 +73,8 @@ enum TargetUpdateOp {
|
||||
/// Connection group: credentials plus endpoint, target bucket, and TLS settings.
|
||||
Credentials,
|
||||
Sync,
|
||||
/// Per-target read-proxy opt-out (`disableProxy`).
|
||||
Proxy,
|
||||
Bandwidth,
|
||||
Path,
|
||||
}
|
||||
@@ -81,12 +83,13 @@ fn parse_remote_target_update_ops(queries: &HashMap<String, String>) -> S3Result
|
||||
const SUPPORTED_OPS: &[(&str, TargetUpdateOp)] = &[
|
||||
("creds", TargetUpdateOp::Credentials),
|
||||
("sync", TargetUpdateOp::Sync),
|
||||
("proxy", TargetUpdateOp::Proxy),
|
||||
("bandwidth", TargetUpdateOp::Bandwidth),
|
||||
("path", TargetUpdateOp::Path),
|
||||
];
|
||||
// Present in the MinIO wire contract, but they drive target fields this
|
||||
// version rejects as unsupported — fail loudly instead of silently ignoring.
|
||||
const UNSUPPORTED_OPS: &[&str] = &["proxy", "healthcheck", "edge", "edgeSyncBeforeExpiry"];
|
||||
const UNSUPPORTED_OPS: &[&str] = &["healthcheck", "edge", "edgeSyncBeforeExpiry"];
|
||||
|
||||
for key in UNSUPPORTED_OPS {
|
||||
if queries.get(*key).is_some_and(|value| value == "true") {
|
||||
@@ -312,11 +315,10 @@ impl RemoteTargetRequest {
|
||||
));
|
||||
}
|
||||
|
||||
for (unsupported, configured) in
|
||||
REMOTE_TARGET_UNSUPPORTED_FIELDS
|
||||
.iter()
|
||||
.copied()
|
||||
.zip([self.disable_proxy, self.edge, self.edge_sync_before_expiry])
|
||||
for (unsupported, configured) in REMOTE_TARGET_UNSUPPORTED_FIELDS
|
||||
.iter()
|
||||
.copied()
|
||||
.zip([self.edge, self.edge_sync_before_expiry])
|
||||
{
|
||||
if configured {
|
||||
return Err(s3_error!(
|
||||
@@ -702,6 +704,7 @@ impl Operation for SetRemoteTargetHandler {
|
||||
target.deployment_id = remote_target.deployment_id.clone();
|
||||
}
|
||||
TargetUpdateOp::Sync => target.replication_sync = remote_target.replication_sync,
|
||||
TargetUpdateOp::Proxy => target.disable_proxy = remote_target.disable_proxy,
|
||||
TargetUpdateOp::Bandwidth => target.bandwidth_limit = remote_target.bandwidth_limit,
|
||||
TargetUpdateOp::Path => target.path = remote_target.path.clone(),
|
||||
}
|
||||
@@ -1520,6 +1523,7 @@ mod tests {
|
||||
("update", "true"),
|
||||
("creds", "true"),
|
||||
("sync", "true"),
|
||||
("proxy", "true"),
|
||||
("bandwidth", "true"),
|
||||
("path", "true"),
|
||||
]))
|
||||
@@ -1529,6 +1533,7 @@ mod tests {
|
||||
vec![
|
||||
TargetUpdateOp::Credentials,
|
||||
TargetUpdateOp::Sync,
|
||||
TargetUpdateOp::Proxy,
|
||||
TargetUpdateOp::Bandwidth,
|
||||
TargetUpdateOp::Path
|
||||
]
|
||||
@@ -2070,7 +2075,6 @@ mod tests {
|
||||
("credentials.session_token", serde_json::json!("session-token")),
|
||||
("credentials.expiration", serde_json::json!("2026-01-01T00:00:00Z")),
|
||||
("api", serde_json::json!("s3v2")),
|
||||
("disableProxy", serde_json::json!(true)),
|
||||
("edge", serde_json::json!(true)),
|
||||
("edgeSyncBeforeExpiry", serde_json::json!(true)),
|
||||
] {
|
||||
@@ -2300,6 +2304,44 @@ mod tests {
|
||||
assert!(!REMOTE_TARGET_UNSUPPORTED_FIELDS.contains(&"healthCheckDuration"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn remote_target_disable_proxy_is_declared_writable_edge_stays_unsupported() {
|
||||
assert!(REMOTE_TARGET_WRITABLE_FIELDS.contains(&"disableProxy"));
|
||||
assert!(!REMOTE_TARGET_UNSUPPORTED_FIELDS.contains(&"disableProxy"));
|
||||
// edge sync has no implementation behind it — it must stay rejected.
|
||||
assert!(REMOTE_TARGET_UNSUPPORTED_FIELDS.contains(&"edge"));
|
||||
assert!(REMOTE_TARGET_UNSUPPORTED_FIELDS.contains(&"edgeSyncBeforeExpiry"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn remote_target_create_accepts_disable_proxy() {
|
||||
let mut request = valid_remote_target_request();
|
||||
request["disableProxy"] = serde_json::json!(true);
|
||||
|
||||
let target = serde_json::from_value::<RemoteTargetRequest>(request)
|
||||
.expect("request should deserialize")
|
||||
.into_bucket_target()
|
||||
.expect("disableProxy is a supported per-target read-proxy opt-out");
|
||||
|
||||
assert!(target.disable_proxy);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn update_body_with_proxy_op_toggles_disable_proxy_without_credentials() {
|
||||
// Mirrors the other partial-update groups: a proxy-only update body may
|
||||
// omit the connection fields entirely.
|
||||
let body = serde_json::json!({
|
||||
"arn": "arn:rustfs:replication:us-east-1:dep:target",
|
||||
"type": "replication",
|
||||
"disableProxy": true
|
||||
});
|
||||
let request: RemoteTargetRequest = serde_json::from_value(body).expect("partial update body should deserialize");
|
||||
let target = request
|
||||
.into_update_bucket_target(&[TargetUpdateOp::Proxy])
|
||||
.expect("proxy-only update must not require credentials");
|
||||
assert!(target.disable_proxy);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn remote_target_capability_fields_do_not_overlap() {
|
||||
for field in REMOTE_TARGET_UNSUPPORTED_FIELDS {
|
||||
|
||||
@@ -1262,7 +1262,9 @@ mod tests {
|
||||
assert_eq!(response.summary.manual_transition_jobs.state, CapabilityState::Supported);
|
||||
assert_eq!(response.replication.contract_version, 1);
|
||||
assert_eq!(response.replication.bucket_replication.contract_version, 1);
|
||||
assert_eq!(response.replication.remote_targets.contract_version, 1);
|
||||
// v2: disableProxy moved from unsupported to writable (per-target
|
||||
// read-proxy opt-out reached the admin API).
|
||||
assert_eq!(response.replication.remote_targets.contract_version, 2);
|
||||
assert_eq!(response.replication.bucket_replication.status.state, CapabilityState::Supported);
|
||||
assert_eq!(response.replication.remote_targets.status.state, CapabilityState::Supported);
|
||||
assert_eq!(
|
||||
@@ -1293,7 +1295,15 @@ mod tests {
|
||||
.remote_targets
|
||||
.fields
|
||||
.iter()
|
||||
.any(|field| field.name == "disableProxy" && field.state == super::ReplicationFieldState::Unsupported)
|
||||
.any(|field| field.name == "disableProxy" && field.state == super::ReplicationFieldState::Supported)
|
||||
);
|
||||
assert!(
|
||||
response
|
||||
.replication
|
||||
.remote_targets
|
||||
.fields
|
||||
.iter()
|
||||
.any(|field| field.name == "edge" && field.state == super::ReplicationFieldState::Unsupported)
|
||||
);
|
||||
assert!(
|
||||
response
|
||||
@@ -1364,7 +1374,7 @@ mod tests {
|
||||
assert_eq!(value["summary"]["manual_transition_jobs"]["state"], "supported");
|
||||
assert_eq!(value["replication"]["contract_version"], 1);
|
||||
assert_eq!(value["replication"]["bucket_replication"]["contract_version"], 1);
|
||||
assert_eq!(value["replication"]["remote_targets"]["contract_version"], 1);
|
||||
assert_eq!(value["replication"]["remote_targets"]["contract_version"], 2);
|
||||
assert_eq!(value["replication"]["bucket_replication"]["status"]["state"], "supported");
|
||||
assert_eq!(value["replication"]["remote_targets"]["status"]["state"], "supported");
|
||||
assert_eq!(
|
||||
@@ -1383,7 +1393,14 @@ mod tests {
|
||||
.as_array()
|
||||
.expect("remote target fields should be an array")
|
||||
.iter()
|
||||
.any(|field| field["name"] == "disableProxy" && field["state"] == "unsupported")
|
||||
.any(|field| field["name"] == "disableProxy" && field["state"] == "supported")
|
||||
);
|
||||
assert!(
|
||||
value["replication"]["remote_targets"]["fields"]
|
||||
.as_array()
|
||||
.expect("remote target fields should be an array")
|
||||
.iter()
|
||||
.any(|field| field["name"] == "edge" && field["state"] == "unsupported")
|
||||
);
|
||||
assert!(
|
||||
value["replication"]["remote_targets"]["fields"]
|
||||
|
||||
@@ -1987,8 +1987,10 @@ impl Node for NodeService {
|
||||
error_info: Some("errServerNotInitialized".to_string()),
|
||||
}));
|
||||
};
|
||||
// Recover missing workers only after the reload merged newer state; a
|
||||
// stale or duplicate reload must not spawn workers for an older generation.
|
||||
match store.reload_pool_meta().await {
|
||||
Ok(_) => match store.spawn_missing_local_decommission_routines().await {
|
||||
Ok(true) => match store.spawn_missing_local_decommission_routines().await {
|
||||
Ok(_) => Ok(Response::new(ReloadPoolMetaResponse {
|
||||
success: true,
|
||||
error_info: None,
|
||||
@@ -1998,6 +2000,10 @@ impl Node for NodeService {
|
||||
error_info: Some(err.to_string()),
|
||||
})),
|
||||
},
|
||||
Ok(false) => Ok(Response::new(ReloadPoolMetaResponse {
|
||||
success: true,
|
||||
error_info: None,
|
||||
})),
|
||||
Err(err) => Ok(Response::new(ReloadPoolMetaResponse {
|
||||
success: false,
|
||||
error_info: Some(err.to_string()),
|
||||
|
||||
Reference in New Issue
Block a user