feat(heal): add scanner-aware bitrot controls (#3297)

* feat(heal): add scanner-aware bitrot controls

* fix(config): register heal defaults

* test(admin): avoid heal config queue conflict

---------

Co-authored-by: Henry Guo <marshawcoco@users.noreply.github.com>
Co-authored-by: houseme <housemecn@gmail.com>
This commit is contained in:
Henry Guo
2026-06-09 21:43:18 +08:00
committed by GitHub
parent 0cdcd1eb7b
commit e79c793803
6 changed files with 210 additions and 61 deletions
+12
View File
@@ -12,6 +12,18 @@
// See the License for the specific language governing permissions and
// limitations under the License.
/// Heal admin config subsystem name.
pub const HEAL_SUB_SYS: &str = "heal";
/// Heal config key setting the scanner-driven periodic deep bitrot scan cycle in seconds.
pub const HEAL_BITROT_CYCLE: &str = "bitrot_cycle";
/// Heal config keys supported by the admin config subsystem.
pub const HEAL_KEYS: &[&str] = &[HEAL_BITROT_CYCLE];
/// Default scanner-driven bitrot scan cycle used by heal/scanner runtime config.
pub const DEFAULT_HEAL_BITROT_CYCLE_SECS: u64 = 30 * 24 * 60 * 60;
/// Environment variable name that enables or disables auto-heal functionality.
/// - Purpose: Control whether the system automatically performs heal operations.
/// - Valid values: "true" or "false" (case insensitive).
+11
View File
@@ -12,10 +12,21 @@
// See the License for the specific language governing permissions and
// limitations under the License.
use crate::config::{KV, KVS};
use crate::error::{Error, Result};
use rustfs_config::{DEFAULT_HEAL_BITROT_CYCLE_SECS, HEAL_BITROT_CYCLE};
use rustfs_utils::string::parse_bool;
use std::sync::LazyLock;
use std::time::Duration;
pub static DEFAULT_KVS: LazyLock<KVS> = LazyLock::new(|| {
KVS(vec![KV {
key: HEAL_BITROT_CYCLE.to_owned(),
value: DEFAULT_HEAL_BITROT_CYCLE_SECS.to_string(),
hidden_if_empty: false,
}])
});
#[derive(Debug, Default)]
pub struct Config {
pub bitrot: String,
+12 -1
View File
@@ -26,6 +26,7 @@ use crate::store::ECStore;
use com::{STORAGE_CLASS_SUB_SYS, lookup_configs, read_config_without_migrate};
use rustfs_config::COMMENT_KEY;
use rustfs_config::DEFAULT_DELIMITER;
use rustfs_config::HEAL_SUB_SYS;
use rustfs_config::audit::{
AUDIT_AMQP_SUB_SYS, AUDIT_KAFKA_SUB_SYS, AUDIT_MQTT_SUB_SYS, AUDIT_MYSQL_SUB_SYS, AUDIT_NATS_SUB_SYS, AUDIT_POSTGRES_SUB_SYS,
AUDIT_PULSAR_SUB_SYS, AUDIT_REDIS_SUB_SYS, AUDIT_WEBHOOK_SUB_SYS,
@@ -256,6 +257,7 @@ pub fn init() {
// Load storageclass default configuration
kvs.insert(STORAGE_CLASS_SUB_SYS.to_owned(), storageclass::DEFAULT_KVS.clone());
kvs.insert(rustfs_config::SCANNER_SUB_SYS.to_owned(), scanner::DEFAULT_KVS.clone());
kvs.insert(HEAL_SUB_SYS.to_owned(), heal::DEFAULT_KVS.clone());
// New: Loading default configurations for notify_webhook and notify_mqtt
// Referring subsystem names through constants to improve the readability and maintainability of the code
kvs.insert(NOTIFY_WEBHOOK_SUB_SYS.to_owned(), notify::DEFAULT_NOTIFY_WEBHOOK_KVS.clone());
@@ -285,7 +287,10 @@ pub fn init() {
#[cfg(test)]
mod tests {
use super::*;
use rustfs_config::{DEFAULT_DELIMITER, DEFAULT_SCANNER_SPEED, SCANNER_CYCLE_MAX_OBJECTS, SCANNER_SPEED, SCANNER_SUB_SYS};
use rustfs_config::{
DEFAULT_DELIMITER, DEFAULT_HEAL_BITROT_CYCLE_SECS, DEFAULT_SCANNER_SPEED, HEAL_BITROT_CYCLE, SCANNER_CYCLE_MAX_OBJECTS,
SCANNER_SPEED, SCANNER_SUB_SYS,
};
#[test]
fn global_server_config_set_and_get_roundtrip() {
@@ -314,5 +319,11 @@ mod tests {
assert_eq!(scanner_kvs.get(SCANNER_SPEED), DEFAULT_SCANNER_SPEED);
assert_eq!(scanner_kvs.get(SCANNER_CYCLE_MAX_OBJECTS), "0");
let heal_kvs = cfg
.get_value(HEAL_SUB_SYS, DEFAULT_DELIMITER)
.expect("heal defaults should exist");
assert_eq!(heal_kvs.get(HEAL_BITROT_CYCLE), DEFAULT_HEAL_BITROT_CYCLE_SECS.to_string());
}
}
+125 -49
View File
@@ -15,8 +15,8 @@
use crate::scanner_budget::ScannerCycleBudgetConfig;
use crate::sleeper::{SCANNER_SLEEPER, scanner_default_speed};
use rustfs_config::{
DEFAULT_DELIMITER, DEFAULT_SCANNER_ALERT_EXCESS_FOLDERS, DEFAULT_SCANNER_ALERT_EXCESS_VERSION_SIZE,
DEFAULT_SCANNER_ALERT_EXCESS_VERSIONS, DEFAULT_SCANNER_BITROT_CYCLE_SECS, DEFAULT_SCANNER_CACHE_SAVE_TIMEOUT_SECS,
DEFAULT_DELIMITER, DEFAULT_HEAL_BITROT_CYCLE_SECS, DEFAULT_SCANNER_ALERT_EXCESS_FOLDERS,
DEFAULT_SCANNER_ALERT_EXCESS_VERSION_SIZE, DEFAULT_SCANNER_ALERT_EXCESS_VERSIONS, DEFAULT_SCANNER_CACHE_SAVE_TIMEOUT_SECS,
DEFAULT_SCANNER_CYCLE_MAX_DIRECTORIES, DEFAULT_SCANNER_CYCLE_MAX_DURATION_SECS, DEFAULT_SCANNER_CYCLE_MAX_OBJECTS,
DEFAULT_SCANNER_IDLE_MODE, DEFAULT_SCANNER_MAX_CONCURRENT_DISK_SCANS, DEFAULT_SCANNER_MAX_CONCURRENT_SET_SCANS,
DEFAULT_SCANNER_SPEED, DEFAULT_SCANNER_YIELD_EVERY_N_OBJECTS, ENV_SCANNER_ALERT_EXCESS_FOLDERS,
@@ -24,9 +24,9 @@ use rustfs_config::{
ENV_SCANNER_CACHE_SAVE_TIMEOUT_SECS, ENV_SCANNER_CYCLE, ENV_SCANNER_CYCLE_MAX_DIRECTORIES,
ENV_SCANNER_CYCLE_MAX_DURATION_SECS, ENV_SCANNER_CYCLE_MAX_OBJECTS, ENV_SCANNER_IDLE_MODE,
ENV_SCANNER_MAX_CONCURRENT_DISK_SCANS, ENV_SCANNER_MAX_CONCURRENT_SET_SCANS, ENV_SCANNER_SPEED, ENV_SCANNER_START_DELAY_SECS,
ENV_SCANNER_YIELD_EVERY_N_OBJECTS, SCANNER_ALERT_EXCESS_FOLDERS, SCANNER_ALERT_EXCESS_VERSION_SIZE,
SCANNER_ALERT_EXCESS_VERSIONS, SCANNER_BITROT_CYCLE, SCANNER_CACHE_SAVE_TIMEOUT, SCANNER_CYCLE,
SCANNER_CYCLE_MAX_DIRECTORIES, SCANNER_CYCLE_MAX_DURATION, SCANNER_CYCLE_MAX_OBJECTS, SCANNER_IDLE_MODE,
ENV_SCANNER_YIELD_EVERY_N_OBJECTS, HEAL_BITROT_CYCLE, HEAL_SUB_SYS, SCANNER_ALERT_EXCESS_FOLDERS,
SCANNER_ALERT_EXCESS_VERSION_SIZE, SCANNER_ALERT_EXCESS_VERSIONS, SCANNER_BITROT_CYCLE, SCANNER_CACHE_SAVE_TIMEOUT,
SCANNER_CYCLE, SCANNER_CYCLE_MAX_DIRECTORIES, SCANNER_CYCLE_MAX_DURATION, SCANNER_CYCLE_MAX_OBJECTS, SCANNER_IDLE_MODE,
SCANNER_MAX_CONCURRENT_DISK_SCANS, SCANNER_MAX_CONCURRENT_SET_SCANS, SCANNER_SPEED, SCANNER_START_DELAY, SCANNER_SUB_SYS,
SCANNER_YIELD_EVERY_N_OBJECTS, ScannerSpeed,
};
@@ -51,6 +51,7 @@ static SCANNER_RUNTIME_CONFIG: LazyLock<RwLock<ScannerRuntimeConfig>> =
pub enum ScannerRuntimeConfigSource {
Env,
Config,
ScannerCompatConfig,
Default,
}
@@ -99,7 +100,7 @@ impl Default for ScannerRuntimeConfig {
.map(Duration::from_secs)
.unwrap_or_else(|| scanner_default_speed().cycle_interval()),
cycle_interval_source: ScannerRuntimeConfigSource::Default,
bitrot_cycle: Some(Duration::from_secs(DEFAULT_SCANNER_BITROT_CYCLE_SECS)),
bitrot_cycle: Some(Duration::from_secs(DEFAULT_HEAL_BITROT_CYCLE_SECS)),
bitrot_cycle_source: ScannerRuntimeConfigSource::Default,
cycle_budget: ScannerCycleBudgetConfig::default(),
cycle_max_duration_source: ScannerRuntimeConfigSource::Default,
@@ -183,6 +184,10 @@ fn scanner_kvs(config: Option<&ServerConfig>) -> Option<KVS> {
config.and_then(|config| config.get_value(SCANNER_SUB_SYS, DEFAULT_DELIMITER))
}
fn heal_kvs(config: Option<&ServerConfig>) -> Option<KVS> {
config.and_then(|config| config.get_value(HEAL_SUB_SYS, DEFAULT_DELIMITER))
}
fn validate_default_scanner_target(config: &ServerConfig) -> Result<(), ScannerRuntimeConfigError> {
let Some(targets) = config.0.get(SCANNER_SUB_SYS) else {
return Ok(());
@@ -196,6 +201,19 @@ fn validate_default_scanner_target(config: &ServerConfig) -> Result<(), ScannerR
Ok(())
}
fn validate_default_heal_target(config: &ServerConfig) -> Result<(), ScannerRuntimeConfigError> {
let Some(targets) = config.0.get(HEAL_SUB_SYS) else {
return Ok(());
};
for target in targets.keys() {
if target != DEFAULT_DELIMITER {
return Err(invalid_value("target", target.clone(), "heal config only supports the default target"));
}
}
Ok(())
}
fn invalid_value(key: &'static str, value: impl Into<String>, reason: &'static str) -> ScannerRuntimeConfigError {
ScannerRuntimeConfigError::InvalidValue {
key,
@@ -228,11 +246,11 @@ fn parse_config_speed(value: String) -> Result<ScannerSpeed, ScannerRuntimeConfi
ScannerSpeed::parse_str(&value).ok_or_else(|| invalid_value(SCANNER_SPEED, value, "expected scanner speed preset"))
}
fn parse_config_bitrot_cycle(value: String) -> Result<Option<Duration>, ScannerRuntimeConfigError> {
fn parse_config_bitrot_cycle(key: &'static str, value: String) -> Result<Option<Duration>, ScannerRuntimeConfigError> {
match value.trim().to_ascii_lowercase().as_str() {
"0" | "true" | "on" | "yes" => Ok(Some(Duration::ZERO)),
"false" | "off" | "no" | "disabled" => Ok(None),
_ => parse_config_u64(SCANNER_BITROT_CYCLE, value).map(|secs| Some(Duration::from_secs(secs))),
_ => parse_config_u64(key, value).map(|secs| Some(Duration::from_secs(secs))),
}
}
@@ -244,10 +262,10 @@ fn parse_env_bitrot_cycle(value: String) -> Option<Duration> {
warn!(
env = ENV_SCANNER_BITROT_CYCLE_SECS,
value,
default_secs = DEFAULT_SCANNER_BITROT_CYCLE_SECS,
default_secs = DEFAULT_HEAL_BITROT_CYCLE_SECS,
"Invalid scanner bitrot cycle, using default"
);
Some(Duration::from_secs(DEFAULT_SCANNER_BITROT_CYCLE_SECS))
Some(Duration::from_secs(DEFAULT_HEAL_BITROT_CYCLE_SECS))
}),
}
}
@@ -276,31 +294,37 @@ fn validate_optional_config_usize(
fn validate_persisted_scanner_runtime_config(config: &ServerConfig) -> Result<(), ScannerRuntimeConfigError> {
validate_default_scanner_target(config)?;
validate_default_heal_target(config)?;
let kvs = scanner_kvs(Some(config));
let kvs = kvs.as_ref();
let scanner_kvs = scanner_kvs(Some(config));
let scanner_kvs = scanner_kvs.as_ref();
let heal_kvs = heal_kvs(Some(config));
let heal_kvs = heal_kvs.as_ref();
if let Some(value) = config_value(kvs, SCANNER_SPEED, DEFAULT_SCANNER_SPEED) {
if let Some(value) = config_value(scanner_kvs, SCANNER_SPEED, DEFAULT_SCANNER_SPEED) {
parse_config_speed(value)?;
}
if let Some(value) = config_value(kvs, SCANNER_IDLE_MODE, DEFAULT_SCANNER_IDLE_MODE) {
if let Some(value) = config_value(scanner_kvs, SCANNER_IDLE_MODE, DEFAULT_SCANNER_IDLE_MODE) {
parse_config_bool(SCANNER_IDLE_MODE, value)?;
}
validate_optional_config_u64(kvs, SCANNER_START_DELAY, "")?;
validate_optional_config_u64(kvs, SCANNER_CYCLE, "")?;
validate_optional_config_u64(kvs, SCANNER_CYCLE_MAX_DURATION, DEFAULT_SCANNER_CYCLE_MAX_DURATION_SECS)?;
validate_optional_config_u64(kvs, SCANNER_CYCLE_MAX_OBJECTS, DEFAULT_SCANNER_CYCLE_MAX_OBJECTS)?;
validate_optional_config_u64(kvs, SCANNER_CYCLE_MAX_DIRECTORIES, DEFAULT_SCANNER_CYCLE_MAX_DIRECTORIES)?;
if let Some(value) = config_value(kvs, SCANNER_BITROT_CYCLE, DEFAULT_SCANNER_BITROT_CYCLE_SECS) {
parse_config_bitrot_cycle(value)?;
validate_optional_config_u64(scanner_kvs, SCANNER_START_DELAY, "")?;
validate_optional_config_u64(scanner_kvs, SCANNER_CYCLE, "")?;
validate_optional_config_u64(scanner_kvs, SCANNER_CYCLE_MAX_DURATION, DEFAULT_SCANNER_CYCLE_MAX_DURATION_SECS)?;
validate_optional_config_u64(scanner_kvs, SCANNER_CYCLE_MAX_OBJECTS, DEFAULT_SCANNER_CYCLE_MAX_OBJECTS)?;
validate_optional_config_u64(scanner_kvs, SCANNER_CYCLE_MAX_DIRECTORIES, DEFAULT_SCANNER_CYCLE_MAX_DIRECTORIES)?;
if let Some(value) = config_value(heal_kvs, HEAL_BITROT_CYCLE, DEFAULT_HEAL_BITROT_CYCLE_SECS) {
parse_config_bitrot_cycle(HEAL_BITROT_CYCLE, value)?;
}
validate_optional_config_u64(kvs, SCANNER_CACHE_SAVE_TIMEOUT, DEFAULT_SCANNER_CACHE_SAVE_TIMEOUT_SECS)?;
validate_optional_config_usize(kvs, SCANNER_MAX_CONCURRENT_SET_SCANS, DEFAULT_SCANNER_MAX_CONCURRENT_SET_SCANS)?;
validate_optional_config_usize(kvs, SCANNER_MAX_CONCURRENT_DISK_SCANS, DEFAULT_SCANNER_MAX_CONCURRENT_DISK_SCANS)?;
validate_optional_config_u64(kvs, SCANNER_YIELD_EVERY_N_OBJECTS, DEFAULT_SCANNER_YIELD_EVERY_N_OBJECTS)?;
validate_optional_config_u64(kvs, SCANNER_ALERT_EXCESS_VERSIONS, DEFAULT_SCANNER_ALERT_EXCESS_VERSIONS)?;
validate_optional_config_u64(kvs, SCANNER_ALERT_EXCESS_VERSION_SIZE, DEFAULT_SCANNER_ALERT_EXCESS_VERSION_SIZE)?;
validate_optional_config_u64(kvs, SCANNER_ALERT_EXCESS_FOLDERS, DEFAULT_SCANNER_ALERT_EXCESS_FOLDERS)?;
if let Some(value) = config_value(scanner_kvs, SCANNER_BITROT_CYCLE, DEFAULT_HEAL_BITROT_CYCLE_SECS) {
parse_config_bitrot_cycle(SCANNER_BITROT_CYCLE, value)?;
}
validate_optional_config_u64(scanner_kvs, SCANNER_CACHE_SAVE_TIMEOUT, DEFAULT_SCANNER_CACHE_SAVE_TIMEOUT_SECS)?;
validate_optional_config_usize(scanner_kvs, SCANNER_MAX_CONCURRENT_SET_SCANS, DEFAULT_SCANNER_MAX_CONCURRENT_SET_SCANS)?;
validate_optional_config_usize(scanner_kvs, SCANNER_MAX_CONCURRENT_DISK_SCANS, DEFAULT_SCANNER_MAX_CONCURRENT_DISK_SCANS)?;
validate_optional_config_u64(scanner_kvs, SCANNER_YIELD_EVERY_N_OBJECTS, DEFAULT_SCANNER_YIELD_EVERY_N_OBJECTS)?;
validate_optional_config_u64(scanner_kvs, SCANNER_ALERT_EXCESS_VERSIONS, DEFAULT_SCANNER_ALERT_EXCESS_VERSIONS)?;
validate_optional_config_u64(scanner_kvs, SCANNER_ALERT_EXCESS_VERSION_SIZE, DEFAULT_SCANNER_ALERT_EXCESS_VERSION_SIZE)?;
validate_optional_config_u64(scanner_kvs, SCANNER_ALERT_EXCESS_FOLDERS, DEFAULT_SCANNER_ALERT_EXCESS_FOLDERS)?;
Ok(())
}
@@ -408,15 +432,18 @@ fn lookup_bool(
pub(crate) fn lookup_scanner_runtime_config(
config: Option<&ServerConfig>,
) -> Result<ScannerRuntimeConfig, ScannerRuntimeConfigError> {
let kvs = scanner_kvs(config);
let kvs = kvs.as_ref();
let (speed, speed_source) = lookup_speed(kvs)?;
let (idle_mode, idle_mode_source) = lookup_bool(kvs, SCANNER_IDLE_MODE, ENV_SCANNER_IDLE_MODE, DEFAULT_SCANNER_IDLE_MODE)?;
let (start_delay, start_delay_source) = lookup_start_delay(kvs)?;
let scanner_kvs = scanner_kvs(config);
let scanner_kvs = scanner_kvs.as_ref();
let heal_kvs = heal_kvs(config);
let heal_kvs = heal_kvs.as_ref();
let (speed, speed_source) = lookup_speed(scanner_kvs)?;
let (idle_mode, idle_mode_source) =
lookup_bool(scanner_kvs, SCANNER_IDLE_MODE, ENV_SCANNER_IDLE_MODE, DEFAULT_SCANNER_IDLE_MODE)?;
let (start_delay, start_delay_source) = lookup_start_delay(scanner_kvs)?;
let (cycle_interval, cycle_interval_source) = if let Some(secs) = rustfs_utils::get_env_opt_u64(ENV_SCANNER_CYCLE) {
(Duration::from_secs(secs), ScannerRuntimeConfigSource::Env)
} else if let Some(value) = config_value(kvs, SCANNER_CYCLE, "") {
} else if let Some(value) = config_value(scanner_kvs, SCANNER_CYCLE, "") {
(
Duration::from_secs(parse_config_u64(SCANNER_CYCLE, value)?),
ScannerRuntimeConfigSource::Config,
@@ -430,19 +457,19 @@ pub(crate) fn lookup_scanner_runtime_config(
};
let (cycle_max_duration, cycle_max_duration_source) = lookup_optional_seconds(
kvs,
scanner_kvs,
SCANNER_CYCLE_MAX_DURATION,
ENV_SCANNER_CYCLE_MAX_DURATION_SECS,
DEFAULT_SCANNER_CYCLE_MAX_DURATION_SECS,
)?;
let (cycle_max_objects, cycle_max_objects_source) = lookup_count_budget(
kvs,
scanner_kvs,
SCANNER_CYCLE_MAX_OBJECTS,
ENV_SCANNER_CYCLE_MAX_OBJECTS,
DEFAULT_SCANNER_CYCLE_MAX_OBJECTS,
)?;
let (cycle_max_directories, cycle_max_directories_source) = lookup_count_budget(
kvs,
scanner_kvs,
SCANNER_CYCLE_MAX_DIRECTORIES,
ENV_SCANNER_CYCLE_MAX_DIRECTORIES,
DEFAULT_SCANNER_CYCLE_MAX_DIRECTORIES,
@@ -455,53 +482,58 @@ pub(crate) fn lookup_scanner_runtime_config(
let (bitrot_cycle, bitrot_cycle_source) = if let Some(value) = rustfs_utils::get_env_opt_str(ENV_SCANNER_BITROT_CYCLE_SECS) {
(parse_env_bitrot_cycle(value), ScannerRuntimeConfigSource::Env)
} else if let Some(value) = config_value(kvs, SCANNER_BITROT_CYCLE, DEFAULT_SCANNER_BITROT_CYCLE_SECS) {
(parse_config_bitrot_cycle(value)?, ScannerRuntimeConfigSource::Config)
} else if let Some(value) = config_value(heal_kvs, HEAL_BITROT_CYCLE, DEFAULT_HEAL_BITROT_CYCLE_SECS) {
(parse_config_bitrot_cycle(HEAL_BITROT_CYCLE, value)?, ScannerRuntimeConfigSource::Config)
} else if let Some(value) = config_value(scanner_kvs, SCANNER_BITROT_CYCLE, DEFAULT_HEAL_BITROT_CYCLE_SECS) {
(
parse_config_bitrot_cycle(SCANNER_BITROT_CYCLE, value)?,
ScannerRuntimeConfigSource::ScannerCompatConfig,
)
} else {
(
Some(Duration::from_secs(DEFAULT_SCANNER_BITROT_CYCLE_SECS)),
Some(Duration::from_secs(DEFAULT_HEAL_BITROT_CYCLE_SECS)),
ScannerRuntimeConfigSource::Default,
)
};
let (cache_save_timeout, cache_save_timeout_source) = lookup_u64(
kvs,
scanner_kvs,
SCANNER_CACHE_SAVE_TIMEOUT,
ENV_SCANNER_CACHE_SAVE_TIMEOUT_SECS,
DEFAULT_SCANNER_CACHE_SAVE_TIMEOUT_SECS,
)?;
let (max_concurrent_set_scans, max_concurrent_set_scans_source) = lookup_usize(
kvs,
scanner_kvs,
SCANNER_MAX_CONCURRENT_SET_SCANS,
ENV_SCANNER_MAX_CONCURRENT_SET_SCANS,
DEFAULT_SCANNER_MAX_CONCURRENT_SET_SCANS,
)?;
let (max_concurrent_disk_scans, max_concurrent_disk_scans_source) = lookup_usize(
kvs,
scanner_kvs,
SCANNER_MAX_CONCURRENT_DISK_SCANS,
ENV_SCANNER_MAX_CONCURRENT_DISK_SCANS,
DEFAULT_SCANNER_MAX_CONCURRENT_DISK_SCANS,
)?;
let (yield_every_n_objects, yield_every_n_objects_source) = lookup_u64(
kvs,
scanner_kvs,
SCANNER_YIELD_EVERY_N_OBJECTS,
ENV_SCANNER_YIELD_EVERY_N_OBJECTS,
DEFAULT_SCANNER_YIELD_EVERY_N_OBJECTS,
)?;
let (alert_excess_versions, alert_excess_versions_source) = lookup_u64(
kvs,
scanner_kvs,
SCANNER_ALERT_EXCESS_VERSIONS,
ENV_SCANNER_ALERT_EXCESS_VERSIONS,
DEFAULT_SCANNER_ALERT_EXCESS_VERSIONS,
)?;
let (alert_excess_version_size, alert_excess_version_size_source) = lookup_u64(
kvs,
scanner_kvs,
SCANNER_ALERT_EXCESS_VERSION_SIZE,
ENV_SCANNER_ALERT_EXCESS_VERSION_SIZE,
DEFAULT_SCANNER_ALERT_EXCESS_VERSION_SIZE,
)?;
let (alert_excess_folders, alert_excess_folders_source) = lookup_u64(
kvs,
scanner_kvs,
SCANNER_ALERT_EXCESS_FOLDERS,
ENV_SCANNER_ALERT_EXCESS_FOLDERS,
DEFAULT_SCANNER_ALERT_EXCESS_FOLDERS,
@@ -685,8 +717,9 @@ pub(crate) fn scanner_alert_excess_folders() -> u64 {
mod tests {
use super::{ScannerRuntimeConfigSource, lookup_scanner_runtime_config, validate_scanner_runtime_config};
use rustfs_config::{
DEFAULT_DELIMITER, ENV_SCANNER_CACHE_SAVE_TIMEOUT_SECS, ENV_SCANNER_CYCLE, ENV_SCANNER_CYCLE_MAX_OBJECTS,
ENV_SCANNER_SPEED, SCANNER_CACHE_SAVE_TIMEOUT, SCANNER_CYCLE, SCANNER_CYCLE_MAX_DIRECTORIES, SCANNER_CYCLE_MAX_DURATION,
DEFAULT_DELIMITER, ENV_SCANNER_BITROT_CYCLE_SECS, ENV_SCANNER_CACHE_SAVE_TIMEOUT_SECS, ENV_SCANNER_CYCLE,
ENV_SCANNER_CYCLE_MAX_OBJECTS, ENV_SCANNER_SPEED, HEAL_BITROT_CYCLE, HEAL_SUB_SYS, SCANNER_BITROT_CYCLE,
SCANNER_CACHE_SAVE_TIMEOUT, SCANNER_CYCLE, SCANNER_CYCLE_MAX_DIRECTORIES, SCANNER_CYCLE_MAX_DURATION,
SCANNER_CYCLE_MAX_OBJECTS, SCANNER_IDLE_MODE, SCANNER_SPEED, SCANNER_SUB_SYS, ScannerSpeed,
};
use rustfs_ecstore::config::{Config as ServerConfig, KVS};
@@ -709,6 +742,19 @@ mod tests {
config
}
fn server_config_with_scanner_and_heal(scanner_entries: &[(&str, &str)], heal_entries: &[(&str, &str)]) -> ServerConfig {
let mut config = server_config_with_scanner(scanner_entries);
let mut kvs = KVS::new();
for (key, value) in heal_entries {
kvs.insert((*key).to_string(), (*value).to_string());
}
config
.0
.insert(HEAL_SUB_SYS.to_string(), HashMap::from([(DEFAULT_DELIMITER.to_string(), kvs)]));
config.set_defaults();
config
}
fn server_config_with_scanner_target(target: &str, entries: &[(&str, &str)]) -> ServerConfig {
let mut config = server_config_with_scanner(&[]);
let mut kvs = KVS::new();
@@ -769,6 +815,36 @@ mod tests {
});
}
#[test]
#[serial]
fn scanner_runtime_config_prefers_heal_bitrot_cycle_over_scanner_compat_config() {
let config = server_config_with_scanner_and_heal(&[(SCANNER_BITROT_CYCLE, "3600")], &[(HEAL_BITROT_CYCLE, "off")]);
with_var_unset(ENV_SCANNER_BITROT_CYCLE_SECS, || {
let resolved = lookup_scanner_runtime_config(Some(&config)).expect("scanner runtime config");
assert_eq!(resolved.bitrot_cycle, None);
assert_eq!(resolved.bitrot_cycle_source, ScannerRuntimeConfigSource::Config);
});
}
#[test]
#[serial]
fn scanner_runtime_config_marks_scanner_bitrot_cycle_as_compat_source() {
let config = server_config_with_scanner(&[(SCANNER_BITROT_CYCLE, "3600")]);
with_var_unset(ENV_SCANNER_BITROT_CYCLE_SECS, || {
let resolved = lookup_scanner_runtime_config(Some(&config)).expect("scanner runtime config");
super::apply_resolved_runtime_config(resolved);
let encoded = serde_json::to_value(super::scanner_runtime_config_status()).expect("scanner status should serialize");
assert_eq!(encoded["bitrot_cycle_seconds"]["value"], 3600);
assert_eq!(encoded["bitrot_cycle_seconds"]["source"], "scanner_compat_config");
});
super::refresh_scanner_runtime_config_for_tests();
}
#[test]
fn scanner_runtime_config_rejects_invalid_persisted_speed() {
let config = server_config_with_scanner(&[(SCANNER_SPEED, "warp")]);
+42 -7
View File
@@ -54,13 +54,14 @@ use rustfs_config::{
ENV_SCANNER_CACHE_SAVE_TIMEOUT_SECS, ENV_SCANNER_CYCLE, ENV_SCANNER_CYCLE_MAX_DIRECTORIES,
ENV_SCANNER_CYCLE_MAX_DURATION_SECS, ENV_SCANNER_CYCLE_MAX_OBJECTS, ENV_SCANNER_IDLE_MODE,
ENV_SCANNER_MAX_CONCURRENT_DISK_SCANS, ENV_SCANNER_MAX_CONCURRENT_SET_SCANS, ENV_SCANNER_SPEED, ENV_SCANNER_START_DELAY_SECS,
ENV_SCANNER_YIELD_EVERY_N_OBJECTS, MAX_ADMIN_REQUEST_BODY_SIZE, MQTT_BROKER, MQTT_KEEP_ALIVE_INTERVAL, MQTT_PASSWORD,
MQTT_QOS, MQTT_QUEUE_DIR, MQTT_QUEUE_LIMIT, MQTT_RECONNECT_INTERVAL, MQTT_TOPIC, MQTT_USERNAME, SCANNER_ALERT_EXCESS_FOLDERS,
SCANNER_ALERT_EXCESS_VERSION_SIZE, SCANNER_ALERT_EXCESS_VERSIONS, SCANNER_BITROT_CYCLE, SCANNER_CACHE_SAVE_TIMEOUT,
SCANNER_CYCLE, SCANNER_CYCLE_MAX_DIRECTORIES, SCANNER_CYCLE_MAX_DURATION, SCANNER_CYCLE_MAX_OBJECTS, SCANNER_IDLE_MODE,
SCANNER_MAX_CONCURRENT_DISK_SCANS, SCANNER_MAX_CONCURRENT_SET_SCANS, SCANNER_SPEED, SCANNER_START_DELAY, SCANNER_SUB_SYS,
SCANNER_YIELD_EVERY_N_OBJECTS, WEBHOOK_AUTH_TOKEN, WEBHOOK_BATCH_SIZE, WEBHOOK_CLIENT_CERT, WEBHOOK_CLIENT_KEY,
WEBHOOK_ENDPOINT, WEBHOOK_HTTP_TIMEOUT, WEBHOOK_MAX_RETRY, WEBHOOK_QUEUE_DIR, WEBHOOK_QUEUE_LIMIT, WEBHOOK_RETRY_INTERVAL,
ENV_SCANNER_YIELD_EVERY_N_OBJECTS, HEAL_BITROT_CYCLE, HEAL_SUB_SYS, MAX_ADMIN_REQUEST_BODY_SIZE, MQTT_BROKER,
MQTT_KEEP_ALIVE_INTERVAL, MQTT_PASSWORD, MQTT_QOS, MQTT_QUEUE_DIR, MQTT_QUEUE_LIMIT, MQTT_RECONNECT_INTERVAL, MQTT_TOPIC,
MQTT_USERNAME, SCANNER_ALERT_EXCESS_FOLDERS, SCANNER_ALERT_EXCESS_VERSION_SIZE, SCANNER_ALERT_EXCESS_VERSIONS,
SCANNER_BITROT_CYCLE, SCANNER_CACHE_SAVE_TIMEOUT, SCANNER_CYCLE, SCANNER_CYCLE_MAX_DIRECTORIES, SCANNER_CYCLE_MAX_DURATION,
SCANNER_CYCLE_MAX_OBJECTS, SCANNER_IDLE_MODE, SCANNER_MAX_CONCURRENT_DISK_SCANS, SCANNER_MAX_CONCURRENT_SET_SCANS,
SCANNER_SPEED, SCANNER_START_DELAY, SCANNER_SUB_SYS, SCANNER_YIELD_EVERY_N_OBJECTS, WEBHOOK_AUTH_TOKEN, WEBHOOK_BATCH_SIZE,
WEBHOOK_CLIENT_CERT, WEBHOOK_CLIENT_KEY, WEBHOOK_ENDPOINT, WEBHOOK_HTTP_TIMEOUT, WEBHOOK_MAX_RETRY, WEBHOOK_QUEUE_DIR,
WEBHOOK_QUEUE_LIMIT, WEBHOOK_RETRY_INTERVAL,
};
use rustfs_credentials::Credentials;
use rustfs_ecstore::config::com::STORAGE_CLASS_SUB_SYS;
@@ -277,6 +278,13 @@ const SCANNER_HELP_KEYS: &[HelpKeyMetadata] = &[
},
];
const HEAL_HELP_KEYS: &[HelpKeyMetadata] = &[HelpKeyMetadata {
key: HEAL_BITROT_CYCLE,
type_name: "seconds|off",
description: "set scanner-driven periodic deep bitrot scan cycle, 0 scans deeply every cycle, off disables periodic deep scans",
optional: true,
}];
const OIDC_HELP_KEYS: &[HelpKeyMetadata] = &[
HelpKeyMetadata {
key: OIDC_CONFIG_URL,
@@ -530,6 +538,12 @@ const HELP_SUBSYSTEMS: &[HelpSubSystemMetadata] = &[
multiple_targets: false,
keys: SCANNER_HELP_KEYS,
},
HelpSubSystemMetadata {
key: HEAL_SUB_SYS,
description: "configure heal and scanner-driven bitrot behavior",
multiple_targets: false,
keys: HEAL_HELP_KEYS,
},
HelpSubSystemMetadata {
key: IDENTITY_OPENID_SUB_SYS,
description: "enable OpenID SSO support",
@@ -1322,6 +1336,7 @@ fn env_help_key(sub_system: &str, key: &str) -> String {
(SCANNER_SUB_SYS, SCANNER_ALERT_EXCESS_VERSIONS) => ENV_SCANNER_ALERT_EXCESS_VERSIONS.to_string(),
(SCANNER_SUB_SYS, SCANNER_ALERT_EXCESS_VERSION_SIZE) => ENV_SCANNER_ALERT_EXCESS_VERSION_SIZE.to_string(),
(SCANNER_SUB_SYS, SCANNER_ALERT_EXCESS_FOLDERS) => ENV_SCANNER_ALERT_EXCESS_FOLDERS.to_string(),
(HEAL_SUB_SYS, HEAL_BITROT_CYCLE) => ENV_SCANNER_BITROT_CYCLE_SECS.to_string(),
(IDENTITY_OPENID_SUB_SYS, ENABLE_KEY) => ENV_IDENTITY_OPENID_ENABLE.to_string(),
(IDENTITY_OPENID_SUB_SYS, OIDC_CONFIG_URL) => ENV_IDENTITY_OPENID_CONFIG_URL.to_string(),
(IDENTITY_OPENID_SUB_SYS, OIDC_CLIENT_ID) => ENV_IDENTITY_OPENID_CLIENT_ID.to_string(),
@@ -1515,6 +1530,7 @@ async fn apply_and_signal_dynamic_subsystems(config: &ServerConfig) {
AUDIT_WEBHOOK_SUB_SYS,
AUDIT_MQTT_SUB_SYS,
SCANNER_SUB_SYS,
HEAL_SUB_SYS,
] {
if apply_dynamic_config_for_subsystem(config, sub_system).await.unwrap_or(false) {
signal_dynamic_config_reload(sub_system).await;
@@ -1974,6 +1990,25 @@ identity_openid config_url="https://issuer.example" client_id="console""#,
assert!(response.keys_help[0].description.contains("cycle"));
}
#[test]
fn validate_config_directives_accepts_heal_bitrot_cycle() {
let input = format!("{HEAL_SUB_SYS} {HEAL_BITROT_CYCLE}=\"off\"");
let directives = parse_config_directives(&input, false).expect("parse heal directive");
validate_config_directives(&directives).expect("heal bitrot_cycle should be a supported config key");
}
#[test]
fn build_help_response_reports_heal_bitrot_cycle() {
let response = build_help_response(Some(HEAL_SUB_SYS), Some(HEAL_BITROT_CYCLE), false).expect("heal help response");
assert_eq!(response.sub_sys, HEAL_SUB_SYS);
assert!(!response.multiple_targets);
assert_eq!(response.keys_help.len(), 1);
assert_eq!(response.keys_help[0].key, HEAL_BITROT_CYCLE);
assert!(response.keys_help[0].description.contains("scanner"));
}
#[test]
fn build_top_level_help_response_uses_empty_type_names() {
let response = build_help_response(None, None, false).expect("top level help response");
+8 -4
View File
@@ -13,12 +13,12 @@
// limitations under the License.
use rustfs_audit::reload_audit_config;
use rustfs_config::SCANNER_SUB_SYS;
use rustfs_config::audit::{AUDIT_MQTT_SUB_SYS, AUDIT_REDIS_DEFAULT_CHANNEL, AUDIT_WEBHOOK_SUB_SYS};
use rustfs_config::notify::{NOTIFY_MQTT_SUB_SYS, NOTIFY_REDIS_DEFAULT_CHANNEL, NOTIFY_WEBHOOK_SUB_SYS};
use rustfs_config::oidc::IDENTITY_OPENID_SUB_SYS;
use rustfs_config::{AUDIT_DEFAULT_DIR, EVENT_DEFAULT_DIR};
use rustfs_config::{DEFAULT_DELIMITER, ENABLE_KEY, EnableState};
use rustfs_config::{HEAL_SUB_SYS, SCANNER_SUB_SYS};
use rustfs_ecstore::StorageAPI;
use rustfs_ecstore::config::com::{STORAGE_CLASS_SUB_SYS, read_config_without_migrate};
use rustfs_ecstore::config::storageclass;
@@ -37,7 +37,7 @@ use url::Url;
pub fn is_dynamic_config_subsystem(sub_system: &str) -> bool {
matches!(
sub_system,
STORAGE_CLASS_SUB_SYS | AUDIT_WEBHOOK_SUB_SYS | AUDIT_MQTT_SUB_SYS | SCANNER_SUB_SYS
STORAGE_CLASS_SUB_SYS | AUDIT_WEBHOOK_SUB_SYS | AUDIT_MQTT_SUB_SYS | SCANNER_SUB_SYS | HEAL_SUB_SYS
)
}
@@ -235,7 +235,7 @@ pub async fn validate_server_config(config: &ServerConfig, sub_system: Option<&s
Some(AUDIT_WEBHOOK_SUB_SYS) => validate_audit_subsystem_config(config, AUDIT_WEBHOOK_SUB_SYS),
Some(AUDIT_MQTT_SUB_SYS) => validate_audit_subsystem_config(config, AUDIT_MQTT_SUB_SYS),
Some(IDENTITY_OPENID_SUB_SYS) => validate_identity_openid_config(config),
Some(SCANNER_SUB_SYS) => rustfs_scanner::validate_scanner_runtime_config(config)
Some(SCANNER_SUB_SYS | HEAL_SUB_SYS) => rustfs_scanner::validate_scanner_runtime_config(config)
.map_err(|err| invalid_request(format!("invalid scanner config: {err}"))),
Some(_) => Ok(()),
None => {
@@ -264,6 +264,8 @@ pub async fn apply_dynamic_config_for_subsystem(config: &ServerConfig, sub_syste
.map_err(|err| internal_error(format!("failed to reload audit config: {err}")))?,
SCANNER_SUB_SYS => rustfs_scanner::apply_scanner_runtime_config(config)
.map_err(|err| internal_error(format!("failed to reload scanner config: {err}")))?,
HEAL_SUB_SYS => rustfs_scanner::apply_scanner_runtime_config(config)
.map_err(|err| internal_error(format!("failed to reload heal scanner controls: {err}")))?,
_ => return Ok(false),
}
@@ -307,6 +309,7 @@ pub async fn reload_runtime_config_snapshot() -> S3Result<()> {
AUDIT_WEBHOOK_SUB_SYS,
AUDIT_MQTT_SUB_SYS,
SCANNER_SUB_SYS,
HEAL_SUB_SYS,
] {
if let Err(err) = apply_dynamic_config_for_subsystem(&config, sub_system).await {
warn!("peer reload_runtime_config_snapshot: failed to apply {sub_system}: {err}");
@@ -348,9 +351,9 @@ pub async fn signal_config_snapshot_reload() {
#[cfg(test)]
mod tests {
use super::*;
use rustfs_config::SCANNER_SUB_SYS;
use rustfs_config::notify::NOTIFY_WEBHOOK_SUB_SYS;
use rustfs_config::oidc::{OIDC_CLIENT_ID, OIDC_CONFIG_URL, OIDC_SCOPES};
use rustfs_config::{HEAL_SUB_SYS, SCANNER_SUB_SYS};
use rustfs_config::{MQTT_BROKER, MQTT_QUEUE_DIR, MQTT_TOPIC, WEBHOOK_ENDPOINT, WEBHOOK_QUEUE_DIR};
#[test]
@@ -358,6 +361,7 @@ mod tests {
assert!(is_dynamic_config_subsystem(AUDIT_WEBHOOK_SUB_SYS));
assert!(is_dynamic_config_subsystem(AUDIT_MQTT_SUB_SYS));
assert!(is_dynamic_config_subsystem(SCANNER_SUB_SYS));
assert!(is_dynamic_config_subsystem(HEAL_SUB_SYS));
assert!(is_dynamic_config_subsystem(STORAGE_CLASS_SUB_SYS));
assert!(!is_dynamic_config_subsystem("identity_openid"));
assert!(!is_dynamic_config_subsystem("notify_webhook"));