From 050ae3ed813ee22658c4304772553172c0be32c1 Mon Sep 17 00:00:00 2001 From: houseme Date: Mon, 7 Sep 2026 01:17:50 +0800 Subject: [PATCH] feat(scanner): send scoped dirty usage acknowledgements Bind dirty usage snapshots to storage owner and bucket incarnation IDs, then route completed dirty-scope scans through scoped peer acknowledgements instead of whole-generation clears. Keep mixed-version and high-cardinality cases conservative by probing scoped ACK capability, capping automatic scoped activation to one protocol batch, and falling back to full-scan generation ACK when the scoped path is not clearly supported. Co-Authored-By: heihutu Co-Authored-By: zhi22915 --- crates/ecstore/src/api/mod.rs | 3 +- crates/ecstore/src/cluster/rpc/mod.rs | 3 +- .../src/cluster/rpc/peer_rest_client.rs | 315 ++++++++++++++++-- .../ecstore/src/services/notification_sys.rs | 119 ++++++- .../src/generated/proto_gen/node_service.rs | 4 + crates/protos/src/lib.rs | 11 + crates/protos/src/node.proto | 2 + crates/scanner/src/lib.rs | 1 + crates/scanner/src/scanner.rs | 12 +- crates/scanner/src/scanner/activity.rs | 31 +- crates/scanner/src/scanner/tests.rs | 6 +- .../scanner/tests/scoped_ack_publication.rs | 2 +- crates/scanner/src/scanner_io.rs | 51 ++- crates/scanner/src/scanner_io/io_cycle.rs | 107 +++++- crates/scanner/src/scanner_io/tests.rs | 70 +++- crates/scanner/src/storage_api.rs | 8 +- crates/storage-api/src/lib.rs | 5 +- rustfs/src/server/http.rs | 7 +- rustfs/src/storage/rpc/node_service.rs | 68 +++- rustfs/src/storage/storage_api.rs | 6 +- rustfs/src/storage_api.rs | 4 +- scripts/error-other-format-baseline.txt | 2 +- 22 files changed, 728 insertions(+), 109 deletions(-) diff --git a/crates/ecstore/src/api/mod.rs b/crates/ecstore/src/api/mod.rs index 9101a5577..acdd95ee5 100644 --- a/crates/ecstore/src/api/mod.rs +++ b/crates/ecstore/src/api/mod.rs @@ -520,7 +520,8 @@ pub mod rpc { pub use crate::cluster::rpc::{ AuthenticatedChannel, KMS_SIGNAL_SUBSYSTEM, LocalPeerS3Client, PEER_RESTDRY_RUN, PEER_RESTSIGNAL, PEER_RESTSUB_SYS, PeerRestClient, PeerS3Client, S3PeerSys, SERVICE_SIGNAL_REFRESH_CONFIG, SERVICE_SIGNAL_RELOAD_DYNAMIC, - ScannerBucketListing, ScannerPeerActivity, ScannerPeerDirtyUsageSnapshot, ScannerPublicationLease, TONIC_RPC_PREFIX, + ScannerBucketListing, ScannerDirtyUsageAcknowledgement, ScannerPeerActivity, ScannerPeerDirtyUsageBucket, + ScannerPeerDirtyUsageSnapshot, ScannerPublicationLease, ScannerScopedDirtyUsageAckEntry, TONIC_RPC_PREFIX, TonicInterceptor, build_put_file_auth_trailer, check_and_record_signed_rpc_nonce, decode_heal_bucket_rpc_options, encode_heal_bucket_rpc_options, gen_signature_headers, gen_tonic_replay_scope_headers, gen_tonic_signature_headers, gen_tonic_signature_interceptor, node_service_time_out_client, node_service_time_out_client_no_auth, diff --git a/crates/ecstore/src/cluster/rpc/mod.rs b/crates/ecstore/src/cluster/rpc/mod.rs index d0987a063..dc7405719 100644 --- a/crates/ecstore/src/cluster/rpc/mod.rs +++ b/crates/ecstore/src/cluster/rpc/mod.rs @@ -48,7 +48,8 @@ pub use internode_data_transport::build_internode_data_transport_from_env; pub(crate) use peer_rest_client::TierConfigReloadOutcome; pub use peer_rest_client::{ KMS_SIGNAL_SUBSYSTEM, PEER_RESTDRY_RUN, PEER_RESTSIGNAL, PEER_RESTSUB_SYS, PeerRestClient, SERVICE_SIGNAL_REFRESH_CONFIG, - SERVICE_SIGNAL_RELOAD_DYNAMIC, ScannerPeerActivity, ScannerPeerDirtyUsageSnapshot, ScannerPublicationLease, + SERVICE_SIGNAL_RELOAD_DYNAMIC, ScannerDirtyUsageAcknowledgement, ScannerPeerActivity, ScannerPeerDirtyUsageBucket, + ScannerPeerDirtyUsageSnapshot, ScannerPublicationLease, ScannerScopedDirtyUsageAckEntry, }; pub(crate) use peer_s3_client::heal_bucket_local_on_disks; pub use peer_s3_client::{ diff --git a/crates/ecstore/src/cluster/rpc/peer_rest_client.rs b/crates/ecstore/src/cluster/rpc/peer_rest_client.rs index f6f5f58e8..c0320b237 100644 --- a/crates/ecstore/src/cluster/rpc/peer_rest_client.rs +++ b/crates/ecstore/src/cluster/rpc/peer_rest_client.rs @@ -49,10 +49,11 @@ use rustfs_protos::proto_gen::node_service::{ LoadRebalanceMetaRequest, LoadServiceAccountRequest, LoadTransitionTierConfigRequest, LoadUserRequest, LocalStorageInfoRequest, Mss, ReloadPoolMetaRequest, ReloadSiteReplicationConfigRequest, ReplacementRecoveryStatusRequest, ScannerActivityRequest, ScannerActivityResponse, ScannerDirtyUsageSnapshotRequest, ScannerDirtyUsageSnapshotResponse, - ScannerPublicationLeaseReleaseRequest, ScannerPublicationLeaseRequest, ScannerPublicationLeaseResponse, ServerInfoRequest, - SignalServiceRequest, SignalServiceResponse, StartDecommissionRequest, StartProfilingRequest, StopRebalanceRequest, - TierDailyStatsRequest, TierMutationAbortRequest, TierMutationCommitRequest, TierMutationControlResponse, - TierMutationFailureClass, TierMutationPeerState, TierMutationPrepareRequest, node_service_client::NodeServiceClient, + ScannerPublicationLeaseReleaseRequest, ScannerPublicationLeaseRequest, ScannerPublicationLeaseResponse, + ScannerScopedDirtyUsageAckRequest, ScannerScopedDirtyUsageEntry, ServerInfoRequest, SignalServiceRequest, + SignalServiceResponse, StartDecommissionRequest, StartProfilingRequest, StopRebalanceRequest, TierDailyStatsRequest, + TierMutationAbortRequest, TierMutationCommitRequest, TierMutationControlResponse, TierMutationFailureClass, + TierMutationPeerState, TierMutationPrepareRequest, node_service_client::NodeServiceClient, tier_mutation_control_service_client::TierMutationControlServiceClient, }; pub use rustfs_protos::{PEER_RESTDRY_RUN, PEER_RESTSIGNAL, PEER_RESTSUB_SYS}; @@ -92,6 +93,7 @@ const HEAL_CONTROL_PAYLOAD_MAX_SIZE: usize = 64 * 1024; const PEER_REST_RECOVERY_MAX_ATTEMPTS: u32 = 60; const PEER_REST_RECOVERY_MAX_BACKOFF: Duration = Duration::from_secs(30); const SCANNER_ACTIVITY_MAX_MESSAGE_SIZE: usize = 1024; +const SCANNER_SCOPED_DIRTY_USAGE_STAGE_TIMEOUT: Duration = Duration::from_secs(5); /// Reserve time for the acquire response's network/clock uncertainty. The /// server owns the real expiry; this local deadline is intentionally earlier /// so a coordinator never starts a bounded persistence operation at the edge @@ -192,12 +194,102 @@ pub struct ScannerPeerActivity { #[derive(Clone, Debug, PartialEq, Eq)] pub struct ScannerPeerDirtyUsageSnapshot { + pub owner_id: String, pub instance_id: String, pub generation: u64, pub pending_bucket_count: u64, pub protocol_version: u32, pub complete: bool, - pub buckets: BTreeMap, + pub buckets: BTreeMap, +} + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub struct ScannerPeerDirtyUsageBucket { + pub bucket_incarnation: Uuid, + pub generation: u64, +} + +#[derive(Clone, Debug, PartialEq, Eq)] +pub struct ScannerScopedDirtyUsageAckEntry { + pub bucket: String, + pub bucket_incarnation: Uuid, + pub generation: u64, +} + +#[derive(Clone, Debug, PartialEq, Eq)] +pub enum ScannerDirtyUsageAcknowledgement { + Generation { + host: String, + instance_id: String, + generation: u64, + }, + Scoped { + host: String, + owner_id: String, + instance_id: String, + entries: Vec, + }, +} + +fn scanner_scoped_dirty_usage_ack_payloads( + owner_id: String, + instance_id: String, + probe_only: bool, + entries: Vec, +) -> Result> { + use rustfs_protos::scoped_dirty_usage::*; + + if entries.is_empty() { + return Err(Error::other("scoped dirty usage acknowledgement entries must be nonempty")); + } + + let mut payloads = Vec::with_capacity(entries.len().div_ceil(SCOPED_DIRTY_USAGE_MAX_ENTRIES as usize)); + let mut batch = Vec::with_capacity(SCOPED_DIRTY_USAGE_MAX_ENTRIES as usize); + for entry in entries { + batch.push(ScannerScopedDirtyUsageEntry { + bucket: entry.bucket, + bucket_incarnation: entry.bucket_incarnation.as_bytes().to_vec().into(), + generation: entry.generation, + }); + if batch.len() == SCOPED_DIRTY_USAGE_MAX_ENTRIES as usize { + payloads.push(scanner_scoped_dirty_usage_ack_payload( + &owner_id, + &instance_id, + probe_only, + std::mem::take(&mut batch), + )?); + } + } + if !batch.is_empty() { + payloads.push(scanner_scoped_dirty_usage_ack_payload(&owner_id, &instance_id, probe_only, batch)?); + } + + Ok(payloads) +} + +fn scanner_scoped_dirty_usage_ack_payload( + owner_id: &str, + instance_id: &str, + probe_only: bool, + entries: Vec, +) -> Result { + use rustfs_protos::scoped_dirty_usage::*; + + let payload = ScannerScopedDirtyUsageAckRequest { + challenge: Uuid::new_v4().as_bytes().to_vec().into(), + protocol_version: SCOPED_DIRTY_USAGE_PROTOCOL_VERSION, + owner_id: owner_id.to_string(), + instance_id: instance_id.to_string(), + scope: SCOPED_DIRTY_USAGE_BUCKET_SCOPE, + probe_only, + entries, + }; + canonical_scoped_dirty_usage_request(&payload).map_err(|err| Error::other(err.to_string()))?; + Ok(payload) +} + +fn scanner_scoped_dirty_usage_ack_reconciled(activity: &ScannerPeerActivity, expected_instance_id: &str) -> bool { + activity.instance_id == expected_instance_id && activity.dirty_usage_pending == Some(false) } fn scanner_instance_id_is_valid(instance_id: &str) -> bool { @@ -351,6 +443,11 @@ fn decode_scanner_dirty_usage_snapshot_with_verifier( if !scanner_instance_id_is_valid(&response.instance_id) { return Err(Error::other("peer returned an invalid scanner dirty usage snapshot instance ID")); } + let owner_id = Uuid::parse_str(&response.owner_id) + .ok() + .filter(|owner_id| !owner_id.is_nil()) + .map(|owner_id| owner_id.to_string()) + .ok_or_else(|| Error::other("peer returned an invalid scanner dirty usage snapshot owner"))?; if response.generation == u64::MAX { return Err(Error::other("peer scanner dirty usage snapshot exhausted its generation")); } @@ -386,9 +483,14 @@ fn decode_scanner_dirty_usage_snapshot_with_verifier( if bucket.generation == 0 || bucket.generation > response.generation { return Err(Error::other("peer scanner dirty usage snapshot contains an invalid bucket generation")); } + Uuid::from_slice(bucket.bucket_incarnation.as_ref()) + .ok() + .filter(|bucket_incarnation| !bucket_incarnation.is_nil()) + .ok_or_else(|| Error::other("peer scanner dirty usage snapshot contains an invalid bucket incarnation"))?; } Ok(ScannerPeerDirtyUsageSnapshot { + owner_id, instance_id: response.instance_id, generation: response.generation, pending_bucket_count: response.pending_bucket_count, @@ -397,7 +499,16 @@ fn decode_scanner_dirty_usage_snapshot_with_verifier( buckets: response .buckets .into_iter() - .map(|bucket| (bucket.bucket, bucket.generation)) + .map(|bucket| { + ( + bucket.bucket, + ScannerPeerDirtyUsageBucket { + bucket_incarnation: Uuid::from_slice(bucket.bucket_incarnation.as_ref()) + .expect("bucket incarnation was validated"), + generation: bucket.generation, + }, + ) + }) .collect(), }) } @@ -2077,19 +2188,10 @@ impl PeerRestClient { &self, owner_id: String, instance_id: String, - entries: Vec, + 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()))?; + let payloads = scanner_scoped_dirty_usage_ack_payloads(owner_id, instance_id, true, entries)?; self.finalize_result( async { let mut client = super::client::scanner_control_time_out_client( @@ -2097,26 +2199,106 @@ impl PeerRestClient { TonicInterceptor::Signature(gen_tonic_signature_interceptor()), ) .await?; + for payload in payloads { + let canonical = + canonical_scoped_dirty_usage_request(&payload).map_err(|err| Error::other(err.to_string()))?; + 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")); + } + if !response.supported { + return Ok(false); + } + } + Ok(true) + } + .await, + ) + .await + } + + pub async fn acknowledge_scanner_scoped_dirty_usage( + &self, + owner_id: String, + instance_id: String, + entries: Vec, + ) -> Result { + use rustfs_protos::scoped_dirty_usage::*; + let payloads = scanner_scoped_dirty_usage_ack_payloads(owner_id, instance_id.clone(), false, entries)?; + let ack_attempt = async { + let mut client = super::client::scanner_control_time_out_client( + &self.grid_host, + TonicInterceptor::Signature(gen_tonic_signature_interceptor()), + ) + .await?; + for payload in payloads { + let canonical = canonical_scoped_dirty_usage_request(&payload).map_err(|err| Error::other(err.to_string()))?; 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"))?; + .map_err(|_| Error::other("scoped dirty usage acknowledgement 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 + || !response.supported { - return Err(Error::other("scoped dirty usage capability response does not match request")); + return Err(Error::other("scoped dirty usage acknowledgement response does not match request")); } - Ok(response.supported) } - .await, - ) - .await + Ok(()) + }; + let result = match timeout(SCANNER_SCOPED_DIRTY_USAGE_STAGE_TIMEOUT, ack_attempt).await { + Ok(result) => self.finalize_result(result).await, + Err(_) => { + self.prepare_retry_with_timeout(SCANNER_SCOPED_DIRTY_USAGE_STAGE_TIMEOUT) + .await; + Err(Error::other("scoped dirty usage acknowledgement deadline elapsed")) + } + }; + + match result { + Ok(()) => { + let activity = self.scanner_scoped_dirty_usage_activity_confirmation().await?; + if activity.instance_id == instance_id { + Ok(activity) + } else { + Err(Error::other( + "scoped dirty usage acknowledgement peer restarted before activity confirmation", + )) + } + } + Err(err) => { + if Self::is_network_like_error(&err) { + self.prepare_retry_with_timeout(SCANNER_SCOPED_DIRTY_USAGE_STAGE_TIMEOUT) + .await; + } + match self.scanner_scoped_dirty_usage_activity_confirmation().await { + Ok(activity) if scanner_scoped_dirty_usage_ack_reconciled(&activity, &instance_id) => Ok(activity), + _ => Err(err), + } + } + } + } + + async fn scanner_scoped_dirty_usage_activity_confirmation(&self) -> Result { + timeout(SCANNER_SCOPED_DIRTY_USAGE_STAGE_TIMEOUT, self.scanner_activity()) + .await + .map_err(|_| Error::other("scoped dirty usage activity confirmation timed out"))? } pub async fn acknowledge_scanner_dirty_usage(&self, instance_id: String, generation: u64) -> Result { @@ -2845,16 +3027,79 @@ mod tests { rustfs_protos::proto_gen::node_service::ScannerDirtyUsageBucket { bucket: "archive".to_string(), generation: 3, + bucket_incarnation: Uuid::from_u128(0x11111111111111111111111111111111).as_bytes().to_vec().into(), }, rustfs_protos::proto_gen::node_service::ScannerDirtyUsageBucket { bucket: "photos".to_string(), generation: 7, + bucket_incarnation: Uuid::from_u128(0x22222222222222222222222222222222).as_bytes().to_vec().into(), }, ], response_proof: b"proof".to_vec().into(), + owner_id: "33333333-3333-3333-3333-333333333333".to_string(), } } + #[test] + fn scanner_scoped_dirty_usage_ack_payloads_split_at_protocol_limit() { + use rustfs_protos::scoped_dirty_usage::{SCOPED_DIRTY_USAGE_MAX_ENTRIES, canonical_scoped_dirty_usage_request}; + + let entries = (0..=SCOPED_DIRTY_USAGE_MAX_ENTRIES) + .map(|index| ScannerScopedDirtyUsageAckEntry { + bucket: format!("bucket-{index:02}"), + bucket_incarnation: Uuid::from_u128(0x11111111111111111111111111111111), + generation: 9, + }) + .collect::>(); + + let payloads = scanner_scoped_dirty_usage_ack_payloads( + "33333333-3333-3333-3333-333333333333".to_string(), + "0123456789abcdef0123456789abcdef".to_string(), + false, + entries, + ) + .expect("33 entries should split into valid scoped dirty usage requests"); + + assert_eq!(payloads.len(), 2); + assert_eq!(payloads[0].entries.len(), SCOPED_DIRTY_USAGE_MAX_ENTRIES as usize); + assert_eq!(payloads[1].entries.len(), 1); + assert_eq!(payloads[0].entries.first().map(|entry| entry.bucket.as_str()), Some("bucket-00")); + assert_eq!(payloads[0].entries.last().map(|entry| entry.bucket.as_str()), Some("bucket-31")); + assert_eq!(payloads[1].entries.first().map(|entry| entry.bucket.as_str()), Some("bucket-32")); + for payload in payloads { + canonical_scoped_dirty_usage_request(&payload).expect("each split scoped ACK payload should be canonical"); + } + } + + #[test] + fn scanner_scoped_dirty_usage_ack_reconciliation_requires_same_clean_instance() { + let activity = |instance_id: &str, pending| ScannerPeerActivity { + instance_id: instance_id.to_string(), + namespace_generation: 1, + maintenance_generation: 1, + protocol_version: SCANNER_ACTIVITY_PROTOCOL_VERSION, + topology_digest: Some([1; 32]), + data_movement_active: Some(false), + dirty_usage_generation: Some(9), + dirty_usage_pending: pending, + movement_generation: Some(1), + publication_blocked: Some(false), + }; + + assert!(scanner_scoped_dirty_usage_ack_reconciled( + &activity("0123456789abcdef0123456789abcdef", Some(false)), + "0123456789abcdef0123456789abcdef" + )); + assert!(!scanner_scoped_dirty_usage_ack_reconciled( + &activity("0123456789abcdef0123456789abcdef", Some(true)), + "0123456789abcdef0123456789abcdef" + )); + assert!(!scanner_scoped_dirty_usage_ack_reconciled( + &activity("fedcba9876543210fedcba9876543210", Some(false)), + "0123456789abcdef0123456789abcdef" + )); + } + #[test] fn scanner_dirty_usage_snapshot_requires_a_complete_authenticated_ordered_view() { let decoded = decode_test_scanner_dirty_usage_snapshot(test_scanner_dirty_usage_snapshot_response()) @@ -2863,9 +3108,18 @@ mod tests { assert_eq!(decoded.generation, 7); assert_eq!(decoded.pending_bucket_count, 2); assert_eq!(decoded.protocol_version, SCANNER_DIRTY_USAGE_SNAPSHOT_PROTOCOL_VERSION); + assert_eq!(decoded.owner_id, "33333333-3333-3333-3333-333333333333"); assert!(decoded.complete); - assert_eq!(decoded.buckets.get("archive"), Some(&3)); - assert_eq!(decoded.buckets.get("photos"), Some(&7)); + assert_eq!( + decoded.buckets.get("archive").map(|bucket| bucket.bucket_incarnation), + Some(Uuid::from_u128(0x11111111111111111111111111111111)) + ); + assert_eq!(decoded.buckets.get("archive").map(|bucket| bucket.generation), Some(3)); + assert_eq!( + decoded.buckets.get("photos").map(|bucket| bucket.bucket_incarnation), + Some(Uuid::from_u128(0x22222222222222222222222222222222)) + ); + assert_eq!(decoded.buckets.get("photos").map(|bucket| bucket.generation), Some(7)); let overflow_count = u64::try_from(SCANNER_DIRTY_USAGE_SNAPSHOT_MAX_ENTRIES + 1).expect("the test snapshot entry limit should fit in u64"); @@ -2916,6 +3170,14 @@ mod tests { empty_bucket.buckets[0].bucket.clear(); cases.push((empty_bucket, "empty bucket name")); + let mut invalid_owner = test_scanner_dirty_usage_snapshot_response(); + invalid_owner.owner_id.clear(); + cases.push((invalid_owner, "snapshot owner")); + + let mut invalid_incarnation = test_scanner_dirty_usage_snapshot_response(); + invalid_incarnation.buckets[0].bucket_incarnation = Uuid::nil().as_bytes().to_vec().into(); + cases.push((invalid_incarnation, "bucket incarnation")); + let mut partial = test_scanner_dirty_usage_snapshot_response(); partial.complete = false; cases.push((partial, "entry-limit overflow")); @@ -2928,6 +3190,7 @@ mod tests { .map(|index| rustfs_protos::proto_gen::node_service::ScannerDirtyUsageBucket { bucket: format!("bucket-{index:04}"), generation: 1, + bucket_incarnation: Uuid::from_u128(0x11111111111111111111111111111111).as_bytes().to_vec().into(), }) .collect(), ..test_scanner_dirty_usage_snapshot_response() diff --git a/crates/ecstore/src/services/notification_sys.rs b/crates/ecstore/src/services/notification_sys.rs index 7769ea1b1..a79f6e858 100644 --- a/crates/ecstore/src/services/notification_sys.rs +++ b/crates/ecstore/src/services/notification_sys.rs @@ -14,7 +14,8 @@ use crate::bucket::lifecycle::tier_last_day_stats::DailyAllTierStats; use crate::cluster::rpc::{ - PeerRestClient, ScannerPeerActivity, ScannerPeerDirtyUsageSnapshot, ScannerPublicationLease, TierConfigReloadOutcome, + PeerRestClient, ScannerDirtyUsageAcknowledgement, ScannerPeerActivity, ScannerPeerDirtyUsageSnapshot, + ScannerPublicationLease, TierConfigReloadOutcome, }; use crate::diagnostics::admin_server_info::get_commit_id; use crate::disk::DiskAPI; @@ -2552,11 +2553,70 @@ impl NotificationSys { Ok(snapshots) } - pub async fn acknowledge_scanner_dirty_usage(&self, acknowledgements: Vec<(String, String, u64)>) -> Result { + pub async fn scanner_scoped_dirty_usage_capabilities( + &self, + acknowledgements: Vec, + ) -> Result { 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}"))); + for acknowledgement in acknowledgements { + let host = match &acknowledgement { + ScannerDirtyUsageAcknowledgement::Scoped { host, .. } => host.clone(), + ScannerDirtyUsageAcknowledgement::Generation { .. } => { + return Err(Error::other("scanner scoped dirty usage capability requires scoped acknowledgements")); + } + }; + if by_host.insert(host.clone(), acknowledgement).is_some() { + return Err(Error::other("duplicate scanner dirty usage acknowledgement target")); + } + } + + let clients = self + .peer_clients + .iter() + .flatten() + .map(|client| (client.grid_host.clone(), client.clone())) + .collect::>(); + let mut futures = Vec::with_capacity(by_host.len()); + for (host, acknowledgement) in by_host { + let Some(client) = clients.get(&host).cloned() else { + return Err(Error::other("scanner scoped dirty usage capability failed: peer is not reachable")); + }; + futures.push(async move { + let ScannerDirtyUsageAcknowledgement::Scoped { + owner_id, + instance_id, + entries, + .. + } = acknowledgement + else { + unreachable!("scoped acknowledgement was validated before probing"); + }; + timeout( + SCANNER_ACTIVITY_PROBE_TIMEOUT, + client.scanner_scoped_dirty_usage_capability(owner_id, instance_id, entries), + ) + .await + .map_err(|_| Error::other("scanner scoped dirty usage capability timed out"))? + }); + } + + for result in join_all(futures).await { + if !result? { + return Ok(false); + } + } + Ok(true) + } + + pub async fn acknowledge_scanner_dirty_usage(&self, acknowledgements: Vec) -> Result { + let mut by_host = HashMap::with_capacity(acknowledgements.len()); + for acknowledgement in acknowledgements { + let host = match &acknowledgement { + ScannerDirtyUsageAcknowledgement::Generation { host, .. } + | ScannerDirtyUsageAcknowledgement::Scoped { host, .. } => host.clone(), + }; + if by_host.insert(host.clone(), acknowledgement).is_some() { + return Err(Error::other("duplicate scanner dirty usage acknowledgement target")); } } @@ -2568,18 +2628,34 @@ impl NotificationSys { .collect::>(); let mut failures = Vec::new(); let mut futures = Vec::with_capacity(by_host.len()); - for (host, (instance_id, generation)) in by_host { + for (host, acknowledgement) 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; + let result = match acknowledgement { + ScannerDirtyUsageAcknowledgement::Generation { + instance_id, generation, .. + } => { + scanner_activity_with_timeout( + SCANNER_ACTIVITY_PROBE_TIMEOUT, + &host, + client.acknowledge_scanner_dirty_usage(instance_id, generation), + ) + .await + } + ScannerDirtyUsageAcknowledgement::Scoped { + owner_id, + instance_id, + entries, + .. + } => { + client + .acknowledge_scanner_scoped_dirty_usage(owner_id, instance_id, entries) + .await + } + }; (host, result) }); } @@ -4744,15 +4820,28 @@ mod tests { peer_topology_hosts: Vec::new(), }; let missing = sys - .acknowledge_scanner_dirty_usage(vec![("peer-1".to_string(), "0123456789abcdef0123456789abcdef".to_string(), 7)]) + .acknowledge_scanner_dirty_usage(vec![ScannerDirtyUsageAcknowledgement::Generation { + host: "peer-1".to_string(), + instance_id: "0123456789abcdef0123456789abcdef".to_string(), + generation: 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), + ScannerDirtyUsageAcknowledgement::Generation { + host: "peer-1".to_string(), + instance_id: "0123456789abcdef0123456789abcdef".to_string(), + generation: 7, + }, + ScannerDirtyUsageAcknowledgement::Scoped { + host: "peer-1".to_string(), + owner_id: "11111111-1111-1111-1111-111111111111".to_string(), + instance_id: "0123456789abcdef0123456789abcdef".to_string(), + entries: Vec::new(), + }, ]) .await .expect_err("duplicate acknowledgement targets must be rejected"); diff --git a/crates/protos/src/generated/proto_gen/node_service.rs b/crates/protos/src/generated/proto_gen/node_service.rs index fe04d93ab..0ee5014ad 100644 --- a/crates/protos/src/generated/proto_gen/node_service.rs +++ b/crates/protos/src/generated/proto_gen/node_service.rs @@ -1257,6 +1257,8 @@ pub struct ScannerDirtyUsageBucket { pub bucket: ::prost::alloc::string::String, #[prost(uint64, tag = "2")] pub generation: u64, + #[prost(bytes = "bytes", tag = "3")] + pub bucket_incarnation: ::prost::bytes::Bytes, } #[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] pub struct ScannerDirtyUsageSnapshotRequest { @@ -1282,6 +1284,8 @@ pub struct ScannerDirtyUsageSnapshotResponse { pub buckets: ::prost::alloc::vec::Vec, #[prost(bytes = "bytes", tag = "7")] pub response_proof: ::prost::bytes::Bytes, + #[prost(string, tag = "8")] + pub owner_id: ::prost::alloc::string::String, } /// Receiver-only protocol. Producers must retain whole-cycle ACK until they /// have a durable per-bucket publication proof. diff --git a/crates/protos/src/lib.rs b/crates/protos/src/lib.rs index 47b7811e1..f0d1a1255 100644 --- a/crates/protos/src/lib.rs +++ b/crates/protos/src/lib.rs @@ -575,11 +575,13 @@ pub fn canonical_scanner_dirty_usage_snapshot_response_body( body.push_u64(response.generation); body.push_u64(response.pending_bucket_count); body.push_u32(response.protocol_version); + body.push_str(&response.owner_id)?; body.push_bool(response.complete); body.push_count(response.buckets.len())?; for bucket in &response.buckets { body.push_str(&bucket.bucket)?; body.push_u64(bucket.generation); + body.push_bytes(bucket.bucket_incarnation.as_ref())?; } Ok(body.finish()) } @@ -1842,13 +1844,16 @@ mod scanner_activity_tests { ScannerDirtyUsageBucket { bucket: "archive".to_string(), generation: 3, + bucket_incarnation: vec![1; 16].into(), }, ScannerDirtyUsageBucket { bucket: "photos".to_string(), generation: 7, + bucket_incarnation: vec![2; 16].into(), }, ], response_proof: vec![9; 32].into(), + owner_id: "11111111-1111-1111-1111-111111111111".to_string(), }; let baseline = canonical_scanner_dirty_usage_snapshot_response_body(&[1; 16], &response) .expect("scanner dirty usage snapshot response should encode"); @@ -1865,6 +1870,9 @@ mod scanner_activity_tests { let mut protocol = response.clone(); protocol.protocol_version = 2; variants.push(protocol); + let mut owner = response.clone(); + owner.owner_id = "22222222-2222-2222-2222-222222222222".to_string(); + variants.push(owner); let mut complete = response.clone(); complete.complete = false; variants.push(complete); @@ -1874,6 +1882,9 @@ mod scanner_activity_tests { let mut bucket_generation = response.clone(); bucket_generation.buckets[0].generation = 4; variants.push(bucket_generation); + let mut bucket_incarnation = response.clone(); + bucket_incarnation.buckets[0].bucket_incarnation = vec![3; 16].into(); + variants.push(bucket_incarnation); let mut bucket_order = response.clone(); bucket_order.buckets.reverse(); variants.push(bucket_order); diff --git a/crates/protos/src/node.proto b/crates/protos/src/node.proto index abc030adc..a35797497 100644 --- a/crates/protos/src/node.proto +++ b/crates/protos/src/node.proto @@ -885,6 +885,7 @@ message ScannerActivityResponse { message ScannerDirtyUsageBucket { string bucket = 1; uint64 generation = 2; + bytes bucket_incarnation = 3; } message ScannerDirtyUsageSnapshotRequest { @@ -901,6 +902,7 @@ message ScannerDirtyUsageSnapshotResponse { bool complete = 5; repeated ScannerDirtyUsageBucket buckets = 6; bytes response_proof = 7; + string owner_id = 8; } // Receiver-only protocol. Producers must retain whole-cycle ACK until they diff --git a/crates/scanner/src/lib.rs b/crates/scanner/src/lib.rs index 7cc6d22e8..696b37332 100644 --- a/crates/scanner/src/lib.rs +++ b/crates/scanner/src/lib.rs @@ -100,6 +100,7 @@ pub use storage_api::ScannerReplicationConfig as ReplicationConfig; pub use storage_api::scan::{ SCANNER_ACTIVITY_PROTOCOL_VERSION, SCANNER_ACTIVITY_V6_PROTOCOL_VERSION, SCANNER_DIRTY_USAGE_SNAPSHOT_MAX_ENTRIES, SCANNER_DIRTY_USAGE_SNAPSHOT_PROTOCOL_VERSION, SCANNER_DIRTY_USAGE_SNAPSHOT_RPC_MAX_MESSAGE_SIZE, + SCANNER_SCOPED_DIRTY_USAGE_ACK_MAX_ENTRIES, }; pub use workload_admission::set_scanner_workload_admission_snapshot_provider; diff --git a/crates/scanner/src/scanner.rs b/crates/scanner/src/scanner.rs index afd54825d..8cc4c481e 100644 --- a/crates/scanner/src/scanner.rs +++ b/crates/scanner/src/scanner.rs @@ -2216,10 +2216,7 @@ where false } else if let Some(notification_system) = storeapi.scanner_notification_system() { let acknowledgement_count = remote_dirty_usage_acknowledgements.len(); - let acknowledgements = remote_dirty_usage_acknowledgements - .into_iter() - .map(|acknowledgement| (acknowledgement.host, acknowledgement.instance_id, acknowledgement.generation)) - .collect(); + let acknowledgements = remote_dirty_usage_acknowledgements.into_iter().map(Into::into).collect(); remote_dirty_usage_acknowledgement_pending( cycle_info.current, acknowledgement_count, @@ -3578,9 +3575,10 @@ use usage_store::*; pub use activity::scanner_topology_digest; pub(crate) use activity::{ - ScannerActivitySnapshot, ScannerDirtyUsageAcknowledgement, probe_scanner_activity, scanner_activity_allows_usage_publication, - scanner_activity_dirty_usage_state_for_host, scanner_activity_publication_lease_targets, scanner_activity_snapshot_digest, - scanner_activity_structural_digest, scanner_dirty_usage_acknowledgements, + ScannerActivitySnapshot, ScannerDirtyUsageAcknowledgement, ScannerDirtyUsageAcknowledgementKind, probe_scanner_activity, + scanner_activity_allows_usage_publication, scanner_activity_dirty_usage_state_for_host, + scanner_activity_publication_lease_targets, scanner_activity_snapshot_digest, scanner_activity_structural_digest, + scanner_dirty_usage_acknowledgements, }; pub(crate) use activity::{ScannerCycleOutcome, scanner_cycle_outcome_with_pending_maintenance}; pub use backlog::{ diff --git a/crates/scanner/src/scanner/activity.rs b/crates/scanner/src/scanner/activity.rs index 58c79f3ed..1b2761bb1 100644 --- a/crates/scanner/src/scanner/activity.rs +++ b/crates/scanner/src/scanner/activity.rs @@ -445,7 +445,34 @@ pub(crate) type ScannerActivitySnapshot = BTreeMap; pub(crate) struct ScannerDirtyUsageAcknowledgement { pub(crate) host: String, pub(crate) instance_id: String, - pub(crate) generation: u64, + pub(crate) kind: ScannerDirtyUsageAcknowledgementKind, +} + +#[derive(Clone, Debug, PartialEq, Eq)] +pub(crate) enum ScannerDirtyUsageAcknowledgementKind { + Generation(u64), + Scoped { + owner_id: String, + entries: Vec, + }, +} + +impl From for crate::storage_api::EcstoreScannerDirtyUsageAcknowledgement { + fn from(acknowledgement: ScannerDirtyUsageAcknowledgement) -> Self { + match acknowledgement.kind { + ScannerDirtyUsageAcknowledgementKind::Generation(generation) => Self::Generation { + host: acknowledgement.host, + instance_id: acknowledgement.instance_id, + generation, + }, + ScannerDirtyUsageAcknowledgementKind::Scoped { owner_id, entries } => Self::Scoped { + host: acknowledgement.host, + owner_id, + instance_id: acknowledgement.instance_id, + entries, + }, + } + } } #[derive(Clone, Copy, Debug, PartialEq, Eq)] @@ -974,7 +1001,7 @@ pub(crate) fn scanner_dirty_usage_acknowledgements(snapshot: &ScannerActivitySna .map(|(host, activity)| ScannerDirtyUsageAcknowledgement { host: host.clone(), instance_id: activity.instance_id.clone(), - generation: activity.dirty_usage_generation, + kind: ScannerDirtyUsageAcknowledgementKind::Generation(activity.dirty_usage_generation), }) .collect() } diff --git a/crates/scanner/src/scanner/tests.rs b/crates/scanner/src/scanner/tests.rs index cfcd37fe9..42a128cd0 100644 --- a/crates/scanner/src/scanner/tests.rs +++ b/crates/scanner/src/scanner/tests.rs @@ -7433,7 +7433,7 @@ fn finalizing_a_saved_enum_without_proof_keeps_dirty_pending() { let remote_acknowledgement = ScannerDirtyUsageAcknowledgement { host: "node-2".to_string(), instance_id: "0123456789abcdef0123456789abcdef".to_string(), - generation: 11, + kind: ScannerDirtyUsageAcknowledgementKind::Generation(11), }; let unsaved = crate::scanner_io::ScannerCycleResult::new(ScannerCycleStatus::Complete, Some(dirty_snapshot.clone())) .with_remote_dirty_usage_acknowledgements(vec![remote_acknowledgement.clone()]); @@ -8982,7 +8982,7 @@ fn post_lease_activity_proof_rejects_a_put_tail_that_finished_before_lease_acqui ScannerDirtyUsageAcknowledgement { host: "node-2".to_string(), instance_id: "epoch-a".to_string(), - generation: 5, + kind: ScannerDirtyUsageAcknowledgementKind::Generation(5), }, ]); let (outcome, _, acknowledgements) = finalize_scanner_cycle_result( @@ -9142,7 +9142,7 @@ fn scanner_dirty_usage_acknowledgements_exclude_local_and_clean_nodes() { vec![ScannerDirtyUsageAcknowledgement { host: "node-3".to_string(), instance_id: "epoch-dirty".to_string(), - generation: 11, + kind: ScannerDirtyUsageAcknowledgementKind::Generation(11), }] ); } diff --git a/crates/scanner/src/scanner/tests/scoped_ack_publication.rs b/crates/scanner/src/scanner/tests/scoped_ack_publication.rs index 39deae1e5..2fa58fff2 100644 --- a/crates/scanner/src/scanner/tests/scoped_ack_publication.rs +++ b/crates/scanner/src/scanner/tests/scoped_ack_publication.rs @@ -463,7 +463,7 @@ async fn scoped_ack_publication_rejects_builder_mutation_after_real_root_publish "remote_ack_target" => scan.with_remote_dirty_usage_acknowledgements(vec![ScannerDirtyUsageAcknowledgement { host: "proof-peer:9000".to_string(), instance_id: crate::scanner_activity_epoch().to_string(), - generation: changed_generation, + kind: crate::scanner::ScannerDirtyUsageAcknowledgementKind::Generation(changed_generation), }]), "publication_epoch" => scan.with_publication_epoch(Some(changed_epoch)), "remote_lease_targets" => scan.with_remote_publication_lease_targets(vec![( diff --git a/crates/scanner/src/scanner_io.rs b/crates/scanner/src/scanner_io.rs index b51e5d2f8..94ffd5011 100644 --- a/crates/scanner/src/scanner_io.rs +++ b/crates/scanner/src/scanner_io.rs @@ -148,19 +148,26 @@ struct ScannerPeerDirtyUsageExpectation { pending: bool, } -fn verified_remote_dirty_usage_buckets( +#[derive(Debug, PartialEq, Eq)] +struct VerifiedRemoteDirtyUsage { + dirty_buckets: HashSet, + acknowledgements: Vec, +} + +fn verified_remote_dirty_usage( expected_peers: &HashMap, peer_snapshots: Vec<(String, EcstoreScannerPeerDirtyUsageSnapshot)>, -) -> Option> { +) -> Option { if expected_peers.is_empty() || peer_snapshots.len() != expected_peers.len() { return None; } let mut received_peers = HashSet::with_capacity(peer_snapshots.len()); let mut dirty_buckets = HashSet::new(); + let mut acknowledgements = Vec::new(); for (host, snapshot) in peer_snapshots { let expected = expected_peers.get(&host)?; - if !received_peers.insert(host) + if !received_peers.insert(host.clone()) || snapshot.instance_id != expected.instance_id || snapshot.generation != expected.generation || snapshot.generation == u64::MAX @@ -171,10 +178,44 @@ fn verified_remote_dirty_usage_buckets( { return None; } - dirty_buckets.extend(snapshot.buckets.into_keys()); + let entries = snapshot + .buckets + .iter() + .map(|(bucket, state)| crate::storage_api::EcstoreScannerScopedDirtyUsageAckEntry { + bucket: bucket.clone(), + bucket_incarnation: state.bucket_incarnation, + generation: state.generation, + }) + .collect::>(); + dirty_buckets.extend(snapshot.buckets.keys().cloned()); + if !entries.is_empty() { + acknowledgements.push(crate::scanner::ScannerDirtyUsageAcknowledgement { + host, + instance_id: snapshot.instance_id, + kind: crate::scanner::ScannerDirtyUsageAcknowledgementKind::Scoped { + owner_id: snapshot.owner_id, + entries, + }, + }); + } } - (received_peers.len() == expected_peers.len()).then_some(dirty_buckets) + (received_peers.len() == expected_peers.len()).then_some(VerifiedRemoteDirtyUsage { + dirty_buckets, + acknowledgements, + }) +} + +fn scanner_scoped_dirty_usage_ack_exceeds_cost_threshold( + acknowledgements: &[crate::scanner::ScannerDirtyUsageAcknowledgement], +) -> bool { + acknowledgements.iter().any(|acknowledgement| { + matches!( + &acknowledgement.kind, + crate::scanner::ScannerDirtyUsageAcknowledgementKind::Scoped { entries, .. } + if entries.len() > crate::SCANNER_SCOPED_DIRTY_USAGE_ACK_MAX_ENTRIES + ) + }) } fn complete_scanner_cache_snapshot_plan_digest( diff --git a/crates/scanner/src/scanner_io/io_cycle.rs b/crates/scanner/src/scanner_io/io_cycle.rs index d4094a311..eb09af055 100644 --- a/crates/scanner/src/scanner_io/io_cycle.rs +++ b/crates/scanner/src/scanner_io/io_cycle.rs @@ -107,23 +107,32 @@ struct ScannerBucketScopeResolution<'a> { requires_full_scan: bool, } +struct ScannerBucketScopeResolutionResult { + scope: ScannerBucketScanScope, + remote_dirty_usage_acknowledgements: Vec, +} + async fn resolve_scanner_bucket_scan_scope( store: &S, distributed: bool, resolution: ScannerBucketScopeResolution<'_>, -) -> ScannerBucketScanScope +) -> ScannerBucketScopeResolutionResult where S: ScannerStorage, { + let default_result = |scope: ScannerBucketScanScope| ScannerBucketScopeResolutionResult { + scope, + remote_dirty_usage_acknowledgements: Vec::new(), + }; if resolution.requires_full_scan { - return ScannerBucketScanScope::default(); + return default_result(ScannerBucketScanScope::default()); } if !resolution.requested_scope.is_default() || !resolution.dirty_usage_snapshot.covers_all_pending || resolution.dirty_usage_snapshot.generation == u64::MAX || resolution.dirty_usage_snapshot.buckets.len() > crate::SCANNER_DIRTY_USAGE_SNAPSHOT_MAX_ENTRIES { - return resolution.requested_scope; + return default_result(resolution.requested_scope); } let mut dirty_buckets = resolution @@ -134,10 +143,10 @@ where .collect::>(); if distributed { let Some(notification_system) = store.scanner_notification_system() else { - return resolution.requested_scope; + return default_result(resolution.requested_scope); }; let Ok(peer_snapshots) = notification_system.scanner_dirty_usage_snapshots().await else { - return resolution.requested_scope; + return default_result(resolution.requested_scope); }; let mut expected_peers = HashMap::new(); for (host, lease_instance_id, _) in crate::scanner::scanner_activity_publication_lease_targets(resolution.activity_before) @@ -145,10 +154,10 @@ where let Some((activity_instance_id, generation, pending)) = crate::scanner::scanner_activity_dirty_usage_state_for_host(resolution.activity_before, &host) else { - return resolution.requested_scope; + return default_result(resolution.requested_scope); }; if activity_instance_id != lease_instance_id || expected_peers.contains_key(&host) { - return resolution.requested_scope; + return default_result(resolution.requested_scope); } expected_peers.insert( host, @@ -159,19 +168,76 @@ where }, ); } - let Some(remote_dirty_buckets) = verified_remote_dirty_usage_buckets(&expected_peers, peer_snapshots) else { - return resolution.requested_scope; + let Some(remote_dirty_usage) = verified_remote_dirty_usage(&expected_peers, peer_snapshots) else { + return default_result(resolution.requested_scope); + }; + dirty_buckets.extend(remote_dirty_usage.dirty_buckets); + let scope = scoped_scan_scope_from_dirty_buckets( + resolution.requested_scope, + dirty_buckets, + true, + resolution.all_buckets, + resolution.baseline_proof, + ); + if scope.is_default() { + return default_result(scope); + } + let Some(selected_buckets) = scope.selected_buckets.as_ref() else { + return default_result(scope); + }; + let mut scoped_acknowledgements = Vec::with_capacity(remote_dirty_usage.acknowledgements.len()); + for acknowledgement in remote_dirty_usage.acknowledgements { + let crate::scanner::ScannerDirtyUsageAcknowledgement { + host, + instance_id, + kind: crate::scanner::ScannerDirtyUsageAcknowledgementKind::Scoped { owner_id, entries }, + } = acknowledgement + else { + return default_result(scope); + }; + let entries = entries + .into_iter() + .filter(|entry| selected_buckets.contains(&entry.bucket)) + .collect::>(); + if !entries.is_empty() { + scoped_acknowledgements.push(crate::scanner::ScannerDirtyUsageAcknowledgement { + host, + instance_id, + kind: crate::scanner::ScannerDirtyUsageAcknowledgementKind::Scoped { owner_id, entries }, + }); + } + } + if super::scanner_scoped_dirty_usage_ack_exceeds_cost_threshold(&scoped_acknowledgements) { + return default_result(ScannerBucketScanScope::default()); + } + if !scoped_acknowledgements.is_empty() { + let capability_acknowledgements = scoped_acknowledgements + .clone() + .into_iter() + .map(Into::into) + .collect::>(); + if !matches!( + notification_system + .scanner_scoped_dirty_usage_capabilities(capability_acknowledgements) + .await, + Ok(true) + ) { + return default_result(ScannerBucketScanScope::default()); + } + } + return ScannerBucketScopeResolutionResult { + scope, + remote_dirty_usage_acknowledgements: scoped_acknowledgements, }; - dirty_buckets.extend(remote_dirty_buckets); } - scoped_scan_scope_from_dirty_buckets( + default_result(scoped_scan_scope_from_dirty_buckets( resolution.requested_scope, dirty_buckets, true, resolution.all_buckets, resolution.baseline_proof, - ) + )) } pub(crate) async fn nsscanner_with_storage_status_scoped(store: &S, request: ScannerCycleRequest) -> Result @@ -294,7 +360,7 @@ where let bucket_coverage_digest = scanner_bucket_plan_digest(&all_buckets, activity_digest); let execution_digest = scanner_bucket_work_digest(bucket_coverage_digest, scan_mode, requires_full_scan); let dirty_usage_snapshot = Arc::new(snapshot_dirty_usage_buckets(&all_buckets, dirty_generation_before_bucket_list)); - let scan_scope = resolve_scanner_bucket_scan_scope( + let scope_resolution = resolve_scanner_bucket_scan_scope( store, distributed, ScannerBucketScopeResolution { @@ -314,6 +380,8 @@ where }, ) .await; + let remote_dirty_usage_acknowledgements = scope_resolution.remote_dirty_usage_acknowledgements; + let scan_scope = scope_resolution.scope; #[cfg(test)] if let Some(observer) = resolved_scope_observer { let _ = observer.send(scan_scope.clone()); @@ -674,11 +742,14 @@ where if cycle_status == ScannerCycleStatus::Complete { complete_tier_registry_cycle(want_cycle, leader_epoch); } - let remote_dirty_usage_acknowledgements = if cycle_status == ScannerCycleStatus::Complete { - crate::scanner::scanner_dirty_usage_acknowledgements(&activity_before) - } else { - Vec::new() - }; + let remote_dirty_usage_acknowledgements = + if cycle_status == ScannerCycleStatus::Complete && !remote_dirty_usage_acknowledgements.is_empty() { + remote_dirty_usage_acknowledgements + } else if cycle_status == ScannerCycleStatus::Complete && scan_scope.is_default() { + crate::scanner::scanner_dirty_usage_acknowledgements(&activity_before) + } else { + Vec::new() + }; Ok(ScannerCycleResult::new(cycle_status, dirty_usage_clear) .with_publication_epoch(publication_epoch) .with_activity_digest(activity_digest) diff --git a/crates/scanner/src/scanner_io/tests.rs b/crates/scanner/src/scanner_io/tests.rs index 1d58dc88e..ab7b3bbfd 100644 --- a/crates/scanner/src/scanner_io/tests.rs +++ b/crates/scanner/src/scanner_io/tests.rs @@ -1561,6 +1561,7 @@ fn peer_dirty_usage_snapshot( buckets: &[(&str, u64)], ) -> EcstoreScannerPeerDirtyUsageSnapshot { EcstoreScannerPeerDirtyUsageSnapshot { + owner_id: uuid::Uuid::from_u128(0x11111111111111111111111111111111).to_string(), instance_id: instance_id.to_string(), generation, pending_bucket_count: u64::try_from(buckets.len()).expect("test bucket count should fit"), @@ -1568,7 +1569,15 @@ fn peer_dirty_usage_snapshot( complete, buckets: buckets .iter() - .map(|(bucket, generation)| ((*bucket).to_string(), *generation)) + .map(|(bucket, generation)| { + ( + (*bucket).to_string(), + crate::storage_api::EcstoreScannerPeerDirtyUsageBucket { + bucket_incarnation: uuid::Uuid::from_u128(0x22222222222222222222222222222222), + generation: *generation, + }, + ) + }) .collect(), } } @@ -1595,7 +1604,7 @@ fn verified_remote_dirty_usage_buckets_merges_only_complete_current_snapshots() ]); assert_eq!( - verified_remote_dirty_usage_buckets( + verified_remote_dirty_usage( &expected_peers, vec![ ( @@ -1608,10 +1617,63 @@ fn verified_remote_dirty_usage_buckets_merges_only_complete_current_snapshots() ), ], ), - Some(HashSet::from(["photos".to_string(), "archive".to_string()])) + Some(VerifiedRemoteDirtyUsage { + dirty_buckets: HashSet::from(["photos".to_string(), "archive".to_string()]), + acknowledgements: vec![ + crate::scanner::ScannerDirtyUsageAcknowledgement { + host: "node-a:9000".to_string(), + instance_id: "instance-a".to_string(), + kind: crate::scanner::ScannerDirtyUsageAcknowledgementKind::Scoped { + owner_id: uuid::Uuid::from_u128(0x11111111111111111111111111111111).to_string(), + entries: vec![crate::storage_api::EcstoreScannerScopedDirtyUsageAckEntry { + bucket: "photos".to_string(), + bucket_incarnation: uuid::Uuid::from_u128(0x22222222222222222222222222222222), + generation: 7, + }], + }, + }, + crate::scanner::ScannerDirtyUsageAcknowledgement { + host: "node-b:9000".to_string(), + instance_id: "instance-b".to_string(), + kind: crate::scanner::ScannerDirtyUsageAcknowledgementKind::Scoped { + owner_id: uuid::Uuid::from_u128(0x11111111111111111111111111111111).to_string(), + entries: vec![crate::storage_api::EcstoreScannerScopedDirtyUsageAckEntry { + bucket: "archive".to_string(), + bucket_incarnation: uuid::Uuid::from_u128(0x22222222222222222222222222222222), + generation: 3, + }], + }, + }, + ], + }) ); } +#[test] +fn scanner_scoped_dirty_usage_ack_cost_threshold_is_single_protocol_batch() { + let acknowledgement = |entry_count: usize| crate::scanner::ScannerDirtyUsageAcknowledgement { + host: "node-a:9000".to_string(), + instance_id: "instance-a".to_string(), + kind: crate::scanner::ScannerDirtyUsageAcknowledgementKind::Scoped { + owner_id: uuid::Uuid::from_u128(0x11111111111111111111111111111111).to_string(), + entries: (0..entry_count) + .map(|index| crate::storage_api::EcstoreScannerScopedDirtyUsageAckEntry { + bucket: format!("bucket-{index:02}"), + bucket_incarnation: uuid::Uuid::from_u128(0x22222222222222222222222222222222), + generation: 7, + }) + .collect(), + }, + }; + + assert!(!scanner_scoped_dirty_usage_ack_exceeds_cost_threshold(&[acknowledgement( + crate::SCANNER_SCOPED_DIRTY_USAGE_ACK_MAX_ENTRIES + )])); + assert!(scanner_scoped_dirty_usage_ack_exceeds_cost_threshold(&[acknowledgement( + crate::SCANNER_SCOPED_DIRTY_USAGE_ACK_MAX_ENTRIES + 1 + )])); +} + #[test] fn verified_remote_dirty_usage_buckets_rejects_incomplete_or_stale_peer_state() { let expected_peers = HashMap::from([( @@ -1630,7 +1692,7 @@ fn verified_remote_dirty_usage_buckets_rejects_incomplete_or_stale_peer_state() peer_dirty_usage_snapshot("instance-a", 7, true, &[]), ] { assert!( - verified_remote_dirty_usage_buckets(&expected_peers, vec![("node-a:9000".to_string(), snapshot)]).is_none(), + verified_remote_dirty_usage(&expected_peers, vec![("node-a:9000".to_string(), snapshot)]).is_none(), "incomplete, stale, mismatched, or empty pending peer state must fall back to a full scan" ); } diff --git a/crates/scanner/src/storage_api.rs b/crates/scanner/src/storage_api.rs index 26a1eb935..92515ec7e 100644 --- a/crates/scanner/src/storage_api.rs +++ b/crates/scanner/src/storage_api.rs @@ -103,8 +103,13 @@ pub(crate) use rustfs_ecstore::api::rebalance::{ RebalStatus as EcstoreRebalStatus, RebalanceInfo as EcstoreRebalanceInfo, RebalanceMeta as EcstoreRebalanceMeta, RebalanceStats as EcstoreRebalanceStats, }; +#[cfg(test)] +pub(crate) use rustfs_ecstore::api::rpc::ScannerPeerDirtyUsageBucket as EcstoreScannerPeerDirtyUsageBucket; pub(crate) use rustfs_ecstore::api::rpc::{ - ScannerBucketListing as EcstoreScannerBucketListing, ScannerPeerDirtyUsageSnapshot as EcstoreScannerPeerDirtyUsageSnapshot, + ScannerBucketListing as EcstoreScannerBucketListing, + ScannerDirtyUsageAcknowledgement as EcstoreScannerDirtyUsageAcknowledgement, + ScannerPeerDirtyUsageSnapshot as EcstoreScannerPeerDirtyUsageSnapshot, + ScannerScopedDirtyUsageAckEntry as EcstoreScannerScopedDirtyUsageAckEntry, }; #[cfg(test)] pub(crate) use rustfs_ecstore::api::runtime::InstanceContext as EcstoreInstanceContext; @@ -315,6 +320,7 @@ pub(crate) mod scan { pub use super::storage_contracts::{ SCANNER_ACTIVITY_PROTOCOL_VERSION, SCANNER_ACTIVITY_V6_PROTOCOL_VERSION, SCANNER_DIRTY_USAGE_SNAPSHOT_MAX_ENTRIES, SCANNER_DIRTY_USAGE_SNAPSHOT_PROTOCOL_VERSION, SCANNER_DIRTY_USAGE_SNAPSHOT_RPC_MAX_MESSAGE_SIZE, + SCANNER_SCOPED_DIRTY_USAGE_ACK_MAX_ENTRIES, }; } diff --git a/crates/storage-api/src/lib.rs b/crates/storage-api/src/lib.rs index 2726939f0..5ba4dc186 100644 --- a/crates/storage-api/src/lib.rs +++ b/crates/storage-api/src/lib.rs @@ -56,8 +56,9 @@ pub const SCANNER_ACTIVITY_PREVIOUS_PROTOCOL_VERSION: u32 = 5; /// state is not authenticated by that version. pub const SCANNER_ACTIVITY_V6_PROTOCOL_VERSION: u32 = 6; pub const SCANNER_ACTIVITY_PROTOCOL_VERSION: u32 = 7; -pub const SCANNER_DIRTY_USAGE_SNAPSHOT_PROTOCOL_VERSION: u32 = 1; -pub const SCANNER_DIRTY_USAGE_SNAPSHOT_MAX_ENTRIES: usize = 4096; +pub const SCANNER_DIRTY_USAGE_SNAPSHOT_PROTOCOL_VERSION: u32 = 2; +pub const SCANNER_SCOPED_DIRTY_USAGE_ACK_MAX_ENTRIES: usize = 32; +pub const SCANNER_DIRTY_USAGE_SNAPSHOT_MAX_ENTRIES: usize = SCANNER_SCOPED_DIRTY_USAGE_ACK_MAX_ENTRIES; pub const SCANNER_DIRTY_USAGE_SNAPSHOT_RPC_MAX_MESSAGE_SIZE: usize = 512 * 1024; #[derive(Debug, serde::Deserialize, serde::Serialize)] diff --git a/rustfs/src/server/http.rs b/rustfs/src/server/http.rs index d82d67f20..4365196aa 100644 --- a/rustfs/src/server/http.rs +++ b/rustfs/src/server/http.rs @@ -2418,6 +2418,7 @@ fn get_default_tcp_keepalive() -> TcpKeepalive { mod tests { use super::*; use crate::server::compress::RequestPathCategory; + use crate::storage_api::server::http::ScannerScopedDirtyUsageAckEntry; use bytes::Bytes; use http::Request as HttpRequest; use http::{HeaderMap, StatusCode}; @@ -3454,9 +3455,9 @@ mod tests { .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(), + vec![ScannerScopedDirtyUsageAckEntry { + bucket: "photos".to_string(), + bucket_incarnation: uuid::Uuid::from_u128(0x11111111111111111111111111111111), generation: 8, }], ) diff --git a/rustfs/src/storage/rpc/node_service.rs b/rustfs/src/storage/rpc/node_service.rs index dd44f134d..4bf94c00a 100644 --- a/rustfs/src/storage/rpc/node_service.rs +++ b/rustfs/src/storage/rpc/node_service.rs @@ -303,25 +303,38 @@ fn scanner_activity_response( } } -fn scanner_dirty_usage_snapshot_response( +async fn scanner_dirty_usage_snapshot_response( + store: &ECStore, snapshot: rustfs_scanner::ScannerDirtyUsageSnapshot, -) -> ScannerDirtyUsageSnapshotResponse { - ScannerDirtyUsageSnapshotResponse { +) -> Result { + if store.id.is_nil() { + return Err(Status::failed_precondition("scanner dirty usage snapshot owner is unavailable")); + } + let mut buckets = Vec::with_capacity(snapshot.buckets.len()); + for bucket in snapshot.buckets { + let bucket_incarnation = store + .bucket_incarnation_id_from_disk(&bucket.bucket) + .await + .map_err(|_| Status::failed_precondition("scanner dirty usage bucket incarnation is unavailable"))?; + if bucket_incarnation.is_nil() { + return Err(Status::failed_precondition("scanner dirty usage bucket incarnation is unavailable")); + } + buckets.push(ScannerDirtyUsageBucket { + bucket: bucket.bucket, + generation: bucket.generation, + bucket_incarnation: bucket_incarnation.as_bytes().to_vec().into(), + }); + } + Ok(ScannerDirtyUsageSnapshotResponse { instance_id: rustfs_scanner::scanner_activity_epoch().to_string(), generation: snapshot.generation, pending_bucket_count: snapshot.pending_bucket_count, protocol_version: rustfs_scanner::SCANNER_DIRTY_USAGE_SNAPSHOT_PROTOCOL_VERSION, complete: snapshot.complete, - buckets: snapshot - .buckets - .into_iter() - .map(|bucket| ScannerDirtyUsageBucket { - bucket: bucket.bucket, - generation: bucket.generation, - }) - .collect(), + buckets, response_proof: Bytes::new(), - } + owner_id: store.id.to_string(), + }) } fn scanner_activity_response_v7( @@ -2219,11 +2232,14 @@ impl Node for NodeService { .as_ref() .try_into() .map_err(|_| Status::invalid_argument("scanner dirty usage snapshot challenge must be 16 bytes"))?; - let snapshot = rustfs_scanner::scanner_dirty_usage_snapshot(rustfs_scanner::SCANNER_DIRTY_USAGE_SNAPSHOT_MAX_ENTRIES); + let store = self + .resolve_object_store() + .ok_or_else(|| Status::unavailable("storage layer is not initialized"))?; + let snapshot = rustfs_scanner::scanner_dirty_usage_snapshot(rustfs_scanner::SCANNER_SCOPED_DIRTY_USAGE_ACK_MAX_ENTRIES); if snapshot.generation == u64::MAX { return Err(Status::resource_exhausted("scanner dirty usage generation is exhausted")); } - let mut response = scanner_dirty_usage_snapshot_response(snapshot); + let mut response = scanner_dirty_usage_snapshot_response(&store, snapshot).await?; let canonical = rustfs_protos::canonical_scanner_dirty_usage_snapshot_response_body(&challenge, &response) .map_err(|_| Status::internal("scanner dirty usage snapshot response is too large to authenticate"))?; response.response_proof = sign_tonic_rpc_response_proof(&canonical) @@ -6510,7 +6526,28 @@ mod tests { #[tokio::test] async fn test_scanner_dirty_usage_snapshot_requires_body_bound_auth_and_signs_a_consistent_view() { let _ = rustfs_credentials::set_global_rpc_secret("scanner-dirty-usage-snapshot-test-secret".to_string()); - let service = create_test_node_service(); + let _ = rustfs_credentials::init_global_action_credentials( + Some("TESTROOTACCESSKEY".to_string()), + Some("TESTROOTSECRET123".to_string()), + ); + let temp_dir = tempfile::tempdir().expect("scanner dirty usage snapshot RPC test directory"); + let env = rustfs_test_utils::TestECStoreEnv::builder() + .base_dir(temp_dir.path()) + .build() + .await; + ObjectStore::new(Arc::clone(&env.ecstore)) + .save_iam_config(serde_json::json!({"version": 1}), format!("{}/format.json", *IAM_CONFIG_PREFIX)) + .await + .expect("seed IAM format"); + let iam = rustfs_iam::build_iam_sys(Arc::clone(&env.ecstore)) + .await + .expect("build isolated IAM"); + let context = Arc::new(crate::runtime_sources::AppContext::with_default_interfaces( + Arc::clone(&env.ecstore), + iam, + Arc::new(KmsServiceManager::new()), + )); + let service = make_server_for_context(Some(context)); let unsigned = service .scanner_dirty_usage_snapshot(Request::new(ScannerDirtyUsageSnapshotRequest { challenge: vec![7; 16].into(), @@ -6564,6 +6601,7 @@ mod tests { .into_inner(); assert_eq!(response.instance_id, rustfs_scanner::scanner_activity_epoch()); assert_eq!(response.protocol_version, rustfs_scanner::SCANNER_DIRTY_USAGE_SNAPSHOT_PROTOCOL_VERSION); + assert!(Uuid::parse_str(&response.owner_id).is_ok_and(|owner_id| !owner_id.is_nil())); let bucket_count = u64::try_from(response.buckets.len()).expect("snapshot bucket count should fit in u64"); assert_eq!(response.complete, response.pending_bucket_count == bucket_count); let canonical = rustfs_protos::canonical_scanner_dirty_usage_snapshot_response_body(&challenge, &response) diff --git a/rustfs/src/storage/storage_api.rs b/rustfs/src/storage/storage_api.rs index ce3cca40d..34210287d 100644 --- a/rustfs/src/storage/storage_api.rs +++ b/rustfs/src/storage/storage_api.rs @@ -550,8 +550,8 @@ pub(crate) mod ecstore_rpc { }; #[cfg(test)] pub(crate) use rustfs_ecstore::api::rpc::{ - build_put_file_auth_trailer, gen_signature_headers, gen_tonic_signature_headers, set_tonic_canonical_body_digest, - verify_put_file_capability, verify_tonic_rpc_response_proof, + ScannerScopedDirtyUsageAckEntry, build_put_file_auth_trailer, gen_signature_headers, gen_tonic_signature_headers, + set_tonic_canonical_body_digest, verify_put_file_capability, verify_tonic_rpc_response_proof, }; } @@ -694,6 +694,8 @@ pub(crate) type ServerContextSlot = crate::storage::runtime_sources::ServerConte pub(crate) type LocalPeerS3Client = ecstore_rpc::LocalPeerS3Client; #[cfg(test)] pub(crate) type PeerRestClient = ecstore_rpc::PeerRestClient; +#[cfg(test)] +pub(crate) type ScannerScopedDirtyUsageAckEntry = ecstore_rpc::ScannerScopedDirtyUsageAckEntry; pub(crate) type MetricType = ecstore_metrics::MetricType; pub(crate) type ObjectPartInfo = rustfs_filemeta::ObjectPartInfo; pub(crate) type ObjectLockBlockReason = ecstore_bucket::object_lock::objectlock_sys::ObjectLockBlockReason; diff --git a/rustfs/src/storage_api.rs b/rustfs/src/storage_api.rs index e74cb689a..62f883363 100644 --- a/rustfs/src/storage_api.rs +++ b/rustfs/src/storage_api.rs @@ -147,8 +147,8 @@ pub(crate) mod server { #[cfg(test)] pub(crate) use crate::storage::storage_api::{ - Endpoint, EndpointServerPools, Endpoints, PeerRestClient, PoolEndpoints, gen_signature_headers, - gen_tonic_signature_headers, + Endpoint, EndpointServerPools, Endpoints, PeerRestClient, PoolEndpoints, ScannerScopedDirtyUsageAckEntry, + gen_signature_headers, gen_tonic_signature_headers, }; pub(crate) mod ecfs { diff --git a/scripts/error-other-format-baseline.txt b/scripts/error-other-format-baseline.txt index 269a7d3df..4a45bab07 100644 --- a/scripts/error-other-format-baseline.txt +++ b/scripts/error-other-format-baseline.txt @@ -45,7 +45,7 @@ 1|crates/ecstore/src/object_api/types.rs 3|crates/ecstore/src/runtime/sources.rs 4|crates/ecstore/src/services/batch_processor.rs -14|crates/ecstore/src/services/notification_sys.rs +13|crates/ecstore/src/services/notification_sys.rs 16|crates/ecstore/src/services/rebalance/control.rs 1|crates/ecstore/src/services/rebalance/entry.rs 8|crates/ecstore/src/services/rebalance/meta.rs