diff --git a/crates/scanner/src/lib.rs b/crates/scanner/src/lib.rs index d4d76308c..5a4b57d1e 100644 --- a/crates/scanner/src/lib.rs +++ b/crates/scanner/src/lib.rs @@ -103,8 +103,9 @@ pub use scanner_io::{ acknowledge_scoped_dirty_usage, clear_dirty_usage_bucket, encode_durable_dirty_usage_producer_replay_record, record_dirty_usage_bucket, record_dirty_usage_bucket_from_producer, record_dirty_usage_bucket_from_producers, record_dirty_usage_object, record_dirty_usage_object_from_producer, record_scanner_maintenance_change, - replay_durable_dirty_usage_producer_record, scanner_activity_epoch, scanner_dirty_usage_snapshot, scanner_dirty_usage_state, - scanner_maintenance_generation, set_scanner_dirty_usage_clear_observer, set_scanner_dirty_usage_mutation_observer, + replay_durable_dirty_usage_producer_record, scanner_activity_epoch, scanner_dirty_usage_bucket_generation, + scanner_dirty_usage_snapshot, scanner_dirty_usage_state, scanner_maintenance_generation, + set_scanner_dirty_usage_clear_observer, set_scanner_dirty_usage_mutation_observer, }; pub use segment_invalidation::SegmentInvalidationProducerIdentity; pub use sleeper::{DynamicSleeper, SCANNER_IDLE_MODE, SCANNER_SLEEPER}; diff --git a/crates/scanner/src/scanner.rs b/crates/scanner/src/scanner.rs index c32d77f29..2377c317c 100644 --- a/crates/scanner/src/scanner.rs +++ b/crates/scanner/src/scanner.rs @@ -2230,6 +2230,7 @@ where cycle = cycle_info.current, required_cycle, state = "cache_cycle_ahead", + partial_cause = "cache_cycle_ahead", "Scanner cycle is recovering to a newer durable cache generation" ); emit_scan_cycle_partial_with_source(cycle_start.elapsed(), ScanCyclePartialReason::Unknown, None); @@ -2341,6 +2342,15 @@ where } usage_publication_result.restrict_outcome(usage_persist_outcome); + let partial_cause = if scan_cycle_result.has_observational_snapshot() + && matches!( + scan_cycle_result.status, + ScannerCycleStatus::Deferred(ScannerCycleDeferReason::ActivityBaselineUnavailable) + ) { + "activity_unverified_observation" + } else { + "incomplete_coverage" + }; let (completion_outcome, scanner_pending_maintenance_work, remote_dirty_usage_acknowledgements) = finalize_scanner_cycle_result(scan_cycle_result, usage_publication_result); let remote_dirty_usage_pending = if remote_dirty_usage_acknowledgements.is_empty() { @@ -2405,6 +2415,7 @@ where subsystem = LOG_SUBSYSTEM_RUNTIME, cycle = cycle_info.current, state = "incomplete", + partial_cause, "Scanner cycle ended without a complete usage snapshot" ); } @@ -3526,8 +3537,26 @@ where // Pending namespace commits invalidate this publication attempt, but only // storage movement creates durable, rate-limited catch-up debt. if storeapi.scanner_data_movement_pause_status().await.paused { + debug!( + target: "rustfs::scanner", + event = EVENT_SCANNER_PERSIST_STATE, + component = LOG_COMPONENT_SCANNER, + subsystem = LOG_SUBSYSTEM_RUNTIME, + stage = "local_barrier", + blocker = "data_movement", + "Scanner usage publication deferred" + ); Some(ScannerCycleDeferReason::DataMovement) } else { + debug!( + target: "rustfs::scanner", + event = EVENT_SCANNER_PERSIST_STATE, + component = LOG_COMPONENT_SCANNER, + subsystem = LOG_SUBSYSTEM_RUNTIME, + stage = "local_barrier", + blocker = "pending_namespace_commit", + "Scanner usage publication deferred" + ); Some(ScannerCycleDeferReason::ActivityBaselineUnavailable) } } @@ -3543,7 +3572,24 @@ fn scanner_post_lease_activity_defer_reason( { None } - Ok(_) | Err(_) => Some(ScannerCycleDeferReason::ActivityBaselineUnavailable), + observed => { + let blocker = match &observed { + Err(_) => "probe_failed", + Ok(snapshot) if !scanner_activity_allows_usage_publication(snapshot) => "publication_blocked", + Ok(_) if expected_digest.is_none() => "baseline_missing", + Ok(_) => "activity_changed", + }; + debug!( + target: "rustfs::scanner", + event = EVENT_SCANNER_PERSIST_STATE, + component = LOG_COMPONENT_SCANNER, + subsystem = LOG_SUBSYSTEM_RUNTIME, + stage = "post_lease_activity", + blocker, + "Scanner usage publication deferred" + ); + Some(ScannerCycleDeferReason::ActivityBaselineUnavailable) + } } } diff --git a/crates/scanner/src/scanner/usage_store.rs b/crates/scanner/src/scanner/usage_store.rs index a378d8a7a..21a8a7cad 100644 --- a/crates/scanner/src/scanner/usage_store.rs +++ b/crates/scanner/src/scanner/usage_store.rs @@ -1097,6 +1097,16 @@ where global_metrics().record_scanner_usage_save_result(ScannerUsageSaveResult::Success); global_metrics().record_scanner_usage_durable_success(); outcome = DataUsagePersistOutcome::Saved; + debug!( + target: "rustfs::scanner", + event = EVENT_SCANNER_PERSIST_STATE, + component = LOG_COMPONENT_SCANNER, + subsystem = LOG_SUBSYSTEM_RUNTIME, + state = "saved", + snapshot_kind = if observational { "observed" } else { "authoritative" }, + cycle = ?data_usage_info.scanner_cycle, + "Scanner usage snapshot saved" + ); } } diff --git a/crates/scanner/src/scanner_io.rs b/crates/scanner/src/scanner_io.rs index 0ff8b2dfc..1af655211 100644 --- a/crates/scanner/src/scanner_io.rs +++ b/crates/scanner/src/scanner_io.rs @@ -1083,6 +1083,7 @@ where component = LOG_COMPONENT_SCANNER, subsystem = LOG_SUBSYSTEM_IO, state = "cycle_activity_probe_failed", + stage = "post_scan_probe", error = %err, "Scanner cycle activity verification failed" ); @@ -1619,8 +1620,9 @@ pub use dirty_usage::{ acknowledge_scoped_dirty_usage, clear_dirty_usage_bucket, encode_durable_dirty_usage_producer_replay_record, record_dirty_usage_bucket, record_dirty_usage_bucket_from_producer, record_dirty_usage_bucket_from_producers, record_dirty_usage_object, record_dirty_usage_object_from_producer, record_scanner_maintenance_change, - replay_durable_dirty_usage_producer_record, scanner_activity_epoch, scanner_dirty_usage_snapshot, scanner_dirty_usage_state, - scanner_maintenance_generation, set_scanner_dirty_usage_clear_observer, set_scanner_dirty_usage_mutation_observer, + replay_durable_dirty_usage_producer_record, scanner_activity_epoch, scanner_dirty_usage_bucket_generation, + scanner_dirty_usage_snapshot, scanner_dirty_usage_state, scanner_maintenance_generation, + set_scanner_dirty_usage_clear_observer, set_scanner_dirty_usage_mutation_observer, }; #[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 2c0e7fe4e..43af3e1df 100644 --- a/crates/scanner/src/scanner_io/dirty_usage.rs +++ b/crates/scanner/src/scanner_io/dirty_usage.rs @@ -243,7 +243,11 @@ pub fn encode_durable_dirty_usage_producer_replay_record( entries, }; validate_durable_dirty_usage_replay_record(&record)?; - serde_json::to_vec(&record).map_err(|_| ScannerDurableDirtyUsageReplayError::InvalidJson) + let bytes = serde_json::to_vec(&record).map_err(|_| ScannerDurableDirtyUsageReplayError::InvalidJson)?; + if bytes.len() > SCANNER_DURABLE_DIRTY_USAGE_REPLAY_MAX_BYTES { + return Err(ScannerDurableDirtyUsageReplayError::ByteLimit); + } + Ok(bytes) } pub fn replay_durable_dirty_usage_producer_record( @@ -661,6 +665,24 @@ mod scoped_dirty_usage_tests { clear_dirty_usage_buckets_for_tests(); } + #[test] + fn durable_dirty_usage_encoder_rejects_unreplayable_byte_size() { + let entries = (0..3) + .map(|bucket| ScannerDurableDirtyUsageReplayEntry { + bucket: format!("bucket-{bucket}"), + generation: 7, + scope: ScannerDurableDirtyUsageReplayScope::TopLevelEntries { + entries: (0..128).map(|entry| format!("{entry:03}{}", "x".repeat(197))).collect(), + }, + producers: BTreeSet::from([SegmentInvalidationProducerIdentity::PutObject]), + }) + .collect(); + assert_eq!( + encode_durable_dirty_usage_producer_replay_record(entries), + Err(ScannerDurableDirtyUsageReplayError::ByteLimit) + ); + } + #[test] #[serial] fn durable_dirty_usage_replay_rejects_invalid_records_without_partial_state() { @@ -1001,6 +1023,11 @@ pub fn scanner_dirty_usage_state() -> ScannerDirtyUsageState { } } +/// Read one bucket generation without allocating a cluster-wide snapshot. +pub fn scanner_dirty_usage_bucket_generation(bucket: &str) -> Option { + dirty_usage_buckets().get(bucket).copied() +} + pub fn scanner_dirty_usage_snapshot(max_entries: usize) -> ScannerDirtyUsageSnapshot { let (generation, pending_bucket_count, complete, mut buckets) = { let dirty_buckets = dirty_usage_buckets(); diff --git a/crates/scanner/src/scanner_io/io_cycle.rs b/crates/scanner/src/scanner_io/io_cycle.rs index 08a177984..b77411e09 100644 --- a/crates/scanner/src/scanner_io/io_cycle.rs +++ b/crates/scanner/src/scanner_io/io_cycle.rs @@ -333,6 +333,7 @@ where component = LOG_COMPONENT_SCANNER, subsystem = LOG_SUBSYSTEM_IO, state = "cycle_activity_baseline_failed", + stage = "baseline_probe", error = %err, "Scanner cycle skipped because cluster activity could not be baselined" ); diff --git a/docs/architecture/scanner-usage-publication.md b/docs/architecture/scanner-usage-publication.md index 0665fc473..12c600d21 100644 --- a/docs/architecture/scanner-usage-publication.md +++ b/docs/architecture/scanner-usage-publication.md @@ -47,6 +47,39 @@ This PUT-tail protection requires every writer node to be upgraded. It does not prove that a failed tail replica has healed, and it does not extend the same in-flight tracking to multipart or other namespace mutation paths. +## Dirty producer replay ownership + +The dirty producer journal is partitioned by the configured node RPC authority. +A domain-separated SHA-256 owner identifier selects +`scanner/durable-dirty-producer-replay/.json`; the record also binds that +owner, which must match before replay. The same endpoint reuses its record after +a process restart. Changing the endpoint selects a new record and retains the +existing full-scan and peer-instance fences. + +The legacy shared `scanner/durable-dirty-producer-replay.json` is retained without +being imported or deleted: it cannot prove which node owns all pending entries. +Mixed-version nodes therefore cannot overwrite or clear upgraded nodes' records. +Rollback leaves the owned records available for a subsequent upgrade. + +Committed mutation callbacks coalesce into bounded node-local state. A single +worker snapshots that state at most once per second and retries failed storage +operations after five seconds. It holds no state lock during storage I/O and +confirms only the snapshot revision it saved. Dropping the last journal handle +releases its worker. This asynchronous hint is not a per-S3-commit durability +receipt; initial full scans and peer-instance changes still fence restart gaps. + +The complete owned JSON record must fit the replay reader's 64 KiB limit. +Oversized prefix scopes become whole-bucket hints; if the resulting complete +record still cannot fit, an explicit unverified marker prevents truncated +coverage from becoming producer authority. Invalid or unknown mutations remain +unverified until a verified dirty acknowledgement clears all pending work. + +Publication logs preserve the existing retry and ACK semantics while exposing +`stage`, `blocker`, `partial_cause`, and `snapshot_kind`. A successful `observed` +save does not imply authoritative publication or dirty acknowledgement. Cold +bucket traversal still requires complete distributed invalidation evidence; +node-local journal ownership does not relax the cluster activity proof. + ## Fences The protocol uses separate fences because they exclude different stale inputs. diff --git a/rustfs/src/storage/scanner_dirty_journal.rs b/rustfs/src/storage/scanner_dirty_journal.rs index 5ab3a1662..ff7d9b2e8 100644 --- a/rustfs/src/storage/scanner_dirty_journal.rs +++ b/rustfs/src/storage/scanner_dirty_journal.rs @@ -17,78 +17,121 @@ use rustfs_scanner::{ ScannerDirtyUsageBucket, ScannerDurableDirtyUsageReplayEntry, ScannerDurableDirtyUsageReplayRecord, ScannerDurableDirtyUsageReplayScope, SegmentInvalidationProducerIdentity, }; +use serde::{Deserialize, Serialize}; +use sha2::{Digest, Sha256}; use std::collections::{BTreeMap, BTreeSet}; -use std::sync::Arc; +use std::sync::{Arc, Mutex, Weak}; use std::time::Duration; -use tokio::sync::mpsc; +use tokio::sync::Notify; use tracing::{debug, warn}; const LOG_COMPONENT_STORAGE: &str = "storage"; const LOG_SUBSYSTEM_SCANNER: &str = "scanner"; const EVENT_SCANNER_DIRTY_USAGE_JOURNAL: &str = "scanner_dirty_usage_journal"; -const DURABLE_DIRTY_USAGE_REPLAY_OBJECT: &str = "scanner/durable-dirty-producer-replay.json"; +const DURABLE_DIRTY_USAGE_REPLAY_PREFIX: &str = "scanner/durable-dirty-producer-replay"; const DURABLE_DIRTY_USAGE_REPLAY_MAX_BYTES: usize = 64 * 1024; const DURABLE_DIRTY_USAGE_REPLAY_MAX_ENTRIES: usize = 1024; const DURABLE_DIRTY_USAGE_REPLAY_MAX_TOP_LEVEL_ENTRIES: usize = 128; -const DURABLE_DIRTY_USAGE_JOURNAL_CHANNEL_DEPTH: usize = 256; +const DURABLE_DIRTY_USAGE_JOURNAL_FLUSH_INTERVAL: Duration = Duration::from_secs(1); const DURABLE_DIRTY_USAGE_JOURNAL_RETRY_INTERVAL: Duration = Duration::from_secs(5); #[derive(Clone)] pub(crate) struct DurableDirtyUsageJournal { - sender: mpsc::Sender, + shared: Option>, } -impl DurableDirtyUsageJournal { - pub(crate) fn record_committed_mutation(&self, bucket: &str, object: &str, producer: SegmentInvalidationProducerIdentity) { - let snapshot = rustfs_scanner::scanner_dirty_usage_snapshot(DURABLE_DIRTY_USAGE_REPLAY_MAX_ENTRIES); - let generation = snapshot - .buckets - .into_iter() - .find(|entry| entry.bucket == bucket) - .map(|entry| entry.generation) - .unwrap_or(0); - self.dispatch(JournalCommand::RecordMutation { - bucket: bucket.to_string(), - object: object.to_string(), - generation, - producer, - }); +struct SharedJournal { + pending: Mutex, + changed: Arc, + #[cfg(test)] + failed_attempts: std::sync::atomic::AtomicU64, +} + +#[derive(Default)] +struct PendingJournal { + state: JournalState, + revision: u64, + persisted_revision: u64, +} + +impl PendingJournal { + fn changed(&mut self) { + self.revision = self.revision.saturating_add(1); } - pub(crate) fn clear_confirmed_buckets(&self, cleared: Vec) { - self.dispatch(JournalCommand::ClearBuckets { cleared }); - } - - fn dispatch(&self, command: JournalCommand) { - match self.sender.try_send(command) { - Ok(()) => {} - Err(mpsc::error::TrySendError::Full(command)) => { - let sender = self.sender.clone(); - tokio::spawn(async move { - let _ = sender.send(command).await; - }); - } - Err(mpsc::error::TrySendError::Closed(_)) => {} + fn confirm_flush(&mut self, revision: u64) { + // A save cannot acknowledge mutations received after its snapshot. + if revision == self.revision && revision != u64::MAX { + self.persisted_revision = revision; } } } -#[derive(Debug)] -enum JournalCommand { - RecordMutation { - bucket: String, - object: String, - generation: u64, - producer: SegmentInvalidationProducerIdentity, - }, - ClearBuckets { - cleared: Vec, - }, +impl DurableDirtyUsageJournal { + pub(crate) fn record_committed_mutation(&self, bucket: &str, object: &str, producer: SegmentInvalidationProducerIdentity) { + let Some(shared) = &self.shared else { return }; + let generation = rustfs_scanner::scanner_dirty_usage_bucket_generation(bucket).unwrap_or(0); + { + // Release the scanner dirty-map lock first; no journal guard spans I/O. + let mut pending = shared.pending.lock().unwrap_or_else(std::sync::PoisonError::into_inner); + pending + .state + .record_mutation(bucket.to_string(), object, generation, producer); + pending.changed(); + } + shared.changed.notify_one(); + } + + pub(crate) fn clear_confirmed_buckets(&self, cleared: Vec) { + let Some(shared) = &self.shared else { return }; + { + let mut pending = shared.pending.lock().unwrap_or_else(std::sync::PoisonError::into_inner); + pending.state.clear_buckets(&cleared); + // Observer callbacks run after scanner locks are released. Recheck under + // the journal lock so an intervening unknown mutation cannot lose its marker. + if pending.state.buckets.is_empty() && !rustfs_scanner::scanner_dirty_usage_state().pending { + pending.state.invalidated = false; + } + pending.changed(); + } + shared.changed.notify_one(); + } +} + +impl Drop for DurableDirtyUsageJournal { + fn drop(&mut self) { + if let Some(shared) = &self.shared { + // An idle worker owns only a Weak reference and exits after its owner drops. + shared.changed.notify_one(); + } + } +} + +#[derive(Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +struct OwnedReplayRecord { + owner: String, + replay: Option, +} + +fn durable_dirty_usage_owner(node_name: &str) -> Option { + (!node_name.is_empty()).then(|| { + let mut hash = Sha256::new(); + hash.update(b"rustfs/scanner/dirty-journal/owner/v1\0"); + hash.update(node_name.as_bytes()); + hex_simd::encode_to_string(hash.finalize(), hex_simd::AsciiCase::Lower) + }) +} + +fn durable_dirty_usage_replay_object(owner: &str) -> String { + format!("{DURABLE_DIRTY_USAGE_REPLAY_PREFIX}/{owner}.json") } #[derive(Clone, Debug, Default, PartialEq, Eq)] struct JournalState { buckets: BTreeMap, + scope_bytes: usize, + invalidated: bool, } #[derive(Clone, Debug, PartialEq, Eq)] @@ -109,22 +152,31 @@ impl JournalState { if bucket.is_empty() || generation == 0 || generation == u64::MAX || !durable_dirty_usage_producer_is_supported(producer) { self.buckets.clear(); + self.scope_bytes = 0; + self.invalidated = true; return; } - let event_scope = durable_dirty_usage_journal_scope(object); - self.buckets - .entry(bucket) - .and_modify(|state| { - state.generation = state.generation.max(generation); - merge_journal_scope(&mut state.scope, event_scope.clone()); - state.producers.insert(producer); - }) - .or_insert_with(|| JournalBucketState { - generation, - scope: event_scope, - producers: BTreeSet::from([producer]), - }); + let state = self.buckets.entry(bucket).or_insert_with(|| JournalBucketState { + generation, + scope: JournalScope::TopLevelEntries(BTreeSet::new()), + producers: BTreeSet::new(), + }); + self.scope_bytes = self.scope_bytes.saturating_sub(journal_scope_bytes(&state.scope)); + state.generation = state.generation.max(generation); + merge_journal_scope(&mut state.scope, durable_dirty_usage_journal_scope(object)); + state.producers.insert(producer); + let bytes = journal_scope_bytes(&state.scope); + if self.scope_bytes.saturating_add(bytes) > DURABLE_DIRTY_USAGE_REPLAY_MAX_BYTES { + state.scope = JournalScope::WholeBucket; + } else { + self.scope_bytes += bytes; + } + if self.buckets.len() > DURABLE_DIRTY_USAGE_REPLAY_MAX_ENTRIES { + self.buckets.clear(); + self.scope_bytes = 0; + self.invalidated = true; + } } fn clear_buckets(&mut self, cleared: &[ScannerDirtyUsageBucket]) { @@ -133,14 +185,15 @@ impl JournalState { .buckets .get(&entry.bucket) .is_some_and(|state| state.generation <= entry.generation) + && let Some(removed) = self.buckets.remove(&entry.bucket) { - self.buckets.remove(&entry.bucket); + self.scope_bytes = self.scope_bytes.saturating_sub(journal_scope_bytes(&removed.scope)); } } } fn replay_entries(&self) -> Option> { - if self.buckets.is_empty() || self.buckets.len() > DURABLE_DIRTY_USAGE_REPLAY_MAX_ENTRIES { + if self.invalidated || self.buckets.is_empty() || self.buckets.len() > DURABLE_DIRTY_USAGE_REPLAY_MAX_ENTRIES { return None; } let mut entries = Vec::with_capacity(self.buckets.len()); @@ -204,6 +257,7 @@ impl JournalState { JournalScope::TopLevelEntries(entries) } }; + state.scope_bytes += journal_scope_bytes(&scope); if state .buckets .insert( @@ -224,14 +278,36 @@ impl JournalState { } pub(crate) async fn start_durable_dirty_usage_journal(store: Arc) -> DurableDirtyUsageJournal { - let initial_state = replay_durable_dirty_usage_journal(store.clone()).await; - let (sender, receiver) = mpsc::channel(DURABLE_DIRTY_USAGE_JOURNAL_CHANNEL_DEPTH); - tokio::spawn(run_durable_dirty_usage_journal(store, receiver, initial_state)); - DurableDirtyUsageJournal { sender } + // ECStore initializes this stable endpoint authority before AppContext exists. + let node_name = rustfs_common::get_global_local_node_name().await; + let Some(owner) = durable_dirty_usage_owner(&node_name) else { + warn!( + event = EVENT_SCANNER_DIRTY_USAGE_JOURNAL, + component = LOG_COMPONENT_STORAGE, + subsystem = LOG_SUBSYSTEM_SCANNER, + state = "owner_unavailable", + "Scanner dirty journal owner is unavailable" + ); + return DurableDirtyUsageJournal { shared: None }; + }; + let state = replay_durable_dirty_usage_journal(store.clone(), &owner).await; + let changed = Arc::new(Notify::new()); + let shared = Arc::new(SharedJournal { + pending: Mutex::new(PendingJournal { + state, + ..Default::default() + }), + changed: changed.clone(), + #[cfg(test)] + failed_attempts: std::sync::atomic::AtomicU64::new(0), + }); + tokio::spawn(run_durable_dirty_usage_journal(store, owner, Arc::downgrade(&shared), changed)); + DurableDirtyUsageJournal { shared: Some(shared) } } -async fn replay_durable_dirty_usage_journal(store: Arc) -> JournalState { - let bytes = match super::read_config(store, DURABLE_DIRTY_USAGE_REPLAY_OBJECT).await { +async fn replay_durable_dirty_usage_journal(store: Arc, owner: &str) -> JournalState { + let path = durable_dirty_usage_replay_object(owner); + let bytes = match super::read_config(store, &path).await { Ok(bytes) if !bytes.is_empty() => bytes, Ok(_) | Err(Error::ConfigNotFound) => return JournalState::default(), Err(err) => { @@ -241,48 +317,39 @@ async fn replay_durable_dirty_usage_journal(store: Arc) -> JournalState subsystem = LOG_SUBSYSTEM_SCANNER, state = "read_failed", error = ?err, - "Durable scanner dirty usage journal could not be read" + "Scanner dirty journal could not be read" ); - return JournalState::default(); + return JournalState { + invalidated: true, + ..Default::default() + }; } }; - if bytes.len() > DURABLE_DIRTY_USAGE_REPLAY_MAX_BYTES { + let Some((state, replay)) = decode_owned_replay_record(owner, &bytes) else { warn!( event = EVENT_SCANNER_DIRTY_USAGE_JOURNAL, component = LOG_COMPONENT_STORAGE, subsystem = LOG_SUBSYSTEM_SCANNER, - state = "oversized", - "Durable scanner dirty usage journal exceeds replay size limit" + state = "replay_unverified", + "Scanner dirty journal requires a full scan" ); - return JournalState::default(); - } - match rustfs_scanner::replay_durable_dirty_usage_producer_record(&bytes) { - Ok(state) => { + return JournalState { + invalidated: true, + ..Default::default() + }; + }; + match rustfs_scanner::replay_durable_dirty_usage_producer_record(&replay) { + Ok(replayed) => { debug!( event = EVENT_SCANNER_DIRTY_USAGE_JOURNAL, component = LOG_COMPONENT_STORAGE, subsystem = LOG_SUBSYSTEM_SCANNER, state = "replayed", - generation = state.generation, - pending = state.pending, - "Durable scanner dirty usage journal replayed" + generation = replayed.generation, + pending = replayed.pending, + "Scanner dirty journal replayed" ); - match serde_json::from_slice::(&bytes) - .ok() - .and_then(JournalState::from_replay_record) - { - Some(state) => state, - None => { - warn!( - event = EVENT_SCANNER_DIRTY_USAGE_JOURNAL, - component = LOG_COMPONENT_STORAGE, - subsystem = LOG_SUBSYSTEM_SCANNER, - state = "hydrate_rejected", - "Durable scanner dirty usage journal could not hydrate writer state" - ); - JournalState::default() - } - } + state } Err(err) => { warn!( @@ -291,76 +358,119 @@ async fn replay_durable_dirty_usage_journal(store: Arc) -> JournalState subsystem = LOG_SUBSYSTEM_SCANNER, state = "replay_rejected", error = ?err, - "Durable scanner dirty usage journal was rejected" + "Scanner dirty journal was rejected" ); - JournalState::default() - } - } -} - -async fn run_durable_dirty_usage_journal( - store: Arc, - mut receiver: mpsc::Receiver, - initial_state: JournalState, -) { - let mut state = initial_state; - let mut dirty = false; - let mut retry = tokio::time::interval(DURABLE_DIRTY_USAGE_JOURNAL_RETRY_INTERVAL); - retry.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay); - - loop { - tokio::select! { - command = receiver.recv() => { - let Some(command) = command else { - break; - }; - apply_journal_command(&mut state, command); - while let Ok(command) = receiver.try_recv() { - apply_journal_command(&mut state, command); - } - dirty = true; + JournalState { + invalidated: true, + ..Default::default() } - _ = retry.tick(), if dirty => {} - } - - if dirty && flush_durable_dirty_usage_journal(store.clone(), &state).await { - dirty = false; } } } -fn apply_journal_command(state: &mut JournalState, command: JournalCommand) { - match command { - JournalCommand::RecordMutation { - bucket, - object, - generation, - producer, - } => state.record_mutation(bucket, &object, generation, producer), - JournalCommand::ClearBuckets { cleared } => state.clear_buckets(&cleared), +fn decode_owned_replay_record(owner: &str, bytes: &[u8]) -> Option<(JournalState, Vec)> { + if bytes.len() > DURABLE_DIRTY_USAGE_REPLAY_MAX_BYTES { + return None; + } + let record: OwnedReplayRecord = serde_json::from_slice(bytes).ok()?; + if record.owner != owner { + return None; + } + let replay = record.replay?; + let bytes = serde_json::to_vec(&replay).ok()?; + let state = JournalState::from_replay_record(replay)?; + Some((state, bytes)) +} + +fn encode_owned_replay_record(owner: &str, state: &JournalState) -> Result, serde_json::Error> { + let mut entries = state.replay_entries(); + for compact in [false, true] { + if let Some(entries) = &mut entries { + if compact { + for entry in entries.iter_mut() { + entry.scope = ScannerDurableDirtyUsageReplayScope::WholeBucket; + } + } + if let Ok(replay) = rustfs_scanner::encode_durable_dirty_usage_producer_replay_record(entries.clone()) + && let Ok(replay) = serde_json::from_slice(&replay) + { + let bytes = serde_json::to_vec(&OwnedReplayRecord { + owner: owner.to_string(), + replay: Some(replay), + })?; + if bytes.len() <= DURABLE_DIRTY_USAGE_REPLAY_MAX_BYTES { + return Ok(bytes); + } + } + } + } + // Never label a truncated record as complete producer coverage. + serde_json::to_vec(&OwnedReplayRecord { + owner: owner.to_string(), + replay: None, + }) +} + +async fn run_durable_dirty_usage_journal(store: Arc, owner: String, shared: Weak, changed: Arc) { + let mut next_flush = tokio::time::Instant::now() + DURABLE_DIRTY_USAGE_JOURNAL_FLUSH_INTERVAL; + loop { + changed.notified().await; + if shared.strong_count() == 0 { + break; + } + tokio::time::sleep_until(next_flush).await; + let Some(current) = shared.upgrade() else { break }; + let snapshot = { + let pending = current.pending.lock().unwrap_or_else(std::sync::PoisonError::into_inner); + (pending.revision != pending.persisted_revision).then(|| (pending.state.clone(), pending.revision)) + }; + drop(current); + let Some((snapshot, revision)) = snapshot else { continue }; + let saved = flush_durable_dirty_usage_journal(store.clone(), &owner, &snapshot).await; + next_flush = tokio::time::Instant::now() + + if saved { + DURABLE_DIRTY_USAGE_JOURNAL_FLUSH_INTERVAL + } else { + DURABLE_DIRTY_USAGE_JOURNAL_RETRY_INTERVAL + }; + let Some(current) = shared.upgrade() else { break }; + #[cfg(test)] + if !saved { + current.failed_attempts.fetch_add(1, std::sync::atomic::Ordering::Release); + } + let mut pending = current.pending.lock().unwrap_or_else(std::sync::PoisonError::into_inner); + if saved { + pending.confirm_flush(revision); + } + if pending.revision != pending.persisted_revision { + changed.notify_one(); + } } } -async fn flush_durable_dirty_usage_journal(store: Arc, state: &JournalState) -> bool { - let Some(entries) = state.replay_entries() else { - return delete_durable_dirty_usage_journal(store).await; +async fn flush_durable_dirty_usage_journal(store: Arc, owner: &str, state: &JournalState) -> bool { + let path = durable_dirty_usage_replay_object(owner); + let result = if state.buckets.is_empty() && !state.invalidated { + ecstore_config::com::delete_config(store, &path).await + } else { + let bytes = match encode_owned_replay_record(owner, state) { + Ok(bytes) => bytes, + Err(err) => { + warn!( + event = EVENT_SCANNER_DIRTY_USAGE_JOURNAL, + component = LOG_COMPONENT_STORAGE, + subsystem = LOG_SUBSYSTEM_SCANNER, + state = "encode_failed", + error = ?err, + "Scanner dirty journal could not be encoded" + ); + return false; + } + }; + ecstore_config::com::save_config(store, &path, bytes).await }; - let bytes = match rustfs_scanner::encode_durable_dirty_usage_producer_replay_record(entries) { - Ok(bytes) => bytes, - Err(err) => { - warn!( - event = EVENT_SCANNER_DIRTY_USAGE_JOURNAL, - component = LOG_COMPONENT_STORAGE, - subsystem = LOG_SUBSYSTEM_SCANNER, - state = "encode_rejected", - error = ?err, - "Durable scanner dirty usage journal could not be encoded" - ); - return delete_durable_dirty_usage_journal(store).await; - } - }; - match ecstore_config::com::save_config(store, DURABLE_DIRTY_USAGE_REPLAY_OBJECT, bytes).await { - Ok(()) => true, + match result { + Ok(()) | Err(Error::ConfigNotFound) => true, Err(err) => { warn!( event = EVENT_SCANNER_DIRTY_USAGE_JOURNAL, @@ -368,27 +478,17 @@ async fn flush_durable_dirty_usage_journal(store: Arc, state: &JournalS subsystem = LOG_SUBSYSTEM_SCANNER, state = "write_failed", error = ?err, - "Durable scanner dirty usage journal write failed" + "Scanner dirty journal persistence failed" ); false } } } -async fn delete_durable_dirty_usage_journal(store: Arc) -> bool { - match ecstore_config::com::delete_config(store, DURABLE_DIRTY_USAGE_REPLAY_OBJECT).await { - Ok(()) | Err(Error::ConfigNotFound) => true, - Err(err) => { - warn!( - event = EVENT_SCANNER_DIRTY_USAGE_JOURNAL, - component = LOG_COMPONENT_STORAGE, - subsystem = LOG_SUBSYSTEM_SCANNER, - state = "delete_failed", - error = ?err, - "Durable scanner dirty usage journal delete failed" - ); - false - } +fn journal_scope_bytes(scope: &JournalScope) -> usize { + match scope { + JournalScope::WholeBucket => 0, + JournalScope::TopLevelEntries(entries) => entries.iter().map(String::len).sum(), } } @@ -444,12 +544,357 @@ fn durable_dirty_usage_top_level_entry(object: &str) -> Option { #[cfg(test)] mod tests { - use super::{DURABLE_DIRTY_USAGE_REPLAY_MAX_TOP_LEVEL_ENTRIES, JournalState}; + use super::{ + DURABLE_DIRTY_USAGE_REPLAY_MAX_BYTES, DURABLE_DIRTY_USAGE_REPLAY_MAX_TOP_LEVEL_ENTRIES, DurableDirtyUsageJournal, + JournalState, OwnedReplayRecord, PendingJournal, SharedJournal, decode_owned_replay_record, durable_dirty_usage_owner, + durable_dirty_usage_replay_object, encode_owned_replay_record, flush_durable_dirty_usage_journal, + replay_durable_dirty_usage_journal, run_durable_dirty_usage_journal, + }; use rustfs_scanner::{ ScannerDirtyUsageBucket, ScannerDurableDirtyUsageReplayEntry, ScannerDurableDirtyUsageReplayRecord, ScannerDurableDirtyUsageReplayScope, SegmentInvalidationProducerIdentity, }; use std::collections::BTreeSet; + use std::sync::{Arc, Mutex}; + use std::time::Duration; + use tokio::sync::Notify; + + async fn journal_test_env() -> &'static rustfs_test_utils::TestECStoreEnv { + // Bucket metadata owns a process-global Weak store reference. Keep its + // fixture alive across serialized journal cases in the same test process. + static ENV: tokio::sync::OnceCell = tokio::sync::OnceCell::const_new(); + ENV.get_or_init(|| rustfs_test_utils::TestECStoreEnv::builder().build()).await + } + + #[test] + fn journal_owner_is_stable_and_never_defaults_to_a_shared_identity() { + assert!(durable_dirty_usage_owner("").is_none()); + let first = durable_dirty_usage_owner("node-a:9000").expect("configured owner"); + assert_eq!(Some(first.clone()), durable_dirty_usage_owner("node-a:9000")); + assert_ne!(Some(first), durable_dirty_usage_owner("node-b:9000")); + } + + #[tokio::test] + #[serial_test::serial] + async fn owned_journals_keep_peer_records_after_interleaved_writes_and_clear() { + let env = journal_test_env().await; + let owner_a = durable_dirty_usage_owner("node-a:9000").expect("owner a"); + let owner_b = durable_dirty_usage_owner("node-b:9000").expect("owner b"); + let mut state_a = JournalState::default(); + let mut state_b = JournalState::default(); + state_a.record_mutation("bucket-a".to_string(), "a/object", 7, SegmentInvalidationProducerIdentity::PutObject); + state_b.record_mutation("bucket-b".to_string(), "b/object", 3, SegmentInvalidationProducerIdentity::PutObject); + + let legacy_path = "scanner/durable-dirty-producer-replay.json"; + let legacy = + rustfs_scanner::encode_durable_dirty_usage_producer_replay_record(state_a.replay_entries().expect("legacy entries")) + .expect("legacy bytes"); + super::ecstore_config::com::save_config(env.ecstore.clone(), legacy_path, legacy.clone()) + .await + .expect("legacy save"); + assert!(flush_durable_dirty_usage_journal(env.ecstore.clone(), &owner_a, &state_a).await); + assert!(flush_durable_dirty_usage_journal(env.ecstore.clone(), &owner_b, &state_b).await); + let path_b = durable_dirty_usage_replay_object(&owner_b); + let bytes_b = super::super::read_config(env.ecstore.clone(), &path_b) + .await + .expect("read owner b"); + assert!(decode_owned_replay_record(&owner_a, &bytes_b).is_none()); + assert!(decode_owned_replay_record(&owner_a, &legacy).is_none()); + + state_a.clear_buckets(&[ScannerDirtyUsageBucket { + bucket: "bucket-a".to_string(), + generation: 7, + }]); + assert!(flush_durable_dirty_usage_journal(env.ecstore.clone(), &owner_a, &state_a).await); + assert_eq!( + super::super::read_config(env.ecstore.clone(), &path_b) + .await + .expect("peer survives"), + bytes_b + ); + assert_eq!( + super::super::read_config(env.ecstore.clone(), legacy_path) + .await + .expect("legacy retained"), + legacy + ); + let replayed = replay_durable_dirty_usage_journal(env.ecstore.clone(), &owner_b).await; + assert_eq!(replayed, state_b); + assert!(matches!( + super::super::read_config(env.ecstore.clone(), &durable_dirty_usage_replay_object(&owner_a)).await, + Err(super::Error::ConfigNotFound) + )); + let dirty = rustfs_scanner::scanner_dirty_usage_state(); + rustfs_scanner::acknowledge_dirty_usage_generation(rustfs_scanner::scanner_activity_epoch(), dirty.generation) + .expect("clear test replay"); + } + + #[test] + fn owned_journal_budget_compacts_scopes_or_records_unverified_coverage() { + let owner = durable_dirty_usage_owner("node-a:9000").expect("owner"); + let mut state = JournalState::default(); + for bucket in 0..3 { + for entry in 0..128 { + state.record_mutation( + format!("bucket-{bucket}"), + &format!("{entry:03}{}/object", "\"".repeat(197)), + 7, + SegmentInvalidationProducerIdentity::PutObject, + ); + } + } + assert!(state.scope_bytes <= DURABLE_DIRTY_USAGE_REPLAY_MAX_BYTES); + let bytes = encode_owned_replay_record(&owner, &state).expect("bounded record"); + assert!(bytes.len() <= DURABLE_DIRTY_USAGE_REPLAY_MAX_BYTES); + let (decoded, _) = decode_owned_replay_record(&owner, &bytes).expect("compacted replay"); + assert_eq!(decoded.buckets.len(), 3); + assert!( + decoded + .replay_entries() + .expect("entries") + .iter() + .all(|entry| entry.scope == ScannerDurableDirtyUsageReplayScope::WholeBucket) + ); + state = JournalState::default(); + for bucket in 0..1024 { + for producer in SegmentInvalidationProducerIdentity::REQUIRED_PRODUCTION { + state.record_mutation(format!("bucket-{bucket:04}-{}", "x".repeat(48)), "object", 9, producer); + } + } + assert!(!state.invalidated); + assert_eq!(state.buckets.len(), 1024); + let bytes = encode_owned_replay_record(&owner, &state).expect("unverified marker"); + assert!(bytes.len() <= DURABLE_DIRTY_USAGE_REPLAY_MAX_BYTES); + let record: OwnedReplayRecord = serde_json::from_slice(&bytes).expect("marker"); + assert!(record.replay.is_none()); + assert!(decode_owned_replay_record(&owner, &bytes).is_none()); + } + + #[test] + fn an_older_flush_cannot_confirm_a_new_mutation_or_exhausted_revision() { + let mut pending = PendingJournal::default(); + pending.changed(); + let saving = pending.revision; + pending.changed(); + pending.confirm_flush(saving); + assert_ne!(pending.revision, pending.persisted_revision); + pending.confirm_flush(pending.revision); + assert_eq!(pending.revision, pending.persisted_revision); + pending.revision = u64::MAX; + pending.changed(); + pending.confirm_flush(u64::MAX); + assert_ne!(pending.revision, pending.persisted_revision); + } + + #[tokio::test] + #[serial_test::serial] + async fn clear_callback_keeps_unknown_coverage_until_all_pending_work_is_confirmed() { + let env = journal_test_env().await; + let owner = durable_dirty_usage_owner("clear:9000").expect("owner"); + let shared = Arc::new(SharedJournal { + pending: Mutex::new(PendingJournal::default()), + changed: Arc::new(Notify::new()), + failed_attempts: std::sync::atomic::AtomicU64::new(0), + }); + let handle = DurableDirtyUsageJournal { + shared: Some(shared.clone()), + }; + rustfs_scanner::record_dirty_usage_bucket("unidentified"); + handle.record_committed_mutation("unidentified", "", SegmentInvalidationProducerIdentity::Unknown); + rustfs_scanner::record_dirty_usage_object_from_producer( + "identified", + "prefix/object", + SegmentInvalidationProducerIdentity::PutObject, + ); + handle.record_committed_mutation("identified", "prefix/object", SegmentInvalidationProducerIdentity::PutObject); + let old = rustfs_scanner::scanner_dirty_usage_bucket_generation("identified").expect("old generation"); + let snapshot = shared.pending.lock().expect("state").state.clone(); + assert!(snapshot.invalidated); + assert!(snapshot.scope_bytes > 0); + assert!(flush_durable_dirty_usage_journal(env.ecstore.clone(), &owner, &snapshot).await); + + rustfs_scanner::clear_dirty_usage_bucket("identified"); + handle.clear_confirmed_buckets(vec![ScannerDirtyUsageBucket { + bucket: "identified".to_string(), + generation: old, + }]); + { + let pending = shared.pending.lock().expect("partial clear"); + assert!(pending.state.invalidated, "the unknown bucket still has pending work"); + assert_eq!(pending.state.scope_bytes, 0); + } + rustfs_scanner::record_dirty_usage_object_from_producer( + "identified", + "new/object", + SegmentInvalidationProducerIdentity::PutObject, + ); + handle.record_committed_mutation("identified", "new/object", SegmentInvalidationProducerIdentity::PutObject); + handle.clear_confirmed_buckets(vec![ScannerDirtyUsageBucket { + bucket: "identified".to_string(), + generation: old, + }]); + let latest = rustfs_scanner::scanner_dirty_usage_bucket_generation("identified").expect("latest generation"); + assert!(latest > old); + { + let pending = shared.pending.lock().expect("stale clear"); + assert!(pending.state.invalidated); + assert_eq!(pending.state.buckets["identified"].generation, latest); + assert!(pending.state.scope_bytes > 0); + } + let dirty = rustfs_scanner::scanner_dirty_usage_state(); + rustfs_scanner::acknowledge_dirty_usage_generation(rustfs_scanner::scanner_activity_epoch(), dirty.generation) + .expect("verified full clear"); + handle.clear_confirmed_buckets(vec![ScannerDirtyUsageBucket { + bucket: "identified".to_string(), + generation: latest, + }]); + let cleared = shared.pending.lock().expect("full clear").state.clone(); + assert!(!cleared.invalidated); + assert!(cleared.buckets.is_empty()); + assert_eq!(cleared.scope_bytes, 0); + assert!(flush_durable_dirty_usage_journal(env.ecstore.clone(), &owner, &cleared).await); + assert!(matches!( + super::super::read_config(env.ecstore.clone(), &durable_dirty_usage_replay_object(&owner)).await, + Err(super::Error::ConfigNotFound) + )); + } + + #[tokio::test] + #[serial_test::serial] + async fn journal_worker_retries_failed_storage_without_a_new_mutation() { + let env = journal_test_env().await; + let owner = durable_dirty_usage_owner("retry:9000").expect("owner"); + let mut blocked = Vec::new(); + for disk in &env.disk_paths { + let path = disk.join(".rustfs.sys/scanner/durable-dirty-producer-replay"); + std::fs::create_dir_all(path.parent().expect("parent")).expect("scanner directory"); + let backup = path.with_extension("retry-backup"); + let existed = path.exists(); + if existed { + std::fs::rename(&path, &backup).expect("retain other journal records during fault injection"); + } + std::fs::write(&path, b"not a directory").expect("block journal writes"); + blocked.push((path, backup, existed)); + } + let mut pending = PendingJournal::default(); + pending + .state + .record_mutation("retry".to_string(), "prefix/object", 7, SegmentInvalidationProducerIdentity::PutObject); + pending.changed(); + let changed = Arc::new(Notify::new()); + let shared = Arc::new(SharedJournal { + pending: Mutex::new(pending), + changed: changed.clone(), + failed_attempts: std::sync::atomic::AtomicU64::new(0), + }); + let handle = DurableDirtyUsageJournal { + shared: Some(shared.clone()), + }; + let worker = tokio::spawn(run_durable_dirty_usage_journal( + env.ecstore.clone(), + owner.clone(), + Arc::downgrade(&shared), + changed.clone(), + )); + changed.notify_one(); + tokio::time::timeout(Duration::from_secs(20), async { + while shared.failed_attempts.load(std::sync::atomic::Ordering::Acquire) == 0 { + tokio::time::sleep(Duration::from_millis(20)).await; + } + }) + .await + .expect("real storage failure must be observed"); + assert_eq!(shared.pending.lock().expect("failed state").persisted_revision, 0); + for (path, backup, existed) in blocked { + std::fs::remove_file(&path).expect("restore journal writes"); + if existed { + std::fs::rename(backup, path).expect("restore other journal records"); + } + } + tokio::time::timeout(Duration::from_secs(30), async { + loop { + let confirmed = { + let pending = shared.pending.lock().expect("retry state"); + pending.persisted_revision == pending.revision + }; + if confirmed { + break; + } + tokio::time::sleep(Duration::from_millis(20)).await; + } + }) + .await + .expect("retry must persist the original mutation without another notification"); + let bytes = super::super::read_config(env.ecstore.clone(), &durable_dirty_usage_replay_object(&owner)) + .await + .expect("retry bytes"); + let (saved, _) = decode_owned_replay_record(&owner, &bytes).expect("retry replay"); + assert_eq!(saved.buckets["retry"].generation, 7); + drop(shared); + drop(handle); + tokio::time::timeout(Duration::from_secs(5), worker) + .await + .expect("worker stops") + .expect("worker task"); + } + + #[tokio::test] + #[serial_test::serial] + async fn journal_worker_coalesces_mutations_and_releases_its_store() { + let env = journal_test_env().await; + let owner = durable_dirty_usage_owner("worker:9000").expect("owner"); + let changed = Arc::new(Notify::new()); + let shared = Arc::new(SharedJournal { + pending: Mutex::new(PendingJournal::default()), + changed: changed.clone(), + failed_attempts: std::sync::atomic::AtomicU64::new(0), + }); + let handle = DurableDirtyUsageJournal { + shared: Some(shared.clone()), + }; + let worker = tokio::spawn(run_durable_dirty_usage_journal( + env.ecstore.clone(), + owner.clone(), + Arc::downgrade(&shared), + changed, + )); + for entry in 0..1000 { + let object = format!("prefix/object-{entry}"); + rustfs_scanner::record_dirty_usage_object_from_producer( + "coalesced", + &object, + SegmentInvalidationProducerIdentity::PutObject, + ); + handle.record_committed_mutation("coalesced", &object, SegmentInvalidationProducerIdentity::PutObject); + } + let expected = rustfs_scanner::scanner_dirty_usage_bucket_generation("coalesced").expect("generation"); + let path = durable_dirty_usage_replay_object(&owner); + tokio::time::timeout(Duration::from_secs(20), async { + loop { + if let Ok(bytes) = super::super::read_config(env.ecstore.clone(), &path).await + && let Some((state, _)) = decode_owned_replay_record(&owner, &bytes) + && state + .buckets + .get("coalesced") + .is_some_and(|bucket| bucket.generation == expected) + { + break; + } + tokio::time::sleep(Duration::from_millis(20)).await; + } + }) + .await + .expect("latest coalesced mutation should persist"); + drop(shared); + drop(handle); + tokio::time::timeout(Duration::from_secs(5), worker) + .await + .expect("worker should stop") + .expect("worker task"); + let dirty = rustfs_scanner::scanner_dirty_usage_state(); + rustfs_scanner::acknowledge_dirty_usage_generation(rustfs_scanner::scanner_activity_epoch(), dirty.generation) + .expect("clear worker test"); + } #[test] fn journal_state_merges_scopes_and_clears_only_confirmed_generations() {