fix(scanner): isolate and coalesce dirty usage journals (#8325)

* fix(scanner): isolate and coalesce dirty usage journals

Bind dirty replay records to a stable node owner so peers cannot overwrite
or delete another node's pending journal. Replace the unbounded command
queue with bounded coalesced state and revision-aware background saves.
Keep legacy records isolated and enforce replay byte limits with conservative
whole-bucket or unverified fallback. Add publication-stage diagnostics and
storage regressions for ownership, clearing, retry, and restart behavior.

Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>

* fix(scanner): collapse confirmed journal clear condition

Use a let-chain to retain the generation fence before removing a confirmed
bucket and updating its scope byte accounting. Satisfy collapsible_if
without changing journal clearing behavior.

Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>

---------

Co-authored-by: heihutu <heihutu@gmail.com>
Co-authored-by: zhi22915 <qiuzgang@gmail.com>
This commit is contained in:
Hauser
2026-10-04 03:02:22 +08:00
committed by GitHub
parent 5aca750a6e
commit 5e1bd498ce
8 changed files with 741 additions and 176 deletions
+3 -2
View File
@@ -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};
+47 -1
View File
@@ -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)
}
}
}
+10
View File
@@ -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"
);
}
}
+4 -2
View File
@@ -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};
+28 -1
View File
@@ -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<u64> {
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();
@@ -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"
);
@@ -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/<owner>.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.
+615 -170
View File
@@ -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<JournalCommand>,
shared: Option<Arc<SharedJournal>>,
}
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<PendingJournal>,
changed: Arc<Notify>,
#[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<ScannerDirtyUsageBucket>) {
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<ScannerDirtyUsageBucket>,
},
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<ScannerDirtyUsageBucket>) {
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<ScannerDurableDirtyUsageReplayRecord>,
}
fn durable_dirty_usage_owner(node_name: &str) -> Option<String> {
(!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<String, JournalBucketState>,
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<Vec<ScannerDurableDirtyUsageReplayEntry>> {
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<ECStore>) -> 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<ECStore>) -> JournalState {
let bytes = match super::read_config(store, DURABLE_DIRTY_USAGE_REPLAY_OBJECT).await {
async fn replay_durable_dirty_usage_journal(store: Arc<ECStore>, 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<ECStore>) -> 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::<ScannerDurableDirtyUsageReplayRecord>(&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<ECStore>) -> 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<ECStore>,
mut receiver: mpsc::Receiver<JournalCommand>,
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<u8>)> {
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<Vec<u8>, 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<ECStore>, owner: String, shared: Weak<SharedJournal>, changed: Arc<Notify>) {
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<ECStore>, 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<ECStore>, 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<ECStore>, 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<ECStore>) -> 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<String> {
#[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<rustfs_test_utils::TestECStoreEnv> = 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() {