mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-21 20:06:37 +00:00
refactor(scanner): split scanner_folder item actions and ledger (#6302)
* refactor(scanner): split scanner_folder item actions and ledger Split the 6345-line scanner_folder.rs (46% inline tests) into a canonical scanner_folder.rs + scanner_folder/ module tree with zero behavior change: - scanner_folder.rs (~2280): scan constants, alert cooldowns, metric accounting, resume ordering, tracing helpers, the FolderScanner struct with failed-object bookkeeping and the scan_folder traversal, and scan_data_folder - scanner_folder/item_actions.rs (~890): CachedFolder, the get-size failure policy, ScannerItem with apply_actions and the heal/ILM admission helpers - scanner_folder/ledger.rs (~280): the pending-scanner-heal ledger methods and their entry helpers (record/prune/clear-for-repaired/ retry) - scanner_folder/tests.rs (~2950): the inline test module as a child module The ScannerItem path used by scanner_io resolves through a root re-export, and every other crate path is unchanged. Cross-module items gain pub(super), whose scope equals the old single-module privacy domain. Code is moved verbatim apart from those markers, per-module import headers, and rustfmt re-wraps. Co-Authored-By: heihutu <heihutu@gmail.com> * fmt --------- Co-authored-by: heihutu <heihutu@gmail.com>
This commit is contained in:
File diff suppressed because it is too large
Load Diff
@@ -0,0 +1,891 @@
|
||||
// Copyright 2024 RustFS Team
|
||||
//
|
||||
// Licensed under the Apache License, Version 2.0 (the "License");
|
||||
// you may not use this file except in compliance with the License.
|
||||
// You may obtain a copy of the License at
|
||||
//
|
||||
// http://www.apache.org/licenses/LICENSE-2.0
|
||||
//
|
||||
// Unless required by applicable law or agreed to in writing, software
|
||||
// distributed under the License is distributed on an "AS IS" BASIS,
|
||||
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
/// Per-object scan actions: ScannerItem, the get-size failure policy, and the heal/ILM admission helpers.
|
||||
use super::*;
|
||||
|
||||
/// Cached folder information for scanning
|
||||
#[derive(Clone, Debug)]
|
||||
pub struct CachedFolder {
|
||||
pub name: String,
|
||||
pub parent: Option<DataUsageHash>,
|
||||
pub object_heal_prob_div: u32,
|
||||
}
|
||||
|
||||
/// Type alias for get size function
|
||||
pub type GetSizeFn = Box<dyn Fn(ScannerItem) -> Result<SizeSummary, StorageError> + Send + Sync>;
|
||||
|
||||
#[derive(Debug, PartialEq, Eq)]
|
||||
pub(super) enum GetSizeFailureAction {
|
||||
Skip,
|
||||
RecordFailed,
|
||||
HealMetadata { object: String },
|
||||
}
|
||||
|
||||
/// How the corrupt-metadata branch records the repair after attempting an
|
||||
/// MRF intent (backlog#1894 axis A).
|
||||
#[derive(Debug, PartialEq, Eq)]
|
||||
pub(super) enum CorruptMetadataRecording {
|
||||
/// Intent accepted: the MRF consumer owns the repair (High Metadata
|
||||
/// heal, durable after the journal's group-commit flush), so the
|
||||
/// immediate heal request is skipped — the manager would otherwise book
|
||||
/// two tasks for one target. A pending-ledger entry stays behind as the
|
||||
/// backstop for what the journal cannot cover on its own (a crash inside
|
||||
/// the flush window, or the consumer exhausting its admission attempts);
|
||||
/// the repaired-notice fanout (axis B) drops the entry once the repair
|
||||
/// lands.
|
||||
LedgerOnly,
|
||||
/// Intent rejected (feature disabled, channel uninitialized, or full):
|
||||
/// the historical immediate heal request plus the ledger entry.
|
||||
ImmediateAndLedger,
|
||||
}
|
||||
|
||||
pub(super) fn corrupt_metadata_recording(mrf_accepted: bool) -> CorruptMetadataRecording {
|
||||
if mrf_accepted {
|
||||
CorruptMetadataRecording::LedgerOnly
|
||||
} else {
|
||||
CorruptMetadataRecording::ImmediateAndLedger
|
||||
}
|
||||
}
|
||||
|
||||
pub(super) fn build_bucket_heal_request(bucket: String, priority: HealChannelPriority) -> HealChannelRequest {
|
||||
HealChannelRequest {
|
||||
bucket,
|
||||
priority,
|
||||
recreate_missing: Some(false),
|
||||
source: HealRequestSource::Scanner,
|
||||
..Default::default()
|
||||
}
|
||||
}
|
||||
|
||||
pub(super) fn build_object_heal_request(
|
||||
bucket: String,
|
||||
object: String,
|
||||
version_id: Option<String>,
|
||||
scan_mode: HealScanMode,
|
||||
priority: HealChannelPriority,
|
||||
) -> HealChannelRequest {
|
||||
HealChannelRequest {
|
||||
bucket,
|
||||
object_prefix: Some(object),
|
||||
object_version_id: version_id,
|
||||
priority,
|
||||
scan_mode: Some(scan_mode),
|
||||
remove_corrupted: Some(HEAL_DELETE_DANGLING),
|
||||
recreate_missing: Some(false),
|
||||
source: HealRequestSource::Scanner,
|
||||
..Default::default()
|
||||
}
|
||||
}
|
||||
|
||||
pub(super) fn resolve_object_heal_entry(
|
||||
entries: &MetaCacheEntries,
|
||||
resolver: MetadataResolutionParams,
|
||||
) -> Option<MetaCacheEntry> {
|
||||
if let Some(entry) = entries.resolve(resolver) {
|
||||
return entry.is_object().then_some(entry);
|
||||
}
|
||||
|
||||
entries
|
||||
.as_ref()
|
||||
.iter()
|
||||
.flatten()
|
||||
.find(|entry| entry.is_object() && !entry.name.ends_with(SLASH_SEPARATOR))
|
||||
.cloned()
|
||||
}
|
||||
|
||||
pub(super) fn is_missing_path_disk_error(err: &DiskError) -> bool {
|
||||
matches!(err, DiskError::FileNotFound | DiskError::FileVersionNotFound | DiskError::VolumeNotFound)
|
||||
}
|
||||
|
||||
pub(super) fn disk_errors_are_only_missing_paths(errs: &[Option<DiskError>]) -> bool {
|
||||
let mut saw_missing_path = false;
|
||||
for err in errs.iter().flatten() {
|
||||
if !is_missing_path_disk_error(err) {
|
||||
return false;
|
||||
}
|
||||
saw_missing_path = true;
|
||||
}
|
||||
saw_missing_path
|
||||
}
|
||||
|
||||
pub(super) fn heal_priority_label(priority: HealChannelPriority) -> &'static str {
|
||||
match priority {
|
||||
HealChannelPriority::Low => "low",
|
||||
HealChannelPriority::Normal => "normal",
|
||||
HealChannelPriority::High => "high",
|
||||
HealChannelPriority::Critical => "critical",
|
||||
}
|
||||
}
|
||||
|
||||
pub(super) fn describe_heal_admission(result: HealAdmissionResult) -> String {
|
||||
match result {
|
||||
HealAdmissionResult::Accepted | HealAdmissionResult::Merged => result.result_label().to_string(),
|
||||
HealAdmissionResult::Full => "queue_full".to_string(),
|
||||
HealAdmissionResult::Dropped(reason) => format!("dropped:{}", reason.as_str()),
|
||||
}
|
||||
}
|
||||
|
||||
pub(super) fn record_high_priority_heal_escalation(
|
||||
candidate_type: &'static str,
|
||||
priority: HealChannelPriority,
|
||||
result: HealAdmissionResult,
|
||||
) {
|
||||
counter!(
|
||||
"rustfs_heal_candidate_priority_reject_total",
|
||||
"type" => candidate_type.to_string(),
|
||||
"priority" => heal_priority_label(priority).to_string(),
|
||||
"result" => result.result_label().to_string(),
|
||||
"reason" => result.reason_label().to_string()
|
||||
)
|
||||
.increment(1);
|
||||
}
|
||||
|
||||
pub(super) fn build_high_priority_heal_admission_error(
|
||||
candidate_type: &'static str,
|
||||
bucket: &str,
|
||||
object: Option<&str>,
|
||||
priority: HealChannelPriority,
|
||||
result: HealAdmissionResult,
|
||||
) -> ScannerError {
|
||||
let object_text = object.map(|object| format!(", object='{object}'")).unwrap_or_default();
|
||||
ScannerError::Other(format!(
|
||||
"high-priority heal request was not admitted: type={candidate_type}, bucket='{bucket}'{object_text}, priority={}, admission={}",
|
||||
heal_priority_label(priority),
|
||||
describe_heal_admission(result)
|
||||
))
|
||||
}
|
||||
|
||||
pub(super) fn record_heal_candidate_admission(
|
||||
candidate_type: &'static str,
|
||||
priority: HealChannelPriority,
|
||||
result: HealAdmissionResult,
|
||||
) {
|
||||
counter!(
|
||||
"rustfs_heal_candidate_enqueue_total",
|
||||
"type" => candidate_type.to_string(),
|
||||
"priority" => heal_priority_label(priority).to_string(),
|
||||
"result" => result.result_label().to_string()
|
||||
)
|
||||
.increment(1);
|
||||
|
||||
if matches!(result, HealAdmissionResult::Merged) {
|
||||
counter!(
|
||||
"rustfs_heal_candidate_merge_total",
|
||||
"type" => candidate_type.to_string()
|
||||
)
|
||||
.increment(1);
|
||||
}
|
||||
|
||||
if let HealAdmissionResult::Dropped(reason) = result {
|
||||
counter!(
|
||||
"rustfs_heal_candidate_drop_total",
|
||||
"type" => candidate_type.to_string(),
|
||||
"reason" => reason.as_str().to_string()
|
||||
)
|
||||
.increment(1);
|
||||
}
|
||||
}
|
||||
|
||||
pub(super) async fn send_scanner_heal_request(
|
||||
candidate_type: &'static str,
|
||||
request: HealChannelRequest,
|
||||
) -> Result<HealAdmissionResult, ScannerError> {
|
||||
let priority = request.priority;
|
||||
let trace_context = scanner_heal_candidate_trace_context(&request);
|
||||
match send_heal_request_with_admission(request).await {
|
||||
Ok(result) => {
|
||||
record_heal_candidate_admission(candidate_type, priority, result);
|
||||
if let Some(trace_context) = trace_context.as_ref() {
|
||||
emit_scanner_heal_candidate_trace(ScannerHealCandidateTrace {
|
||||
candidate_type,
|
||||
bucket: &trace_context.bucket,
|
||||
object: trace_context.object.as_deref(),
|
||||
version_id: trace_context.version_id.as_deref(),
|
||||
priority,
|
||||
scan_mode: trace_context.scan_mode,
|
||||
result: Ok(result),
|
||||
started_at: trace_context.started_at,
|
||||
});
|
||||
}
|
||||
Ok(result)
|
||||
}
|
||||
Err(err) => {
|
||||
counter!(
|
||||
"rustfs_heal_candidate_enqueue_total",
|
||||
"type" => candidate_type.to_string(),
|
||||
"priority" => heal_priority_label(priority).to_string(),
|
||||
"result" => "channel_error".to_string()
|
||||
)
|
||||
.increment(1);
|
||||
if let Some(trace_context) = trace_context.as_ref() {
|
||||
emit_scanner_heal_candidate_trace(ScannerHealCandidateTrace {
|
||||
candidate_type,
|
||||
bucket: &trace_context.bucket,
|
||||
object: trace_context.object.as_deref(),
|
||||
version_id: trace_context.version_id.as_deref(),
|
||||
priority,
|
||||
scan_mode: trace_context.scan_mode,
|
||||
result: Err(err.as_str()),
|
||||
started_at: trace_context.started_at,
|
||||
});
|
||||
}
|
||||
Err(ScannerError::Other(err))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Scanner item representing a file during scanning
|
||||
#[derive(Clone, Debug)]
|
||||
pub struct ScannerItem {
|
||||
pub path: String,
|
||||
pub bucket: String,
|
||||
pub prefix: String,
|
||||
pub object_name: String,
|
||||
pub file_type: FileType,
|
||||
pub lifecycle: Option<Arc<BucketLifecycleConfiguration>>,
|
||||
pub object_lock: Option<Arc<ObjectLockConfiguration>>,
|
||||
pub replication: Option<Arc<ReplicationConfig>>,
|
||||
pub heal_enabled: bool,
|
||||
pub heal_bitrot: bool,
|
||||
pub debug: bool,
|
||||
}
|
||||
|
||||
impl ScannerItem {
|
||||
/// Get the object path (prefix + object_name)
|
||||
pub fn object_path(&self) -> String {
|
||||
if self.prefix.is_empty() {
|
||||
self.object_name.clone()
|
||||
} else {
|
||||
path_join_buf(&[&self.prefix, &self.object_name])
|
||||
}
|
||||
}
|
||||
|
||||
/// Transform meta directory by splitting prefix and extracting object name
|
||||
/// This converts a directory path like "bucket/dir1/dir2/file" to prefix="bucket/dir1/dir2" and object_name="file"
|
||||
pub fn transform_meta_dir(&mut self) {
|
||||
let prefix = self.prefix.clone(); // Clone to avoid borrow checker issues
|
||||
let split: Vec<&str> = prefix.split(SLASH_SEPARATOR).collect();
|
||||
|
||||
if split.len() > 1 {
|
||||
let prefix_parts: Vec<&str> = split[..split.len() - 1].to_vec();
|
||||
self.prefix = path_join_buf(&prefix_parts);
|
||||
} else {
|
||||
self.prefix = String::new();
|
||||
}
|
||||
|
||||
// Object name is the last element
|
||||
self.object_name = split.last().unwrap_or(&"").to_string();
|
||||
}
|
||||
|
||||
pub(super) fn metadata_object_path(&self) -> String {
|
||||
let mut item = self.clone();
|
||||
item.transform_meta_dir();
|
||||
item.object_path()
|
||||
}
|
||||
|
||||
pub async fn apply_actions(
|
||||
&mut self,
|
||||
object_infos: Vec<ObjectInfo>,
|
||||
lock_retention: Option<Arc<ObjectLockConfiguration>>,
|
||||
versioning_config: VersioningConfiguration,
|
||||
size_summary: &mut SizeSummary,
|
||||
) {
|
||||
if object_infos.is_empty() {
|
||||
debug!(
|
||||
target: "rustfs::scanner::folder",
|
||||
event = EVENT_SCANNER_LIFECYCLE_ACTION,
|
||||
component = LOG_COMPONENT_SCANNER,
|
||||
subsystem = LOG_SUBSYSTEM_LIFECYCLE,
|
||||
object_path = %self.object_path(),
|
||||
state = "no_object_versions",
|
||||
"Scanner lifecycle action skipped"
|
||||
);
|
||||
return;
|
||||
}
|
||||
debug!(
|
||||
target: "rustfs::scanner::folder",
|
||||
event = EVENT_SCANNER_LIFECYCLE_ACTION,
|
||||
component = LOG_COMPONENT_SCANNER,
|
||||
subsystem = LOG_SUBSYSTEM_LIFECYCLE,
|
||||
object_path = %self.object_path(),
|
||||
state = "started",
|
||||
"Scanner lifecycle evaluation started"
|
||||
);
|
||||
|
||||
// `versioning_config` is resolved once per object by the caller
|
||||
// (`get_size`) and handed in; only `prefix_enabled` is consulted here.
|
||||
|
||||
let Some(lifecycle) = self.lifecycle.as_ref() else {
|
||||
let mut cumulative_size = 0;
|
||||
for oi in object_infos.iter() {
|
||||
let actual_size = match oi.get_actual_size() {
|
||||
Ok(size) => size,
|
||||
Err(_) => {
|
||||
warn!(
|
||||
target: "rustfs::scanner::folder",
|
||||
event = EVENT_SCANNER_LIFECYCLE_ACTION,
|
||||
component = LOG_COMPONENT_SCANNER,
|
||||
subsystem = LOG_SUBSYSTEM_LIFECYCLE,
|
||||
bucket = %self.bucket,
|
||||
object = %oi.name,
|
||||
state = "size_lookup_failed",
|
||||
"Scanner lifecycle action used fallback size"
|
||||
);
|
||||
continue;
|
||||
}
|
||||
};
|
||||
|
||||
let size = self.heal_actions(oi, actual_size, size_summary).await;
|
||||
|
||||
size_summary.actions_accounting(oi, size, actual_size);
|
||||
|
||||
cumulative_size += size;
|
||||
}
|
||||
|
||||
self.alert_excessive_versions(object_infos.len(), cumulative_size);
|
||||
|
||||
debug!(
|
||||
target: "rustfs::scanner::folder",
|
||||
event = EVENT_SCANNER_LIFECYCLE_ACTION,
|
||||
component = LOG_COMPONENT_SCANNER,
|
||||
subsystem = LOG_SUBSYSTEM_LIFECYCLE,
|
||||
object_path = %self.object_path(),
|
||||
state = "no_lifecycle_config",
|
||||
"Scanner lifecycle action finished without lifecycle rules"
|
||||
);
|
||||
return;
|
||||
};
|
||||
|
||||
let object_opts = object_infos
|
||||
.iter()
|
||||
.map(crate::ecstore_object_opts_from_object_info)
|
||||
.collect::<Vec<ObjectOpts>>();
|
||||
|
||||
let events = match Evaluator::new(lifecycle.clone())
|
||||
.with_lock_retention(lock_retention)
|
||||
.with_replication_config(scanner_replication_config_for_lifecycle_eval(self.replication.clone()))
|
||||
.eval(&object_opts)
|
||||
.await
|
||||
{
|
||||
Ok(events) => events,
|
||||
Err(e) => {
|
||||
warn!(
|
||||
target: "rustfs::scanner::folder",
|
||||
event = EVENT_SCANNER_LIFECYCLE_ACTION,
|
||||
component = LOG_COMPONENT_SCANNER,
|
||||
subsystem = LOG_SUBSYSTEM_LIFECYCLE,
|
||||
object_path = %self.object_path(),
|
||||
state = "evaluate_failed",
|
||||
error = %e,
|
||||
"Scanner lifecycle action evaluation failed"
|
||||
);
|
||||
return;
|
||||
}
|
||||
};
|
||||
|
||||
// Every version handed to the evaluator counts as an ILM-checked
|
||||
// version, whether or not a rule matched (NoneAction included). This is
|
||||
// reached only for buckets with lifecycle rules; it feeds
|
||||
// rustfs_ilm_versions_scanned_total via the Lifecycle source's checked
|
||||
// counter.
|
||||
global_metrics().record_scanner_source_checked(ScannerWorkSource::Lifecycle, object_opts.len() as u64);
|
||||
|
||||
let mut to_delete_objs: Vec<ObjectToDelete> = Vec::new();
|
||||
let mut noncurrent_events: Vec<Event> = Vec::new();
|
||||
let mut noncurrent_accounting: Vec<PendingScannerAccounting<'_>> = Vec::new();
|
||||
let mut cumulative_size = 0;
|
||||
let mut remaining_versions = object_infos.len();
|
||||
'eventLoop: {
|
||||
for (i, event) in events.iter().enumerate() {
|
||||
let oi = &object_infos[i];
|
||||
let actual_size = match oi.get_actual_size() {
|
||||
Ok(size) => size,
|
||||
Err(_) => {
|
||||
warn!(
|
||||
target: "rustfs::scanner::folder",
|
||||
event = EVENT_SCANNER_LIFECYCLE_ACTION,
|
||||
component = LOG_COMPONENT_SCANNER,
|
||||
subsystem = LOG_SUBSYSTEM_LIFECYCLE,
|
||||
bucket = %self.bucket,
|
||||
object = %oi.name,
|
||||
state = "size_lookup_failed",
|
||||
"Scanner lifecycle action used fallback size"
|
||||
);
|
||||
0
|
||||
}
|
||||
};
|
||||
|
||||
let mut size = actual_size;
|
||||
let mut account_now = true;
|
||||
|
||||
match event.action {
|
||||
IlmAction::DeleteAllVersionsAction | IlmAction::DelMarkerDeleteAllVersionsAction => {
|
||||
debug!(
|
||||
target: "rustfs::scanner::folder",
|
||||
event = EVENT_SCANNER_LIFECYCLE_ACTION,
|
||||
component = LOG_COMPONENT_SCANNER,
|
||||
subsystem = LOG_SUBSYSTEM_LIFECYCLE,
|
||||
bucket = %self.bucket,
|
||||
object = %oi.name,
|
||||
action = %event.action,
|
||||
state = "apply_expiry_rule",
|
||||
"Scanner lifecycle action dispatched"
|
||||
);
|
||||
let done_ilm = Metrics::time_ilm(event.action);
|
||||
let trace_started_at = trace_start_instant();
|
||||
let queued = apply_expiry_rule(event, &LcEventSrc::Scanner, oi).await;
|
||||
emit_scanner_ilm_action_trace(&self.bucket, &oi.name, event.action, 1, queued, trace_started_at);
|
||||
if record_scanner_ilm_action_if_queued(global_metrics(), event.action, 1, queued) {
|
||||
done_ilm(1)();
|
||||
remaining_versions = 0;
|
||||
} else {
|
||||
PendingScannerAccounting {
|
||||
object: oi,
|
||||
retained_size: actual_size,
|
||||
expired_size: 0,
|
||||
}
|
||||
.apply(size_summary, &mut cumulative_size, false);
|
||||
for retained in object_infos.iter().skip(i + 1) {
|
||||
let retained_size = match retained.get_actual_size() {
|
||||
Ok(size) => size,
|
||||
Err(_) => {
|
||||
warn!(
|
||||
target: "rustfs::scanner::folder",
|
||||
event = EVENT_SCANNER_LIFECYCLE_ACTION,
|
||||
component = LOG_COMPONENT_SCANNER,
|
||||
subsystem = LOG_SUBSYSTEM_LIFECYCLE,
|
||||
bucket = %self.bucket,
|
||||
object = %retained.name,
|
||||
state = "size_lookup_failed",
|
||||
"Scanner lifecycle action used fallback size"
|
||||
);
|
||||
0
|
||||
}
|
||||
};
|
||||
PendingScannerAccounting {
|
||||
object: retained,
|
||||
retained_size,
|
||||
expired_size: 0,
|
||||
}
|
||||
.apply(size_summary, &mut cumulative_size, false);
|
||||
}
|
||||
}
|
||||
break 'eventLoop;
|
||||
}
|
||||
|
||||
IlmAction::DeleteAction | IlmAction::DeleteRestoredAction | IlmAction::DeleteRestoredVersionAction => {
|
||||
debug!(
|
||||
target: "rustfs::scanner::folder",
|
||||
event = EVENT_SCANNER_LIFECYCLE_ACTION,
|
||||
component = LOG_COMPONENT_SCANNER,
|
||||
subsystem = LOG_SUBSYSTEM_LIFECYCLE,
|
||||
bucket = %self.bucket,
|
||||
object = %oi.name,
|
||||
action = %event.action,
|
||||
state = "apply_expiry_rule",
|
||||
"Scanner lifecycle action dispatched"
|
||||
);
|
||||
let done_ilm = Metrics::time_ilm(event.action);
|
||||
let trace_started_at = trace_start_instant();
|
||||
let queued = apply_expiry_rule(event, &LcEventSrc::Scanner, oi).await;
|
||||
emit_scanner_ilm_action_trace(&self.bucket, &oi.name, event.action, 1, queued, trace_started_at);
|
||||
if record_scanner_ilm_action_if_queued(global_metrics(), event.action, 1, queued) {
|
||||
done_ilm(1)();
|
||||
if !versioning_config.prefix_enabled(&self.object_path()) && event.action == IlmAction::DeleteAction {
|
||||
remaining_versions -= 1;
|
||||
size = 0;
|
||||
}
|
||||
}
|
||||
}
|
||||
IlmAction::DeleteVersionAction => {
|
||||
if let Some(opt) = object_opts.get(i) {
|
||||
to_delete_objs.push(ObjectToDelete {
|
||||
object_name: opt.name.clone(),
|
||||
version_id: opt.version_id,
|
||||
..Default::default()
|
||||
});
|
||||
noncurrent_accounting.push(PendingScannerAccounting {
|
||||
object: oi,
|
||||
retained_size: actual_size,
|
||||
expired_size: 0,
|
||||
});
|
||||
account_now = false;
|
||||
}
|
||||
noncurrent_events.push(event.clone());
|
||||
}
|
||||
IlmAction::TransitionAction | IlmAction::TransitionVersionAction => {
|
||||
debug!(
|
||||
target: "rustfs::scanner::folder",
|
||||
event = EVENT_SCANNER_LIFECYCLE_ACTION,
|
||||
component = LOG_COMPONENT_SCANNER,
|
||||
subsystem = LOG_SUBSYSTEM_LIFECYCLE,
|
||||
bucket = %self.bucket,
|
||||
object = %oi.name,
|
||||
action = %event.action,
|
||||
state = "apply_transition_rule",
|
||||
"Scanner lifecycle action dispatched"
|
||||
);
|
||||
let done_ilm = Metrics::time_ilm(event.action);
|
||||
let trace_started_at = trace_start_instant();
|
||||
let queued = apply_transition_rule(event, &LcEventSrc::Scanner, oi).await;
|
||||
emit_scanner_ilm_action_trace(&self.bucket, &oi.name, event.action, 1, queued, trace_started_at);
|
||||
if record_scanner_ilm_action_if_queued(global_metrics(), event.action, 1, queued) {
|
||||
done_ilm(1)();
|
||||
}
|
||||
}
|
||||
|
||||
IlmAction::NoneAction | IlmAction::ActionCount => {
|
||||
size = self.heal_actions(oi, actual_size, size_summary).await;
|
||||
}
|
||||
}
|
||||
|
||||
if account_now {
|
||||
size_summary.actions_accounting(oi, size, actual_size);
|
||||
cumulative_size += size;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if !to_delete_objs.is_empty()
|
||||
&& let Some(event) = noncurrent_events.first().cloned()
|
||||
{
|
||||
let action = event.action;
|
||||
let count = u64::try_from(to_delete_objs.len()).unwrap_or(u64::MAX);
|
||||
let done_ilm = Metrics::time_ilm(action);
|
||||
let trace_started_at = trace_start_instant();
|
||||
let queued = enqueue_runtime_newer_noncurrent(&self.bucket, to_delete_objs, event, &LcEventSrc::Scanner).await;
|
||||
if let Some(trace_started_at) = trace_started_at {
|
||||
let state = if queued { "queued" } else { "not_queued" };
|
||||
trace_emit(|| {
|
||||
TraceEvent::new(TraceKind::Scanner, TraceFunc::ScannerIlmAction)
|
||||
.with_bucket(self.bucket.as_str())
|
||||
.with_object(self.object_path())
|
||||
.with_duration(trace_started_at.elapsed())
|
||||
.with_attr("state", state)
|
||||
.with_attr("action", action.as_str())
|
||||
.with_attr("count", count)
|
||||
.with_attr("queued", queued)
|
||||
});
|
||||
}
|
||||
if record_scanner_ilm_action_if_queued(global_metrics(), action, count, queued) {
|
||||
done_ilm(count)();
|
||||
remaining_versions = remaining_versions.saturating_sub(noncurrent_accounting.len());
|
||||
}
|
||||
for pending in noncurrent_accounting {
|
||||
pending.apply(size_summary, &mut cumulative_size, queued);
|
||||
}
|
||||
}
|
||||
self.alert_excessive_versions(remaining_versions, cumulative_size);
|
||||
}
|
||||
|
||||
pub(super) async fn heal_actions(&mut self, oi: &ObjectInfo, actual_size: i64, size_summary: &mut SizeSummary) -> i64 {
|
||||
if self.heal_enabled {
|
||||
self.enqueue_heal(oi).await;
|
||||
}
|
||||
|
||||
self.heal_replication(oi, size_summary).await;
|
||||
|
||||
actual_size
|
||||
}
|
||||
|
||||
pub(super) async fn heal_replication(&mut self, oi: &ObjectInfo, size_summary: &mut SizeSummary) {
|
||||
if oi.version_id.is_none_or(|version| version.is_nil()) && !oi.delete_marker && oi.version_purge_status.is_empty() {
|
||||
return;
|
||||
}
|
||||
|
||||
let Some(replication) = self.replication.clone() else {
|
||||
return;
|
||||
};
|
||||
|
||||
let done_replication = Metrics::time(Metric::CheckReplication);
|
||||
let replication_result = queue_replication_heal(&oi.bucket, oi.clone(), (*replication).clone(), 0).await;
|
||||
done_replication();
|
||||
let roi = replication_result.object_info;
|
||||
record_scanner_replication_admission(global_metrics(), &roi, replication_result.admission);
|
||||
if !Self::should_account_replication_stats(oi) {
|
||||
return;
|
||||
}
|
||||
|
||||
for (arn, target_status) in roi.target_statuses.iter() {
|
||||
if !size_summary.repl_target_stats.contains_key(arn.as_str()) {
|
||||
size_summary
|
||||
.repl_target_stats
|
||||
.insert(arn.clone(), ReplTargetSizeSummary::default());
|
||||
}
|
||||
|
||||
if let Some(repl_target_size_summary) = size_summary.repl_target_stats.get_mut(arn.as_str()) {
|
||||
match target_status {
|
||||
ReplicationStatusType::Pending => {
|
||||
repl_target_size_summary.pending_size = repl_target_size_summary.pending_size.saturating_add(roi.size);
|
||||
repl_target_size_summary.pending_count = repl_target_size_summary.pending_count.saturating_add(1);
|
||||
size_summary.pending_size = size_summary.pending_size.saturating_add(roi.size);
|
||||
size_summary.pending_count = size_summary.pending_count.saturating_add(1);
|
||||
}
|
||||
ReplicationStatusType::Failed => {
|
||||
repl_target_size_summary.failed_size = repl_target_size_summary.failed_size.saturating_add(roi.size);
|
||||
repl_target_size_summary.failed_count = repl_target_size_summary.failed_count.saturating_add(1);
|
||||
size_summary.failed_size = size_summary.failed_size.saturating_add(roi.size);
|
||||
size_summary.failed_count = size_summary.failed_count.saturating_add(1);
|
||||
}
|
||||
ReplicationStatusType::Completed | ReplicationStatusType::CompletedLegacy => {
|
||||
repl_target_size_summary.replicated_size =
|
||||
repl_target_size_summary.replicated_size.saturating_add(roi.size);
|
||||
repl_target_size_summary.replicated_count = repl_target_size_summary.replicated_count.saturating_add(1);
|
||||
size_summary.replicated_size = size_summary.replicated_size.saturating_add(roi.size);
|
||||
size_summary.replicated_count = size_summary.replicated_count.saturating_add(1);
|
||||
}
|
||||
_ => {}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if oi.replication_status == ReplicationStatusType::Replica {
|
||||
size_summary.replica_size = size_summary.replica_size.saturating_add(roi.size);
|
||||
size_summary.replica_count = size_summary.replica_count.saturating_add(1);
|
||||
}
|
||||
}
|
||||
|
||||
pub(super) fn should_account_replication_stats(oi: &ObjectInfo) -> bool {
|
||||
!oi.delete_marker && oi.version_purge_status.is_empty()
|
||||
}
|
||||
|
||||
pub(super) async fn enqueue_heal(&mut self, oi: &ObjectInfo) {
|
||||
let done_heal = Metrics::time(Metric::HealAbandonedObject);
|
||||
let object = if oi.name.is_empty() {
|
||||
self.object_path()
|
||||
} else {
|
||||
oi.name.clone()
|
||||
};
|
||||
debug!(
|
||||
target: "rustfs::scanner::folder",
|
||||
event = EVENT_SCANNER_HEAL_ADMISSION,
|
||||
component = LOG_COMPONENT_SCANNER,
|
||||
subsystem = LOG_SUBSYSTEM_HEAL,
|
||||
bucket = %self.bucket,
|
||||
object = %object,
|
||||
version_id = %oi.version_id.unwrap_or_default(),
|
||||
state = "request_started",
|
||||
"Scanner heal admission started"
|
||||
);
|
||||
|
||||
let now = OffsetDateTime::now_utc();
|
||||
let scan_mode = effective_object_heal_scan_mode(self.heal_bitrot, oi.mod_time, now);
|
||||
if self.heal_bitrot && scan_mode != HealScanMode::Deep {
|
||||
let cooldown = deep_verify_cooldown();
|
||||
let age_secs = oi.mod_time.map(|mod_time| {
|
||||
let age = now - mod_time;
|
||||
age.whole_seconds().max(0)
|
||||
});
|
||||
debug!(
|
||||
target: "rustfs::scanner::folder",
|
||||
event = EVENT_SCANNER_HEAL_ADMISSION,
|
||||
component = LOG_COMPONENT_SCANNER,
|
||||
subsystem = LOG_SUBSYSTEM_HEAL,
|
||||
bucket = %self.bucket,
|
||||
object = %object,
|
||||
version_id = %oi.version_id.unwrap_or_default(),
|
||||
object_age_secs = age_secs.unwrap_or_default(),
|
||||
cooldown_secs = cooldown.as_secs(),
|
||||
original_scan_mode = %HealScanMode::Deep.as_str(),
|
||||
effective_scan_mode = %scan_mode.as_str(),
|
||||
state = "downgraded_to_normal",
|
||||
"Scanner heal deep scan downgraded"
|
||||
);
|
||||
}
|
||||
|
||||
let result = send_scanner_heal_request(
|
||||
"object",
|
||||
build_object_heal_request(
|
||||
self.bucket.clone(),
|
||||
object.clone(),
|
||||
oi.version_id
|
||||
.and_then(|v| if v.is_nil() { None } else { Some(v.to_string()) }),
|
||||
scan_mode,
|
||||
HealChannelPriority::Low,
|
||||
),
|
||||
)
|
||||
.await;
|
||||
|
||||
let admission = result.as_ref().copied().map_err(|_| ());
|
||||
let admitted = record_scanner_heal_admission(global_metrics(), scan_mode, admission);
|
||||
match result {
|
||||
Ok(HealAdmissionResult::Accepted | HealAdmissionResult::Merged) => {}
|
||||
Ok(result @ (HealAdmissionResult::Full | HealAdmissionResult::Dropped(_))) => {
|
||||
warn!(
|
||||
target: "rustfs::scanner::folder",
|
||||
event = EVENT_SCANNER_HEAL_ADMISSION,
|
||||
component = LOG_COMPONENT_SCANNER,
|
||||
subsystem = LOG_SUBSYSTEM_HEAL,
|
||||
bucket = %self.bucket,
|
||||
object = %object,
|
||||
admission = %describe_heal_admission(result),
|
||||
state = "not_admitted",
|
||||
"Scanner heal admission rejected low-priority request"
|
||||
);
|
||||
}
|
||||
Err(e) => warn!(
|
||||
target: "rustfs::scanner::folder",
|
||||
event = EVENT_SCANNER_HEAL_ADMISSION,
|
||||
component = LOG_COMPONENT_SCANNER,
|
||||
subsystem = LOG_SUBSYSTEM_HEAL,
|
||||
bucket = %self.bucket,
|
||||
object = %object,
|
||||
state = "submit_failed",
|
||||
error = %e,
|
||||
"Scanner heal admission submission failed"
|
||||
),
|
||||
}
|
||||
if admitted {
|
||||
done_heal();
|
||||
}
|
||||
}
|
||||
|
||||
pub(super) fn alert_excessive_versions(&self, remaining_versions: usize, cumulative_size: i64) {
|
||||
ensure_scanner_alert_metrics_registered();
|
||||
let (too_many_versions, too_large_versions) = should_alert_excessive_versions(remaining_versions, cumulative_size);
|
||||
// Threshold check first so healthy objects never pay for the
|
||||
// object-path allocation below.
|
||||
if !too_many_versions && !too_large_versions {
|
||||
return;
|
||||
}
|
||||
let object_path = self.object_path();
|
||||
if too_many_versions {
|
||||
global_metrics().record_scanner_source_executed(ScannerWorkSource::Alerts, 1);
|
||||
counter!(
|
||||
METRIC_SCANNER_EXCESS_OBJECT_VERSIONS_TOTAL,
|
||||
"bucket" => self.bucket.clone()
|
||||
)
|
||||
.increment(1);
|
||||
if scanner_alert_emission_allows(ScannerAlertKind::ManyVersions, &self.bucket, &object_path, scanner_alert_cooldown())
|
||||
{
|
||||
emit_scanner_alert_event(
|
||||
EVENT_SCANNER_MANY_VERSIONS,
|
||||
&self.bucket,
|
||||
&object_path,
|
||||
cumulative_size,
|
||||
&[
|
||||
("versions", remaining_versions.to_string()),
|
||||
("threshold", scanner_excess_versions_threshold().to_string()),
|
||||
],
|
||||
);
|
||||
}
|
||||
warn!(
|
||||
target: "rustfs::scanner::folder",
|
||||
event = EVENT_SCANNER_ALERT_STATE,
|
||||
component = LOG_COMPONENT_SCANNER,
|
||||
subsystem = LOG_SUBSYSTEM_FOLDER,
|
||||
bucket = %self.bucket,
|
||||
object = %object_path,
|
||||
versions = remaining_versions,
|
||||
threshold = scanner_excess_versions_threshold(),
|
||||
state = "excess_versions",
|
||||
"Scanner alert recorded excessive retained versions"
|
||||
);
|
||||
}
|
||||
if too_large_versions {
|
||||
global_metrics().record_scanner_source_executed(ScannerWorkSource::Alerts, 1);
|
||||
counter!(
|
||||
METRIC_SCANNER_EXCESS_OBJECT_VERSION_SIZE_TOTAL,
|
||||
"bucket" => self.bucket.clone()
|
||||
)
|
||||
.increment(1);
|
||||
if scanner_alert_emission_allows(
|
||||
ScannerAlertKind::LargeVersions,
|
||||
&self.bucket,
|
||||
&object_path,
|
||||
scanner_alert_cooldown(),
|
||||
) {
|
||||
emit_scanner_alert_event(
|
||||
EVENT_SCANNER_LARGE_VERSIONS,
|
||||
&self.bucket,
|
||||
&object_path,
|
||||
cumulative_size,
|
||||
&[
|
||||
("versions", remaining_versions.to_string()),
|
||||
("cumulativeSize", cumulative_size.to_string()),
|
||||
("threshold", scanner_excess_version_size_threshold().to_string()),
|
||||
],
|
||||
);
|
||||
}
|
||||
warn!(
|
||||
target: "rustfs::scanner::folder",
|
||||
event = EVENT_SCANNER_ALERT_STATE,
|
||||
component = LOG_COMPONENT_SCANNER,
|
||||
subsystem = LOG_SUBSYSTEM_FOLDER,
|
||||
bucket = %self.bucket,
|
||||
object = %object_path,
|
||||
versions = remaining_versions,
|
||||
cumulative_size,
|
||||
threshold = scanner_excess_version_size_threshold(),
|
||||
state = "excess_version_size",
|
||||
"Scanner alert recorded excessive retained version size"
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
pub(super) fn classify_get_size_failure(item: &ScannerItem, err: &StorageError) -> GetSizeFailureAction {
|
||||
if matches!(err, StorageError::Io(io) if io.to_string() == SCANNER_SKIP_FILE_ERROR) {
|
||||
return GetSizeFailureAction::Skip;
|
||||
}
|
||||
|
||||
if is_scanner_metadata_corrupt_error(err) {
|
||||
return GetSizeFailureAction::HealMetadata {
|
||||
object: item.metadata_object_path(),
|
||||
};
|
||||
}
|
||||
|
||||
if is_scanner_metadata_transient_error(err) {
|
||||
return GetSizeFailureAction::RecordFailed;
|
||||
}
|
||||
|
||||
GetSizeFailureAction::RecordFailed
|
||||
}
|
||||
|
||||
pub(super) async fn contains_erasure_part_file(path: &str) -> Result<bool, ScannerError> {
|
||||
let mut entries = match tokio::fs::read_dir(path).await {
|
||||
Ok(entries) => entries,
|
||||
Err(err) if matches!(err.kind(), ErrorKind::NotFound | ErrorKind::NotADirectory) => return Ok(false),
|
||||
Err(err) => return Err(ScannerError::Io(err)),
|
||||
};
|
||||
|
||||
for _ in 0..ERASURE_DATA_DIR_PROBE_ENTRY_LIMIT {
|
||||
let entry = match entries.next_entry().await {
|
||||
Ok(Some(entry)) => entry,
|
||||
Ok(None) => return Ok(false),
|
||||
Err(err) if matches!(err.kind(), ErrorKind::NotFound | ErrorKind::NotADirectory) => return Ok(false),
|
||||
Err(err) => return Err(ScannerError::Io(err)),
|
||||
};
|
||||
let file_name = entry.file_name();
|
||||
let Some(part_number) = file_name
|
||||
.to_str()
|
||||
.and_then(|name| name.strip_prefix("part."))
|
||||
.and_then(|number| number.parse::<u32>().ok())
|
||||
else {
|
||||
continue;
|
||||
};
|
||||
if part_number == 0 {
|
||||
continue;
|
||||
}
|
||||
|
||||
match entry.file_type().await {
|
||||
Ok(file_type) if file_type.is_file() => return Ok(true),
|
||||
Ok(_) => {}
|
||||
Err(err) if matches!(err.kind(), ErrorKind::NotFound | ErrorKind::TooManyLinks) => {}
|
||||
Err(err) => return Err(ScannerError::Io(err)),
|
||||
}
|
||||
}
|
||||
|
||||
Ok(false)
|
||||
}
|
||||
@@ -0,0 +1,282 @@
|
||||
// Copyright 2024 RustFS Team
|
||||
//
|
||||
// Licensed under the Apache License, Version 2.0 (the "License");
|
||||
// you may not use this file except in compliance with the License.
|
||||
// You may obtain a copy of the License at
|
||||
//
|
||||
// http://www.apache.org/licenses/LICENSE-2.0
|
||||
//
|
||||
// Unless required by applicable law or agreed to in writing, software
|
||||
// distributed under the License is distributed on an "AS IS" BASIS,
|
||||
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
/// The pending-scanner-heal ledger: durable heal intents recorded during scans and retried after MRF consumption.
|
||||
use super::*;
|
||||
|
||||
impl FolderScanner {
|
||||
pub(super) fn sync_pending_heals(&mut self) {
|
||||
self.update_cache.info.pending_heals = self.new_cache.info.pending_heals.clone();
|
||||
self.pending_heals_changed = true;
|
||||
}
|
||||
|
||||
pub(super) fn clear_pending_scanner_heal(
|
||||
&mut self,
|
||||
kind: PendingScannerHealKind,
|
||||
bucket: &str,
|
||||
object: Option<&str>,
|
||||
version_id: Option<&str>,
|
||||
) {
|
||||
let before = self.new_cache.info.pending_heals.len();
|
||||
self.new_cache
|
||||
.info
|
||||
.pending_heals
|
||||
.retain(|entry| !pending_scanner_heal_matches(entry, kind, bucket, object, version_id));
|
||||
if self.new_cache.info.pending_heals.len() != before {
|
||||
self.sync_pending_heals();
|
||||
}
|
||||
}
|
||||
|
||||
/// Batched variant of [`Self::clear_pending_scanner_heal`] for repaired
|
||||
/// notices (backlog#1894 axis B): one retain pass and one ledger sync
|
||||
/// for the whole notice set, so a mass-recovery first sweep cannot turn
|
||||
/// into thousands of full-table clones on the scan task. Only Object
|
||||
/// entries match — bucket-level heals are never the MRF consumer's work.
|
||||
pub(super) fn clear_pending_scanner_heals_for_repaired(&mut self, events: &[rustfs_common::mrf_channel::MrfRepairedEvent]) {
|
||||
// Pre-resolve the notice version strings once; each ledger entry then
|
||||
// compares against plain Option<&str>.
|
||||
let targets: Vec<(&str, &str, Option<String>)> = events
|
||||
.iter()
|
||||
.map(|event| (event.bucket.as_ref(), event.object.as_ref(), mrf_repaired_version_id(event.version_id)))
|
||||
.collect();
|
||||
let before = self.new_cache.info.pending_heals.len();
|
||||
self.new_cache.info.pending_heals.retain(|entry| {
|
||||
entry.kind != PendingScannerHealKind::Object
|
||||
|| !targets.iter().any(|(bucket, object, version)| {
|
||||
entry.bucket.as_str() == *bucket
|
||||
&& entry.object.as_deref() == Some(*object)
|
||||
&& entry.version_id.as_deref() == version.as_deref()
|
||||
})
|
||||
});
|
||||
if self.new_cache.info.pending_heals.len() != before {
|
||||
self.sync_pending_heals();
|
||||
}
|
||||
}
|
||||
|
||||
pub(super) fn record_pending_scanner_heal(
|
||||
&mut self,
|
||||
kind: PendingScannerHealKind,
|
||||
bucket: &str,
|
||||
object: Option<&str>,
|
||||
version_id: Option<&str>,
|
||||
scan_mode: HealScanMode,
|
||||
result: HealAdmissionResult,
|
||||
) {
|
||||
let now = Self::now_secs();
|
||||
if let Some(entry) = self
|
||||
.new_cache
|
||||
.info
|
||||
.pending_heals
|
||||
.iter_mut()
|
||||
.find(|entry| pending_scanner_heal_matches(entry, kind, bucket, object, version_id))
|
||||
{
|
||||
entry.last_attempt = now;
|
||||
entry.attempts = entry.attempts.saturating_add(1);
|
||||
entry.last_admission_result = result.result_label().to_string();
|
||||
entry.last_admission_reason = result.reason_label().to_string();
|
||||
self.sync_pending_heals();
|
||||
return;
|
||||
}
|
||||
|
||||
self.new_cache.info.pending_heals.push(PendingScannerHeal {
|
||||
kind,
|
||||
bucket: bucket.to_string(),
|
||||
object: object.map(ToOwned::to_owned),
|
||||
version_id: version_id.map(ToOwned::to_owned),
|
||||
scan_mode,
|
||||
first_seen: now,
|
||||
last_attempt: now,
|
||||
attempts: 1,
|
||||
last_admission_result: result.result_label().to_string(),
|
||||
last_admission_reason: result.reason_label().to_string(),
|
||||
});
|
||||
self.prune_pending_scanner_heals();
|
||||
self.sync_pending_heals();
|
||||
}
|
||||
|
||||
pub(super) fn prune_pending_scanner_heals(&mut self) {
|
||||
let len = self.new_cache.info.pending_heals.len();
|
||||
if len <= MAX_PENDING_SCANNER_HEALS_PER_BUCKET {
|
||||
return;
|
||||
}
|
||||
|
||||
sort_pending_scanner_heals_for_retry(&mut self.new_cache.info.pending_heals);
|
||||
let remove_count = len.saturating_sub(MAX_PENDING_SCANNER_HEALS_PER_BUCKET);
|
||||
self.new_cache.info.pending_heals.drain(..remove_count);
|
||||
counter!(
|
||||
METRIC_SCANNER_PENDING_HEAL_PRUNE_TOTAL,
|
||||
"bucket" => self.new_cache.info.name.clone()
|
||||
)
|
||||
.increment(u64::try_from(remove_count).unwrap_or(u64::MAX));
|
||||
warn!(
|
||||
target: "rustfs::scanner::folder",
|
||||
event = EVENT_SCANNER_HEAL_ADMISSION,
|
||||
component = LOG_COMPONENT_SCANNER,
|
||||
subsystem = LOG_SUBSYSTEM_HEAL,
|
||||
bucket = %self.new_cache.info.name,
|
||||
pruned = remove_count,
|
||||
remaining = self.new_cache.info.pending_heals.len(),
|
||||
state = "pending_heal_pruned",
|
||||
"Scanner pending heal ledger pruned oldest entries"
|
||||
);
|
||||
}
|
||||
|
||||
pub(super) fn update_pending_scanner_heal_after_admission(
|
||||
&mut self,
|
||||
kind: PendingScannerHealKind,
|
||||
bucket: &str,
|
||||
object: Option<&str>,
|
||||
version_id: Option<&str>,
|
||||
scan_mode: HealScanMode,
|
||||
result: HealAdmissionResult,
|
||||
) {
|
||||
match result {
|
||||
HealAdmissionResult::Accepted | HealAdmissionResult::Merged => {
|
||||
self.clear_pending_scanner_heal(kind, bucket, object, version_id);
|
||||
}
|
||||
HealAdmissionResult::Full | HealAdmissionResult::Dropped(HealAdmissionDropReason::QueueFull) => {
|
||||
self.record_pending_scanner_heal(kind, bucket, object, version_id, scan_mode, result);
|
||||
}
|
||||
HealAdmissionResult::Dropped(HealAdmissionDropReason::PolicyDropped) => {
|
||||
self.clear_pending_scanner_heal(kind, bucket, object, version_id);
|
||||
}
|
||||
// Admin-only overlap rejections (HS-06); the scanner never sees
|
||||
// them, but if it ever does, treat them as terminal like any
|
||||
// other policy drop rather than endlessly retrying.
|
||||
HealAdmissionResult::Dropped(HealAdmissionDropReason::AlreadyRunning)
|
||||
| HealAdmissionResult::Dropped(HealAdmissionDropReason::OverlappingPaths) => {
|
||||
self.clear_pending_scanner_heal(kind, bucket, object, version_id);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
pub(super) async fn retry_pending_scanner_heals(&mut self) -> Result<(), ScannerError> {
|
||||
if !self.should_heal().await {
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
let bucket = self.new_cache.info.name.clone();
|
||||
// Backlog#1894 axis B: repairs the MRF consumer landed hand the
|
||||
// manager the heal task, so the matching pending-ledger entries are
|
||||
// retried nowhere — drop them here. Best-effort: a lost notice just
|
||||
// leaves the entry to expire through its own attempts/age limits.
|
||||
let repaired = rustfs_common::mrf_channel::take_mrf_repaired_events_for(&bucket);
|
||||
if !repaired.is_empty() {
|
||||
self.clear_pending_scanner_heals_for_repaired(&repaired);
|
||||
}
|
||||
for pending in pending_scanner_heal_retry_candidates(&self.new_cache.info.pending_heals, &bucket) {
|
||||
if !self.should_heal().await {
|
||||
break;
|
||||
}
|
||||
|
||||
let Some(request) = build_pending_scanner_heal_request(&pending) else {
|
||||
self.clear_pending_scanner_heal(pending.kind, &pending.bucket, None, pending.version_id.as_deref());
|
||||
counter!(
|
||||
METRIC_SCANNER_PENDING_HEAL_MALFORMED_TOTAL,
|
||||
"bucket" => pending.bucket.clone(),
|
||||
"type" => pending_scanner_heal_candidate_type(pending.kind).to_string()
|
||||
)
|
||||
.increment(1);
|
||||
warn!(
|
||||
target: "rustfs::scanner::folder",
|
||||
event = EVENT_SCANNER_HEAL_ADMISSION,
|
||||
component = LOG_COMPONENT_SCANNER,
|
||||
subsystem = LOG_SUBSYSTEM_HEAL,
|
||||
bucket = %pending.bucket,
|
||||
state = "pending_heal_malformed",
|
||||
"Scanner dropped malformed pending heal entry"
|
||||
);
|
||||
continue;
|
||||
};
|
||||
|
||||
self.send_required_scanner_heal_request(
|
||||
pending.kind,
|
||||
pending.bucket.clone(),
|
||||
pending.object.clone(),
|
||||
pending.version_id.clone(),
|
||||
request,
|
||||
)
|
||||
.await?;
|
||||
}
|
||||
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
pub(super) fn pending_scanner_heal_candidate_type(kind: PendingScannerHealKind) -> &'static str {
|
||||
match kind {
|
||||
PendingScannerHealKind::Bucket => "bucket",
|
||||
PendingScannerHealKind::Object => "object",
|
||||
}
|
||||
}
|
||||
|
||||
pub(super) fn pending_scanner_heal_matches(
|
||||
entry: &PendingScannerHeal,
|
||||
kind: PendingScannerHealKind,
|
||||
bucket: &str,
|
||||
object: Option<&str>,
|
||||
version_id: Option<&str>,
|
||||
) -> bool {
|
||||
entry.kind == kind && entry.bucket == bucket && entry.object.as_deref() == object && entry.version_id.as_deref() == version_id
|
||||
}
|
||||
|
||||
pub(super) fn pending_scanner_heal_identity(entry: &PendingScannerHeal) -> (u8, &str, Option<&str>, Option<&str>) {
|
||||
let kind = match entry.kind {
|
||||
PendingScannerHealKind::Bucket => 0,
|
||||
PendingScannerHealKind::Object => 1,
|
||||
};
|
||||
(kind, entry.bucket.as_str(), entry.object.as_deref(), entry.version_id.as_deref())
|
||||
}
|
||||
|
||||
/// Decode an MRF repaired-notice version id for ledger matching. A nil UUID
|
||||
/// means "no value" per the repo-wide defensive-UUID invariant, so it maps
|
||||
/// to `None` and matches unversioned ledger entries only.
|
||||
pub(super) fn mrf_repaired_version_id(version_id: Option<[u8; 16]>) -> Option<String> {
|
||||
version_id
|
||||
.map(uuid::Uuid::from_bytes)
|
||||
.filter(|uuid| !uuid.is_nil())
|
||||
.map(|uuid| uuid.to_string())
|
||||
}
|
||||
|
||||
pub(super) fn sort_pending_scanner_heals_for_retry(entries: &mut [PendingScannerHeal]) {
|
||||
entries.sort_by(|a, b| {
|
||||
a.last_attempt
|
||||
.cmp(&b.last_attempt)
|
||||
.then_with(|| a.attempts.cmp(&b.attempts))
|
||||
.then_with(|| pending_scanner_heal_identity(a).cmp(&pending_scanner_heal_identity(b)))
|
||||
});
|
||||
}
|
||||
|
||||
pub(super) fn pending_scanner_heal_retry_candidates(
|
||||
pending_heals: &[PendingScannerHeal],
|
||||
bucket: &str,
|
||||
) -> Vec<PendingScannerHeal> {
|
||||
let mut entries: Vec<PendingScannerHeal> = pending_heals.iter().filter(|entry| entry.bucket == bucket).cloned().collect();
|
||||
sort_pending_scanner_heals_for_retry(&mut entries);
|
||||
entries.truncate(MAX_PENDING_SCANNER_HEAL_RETRIES_PER_BUCKET);
|
||||
entries
|
||||
}
|
||||
|
||||
pub(super) fn build_pending_scanner_heal_request(entry: &PendingScannerHeal) -> Option<HealChannelRequest> {
|
||||
match entry.kind {
|
||||
PendingScannerHealKind::Bucket => Some(build_bucket_heal_request(entry.bucket.clone(), HealChannelPriority::High)),
|
||||
PendingScannerHealKind::Object => entry.object.as_ref().map(|object| {
|
||||
build_object_heal_request(
|
||||
entry.bucket.clone(),
|
||||
object.clone(),
|
||||
entry.version_id.clone(),
|
||||
entry.scan_mode,
|
||||
HealChannelPriority::High,
|
||||
)
|
||||
}),
|
||||
}
|
||||
}
|
||||
File diff suppressed because it is too large
Load Diff
Reference in New Issue
Block a user