diff --git a/crates/ecstore/src/api/mod.rs b/crates/ecstore/src/api/mod.rs index 34e784288..437f5c458 100644 --- a/crates/ecstore/src/api/mod.rs +++ b/crates/ecstore/src/api/mod.rs @@ -196,15 +196,16 @@ pub mod bucket { pub use crate::bucket::metadata_sys::ConfigWriteLockProbe; pub use crate::bucket::metadata_sys::{ BucketMetadataMutationGuard, BucketMetadataSys, ObjectLockConfigState, acquire_bucket_metadata_transaction_lock, - acquire_bucket_metadata_transaction_lock_for_incarnation, capture_bucket_metadata_incarnation, delete, - delete_if_incarnation, delete_under_transaction_lock, get, get_accelerate_config, get_bucket_policy, - get_bucket_policy_raw, get_bucket_targets_config, get_config_from_disk, get_cors_config, get_durability_config, - get_global_bucket_metadata_sys, get_lifecycle_config, get_logging_config, get_notification_config, - get_object_lock_config, get_object_lock_config_state, get_on_demand_migration_config, get_public_access_block_config, - get_quota_config, get_replication_config, get_request_payment_config, get_sse_config, get_tagging_config, - get_versioning_config, get_website_config, init_bucket_metadata_sys, list_bucket_targets, reload_bucket_metadata, - remove_bucket_metadata, set_bucket_metadata, update, update_bucket_targets_under_transaction_lock, - update_config_with, update_if_incarnation, update_quota_if_incarnation, update_under_transaction_lock, + acquire_bucket_metadata_transaction_lock_for_incarnation, acquire_scanner_bucket_incarnation_fence, + capture_bucket_metadata_incarnation, delete, delete_if_incarnation, delete_under_transaction_lock, get, + get_accelerate_config, get_bucket_policy, get_bucket_policy_raw, get_bucket_targets_config, get_config_from_disk, + get_cors_config, get_durability_config, get_global_bucket_metadata_sys, get_lifecycle_config, get_logging_config, + get_notification_config, get_object_lock_config, get_object_lock_config_state, get_on_demand_migration_config, + get_public_access_block_config, get_quota_config, get_replication_config, get_request_payment_config, get_sse_config, + get_tagging_config, get_versioning_config, get_website_config, init_bucket_metadata_sys, list_bucket_targets, + reload_bucket_metadata, remove_bucket_metadata, set_bucket_metadata, update, + update_bucket_targets_under_transaction_lock, update_config_with, update_if_incarnation, update_quota_if_incarnation, + update_under_transaction_lock, }; } diff --git a/crates/ecstore/src/bucket/metadata_sys.rs b/crates/ecstore/src/bucket/metadata_sys.rs index eb40ed8c9..4a53927f9 100644 --- a/crates/ecstore/src/bucket/metadata_sys.rs +++ b/crates/ecstore/src/bucket/metadata_sys.rs @@ -655,6 +655,12 @@ pub struct BucketMetadataMutationGuard { } impl BucketMetadataMutationGuard { + /// Returns the storage-verified identity while both incarnation fences remain valid. + pub fn checked_bucket_incarnation(&self) -> Result<(&str, Uuid)> { + self.ensure_valid(&self.bucket)?; + Ok((&self.bucket, self.incarnation_id)) + } + fn ensure_valid(&self, bucket: &str) -> Result<()> { if self.bucket != bucket { return Err(Error::other("bucket metadata mutation guard does not match bucket")); @@ -674,6 +680,29 @@ async fn acquire_config_write_guard_for_incarnation( sys: Arc>, bucket: &str, expected_incarnation_id: Option, +) -> Result { + acquire_config_write_guard_with_migration(sys, bucket, expected_incarnation_id, true).await +} + +/// Scanner probes must not create an incarnation to make a capability available. +pub async fn acquire_scanner_bucket_incarnation_fence( + bucket: &str, + expected_incarnation_id: Uuid, + expected_owner_id: Uuid, +) -> Result { + super::utils::check_valid_bucket_name(bucket)?; + let sys = get_bucket_metadata_sys()?; + if expected_owner_id.is_nil() || sys.read().await.api.id != expected_owner_id || expected_incarnation_id.is_nil() { + return Err(Error::other("scanner bucket incarnation owner does not match")); + } + acquire_config_write_guard_with_migration(sys, bucket, Some(expected_incarnation_id), false).await +} + +async fn acquire_config_write_guard_with_migration( + sys: Arc>, + bucket: &str, + expected_incarnation_id: Option, + migrate: bool, ) -> Result { let metadata_sys = sys.read().await.clone(); let lifecycle_guard = metadata_sys.api.acquire_bucket_lifecycle_read_lock(bucket).await?; @@ -681,13 +710,15 @@ async fn acquire_config_write_guard_for_incarnation( // Legacy buckets are migrated while the lifecycle fence prevents a // same-name replacement. The second read under the write transaction is // the CAS source of truth for the actual rewrite. - await_bucket_namespace_operation( - Some(&lifecycle_guard), - bucket, - "bucket config incarnation migration", - metadata_sys.get_bucket_incarnation_id(bucket), - ) - .await?; + if migrate { + await_bucket_namespace_operation( + Some(&lifecycle_guard), + bucket, + "bucket config incarnation migration", + metadata_sys.get_bucket_incarnation_id(bucket), + ) + .await?; + } let transaction_guard = await_bucket_namespace_operation( Some(&lifecycle_guard), bucket, @@ -3176,6 +3207,82 @@ mod tests { ); } + #[tokio::test] + async fn scoped_dirty_usage_incarnation_probe_does_not_migrate_legacy_metadata() { + let (dirs, store) = isolated_store_over_temp_disks().await; + let sys = Arc::new(RwLock::new(BucketMetadataSys::new(store.clone()))); + let bucket = "scoped-ack-legacy"; + for dir in &dirs { + std::fs::create_dir_all(dir.path().join(bucket)).expect("create legacy bucket"); + } + let mut metadata = BucketMetadata::new(bucket); + metadata.bucket_incarnation_id = Uuid::nil(); + sys.read() + .await + .persist_and_set(metadata) + .await + .expect("persist legacy metadata"); + assert!( + acquire_config_write_guard_with_migration(sys.clone(), bucket, Some(Uuid::new_v4()), false) + .await + .is_err() + ); + assert!(load_bucket_incarnation(store, bucket).await.expect("read sidecar").is_none()); + assert!( + sys.read() + .await + .get_config_from_disk(bucket) + .await + .expect("read metadata") + .bucket_incarnation_id + .is_nil() + ); + } + + #[tokio::test(flavor = "multi_thread", worker_threads = 4)] + #[serial] + async fn scoped_dirty_usage_incarnation_rejects_deleted_and_recreated_bucket() { + let (_dirs, store) = isolated_store_over_temp_disks().await; + init_bucket_metadata_sys(store.clone(), Vec::new()).await; + let sys = bucket_metadata_sys_of(&store.ctx).expect("metadata owner"); + let bucket = "scoped-ack-recreated"; + store + .make_bucket(bucket, &MakeBucketOptions::default()) + .await + .expect("create bucket"); + let old = store.bucket_incarnation_id_from_disk(bucket).await.expect("old incarnation"); + let guard = acquire_config_write_guard_with_migration(sys.clone(), bucket, Some(old), false) + .await + .expect("trusted incarnation fence"); + assert_eq!(guard.checked_bucket_incarnation().expect("valid fences"), (bucket, old)); + drop(guard); + store + .delete_bucket(bucket, &DeleteBucketOptions::default()) + .await + .expect("delete bucket"); + assert!( + acquire_config_write_guard_with_migration(sys.clone(), bucket, Some(old), false) + .await + .is_err() + ); + store + .make_bucket(bucket, &MakeBucketOptions::default()) + .await + .expect("recreate bucket"); + let new = store.bucket_incarnation_id_from_disk(bucket).await.expect("new incarnation"); + assert_ne!(old, new); + assert!( + acquire_config_write_guard_with_migration(sys.clone(), bucket, Some(old), false) + .await + .is_err() + ); + assert!( + acquire_config_write_guard_with_migration(sys, bucket, Some(new), false) + .await + .is_ok() + ); + } + #[tokio::test] async fn old_node_metadata_rewrite_cannot_replace_bucket_incarnation_sidecar() { let (dirs, ecstore) = isolated_store_over_temp_disks().await; diff --git a/crates/ecstore/src/cluster/rpc/client.rs b/crates/ecstore/src/cluster/rpc/client.rs index 435cbf64c..c4bf786f9 100644 --- a/crates/ecstore/src/cluster/rpc/client.rs +++ b/crates/ecstore/src/cluster/rpc/client.rs @@ -30,6 +30,7 @@ use rustfs_protos::{ ChannelClass, create_new_channel, get_channel_for_class, proto_gen::node_service::{ heal_control_service_client::HealControlServiceClient, node_service_client::NodeServiceClient, + scanner_control_service_client::ScannerControlServiceClient, tier_mutation_control_service_client::TierMutationControlServiceClient, }, }; @@ -60,6 +61,24 @@ pub async fn node_service_time_out_client( node_service_time_out_client_for_class(addr, interceptor, ChannelClass::Control).await } +pub(crate) async fn scanner_control_time_out_client( + addr: &str, + interceptor: TonicInterceptor, +) -> crate::error::Result>> { + let interceptor = interceptor.with_rpc_audience(addr)?; + let channel = match runtime_sources::cached_node_channel(addr).await { + Some(channel) => channel, + None => create_new_channel(addr) + .await + .map_err(|err| crate::error::Error::other(err.to_string()))?, + }; + let channel = ReplayScopeChannel::new(channel, interceptor.replay_scope_audience()); + let limit = rustfs_protos::scoped_dirty_usage::SCOPED_DIRTY_USAGE_MAX_REQUEST_BYTES as usize; + Ok(ScannerControlServiceClient::with_interceptor(channel, interceptor) + .max_decoding_message_size(limit) + .max_encoding_message_size(limit)) +} + pub async fn heal_control_time_out_client( addr: &str, interceptor: TonicInterceptor, diff --git a/crates/ecstore/src/cluster/rpc/peer_rest_client.rs b/crates/ecstore/src/cluster/rpc/peer_rest_client.rs index 5316199e3..9d57c4647 100644 --- a/crates/ecstore/src/cluster/rpc/peer_rest_client.rs +++ b/crates/ecstore/src/cluster/rpc/peer_rest_client.rs @@ -2050,6 +2050,53 @@ impl PeerRestClient { .await } + /// Probe only: scoped ACK production requires a durable per-bucket proof. + pub async fn scanner_scoped_dirty_usage_capability( + &self, + owner_id: String, + instance_id: String, + entries: Vec, + ) -> Result { + use rustfs_protos::scoped_dirty_usage::*; + let payload = rustfs_protos::proto_gen::node_service::ScannerScopedDirtyUsageAckRequest { + challenge: Uuid::new_v4().as_bytes().to_vec().into(), + protocol_version: SCOPED_DIRTY_USAGE_PROTOCOL_VERSION, + owner_id, + instance_id, + scope: SCOPED_DIRTY_USAGE_BUCKET_SCOPE, + probe_only: true, + entries, + }; + let canonical = canonical_scoped_dirty_usage_request(&payload).map_err(|err| Error::other(err.to_string()))?; + self.finalize_result( + async { + let mut client = super::client::scanner_control_time_out_client( + &self.grid_host, + TonicInterceptor::Signature(gen_tonic_signature_interceptor()), + ) + .await?; + let mut request = Request::new(payload.clone()); + set_tonic_canonical_body_digest(&mut request, &canonical)?; + let response = client.scanner_scoped_dirty_usage_ack(request).await?.into_inner(); + let body = canonical_scoped_dirty_usage_response(&canonical, &response) + .map_err(|_| Error::other("scoped dirty usage capability response is too large"))?; + verify_tonic_rpc_response_proof(&body, response.response_proof.as_ref())?; + if response.protocol_version != SCOPED_DIRTY_USAGE_PROTOCOL_VERSION + || response.owner_id != payload.owner_id + || response.instance_id != payload.instance_id + || response.max_entries != SCOPED_DIRTY_USAGE_MAX_ENTRIES + || response.max_request_bytes != SCOPED_DIRTY_USAGE_MAX_REQUEST_BYTES + || response.cleared != 0 + { + return Err(Error::other("scoped dirty usage capability response does not match request")); + } + Ok(response.supported) + } + .await, + ) + .await + } + pub async fn acknowledge_scanner_dirty_usage(&self, instance_id: String, generation: u64) -> Result { let result = self .scanner_activity_request_with_protocol(instance_id.clone(), generation, SCANNER_ACTIVITY_PROTOCOL_VERSION) diff --git a/crates/protos/src/generated/proto_gen/node_service.rs b/crates/protos/src/generated/proto_gen/node_service.rs index f015305f0..fe04d93ab 100644 --- a/crates/protos/src/generated/proto_gen/node_service.rs +++ b/crates/protos/src/generated/proto_gen/node_service.rs @@ -1283,6 +1283,54 @@ pub struct ScannerDirtyUsageSnapshotResponse { #[prost(bytes = "bytes", tag = "7")] pub response_proof: ::prost::bytes::Bytes, } +/// Receiver-only protocol. Producers must retain whole-cycle ACK until they +/// have a durable per-bucket publication proof. +#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] +pub struct ScannerScopedDirtyUsageEntry { + #[prost(string, tag = "1")] + pub bucket: ::prost::alloc::string::String, + #[prost(bytes = "bytes", tag = "2")] + pub bucket_incarnation: ::prost::bytes::Bytes, + #[prost(uint64, tag = "3")] + pub generation: u64, +} +#[derive(Clone, PartialEq, ::prost::Message)] +pub struct ScannerScopedDirtyUsageAckRequest { + #[prost(bytes = "bytes", tag = "1")] + pub challenge: ::prost::bytes::Bytes, + #[prost(uint32, tag = "2")] + pub protocol_version: u32, + #[prost(string, tag = "3")] + pub owner_id: ::prost::alloc::string::String, + #[prost(string, tag = "4")] + pub instance_id: ::prost::alloc::string::String, + /// Only scope 1 (a complete bucket) is supported; zero is invalid. + #[prost(uint32, tag = "5")] + pub scope: u32, + #[prost(bool, tag = "6")] + pub probe_only: bool, + #[prost(message, repeated, tag = "7")] + pub entries: ::prost::alloc::vec::Vec, +} +#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] +pub struct ScannerScopedDirtyUsageAckResponse { + #[prost(uint32, tag = "1")] + pub protocol_version: u32, + #[prost(string, tag = "2")] + pub owner_id: ::prost::alloc::string::String, + #[prost(string, tag = "3")] + pub instance_id: ::prost::alloc::string::String, + #[prost(bool, tag = "4")] + pub supported: bool, + #[prost(uint32, tag = "5")] + pub max_entries: u32, + #[prost(uint32, tag = "6")] + pub max_request_bytes: u32, + #[prost(uint64, tag = "7")] + pub cleared: u64, + #[prost(bytes = "bytes", tag = "8")] + pub response_proof: ::prost::bytes::Bytes, +} /// A short-lived storage-owned read admission used only around a final /// scanner metadata publication. It is intentionally separate from the /// ScannerActivity observation wire so v6/v7 rolling compatibility remains @@ -6282,6 +6330,244 @@ pub mod node_service_server { } } /// Generated client implementations. +pub mod scanner_control_service_client { + #![allow(unused_variables, dead_code, missing_docs, clippy::wildcard_imports, clippy::let_unit_value)] + use tonic::codegen::http::Uri; + use tonic::codegen::*; + #[derive(Debug, Clone)] + pub struct ScannerControlServiceClient { + inner: tonic::client::Grpc, + } + impl ScannerControlServiceClient { + /// Attempt to create a new client by connecting to a given endpoint. + pub async fn connect(dst: D) -> Result + where + D: TryInto, + D::Error: Into, + { + let conn = tonic::transport::Endpoint::new(dst)?.connect().await?; + Ok(Self::new(conn)) + } + } + impl ScannerControlServiceClient + where + T: tonic::client::GrpcService, + T::Error: Into, + T::ResponseBody: Body + std::marker::Send + 'static, + ::Error: Into + std::marker::Send, + { + pub fn new(inner: T) -> Self { + let inner = tonic::client::Grpc::new(inner); + Self { inner } + } + pub fn with_origin(inner: T, origin: Uri) -> Self { + let inner = tonic::client::Grpc::with_origin(inner, origin); + Self { inner } + } + pub fn with_interceptor(inner: T, interceptor: F) -> ScannerControlServiceClient> + where + F: tonic::service::Interceptor, + T::ResponseBody: Default, + T: tonic::codegen::Service< + http::Request, + Response = http::Response<>::ResponseBody>, + >, + >>::Error: + Into + std::marker::Send + std::marker::Sync, + { + ScannerControlServiceClient::new(InterceptedService::new(inner, interceptor)) + } + /// Compress requests with the given encoding. + /// + /// This requires the server to support it otherwise it might respond with an + /// error. + #[must_use] + pub fn send_compressed(mut self, encoding: CompressionEncoding) -> Self { + self.inner = self.inner.send_compressed(encoding); + self + } + /// Enable decompressing responses. + #[must_use] + pub fn accept_compressed(mut self, encoding: CompressionEncoding) -> Self { + self.inner = self.inner.accept_compressed(encoding); + self + } + /// Limits the maximum size of a decoded message. + /// + /// Default: `4MB` + #[must_use] + pub fn max_decoding_message_size(mut self, limit: usize) -> Self { + self.inner = self.inner.max_decoding_message_size(limit); + self + } + /// Limits the maximum size of an encoded message. + /// + /// Default: `usize::MAX` + #[must_use] + pub fn max_encoding_message_size(mut self, limit: usize) -> Self { + self.inner = self.inner.max_encoding_message_size(limit); + self + } + pub async fn scanner_scoped_dirty_usage_ack( + &mut self, + request: impl tonic::IntoRequest, + ) -> std::result::Result, tonic::Status> { + self.inner + .ready() + .await + .map_err(|e| tonic::Status::unknown(format!("Service was not ready: {}", e.into())))?; + let codec = tonic_prost::ProstCodec::default(); + let path = http::uri::PathAndQuery::from_static("/node_service.ScannerControlService/ScannerScopedDirtyUsageAck"); + let mut req = request.into_request(); + req.extensions_mut() + .insert(GrpcMethod::new("node_service.ScannerControlService", "ScannerScopedDirtyUsageAck")); + self.inner.unary(req, path, codec).await + } + } +} +/// Generated server implementations. +pub mod scanner_control_service_server { + #![allow(unused_variables, dead_code, missing_docs, clippy::wildcard_imports, clippy::let_unit_value)] + use tonic::codegen::*; + /// Generated trait containing gRPC methods that should be implemented for use with ScannerControlServiceServer. + #[async_trait] + pub trait ScannerControlService: std::marker::Send + std::marker::Sync + 'static { + async fn scanner_scoped_dirty_usage_ack( + &self, + request: tonic::Request, + ) -> std::result::Result, tonic::Status>; + } + #[derive(Debug)] + pub struct ScannerControlServiceServer { + inner: Arc, + accept_compression_encodings: EnabledCompressionEncodings, + send_compression_encodings: EnabledCompressionEncodings, + max_decoding_message_size: Option, + max_encoding_message_size: Option, + } + impl ScannerControlServiceServer { + pub fn new(inner: T) -> Self { + Self::from_arc(Arc::new(inner)) + } + pub fn from_arc(inner: Arc) -> Self { + Self { + inner, + accept_compression_encodings: Default::default(), + send_compression_encodings: Default::default(), + max_decoding_message_size: None, + max_encoding_message_size: None, + } + } + pub fn with_interceptor(inner: T, interceptor: F) -> InterceptedService + where + F: tonic::service::Interceptor, + { + InterceptedService::new(Self::new(inner), interceptor) + } + /// Enable decompressing requests with the given encoding. + #[must_use] + pub fn accept_compressed(mut self, encoding: CompressionEncoding) -> Self { + self.accept_compression_encodings.enable(encoding); + self + } + /// Compress responses with the given encoding, if the client supports it. + #[must_use] + pub fn send_compressed(mut self, encoding: CompressionEncoding) -> Self { + self.send_compression_encodings.enable(encoding); + self + } + /// Limits the maximum size of a decoded message. + /// + /// Default: `4MB` + #[must_use] + pub fn max_decoding_message_size(mut self, limit: usize) -> Self { + self.max_decoding_message_size = Some(limit); + self + } + /// Limits the maximum size of an encoded message. + /// + /// Default: `usize::MAX` + #[must_use] + pub fn max_encoding_message_size(mut self, limit: usize) -> Self { + self.max_encoding_message_size = Some(limit); + self + } + } + impl tonic::codegen::Service> for ScannerControlServiceServer + where + T: ScannerControlService, + B: Body + std::marker::Send + 'static, + B::Error: Into + std::marker::Send + 'static, + { + type Response = http::Response; + type Error = std::convert::Infallible; + type Future = BoxFuture; + fn poll_ready(&mut self, _cx: &mut Context<'_>) -> Poll> { + Poll::Ready(Ok(())) + } + fn call(&mut self, req: http::Request) -> Self::Future { + match req.uri().path() { + "/node_service.ScannerControlService/ScannerScopedDirtyUsageAck" => { + #[allow(non_camel_case_types)] + struct ScannerScopedDirtyUsageAckSvc(pub Arc); + impl tonic::server::UnaryService + for ScannerScopedDirtyUsageAckSvc + { + type Response = super::ScannerScopedDirtyUsageAckResponse; + type Future = BoxFuture, tonic::Status>; + fn call(&mut self, request: tonic::Request) -> Self::Future { + let inner = Arc::clone(&self.0); + let fut = async move { + ::scanner_scoped_dirty_usage_ack(&inner, request).await + }; + Box::pin(fut) + } + } + let accept_compression_encodings = self.accept_compression_encodings; + let send_compression_encodings = self.send_compression_encodings; + let max_decoding_message_size = self.max_decoding_message_size; + let max_encoding_message_size = self.max_encoding_message_size; + let inner = self.inner.clone(); + let fut = async move { + let method = ScannerScopedDirtyUsageAckSvc(inner); + let codec = tonic_prost::ProstCodec::default(); + let mut grpc = tonic::server::Grpc::new(codec) + .apply_compression_config(accept_compression_encodings, send_compression_encodings) + .apply_max_message_size_config(max_decoding_message_size, max_encoding_message_size); + let res = grpc.unary(method, req).await; + Ok(res) + }; + Box::pin(fut) + } + _ => Box::pin(async move { + let mut response = http::Response::new(tonic::body::Body::default()); + let headers = response.headers_mut(); + headers.insert(tonic::Status::GRPC_STATUS, (tonic::Code::Unimplemented as i32).into()); + headers.insert(http::header::CONTENT_TYPE, tonic::metadata::GRPC_CONTENT_TYPE); + Ok(response) + }), + } + } + } + impl Clone for ScannerControlServiceServer { + fn clone(&self) -> Self { + let inner = self.inner.clone(); + Self { + inner, + accept_compression_encodings: self.accept_compression_encodings, + send_compression_encodings: self.send_compression_encodings, + max_decoding_message_size: self.max_decoding_message_size, + max_encoding_message_size: self.max_encoding_message_size, + } + } + } + /// Generated gRPC service name + pub const SERVICE_NAME: &str = "node_service.ScannerControlService"; + impl tonic::server::NamedService for ScannerControlServiceServer { + const NAME: &'static str = SERVICE_NAME; + } +} +/// Generated client implementations. pub mod heal_control_service_client { #![allow(unused_variables, dead_code, missing_docs, clippy::wildcard_imports, clippy::let_unit_value)] use tonic::codegen::http::Uri; diff --git a/crates/protos/src/lib.rs b/crates/protos/src/lib.rs index fd8704382..81ead7d65 100644 --- a/crates/protos/src/lib.rs +++ b/crates/protos/src/lib.rs @@ -541,6 +541,8 @@ pub fn canonical_scanner_activity_v7_response_body( Ok(body) } +pub mod scoped_dirty_usage; + pub fn canonical_scanner_dirty_usage_snapshot_request_body( request: &proto_gen::node_service::ScannerDirtyUsageSnapshotRequest, ) -> Result, std::num::TryFromIntError> { diff --git a/crates/protos/src/node.proto b/crates/protos/src/node.proto index d3796991c..abc030adc 100644 --- a/crates/protos/src/node.proto +++ b/crates/protos/src/node.proto @@ -903,6 +903,36 @@ message ScannerDirtyUsageSnapshotResponse { bytes response_proof = 7; } +// Receiver-only protocol. Producers must retain whole-cycle ACK until they +// have a durable per-bucket publication proof. +message ScannerScopedDirtyUsageEntry { + string bucket = 1; + bytes bucket_incarnation = 2; + uint64 generation = 3; +} + +message ScannerScopedDirtyUsageAckRequest { + bytes challenge = 1; + uint32 protocol_version = 2; + string owner_id = 3; + string instance_id = 4; + // Only scope 1 (a complete bucket) is supported; zero is invalid. + uint32 scope = 5; + bool probe_only = 6; + repeated ScannerScopedDirtyUsageEntry entries = 7; +} + +message ScannerScopedDirtyUsageAckResponse { + uint32 protocol_version = 1; + string owner_id = 2; + string instance_id = 3; + bool supported = 4; + uint32 max_entries = 5; + uint32 max_request_bytes = 6; + uint64 cleared = 7; + bytes response_proof = 8; +} + // A short-lived storage-owned read admission used only around a final // scanner metadata publication. It is intentionally separate from the // ScannerActivity observation wire so v6/v7 rolling compatibility remains @@ -1245,6 +1275,10 @@ service NodeService { rpc GetLiveEvents(GetLiveEventsRequest) returns (GetLiveEventsResponse) {}; // auth-policy: read-only } +service ScannerControlService { + rpc ScannerScopedDirtyUsageAck(ScannerScopedDirtyUsageAckRequest) returns (ScannerScopedDirtyUsageAckResponse) {}; // auth-policy: body-bound +} + service HealControlService { rpc HealControl(HealControlRequest) returns (HealControlResponse) {}; } diff --git a/crates/protos/src/scoped_dirty_usage.rs b/crates/protos/src/scoped_dirty_usage.rs new file mode 100644 index 000000000..ac3effaf7 --- /dev/null +++ b/crates/protos/src/scoped_dirty_usage.rs @@ -0,0 +1,213 @@ +// Copyright 2024 RustFS Team +// Licensed under the Apache License, Version 2.0. + +//! Bounded, authenticated receiver contract for per-bucket dirty acknowledgements. + +use crate::CanonicalBodyBuilder; +use crate::proto_gen::node_service::{ScannerScopedDirtyUsageAckRequest, ScannerScopedDirtyUsageAckResponse}; +use prost::Message; + +pub const SCOPED_DIRTY_USAGE_PROTOCOL_VERSION: u32 = 1; +pub const SCOPED_DIRTY_USAGE_BUCKET_SCOPE: u32 = 1; +pub const SCOPED_DIRTY_USAGE_MAX_ENTRIES: u32 = 32; +pub const SCOPED_DIRTY_USAGE_MAX_REQUEST_BYTES: u32 = 8192; + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum ScopedDirtyUsageRequestError { + UnsupportedProtocol, + UnsupportedScope, + InvalidIdentity, + InvalidGeneration, + InvalidEntries, + TooLarge, +} + +impl std::fmt::Display for ScopedDirtyUsageRequestError { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.write_str(match self { + Self::UnsupportedProtocol => "unsupported scoped dirty usage protocol", + Self::UnsupportedScope => "unsupported scoped dirty usage scope", + Self::InvalidIdentity => "invalid scoped dirty usage identity", + Self::InvalidGeneration => "invalid scoped dirty usage generation", + Self::InvalidEntries => "scoped dirty usage entries must be nonempty and strictly ordered", + Self::TooLarge => "scoped dirty usage request exceeds its budget", + }) + } +} + +impl std::error::Error for ScopedDirtyUsageRequestError {} + +pub fn validate_scoped_dirty_usage_request( + request: &ScannerScopedDirtyUsageAckRequest, +) -> Result<(), ScopedDirtyUsageRequestError> { + use ScopedDirtyUsageRequestError as E; + if request.entries.len() > SCOPED_DIRTY_USAGE_MAX_ENTRIES as usize + || request.encoded_len() > SCOPED_DIRTY_USAGE_MAX_REQUEST_BYTES as usize + { + return Err(E::TooLarge); + } + if request.protocol_version != SCOPED_DIRTY_USAGE_PROTOCOL_VERSION { + return Err(E::UnsupportedProtocol); + } + if request.scope != SCOPED_DIRTY_USAGE_BUCKET_SCOPE { + return Err(E::UnsupportedScope); + } + if request.challenge.len() != 16 || request.owner_id.len() != 36 || request.instance_id.len() != 32 { + return Err(E::InvalidIdentity); + } + if request.entries.is_empty() || request.entries.windows(2).any(|pair| pair[0].bucket >= pair[1].bucket) { + return Err(E::InvalidEntries); + } + for entry in &request.entries { + if entry.bucket.is_empty() + || entry.bucket.len() > 63 + || entry.bucket_incarnation.len() != 16 + || entry.bucket_incarnation.iter().all(|byte| *byte == 0) + { + return Err(E::InvalidIdentity); + } + if entry.generation == 0 || entry.generation == u64::MAX { + return Err(E::InvalidGeneration); + } + } + Ok(()) +} + +pub fn canonical_scoped_dirty_usage_request( + request: &ScannerScopedDirtyUsageAckRequest, +) -> Result, ScopedDirtyUsageRequestError> { + validate_scoped_dirty_usage_request(request)?; + let mut body = CanonicalBodyBuilder::new(b"rustfs-scoped-dirty-usage-ack-request-v1\0"); + let encode = |_: std::num::TryFromIntError| ScopedDirtyUsageRequestError::TooLarge; + body.push_bytes(request.challenge.as_ref()).map_err(encode)?; + body.push_u32(request.protocol_version); + body.push_str(&request.owner_id).map_err(encode)?; + body.push_str(&request.instance_id).map_err(encode)?; + body.push_u32(request.scope); + body.push_bool(request.probe_only); + body.push_count(request.entries.len()).map_err(encode)?; + for entry in &request.entries { + body.push_str(&entry.bucket).map_err(encode)?; + body.push_bytes(entry.bucket_incarnation.as_ref()).map_err(encode)?; + body.push_u64(entry.generation); + } + Ok(body.finish()) +} + +pub fn canonical_scoped_dirty_usage_response( + request_body: &[u8], + response: &ScannerScopedDirtyUsageAckResponse, +) -> Result, std::num::TryFromIntError> { + let mut body = CanonicalBodyBuilder::new(b"rustfs-scoped-dirty-usage-ack-response-v1\0"); + body.push_bytes(request_body)?; + body.push_u32(response.protocol_version); + body.push_str(&response.owner_id)?; + body.push_str(&response.instance_id)?; + body.push_bool(response.supported); + body.push_u32(response.max_entries); + body.push_u32(response.max_request_bytes); + body.push_u64(response.cleared); + Ok(body.finish()) +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::proto_gen::node_service::ScannerScopedDirtyUsageEntry; + + fn request() -> ScannerScopedDirtyUsageAckRequest { + ScannerScopedDirtyUsageAckRequest { + challenge: vec![1; 16].into(), + protocol_version: 1, + owner_id: "11111111-1111-1111-1111-111111111111".into(), + instance_id: "a".repeat(32), + scope: 1, + probe_only: false, + entries: vec![ScannerScopedDirtyUsageEntry { + bucket: "photos".into(), + bucket_incarnation: vec![2; 16].into(), + generation: 8, + }], + } + } + + #[test] + fn scoped_dirty_usage_binds_every_request_field() { + let base = request(); + let baseline = canonical_scoped_dirty_usage_request(&base).expect("valid request"); + for field in 0..9 { + let mut changed = base.clone(); + match field { + 0 => changed.challenge = vec![3; 16].into(), + 1 => changed.protocol_version += 1, + 2 => changed.owner_id = "22222222-2222-2222-2222-222222222222".into(), + 3 => changed.instance_id = "b".repeat(32), + 4 => changed.scope += 1, + 5 => changed.probe_only = true, + 6 => changed.entries[0].bucket = "videos".into(), + 7 => changed.entries[0].bucket_incarnation = vec![3; 16].into(), + _ => changed.entries[0].generation += 1, + } + assert!(canonical_scoped_dirty_usage_request(&changed).map_or(true, |body| body != baseline)); + } + } + + #[test] + fn scoped_dirty_usage_binds_capability_and_ack_to_exact_request() { + let request = canonical_scoped_dirty_usage_request(&request()).expect("valid request"); + let response = ScannerScopedDirtyUsageAckResponse { + protocol_version: 1, + owner_id: "owner".into(), + instance_id: "process".into(), + supported: true, + max_entries: 32, + max_request_bytes: 8192, + cleared: 1, + response_proof: vec![1; 32].into(), + }; + let baseline = canonical_scoped_dirty_usage_response(&request, &response).expect("valid response"); + for field in 0..7 { + let mut changed = response.clone(); + match field { + 0 => changed.protocol_version += 1, + 1 => changed.owner_id.push('x'), + 2 => changed.instance_id.push('x'), + 3 => changed.supported = false, + 4 => changed.max_entries += 1, + 5 => changed.max_request_bytes += 1, + _ => changed.cleared += 1, + } + assert_ne!( + canonical_scoped_dirty_usage_response(&request, &changed).expect("response variant"), + baseline + ); + } + assert_ne!( + canonical_scoped_dirty_usage_response(b"another request", &response).expect("request variant"), + baseline + ); + } + + #[test] + fn scoped_dirty_usage_rejects_overflow_unknown_and_duplicate_entries() { + let base = request(); + let mut invalid = base.clone(); + invalid.entries = vec![base.entries[0].clone(); SCOPED_DIRTY_USAGE_MAX_ENTRIES as usize + 1]; + assert_eq!(validate_scoped_dirty_usage_request(&invalid), Err(ScopedDirtyUsageRequestError::TooLarge)); + invalid = base.clone(); + invalid.entries[0].bucket = "x".repeat(SCOPED_DIRTY_USAGE_MAX_REQUEST_BYTES as usize); + assert_eq!(validate_scoped_dirty_usage_request(&invalid), Err(ScopedDirtyUsageRequestError::TooLarge)); + invalid = base.clone(); + invalid.entries.push(base.entries[0].clone()); + assert_eq!( + validate_scoped_dirty_usage_request(&invalid), + Err(ScopedDirtyUsageRequestError::InvalidEntries) + ); + invalid = base; + invalid.entries[0].bucket_incarnation = vec![0; 16].into(); + assert_eq!( + validate_scoped_dirty_usage_request(&invalid), + Err(ScopedDirtyUsageRequestError::InvalidIdentity) + ); + } +} diff --git a/crates/scanner/src/lib.rs b/crates/scanner/src/lib.rs index f85a2e0b6..28cf638e9 100644 --- a/crates/scanner/src/lib.rs +++ b/crates/scanner/src/lib.rs @@ -90,8 +90,9 @@ pub use scanner::{ }; pub use scanner_io::{ ScannerDirtyUsageAckError, ScannerDirtyUsageBucket, ScannerDirtyUsageSnapshot, ScannerDirtyUsageState, - acknowledge_dirty_usage_generation, clear_dirty_usage_bucket, record_dirty_usage_bucket, record_scanner_maintenance_change, - scanner_activity_epoch, scanner_dirty_usage_snapshot, scanner_dirty_usage_state, scanner_maintenance_generation, + acknowledge_dirty_usage_generation, acknowledge_scoped_dirty_usage, clear_dirty_usage_bucket, record_dirty_usage_bucket, + record_scanner_maintenance_change, scanner_activity_epoch, scanner_dirty_usage_snapshot, scanner_dirty_usage_state, + scanner_maintenance_generation, }; pub use sleeper::{DynamicSleeper, SCANNER_IDLE_MODE, SCANNER_SLEEPER}; use std::sync::atomic::{AtomicU64, Ordering}; diff --git a/crates/scanner/src/scanner_io.rs b/crates/scanner/src/scanner_io.rs index c8353bae9..0d5e82eeb 100644 --- a/crates/scanner/src/scanner_io.rs +++ b/crates/scanner/src/scanner_io.rs @@ -883,8 +883,9 @@ pub(crate) use cache::{ }; pub use dirty_usage::{ ScannerDirtyUsageAckError, ScannerDirtyUsageBucket, ScannerDirtyUsageSnapshot, ScannerDirtyUsageState, - acknowledge_dirty_usage_generation, clear_dirty_usage_bucket, record_dirty_usage_bucket, record_scanner_maintenance_change, - scanner_activity_epoch, scanner_dirty_usage_snapshot, scanner_dirty_usage_state, scanner_maintenance_generation, + acknowledge_dirty_usage_generation, acknowledge_scoped_dirty_usage, clear_dirty_usage_bucket, record_dirty_usage_bucket, + record_scanner_maintenance_change, scanner_activity_epoch, scanner_dirty_usage_snapshot, scanner_dirty_usage_state, + scanner_maintenance_generation, }; #[cfg(test)] pub(crate) use dirty_usage::{clear_dirty_usage_buckets_for_tests, dirty_usage_buckets_for_tests}; diff --git a/crates/scanner/src/scanner_io/dirty_usage.rs b/crates/scanner/src/scanner_io/dirty_usage.rs index a5978e263..18db78a81 100644 --- a/crates/scanner/src/scanner_io/dirty_usage.rs +++ b/crates/scanner/src/scanner_io/dirty_usage.rs @@ -52,6 +52,112 @@ pub enum ScannerDirtyUsageAckError { ProcessChanged, #[error("scanner dirty usage generation cannot be acknowledged")] InvalidGeneration, + #[error("scanner dirty usage bucket incarnation fence is unavailable")] + IncarnationUnavailable, +} + +/// A scoped ACK requires storage-owned lifecycle and incarnation fences. +/// Callers must only send ACKs backed by durable per-bucket publication. +pub fn acknowledge_scoped_dirty_usage( + instance_id: &str, + entries: &[(&crate::storage_api::EcstoreBucketMetadataMutationGuard, u64)], + probe_only: bool, +) -> std::result::Result { + // Lock order: sorted bucket lifecycle/metadata fences (caller), then dirty map. + // No await or storage operation occurs while the dirty map is locked. + let (cleared, pending) = { + let mut dirty = dirty_usage_buckets(); + let checked = entries + .iter() + .map(|(guard, generation)| { + guard + .checked_bucket_incarnation() + .map(|(bucket, _)| (bucket, *generation)) + .map_err(|_| ScannerDirtyUsageAckError::IncarnationUnavailable) + }) + .collect::, _>>()?; + let cleared = apply_scoped_dirty_usage_ack( + instance_id, + scanner_activity_epoch(), + DIRTY_USAGE_BUCKET_GENERATION.load(Ordering::Acquire), + &mut dirty, + &checked, + probe_only, + )?; + if cleared > 0 { + advance_generation(&DIRTY_USAGE_BUCKET_GENERATION); + } + (cleared, dirty.len()) + }; + if !probe_only { + global_metrics().record_scanner_dirty_usage_cycle_clear(usize_to_u64_saturated(cleared), usize_to_u64_saturated(pending)); + } + Ok(usize_to_u64_saturated(cleared)) +} + +fn apply_scoped_dirty_usage_ack( + instance_id: &str, + current_instance: &str, + current_generation: u64, + dirty: &mut DirtyUsageBuckets, + entries: &[(&str, u64)], + probe_only: bool, +) -> std::result::Result { + if instance_id != current_instance { + return Err(ScannerDirtyUsageAckError::ProcessChanged); + } + if current_generation == u64::MAX + || entries + .iter() + .any(|(_, generation)| *generation == 0 || *generation == u64::MAX || *generation > current_generation) + { + return Err(ScannerDirtyUsageAckError::InvalidGeneration); + } + let mut cleared = 0; + if !probe_only { + for (bucket, generation) in entries { + if dirty.get(*bucket) == Some(generation) { + dirty.remove(*bucket); + cleared += 1; + } + } + } + Ok(cleared) +} + +#[cfg(test)] +mod scoped_dirty_usage_tests { + use super::*; + + #[test] + fn scoped_dirty_usage_preserves_uncovered_newer_and_replayed_generations() { + let mut dirty = HashMap::from([("hot".to_string(), 7), ("cold".to_string(), 8)]); + assert_eq!(apply_scoped_dirty_usage_ack("p", "p", 8, &mut dirty, &[("cold", 8)], true), Ok(0)); + assert_eq!(dirty.len(), 2); + assert_eq!(apply_scoped_dirty_usage_ack("p", "p", 8, &mut dirty, &[("cold", 8)], false), Ok(1)); + assert_eq!(dirty.get("hot"), Some(&7)); + assert_eq!(apply_scoped_dirty_usage_ack("p", "p", 8, &mut dirty, &[("cold", 8)], false), Ok(0)); + dirty.insert("cold".to_string(), 9); + assert_eq!(apply_scoped_dirty_usage_ack("p", "p", 9, &mut dirty, &[("cold", 8)], false), Ok(0)); + assert_eq!(dirty.get("cold"), Some(&9)); + } + + #[test] + fn scoped_dirty_usage_rejects_restart_and_invalid_batch_before_clearing() { + let original = HashMap::from([("hot".to_string(), 7), ("cold".to_string(), 8)]); + let mut dirty = original.clone(); + assert_eq!( + apply_scoped_dirty_usage_ack("old", "new", 8, &mut dirty, &[("cold", 8)], false), + Err(ScannerDirtyUsageAckError::ProcessChanged) + ); + for generation in [0, 9, u64::MAX] { + assert_eq!( + apply_scoped_dirty_usage_ack("p", "p", 8, &mut dirty, &[("cold", 8), ("hot", generation)], false), + Err(ScannerDirtyUsageAckError::InvalidGeneration) + ); + assert_eq!(dirty, original); + } + } } pub(super) fn dirty_usage_buckets() -> MutexGuard<'static, DirtyUsageBuckets> { diff --git a/crates/scanner/src/storage_api.rs b/crates/scanner/src/storage_api.rs index f81065d59..7b2cda493 100644 --- a/crates/scanner/src/storage_api.rs +++ b/crates/scanner/src/storage_api.rs @@ -38,8 +38,8 @@ pub(crate) use rustfs_ecstore::api::bucket::lifecycle::lifecycle::object_opts_fr #[cfg(test)] pub(crate) use rustfs_ecstore::api::bucket::metadata_sys::init_bucket_metadata_sys as ecstore_init_bucket_metadata_sys; pub(crate) use rustfs_ecstore::api::bucket::metadata_sys::{ - get_lifecycle_config as ecstore_get_lifecycle_config, get_object_lock_config as ecstore_get_object_lock_config, - get_replication_config as ecstore_get_replication_config, + BucketMetadataMutationGuard as EcstoreBucketMetadataMutationGuard, get_lifecycle_config as ecstore_get_lifecycle_config, + get_object_lock_config as ecstore_get_object_lock_config, get_replication_config as ecstore_get_replication_config, }; pub(crate) use rustfs_ecstore::api::bucket::replication::{ ReplicateObjectInfo, ReplicationConfig as EcstoreReplicationConfig, diff --git a/rustfs/src/server/http.rs b/rustfs/src/server/http.rs index 78ce81e5b..d82d67f20 100644 --- a/rustfs/src/server/http.rs +++ b/rustfs/src/server/http.rs @@ -208,6 +208,7 @@ const EVENT_PEER_ADDR_UNAVAILABLE: &str = "peer_addr_unavailable"; const EVENT_RPC_SIGNATURE_VERIFICATION_FAILED: &str = "rpc_signature_verification_failed"; const EVENT_GRPC_TRACE_CONTEXT_PROPAGATION_FAILED: &str = "grpc_trace_context_propagation_failed"; const HEAL_CONTROL_TONIC_RPC_PATH: &str = "/node_service.HealControlService/HealControl"; +const SCANNER_SCOPED_DIRTY_USAGE_ACK_TONIC_RPC_PATH: &str = "/node_service.ScannerControlService/ScannerScopedDirtyUsageAck"; const TIER_MUTATION_PREPARE_TONIC_RPC_PATH: &str = "/node_service.TierMutationControlService/PrepareTierMutation"; const TIER_MUTATION_COMMIT_TONIC_RPC_PATH: &str = "/node_service.TierMutationControlService/CommitTierMutation"; const TIER_MUTATION_ABORT_TONIC_RPC_PATH: &str = "/node_service.TierMutationControlService/AbortTierMutation"; @@ -1856,6 +1857,7 @@ fn process_connection( ); let rpc_service = RpcRequestPathService::new( Routes::new(node_service) + .add_service(InterceptedService::new(storage::tonic_service::make_scanner_control_server(), check_auth)) .add_service(heal_control_service) .add_service(tier_mutation_control_service) .prepare(), @@ -2259,6 +2261,7 @@ fn check_auth(req: Request<()>) -> std::result::Result, Status> { .strip_prefix(TONIC_RPC_PREFIX) .and_then(|suffix| suffix.strip_prefix('/')) .or_else(|| (target.uri.path() == HEAL_CONTROL_TONIC_RPC_PATH).then_some("HealControl")) + .or_else(|| (target.uri.path() == SCANNER_SCOPED_DIRTY_USAGE_ACK_TONIC_RPC_PATH).then_some("ScannerScopedDirtyUsageAck")) .or_else(|| (target.uri.path() == TIER_MUTATION_PREPARE_TONIC_RPC_PATH).then_some("PrepareTierMutation")) .or_else(|| (target.uri.path() == TIER_MUTATION_COMMIT_TONIC_RPC_PATH).then_some("CommitTierMutation")) .or_else(|| (target.uri.path() == TIER_MUTATION_ABORT_TONIC_RPC_PATH).then_some("AbortTierMutation")) @@ -3427,6 +3430,50 @@ mod tests { rustfs_common::set_global_local_node_name(&previous_node_name).await; } + #[tokio::test] + #[serial_test::serial] + async fn scoped_dirty_usage_peer_probe_reaches_handler_through_production_auth() { + let _ = rustfs_credentials::set_global_rpc_secret("rpc-http-test-secret".to_string()); + let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind scoped ACK auth test"); + let addr = listener.local_addr().expect("listener address"); + let previous_node_name = rustfs_common::get_global_local_node_name().await; + rustfs_common::set_global_local_node_name(&addr.to_string()).await; + let node = InterceptedService::new(NodeServiceServer::new(make_server()), check_auth); + let scanner = InterceptedService::new(storage::tonic_service::make_scanner_control_server(), check_auth); + let service = RpcRequestPathService::new(Routes::new(node).add_service(scanner).prepare()); + let server = tokio::spawn(async move { + let (socket, _) = listener.accept().await.expect("accept test connection"); + ConnBuilder::new(TokioExecutor::new()) + .serve_connection(TokioIo::new(socket), TowerToHyperService::new(service)) + .await + .expect("serve scoped ACK auth test"); + }); + let host = rustfs_utils::XHost::try_from(addr.to_string()).expect("peer address"); + let client = storage::PeerRestClient::new(host, format!("http://{addr}")); + let result = client + .scanner_scoped_dirty_usage_capability( + "11111111-1111-1111-1111-111111111111".to_string(), + "a".repeat(32), + vec![rustfs_protos::proto_gen::node_service::ScannerScopedDirtyUsageEntry { + bucket: "photos".into(), + bucket_incarnation: vec![1; 16].into(), + generation: 8, + }], + ) + .await; + client.evict_connection().await; + server.abort(); + let _ = server.await; + rustfs_common::set_global_local_node_name(&previous_node_name).await; + let error = result + .expect_err("probe must fail closed without the requested storage owner") + .to_string(); + assert!( + error.contains("storage layer is not initialized") || error.contains("scoped dirty usage peer or process changed"), + "signed probe must pass production path authentication and reach owner validation: {error}" + ); + } + #[tokio::test] #[serial_test::serial] async fn peer_rest_heal_control_uses_production_auth_and_keeps_validation_errors_online() { diff --git a/rustfs/src/storage/rpc/node_service.rs b/rustfs/src/storage/rpc/node_service.rs index 0dd18424e..111726af3 100644 --- a/rustfs/src/storage/rpc/node_service.rs +++ b/rustfs/src/storage/rpc/node_service.rs @@ -493,6 +493,13 @@ impl std::fmt::Debug for NodeService { } } +pub(crate) fn make_scanner_control_server() -> scanner_control_service_server::ScannerControlServiceServer { + let limit = rustfs_protos::scoped_dirty_usage::SCOPED_DIRTY_USAGE_MAX_REQUEST_BYTES as usize; + scanner_control_service_server::ScannerControlServiceServer::new(make_server()) + .max_decoding_message_size(limit) + .max_encoding_message_size(limit) +} + pub fn make_server() -> NodeService { let context = runtime_sources::current_app_context(); make_server_for_context(context) @@ -1087,6 +1094,74 @@ impl NodeService { } } +#[tonic::async_trait] +impl scanner_control_service_server::ScannerControlService for NodeService { + async fn scanner_scoped_dirty_usage_ack( + &self, + request: Request, + ) -> Result, Status> { + use rustfs_protos::scoped_dirty_usage::*; + static ADMISSION: tokio::sync::Semaphore = tokio::sync::Semaphore::const_new(4); + + let canonical = + canonical_scoped_dirty_usage_request(request.get_ref()).map_err(|err| Status::invalid_argument(err.to_string()))?; + verify_tonic_canonical_body_digest(&request, &canonical) + .map_err(|_| Status::permission_denied("scoped dirty usage authentication failed"))?; + let _admission = ADMISSION + .try_acquire() + .map_err(|_| Status::resource_exhausted("scoped dirty usage receiver is busy"))?; + let request = request.into_inner(); + let store = self + .resolve_object_store() + .ok_or_else(|| Status::unavailable("storage layer is not initialized"))?; + if store.id.is_nil() + || request.owner_id != store.id.to_string() + || request.instance_id != rustfs_scanner::scanner_activity_epoch() + { + return Err(Status::failed_precondition("scoped dirty usage peer or process changed")); + } + let cleared = timeout(Duration::from_secs(30), async { + // Strict bucket order is validated before admission. Acquire every + // lifecycle/metadata fence before clearing any dirty record. + let mut guards = Vec::with_capacity(request.entries.len()); + for entry in &request.entries { + let incarnation = Uuid::from_slice(entry.bucket_incarnation.as_ref()) + .map_err(|_| Status::invalid_argument("invalid bucket incarnation"))?; + let guard = + crate::storage::storage_api::acquire_scanner_bucket_incarnation_fence(&entry.bucket, incarnation, store.id) + .await + .map_err(|_| Status::failed_precondition("trusted bucket incarnation is unavailable"))?; + guards.push(guard); + } + let entries = guards + .iter() + .zip(&request.entries) + .map(|(guard, entry)| (guard, entry.generation)) + .collect::>(); + rustfs_scanner::acknowledge_scoped_dirty_usage(&request.instance_id, &entries, request.probe_only) + .map_err(|err| Status::failed_precondition(err.to_string())) + }) + .await + .map_err(|_| Status::deadline_exceeded("scoped dirty usage incarnation validation timed out"))??; + let mut response = ScannerScopedDirtyUsageAckResponse { + protocol_version: SCOPED_DIRTY_USAGE_PROTOCOL_VERSION, + owner_id: request.owner_id, + instance_id: request.instance_id, + supported: true, + max_entries: SCOPED_DIRTY_USAGE_MAX_ENTRIES, + max_request_bytes: SCOPED_DIRTY_USAGE_MAX_REQUEST_BYTES, + cleared, + response_proof: Bytes::new(), + }; + let body = canonical_scoped_dirty_usage_response(&canonical, &response) + .map_err(|_| Status::internal("scoped dirty usage response is too large"))?; + response.response_proof = sign_tonic_rpc_response_proof(&body) + .map_err(|_| Status::unavailable("scoped dirty usage response authentication is unavailable"))? + .into(); + Ok(Response::new(response)) + } +} + #[tonic::async_trait] impl Node for NodeService { async fn ping(&self, request: Request) -> Result, Status> { @@ -2623,6 +2698,7 @@ mod tests { use rustfs_kms::KmsServiceManager; use rustfs_protos::CanonicalMutationBody as _; use rustfs_protos::models::PingBodyBuilder; + use rustfs_protos::proto_gen::node_service::scanner_control_service_server::ScannerControlService as _; use rustfs_protos::proto_gen::node_service::{ BackgroundHealStatusRequest, BatchGenerallyLockRequest, CancelDecommissionRequest, CheckPartsRequest, ClearDecommissionRequest, ControlPlaneErrorCode, DeleteBucketMetadataRequest, DeleteBucketRequest, DeletePathsRequest, @@ -5990,6 +6066,74 @@ mod tests { ); } + fn scoped_dirty_usage_request() -> rustfs_protos::proto_gen::node_service::ScannerScopedDirtyUsageAckRequest { + rustfs_protos::proto_gen::node_service::ScannerScopedDirtyUsageAckRequest { + challenge: vec![7; 16].into(), + protocol_version: 1, + owner_id: "11111111-1111-1111-1111-111111111111".into(), + instance_id: "a".repeat(32), + scope: 1, + probe_only: false, + entries: vec![rustfs_protos::proto_gen::node_service::ScannerScopedDirtyUsageEntry { + bucket: "photos".into(), + bucket_incarnation: vec![1; 16].into(), + generation: 8, + }], + } + } + + #[tokio::test] + async fn scoped_dirty_usage_authenticates_before_storage_and_rejects_tampering() { + use rustfs_protos::scoped_dirty_usage::canonical_scoped_dirty_usage_request; + let service = create_test_node_service(); + let unsigned = service + .scanner_scoped_dirty_usage_ack(Request::new(scoped_dirty_usage_request())) + .await + .expect_err("unsigned ACK must not access storage"); + assert_eq!(unsigned.code(), tonic::Code::PermissionDenied); + for field in 0..9 { + let mut signed = Request::new(scoped_dirty_usage_request()); + let canonical = canonical_scoped_dirty_usage_request(signed.get_ref()).expect("canonical request"); + set_tonic_canonical_body_digest(&mut signed, &canonical).expect("digest"); + mark_v2_authenticated(&mut signed); + match field { + 0 => signed.get_mut().challenge = vec![3; 16].into(), + 1 => signed.get_mut().owner_id = "22222222-2222-2222-2222-222222222222".into(), + 2 => signed.get_mut().instance_id = "b".repeat(32), + 3 => signed.get_mut().probe_only = true, + 4 => signed.get_mut().entries[0].bucket = "videos".into(), + 5 => signed.get_mut().entries[0].bucket_incarnation = vec![2; 16].into(), + 6 => signed.get_mut().entries[0].generation += 1, + 7 => signed.get_mut().scope += 1, + _ => signed.get_mut().protocol_version += 1, + } + let error = service + .scanner_scoped_dirty_usage_ack(signed) + .await + .expect_err("tampered ACK must fail"); + assert_eq!( + error.code(), + if field < 7 { + tonic::Code::PermissionDenied + } else { + tonic::Code::InvalidArgument + } + ); + } + let mut signed = Request::new(scoped_dirty_usage_request()); + let canonical = canonical_scoped_dirty_usage_request(signed.get_ref()).expect("canonical request"); + set_tonic_canonical_body_digest(&mut signed, &canonical).expect("digest"); + mark_v2_authenticated(&mut signed); + assert_eq!( + service + .scanner_scoped_dirty_usage_ack(signed) + .await + .expect_err("missing owner cannot advertise capability") + .code(), + tonic::Code::Unavailable + ); + } + #[tokio::test] async fn test_scanner_activity_requires_body_bound_auth_before_storage_lookup() { let service = create_test_node_service(); @@ -6485,6 +6629,62 @@ mod tests { ) } + #[tokio::test] + async fn scoped_dirty_usage_transport_rejects_oversized_unknown_and_duplicate_fields() { + let listener = TcpListener::bind("127.0.0.1:0") + .await + .expect("bind scoped ACK transport test"); + let addr = listener.local_addr().expect("test listener address"); + let (shutdown, stopped) = tokio::sync::oneshot::channel(); + let server = tokio::spawn(async move { + tonic::transport::Server::builder() + .add_service(super::make_scanner_control_server()) + .serve_with_incoming_shutdown(TcpListenerStream::new(listener), async { + let _ = stopped.await; + }) + .await + .expect("scoped ACK transport server"); + }); + let client = reqwest::Client::builder() + .no_proxy() + .http2_prior_knowledge() + .build() + .expect("HTTP/2 client"); + let limit = rustfs_protos::scoped_dirty_usage::SCOPED_DIRTY_USAGE_MAX_REQUEST_BYTES as usize; + for tag in [0x78, 0x0a] { + // Unknown varint field 15, or repeated empty singular challenge: + // both decode to a tiny default struct despite the large wire body. + for oversized in [false, true] { + let mut payload = [tag, 0].repeat(if oversized { (limit - 4) / 2 } else { limit / 2 }); + if oversized { + // Unknown fixed32 field 15 makes a valid cap+1 protobuf. + payload.extend_from_slice(&[0x7d, 0, 0, 0, 0]); + } + assert_eq!(payload.len(), limit + usize::from(oversized)); + let mut frame = vec![0]; + frame.extend_from_slice(&u32::try_from(payload.len()).expect("bounded test payload").to_be_bytes()); + frame.extend_from_slice(&payload); + let response = client + .post(format!("http://{addr}/node_service.ScannerControlService/ScannerScopedDirtyUsageAck")) + .header("content-type", "application/grpc") + .header("te", "trailers") + .body(frame) + .send() + .await + .expect("send raw protobuf frame"); + let status = response.headers().get("grpc-status").expect("gRPC failure status"); + assert_eq!( + status.to_str().expect("status text"), + if oversized { "11" } else { "3" }, + "cap+1 must fail in the codec, while cap bytes reach request validation" + ); + } + } + drop(client); + shutdown.send(()).expect("stop test server"); + server.await.expect("join test server"); + } + #[tokio::test] async fn heal_control_transport_enforces_codec_limit_and_fails_closed() { let Some(mut client) = connect_test_heal_control_client().await else { diff --git a/rustfs/src/storage/storage_api.rs b/rustfs/src/storage/storage_api.rs index 9e1e92dc4..2216db7f4 100644 --- a/rustfs/src/storage/storage_api.rs +++ b/rustfs/src/storage/storage_api.rs @@ -379,7 +379,7 @@ pub(crate) mod tonic_service_consumer { #[cfg(test)] pub(crate) use super::super::tonic_service::{heal_topology_fingerprint, make_heal_control_server_for_source}; pub(crate) use super::super::tonic_service::{ - make_heal_control_server_with_cache, make_server, make_tier_mutation_control_server, + make_heal_control_server_with_cache, make_scanner_control_server, make_server, make_tier_mutation_control_server, }; } @@ -1704,6 +1704,14 @@ pub(crate) async fn acquire_bucket_metadata_transaction_lock( ecstore_bucket::metadata_sys::acquire_bucket_metadata_transaction_lock(bucket).await } +pub(crate) async fn acquire_scanner_bucket_incarnation_fence( + bucket: &str, + incarnation: uuid::Uuid, + owner_id: uuid::Uuid, +) -> Result { + ecstore_bucket::metadata_sys::acquire_scanner_bucket_incarnation_fence(bucket, incarnation, owner_id).await +} + pub(crate) async fn update_bucket_targets_under_transaction_lock( guard: &ecstore_bucket::metadata_sys::BucketMetadataMutationGuard, bucket: &str, diff --git a/rustfs/src/storage/tonic_service.rs b/rustfs/src/storage/tonic_service.rs index 539361507..c7d2ee28a 100644 --- a/rustfs/src/storage/tonic_service.rs +++ b/rustfs/src/storage/tonic_service.rs @@ -13,6 +13,7 @@ // limitations under the License. pub(crate) use crate::storage::rpc::node_service::make_heal_control_server_with_cache; +pub(crate) use crate::storage::rpc::node_service::make_scanner_control_server; #[cfg(test)] pub(crate) use crate::storage::rpc::node_service::{heal::heal_topology_fingerprint, make_heal_control_server_for_source}; pub use crate::storage::rpc::{make_heal_control_server, make_server, make_tier_mutation_control_server}; diff --git a/rustfs/src/storage_api.rs b/rustfs/src/storage_api.rs index db459b72d..4f824f8da 100644 --- a/rustfs/src/storage_api.rs +++ b/rustfs/src/storage_api.rs @@ -176,7 +176,7 @@ pub(crate) mod server { heal_topology_fingerprint, make_heal_control_server_for_source, }; pub(crate) use crate::storage::storage_api::tonic_service_consumer::{ - make_heal_control_server_with_cache, make_server, make_tier_mutation_control_server, + make_heal_control_server_with_cache, make_scanner_control_server, make_server, make_tier_mutation_control_server, }; } }