fix(scanner): make distributed usage convergence authoritative (#5151)

* fix(scanner): make distributed usage cycles authoritative

* fix(scanner): close distributed refresh races

* fix(config): align scanner reload integration

* fix(admin): scope config test helpers

* fix(scanner): harden distributed usage convergence

* fix(scanner): preserve rolling activity compatibility

* fix(admin): expose non-secret optional config values

* fix(scanner): acknowledge distributed dirty usage

* fix(ecstore): make bucket mutations cancellation safe

* fix(scanner): preserve pending dirty acknowledgements

* test(obs): account for superseded scanner metric

* fix(api): reject excess detached bucket mutations

* test: close scanner convergence coverage gaps

* fix(scanner): make path tracking cleanup one-shot

---------

Co-authored-by: Henry Guo <marshawcoco@users.noreply.github.com>
Co-authored-by: houseme <housemecn@gmail.com>
This commit is contained in:
Henry Guo
2026-07-25 18:45:16 +08:00
committed by GitHub
parent 0364523dad
commit a63b79004c
69 changed files with 18073 additions and 1585 deletions
+15 -14
View File
@@ -244,11 +244,12 @@ pub mod config {
pub mod com {
pub use crate::config::com::{
COMMA_SEPARATED_LISTS, CONFIG_PREFIX, ENV_CONFIG_RECOVER_ON_CORRUPTION, STORAGE_CLASS_SUB_SYS,
ServerConfigCorruptError, delete_config, is_server_config_corrupt_error, lookup_configs, read_config,
read_config_no_lock, read_config_with_metadata, read_config_without_migrate, read_config_without_migrate_no_lock,
read_existing_server_config_no_lock, save_config, save_config_no_lock, save_config_with_opts, save_server_config,
save_server_config_no_lock, try_migrate_server_config, with_config_object_read_lock, with_config_object_write_lock,
with_server_config_read_lock, with_server_config_write_lock,
ServerConfigCorruptError, ServerConfigSnapshot, delete_config, is_server_config_corrupt_error, lookup_configs,
read_config, read_config_no_lock, read_config_with_metadata, read_config_without_migrate,
read_config_without_migrate_no_lock, read_existing_server_config_no_lock, read_server_config_snapshot, save_config,
save_config_no_lock, save_config_with_opts, save_server_config, save_server_config_no_lock,
save_server_config_snapshot, server_config_path, try_migrate_server_config, with_config_object_read_lock,
with_config_object_write_lock, with_server_config_read_lock, with_server_config_write_lock,
};
}
@@ -287,9 +288,9 @@ pub mod disk {
pub use crate::disk::{
BATCH_READ_VERSION_MAX_ITEMS, BUCKET_META_PREFIX, BatchReadVersionItem, BatchReadVersionReq, BatchReadVersionResp,
CheckPartsResp, DeleteOptions, Disk, DiskAPI, DiskInfo, DiskInfoOptions, DiskLocation, DiskOption, DiskStore,
FileInfoVersions, FileReader, FileWriter, HEALING_MARKER_PATH, OldCurrentSize, RUSTFS_META_BUCKET, ReadMultipleReq,
ReadMultipleResp, ReadOptions, RenameDataResp, STORAGE_FORMAT_FILE, UpdateMetadataOpts, VolumeInfo, WalkDirOptions,
new_disk, validate_batch_read_version_item_count,
FileInfoVersions, FileReader, FileWriter, HEALING_MARKER_PATH, NsScannerOpenRequest, OldCurrentSize, RUSTFS_META_BUCKET,
ReadMultipleReq, ReadMultipleResp, ReadOptions, RenameDataResp, STORAGE_FORMAT_FILE, UpdateMetadataOpts, VolumeInfo,
WalkDirOptions, new_disk, validate_batch_read_version_item_count,
};
pub use bytes::Bytes;
pub use endpoint::Endpoint;
@@ -390,12 +391,12 @@ pub mod rio {
pub mod rpc {
pub use crate::cluster::rpc::{
LocalPeerS3Client, PEER_RESTDRY_RUN, PEER_RESTSIGNAL, PEER_RESTSUB_SYS, PeerRestClient, PeerS3Client,
SERVICE_SIGNAL_REFRESH_CONFIG, SERVICE_SIGNAL_RELOAD_DYNAMIC, ScannerPeerActivity, TONIC_RPC_PREFIX, TonicInterceptor,
gen_signature_headers, gen_tonic_signature_headers, gen_tonic_signature_interceptor, node_service_time_out_client,
node_service_time_out_client_no_auth, normalize_tonic_rpc_audience, set_tonic_canonical_body_digest,
sign_tonic_rpc_response_proof, verify_rpc_signature, verify_tonic_canonical_body_digest, verify_tonic_rpc_response_proof,
verify_tonic_rpc_signature,
LocalPeerS3Client, PEER_RESTDRY_RUN, PEER_RESTSIGNAL, PEER_RESTSUB_SYS, PeerRestClient, PeerS3Client, S3PeerSys,
SERVICE_SIGNAL_REFRESH_CONFIG, SERVICE_SIGNAL_RELOAD_DYNAMIC, ScannerBucketListing, ScannerPeerActivity,
TONIC_RPC_PREFIX, TonicInterceptor, gen_signature_headers, gen_tonic_signature_headers, gen_tonic_signature_interceptor,
node_service_time_out_client, node_service_time_out_client_no_auth, normalize_tonic_rpc_audience,
set_tonic_canonical_body_digest, sign_ns_scanner_capability, sign_tonic_rpc_response_proof, verify_rpc_signature,
verify_tonic_canonical_body_digest, verify_tonic_rpc_response_proof, verify_tonic_rpc_signature,
};
}
@@ -27,6 +27,7 @@
//! Advisory: <https://github.com/rustfs/rustfs/security/advisories/GHSA-r5qv-rc46-hv8q>
use crate::cluster::rpc::context_propagation::{inject_request_id_into_http_headers, inject_trace_context_into_http_headers};
use crate::storage_api_contracts::internode::NS_SCANNER_PROTOCOL_VERSION;
use base64::Engine as _;
use base64::engine::general_purpose;
use hmac::{Hmac, KeyInit, Mac};
@@ -61,6 +62,7 @@ const UNSIGNED_PAYLOAD_NONCE: &str = "unsigned";
const SIGNATURE_VALID_DURATION: i64 = 300; // 5 minutes
const REPLAY_CACHE_RETENTION: Duration = Duration::from_secs(601);
const MAX_REPLAY_PROTECTED_NONCES: usize = 65_536;
const NS_SCANNER_CAPABILITY_AUTH_DOMAIN: &[u8] = b"rustfs-ns-scanner-capability-v3";
pub const TONIC_RPC_PREFIX: &str = "/node_service.NodeService";
static INTERNODE_RPC_SIGNATURE_STRICT: LazyLock<bool> = LazyLock::new(|| {
get_env_bool(
@@ -209,6 +211,42 @@ fn verify_signature(secret: &str, url: &str, method: &Method, timestamp: i64, si
mac.verify_slice(&signature).is_ok()
}
fn update_ns_scanner_capability_mac(mac: &mut HmacSha256, challenge: Uuid, server_epoch: Uuid) {
mac.update(NS_SCANNER_CAPABILITY_AUTH_DOMAIN);
mac.update(&NS_SCANNER_PROTOCOL_VERSION.to_be_bytes());
mac.update(challenge.as_bytes());
mac.update(server_epoch.as_bytes());
}
fn generate_ns_scanner_capability_proof(secret: &str, challenge: Uuid, server_epoch: Uuid) -> std::io::Result<Vec<u8>> {
if challenge.is_nil() || server_epoch.is_nil() {
return Err(std::io::Error::other("Invalid namespace scanner capability scope"));
}
let mut mac =
<HmacSha256 as KeyInit>::new_from_slice(secret.as_bytes()).map_err(|_| std::io::Error::other("Invalid RPC HMAC key"))?;
update_ns_scanner_capability_mac(&mut mac, challenge, server_epoch);
Ok(mac.finalize().into_bytes().to_vec())
}
fn verify_ns_scanner_capability_proof(secret: &str, challenge: Uuid, server_epoch: Uuid, proof: &[u8]) -> std::io::Result<()> {
if challenge.is_nil() || server_epoch.is_nil() {
return Err(std::io::Error::other("Invalid namespace scanner capability scope"));
}
let mut mac =
<HmacSha256 as KeyInit>::new_from_slice(secret.as_bytes()).map_err(|_| std::io::Error::other("Invalid RPC HMAC key"))?;
update_ns_scanner_capability_mac(&mut mac, challenge, server_epoch);
mac.verify_slice(proof)
.map_err(|_| std::io::Error::new(std::io::ErrorKind::PermissionDenied, "Invalid namespace scanner capability proof"))
}
pub fn sign_ns_scanner_capability(challenge: Uuid, server_epoch: Uuid) -> std::io::Result<Vec<u8>> {
generate_ns_scanner_capability_proof(&get_shared_secret()?, challenge, server_epoch)
}
pub fn verify_ns_scanner_capability(challenge: Uuid, server_epoch: Uuid, proof: &[u8]) -> std::io::Result<()> {
verify_ns_scanner_capability_proof(&get_shared_secret()?, challenge, server_epoch, proof)
}
#[derive(Clone, Copy)]
struct SignatureV2Scope<'a> {
audience: &'a str,
@@ -629,6 +667,20 @@ mod tests {
runtime_sources::ensure_test_rpc_secret();
}
#[test]
fn namespace_scanner_capability_proof_binds_challenge_and_server_epoch() {
let secret = "test-scanner-capability-secret";
let challenge = Uuid::new_v4();
let server_epoch = Uuid::new_v4();
let proof =
generate_ns_scanner_capability_proof(secret, challenge, server_epoch).expect("capability proof should be generated");
assert!(verify_ns_scanner_capability_proof(secret, challenge, server_epoch, &proof).is_ok());
assert!(verify_ns_scanner_capability_proof(secret, Uuid::new_v4(), server_epoch, &proof).is_err());
assert!(verify_ns_scanner_capability_proof(secret, challenge, Uuid::new_v4(), &proof).is_err());
assert!(verify_ns_scanner_capability_proof("different-secret", challenge, server_epoch, &proof).is_err());
}
/// Security regression for GHSA-r5qv-rc46-hv8q (internode RPC fail-closed,
/// fixed in rustfs/rustfs#4402): secret resolution must never silently fall
/// back to a default/empty shared secret. Missing and default secrets both
@@ -12,11 +12,14 @@
// See the License for the specific language governing permissions and
// limitations under the License.
use crate::cluster::rpc::build_auth_headers;
use crate::cluster::rpc::{build_auth_headers, verify_ns_scanner_capability};
use crate::disk::error::{Error, Result};
use crate::disk::{FileReader, FileWriter};
use crate::storage_api_contracts::internode::{
WALK_DIR_BODY_SHA256_QUERY, WALK_DIR_STREAM_COMPLETION_QUERY, WALK_DIR_STREAM_COMPLETION_V1,
NS_SCANNER_BODY_SHA256_QUERY, NS_SCANNER_CAPABILITY_CHALLENGE_QUERY, NS_SCANNER_CYCLE_QUERY, NS_SCANNER_LEADER_EPOCH_QUERY,
NS_SCANNER_PROTOCOL_VERSION, NS_SCANNER_PROTOCOL_VERSION_QUERY, NS_SCANNER_REQUEST_ID_QUERY, NS_SCANNER_SERVER_EPOCH_QUERY,
NS_SCANNER_SESSION_ID_QUERY, NS_SCANNER_SESSION_SEQUENCE_QUERY, NsScannerCapabilityResponse, WALK_DIR_BODY_SHA256_QUERY,
WALK_DIR_STREAM_COMPLETION_QUERY, WALK_DIR_STREAM_COMPLETION_V1,
};
use async_trait::async_trait;
use http::{HeaderMap, HeaderValue, Method, header::CONTENT_TYPE};
@@ -28,13 +31,18 @@ use rustfs_rio::{HttpReader, HttpWriter};
use sha2::{Digest, Sha256};
use std::sync::{Arc, OnceLock};
use std::time::Duration;
use tokio::io::AsyncReadExt;
use uuid::Uuid;
static INTERNODE_DATA_TRANSPORT: OnceLock<std::result::Result<Arc<dyn InternodeDataTransport>, String>> = OnceLock::new();
const READ_FILE_STREAM_PATH: &str = "/rustfs/rpc/read_file_stream";
const PUT_FILE_STREAM_PATH: &str = "/rustfs/rpc/put_file_stream";
const WALK_DIR_PATH: &str = "/rustfs/rpc/walk_dir";
const NS_SCANNER_PATH: &str = "/rustfs/rpc/ns_scanner";
const NS_SCANNER_MAX_CAPABILITY_RESPONSE_SIZE: usize = 1024;
const CONTENT_TYPE_JSON: &str = "application/json";
const CONTENT_TYPE_MSGPACK: &str = "application/msgpack";
fn unsupported_transport_message(transport: &str) -> String {
format!(
@@ -101,6 +109,25 @@ pub struct WalkDirStreamRequest {
pub stall_timeout: Option<Duration>,
}
#[derive(Debug, Clone)]
pub struct NsScannerStreamRequest {
pub endpoint: String,
pub disk: String,
pub request_id: Uuid,
pub server_epoch: Uuid,
pub session_id: Uuid,
pub session_sequence: u64,
pub next_cycle: u64,
pub leader_epoch: u64,
pub body: Vec<u8>,
pub stall_timeout: Option<Duration>,
}
#[derive(Debug, Clone)]
pub struct NsScannerCapabilityRequest {
pub endpoint: String,
}
/// Data-plane stream opener used by `RemoteDisk`.
///
/// This boundary is limited to remote disk streams that can move large payloads.
@@ -114,6 +141,12 @@ pub trait InternodeDataTransport: Send + Sync + std::fmt::Debug {
async fn open_read(&self, request: ReadStreamRequest) -> Result<FileReader>;
async fn open_write(&self, request: WriteStreamRequest) -> Result<FileWriter>;
async fn open_walk_dir(&self, request: WalkDirStreamRequest) -> Result<FileReader>;
async fn open_ns_scanner(&self, _request: NsScannerStreamRequest) -> Result<FileReader> {
Err(Error::MethodNotAllowed)
}
async fn probe_ns_scanner(&self, _request: NsScannerCapabilityRequest) -> Result<Uuid> {
Err(Error::MethodNotAllowed)
}
fn name(&self) -> &'static str;
fn capabilities(&self) -> InternodeDataTransportCapabilities;
}
@@ -148,6 +181,39 @@ impl InternodeDataTransport for TcpHttpInternodeDataTransport {
))
}
async fn open_ns_scanner(&self, request: NsScannerStreamRequest) -> Result<FileReader> {
let url = build_ns_scanner_url(&request);
let mut headers = msgpack_headers();
build_auth_headers(&url, &Method::POST, &mut headers)?;
Ok(Box::new(
HttpReader::new_with_stall_timeout(url, Method::POST, headers, Some(request.body), request.stall_timeout).await?,
))
}
async fn probe_ns_scanner(&self, request: NsScannerCapabilityRequest) -> Result<Uuid> {
let challenge = Uuid::new_v4();
let url = build_ns_scanner_capability_url(&request, challenge);
let mut headers = msgpack_headers();
build_auth_headers(&url, &Method::GET, &mut headers)?;
let reader = HttpReader::new(url, Method::GET, headers, None).await?;
let mut body = Vec::new();
reader
.take(u64::try_from(NS_SCANNER_MAX_CAPABILITY_RESPONSE_SIZE + 1).unwrap_or(u64::MAX))
.read_to_end(&mut body)
.await?;
if body.is_empty() || body.len() > NS_SCANNER_MAX_CAPABILITY_RESPONSE_SIZE {
return Err(Error::other("invalid remote namespace scanner capability response size"));
}
let response: NsScannerCapabilityResponse =
rmp_serde::from_slice(&body).map_err(|_| Error::other("invalid remote namespace scanner capability response"))?;
if response.version != NS_SCANNER_PROTOCOL_VERSION || response.server_epoch.is_nil() {
return Err(Error::other("incompatible remote namespace scanner capability response"));
}
verify_ns_scanner_capability(challenge, response.server_epoch, &response.proof)
.map_err(|err| Error::other(format!("remote namespace scanner capability authentication failed: {err}")))?;
Ok(response.server_epoch)
}
fn name(&self) -> &'static str {
DEFAULT_INTERNODE_DATA_TRANSPORT
}
@@ -197,12 +263,54 @@ fn build_walk_dir_url(request: &WalkDirStreamRequest) -> String {
)
}
fn build_ns_scanner_url(request: &NsScannerStreamRequest) -> String {
let body_sha256 = hex_simd::encode_to_string(Sha256::digest(&request.body), hex_simd::AsciiCase::Lower);
format!(
"{}{}?disk={}&{}={}&{}={}&{}={}&{}={}&{}={}&{}={}&{}={}",
request.endpoint,
NS_SCANNER_PATH,
urlencoding::encode(&request.disk),
NS_SCANNER_REQUEST_ID_QUERY,
request.request_id,
NS_SCANNER_SERVER_EPOCH_QUERY,
request.server_epoch,
NS_SCANNER_SESSION_ID_QUERY,
request.session_id,
NS_SCANNER_SESSION_SEQUENCE_QUERY,
request.session_sequence,
NS_SCANNER_CYCLE_QUERY,
request.next_cycle,
NS_SCANNER_LEADER_EPOCH_QUERY,
request.leader_epoch,
NS_SCANNER_BODY_SHA256_QUERY,
body_sha256
)
}
fn build_ns_scanner_capability_url(request: &NsScannerCapabilityRequest, challenge: Uuid) -> String {
format!(
"{}{}?{}={}&{}={}",
request.endpoint,
NS_SCANNER_PATH,
NS_SCANNER_PROTOCOL_VERSION_QUERY,
NS_SCANNER_PROTOCOL_VERSION,
NS_SCANNER_CAPABILITY_CHALLENGE_QUERY,
challenge
)
}
fn json_headers() -> HeaderMap {
let mut headers = HeaderMap::new();
headers.insert(CONTENT_TYPE, HeaderValue::from_static(CONTENT_TYPE_JSON));
headers
}
fn msgpack_headers() -> HeaderMap {
let mut headers = HeaderMap::new();
headers.insert(CONTENT_TYPE, HeaderValue::from_static(CONTENT_TYPE_MSGPACK));
headers
}
fn build_internode_data_transport_result(
configured_transport: Option<&str>,
) -> std::result::Result<Arc<dyn InternodeDataTransport>, String> {
@@ -241,6 +349,65 @@ pub fn build_internode_data_transport_from_env() -> Result<Arc<dyn InternodeData
mod tests {
use super::*;
#[derive(Debug)]
struct LegacyTestTransport;
#[async_trait::async_trait]
impl InternodeDataTransport for LegacyTestTransport {
async fn open_read(&self, _request: ReadStreamRequest) -> Result<FileReader> {
Ok(Box::new(tokio::io::empty()))
}
async fn open_write(&self, _request: WriteStreamRequest) -> Result<FileWriter> {
Ok(Box::new(tokio::io::sink()))
}
async fn open_walk_dir(&self, _request: WalkDirStreamRequest) -> Result<FileReader> {
Ok(Box::new(tokio::io::empty()))
}
fn name(&self) -> &'static str {
"legacy-test"
}
fn capabilities(&self) -> InternodeDataTransportCapabilities {
InternodeDataTransportCapabilities::tcp_http()
}
}
#[tokio::test]
async fn legacy_transport_defaults_namespace_scanner_to_unsupported() {
let transport = LegacyTestTransport;
let probe_err = transport
.probe_ns_scanner(NsScannerCapabilityRequest {
endpoint: "http://node1:9000".to_string(),
})
.await
.expect_err("legacy transport should report namespace scanner as unsupported");
assert!(matches!(probe_err, Error::MethodNotAllowed));
let open_result = transport
.open_ns_scanner(NsScannerStreamRequest {
endpoint: "http://node1:9000".to_string(),
disk: "http://node1:9000/data/rustfs0".to_string(),
request_id: Uuid::new_v4(),
server_epoch: Uuid::new_v4(),
session_id: Uuid::new_v4(),
session_sequence: 0,
next_cycle: 7,
leader_epoch: 9,
body: Vec::new(),
stall_timeout: None,
})
.await;
let open_err = match open_result {
Ok(_) => panic!("legacy transport should not open namespace scanner streams"),
Err(err) => err,
};
assert!(matches!(open_err, Error::MethodNotAllowed));
}
#[test]
fn tcp_http_capabilities_are_behavior_preserving() {
let transport = TcpHttpInternodeDataTransport;
@@ -322,6 +489,57 @@ mod tests {
);
}
#[test]
fn ns_scanner_url_binds_body_and_encodes_disk_ref() {
let request_id = Uuid::parse_str("11111111-2222-4333-8444-555555555555").expect("request ID");
let server_epoch = Uuid::parse_str("aaaaaaaa-bbbb-4ccc-8ddd-eeeeeeeeeeee").expect("server epoch");
let session_id = Uuid::parse_str("99999999-8888-4777-8666-555555555555").expect("session ID");
let url = build_ns_scanner_url(&NsScannerStreamRequest {
endpoint: "http://node1:9000".to_string(),
disk: "http://node1:9000/data/rustfs0".to_string(),
request_id,
server_epoch,
session_id,
session_sequence: 3,
next_cycle: 7,
leader_epoch: 9,
body: b"scanner-request".to_vec(),
stall_timeout: None,
});
assert_eq!(
url,
concat!(
"http://node1:9000/rustfs/rpc/ns_scanner?disk=http%3A%2F%2Fnode1%3A9000%2Fdata%2Frustfs0",
"&ns_scanner_request_id=11111111-2222-4333-8444-555555555555",
"&ns_scanner_server_epoch=aaaaaaaa-bbbb-4ccc-8ddd-eeeeeeeeeeee",
"&ns_scanner_session_id=99999999-8888-4777-8666-555555555555",
"&ns_scanner_session_sequence=3",
"&ns_scanner_cycle=7",
"&ns_scanner_leader_epoch=9",
"&ns_scanner_body_sha256=c958f15ca28422275c1245399f4c44eaba628ca453fcd77d6b3d4484573e4387"
)
);
}
#[test]
fn ns_scanner_capability_url_binds_version_and_challenge() {
let challenge = Uuid::parse_str("12345678-1234-4234-8234-123456789abc").expect("challenge");
let url = build_ns_scanner_capability_url(
&NsScannerCapabilityRequest {
endpoint: "http://node1:9000".to_string(),
},
challenge,
);
assert_eq!(
url,
format!(
"http://node1:9000/rustfs/rpc/ns_scanner?ns_scanner_protocol={NS_SCANNER_PROTOCOL_VERSION}&ns_scanner_challenge={challenge}"
)
);
}
#[test]
fn transport_config_defaults_to_tcp_http() {
let transport = build_internode_data_transport(None).unwrap();
+3 -3
View File
@@ -30,8 +30,8 @@ pub use client::{
};
pub use http_auth::{
TONIC_RPC_PREFIX, build_auth_headers, gen_signature_headers, gen_tonic_signature_headers, normalize_tonic_rpc_audience,
set_tonic_canonical_body_digest, sign_tonic_rpc_response_proof, verify_rpc_signature, verify_tonic_canonical_body_digest,
verify_tonic_rpc_response_proof, verify_tonic_rpc_signature,
set_tonic_canonical_body_digest, sign_ns_scanner_capability, sign_tonic_rpc_response_proof, verify_ns_scanner_capability,
verify_rpc_signature, verify_tonic_canonical_body_digest, verify_tonic_rpc_response_proof, verify_tonic_rpc_signature,
};
#[cfg(test)]
pub(crate) use internode_data_transport::TcpHttpInternodeDataTransport;
@@ -41,6 +41,6 @@ pub use peer_rest_client::{
SERVICE_SIGNAL_RELOAD_DYNAMIC, ScannerPeerActivity,
};
pub(crate) use peer_s3_client::heal_bucket_local_on_disks;
pub use peer_s3_client::{LocalPeerS3Client, PeerS3Client, S3PeerSys};
pub use peer_s3_client::{LocalPeerS3Client, PeerS3Client, S3PeerSys, ScannerBucketListing, ScannerSetBucketListing};
pub use remote_disk::RemoteDisk;
pub use remote_locker::RemoteClient;
@@ -18,6 +18,9 @@ use crate::cluster::rpc::client::{
};
use crate::cluster::rpc::{set_tonic_canonical_body_digest, verify_tonic_rpc_response_proof};
use crate::error::{Error, Result};
use crate::storage_api_contracts::internode::{
SCANNER_ACTIVITY_LEGACY_PROTOCOL_VERSION, SCANNER_ACTIVITY_PREVIOUS_PROTOCOL_VERSION, SCANNER_ACTIVITY_PROTOCOL_VERSION,
};
use crate::{
bucket::replication::BucketStats,
disk::disk_store::{get_drive_active_check_interval, get_drive_active_check_timeout},
@@ -27,6 +30,7 @@ use crate::{
};
use bytes::Bytes;
use rmp_serde::{Deserializer, Serializer};
use rustfs_config::{HEAL_SUB_SYS, SCANNER_SUB_SYS};
use rustfs_madmin::{
ServerProperties,
health::{Cpus, MemInfo, OsInfo, Partitions, ProcInfo, SysConfig, SysErrors, SysServices},
@@ -97,15 +101,34 @@ fn decode_bucket_stats_response(response: GetBucketStatsDataResponse) -> Result<
Ok(stats)
}
fn validate_signal_service_protocol(sig: u64, sub_sys: &str, protocol_version: u32) -> Result<()> {
if sig == SERVICE_SIGNAL_RELOAD_DYNAMIC
&& matches!(sub_sys, SCANNER_SUB_SYS | HEAL_SUB_SYS)
&& protocol_version < rustfs_protos::DYNAMIC_CONFIG_PROTOCOL_VERSION
{
return Err(Error::other(format!("peer does not support dynamic {sub_sys} config convergence")));
}
Ok(())
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct ScannerPeerActivity {
pub instance_id: String,
pub namespace_generation: u64,
pub maintenance_generation: u64,
pub protocol_version: u32,
pub topology_digest: Option<[u8; 32]>,
pub data_movement_active: Option<bool>,
pub dirty_usage_generation: Option<u64>,
pub dirty_usage_pending: Option<bool>,
}
fn decode_scanner_activity(response: ScannerActivityResponse) -> Result<ScannerPeerActivity> {
let instance_id = response.instance_id;
fn decode_scanner_activity_with_verifier(
response: ScannerActivityResponse,
challenge: &[u8; 16],
verify_proof: impl FnOnce(&[u8], &[u8]) -> Result<()>,
) -> Result<ScannerPeerActivity> {
let instance_id = &response.instance_id;
if instance_id.len() != 32
|| !instance_id
.as_bytes()
@@ -114,10 +137,80 @@ fn decode_scanner_activity(response: ScannerActivityResponse) -> Result<ScannerP
{
return Err(Error::other("peer returned an invalid scanner activity instance ID"));
}
let (topology_digest, data_movement_active, dirty_usage_generation, dirty_usage_pending) = match response.protocol_version {
// RUSTFS_COMPAT_TODO(ns-scanner-rpc-v3): legacy response fields are unauthenticated. Remove after protocol v0 peers are unsupported.
SCANNER_ACTIVITY_LEGACY_PROTOCOL_VERSION
if response.topology_digest.is_empty()
&& response.response_proof.is_empty()
&& !response.data_movement_active
&& response.dirty_usage_generation == 0
&& !response.dirty_usage_pending =>
{
(None, None, None, None)
}
SCANNER_ACTIVITY_LEGACY_PROTOCOL_VERSION => {
return Err(Error::other("legacy scanner activity peer returned unexpected extended fields"));
}
SCANNER_ACTIVITY_PREVIOUS_PROTOCOL_VERSION => {
if response.dirty_usage_generation != 0 || response.dirty_usage_pending {
return Err(Error::other("scanner activity protocol v4 peer returned unauthenticated v5 fields"));
}
let canonical = rustfs_protos::canonical_scanner_activity_v4_response_body(challenge, &response)
.map_err(|_| Error::other("peer scanner activity response is too large to authenticate"))?;
verify_proof(&canonical, &response.response_proof)?;
(
Some(
response
.topology_digest
.as_ref()
.try_into()
.map_err(|_| Error::other("peer returned an invalid scanner topology digest"))?,
),
Some(response.data_movement_active),
None,
None,
)
}
SCANNER_ACTIVITY_PROTOCOL_VERSION => {
if response.dirty_usage_pending && response.dirty_usage_generation == 0 {
return Err(Error::other("scanner activity peer returned pending dirty usage without a generation"));
}
let canonical = rustfs_protos::canonical_scanner_activity_response_body(challenge, &response)
.map_err(|_| Error::other("peer scanner activity response is too large to authenticate"))?;
verify_proof(&canonical, &response.response_proof)?;
(
Some(
response
.topology_digest
.as_ref()
.try_into()
.map_err(|_| Error::other("peer returned an invalid scanner topology digest"))?,
),
Some(response.data_movement_active),
Some(response.dirty_usage_generation),
Some(response.dirty_usage_pending),
)
}
version => {
return Err(Error::other(format!("peer returned unsupported scanner activity protocol {version}")));
}
};
Ok(ScannerPeerActivity {
instance_id,
instance_id: response.instance_id,
namespace_generation: response.namespace_generation,
maintenance_generation: response.maintenance_generation,
protocol_version: response.protocol_version,
topology_digest,
data_movement_active,
dirty_usage_generation,
dirty_usage_pending,
})
}
fn decode_scanner_activity(response: ScannerActivityResponse, challenge: &[u8; 16]) -> Result<ScannerPeerActivity> {
decode_scanner_activity_with_verifier(response, challenge, |canonical, proof| {
verify_tonic_rpc_response_proof(canonical, proof)
.map_err(|_| Error::other("peer returned an invalid scanner activity response proof"))
})
}
@@ -1331,6 +1424,7 @@ impl PeerRestClient {
}
return Err(Error::other(""));
}
validate_signal_service_protocol(sig, sub_sys, response.protocol_version)?;
Ok(())
}
.await,
@@ -1338,25 +1432,44 @@ impl PeerRestClient {
.await
}
pub async fn scanner_activity(&self) -> Result<ScannerPeerActivity> {
async fn scanner_activity_request(
&self,
acknowledge_instance_id: String,
acknowledge_dirty_usage_generation: u64,
) -> Result<ScannerPeerActivity> {
self.finalize_result(
async {
let challenge = Uuid::new_v4();
let mut client = self
.get_client()
.await?
.max_decoding_message_size(SCANNER_ACTIVITY_MAX_MESSAGE_SIZE)
.max_encoding_message_size(SCANNER_ACTIVITY_MAX_MESSAGE_SIZE);
let response = client
.scanner_activity(Request::new(ScannerActivityRequest {}))
.await?
.into_inner();
decode_scanner_activity(response)
let mut request = Request::new(ScannerActivityRequest {
challenge: challenge.as_bytes().to_vec().into(),
protocol_version: SCANNER_ACTIVITY_PROTOCOL_VERSION,
acknowledge_instance_id,
acknowledge_dirty_usage_generation,
});
let canonical = rustfs_protos::canonical_scanner_activity_request_body(request.get_ref())
.map_err(|_| Error::other("scanner activity request is too large to authenticate"))?;
set_tonic_canonical_body_digest(&mut request, &canonical)?;
let response = client.scanner_activity(request).await?.into_inner();
decode_scanner_activity(response, challenge.as_bytes())
}
.await,
)
.await
}
pub async fn scanner_activity(&self) -> Result<ScannerPeerActivity> {
self.scanner_activity_request(String::new(), 0).await
}
pub async fn acknowledge_scanner_dirty_usage(&self, instance_id: String, generation: u64) -> Result<ScannerPeerActivity> {
self.scanner_activity_request(instance_id, generation).await
}
pub async fn get_metacache_listing(&self) -> Result<()> {
warn!("get_metacache_listing is not implemented in PeerRestClient");
Err(Error::NotImplemented)
@@ -1542,6 +1655,7 @@ impl PeerRestClient {
#[cfg(test)]
mod tests {
use super::*;
use crate::config::com::STORAGE_CLASS_SUB_SYS;
use serde_json::Value;
use std::io::{self, Write};
use std::sync::{Arc, Mutex};
@@ -1650,6 +1764,14 @@ mod tests {
)
}
fn decode_test_scanner_activity(response: ScannerActivityResponse) -> Result<ScannerPeerActivity> {
decode_scanner_activity_with_verifier(response, &[9; 16], |_canonical, proof| {
(proof == b"proof")
.then_some(())
.ok_or_else(|| Error::other("peer returned an invalid scanner activity response proof"))
})
}
#[test]
fn build_clients_from_slots_preserves_missing_remote_topology_slots() {
let slots = vec![
@@ -1682,13 +1804,89 @@ mod tests {
#[test]
fn scanner_activity_requires_restart_safe_peer_identity() {
let legacy = decode_test_scanner_activity(ScannerActivityResponse {
instance_id: "0123456789abcdef0123456789abcdef".to_string(),
namespace_generation: 7,
maintenance_generation: 3,
protocol_version: SCANNER_ACTIVITY_LEGACY_PROTOCOL_VERSION,
topology_digest: Vec::new().into(),
data_movement_active: false,
response_proof: Vec::new().into(),
dirty_usage_generation: 0,
dirty_usage_pending: false,
})
.expect("legacy peers should retain their activity generations during a rolling upgrade");
assert_eq!(
legacy,
ScannerPeerActivity {
instance_id: "0123456789abcdef0123456789abcdef".to_string(),
namespace_generation: 7,
maintenance_generation: 3,
protocol_version: SCANNER_ACTIVITY_LEGACY_PROTOCOL_VERSION,
topology_digest: None,
data_movement_active: None,
dirty_usage_generation: None,
dirty_usage_pending: None,
}
);
let previous = decode_test_scanner_activity(ScannerActivityResponse {
instance_id: "0123456789abcdef0123456789abcdef".to_string(),
namespace_generation: 7,
maintenance_generation: 3,
protocol_version: SCANNER_ACTIVITY_PREVIOUS_PROTOCOL_VERSION,
topology_digest: vec![7; 32].into(),
data_movement_active: true,
response_proof: b"proof".to_vec().into(),
dirty_usage_generation: 0,
dirty_usage_pending: false,
})
.expect("protocol v4 peers should remain observable during a rolling upgrade");
assert_eq!(
previous,
ScannerPeerActivity {
instance_id: "0123456789abcdef0123456789abcdef".to_string(),
namespace_generation: 7,
maintenance_generation: 3,
protocol_version: SCANNER_ACTIVITY_PREVIOUS_PROTOCOL_VERSION,
topology_digest: Some([7; 32]),
data_movement_active: Some(true),
dirty_usage_generation: None,
dirty_usage_pending: None,
}
);
let malformed_topology = ScannerActivityResponse {
instance_id: "0123456789abcdef0123456789abcdef".to_string(),
namespace_generation: 7,
maintenance_generation: 3,
protocol_version: SCANNER_ACTIVITY_PROTOCOL_VERSION,
topology_digest: vec![7; 31].into(),
data_movement_active: false,
response_proof: b"proof".to_vec().into(),
dirty_usage_generation: 11,
dirty_usage_pending: true,
};
assert!(
decode_test_scanner_activity(malformed_topology)
.expect_err("activity topology digests must have the protocol-defined length")
.to_string()
.contains("topology digest")
);
let missing_instance = ScannerActivityResponse {
instance_id: String::new(),
namespace_generation: 7,
maintenance_generation: 3,
protocol_version: SCANNER_ACTIVITY_PROTOCOL_VERSION,
topology_digest: vec![7; 32].into(),
data_movement_active: false,
response_proof: b"proof".to_vec().into(),
dirty_usage_generation: 11,
dirty_usage_pending: true,
};
assert!(
decode_scanner_activity(missing_instance)
decode_test_scanner_activity(missing_instance)
.expect_err("an empty instance ID is not restart safe")
.to_string()
.contains("instance ID")
@@ -1698,18 +1896,30 @@ mod tests {
instance_id: "ABCDEF0123456789ABCDEF0123456789".to_string(),
namespace_generation: 7,
maintenance_generation: 3,
protocol_version: SCANNER_ACTIVITY_PROTOCOL_VERSION,
topology_digest: vec![7; 32].into(),
data_movement_active: false,
response_proof: b"proof".to_vec().into(),
dirty_usage_generation: 11,
dirty_usage_pending: true,
};
assert!(
decode_scanner_activity(malformed_instance)
decode_test_scanner_activity(malformed_instance)
.expect_err("activity instance IDs must use the canonical lowercase hex form")
.to_string()
.contains("instance ID")
);
let activity = decode_scanner_activity(ScannerActivityResponse {
let activity = decode_test_scanner_activity(ScannerActivityResponse {
instance_id: "0123456789abcdef0123456789abcdef".to_string(),
namespace_generation: 7,
maintenance_generation: 3,
protocol_version: SCANNER_ACTIVITY_PROTOCOL_VERSION,
topology_digest: vec![7; 32].into(),
data_movement_active: true,
response_proof: b"proof".to_vec().into(),
dirty_usage_generation: 11,
dirty_usage_pending: true,
})
.expect("complete activity responses should be accepted");
assert_eq!(
@@ -1718,8 +1928,123 @@ mod tests {
instance_id: "0123456789abcdef0123456789abcdef".to_string(),
namespace_generation: 7,
maintenance_generation: 3,
protocol_version: SCANNER_ACTIVITY_PROTOCOL_VERSION,
topology_digest: Some([7; 32]),
data_movement_active: Some(true),
dirty_usage_generation: Some(11),
dirty_usage_pending: Some(true),
}
);
let pending_without_generation = ScannerActivityResponse {
instance_id: "0123456789abcdef0123456789abcdef".to_string(),
namespace_generation: 7,
maintenance_generation: 3,
protocol_version: SCANNER_ACTIVITY_PROTOCOL_VERSION,
topology_digest: vec![7; 32].into(),
data_movement_active: false,
response_proof: b"proof".to_vec().into(),
dirty_usage_generation: 0,
dirty_usage_pending: true,
};
assert!(
decode_test_scanner_activity(pending_without_generation)
.expect_err("pending dirty usage must carry a nonzero generation")
.to_string()
.contains("without a generation")
);
let previous_with_dirty_usage = ScannerActivityResponse {
instance_id: "0123456789abcdef0123456789abcdef".to_string(),
namespace_generation: 7,
maintenance_generation: 3,
protocol_version: SCANNER_ACTIVITY_PREVIOUS_PROTOCOL_VERSION,
topology_digest: vec![7; 32].into(),
data_movement_active: false,
response_proof: b"proof".to_vec().into(),
dirty_usage_generation: 11,
dirty_usage_pending: true,
};
assert!(
decode_test_scanner_activity(previous_with_dirty_usage)
.expect_err("protocol v4 responses must not claim unauthenticated dirty usage fields")
.to_string()
.contains("unauthenticated v5 fields")
);
let legacy_with_topology = ScannerActivityResponse {
instance_id: "0123456789abcdef0123456789abcdef".to_string(),
namespace_generation: 7,
maintenance_generation: 3,
protocol_version: SCANNER_ACTIVITY_LEGACY_PROTOCOL_VERSION,
topology_digest: vec![7; 32].into(),
data_movement_active: false,
response_proof: b"proof".to_vec().into(),
dirty_usage_generation: 0,
dirty_usage_pending: false,
};
assert!(
decode_test_scanner_activity(legacy_with_topology)
.expect_err("legacy protocol responses must not claim extended fields")
.to_string()
.contains("unexpected extended fields")
);
let unsupported_protocol = ScannerActivityResponse {
instance_id: "0123456789abcdef0123456789abcdef".to_string(),
namespace_generation: 7,
maintenance_generation: 3,
protocol_version: SCANNER_ACTIVITY_PROTOCOL_VERSION + 1,
topology_digest: vec![7; 32].into(),
data_movement_active: false,
response_proof: b"proof".to_vec().into(),
dirty_usage_generation: 11,
dirty_usage_pending: true,
};
assert!(
decode_test_scanner_activity(unsupported_protocol)
.expect_err("unknown activity protocols must fail closed")
.to_string()
.contains("unsupported scanner activity protocol")
);
let missing_proof = ScannerActivityResponse {
instance_id: "0123456789abcdef0123456789abcdef".to_string(),
namespace_generation: 7,
maintenance_generation: 3,
protocol_version: SCANNER_ACTIVITY_PROTOCOL_VERSION,
topology_digest: vec![7; 32].into(),
data_movement_active: false,
response_proof: Vec::new().into(),
dirty_usage_generation: 11,
dirty_usage_pending: true,
};
assert!(
decode_test_scanner_activity(missing_proof)
.expect_err("unsigned scanner activity responses must fail closed")
.to_string()
.contains("response proof")
);
}
#[test]
fn dynamic_scanner_config_requires_versioned_peer_acknowledgement() {
for sub_system in [SCANNER_SUB_SYS, HEAL_SUB_SYS] {
let err = validate_signal_service_protocol(SERVICE_SIGNAL_RELOAD_DYNAMIC, sub_system, 0)
.expect_err("an unversioned peer must not claim scanner config convergence");
assert!(err.to_string().contains("does not support dynamic"));
validate_signal_service_protocol(
SERVICE_SIGNAL_RELOAD_DYNAMIC,
sub_system,
rustfs_protos::DYNAMIC_CONFIG_PROTOCOL_VERSION,
)
.expect("a current peer should support dynamic scanner config");
}
validate_signal_service_protocol(SERVICE_SIGNAL_RELOAD_DYNAMIC, STORAGE_CLASS_SUB_SYS, 0)
.expect("unrelated dynamic config keeps its existing compatibility contract");
validate_signal_service_protocol(SERVICE_SIGNAL_REFRESH_CONFIG, SCANNER_SUB_SYS, 0)
.expect("full refresh compatibility is guarded by its scanner preflight");
}
#[test]
@@ -49,6 +49,20 @@ use tracing::{debug, info, warn};
type Client = Arc<Box<dyn PeerS3Client>>;
#[derive(Clone, Debug)]
pub struct ScannerBucketListing {
pub buckets: Vec<BucketInfo>,
pub set_buckets: Vec<ScannerSetBucketListing>,
pub topology_complete: bool,
}
#[derive(Clone, Debug)]
pub struct ScannerSetBucketListing {
pub pool_index: usize,
pub set_index: usize,
pub buckets: Vec<BucketInfo>,
}
fn pool_participant_errors(clients: &[Client], errors: &[Option<Error>], pool_idx: usize) -> Vec<Option<Error>> {
clients
.iter()
@@ -216,6 +230,10 @@ impl S3PeerSys {
Ok(())
}
pub async fn list_bucket(&self, opts: &BucketOptions) -> Result<Vec<BucketInfo>> {
Ok(self.list_bucket_for_scanner(opts).await?.buckets)
}
pub async fn list_bucket_for_scanner(&self, opts: &BucketOptions) -> Result<ScannerBucketListing> {
let mut futures = Vec::with_capacity(self.clients.len());
for cli in self.clients.iter() {
futures.push(cli.list_bucket(opts));
@@ -239,9 +257,12 @@ impl S3PeerSys {
}
let mut result_map: HashMap<&String, BucketInfo> = HashMap::new();
let mut topology_complete = true;
for i in 0..self.pools_count {
let per_pool_errs = pool_participant_errors(&self.clients, &errors, i);
let quorum = pool_write_quorum(per_pool_errs.len());
topology_complete &=
!per_pool_errs.is_empty() && per_pool_errs.iter().all(|participant_error| participant_error.is_none());
if let Some(pool_err) = reduce_pool_write_quorum_errs(&per_pool_errs) {
tracing::error!("list_bucket per_pool_errs: {per_pool_errs:?}");
@@ -261,20 +282,17 @@ impl S3PeerSys {
}
for bucket in buckets.iter() {
if result_map.contains_key(&bucket.name) {
continue;
}
// incr bucket_map count create if not exists
let count = bucket_map.entry(&bucket.name).or_insert(0usize);
*count += 1;
if *count >= quorum {
result_map.insert(&bucket.name, bucket.clone());
result_map.entry(&bucket.name).or_insert_with(|| bucket.clone());
}
}
}
}
topology_complete &= bucket_map.values().all(|count| *count >= quorum);
// TODO: MRF
}
@@ -282,7 +300,11 @@ impl S3PeerSys {
buckets.sort_by_key(|b| b.name.clone());
Ok(buckets)
Ok(ScannerBucketListing {
buckets,
set_buckets: Vec::new(),
topology_complete,
})
}
pub async fn delete_bucket(&self, bucket: &str, opts: &DeleteBucketOptions) -> Result<()> {
let mut futures = Vec::with_capacity(self.clients.len());
@@ -1645,6 +1667,86 @@ mod tests {
assert_eq!(buckets[0].name, bucket.name);
}
#[tokio::test]
async fn scanner_bucket_listing_marks_quorum_result_incomplete_when_a_peer_is_missing() {
let bucket = BucketInfo {
name: "bucket-hidden-by-quorum".to_string(),
..Default::default()
};
let peer_sys = S3PeerSys {
clients: vec![
test_peer_with_list_bucket(&[0], Ok(vec![bucket])),
test_peer_with_list_bucket(&[0], Ok(Vec::new())),
test_peer_with_list_bucket(&[0], Ok(Vec::new())),
test_peer_with_list_bucket(&[0], Err(Error::DiskAccessDenied)),
],
pools_count: 1,
};
let listing = peer_sys
.list_bucket_for_scanner(&BucketOptions::default())
.await
.expect("peer quorum should still produce a scanner candidate listing");
assert!(listing.buckets.is_empty());
assert!(!listing.topology_complete);
}
#[tokio::test]
async fn scanner_bucket_listing_marks_divergent_successful_peers_incomplete() {
let bucket = BucketInfo {
name: "bucket-below-quorum".to_string(),
..Default::default()
};
let peer_sys = S3PeerSys {
clients: vec![
test_peer_with_list_bucket(&[0], Ok(vec![bucket.clone()])),
test_peer_with_list_bucket(&[0], Ok(vec![bucket])),
test_peer_with_list_bucket(&[0], Ok(Vec::new())),
test_peer_with_list_bucket(&[0], Ok(Vec::new())),
],
pools_count: 1,
};
let listing = peer_sys
.list_bucket_for_scanner(&BucketOptions::default())
.await
.expect("successful peer responses should still produce a scanner candidate listing");
assert!(listing.buckets.is_empty());
assert!(!listing.topology_complete);
}
#[tokio::test]
async fn scanner_bucket_listing_checks_same_bucket_in_every_pool() {
let bucket = BucketInfo {
name: "shared-bucket".to_string(),
..Default::default()
};
let peer_sys = S3PeerSys {
clients: vec![
test_peer_with_list_bucket(&[0], Ok(vec![bucket.clone()])),
test_peer_with_list_bucket(&[0], Ok(vec![bucket.clone()])),
test_peer_with_list_bucket(&[0], Ok(vec![bucket.clone()])),
test_peer_with_list_bucket(&[0], Ok(vec![bucket.clone()])),
test_peer_with_list_bucket(&[1], Ok(vec![bucket.clone()])),
test_peer_with_list_bucket(&[1], Ok(vec![bucket.clone()])),
test_peer_with_list_bucket(&[1], Ok(Vec::new())),
test_peer_with_list_bucket(&[1], Ok(Vec::new())),
],
pools_count: 2,
};
let listing = peer_sys
.list_bucket_for_scanner(&BucketOptions::default())
.await
.expect("a bucket visible in one pool should remain a scan candidate");
assert_eq!(listing.buckets.len(), 1);
assert_eq!(listing.buckets[0].name, bucket.name);
assert!(!listing.topology_complete);
}
#[tokio::test]
async fn test_delete_bucket_fails_when_any_pool_misses_write_quorum() {
let peer_sys = S3PeerSys {
+279 -6
View File
@@ -17,7 +17,8 @@ use crate::cluster::rpc::client::{
node_service_time_out_client_for_class, node_service_time_out_client_no_auth,
};
use crate::cluster::rpc::internode_data_transport::{
InternodeDataTransport, ReadStreamRequest, WalkDirStreamRequest, WriteStreamRequest,
InternodeDataTransport, NsScannerCapabilityRequest, NsScannerStreamRequest, ReadStreamRequest, WalkDirStreamRequest,
WriteStreamRequest,
};
use crate::disk::error::{Error, Result};
use crate::disk::{
@@ -81,6 +82,7 @@ const REMOTE_DISK_OPEN_WRITE_MAX_ATTEMPTS: usize = 2;
const REMOTE_DISK_OPEN_WRITE_RETRY_BACKOFF: Duration = Duration::from_millis(20);
const REMOTE_DISK_OPEN_READ_MAX_ATTEMPTS: usize = 2;
const REMOTE_DISK_OPEN_READ_RETRY_BACKOFF: Duration = Duration::from_millis(20);
const NS_SCANNER_CAPABILITY_PROBE_TIMEOUT: Duration = Duration::from_secs(5);
/// Base backoff for idempotent read-only RPC retries (grpc-optimization P3-3); doubles per attempt.
const REMOTE_DISK_READ_RETRY_BASE_BACKOFF: Duration = Duration::from_millis(50);
const ENV_RUSTFS_METADATA_BATCH_READ: &str = "RUSTFS_METADATA_BATCH_READ";
@@ -97,6 +99,17 @@ const LOG_SUBSYSTEM_REMOTE_DISK: &str = "remote_disk";
const EVENT_REMOTE_DISK_HEALTH: &str = "remote_disk_health";
const EVENT_REMOTE_DISK_RPC: &str = "remote_disk_rpc";
fn decode_volume_infos(volume_infos: Vec<String>) -> Result<Vec<VolumeInfo>> {
volume_infos
.into_iter()
.enumerate()
.map(|(index, json)| {
serde_json::from_str::<VolumeInfo>(&json)
.map_err(|err| Error::other(format!("decode list volumes entry {index} failed: {err}")))
})
.collect()
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
enum BatchMetadataRpcMode {
Off,
@@ -264,6 +277,71 @@ fn spawn_control_channel_prewarm(addr: String) {
}
impl RemoteDisk {
pub(crate) async fn ns_scanner_server_epoch(&self) -> Result<Option<Uuid>> {
if self.health.is_faulty() {
return Err(DiskError::FaultyDisk);
}
let probe = self.data_transport.probe_ns_scanner(NsScannerCapabilityRequest {
endpoint: self.endpoint.grid_host(),
});
let result = timeout(NS_SCANNER_CAPABILITY_PROBE_TIMEOUT, probe)
.await
.map_err(|_| DiskError::other("remote namespace scanner capability probe timed out"))?;
match result {
Ok(server_epoch) => Ok(Some(server_epoch)),
// RUSTFS_COMPAT_TODO(ns-scanner-rpc-v3): old peers and legacy transports lack the authenticated startup-epoch handshake. Remove after every supported peer implements namespace scanner protocol v3.
Err(DiskError::MethodNotAllowed) => Ok(None),
Err(err)
if matches!(
err.internode_http_error_kind(),
Some(rustfs_rio::InternodeHttpErrorKind::HttpStatus(status))
if matches!(status.as_u16(), 404 | 405)
) =>
{
Ok(None)
}
Err(err)
if matches!(
err.internode_http_error_kind(),
Some(rustfs_rio::InternodeHttpErrorKind::HttpStatus(status)) if status.as_u16() == 426
) =>
{
Ok(None)
}
Err(err) => Err(err),
}
}
pub(crate) async fn open_ns_scanner_stream(&self, request: crate::disk::NsScannerOpenRequest) -> Result<FileReader> {
if self.health.is_faulty() {
return Err(DiskError::FaultyDisk);
}
let crate::disk::NsScannerOpenRequest {
request_id,
server_epoch,
session_id,
session_sequence,
next_cycle,
leader_epoch,
body,
stall_timeout,
} = request;
self.data_transport
.open_ns_scanner(NsScannerStreamRequest {
endpoint: self.endpoint.grid_host(),
disk: self.disk_ref().await,
request_id,
server_epoch,
session_id,
session_sequence,
next_cycle,
leader_epoch,
body,
stall_timeout,
})
.await
}
fn recovery_monitor_span(addr: &str, endpoint: &Endpoint) -> tracing::Span {
tracing::info_span!(
"recovery-monitor",
@@ -1335,11 +1413,7 @@ impl DiskAPI for RemoteDisk {
return Err(response.error.unwrap_or_default().into());
}
let infos = response
.volume_infos
.into_iter()
.filter_map(|json_str| serde_json::from_str::<VolumeInfo>(&json_str).ok())
.collect();
let infos = decode_volume_infos(response.volume_infos)?;
Ok(infos)
},
@@ -2722,6 +2796,19 @@ mod tests {
static INIT: Once = Once::new();
#[test]
fn list_volumes_decode_rejects_a_malformed_entry() {
let valid = serde_json::to_string(&VolumeInfo {
name: "bucket".to_string(),
created: None,
})
.expect("volume info should serialize");
let err = decode_volume_infos(vec![valid, "{".to_string()])
.expect_err("a malformed volume entry must fail the complete response");
assert!(err.to_string().contains("entry 1"));
}
#[test]
fn decoded_remote_metadata_rejects_default_like_delete_marker() {
let forged = FileInfo {
@@ -2794,14 +2881,31 @@ mod tests {
Read(ReadStreamRequest),
Write(WriteStreamRequest),
WalkDir(WalkDirStreamRequest),
NsScanner(NsScannerStreamRequest),
NsScannerProbe(NsScannerCapabilityRequest),
}
#[derive(Debug, Clone, Default)]
struct RecordingInternodeDataTransport {
calls: Arc<StdMutex<Vec<RecordedTransportCall>>>,
ns_scanner_probe_status: Arc<StdMutex<Option<u16>>>,
}
impl RecordingInternodeDataTransport {
fn with_ns_scanner_probe_status(status: u16) -> Self {
Self {
calls: Arc::default(),
ns_scanner_probe_status: Arc::new(StdMutex::new(Some(status))),
}
}
fn set_ns_scanner_probe_status(&self, status: Option<u16>) {
*self
.ns_scanner_probe_status
.lock()
.expect("namespace scanner probe status lock poisoned") = status;
}
fn calls(&self) -> Vec<RecordedTransportCall> {
self.calls.lock().expect("recorded transport calls lock poisoned").clone()
}
@@ -3369,6 +3473,26 @@ mod tests {
Ok(Box::new(EmptyTestReader))
}
async fn open_ns_scanner(&self, request: NsScannerStreamRequest) -> Result<FileReader> {
self.record(RecordedTransportCall::NsScanner(request));
Ok(Box::new(EmptyTestReader))
}
async fn probe_ns_scanner(&self, request: NsScannerCapabilityRequest) -> Result<Uuid> {
self.record(RecordedTransportCall::NsScannerProbe(request));
if let Some(status) = *self
.ns_scanner_probe_status
.lock()
.expect("namespace scanner probe status lock poisoned")
{
let status = reqwest::StatusCode::from_u16(status).expect("test status code should be valid");
return Err(
rustfs_rio::new_test_internode_http_io_error(rustfs_rio::InternodeHttpErrorKind::HttpStatus(status)).into(),
);
}
Ok(Uuid::from_u128(1))
}
fn name(&self) -> &'static str {
"recording"
}
@@ -3401,6 +3525,14 @@ mod tests {
}
}
async fn open_ns_scanner(&self, _request: NsScannerStreamRequest) -> Result<FileReader> {
panic!("open_ns_scanner should not be used in walk_dir retry test");
}
async fn probe_ns_scanner(&self, _request: NsScannerCapabilityRequest) -> Result<Uuid> {
Ok(Uuid::from_u128(1))
}
fn name(&self) -> &'static str {
"retrying-walk-dir"
}
@@ -3429,6 +3561,14 @@ mod tests {
panic!("open_walk_dir should not be used in open_write retry test");
}
async fn open_ns_scanner(&self, _request: NsScannerStreamRequest) -> Result<FileReader> {
panic!("open_ns_scanner should not be used in open_write retry test");
}
async fn probe_ns_scanner(&self, _request: NsScannerCapabilityRequest) -> Result<Uuid> {
Ok(Uuid::from_u128(1))
}
fn name(&self) -> &'static str {
"retrying-open-write"
}
@@ -4092,6 +4232,139 @@ mod tests {
}
}
#[tokio::test]
async fn test_remote_disk_namespace_scanner_uses_configured_data_transport() {
let transport = RecordingInternodeDataTransport::default();
let remote_disk = new_remote_disk_with_transport(Arc::new(transport.clone())).await;
let expected_disk = remote_disk.disk_ref().await;
let expected_body = b"namespace-scanner-request".to_vec();
let expected_request_id = Uuid::new_v4();
let expected_server_epoch = Uuid::new_v4();
let expected_session_id = Uuid::new_v4();
let _reader = remote_disk
.open_ns_scanner_stream(crate::disk::NsScannerOpenRequest {
request_id: expected_request_id,
server_epoch: expected_server_epoch,
session_id: expected_session_id,
session_sequence: 3,
next_cycle: 7,
leader_epoch: 9,
body: expected_body.clone(),
stall_timeout: Some(Duration::from_secs(15)),
})
.await
.expect("namespace scanner stream should open");
let calls = transport.calls();
assert_eq!(calls.len(), 1);
match &calls[0] {
RecordedTransportCall::NsScanner(request) => {
assert_eq!(request.endpoint, "http://remote-node:9000");
assert_eq!(request.disk, expected_disk);
assert_eq!(request.request_id, expected_request_id);
assert_eq!(request.server_epoch, expected_server_epoch);
assert_eq!(request.session_id, expected_session_id);
assert_eq!(request.session_sequence, 3);
assert_eq!(request.next_cycle, 7);
assert_eq!(request.leader_epoch, 9);
assert_eq!(request.body, expected_body);
assert_eq!(request.stall_timeout, Some(Duration::from_secs(15)));
}
other => panic!("expected namespace scanner transport call, got {other:?}"),
}
}
#[tokio::test]
async fn test_remote_disk_namespace_scanner_capability_uses_configured_data_transport() {
let transport = RecordingInternodeDataTransport::default();
let remote_disk = new_remote_disk_with_transport(Arc::new(transport.clone())).await;
assert_eq!(
remote_disk
.ns_scanner_server_epoch()
.await
.expect("namespace scanner capability probe should succeed"),
Some(Uuid::from_u128(1))
);
let calls = transport.calls();
assert_eq!(calls.len(), 1);
match &calls[0] {
RecordedTransportCall::NsScannerProbe(request) => {
assert_eq!(request.endpoint, "http://remote-node:9000");
}
other => panic!("expected namespace scanner capability probe, got {other:?}"),
}
}
#[tokio::test]
async fn test_remote_disk_namespace_scanner_capability_rejects_old_and_incompatible_peers() {
for status in [404, 405, 426] {
let transport = RecordingInternodeDataTransport::with_ns_scanner_probe_status(status);
let remote_disk = new_remote_disk_with_transport(Arc::new(transport)).await;
assert_eq!(
remote_disk
.ns_scanner_server_epoch()
.await
.expect("unsupported namespace scanner response should be classified"),
None
);
}
}
#[tokio::test]
async fn test_remote_disk_namespace_scanner_capability_rejects_legacy_transport() {
let remote_disk = new_remote_disk_with_transport(Arc::new(RetryingOpenReadInternodeDataTransport::default())).await;
assert_eq!(
remote_disk
.ns_scanner_server_epoch()
.await
.expect("legacy transport should be classified as unsupported"),
None
);
}
#[tokio::test]
async fn test_remote_disk_namespace_scanner_capability_reprobes_after_peer_upgrade() {
let transport = RecordingInternodeDataTransport::with_ns_scanner_probe_status(404);
let remote_disk = new_remote_disk_with_transport(Arc::new(transport.clone())).await;
assert_eq!(
remote_disk
.ns_scanner_server_epoch()
.await
.expect("old peer should be classified as unsupported"),
None
);
transport.set_ns_scanner_probe_status(None);
assert_eq!(
remote_disk
.ns_scanner_server_epoch()
.await
.expect("upgraded peer should be re-probed"),
Some(Uuid::from_u128(1))
);
assert_eq!(transport.calls().len(), 2);
}
#[tokio::test]
async fn test_remote_disk_namespace_scanner_capability_propagates_transient_failure() {
let transport = RecordingInternodeDataTransport::with_ns_scanner_probe_status(503);
let remote_disk = new_remote_disk_with_transport(Arc::new(transport)).await;
let err = remote_disk
.ns_scanner_server_epoch()
.await
.expect_err("transient capability failure must not be reported as unsupported");
assert!(matches!(
err.internode_http_error_kind(),
Some(rustfs_rio::InternodeHttpErrorKind::HttpStatus(status)) if status.as_u16() == 503
));
}
#[tokio::test]
async fn test_remote_disk_walk_dir_preserves_skip_total_timeout_option() {
let transport = RecordingInternodeDataTransport::default();
File diff suppressed because it is too large Load Diff
+16 -11
View File
@@ -560,17 +560,22 @@ pub(crate) async fn cleanup_source_entry_if_unchanged(
ensure_source_cleanup_versions_unchanged(set.clone(), bucket, object, expected, allowed_missing, op_label).await?;
set.delete_object(
bucket,
cleanup_key.as_str(),
ObjectOptions {
delete_prefix: true,
delete_prefix_object: true,
no_lock: true,
..Default::default()
},
)
.await
let result = set
.delete_object(
bucket,
cleanup_key.as_str(),
ObjectOptions {
delete_prefix: true,
delete_prefix_object: true,
no_lock: true,
..Default::default()
},
)
.await;
if result.is_ok() {
crate::store::list_objects::observe_scanner_namespace_mutations(bucket, 1);
}
result
}
fn should_check_data_movement_resume_target(src_pool_idx: usize, target_pool_idx: usize) -> bool {
File diff suppressed because it is too large Load Diff
+18
View File
@@ -231,6 +231,13 @@ impl DiskError {
}
}
pub fn is_internode_http_status(&self, status: u16) -> bool {
matches!(
self.internode_http_error_kind(),
Some(InternodeHttpErrorKind::HttpStatus(actual)) if actual.as_u16() == status
)
}
// /// If all errors are of the same fatal disk error type, returns the corresponding error.
// /// Otherwise, returns Ok.
// pub fn check_disk_fatal_errs(errs: &[Option<Error>]) -> Result<()> {
@@ -981,6 +988,17 @@ mod tests {
assert!(!disk_error.contains_io_error_kind(std::io::ErrorKind::TimedOut));
}
#[test]
fn test_internode_http_status_classification() {
let too_many_requests = DiskError::from(rustfs_rio::new_test_internode_http_io_error(
rustfs_rio::InternodeHttpErrorKind::HttpStatus(http::StatusCode::TOO_MANY_REQUESTS),
));
assert!(too_many_requests.is_internode_http_status(429));
assert!(!too_many_requests.is_internode_http_status(500));
assert!(!DiskError::FileNotFound.is_internode_http_status(429));
}
#[test]
fn test_metacache_output_stream_closed_classification_survives_clone() {
let disk_error = DiskError::metacache_output_stream_closed();
+27 -1
View File
@@ -52,7 +52,7 @@ use local::LocalDisk;
use rustfs_filemeta::{FileInfo, ObjectPartInfo, RawFileInfo};
use rustfs_madmin::info_commands::DiskMetrics;
use serde::{Deserialize, Serialize};
use std::{fmt::Debug, path::PathBuf, sync::Arc};
use std::{fmt::Debug, path::PathBuf, sync::Arc, time::Duration};
use time::OffsetDateTime;
use tokio::io::{AsyncRead, AsyncWrite};
use uuid::Uuid;
@@ -453,6 +453,20 @@ impl DiskAPI for Disk {
}
impl Disk {
pub async fn ns_scanner_server_epoch(&self) -> Result<Option<Uuid>> {
match self {
Disk::Local(_) => Ok(None),
Disk::Remote(remote_disk) => remote_disk.ns_scanner_server_epoch().await,
}
}
pub async fn open_ns_scanner_stream(&self, request: NsScannerOpenRequest) -> Result<FileReader> {
match self {
Disk::Remote(remote_disk) => remote_disk.open_ns_scanner_stream(request).await,
Disk::Local(_) => Err(Error::other("namespace scanner stream requires a remote disk")),
}
}
pub fn runtime_state(&self) -> RuntimeDriveHealthState {
match self {
Disk::Local(local_disk) => local_disk.runtime_state(),
@@ -498,6 +512,18 @@ impl Disk {
}
}
#[derive(Debug)]
pub struct NsScannerOpenRequest {
pub request_id: Uuid,
pub server_epoch: Uuid,
pub session_id: Uuid,
pub session_sequence: u64,
pub next_cycle: u64,
pub leader_epoch: u64,
pub body: Vec<u8>,
pub stall_timeout: Option<Duration>,
}
impl Disk {
/// Reset drive health so `connect_load_init_formats` retries are not blocked by a prior
/// transient mark-faulty (same disk handles are reused across retries).
+8
View File
@@ -471,6 +471,14 @@ impl PutObjReader {
}
}
pub fn from_prehashed_bytes(data: Bytes, sha256hex: Option<String>) -> std::io::Result<Self> {
let content_length =
i64::try_from(data.len()).map_err(|_| std::io::Error::other("prehashed object payload exceeds i64 length"))?;
Ok(PutObjReader {
stream: HashReader::from_stream(Cursor::new(data), content_length, content_length, None, sha256hex, false)?,
})
}
pub fn size(&self) -> i64 {
self.stream.size()
}
@@ -154,6 +154,7 @@ fn to_madmin_scanner_metrics(metrics: rustfs_common::metrics::ScannerMetricsRepo
last_cycle_replication_checks: metrics.last_cycle_replication_checks,
last_cycle_usage_saves: metrics.last_cycle_usage_saves,
failed_cycles: metrics.failed_cycles,
superseded_cycles: metrics.superseded_cycles,
partial_cycles_unknown: metrics.partial_cycles_unknown,
partial_cycles_runtime: metrics.partial_cycles_runtime,
partial_cycles_objects: metrics.partial_cycles_objects,
@@ -796,6 +797,7 @@ mod test {
last_cycle_replication_checks: 27,
last_cycle_usage_saves: 28,
failed_cycles: 29,
superseded_cycles: 30,
partial_cycles_unknown: 30,
partial_cycles_runtime: 31,
partial_cycles_objects: 32,
@@ -927,6 +929,7 @@ mod test {
assert_eq!(scanner.last_cycle_replication_checks, 27);
assert_eq!(scanner.last_cycle_usage_saves, 28);
assert_eq!(scanner.failed_cycles, 29);
assert_eq!(scanner.superseded_cycles, 30);
assert_eq!(scanner.partial_cycles_unknown, 30);
assert_eq!(scanner.partial_cycles_runtime, 31);
assert_eq!(scanner.partial_cycles_objects, 32);
+124 -1
View File
@@ -29,7 +29,7 @@ use rustfs_madmin::metrics::RealtimeMetrics;
use rustfs_madmin::net::NetInfo;
use rustfs_madmin::{ItemState, ServerProperties, StorageInfo};
use rustfs_utils::XHost;
use std::collections::hash_map::DefaultHasher;
use std::collections::{HashMap, hash_map::DefaultHasher};
use std::future::Future;
use std::hash::{Hash, Hasher};
use std::sync::{Arc, Mutex, OnceLock};
@@ -1041,6 +1041,40 @@ impl NotificationSys {
Ok(generations)
}
pub async fn acknowledge_scanner_dirty_usage(&self, acknowledgements: Vec<(String, String, u64)>) -> Result<bool> {
let mut by_host = HashMap::with_capacity(acknowledgements.len());
for (host, instance_id, generation) in acknowledgements {
if by_host.insert(host.clone(), (instance_id, generation)).is_some() {
return Err(Error::other(format!("duplicate scanner dirty usage acknowledgement target: {host}")));
}
}
let clients = self
.peer_clients
.iter()
.flatten()
.map(|client| (client.grid_host.clone(), client.clone()))
.collect::<HashMap<_, _>>();
let mut failures = Vec::new();
let mut futures = Vec::with_capacity(by_host.len());
for (host, (instance_id, generation)) in by_host {
let Some(client) = clients.get(&host).cloned() else {
failures.push(format!("peer {host} scanner dirty usage acknowledgement failed: peer is not reachable"));
continue;
};
futures.push(async move {
let result = scanner_activity_with_timeout(
SCANNER_ACTIVITY_PROBE_TIMEOUT,
&host,
client.acknowledge_scanner_dirty_usage(instance_id, generation),
)
.await;
(host, result)
});
}
aggregate_scanner_dirty_usage_acknowledgement_results(join_all(futures).await, failures)
}
pub async fn reload_site_replication_config(&self) -> Vec<NotificationPeerErr> {
let mut futures = Vec::with_capacity(self.peer_clients.len());
for client in self.peer_clients.iter() {
@@ -1591,6 +1625,23 @@ fn aggregate_notification_failures(operation: &str, failures: Vec<String>) -> Re
)))
}
fn aggregate_scanner_dirty_usage_acknowledgement_results(
results: Vec<(String, Result<ScannerPeerActivity>)>,
mut failures: Vec<String>,
) -> Result<bool> {
let mut dirty_usage_pending = false;
for (host, result) in results {
match result {
Ok(activity) => {
dirty_usage_pending |= activity.dirty_usage_pending != Some(false);
}
Err(err) => failures.push(format!("peer {host} scanner dirty usage acknowledgement failed: {err}")),
}
}
aggregate_notification_failures("acknowledge_scanner_dirty_usage", failures)?;
Ok(dirty_usage_pending)
}
#[cfg(test)]
mod tests {
use super::*;
@@ -1842,6 +1893,78 @@ mod tests {
assert!(err.to_string().contains("peer-1"));
}
#[tokio::test]
async fn scanner_dirty_usage_acknowledgement_rejects_missing_and_duplicate_targets() {
let sys = NotificationSys {
peer_clients: Vec::new(),
all_peer_clients: Vec::new(),
peer_admin_caches: Vec::new(),
peer_topology_hosts: Vec::new(),
};
let missing = sys
.acknowledge_scanner_dirty_usage(vec![("peer-1".to_string(), "0123456789abcdef0123456789abcdef".to_string(), 7)])
.await
.expect_err("a missing acknowledgement target must remain pending");
assert!(missing.to_string().contains("peer is not reachable"));
let duplicate = sys
.acknowledge_scanner_dirty_usage(vec![
("peer-1".to_string(), "0123456789abcdef0123456789abcdef".to_string(), 7),
("peer-1".to_string(), "0123456789abcdef0123456789abcdef".to_string(), 7),
])
.await
.expect_err("duplicate acknowledgement targets must be rejected");
assert!(
duplicate
.to_string()
.contains("duplicate scanner dirty usage acknowledgement target")
);
}
#[test]
fn scanner_dirty_usage_acknowledgement_preserves_newer_pending_work() {
let activity = |dirty_usage_pending| ScannerPeerActivity {
instance_id: "0123456789abcdef0123456789abcdef".to_string(),
namespace_generation: 1,
maintenance_generation: 1,
protocol_version: crate::storage_api_contracts::internode::SCANNER_ACTIVITY_PROTOCOL_VERSION,
topology_digest: Some([0; 32]),
data_movement_active: Some(false),
dirty_usage_generation: Some(2),
dirty_usage_pending,
};
let pending = aggregate_scanner_dirty_usage_acknowledgement_results(
vec![
("peer-1".to_string(), Ok(activity(Some(false)))),
("peer-2".to_string(), Ok(activity(Some(true)))),
],
Vec::new(),
)
.expect("successful acknowledgements should return their pending state");
assert!(pending, "new dirty usage reported by an acknowledged peer must remain pending");
let cleared = aggregate_scanner_dirty_usage_acknowledgement_results(
vec![("peer-1".to_string(), Ok(activity(Some(false))))],
Vec::new(),
)
.expect("a cleared acknowledgement should succeed");
assert!(!cleared, "an explicitly cleared peer must not remain pending");
let unknown =
aggregate_scanner_dirty_usage_acknowledgement_results(vec![("peer-1".to_string(), Ok(activity(None)))], Vec::new())
.expect("an acknowledgement without a pending field should remain retryable");
assert!(unknown, "a peer that cannot prove its dirty state is clear must remain pending");
let err = aggregate_scanner_dirty_usage_acknowledgement_results(
vec![("peer-1".to_string(), Err(Error::other("injected acknowledgement failure")))],
Vec::new(),
)
.expect_err("a reachable peer acknowledgement failure must be reported");
assert!(err.to_string().contains("peer-1"));
assert!(err.to_string().contains("injected acknowledgement failure"));
}
#[tokio::test]
async fn load_bucket_metadata_reports_unreachable_peers() {
let sys = NotificationSys {
@@ -3862,6 +3862,24 @@ impl SetDisks {
Ok(removed)
}
fn write_precondition_lookup_error(
error: StorageError,
http_preconditions: &HTTPPreconditions,
bucket: &str,
object: &str,
) -> Option<StorageError> {
match error {
StorageError::VersionNotFound(_, _, _) | StorageError::ObjectNotFound(_, _) => {
if http_preconditions.if_match_value().is_some() {
Some(StorageError::ObjectNotFound(bucket.to_string(), object.to_string()))
} else {
None
}
}
error => Some(error),
}
}
pub(in crate::set_disk) async fn check_write_precondition(
&self,
bucket: &str,
@@ -3892,19 +3910,8 @@ impl SetDisks {
}
}
Err(StorageError::VersionNotFound(_, _, _))
| Err(StorageError::ObjectNotFound(_, _))
| Err(StorageError::ErasureReadQuorum) => {
// When the object is not found,
// - if If-Match is set, we should return 404 NotFound
// - if If-None-Match is set, we should be able to proceed with the request
if http_preconditions.if_match_value().is_some() {
return Some(StorageError::ObjectNotFound(bucket.to_string(), object.to_string()));
}
}
Err(e) => {
return Some(e);
Err(error) => {
return Self::write_precondition_lookup_error(error, &http_preconditions, bucket, object);
}
}
@@ -4407,6 +4414,41 @@ mod tests {
use tempfile::TempDir;
use tokio::io::AsyncReadExt;
#[test]
fn write_precondition_lookup_errors_fail_closed_unless_absence_is_known() {
let create_only = HTTPPreconditions {
if_none_match: Some("*".to_string()),
..Default::default()
};
let replace_only = HTTPPreconditions {
if_match: Some("etag".to_string()),
..Default::default()
};
assert!(matches!(
SetDisks::write_precondition_lookup_error(StorageError::ErasureReadQuorum, &create_only, "bucket", "object",),
Some(StorageError::ErasureReadQuorum)
));
assert!(
SetDisks::write_precondition_lookup_error(
StorageError::ObjectNotFound("bucket".to_string(), "object".to_string()),
&create_only,
"bucket",
"object",
)
.is_none()
);
assert!(matches!(
SetDisks::write_precondition_lookup_error(
StorageError::ObjectNotFound("bucket".to_string(), "object".to_string()),
&replace_only,
"bucket",
"object",
),
Some(StorageError::ObjectNotFound(_, _))
));
}
fn metadata_test_fileinfo(object: &str) -> FileInfo {
let mut fi = FileInfo::new(object, 2, 2);
fi.volume = "bucket".to_string();
+42
View File
@@ -9275,6 +9275,48 @@ mod tests {
));
}
#[tokio::test]
async fn set_level_if_none_match_fails_closed_without_read_quorum() {
let set_disks = make_local_bucket_test_set_disks_with_drive_count(4).await;
let bucket = "bucket-write-precondition-quorum";
let object = "existing-object.txt";
set_disks
.make_bucket(bucket, &MakeBucketOptions::default())
.await
.expect("bucket should be created before disk loss");
let mut reader = PutObjReader::from_vec(b"existing object body".to_vec());
set_disks
.put_object(
bucket,
object,
&mut reader,
&ObjectOptions {
no_lock: true,
..Default::default()
},
)
.await
.expect("object should be written before disk loss");
{
let mut disks = set_disks.disks.write().await;
disks[1..].fill(None);
}
let create_only = ObjectOptions {
http_preconditions: Some(HTTPPreconditions {
if_none_match: Some("*".to_string()),
..Default::default()
}),
..Default::default()
};
let result = set_disks.check_write_precondition(bucket, object, &create_only).await;
assert!(
matches!(result, Some(StorageError::ErasureReadQuorum | StorageError::InsufficientReadQuorum(_, _))),
"expected read-quorum failure, got {result:?}"
);
}
#[tokio::test]
async fn set_level_versioned_delete_marker_hides_object_without_corrupting_version_metadata() {
let set_disks = make_local_bucket_test_set_disks_with_drive_count(4).await;
+67 -59
View File
@@ -21,6 +21,71 @@
use super::super::*;
impl SetDisks {
pub(crate) async fn list_bucket_for_scanner(&self, _opts: &BucketOptions) -> Result<(Vec<BucketInfo>, bool)> {
let disks = self.disk_inventory().await;
let write_quorum = (disks.len() / 2) + 1;
let mut futures = Vec::with_capacity(disks.len());
for disk in disks {
futures.push(async move {
match disk {
Some(disk) => disk.list_volumes().await,
None => Err(DiskError::DiskNotFound),
}
});
}
let results = join_all(futures).await;
let mut topology_complete = results.iter().all(|result| result.is_ok());
let mut infos = Vec::with_capacity(results.len());
let mut errs = Vec::with_capacity(results.len());
for result in results {
match result {
Ok(volumes) => {
infos.push(Some(volumes));
errs.push(None);
}
Err(err) => {
infos.push(None);
errs.push(Some(err));
}
}
}
if let Some(err) = reduce_write_quorum_errs(&errs, BUCKET_OP_IGNORED_ERRS, write_quorum) {
return Err(err.into());
}
let mut counts: HashMap<String, (usize, BucketInfo)> = HashMap::new();
for volumes in infos.into_iter().flatten() {
for volume in volumes {
if is_reserved_or_invalid_bucket(&volume.name, false) {
continue;
}
let entry = counts.entry(volume.name.clone()).or_insert((
0,
BucketInfo {
name: volume.name.clone(),
created: volume.created,
..Default::default()
},
));
entry.0 += 1;
}
}
topology_complete &= counts.values().all(|(count, _)| *count >= write_quorum);
let mut buckets = counts
.into_values()
.filter_map(|(count, bucket)| (count >= write_quorum).then_some(bucket))
.collect::<Vec<_>>();
buckets.sort_by(|left, right| left.name.cmp(&right.name));
Ok((buckets, topology_complete))
}
}
#[async_trait::async_trait]
impl BucketOperations for SetDisks {
type Error = Error;
@@ -117,65 +182,8 @@ impl BucketOperations for SetDisks {
}
#[tracing::instrument(skip(self))]
async fn list_bucket(&self, _opts: &BucketOptions) -> Result<Vec<BucketInfo>> {
let disks = self.disk_inventory().await;
let write_quorum = (disks.len() / 2) + 1;
let mut futures = Vec::with_capacity(disks.len());
for disk in disks {
futures.push(async move {
match disk {
Some(disk) => disk.list_volumes().await,
None => Err(DiskError::DiskNotFound),
}
});
}
let results = join_all(futures).await;
let mut infos = Vec::with_capacity(results.len());
let mut errs = Vec::with_capacity(results.len());
for result in results {
match result {
Ok(volumes) => {
infos.push(Some(volumes));
errs.push(None);
}
Err(err) => {
infos.push(None);
errs.push(Some(err));
}
}
}
if let Some(err) = reduce_write_quorum_errs(&errs, BUCKET_OP_IGNORED_ERRS, write_quorum) {
return Err(err.into());
}
let mut counts: HashMap<String, (usize, BucketInfo)> = HashMap::new();
for volumes in infos.into_iter().flatten() {
for volume in volumes {
if is_reserved_or_invalid_bucket(&volume.name, false) {
continue;
}
let entry = counts.entry(volume.name.clone()).or_insert((
0,
BucketInfo {
name: volume.name.clone(),
created: volume.created,
..Default::default()
},
));
entry.0 += 1;
}
}
let mut buckets = counts
.into_values()
.filter_map(|(count, bucket)| (count >= write_quorum).then_some(bucket))
.collect::<Vec<_>>();
buckets.sort_by(|left, right| left.name.cmp(&right.name));
Ok(buckets)
async fn list_bucket(&self, opts: &BucketOptions) -> Result<Vec<BucketInfo>> {
Ok(self.list_bucket_for_scanner(opts).await?.0)
}
#[tracing::instrument(skip(self))]
@@ -23,7 +23,12 @@ pub(crate) mod heal {
pub(crate) mod internode {
pub(crate) use rustfs_storage_api::{
WALK_DIR_BODY_SHA256_QUERY, WALK_DIR_STREAM_COMPLETION_QUERY, WALK_DIR_STREAM_COMPLETION_V1,
NS_SCANNER_BODY_SHA256_QUERY, NS_SCANNER_CAPABILITY_CHALLENGE_QUERY, NS_SCANNER_CYCLE_QUERY,
NS_SCANNER_LEADER_EPOCH_QUERY, NS_SCANNER_PROTOCOL_VERSION, NS_SCANNER_PROTOCOL_VERSION_QUERY,
NS_SCANNER_REQUEST_ID_QUERY, NS_SCANNER_SERVER_EPOCH_QUERY, NS_SCANNER_SESSION_ID_QUERY,
NS_SCANNER_SESSION_SEQUENCE_QUERY, NsScannerCapabilityResponse, SCANNER_ACTIVITY_LEGACY_PROTOCOL_VERSION,
SCANNER_ACTIVITY_PREVIOUS_PROTOCOL_VERSION, SCANNER_ACTIVITY_PROTOCOL_VERSION, WALK_DIR_BODY_SHA256_QUERY,
WALK_DIR_STREAM_COMPLETION_QUERY, WALK_DIR_STREAM_COMPLETION_V1,
};
}
+746 -30
View File
@@ -17,12 +17,21 @@ use crate::bucket::{
metadata::{BUCKET_TABLE_RESERVED_PREFIX, table_bucket_catalog_metadata_prefix},
utils::is_meta_bucketname,
};
use crate::error::is_err_bucket_not_found;
use crate::runtime::sources as runtime_sources;
use crate::set_disk::get_lock_acquire_timeout;
use crate::storage_api_contracts::bucket::SRBucketDeleteOp;
use crate::storage_api_contracts::namespace::NamespaceLocking as _;
use futures::stream::{self, StreamExt};
use std::collections::BTreeMap;
use std::future::Future;
const DELETED_BUCKETS_PREFIX: &str = ".deleted";
const SCANNER_BUCKET_LIST_SET_CONCURRENCY: usize = 4;
fn scanner_bucket_list_set_concurrency(set_count: usize) -> usize {
set_count.clamp(1, SCANNER_BUCKET_LIST_SET_CONCURRENCY)
}
fn should_override_created_from_metadata(created: OffsetDateTime) -> bool {
created != OffsetDateTime::UNIX_EPOCH
@@ -78,6 +87,49 @@ fn bucket_deleted_marker_volume(bucket: &str) -> String {
format!("{RUSTFS_META_BUCKET}/{}", bucket_deleted_marker_prefix(bucket))
}
async fn await_bucket_namespace_operation<T, F>(
guard: Option<&rustfs_lock::NamespaceLockGuard>,
bucket: &str,
operation: &'static str,
future: F,
) -> Result<T>
where
F: Future<Output = Result<T>>,
{
let Some(guard) = guard else {
return future.await;
};
if guard.is_lock_lost() {
return Err(StorageError::other(format!(
"bucket namespace lock was lost before {operation}: {bucket}"
)));
}
tokio::select! {
biased;
_ = guard.lock_lost_notified() => Err(StorageError::other(format!(
"bucket namespace lock was lost during {operation}: {bucket}"
))),
result = future => result,
}
}
async fn run_bucket_usage_cleanup<F>(guard: Option<&rustfs_lock::NamespaceLockGuard>, bucket: &str, future: F) -> Result<()>
where
F: Future<Output = Result<()>>,
{
await_bucket_namespace_operation(guard, bucket, "bucket usage cleanup", future).await
}
async fn run_physical_bucket_deletion<F>(guard: Option<&rustfs_lock::NamespaceLockGuard>, bucket: &str, future: F) -> Result<()>
where
F: Future<Output = Result<()>>,
{
// Fence before polling deletion: the physical namespace may become
// invisible at any await point inside the storage operation.
list_objects::observe_scanner_namespace_mutations(bucket, 1);
await_bucket_namespace_operation(guard, bucket, "physical bucket deletion", future).await
}
impl ECStore {
async fn mark_bucket_deleted(&self, bucket: &str) -> Result<()> {
let marker_volume = bucket_deleted_marker_volume(bucket);
@@ -96,21 +148,84 @@ impl ECStore {
Ok(())
}
async fn cleanup_deleted_bucket_metadata(&self, bucket: &str, include_deleted_marker: bool) -> Result<()> {
async fn cleanup_deleted_bucket_metadata(
&self,
bucket: &str,
include_deleted_marker: bool,
guard: Option<&rustfs_lock::NamespaceLockGuard>,
) -> Result<()> {
for prefix in bucket_delete_metadata_cleanup_prefixes(bucket) {
self.delete_all(RUSTFS_META_BUCKET, prefix.as_str()).await?;
await_bucket_namespace_operation(
guard,
bucket,
"deleted bucket metadata cleanup",
self.delete_all(RUSTFS_META_BUCKET, prefix.as_str()),
)
.await?;
}
if include_deleted_marker {
let marker_prefix = bucket_deleted_marker_prefix(bucket);
self.delete_all(RUSTFS_META_BUCKET, marker_prefix.as_str()).await?;
await_bucket_namespace_operation(
guard,
bucket,
"deleted bucket marker cleanup",
self.delete_all(RUSTFS_META_BUCKET, marker_prefix.as_str()),
)
.await?;
}
metadata_sys::remove_bucket_metadata_in(&self.ctx, bucket).await?;
await_bucket_namespace_operation(
guard,
bucket,
"deleted bucket metadata cache cleanup",
metadata_sys::remove_bucket_metadata_in(&self.ctx, bucket),
)
.await?;
runtime_sources::delete_bucket_monitor_entry(bucket);
Ok(())
}
async fn cleanup_bucket_usage(&self, bucket: &str, guard: Option<&rustfs_lock::NamespaceLockGuard>) -> Result<()> {
run_bucket_usage_cleanup(guard, bucket, async {
crate::data_usage::prepare_bucket_usage_for_namespace_change(bucket, guard).await?;
crate::data_usage::remove_bucket_usage_from_backend_with_guard(self, bucket, guard).await
})
.await
}
async fn cleanup_bucket_usage_best_effort(&self, bucket: &str, guard: Option<&rustfs_lock::NamespaceLockGuard>) {
if let Err(err) = self.cleanup_bucket_usage(bucket, guard).await {
warn!(
bucket = %bucket,
error = ?err,
"bucket data usage cleanup deferred to scanner reconciliation"
);
}
}
async fn rollback_failed_bucket_creation(&self, bucket: &str, guard: Option<&rustfs_lock::NamespaceLockGuard>) {
let rollback_opts = DeleteBucketOptions {
no_lock: true,
no_recreate: true,
..Default::default()
};
if let Err(err) = await_bucket_namespace_operation(guard, bucket, "failed bucket creation rollback", async {
self.peer_sys
.delete_bucket(bucket, &rollback_opts)
.await
.map_err(|rollback_err| to_object_err(rollback_err.into(), vec![bucket]))
})
.await
{
warn!(
bucket = %bucket,
error = ?err,
"failed bucket creation rollback did not remove every physical bucket volume"
);
}
}
#[instrument(skip(self))]
pub(super) async fn handle_make_bucket(&self, bucket: &str, opts: &MakeBucketOptions) -> Result<()> {
if !is_meta_bucketname(bucket)
@@ -119,7 +234,7 @@ impl ECStore {
return Err(StorageError::BucketNameInvalid(err.to_string()));
}
let _ns_guard = if !opts.no_lock {
let ns_guard = if !opts.no_lock {
let ns_lock = self.new_ns_lock(bucket, bucket).await?;
Some(
ns_lock
@@ -142,8 +257,34 @@ impl ECStore {
None
};
if let Err(err) = self.peer_sys.make_bucket(bucket, opts).await {
let err = to_object_err(err.into(), vec![bucket]);
let confirmed_missing = match self.peer_sys.get_bucket_info(bucket, &BucketOptions::default()).await {
Ok(_) => false,
Err(err) => {
let err: StorageError = err.into();
if is_err_bucket_not_found(&err) {
true
} else {
return Err(to_object_err(err, vec![bucket]));
}
}
};
if confirmed_missing && !is_meta_bucketname(bucket) {
// Fence every scanner cycle that could have observed the namespace
// before physical creation. Creation may become visible even when a
// later metadata write or namespace-lock check fails.
crate::store::list_objects::observe_scanner_namespace_mutations(bucket, 1);
self.cleanup_bucket_usage(bucket, ns_guard.as_ref()).await?;
}
if let Err(err) = await_bucket_namespace_operation(ns_guard.as_ref(), bucket, "physical bucket creation", async {
self.peer_sys
.make_bucket(bucket, opts)
.await
.map_err(|err| to_object_err(err.into(), vec![bucket]))
})
.await
{
if is_err_bucket_exists(&err)
&& let Err(heal_err) = self
.handle_heal_bucket(
@@ -157,18 +298,9 @@ impl ECStore {
{
warn!("best-effort bucket heal after BucketExists failed: {heal_err}");
}
if !is_err_bucket_exists(&err) {
if !is_err_bucket_exists(&err) && ns_guard.as_ref().is_none_or(|guard| !guard.is_lock_lost()) {
error!("make bucket failed: {err}");
let _ = self
.delete_bucket(
bucket,
&DeleteBucketOptions {
no_lock: true,
no_recreate: true,
..Default::default()
},
)
.await;
self.rollback_failed_bucket_creation(bucket, ns_guard.as_ref()).await;
}
return Err(err);
};
@@ -186,7 +318,13 @@ impl ECStore {
meta.versioning_config_xml = crate::bucket::utils::serialize::<VersioningConfiguration>(&enableVersioningConfig)?;
}
metadata_sys::set_bucket_metadata_in(&self.ctx, meta).await?;
await_bucket_namespace_operation(
ns_guard.as_ref(),
bucket,
"bucket metadata initialization",
metadata_sys::set_bucket_metadata_in(&self.ctx, meta),
)
.await?;
Ok(())
}
@@ -224,6 +362,80 @@ impl ECStore {
Ok(buckets)
}
pub async fn list_bucket_for_scanner(&self, opts: &BucketOptions) -> Result<crate::cluster::rpc::ScannerBucketListing> {
let sets = self
.pools
.iter()
.flat_map(|pool| {
pool.disk_set
.iter()
.map(|set| (set.pool_index, set.set_index, Arc::clone(set)))
})
.collect::<Vec<_>>();
let set_count = sets.len();
let deleted = opts.deleted;
let cached = opts.cached;
let no_metadata = opts.no_metadata;
let mut set_listings = stream::iter(sets.into_iter().map(move |(pool_index, set_index, set)| {
let opts = BucketOptions {
deleted,
cached,
no_metadata,
};
async move {
set.list_bucket_for_scanner(&opts)
.await
.map(|(buckets, complete)| (pool_index, set_index, buckets, complete))
}
}))
.buffer_unordered(scanner_bucket_list_set_concurrency(set_count));
let mut topology_complete = set_count != 0;
let mut bucket_map = BTreeMap::<String, BucketInfo>::new();
let mut scoped_buckets = Vec::with_capacity(set_count);
while let Some(set_listing) = set_listings.next().await {
let (pool_index, set_index, buckets, set_complete) = set_listing?;
topology_complete &= set_complete;
for bucket in &buckets {
bucket_map.entry(bucket.name.clone()).or_insert_with(|| bucket.clone());
}
scoped_buckets.push(crate::cluster::rpc::ScannerSetBucketListing {
pool_index,
set_index,
buckets,
});
}
scoped_buckets.sort_unstable_by_key(|scope| (scope.pool_index, scope.set_index));
let mut listing = crate::cluster::rpc::ScannerBucketListing {
buckets: bucket_map.into_values().collect(),
set_buckets: scoped_buckets,
topology_complete,
};
if !opts.no_metadata {
for bucket in &mut listing.buckets {
if let Ok(created) = metadata_sys::created_at_in(&self.ctx, &bucket.name).await
&& should_override_created_from_metadata(created)
{
bucket.created = Some(created);
}
}
let created_by_bucket = listing
.buckets
.iter()
.map(|bucket| (bucket.name.as_str(), bucket.created))
.collect::<BTreeMap<_, _>>();
for scope in &mut listing.set_buckets {
for bucket in &mut scope.buckets {
if let Some(created) = created_by_bucket.get(bucket.name.as_str()) {
bucket.created = *created;
}
}
}
}
Ok(listing)
}
#[instrument(skip(self))]
pub(super) async fn handle_delete_bucket(&self, bucket: &str, opts: &DeleteBucketOptions) -> Result<()> {
if is_meta_bucketname(bucket) {
@@ -234,7 +446,7 @@ impl ECStore {
return Err(StorageError::BucketNameInvalid(err.to_string()));
}
let _ns_guard = if !opts.no_lock {
let ns_guard = if !opts.no_lock {
let ns_lock = self.new_ns_lock(bucket, bucket).await?;
Some(
ns_lock
@@ -299,17 +511,40 @@ impl ECStore {
}
if sr_mark_delete {
self.mark_bucket_deleted(bucket).await?;
await_bucket_namespace_operation(
ns_guard.as_ref(),
bucket,
"bucket delete marker creation",
self.mark_bucket_deleted(bucket),
)
.await?;
}
if let Err(err) = self.peer_sys.delete_bucket(bucket, &delete_opts).await {
let storage_err = to_object_err(err.into(), vec![bucket]);
if !sr_delete || !is_err_strict_volume_not_found(&storage_err) {
return Err(storage_err);
}
let delete_result = run_physical_bucket_deletion(ns_guard.as_ref(), bucket, async {
self.peer_sys
.delete_bucket(bucket, &delete_opts)
.await
.map_err(|err| to_object_err(err.into(), vec![bucket]))
})
.await;
if let Err(err) = delete_result
&& (!sr_delete || !is_err_strict_volume_not_found(&err))
{
return Err(err);
}
self.cleanup_deleted_bucket_metadata(bucket, sr_purge).await?;
self.cleanup_bucket_usage_best_effort(bucket, ns_guard.as_ref()).await;
if let Err(err) = self
.cleanup_deleted_bucket_metadata(bucket, sr_purge, ns_guard.as_ref())
.await
{
warn!(
bucket = %bucket,
error = ?err,
"physical bucket deletion succeeded but metadata cleanup remains pending"
);
}
Ok(())
}
}
@@ -317,8 +552,9 @@ impl ECStore {
#[cfg(test)]
mod tests {
use super::{
bucket_delete_metadata_cleanup_prefixes, bucket_deleted_marker_prefix, bucket_deleted_marker_volume,
should_override_created_from_metadata, validate_table_bucket_delete_allowed,
SCANNER_BUCKET_LIST_SET_CONCURRENCY, await_bucket_namespace_operation, bucket_delete_metadata_cleanup_prefixes,
bucket_deleted_marker_prefix, bucket_deleted_marker_volume, run_bucket_usage_cleanup, run_physical_bucket_deletion,
scanner_bucket_list_set_concurrency, should_override_created_from_metadata, validate_table_bucket_delete_allowed,
};
use crate::bucket::metadata::table_bucket_catalog_metadata_prefix;
use crate::bucket::metadata_sys;
@@ -335,16 +571,146 @@ mod tests {
disk::endpoint::Endpoint,
layout::endpoints::{EndpointServerPools, Endpoints, PoolEndpoints},
};
use rustfs_data_usage::{BucketUsageInfo, DATA_USAGE_OBJECT_NAME, DataUsageInfo};
use rustfs_lock::{LocalClient, LockRequest, LockType, NamespaceLock, ObjectKey};
use serial_test::serial;
use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::time::{Duration, SystemTime};
use time::OffsetDateTime;
use tokio::sync::OnceCell;
use tokio::sync::{Notify, OnceCell};
use tokio_util::sync::CancellationToken;
use uuid::Uuid;
static BUCKET_DELETE_TEST_ENV: OnceCell<(Vec<PathBuf>, Arc<ECStore>)> = OnceCell::const_new();
#[tokio::test(start_paused = true)]
async fn bucket_namespace_operation_fails_closed_after_lease_expiry() {
let ttl = Duration::from_millis(20);
let lock = NamespaceLock::new("bucket-operation-test".to_string(), Arc::new(LocalClient::new()));
let request = LockRequest::new(ObjectKey::new("bucket", ""), LockType::Exclusive, "test-owner")
.with_acquire_timeout(Duration::from_secs(1))
.with_ttl(ttl)
.with_refresh_interval(ttl);
let guard = lock
.acquire_guard(&request)
.await
.expect("namespace lock acquisition should not fail")
.expect("namespace lock should be acquired");
tokio::time::advance(ttl + Duration::from_millis(1)).await;
let operation_ran = Arc::new(AtomicBool::new(false));
let operation_ran_for_future = operation_ran.clone();
let result = await_bucket_namespace_operation(Some(&guard), "bucket", "test operation", async move {
operation_ran_for_future.store(true, Ordering::SeqCst);
Ok(())
})
.await;
assert!(result.is_err(), "an expired namespace lease must fence the operation");
assert!(
!operation_ran.load(Ordering::SeqCst),
"a fenced namespace operation must not poll its mutation future"
);
}
#[tokio::test]
#[serial]
async fn physical_bucket_delete_fences_scanner_before_polling_storage() {
let generation_before = crate::store::list_objects::scanner_namespace_mutation_generation();
let storage_polled = Arc::new(AtomicBool::new(false));
let storage_polled_for_future = storage_polled.clone();
run_physical_bucket_deletion(None, "generation-order-bucket", async move {
assert!(
crate::store::list_objects::scanner_namespace_mutation_generation() > generation_before,
"scanner generation must advance before physical deletion is polled"
);
storage_polled_for_future.store(true, Ordering::SeqCst);
Ok(())
})
.await
.expect("synthetic physical deletion should succeed");
assert!(storage_polled.load(Ordering::SeqCst));
}
#[tokio::test]
async fn bucket_usage_cleanup_stops_after_parent_cancellation() {
let started = Arc::new(Notify::new());
let started_wait = started.notified();
let release = Arc::new(Notify::new());
let completed = Arc::new(AtomicBool::new(false));
let started_for_cleanup = started.clone();
let release_for_cleanup = release.clone();
let completed_for_cleanup = completed.clone();
let parent = tokio::spawn(run_bucket_usage_cleanup(None, "bucket", async move {
started_for_cleanup.notify_one();
release_for_cleanup.notified().await;
completed_for_cleanup.store(true, Ordering::SeqCst);
Ok(())
}));
started_wait.await;
parent.abort();
let _ = parent.await;
tokio::task::yield_now().await;
release.notify_waiters();
tokio::task::yield_now().await;
assert!(
!completed.load(Ordering::SeqCst),
"a cancelled cleanup future must not continue in a detached task"
);
}
#[tokio::test(start_paused = true)]
#[serial]
async fn bucket_namespace_operation_stops_in_flight_work_after_lock_loss() {
let ttl = Duration::from_millis(20);
let lock = NamespaceLock::new("bucket-operation-in-flight-loss-test".to_string(), Arc::new(LocalClient::new()));
let request = LockRequest::new(ObjectKey::new("bucket", ""), LockType::Exclusive, "test-owner")
.with_acquire_timeout(Duration::from_secs(1))
.with_ttl(ttl)
.with_refresh_interval(ttl);
let guard = lock
.acquire_guard(&request)
.await
.expect("namespace lock acquisition should not fail")
.expect("namespace lock should be acquired");
let operation_started = Arc::new(Notify::new());
let started_wait = operation_started.notified();
let operation_release = Arc::new(Notify::new());
let operation_completed = Arc::new(AtomicBool::new(false));
let started_for_operation = operation_started.clone();
let release_for_operation = operation_release.clone();
let completed_for_operation = operation_completed.clone();
let task = tokio::spawn(async move {
await_bucket_namespace_operation(Some(&guard), "bucket", "test operation", async move {
started_for_operation.notify_one();
release_for_operation.notified().await;
completed_for_operation.store(true, Ordering::SeqCst);
Ok(())
})
.await
});
started_wait.await;
tokio::time::advance(ttl + Duration::from_millis(1)).await;
let err = task
.await
.expect("operation task should join")
.expect_err("an operation still running after lease loss must be fenced");
assert!(err.to_string().contains("namespace lock was lost during test operation"));
operation_release.notify_waiters();
tokio::task::yield_now().await;
assert!(
!operation_completed.load(Ordering::SeqCst),
"lock loss must stop the old owner before a successor can acquire the namespace"
);
}
async fn setup_bucket_delete_test_env() -> (Vec<PathBuf>, Arc<ECStore>) {
BUCKET_DELETE_TEST_ENV
.get_or_init(|| async {
@@ -413,6 +779,58 @@ mod tests {
.clone()
}
async fn setup_multi_pool_scanner_listing_test_env() -> (tempfile::TempDir, Arc<ECStore>) {
let temp_dir = tempfile::tempdir().expect("multi-pool scanner test directory should be created");
let mut pools = 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 scanner test disk should be created");
let mut endpoint =
Endpoint::try_from(disk_path.to_str().expect("disk path should be utf8")).expect("endpoint should parse");
endpoint.set_pool_index(pool_index);
endpoint.set_set_index(0);
endpoint.set_disk_index(disk_index);
endpoints.push(endpoint);
}
pools.push(PoolEndpoints {
legacy: false,
set_count: 1,
drives_per_set: 4,
endpoints: Endpoints::from(endpoints),
cmd_line: format!("scanner-listing-pool-{pool_index}"),
platform: format!("OS: {} | Arch: {}", std::env::consts::OS, std::env::consts::ARCH),
});
}
let endpoint_pools = EndpointServerPools(pools);
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 ecstore = ECStore::new_with_instance_ctx(
"127.0.0.1:0".parse().expect("test address"),
endpoint_pools,
CancellationToken::new(),
instance_ctx,
)
.await
.expect("multi-pool ECStore should initialize");
let storage_class =
crate::config::storageclass::lookup_config_for_pools_without_env(&rustfs_config::server_config::KVS::new(), &[4, 4])
.expect("multi-pool storage class should match both four-disk pools");
for pool in &ecstore.pools {
for set in &pool.disk_set {
set.set_test_storage_class_config(storage_class.clone());
}
}
(temp_dir, ecstore)
}
async fn create_bucket_with_object(ecstore: &Arc<ECStore>, bucket: &str, object: &str) {
let generation_before_make = ecstore.scanner_namespace_mutation_generation();
ecstore
@@ -515,6 +933,98 @@ mod tests {
);
}
#[test]
fn scanner_bucket_listing_bounds_set_fanout() {
assert_eq!(scanner_bucket_list_set_concurrency(0), 1);
assert_eq!(scanner_bucket_list_set_concurrency(2), 2);
assert_eq!(scanner_bucket_list_set_concurrency(100), SCANNER_BUCKET_LIST_SET_CONCURRENCY);
}
#[tokio::test]
#[serial]
async fn scanner_bucket_listing_unions_every_erasure_set() {
let (_temp_dir, ecstore) = setup_multi_pool_scanner_listing_test_env().await;
let bucket = format!("second-pool-only-{}", Uuid::new_v4().simple());
ecstore.pools[1].disk_set[0]
.make_bucket(&bucket, &MakeBucketOptions::default())
.await
.expect("bucket should be created in the second pool only");
let listing = ecstore
.list_bucket_for_scanner(&crate::storage_api_contracts::bucket::BucketOptions {
no_metadata: true,
..Default::default()
})
.await
.expect("scanner should enumerate every pool and set");
assert!(listing.topology_complete);
assert!(listing.buckets.iter().any(|entry| entry.name == bucket));
assert_eq!(listing.set_buckets.len(), 2);
assert!(
listing
.set_buckets
.iter()
.find(|scope| scope.pool_index == 0 && scope.set_index == 0)
.is_some_and(|scope| scope.buckets.is_empty())
);
assert!(
listing
.set_buckets
.iter()
.find(|scope| scope.pool_index == 1 && scope.set_index == 0)
.is_some_and(|scope| scope.buckets.iter().any(|entry| entry.name == bucket))
);
}
#[tokio::test]
async fn scanner_bucket_listing_marks_degraded_set_incomplete() {
let (_temp_dir, ecstore) = setup_multi_pool_scanner_listing_test_env().await;
let bucket = format!("degraded-set-{}", Uuid::new_v4().simple());
let set = &ecstore.pools[0].disk_set[0];
set.make_bucket(&bucket, &MakeBucketOptions::default())
.await
.expect("bucket should be created before a disk is removed");
set.disks.write().await[0] = None;
let listing = ecstore
.list_bucket_for_scanner(&crate::storage_api_contracts::bucket::BucketOptions {
no_metadata: true,
..Default::default()
})
.await
.expect("a degraded set with quorum should still return candidate buckets");
assert!(listing.buckets.iter().any(|entry| entry.name == bucket));
assert!(!listing.topology_complete);
}
#[tokio::test]
async fn scanner_bucket_listing_marks_divergent_disk_views_incomplete() {
let (temp_dir, ecstore) = setup_multi_pool_scanner_listing_test_env().await;
let bucket = format!("divergent-set-{}", Uuid::new_v4().simple());
ecstore.pools[0].disk_set[0]
.make_bucket(&bucket, &MakeBucketOptions::default())
.await
.expect("bucket should be created before disk views diverge");
for disk_index in 0..2 {
tokio::fs::remove_dir_all(temp_dir.path().join(format!("pool0-disk{disk_index}")).join(&bucket))
.await
.expect("test bucket directory should be removed from a minority disk view");
}
let listing = ecstore
.list_bucket_for_scanner(&crate::storage_api_contracts::bucket::BucketOptions {
no_metadata: true,
..Default::default()
})
.await
.expect("responsive disks should still produce a scanner candidate listing");
assert!(listing.buckets.iter().all(|entry| entry.name != bucket));
assert!(!listing.topology_complete);
}
// These tests share one isolated instance and mutate its bucket metadata;
// serialize them so their assertions cannot observe each other's operations.
#[tokio::test]
@@ -736,4 +1246,210 @@ mod tests {
"failed default S3 DeleteBucket must keep metadata cache"
);
}
#[tokio::test]
#[serial]
async fn bucket_delete_finishes_usage_cleanup_before_same_name_recreation() {
let (_, ecstore) = setup_bucket_delete_test_env().await;
let bucket = format!("bucket-usage-generation-{}", Uuid::new_v4().simple());
ecstore
.make_bucket(&bucket, &MakeBucketOptions::default())
.await
.expect("bucket should be created");
let mut snapshot = DataUsageInfo {
last_update: Some(SystemTime::now()),
buckets_count: 1,
..Default::default()
};
snapshot.buckets_usage.insert(
bucket.clone(),
BucketUsageInfo {
size: 42,
objects_count: 1,
versions_count: 1,
..Default::default()
},
);
snapshot.bucket_sizes.insert(bucket.clone(), 42);
snapshot.calculate_totals();
crate::data_usage::store_data_usage_in_backend(snapshot, ecstore.clone())
.await
.expect("usage snapshot should be stored");
ecstore
.delete_bucket(&bucket, &DeleteBucketOptions::default())
.await
.expect("empty bucket should be deleted");
let deleted = crate::data_usage::load_data_usage_from_backend(ecstore.clone())
.await
.expect("usage snapshot should remain readable");
assert!(!deleted.buckets_usage.contains_key(&bucket));
ecstore
.make_bucket(&bucket, &MakeBucketOptions::default())
.await
.expect("same bucket name should be recreated after delete returns");
crate::data_usage::record_bucket_object_write_memory(&bucket, None, 84).await;
let mut recreated = crate::data_usage::load_data_usage_from_backend(ecstore.clone())
.await
.expect("recreated bucket usage base should load");
crate::data_usage::apply_bucket_usage_memory_overlay(&mut recreated).await;
assert_eq!(
recreated
.buckets_usage
.get(&bucket)
.map(|usage| (usage.objects_count, usage.versions_count, usage.size)),
Some((1, 1, 84))
);
}
#[tokio::test]
#[serial]
async fn bucket_create_removes_stale_usage_before_physical_creation() {
let (_, ecstore) = setup_bucket_delete_test_env().await;
let bucket = format!("bucket-create-stale-usage-{}", Uuid::new_v4().simple());
let mut snapshot = DataUsageInfo {
last_update: Some(SystemTime::now()),
buckets_count: 1,
..Default::default()
};
snapshot.buckets_usage.insert(
bucket.clone(),
BucketUsageInfo {
size: 42,
objects_count: 1,
versions_count: 1,
..Default::default()
},
);
snapshot.bucket_sizes.insert(bucket.clone(), 42);
snapshot.calculate_totals();
crate::data_usage::store_data_usage_in_backend(snapshot, ecstore.clone())
.await
.expect("stale usage fixture should be stored");
ecstore
.make_bucket(&bucket, &MakeBucketOptions::default())
.await
.expect("CreateBucket should succeed");
let persisted = crate::data_usage::load_data_usage_from_backend(ecstore.clone())
.await
.expect("usage snapshot should remain readable");
assert!(
!persisted.buckets_usage.contains_key(&bucket),
"a newly created bucket must not inherit the predecessor generation's usage"
);
assert!(
ecstore.get_bucket_info(&bucket, &BucketOptions::default()).await.is_ok(),
"physical creation should happen after the usage fence succeeds"
);
}
#[tokio::test]
#[serial]
async fn failed_create_rollback_does_not_run_unfenced_usage_cleanup() {
let (_, ecstore) = setup_bucket_delete_test_env().await;
let bucket = format!("bucket-create-rollback-{}", Uuid::new_v4().simple());
ecstore
.peer_sys
.make_bucket(&bucket, &MakeBucketOptions::default())
.await
.expect("the partial-create fixture should expose a physical bucket");
let mut snapshot = DataUsageInfo {
last_update: Some(SystemTime::now()),
buckets_count: 1,
..Default::default()
};
snapshot.buckets_usage.insert(
bucket.clone(),
BucketUsageInfo {
size: 42,
objects_count: 1,
versions_count: 1,
..Default::default()
},
);
snapshot.bucket_sizes.insert(bucket.clone(), 42);
snapshot.calculate_totals();
crate::data_usage::store_data_usage_in_backend(snapshot, ecstore.clone())
.await
.expect("the usage fixture should be stored");
ecstore.rollback_failed_bucket_creation(&bucket, None).await;
assert!(
ecstore
.peer_sys
.get_bucket_info(&bucket, &BucketOptions::default())
.await
.is_err(),
"failed-create rollback should remove the partial physical bucket"
);
let persisted = crate::data_usage::load_data_usage_from_backend(ecstore.clone())
.await
.expect("the usage snapshot should remain readable");
assert!(
persisted.buckets_usage.contains_key(&bucket),
"failed-create rollback must not start an unfenced usage cleanup"
);
crate::data_usage::store_data_usage_in_backend(DataUsageInfo::default(), ecstore.clone())
.await
.expect("the rollback usage fixture should be cleared");
}
#[tokio::test]
#[serial]
async fn bucket_create_fails_closed_when_usage_snapshot_cannot_be_fenced() {
let (_, ecstore) = setup_bucket_delete_test_env().await;
let deleted_bucket = format!("bucket-delete-corrupt-usage-{}", Uuid::new_v4().simple());
ecstore
.make_bucket(&deleted_bucket, &MakeBucketOptions::default())
.await
.expect("bucket should be created before corrupting usage");
let usage_path = format!("{BUCKET_META_PREFIX}/{DATA_USAGE_OBJECT_NAME}");
crate::config::com::save_config(ecstore.clone(), &usage_path, b"{".to_vec())
.await
.expect("corrupt usage fixture should be stored");
ecstore
.delete_bucket(&deleted_bucket, &DeleteBucketOptions::default())
.await
.expect("usage snapshot corruption must not block DeleteBucket");
assert!(
ecstore
.get_bucket_info(&deleted_bucket, &crate::storage_api_contracts::bucket::BucketOptions::default())
.await
.is_err(),
"successful DeleteBucket must remove the physical bucket"
);
let new_bucket = format!("bucket-create-corrupt-usage-{}", Uuid::new_v4().simple());
let create_err = ecstore
.make_bucket(&new_bucket, &MakeBucketOptions::default())
.await
.expect_err("MakeBucket must fail when stale usage cannot be fenced");
assert!(!create_err.to_string().is_empty(), "the usage snapshot failure should be surfaced");
assert!(
ecstore
.get_bucket_info(&new_bucket, &crate::storage_api_contracts::bucket::BucketOptions::default())
.await
.is_err(),
"failed MakeBucket must not expose a new physical bucket"
);
let restored = serde_json::to_vec(&DataUsageInfo {
last_update: Some(SystemTime::now()),
..Default::default()
})
.expect("restored usage fixture should encode");
crate::config::com::save_config(ecstore, &usage_path, restored)
.await
.expect("usage fixture should be restored after the failure-path test");
}
}
+1 -1
View File
@@ -619,7 +619,7 @@ pub(super) fn scanner_namespace_mutation_generation() -> u64 {
SCANNER_NAMESPACE_MUTATION_GENERATION.load(Ordering::Acquire)
}
pub(super) fn observe_scanner_namespace_mutations(bucket: &str, delta: u64) {
pub(crate) fn observe_scanner_namespace_mutations(bucket: &str, delta: u64) {
if bucket == RUSTFS_META_BUCKET {
return;
}
+10 -13
View File
@@ -322,6 +322,11 @@ impl ECStore {
pub fn scanner_namespace_mutation_generation(&self) -> u64 {
list_objects::scanner_namespace_mutation_generation()
}
pub async fn scanner_data_movement_active(&self) -> bool {
let (decommission, rebalance) = tokio::join!(self.is_decommission_running(), self.is_rebalance_started());
decommission || rebalance
}
}
// impl Clone for ECStore {
@@ -440,11 +445,7 @@ impl BucketOperations for ECStore {
#[instrument(skip(self))]
async fn make_bucket(&self, bucket: &str, opts: &MakeBucketOptions) -> Result<()> {
let result = self.handle_make_bucket(bucket, opts).await;
if result.is_ok() {
list_objects::observe_scanner_namespace_mutations(bucket, 1);
}
result
Box::pin(self.handle_make_bucket(bucket, opts)).await
}
#[instrument(skip(self))]
@@ -457,11 +458,7 @@ impl BucketOperations for ECStore {
}
#[instrument(skip(self))]
async fn delete_bucket(&self, bucket: &str, opts: &DeleteBucketOptions) -> Result<()> {
let result = self.handle_delete_bucket(bucket, opts).await;
if result.is_ok() {
list_objects::observe_scanner_namespace_mutations(bucket, 1);
}
result
Box::pin(self.handle_delete_bucket(bucket, opts)).await
}
}
@@ -871,9 +868,9 @@ mod tests {
// Build a minimal ECStore carrying an explicit instance context. Empty
// pools/disks are sufficient: the Phase 5 accessors read only `self.ctx`.
fn build_store_with_ctx(ctx: Arc<InstanceContext>) -> ECStore {
fn build_store_with_ctx(ctx: Arc<InstanceContext>) -> Arc<ECStore> {
let endpoint_pools = EndpointServerPools::default();
ECStore {
Arc::new(ECStore {
id: uuid::Uuid::new_v4(),
disk_map: std::collections::HashMap::new(),
pools: Vec::new(),
@@ -884,7 +881,7 @@ mod tests {
start_gate: Mutex::new(()),
pool_meta_save_gate: Mutex::new(()),
ctx,
}
})
}
// The object graph is the isolation carrier: two ECStore instances holding