mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-19 11:06:17 +00:00
f1f86ee9d0
* chore(ecstore): drop the set_disk dead_code blanket Removing the blanket exposes 39 items; exactly one is deleted. The low share is a finding, not caution: unlike the disk root, where platform gating made local adjudication impossible, here the items were checked and nearly all of them are live. Deleted: HealEntryResult, the only item with no reference anywhere. What the checks turned up, in the order the warnings suggest deleting them: SetDisks::rename_data looked like the head of a dead chain feeding into_legacy_tuple and RenameDataLegacyTuple. It is not: production goes through rename_data_owned, and rename_data itself has test callers at mod.rs:5809 and 5880. The chain below it is therefore live through the tests, and inferring "this is dead, so its callee is dead" would have removed three working items. create_bitrot_readers_until_quorum, read_multiple_files and map_cleanup_join_result all have callers inside their files' test modules, so they only look dead in the lib target. TransitionCommitBarrier and TransitionUploadedSaveProbe, with their install/wait_until_paused/release surfaces, are installed by tests behind #[cfg(all(test, feature = "test-util"))]. ctx.rs's SetDisksCtx accessors are the split seam left by the SetDisks god-object break-up (backlog#815). heal_object_dir's two apparent references are comments, and they document an index-alignment contract that live code maintains for it, so they stay as they are. Worth a maintainer decision: the metadata early-stop switch has a complete percentage-rollout facet — ENV_RUSTFS_GET_METADATA_EARLY_STOP_ROLLOUT_PCT, get_metadata_early_stop_rollout_pct and should_use_metadata_early_stop — with no caller, no test and no documentation, while its sibling enable flag is live. It is kept with an allow that says so rather than removed, since a rollout knob is a product call. One placement note for anyone adding allows near heal code: check_logging_guardrails.sh requires #[instrument(level = "trace")] to sit immediately before async fn heal_object_dir, so the allow goes above the instrument attribute. Putting it between the two drops the guard's match count and fails the check. Verification, four lanes warning-free: default, --tests, --features rio-v2 --tests, --features test-util --tests. cargo nextest run -p rustfs-ecstore 4096 passed; clippy --lib --tests -D warnings clean; make pre-commit exit 0. Ref rustfs/backlog#1823 (step 2). * chore(ecstore): fix duplicated and inaccurate dead_code reasons in set_disk format_lock_error carried the same #[allow] twice. Five items in the locking/heal roots were labelled 'asserted by this file's tests' while having no reference at all - heal_object_dir's only two references are comments, as this branch's own notes point out. Say what each item actually is instead, so the next reader does not assume test coverage that is not there. Ref rustfs/backlog#1823. * chore(ecstore): correct the bounded_spare_disk_index dead_code reason The mod.rs copy is an unused test fixture, not something this module's tests assert; the namesake that is exercised lives in the io_primitives test module. Ref rustfs/backlog#1823.
937 lines
34 KiB
Rust
937 lines
34 KiB
Rust
// Copyright 2024 RustFS Team
|
|
//
|
|
// Licensed under the Apache License, Version 2.0 (the "License");
|
|
// you may not use this file except in compliance with the License.
|
|
// You may obtain a copy of the License at
|
|
//
|
|
// http://www.apache.org/licenses/LICENSE-2.0
|
|
//
|
|
// Unless required by applicable law or agreed to in writing, software
|
|
// distributed under the License is distributed on an "AS IS" BASIS,
|
|
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
|
// See the License for the specific language governing permissions and
|
|
// limitations under the License.
|
|
|
|
//! `NamespaceLocking` contract impl for `SetDisks` plus the set-disk locking
|
|
//! helpers (P7 of the God-Object split, tracking backlog#815, issue
|
|
//! backlog#822). The `NamespaceLocking` impl is relocated from set_disk/mod.rs
|
|
//! and the lock formatting/mapping helpers from set_disk/lock.rs are collected
|
|
//! here; the contract stays implemented `for SetDisks`, so its associated-type
|
|
//! bounds are unchanged and helper access is via inherent calls.
|
|
|
|
use super::super::*;
|
|
use crate::disk::health_state::DriveMembershipSnapshot;
|
|
use crate::runtime::sources as runtime_sources;
|
|
|
|
#[async_trait::async_trait]
|
|
impl crate::storage_api_contracts::namespace::NamespaceLocking for SetDisks {
|
|
type Error = Error;
|
|
type NamespaceLock = NamespaceLockWrapper;
|
|
|
|
#[tracing::instrument(level = "trace", skip(self))]
|
|
async fn new_ns_lock(&self, bucket: &str, object: &str) -> Result<NamespaceLockWrapper> {
|
|
// Resolved from this set's own instance context (backlog#1052), not the
|
|
// ambient facade: the facade tracks whichever context is currently
|
|
// published process-wide, so a second instance (or, in tests, another
|
|
// test's transient DistErasure window) would push this set's namespace
|
|
// locking onto its own — possibly empty — dist locker list.
|
|
let set_lock = if self.ctx.is_dist_erasure().await {
|
|
let lockers = if self.lockers.len() == self.shared_lockers.len()
|
|
&& self
|
|
.lockers
|
|
.iter()
|
|
.zip(self.shared_lockers.iter())
|
|
.all(|(current, shared)| Arc::ptr_eq(current, shared))
|
|
{
|
|
self.shared_lockers.clone()
|
|
} else {
|
|
Arc::from(self.lockers.clone())
|
|
};
|
|
// Calculate quorum from the exact client domain used by this lock.
|
|
let lockers_count = lockers.len();
|
|
let write_quorum = if lockers_count > 1 { (lockers_count / 2) + 1 } else { 1 };
|
|
NamespaceLock::with_clients_and_quorum_shared(self.set_lock_namespace.clone(), lockers, write_quorum)
|
|
} else {
|
|
NamespaceLock::with_local_manager_shared(self.set_lock_namespace.clone(), self.local_lock_manager.clone())
|
|
};
|
|
|
|
let resource = ObjectKey {
|
|
bucket: Arc::from(bucket),
|
|
object: Arc::from(object),
|
|
version: None,
|
|
};
|
|
|
|
Ok(NamespaceLockWrapper::new(set_lock, resource, self.locker_owner.clone()))
|
|
}
|
|
}
|
|
|
|
impl SetDisks {
|
|
#[allow(dead_code, reason = "lock diagnostics formatter with no caller in this port (backlog#1823)")]
|
|
pub(in crate::set_disk) fn format_lock_error(&self, bucket: &str, object: &str, mode: &str, err: &LockResult) -> String {
|
|
match err {
|
|
LockResult::Timeout => {
|
|
format!("{mode} lock acquisition timed out on {bucket}/{object} (owner={})", self.locker_owner)
|
|
}
|
|
LockResult::Conflict {
|
|
current_owner,
|
|
current_mode,
|
|
} => format!("{mode} lock conflicted on {bucket}/{object}: held by {current_owner} as {current_mode:?}"),
|
|
LockResult::Acquired => format!("unexpected lock state while acquiring {mode} lock on {bucket}/{object}"),
|
|
}
|
|
}
|
|
|
|
#[allow(dead_code, reason = "lock diagnostics formatter with no caller in this port (backlog#1823)")]
|
|
pub(in crate::set_disk) fn format_lock_error_from_error(
|
|
&self,
|
|
bucket: &str,
|
|
object: &str,
|
|
mode: &str,
|
|
err: &rustfs_lock::error::LockError,
|
|
) -> String {
|
|
match err {
|
|
rustfs_lock::error::LockError::Timeout { .. } => {
|
|
format!(
|
|
"ns_loc: {mode} lock acquisition timed out on {bucket}/{object} (owner={})",
|
|
self.locker_owner
|
|
)
|
|
}
|
|
rustfs_lock::error::LockError::AlreadyLocked { owner, .. } => {
|
|
format!("ns_loc: {mode} lock conflicted on {bucket}/{object}: held by {owner}")
|
|
}
|
|
_ => format!("ns_loc: {mode} lock acquisition failed on {bucket}/{object}: {}", err),
|
|
}
|
|
}
|
|
|
|
pub(in crate::set_disk) fn map_namespace_lock_error(
|
|
&self,
|
|
bucket: &str,
|
|
object: &str,
|
|
mode: &'static str,
|
|
err: rustfs_lock::error::LockError,
|
|
) -> StorageError {
|
|
match err {
|
|
rustfs_lock::error::LockError::QuorumNotReached { required, achieved } => {
|
|
StorageError::NamespaceLockQuorumUnavailable {
|
|
mode,
|
|
bucket: bucket.to_string(),
|
|
object: object.to_string(),
|
|
required,
|
|
achieved,
|
|
}
|
|
}
|
|
other => StorageError::Lock(other),
|
|
}
|
|
}
|
|
|
|
pub(in crate::set_disk) async fn get_disks_internal(&self) -> Vec<Option<DiskStore>> {
|
|
let rl = self.disks.read().await;
|
|
|
|
rl.clone()
|
|
}
|
|
|
|
pub async fn get_local_disks(&self) -> Vec<Option<DiskStore>> {
|
|
let rl = self.disks.read().await;
|
|
|
|
let mut disks: Vec<Option<DiskStore>> = rl
|
|
.clone()
|
|
.into_iter()
|
|
.filter(|v| v.as_ref().is_some_and(|d| d.is_local()))
|
|
.collect();
|
|
|
|
let mut rng = rand::rng();
|
|
|
|
disks.shuffle(&mut rng);
|
|
|
|
disks
|
|
}
|
|
|
|
#[allow(dead_code, reason = "asserted by this file's tests (backlog#1823)")]
|
|
pub(in crate::set_disk) async fn get_online_disks(&self) -> Vec<Option<DiskStore>> {
|
|
let snapshot = self.drive_membership_snapshot().await;
|
|
let mut disks = snapshot.strict_online_candidates().into_iter().map(Some).collect::<Vec<_>>();
|
|
|
|
let mut rng = rand::rng();
|
|
disks.shuffle(&mut rng);
|
|
|
|
disks
|
|
}
|
|
|
|
#[allow(
|
|
dead_code,
|
|
reason = "local-only sibling of the test-covered get_online_disks; no caller in this port (backlog#1823)"
|
|
)]
|
|
pub(in crate::set_disk) async fn get_online_local_disks(&self) -> Vec<Option<DiskStore>> {
|
|
let snapshot = self.drive_membership_snapshot().await;
|
|
let mut disks = snapshot
|
|
.strict_online_local_candidates()
|
|
.into_iter()
|
|
.map(Some)
|
|
.collect::<Vec<_>>();
|
|
|
|
let mut rng = rand::rng();
|
|
|
|
disks.shuffle(&mut rng);
|
|
|
|
disks
|
|
}
|
|
|
|
pub async fn get_online_disks_with_healing(&self, incl_healing: bool) -> (Vec<DiskStore>, bool) {
|
|
let (disks, _, healing) = self.get_online_disks_with_healing_and_info(incl_healing).await;
|
|
(disks, healing > 0)
|
|
}
|
|
|
|
pub async fn drive_membership_snapshot(&self) -> DriveMembershipSnapshot {
|
|
let disks = self.get_disks_internal().await;
|
|
DriveMembershipSnapshot::from_optional_disks(&disks)
|
|
}
|
|
|
|
fn reprobe_runtime_candidates_once(&self, disks: &[DiskStore]) {
|
|
for disk in disks {
|
|
if disk.runtime_state() != disk::health_state::RuntimeDriveHealthState::Online {
|
|
disk.reset_health_for_store_init_retry();
|
|
}
|
|
}
|
|
}
|
|
|
|
pub async fn get_online_disks_with_healing_and_info(&self, incl_healing: bool) -> (Vec<DiskStore>, Vec<DiskInfo>, usize) {
|
|
let snapshot = self.drive_membership_snapshot().await;
|
|
let mut membership_candidates = snapshot.scanner_heal_candidates();
|
|
let mut reprobed = false;
|
|
|
|
loop {
|
|
let mut disks = membership_candidates.clone();
|
|
let mut infos: Vec<Option<DiskInfo>> = vec![None; disks.len()];
|
|
|
|
let mut futures = Vec::with_capacity(disks.len());
|
|
{
|
|
let mut rng = rand::rng();
|
|
disks.shuffle(&mut rng);
|
|
}
|
|
|
|
for (i, disk) in disks.iter().cloned().enumerate() {
|
|
futures.push(async move {
|
|
let info = match disk.disk_info(&DiskInfoOptions::default()).await {
|
|
Ok(info) => info,
|
|
Err(err) => DiskInfo {
|
|
error: err.to_string(),
|
|
..Default::default()
|
|
},
|
|
};
|
|
|
|
Ok((i, info))
|
|
});
|
|
}
|
|
|
|
let processor = runtime_sources::batch_processors().metadata_processor();
|
|
let results = processor.execute_batch(futures).await;
|
|
|
|
for (submitted_idx, result) in results.into_iter().enumerate() {
|
|
match result {
|
|
Ok((disk_idx, info)) => {
|
|
infos[disk_idx] = Some(info);
|
|
}
|
|
Err(err) => {
|
|
infos[submitted_idx] = Some(DiskInfo {
|
|
error: err.to_string(),
|
|
..Default::default()
|
|
});
|
|
}
|
|
}
|
|
}
|
|
|
|
let mut healing: usize = 0;
|
|
|
|
let mut scanning_disks = Vec::new();
|
|
let mut healing_disks = Vec::new();
|
|
let mut scanning_infos = Vec::new();
|
|
let mut healing_infos = Vec::new();
|
|
|
|
let mut new_disks = Vec::new();
|
|
let mut new_infos = Vec::new();
|
|
|
|
for (disk, info) in disks.into_iter().zip(infos) {
|
|
let Some(info) = info else {
|
|
continue;
|
|
};
|
|
|
|
if !info.error.is_empty() {
|
|
continue;
|
|
}
|
|
|
|
if info.healing {
|
|
healing += 1;
|
|
if incl_healing {
|
|
healing_disks.push(disk);
|
|
healing_infos.push(info);
|
|
}
|
|
|
|
continue;
|
|
}
|
|
|
|
if !info.scanning {
|
|
new_disks.push(disk);
|
|
new_infos.push(info);
|
|
} else {
|
|
scanning_disks.push(disk);
|
|
scanning_infos.push(info);
|
|
}
|
|
}
|
|
|
|
new_disks.extend(scanning_disks);
|
|
new_infos.extend(scanning_infos);
|
|
new_disks.extend(healing_disks);
|
|
new_infos.extend(healing_infos);
|
|
|
|
if !new_disks.is_empty() || membership_candidates.is_empty() || reprobed {
|
|
return (new_disks, new_infos, healing);
|
|
}
|
|
|
|
reprobed = true;
|
|
self.reprobe_runtime_candidates_once(&membership_candidates);
|
|
membership_candidates = self.drive_membership_snapshot().await.scanner_heal_candidates();
|
|
}
|
|
}
|
|
|
|
pub(in crate::set_disk) async fn _get_local_disks(&self) -> Vec<Option<DiskStore>> {
|
|
let mut disks = self.get_disks_internal().await;
|
|
|
|
let mut rng = rand::rng();
|
|
|
|
disks.shuffle(&mut rng);
|
|
|
|
disks
|
|
.into_iter()
|
|
.filter(|v| v.as_ref().is_some_and(|d| d.is_local()))
|
|
.collect()
|
|
}
|
|
|
|
pub async fn connect_disks(&self) {
|
|
let rl = self.disks.read().await;
|
|
|
|
let disks = rl.clone();
|
|
|
|
// Explicitly release the lock
|
|
drop(rl);
|
|
|
|
for (i, opdisk) in disks.iter().enumerate() {
|
|
if let Some(disk) = opdisk {
|
|
if disk.is_online().await && disk.get_disk_location().set_idx.is_some() {
|
|
info!("Disk {:?} is online", disk.to_string());
|
|
continue;
|
|
}
|
|
|
|
let _ = disk.close().await;
|
|
}
|
|
|
|
if let Some(endpoint) = self.set_endpoints.get(i) {
|
|
info!("will renew disk, opdisk: {:?}", opdisk);
|
|
self.renew_disk(endpoint).await;
|
|
}
|
|
}
|
|
}
|
|
|
|
pub async fn renew_disk(&self, ep: &Endpoint) {
|
|
debug!("renew_disk: start {:?}", ep);
|
|
|
|
let previous_health = {
|
|
let disks = self.disks.read().await;
|
|
disks
|
|
.iter()
|
|
.filter_map(|disk| disk.as_ref())
|
|
.find(|disk| disk.endpoint() == *ep)
|
|
.and_then(|disk| disk.local_health_tracker_epoch_for_reconnect())
|
|
};
|
|
|
|
let (new_disk, fm) = match Self::connect_endpoint(ep, previous_health).await {
|
|
Ok(res) => res,
|
|
Err(e) => {
|
|
warn!("renew_disk: connect_endpoint err {:?}", &e);
|
|
if ep.is_local && e == DiskError::UnformattedDisk {
|
|
info!("renew_disk unformatteddisk will trigger heal_disk, {:?}", ep);
|
|
let set_disk_id = format!("pool_{}_set_{}", ep.pool_idx, ep.set_idx);
|
|
let _ = send_heal_disk(set_disk_id, Some(HealChannelPriority::Normal)).await;
|
|
}
|
|
return;
|
|
}
|
|
};
|
|
|
|
let (set_idx, disk_idx) = match self.find_disk_index(&fm) {
|
|
Ok(res) => res,
|
|
Err(e) => {
|
|
warn!("renew_disk: find_disk_index err {:?}", e);
|
|
return;
|
|
}
|
|
};
|
|
|
|
// Claiming a misplaced drive into `self.disks` would let two slots or
|
|
// sets manage the same drive and degrade together (backlog#799 B19).
|
|
if set_idx != self.set_index || self.set_endpoints.get(disk_idx) != Some(ep) {
|
|
warn!(
|
|
endpoint = %ep,
|
|
format_set_index = set_idx,
|
|
format_disk_index = disk_idx,
|
|
endpoint_pool_index = ep.pool_idx,
|
|
endpoint_set_index = ep.set_idx,
|
|
endpoint_disk_index = ep.disk_idx,
|
|
expected_pool_index = self.pool_index,
|
|
expected_set_index = self.set_index,
|
|
"renew_disk rejected a drive whose endpoint and format do not identify the same topology slot"
|
|
);
|
|
return;
|
|
}
|
|
|
|
let _ = new_disk.set_disk_id(Some(fm.erasure.this)).await;
|
|
new_disk.enable_health_check();
|
|
|
|
if new_disk.is_local() {
|
|
let local_disk_map = runtime_sources::local_disk_map_handle();
|
|
let mut global_local_disk_map = local_disk_map.write().await;
|
|
let path = new_disk.endpoint().to_string();
|
|
global_local_disk_map.insert(path, Some(new_disk.clone()));
|
|
|
|
if runtime_sources::setup_is_dist_erasure().await {
|
|
let local_disk_set_drives = runtime_sources::local_disk_set_drives_handle();
|
|
let mut local_set_drives = local_disk_set_drives.write().await;
|
|
local_set_drives[self.pool_index][set_idx][disk_idx] = Some(new_disk.clone());
|
|
}
|
|
}
|
|
|
|
debug!("renew_disk: update {:?}", fm.erasure.this);
|
|
|
|
let mut disk_lock = self.disks.write().await;
|
|
disk_lock[disk_idx] = Some(new_disk);
|
|
}
|
|
|
|
pub(in crate::set_disk) fn find_disk_index(&self, fm: &FormatV3) -> Result<(usize, usize)> {
|
|
self.format.check_other(fm)?;
|
|
|
|
if fm.erasure.this.is_nil() {
|
|
return Err(Error::other("DriveID: offline"));
|
|
}
|
|
|
|
for i in 0..self.format.erasure.sets.len() {
|
|
for j in 0..self.format.erasure.sets[0].len() {
|
|
if fm.erasure.this == self.format.erasure.sets[i][j] {
|
|
return Ok((i, j));
|
|
}
|
|
}
|
|
}
|
|
|
|
Err(Error::other("DriveID: not found"))
|
|
}
|
|
|
|
pub(in crate::set_disk) async fn connect_endpoint(
|
|
ep: &Endpoint,
|
|
reconnect: Option<disk::disk_store::ReconnectDiskHealthState>,
|
|
) -> disk::error::Result<(DiskStore, FormatV3)> {
|
|
let disk = crate::disk::new_disk_with_health_tracker(
|
|
ep,
|
|
&DiskOption {
|
|
cleanup: false,
|
|
health_check: true,
|
|
},
|
|
reconnect,
|
|
)
|
|
.await?;
|
|
|
|
let fm = load_format_erasure(&disk, false).await?;
|
|
|
|
Ok((disk, fm))
|
|
}
|
|
|
|
#[allow(
|
|
dead_code,
|
|
reason = "MinIO-parity healing-disk accessor with no caller in this port (backlog#1823)"
|
|
)]
|
|
pub(in crate::set_disk) async fn get_online_disk_with_healing(
|
|
&self,
|
|
incl_healing: bool,
|
|
) -> Result<(Vec<Option<DiskStore>>, bool)> {
|
|
let (new_disks, _, healing) = self.get_online_disk_with_healing_and_info(incl_healing).await?;
|
|
Ok((new_disks, healing > 0))
|
|
}
|
|
|
|
#[allow(
|
|
dead_code,
|
|
reason = "reached only from get_online_disk_with_healing, itself uncalled in this port (backlog#1823)"
|
|
)]
|
|
pub(in crate::set_disk) async fn get_online_disk_with_healing_and_info(
|
|
&self,
|
|
incl_healing: bool,
|
|
) -> Result<(Vec<Option<DiskStore>>, Vec<DiskInfo>, usize)> {
|
|
let mut infos = vec![DiskInfo::default(); self.disks.read().await.len()];
|
|
for (idx, disk) in self.disks.write().await.iter().enumerate() {
|
|
if let Some(disk) = disk {
|
|
match disk.disk_info(&DiskInfoOptions::default()).await {
|
|
Ok(disk_info) => infos[idx] = disk_info,
|
|
Err(err) => infos[idx].error = err.to_string(),
|
|
}
|
|
} else {
|
|
infos[idx].error = "disk not found".to_string();
|
|
}
|
|
}
|
|
|
|
let mut new_disks = Vec::new();
|
|
let mut healing_disks = Vec::new();
|
|
let mut scanning_disks = Vec::new();
|
|
let mut new_infos = Vec::new();
|
|
let mut healing_infos = Vec::new();
|
|
let mut scanning_infos = Vec::new();
|
|
let mut healing = 0;
|
|
|
|
infos.iter().zip(self.disks.write().await.iter()).for_each(|(info, disk)| {
|
|
if info.error.is_empty() {
|
|
if info.healing {
|
|
healing += 1;
|
|
if incl_healing {
|
|
healing_disks.push(disk.clone());
|
|
healing_infos.push(info.clone());
|
|
}
|
|
} else if !info.scanning {
|
|
new_disks.push(disk.clone());
|
|
new_infos.push(info.clone());
|
|
} else {
|
|
scanning_disks.push(disk.clone());
|
|
scanning_infos.push(info.clone());
|
|
}
|
|
}
|
|
});
|
|
|
|
// Prefer non-scanning disks over disks which are currently being scanned.
|
|
new_disks.extend(scanning_disks);
|
|
new_infos.extend(scanning_infos);
|
|
|
|
// Then add healing disks.
|
|
new_disks.extend(healing_disks);
|
|
new_infos.extend(healing_infos);
|
|
|
|
Ok((new_disks, new_infos, healing))
|
|
}
|
|
}
|
|
|
|
#[cfg(test)]
|
|
mod tests {
|
|
use super::*;
|
|
use crate::store::init_format::save_format_file;
|
|
use tempfile::TempDir;
|
|
use tokio::sync::RwLock;
|
|
|
|
async fn make_formatted_local_disk(disk_idx: usize, format: &FormatV3) -> (TempDir, Endpoint, DiskStore) {
|
|
let dir = tempfile::tempdir().expect("tempdir should be created");
|
|
let mut endpoint =
|
|
Endpoint::try_from(dir.path().to_str().expect("tempdir path should be utf8")).expect("endpoint should parse");
|
|
endpoint.set_pool_index(0);
|
|
endpoint.set_set_index(0);
|
|
endpoint.set_disk_index(disk_idx);
|
|
|
|
let disk = new_disk(
|
|
&endpoint,
|
|
&DiskOption {
|
|
cleanup: false,
|
|
health_check: false,
|
|
},
|
|
)
|
|
.await
|
|
.expect("local disk should be created");
|
|
|
|
let mut disk_format = format.clone();
|
|
disk_format.erasure.this = format.erasure.sets[0][disk_idx];
|
|
save_format_file(&Some(disk.clone()), &Some(disk_format))
|
|
.await
|
|
.expect("format should be saved");
|
|
|
|
(dir, endpoint, disk)
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn get_online_disks_with_healing_and_info_keeps_disk_and_info_aligned() {
|
|
let disk_count = 8;
|
|
let format = FormatV3::new(1, disk_count);
|
|
|
|
let mut temp_dirs = Vec::with_capacity(disk_count);
|
|
let mut endpoints = Vec::with_capacity(disk_count);
|
|
let mut disks = Vec::with_capacity(disk_count);
|
|
|
|
for disk_idx in 0..disk_count {
|
|
let (temp_dir, endpoint, disk) = make_formatted_local_disk(disk_idx, &format).await;
|
|
temp_dirs.push(temp_dir);
|
|
endpoints.push(endpoint);
|
|
disks.push(Some(disk));
|
|
}
|
|
|
|
let set_disks = SetDisks::new(
|
|
"test-owner".to_string(),
|
|
Arc::new(RwLock::new(disks)),
|
|
disk_count,
|
|
disk_count / 2,
|
|
0,
|
|
0,
|
|
endpoints,
|
|
format,
|
|
Vec::new(),
|
|
)
|
|
.await;
|
|
|
|
for _ in 0..32 {
|
|
let (online_disks, infos, healing) = set_disks.get_online_disks_with_healing_and_info(false).await;
|
|
assert_eq!(healing, 0);
|
|
assert_eq!(online_disks.len(), disk_count);
|
|
assert_eq!(infos.len(), disk_count);
|
|
|
|
for (disk, info) in online_disks.iter().zip(infos.iter()) {
|
|
assert!(
|
|
info.error.is_empty(),
|
|
"unexpected disk_info error for {}: {}",
|
|
disk.endpoint(),
|
|
info.error
|
|
);
|
|
assert_eq!(info.endpoint, disk.endpoint().to_string());
|
|
assert_eq!(
|
|
info.id,
|
|
disk.get_disk_id().await.expect("disk id lookup should succeed"),
|
|
"disk info should stay aligned with disk {}",
|
|
disk.endpoint()
|
|
);
|
|
}
|
|
}
|
|
|
|
drop(temp_dirs);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn drive_membership_snapshot_filters_offline_disks_from_candidates() {
|
|
let disk_count = 4;
|
|
let format = FormatV3::new(1, disk_count);
|
|
|
|
let mut temp_dirs = Vec::with_capacity(disk_count);
|
|
let mut endpoints = Vec::with_capacity(disk_count);
|
|
let mut disks = Vec::with_capacity(disk_count);
|
|
|
|
for disk_idx in 0..disk_count {
|
|
let (temp_dir, endpoint, disk) = make_formatted_local_disk(disk_idx, &format).await;
|
|
temp_dirs.push(temp_dir);
|
|
endpoints.push(endpoint);
|
|
disks.push(Some(disk));
|
|
}
|
|
|
|
let set_disks = SetDisks::new(
|
|
"test-owner".to_string(),
|
|
Arc::new(RwLock::new(disks)),
|
|
disk_count,
|
|
disk_count / 2,
|
|
0,
|
|
0,
|
|
endpoints,
|
|
format,
|
|
Vec::new(),
|
|
)
|
|
.await;
|
|
|
|
let all_disks = set_disks.get_disks_internal().await;
|
|
all_disks[1]
|
|
.as_ref()
|
|
.expect("disk 1 should exist")
|
|
.force_runtime_state_for_test(disk::health_state::RuntimeDriveHealthState::Suspect);
|
|
all_disks[2]
|
|
.as_ref()
|
|
.expect("disk 2 should exist")
|
|
.force_runtime_state_for_test(disk::health_state::RuntimeDriveHealthState::Returning);
|
|
all_disks[3]
|
|
.as_ref()
|
|
.expect("disk 3 should exist")
|
|
.force_runtime_state_for_test(disk::health_state::RuntimeDriveHealthState::Offline);
|
|
|
|
let snapshot = set_disks.drive_membership_snapshot().await;
|
|
assert_eq!(snapshot.online.len(), 1);
|
|
assert_eq!(snapshot.suspect.len(), 1);
|
|
assert_eq!(snapshot.returning.len(), 1);
|
|
assert_eq!(snapshot.offline.len(), 1);
|
|
assert_eq!(snapshot.scanner_heal_candidates().len(), 3);
|
|
|
|
let strict_online = set_disks.get_online_disks().await;
|
|
assert_eq!(strict_online.len(), 1, "strict online selection should exclude suspect/returning/offline");
|
|
|
|
let (online_disks, infos, healing) = set_disks.get_online_disks_with_healing_and_info(false).await;
|
|
assert_eq!(healing, 0);
|
|
assert_eq!(online_disks.len(), 3);
|
|
assert_eq!(infos.len(), 3);
|
|
assert!(
|
|
online_disks
|
|
.iter()
|
|
.all(|disk| { disk.runtime_state() != disk::health_state::RuntimeDriveHealthState::Offline }),
|
|
"offline disks should be filtered by membership snapshot"
|
|
);
|
|
|
|
drop(temp_dirs);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn get_online_disks_with_healing_and_info_reprobes_runtime_candidates_once() {
|
|
let disk_count = 4;
|
|
let format = FormatV3::new(1, disk_count);
|
|
|
|
let mut temp_dirs = Vec::with_capacity(disk_count);
|
|
let mut endpoints = Vec::with_capacity(disk_count);
|
|
let mut disks = Vec::with_capacity(disk_count);
|
|
|
|
for disk_idx in 0..disk_count {
|
|
let (temp_dir, endpoint, disk) = make_formatted_local_disk(disk_idx, &format).await;
|
|
temp_dirs.push(temp_dir);
|
|
endpoints.push(endpoint);
|
|
disks.push(Some(disk));
|
|
}
|
|
|
|
let set_disks = SetDisks::new(
|
|
"test-owner".to_string(),
|
|
Arc::new(RwLock::new(disks)),
|
|
disk_count,
|
|
disk_count / 2,
|
|
0,
|
|
0,
|
|
endpoints,
|
|
format,
|
|
Vec::new(),
|
|
)
|
|
.await;
|
|
|
|
let all_disks = set_disks.get_disks_internal().await;
|
|
for disk in all_disks.iter().flatten() {
|
|
disk.force_runtime_state_for_test(disk::health_state::RuntimeDriveHealthState::Returning);
|
|
}
|
|
|
|
let (online_disks, infos, healing) = set_disks.get_online_disks_with_healing_and_info(false).await;
|
|
assert_eq!(healing, 0);
|
|
assert_eq!(online_disks.len(), disk_count);
|
|
assert_eq!(infos.len(), disk_count);
|
|
assert!(
|
|
infos.iter().all(|info| info.error.is_empty()),
|
|
"runtime reprobe should recover a usable candidate set without probe errors"
|
|
);
|
|
|
|
drop(temp_dirs);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn renew_disk_enables_health_monitoring_for_recovered_disk() {
|
|
let disk_count = 4;
|
|
let format = FormatV3::new(1, disk_count);
|
|
|
|
let mut temp_dirs = Vec::with_capacity(disk_count);
|
|
let mut endpoints = Vec::with_capacity(disk_count);
|
|
|
|
for disk_idx in 0..disk_count {
|
|
let (temp_dir, endpoint, _) = make_formatted_local_disk(disk_idx, &format).await;
|
|
temp_dirs.push(temp_dir);
|
|
endpoints.push(endpoint);
|
|
}
|
|
|
|
let set_disks = SetDisks::new(
|
|
"test-owner".to_string(),
|
|
Arc::new(RwLock::new(vec![None; disk_count])),
|
|
disk_count,
|
|
disk_count / 2,
|
|
0,
|
|
0,
|
|
endpoints.clone(),
|
|
format,
|
|
Vec::new(),
|
|
)
|
|
.await;
|
|
|
|
set_disks.renew_disk(&endpoints[0]).await;
|
|
|
|
let disks = set_disks.get_disks_internal().await;
|
|
let renewed_disk = disks[0].as_ref().expect("renew_disk should attach the recovered disk");
|
|
assert!(
|
|
renewed_disk.health_check_enabled_for_test(),
|
|
"renewed disks must keep health monitoring enabled so later faulty marks can recover"
|
|
);
|
|
renewed_disk
|
|
.disk_info(&DiskInfoOptions::default())
|
|
.await
|
|
.expect("renewed disk_info should record a drive API metric");
|
|
renewed_disk.force_runtime_state_for_test(disk::health_state::RuntimeDriveHealthState::Offline);
|
|
|
|
set_disks.renew_disk(&endpoints[0]).await;
|
|
|
|
let disks = set_disks.get_disks_internal().await;
|
|
let renewed_again = disks[0]
|
|
.as_ref()
|
|
.expect("second renew_disk should keep the recovered disk attached");
|
|
assert_eq!(
|
|
renewed_again
|
|
.metrics_snapshot()
|
|
.and_then(|metrics| metrics.api_calls.get("disk_info").copied()),
|
|
Some(1),
|
|
"disk reconnect must preserve the local drive metrics tracker epoch"
|
|
);
|
|
|
|
drop(temp_dirs);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn renew_disk_rejects_a_format_from_another_slot_or_cluster() {
|
|
let disk_count = 3;
|
|
let format = FormatV3::new(1, disk_count);
|
|
let mut temp_dirs = Vec::with_capacity(disk_count);
|
|
let mut endpoints = Vec::with_capacity(disk_count);
|
|
let mut fixture_disks = Vec::with_capacity(disk_count);
|
|
|
|
for disk_idx in 0..disk_count {
|
|
let (temp_dir, endpoint, disk) = make_formatted_local_disk(disk_idx, &format).await;
|
|
temp_dirs.push(temp_dir);
|
|
endpoints.push(endpoint);
|
|
fixture_disks.push(disk);
|
|
}
|
|
|
|
let set_disks = SetDisks::new(
|
|
"test-owner".to_string(),
|
|
Arc::new(RwLock::new(vec![Some(fixture_disks[0].clone()), None, None])),
|
|
disk_count,
|
|
disk_count / 2,
|
|
0,
|
|
0,
|
|
endpoints.clone(),
|
|
format.clone(),
|
|
Vec::new(),
|
|
)
|
|
.await;
|
|
|
|
let mut other_cluster_format = format.clone();
|
|
other_cluster_format.id = Uuid::new_v4();
|
|
other_cluster_format.erasure.this = format.erasure.sets[0][2];
|
|
save_format_file(&Some(fixture_disks[2].clone()), &Some(other_cluster_format))
|
|
.await
|
|
.expect("other-cluster format should be written for the rejection test");
|
|
|
|
set_disks.renew_disk(&endpoints[2]).await;
|
|
|
|
let disks = set_disks.get_disks_internal().await;
|
|
assert_eq!(
|
|
disks[0]
|
|
.as_ref()
|
|
.expect("the canonical first slot must remain attached")
|
|
.endpoint(),
|
|
endpoints[0]
|
|
);
|
|
assert!(
|
|
disks[2].is_none(),
|
|
"a disk from another deployment must remain detached even when its slot UUID matches"
|
|
);
|
|
|
|
let mut correct_format = format.clone();
|
|
correct_format.erasure.this = format.erasure.sets[0][2];
|
|
let replacement_disk = new_disk(
|
|
&endpoints[2],
|
|
&DiskOption {
|
|
cleanup: false,
|
|
health_check: false,
|
|
},
|
|
)
|
|
.await
|
|
.expect("third endpoint should reopen after other-cluster rejection");
|
|
save_format_file(&Some(replacement_disk), &Some(correct_format))
|
|
.await
|
|
.expect("correct slot format should be restored");
|
|
|
|
set_disks.renew_disk(&endpoints[2]).await;
|
|
|
|
let disks = set_disks.get_disks_internal().await;
|
|
assert_eq!(
|
|
disks[0]
|
|
.as_ref()
|
|
.expect("the canonical first slot must remain attached")
|
|
.endpoint(),
|
|
endpoints[0]
|
|
);
|
|
assert_eq!(disks[2].as_ref().expect("the restored third slot should attach").endpoint(), endpoints[2]);
|
|
|
|
let third_disk = disks[2].clone();
|
|
let mut wrong_slot_format = format.clone();
|
|
wrong_slot_format.erasure.this = format.erasure.sets[0][0];
|
|
save_format_file(&third_disk, &Some(wrong_slot_format))
|
|
.await
|
|
.expect("wrong-slot format should be written for the rejection test");
|
|
set_disks.disks.write().await[2] = None;
|
|
|
|
let mut misplaced_endpoint = endpoints[2].clone();
|
|
misplaced_endpoint.set_disk_index(0);
|
|
set_disks.renew_disk(&misplaced_endpoint).await;
|
|
|
|
let disks = set_disks.get_disks_internal().await;
|
|
assert_eq!(
|
|
disks[0]
|
|
.as_ref()
|
|
.expect("the canonical first slot must remain attached")
|
|
.endpoint(),
|
|
endpoints[0]
|
|
);
|
|
assert!(disks[2].is_none(), "a disk claiming another endpoint's slot must remain detached");
|
|
|
|
drop(temp_dirs);
|
|
}
|
|
|
|
// SetDisks split P0 (#816): the borrow handle must mirror the core state and
|
|
// the List operation family must run identically through it.
|
|
#[tokio::test]
|
|
async fn set_disks_ctx_mirrors_core_and_drives_list_operations() {
|
|
use crate::set_disk::ops::list::ListOperations;
|
|
|
|
let disk_count = 4;
|
|
let format = FormatV3::new(1, disk_count);
|
|
|
|
let mut temp_dirs = Vec::with_capacity(disk_count);
|
|
let mut endpoints = Vec::with_capacity(disk_count);
|
|
let mut disks = Vec::with_capacity(disk_count);
|
|
|
|
for disk_idx in 0..disk_count {
|
|
let (temp_dir, endpoint, disk) = make_formatted_local_disk(disk_idx, &format).await;
|
|
temp_dirs.push(temp_dir);
|
|
endpoints.push(endpoint);
|
|
disks.push(Some(disk));
|
|
}
|
|
|
|
let set_disks = SetDisks::new(
|
|
"test-owner".to_string(),
|
|
Arc::new(RwLock::new(disks)),
|
|
disk_count,
|
|
disk_count / 2,
|
|
0,
|
|
0,
|
|
endpoints.clone(),
|
|
format,
|
|
Vec::new(),
|
|
)
|
|
.await;
|
|
|
|
// The handle reads through to the same core state, not a copy.
|
|
let ctx = set_disks.ctx();
|
|
assert_eq!(ctx.locker_owner(), "test-owner");
|
|
assert_eq!(ctx.set_drive_count(), disk_count);
|
|
assert_eq!(ctx.default_parity_count(), disk_count / 2);
|
|
assert_eq!(ctx.set_index(), 0);
|
|
assert_eq!(ctx.pool_index(), 0);
|
|
assert_eq!(ctx.set_endpoints().len(), endpoints.len());
|
|
assert!(ctx.lockers().is_empty());
|
|
assert!(
|
|
Arc::ptr_eq(ctx.disks(), &set_disks.disks),
|
|
"ctx must borrow the core disks, not clone them"
|
|
);
|
|
assert!(std::ptr::eq(ctx.core(), &*set_disks), "ctx must borrow the core, not clone it");
|
|
assert!(std::ptr::eq(ctx.format(), &set_disks.format), "ctx must borrow the core format");
|
|
|
|
// The List family runs through the borrow handle with unchanged
|
|
// behavior: delete_all reports success even when the prefix is absent.
|
|
ListOperations::new(set_disks.ctx())
|
|
.delete_all("nonexistent-bucket", "nonexistent-prefix")
|
|
.await
|
|
.expect("delete_all via borrow handle should succeed");
|
|
set_disks
|
|
.delete_all("nonexistent-bucket", "nonexistent-prefix")
|
|
.await
|
|
.expect("delete_all public entry should stay in sync");
|
|
|
|
drop(temp_dirs);
|
|
}
|
|
}
|