Compare commits

..

1 Commits

Author SHA1 Message Date
唐小鸭 21c2fb42bb refactor(replication): split four oversized hot-path functions into focused helpers
Pure-move decomposition of the four oversized functions flagged by the
replication compatibility review (P1-18), unblocking migration milestone
M2 which requires resyncer moves to stay mechanical:

- resync_bucket (522 lines -> 61-line step sequence): leader lock,
  target resolution, walk/collector/worker spawning, and dispatch loop
  extracted into focused helpers; pure decision helpers (DTO builders,
  HEAD-result classification) separated from IO orchestration.
- replicate_all (411 lines -> 113-line main body): initial target-info
  seeding, read/stat option builders, skip-path notes, target HEAD
  action resolution, and the multipart/single-put payload transport
  extracted as private free functions.
- start_mrf_processor (306 lines -> 46-line spawn body): recovery guard,
  ledger load, per-entry replay (delete/object/metadata), and retained
  entry resolution extracted; retry bookkeeping semantics preserved
  exactly (inner continue-paths push inside helpers, outer Missed push
  stays in the loop).
- apply_iam_item (255 lines -> match dispatch skeleton): one helper per
  IAM item type.

No behavior change: log texts, error paths, event emissions, and metric
counts are byte-identical; existing tests unchanged and green (238
ecstore replication/mrf/resync + 232 rustfs site-replication).
2026-08-17 17:08:45 +08:00
17 changed files with 1573 additions and 2532 deletions
Generated
-5
View File
@@ -9825,19 +9825,14 @@ name = "rustfs-madmin"
version = "1.0.0-rc.2"
dependencies = [
"hotpath",
"http 1.5.0",
"humantime",
"hyper",
"jiff",
"reqwest",
"rmp-serde",
"rustfs-signer",
"s3s",
"serde",
"serde_json",
"sysinfo",
"time",
"tokio",
]
[[package]]
-1
View File
@@ -40,7 +40,6 @@ mak = "mak"
gae = "gae"
GAE = "GAE"
thr = "thr"
mis = "mis"
# s3-tests original test names (cannot be changed)
nonexisted = "nonexisted"
consts = "consts"
+2 -2
View File
@@ -373,8 +373,8 @@ pub mod error {
pub mod erasure {
pub use crate::erasure::coding::{
BitrotReader, BitrotSelfTestError, BitrotWriter, BitrotWriterWrapper, CustomWriter, Erasure, ErasureConstructionError,
ReedSolomonEncoder, bitrot_self_test, calc_shard_size, calc_shard_size_legacy,
BitrotReader, BitrotWriter, BitrotWriterWrapper, CustomWriter, Erasure, ErasureConstructionError, ReedSolomonEncoder,
calc_shard_size, calc_shard_size_legacy,
};
}
@@ -667,6 +667,368 @@ async fn acknowledge_mrf_recovery<S: ReplicationStorage>(
Err(EcstoreError::PreconditionFailed)
}
/// Acquires the MRF recovery leader lock for the startup replay.
/// Returns `None` (after logging) when the lock cannot be created or another
/// node is already processing the backlog.
async fn acquire_mrf_recovery_guard<S: ReplicationStorage>(storage: &Arc<S>) -> Option<rustfs_lock::NamespaceLockGuard> {
let recovery_lock = match storage
.new_ns_lock(
ReplicationMetadataStore::rustfs_meta_bucket(),
ReplicationMetadataStore::MRF_REPLICATION_RECOVERY_LOCK,
)
.await
{
Ok(lock) => lock,
Err(error) => {
warn!(
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_REPLICATION,
error = %error,
"Failed to create the MRF recovery leader lock"
);
return None;
}
};
match recovery_lock
.get_write_lock_quiet(ReplicationLockTiming::acquire_timeout())
.await
{
Ok(guard) => Some(guard),
Err(_) => {
debug!(
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_REPLICATION,
"Another node is already processing the MRF recovery backlog"
);
None
}
}
}
/// Reads and decodes the on-disk MRF recovery file.
/// Returns `None` when there is nothing to replay: missing file (publishes an
/// empty available summary), read failure, or corrupt data (quarantined).
async fn load_mrf_recovery_entries<S: ReplicationStorage>(storage: &Arc<S>) -> Option<Vec<MrfReplicateEntry>> {
let data = match ReplicationConfigStore::read(storage.clone(), ReplicationMetadataStore::MRF_REPLICATION_FILE).await {
Ok(d) => d,
Err(EcstoreError::ConfigNotFound) => {
set_durable_mrf_backlog_summary(DurableMrfBacklogSummary {
available: true,
buckets: Vec::new(),
});
return None;
}
Err(e) => {
warn!(
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_REPLICATION,
error = %e,
"Failed to load MRF recovery file"
);
return None;
}
};
match decode_mrf_file(&data) {
Ok(v) => Some(v),
Err(e) => {
warn!(
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_REPLICATION,
error = %e,
"Failed to decode MRF recovery file — preserving corrupt data"
);
quarantine_mrf_file(storage, &data).await;
None
}
}
}
/// Replays one MRF recovery entry by operation kind.
/// Returns `None` when the entry is skipped entirely (no admission outcome);
/// entries that must be retried later are pushed onto `retry_entries`.
async fn replay_mrf_entry<S: ReplicationStorage>(
entry: &MrfReplicateEntry,
storage: &Arc<S>,
retry_entries: &mut Vec<MrfReplicateEntry>,
) -> Option<ReplicationQueueAdmission> {
match entry.op {
MrfOpKind::Delete => replay_mrf_delete_entry(entry, storage, retry_entries).await,
MrfOpKind::Object | MrfOpKind::Heal | MrfOpKind::ExistingObject => {
replay_mrf_object_entry(entry, storage, retry_entries).await
}
MrfOpKind::Metadata => replay_mrf_metadata_entry(entry, storage, retry_entries).await,
}
}
/// Replays a delete-kind MRF entry: force-delete intents replay directly,
/// stale force-delete generations are skipped, and plain deletes are
/// reconstructed as heal deletes.
async fn replay_mrf_delete_entry<S: ReplicationStorage>(
entry: &MrfReplicateEntry,
storage: &Arc<S>,
retry_entries: &mut Vec<MrfReplicateEntry>,
) -> Option<ReplicationQueueAdmission> {
if should_replay_force_delete_intent(entry) {
let operation_id = entry.force_delete_id?;
let delete = force_delete_heal_replication_info(entry, operation_id);
if replicate_delete_with_outcome(delete, storage.clone()).await {
Some(ReplicationQueueAdmission::Queued)
} else {
Some(ReplicationQueueAdmission::Missed)
}
} else if entry.force_delete_id.is_some() {
Some(ReplicationQueueAdmission::Skipped)
} else {
replay_mrf_reconstructed_delete(entry, storage, retry_entries).await
}
}
/// Pure DTO construction: heal replication info for a replayed force-delete intent.
fn force_delete_heal_replication_info(entry: &MrfReplicateEntry, operation_id: uuid::Uuid) -> DeletedObjectReplicationInfo {
DeletedObjectReplicationInfo {
delete_object: ReplicationDeletedObject {
object_name: entry.object.clone(),
force_delete: true,
force_delete_id: Some(operation_id),
force_delete_target_arns: entry.target_arns.clone(),
force_delete_generation: entry.force_delete_generation,
..Default::default()
},
bucket: entry.bucket.clone(),
op_type: ReplicationType::Heal,
event_type: REPLICATE_HEAL_DELETE.to_string(),
..Default::default()
}
}
/// Reconstruct a heal delete and re-queue it. We do NOT call
/// get_object_info here because the delete-marker or version may
/// already be absent from the local store — that is expected.
async fn replay_mrf_reconstructed_delete<S: ReplicationStorage>(
entry: &MrfReplicateEntry,
storage: &Arc<S>,
retry_entries: &mut Vec<MrfReplicateEntry>,
) -> Option<ReplicationQueueAdmission> {
let versioned = ReplicationVersioningStore::prefix_enabled(&entry.bucket, &entry.object).await;
let oi = ObjectInfo {
bucket: entry.bucket.clone(),
name: entry.object.clone(),
version_id: entry.version_id,
delete_marker: entry.delete_marker,
..Default::default()
};
let dsc = resolve_mrf_delete_replicate_decision(entry, &oi, versioned, retry_entries).await?;
let dv = reconstructed_heal_delete_info(entry, &oi, &dsc);
if replicate_delete_with_outcome(dv, storage.clone()).await {
Some(ReplicationQueueAdmission::Queued)
} else {
Some(ReplicationQueueAdmission::Missed)
}
}
/// The MRF entry does not persist the replication decision and the
/// source object is gone, so re-derive the decision from the live
/// bucket config (mirroring get_heal_replicate_object_info) and set
/// it on the reconstructed delete. Without this the decision string
/// is empty and the delete replicates to zero targets — a silent
/// no-op that leaves replicas diverged (backlog#858 / #799 B9).
async fn resolve_mrf_delete_replicate_decision(
entry: &MrfReplicateEntry,
oi: &ObjectInfo,
versioned: bool,
retry_entries: &mut Vec<MrfReplicateEntry>,
) -> Option<ReplicateDecision> {
if entry.target_arns.is_empty() {
match ReplicationMetadataStore::optional_replication_config(&entry.bucket).await {
Ok(None) => None,
Err(_) => {
retry_entries.push(entry.clone());
None
}
Ok(Some(_)) => match check_replicate_delete_strict(
&entry.bucket,
&ObjectToDelete {
object_name: entry.object.clone(),
version_id: entry.version_id,
..Default::default()
},
oi,
&ObjectOptions {
versioned,
..Default::default()
},
None,
)
.await
{
Ok(dsc) => Some(dsc),
Err(_) => {
retry_entries.push(entry.clone());
None
}
},
}
} else {
Some(replicate_decision_for_admitted_targets(&entry.target_arns))
}
}
/// Pure DTO construction: reconstructed heal delete carrying the re-derived
/// replication decision.
fn reconstructed_heal_delete_info(
entry: &MrfReplicateEntry,
oi: &ObjectInfo,
dsc: &ReplicateDecision,
) -> DeletedObjectReplicationInfo {
let mut rstate = oi.replication_state();
rstate.replicate_decision_str = dsc.to_string();
let delete_marker_mtime = entry
.delete_marker_mtime
.and_then(|nanos| OffsetDateTime::from_unix_timestamp_nanos(i128::from(nanos)).ok());
DeletedObjectReplicationInfo {
delete_object: ReplicationDeletedObject {
object_name: entry.object.clone(),
version_id: entry.version_id,
delete_marker_version_id: entry.delete_marker_version_id,
delete_marker: entry.delete_marker,
delete_marker_mtime,
force_delete: entry.force_delete,
replication_state: Some(rstate),
..Default::default()
},
bucket: entry.bucket.clone(),
op_type: ReplicationType::Heal,
event_type: REPLICATE_HEAL_DELETE.to_string(),
..Default::default()
}
}
/// Replays an Object/Heal/ExistingObject MRF entry against the live source object.
async fn replay_mrf_object_entry<S: ReplicationStorage>(
entry: &MrfReplicateEntry,
storage: &Arc<S>,
retry_entries: &mut Vec<MrfReplicateEntry>,
) -> Option<ReplicationQueueAdmission> {
let opts = ObjectOptions {
version_id: entry.version_id.map(|u| u.to_string()),
..Default::default()
};
let oi = match storage.get_object_info(&entry.bucket, &entry.object, &opts).await {
Ok(oi) => oi,
Err(e) => {
debug!(
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_REPLICATION,
bucket = %entry.bucket,
object = %entry.object,
error = %e,
"MRF recovery: source object lookup failed"
);
if should_retry_mrf_source_lookup(&e) {
retry_entries.push(entry.clone());
}
return None;
}
};
if entry.target_arns.is_empty() {
// Legacy entries predate target admission persistence. They cannot
// be safely attributed, so retain the old live-config fallback.
Some(queue_replication_heal(&entry.bucket, oi, entry.retry_count.max(0) as u32).await)
} else {
let roi = admitted_mrf_replicate_object(oi, entry, entry.op.replication_type());
if replicate_object_with_outcome(roi, storage.clone()).await.1 {
Some(ReplicationQueueAdmission::Queued)
} else {
Some(ReplicationQueueAdmission::Missed)
}
}
}
/// Replays a metadata-kind MRF entry against the live source object.
async fn replay_mrf_metadata_entry<S: ReplicationStorage>(
entry: &MrfReplicateEntry,
storage: &Arc<S>,
retry_entries: &mut Vec<MrfReplicateEntry>,
) -> Option<ReplicationQueueAdmission> {
let opts = ObjectOptions {
version_id: entry.version_id.map(|u| u.to_string()),
..Default::default()
};
let oi = match storage.get_object_info(&entry.bucket, &entry.object, &opts).await {
Ok(oi) => oi,
Err(e) => {
debug!(
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_REPLICATION,
bucket = %entry.bucket,
object = %entry.object,
error = %e,
"MRF metadata recovery: source object lookup failed"
);
if should_retry_mrf_source_lookup(&e) {
retry_entries.push(entry.clone());
}
return None;
}
};
if entry.target_arns.is_empty() {
Some(queue_replication_metadata(&entry.bucket, oi, entry.retry_count.max(0) as u32).await)
} else {
let roi = admitted_mrf_replicate_object(oi, entry, ReplicationType::Metadata);
if replicate_object_with_outcome(roi, storage.clone()).await.1 {
Some(ReplicationQueueAdmission::Queued)
} else {
Some(ReplicationQueueAdmission::Missed)
}
}
}
/// Pure DTO construction: replicate-object info for an entry with persisted
/// admitted targets, carrying over the entry's retry count.
fn admitted_mrf_replicate_object(oi: ObjectInfo, entry: &MrfReplicateEntry, op_type: ReplicationType) -> ReplicateObjectInfo {
let dsc = replicate_decision_for_admitted_targets(&entry.target_arns);
let mut roi = replicate_object_info_from_object_info(oi, dsc, op_type);
roi.retry_count = entry.retry_count.max(0) as u32;
roi
}
/// Acknowledges the replayed MRF prefix and returns the retained backlog.
/// On acknowledgement failure the backlog is preserved for the next startup and
/// re-read (falling back to the replayed snapshot) so the published summary stays accurate.
async fn resolve_retained_mrf_entries<S: ReplicationStorage>(
storage: &Arc<S>,
recovery_guard: &rustfs_lock::NamespaceLockGuard,
entries: &[MrfReplicateEntry],
retry_entries: &[MrfReplicateEntry],
) -> Vec<MrfReplicateEntry> {
match acknowledge_mrf_recovery(storage.clone(), recovery_guard, entries, retry_entries).await {
Ok(retained) => retained,
Err(error) => {
warn!(
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_REPLICATION,
error = %error,
"Failed to acknowledge the MRF recovery prefix; preserving it for the next startup"
);
match read_mrf_entries(storage.clone()).await {
Ok(current) => current,
Err(read_error) => {
warn!(
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_REPLICATION,
error = %read_error,
"Failed to refresh the MRF backlog after acknowledgement failure"
);
entries.to_vec()
}
}
}
}
}
#[derive(Debug, thiserror::Error)]
#[error("replication resync {active_resync_id} is already active for {bucket}/{arn}")]
struct ResyncActiveConflictError {
@@ -1221,71 +1583,12 @@ impl<S: ReplicationStorage> ReplicationPool<S> {
let storage = self.storage.clone();
let handle = tokio::spawn(async move {
let recovery_lock = match storage
.new_ns_lock(
ReplicationMetadataStore::rustfs_meta_bucket(),
ReplicationMetadataStore::MRF_REPLICATION_RECOVERY_LOCK,
)
.await
{
Ok(lock) => lock,
Err(error) => {
warn!(
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_REPLICATION,
error = %error,
"Failed to create the MRF recovery leader lock"
);
return;
}
};
let recovery_guard = match recovery_lock
.get_write_lock_quiet(ReplicationLockTiming::acquire_timeout())
.await
{
Ok(guard) => guard,
Err(_) => {
debug!(
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_REPLICATION,
"Another node is already processing the MRF recovery backlog"
);
return;
}
let Some(recovery_guard) = acquire_mrf_recovery_guard(&storage).await else {
return;
};
let data = match ReplicationConfigStore::read(storage.clone(), ReplicationMetadataStore::MRF_REPLICATION_FILE).await {
Ok(d) => d,
Err(EcstoreError::ConfigNotFound) => {
set_durable_mrf_backlog_summary(DurableMrfBacklogSummary {
available: true,
buckets: Vec::new(),
});
return;
}
Err(e) => {
warn!(
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_REPLICATION,
error = %e,
"Failed to load MRF recovery file"
);
return;
}
};
let entries = match decode_mrf_file(&data) {
Ok(v) => v,
Err(e) => {
warn!(
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_REPLICATION,
error = %e,
"Failed to decode MRF recovery file — preserving corrupt data"
);
quarantine_mrf_file(&storage, &data).await;
return;
}
let Some(entries) = load_mrf_recovery_entries(&storage).await else {
return;
};
set_durable_mrf_backlog_snapshot(durable_mrf_backlog_summary_from_entries(&entries));
@@ -1294,187 +1597,8 @@ impl<S: ReplicationStorage> ReplicationPool<S> {
let mut retry_entries = Vec::new();
for entry in entries.iter() {
let admission = match entry.op {
MrfOpKind::Delete => {
if should_replay_force_delete_intent(entry) {
let Some(operation_id) = entry.force_delete_id else {
continue;
};
let delete = DeletedObjectReplicationInfo {
delete_object: ReplicationDeletedObject {
object_name: entry.object.clone(),
force_delete: true,
force_delete_id: Some(operation_id),
force_delete_target_arns: entry.target_arns.clone(),
force_delete_generation: entry.force_delete_generation,
..Default::default()
},
bucket: entry.bucket.clone(),
op_type: ReplicationType::Heal,
event_type: REPLICATE_HEAL_DELETE.to_string(),
..Default::default()
};
if replicate_delete_with_outcome(delete, storage.clone()).await {
ReplicationQueueAdmission::Queued
} else {
ReplicationQueueAdmission::Missed
}
} else if entry.force_delete_id.is_some() {
ReplicationQueueAdmission::Skipped
} else {
// Reconstruct a heal delete and re-queue it. We do NOT call
// get_object_info here because the delete-marker or version may
// already be absent from the local store — that is expected.
//
// The MRF entry does not persist the replication decision and the
// source object is gone, so re-derive the decision from the live
// bucket config (mirroring get_heal_replicate_object_info) and set
// it on the reconstructed delete. Without this the decision string
// is empty and the delete replicates to zero targets — a silent
// no-op that leaves replicas diverged (backlog#858 / #799 B9).
let versioned = ReplicationVersioningStore::prefix_enabled(&entry.bucket, &entry.object).await;
let oi = ObjectInfo {
bucket: entry.bucket.clone(),
name: entry.object.clone(),
version_id: entry.version_id,
delete_marker: entry.delete_marker,
..Default::default()
};
let dsc = if entry.target_arns.is_empty() {
match ReplicationMetadataStore::optional_replication_config(&entry.bucket).await {
Ok(None) => continue,
Err(_) => {
retry_entries.push(entry.clone());
continue;
}
Ok(Some(_)) => match check_replicate_delete_strict(
&entry.bucket,
&ObjectToDelete {
object_name: entry.object.clone(),
version_id: entry.version_id,
..Default::default()
},
&oi,
&ObjectOptions {
versioned,
..Default::default()
},
None,
)
.await
{
Ok(dsc) => dsc,
Err(_) => {
retry_entries.push(entry.clone());
continue;
}
},
}
} else {
replicate_decision_for_admitted_targets(&entry.target_arns)
};
let mut rstate = oi.replication_state();
rstate.replicate_decision_str = dsc.to_string();
let delete_marker_mtime = entry
.delete_marker_mtime
.and_then(|nanos| OffsetDateTime::from_unix_timestamp_nanos(i128::from(nanos)).ok());
let dv = DeletedObjectReplicationInfo {
delete_object: ReplicationDeletedObject {
object_name: entry.object.clone(),
version_id: entry.version_id,
delete_marker_version_id: entry.delete_marker_version_id,
delete_marker: entry.delete_marker,
delete_marker_mtime,
force_delete: entry.force_delete,
replication_state: Some(rstate),
..Default::default()
},
bucket: entry.bucket.clone(),
op_type: ReplicationType::Heal,
event_type: REPLICATE_HEAL_DELETE.to_string(),
..Default::default()
};
if replicate_delete_with_outcome(dv, storage.clone()).await {
ReplicationQueueAdmission::Queued
} else {
ReplicationQueueAdmission::Missed
}
}
}
MrfOpKind::Object | MrfOpKind::Heal | MrfOpKind::ExistingObject => {
let opts = ObjectOptions {
version_id: entry.version_id.map(|u| u.to_string()),
..Default::default()
};
let oi = match storage.get_object_info(&entry.bucket, &entry.object, &opts).await {
Ok(oi) => oi,
Err(e) => {
debug!(
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_REPLICATION,
bucket = %entry.bucket,
object = %entry.object,
error = %e,
"MRF recovery: source object lookup failed"
);
if should_retry_mrf_source_lookup(&e) {
retry_entries.push(entry.clone());
}
continue;
}
};
if entry.target_arns.is_empty() {
// Legacy entries predate target admission persistence. They cannot
// be safely attributed, so retain the old live-config fallback.
queue_replication_heal(&entry.bucket, oi, entry.retry_count.max(0) as u32).await
} else {
let dsc = replicate_decision_for_admitted_targets(&entry.target_arns);
let mut roi = replicate_object_info_from_object_info(oi, dsc, entry.op.replication_type());
roi.retry_count = entry.retry_count.max(0) as u32;
if replicate_object_with_outcome(roi, storage.clone()).await.1 {
ReplicationQueueAdmission::Queued
} else {
ReplicationQueueAdmission::Missed
}
}
}
MrfOpKind::Metadata => {
let opts = ObjectOptions {
version_id: entry.version_id.map(|u| u.to_string()),
..Default::default()
};
let oi = match storage.get_object_info(&entry.bucket, &entry.object, &opts).await {
Ok(oi) => oi,
Err(e) => {
debug!(
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_REPLICATION,
bucket = %entry.bucket,
object = %entry.object,
error = %e,
"MRF metadata recovery: source object lookup failed"
);
if should_retry_mrf_source_lookup(&e) {
retry_entries.push(entry.clone());
}
continue;
}
};
if entry.target_arns.is_empty() {
queue_replication_metadata(&entry.bucket, oi, entry.retry_count.max(0) as u32).await
} else {
let dsc = replicate_decision_for_admitted_targets(&entry.target_arns);
let mut roi = replicate_object_info_from_object_info(oi, dsc, ReplicationType::Metadata);
roi.retry_count = entry.retry_count.max(0) as u32;
if replicate_object_with_outcome(roi, storage.clone()).await.1 {
ReplicationQueueAdmission::Queued
} else {
ReplicationQueueAdmission::Missed
}
}
}
let Some(admission) = replay_mrf_entry(entry, &storage, &mut retry_entries).await else {
continue;
};
if admission == ReplicationQueueAdmission::Missed {
@@ -1484,29 +1608,7 @@ impl<S: ReplicationStorage> ReplicationPool<S> {
}
}
let retained = match acknowledge_mrf_recovery(storage.clone(), &recovery_guard, &entries, &retry_entries).await {
Ok(retained) => retained,
Err(error) => {
warn!(
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_REPLICATION,
error = %error,
"Failed to acknowledge the MRF recovery prefix; preserving it for the next startup"
);
match read_mrf_entries(storage.clone()).await {
Ok(current) => current,
Err(read_error) => {
warn!(
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_REPLICATION,
error = %read_error,
"Failed to refresh the MRF backlog after acknowledgement failure"
);
entries.clone()
}
}
}
};
let retained = resolve_retained_mrf_entries(&storage, &recovery_guard, &entries, &retry_entries).await;
let retained_count = retained.len();
set_durable_mrf_backlog_snapshot(durable_mrf_backlog_summary_from_entries(&retained));
File diff suppressed because it is too large Load Diff
+12 -291
View File
@@ -820,263 +820,10 @@ impl BitrotWriterWrapper {
}
}
// --- startup bitrot self-test (rustfs/backlog#1873, MinIO bitrotSelfTest parity) ---
//
// A broken hash implementation (bad SIMD feature combination, platform drift, a
// key-handling regression) fails silently: every shard reads back "corrupt",
// heal rewrites data that was fine, and cross-platform clusters disagree about
// which copy is healthy. The self-test below pins the algorithms the moment a
// process starts, so a drifted build announces itself instead of quietly
// rewriting objects. See docs/rustfs-heal-scanner-vs-minio-comprehensive-
// analysis-2026-08-16.md §6 HS-11.
/// Length of the deterministic self-test payload.
pub const BITROT_SELF_TEST_PAYLOAD_LEN: usize = 4096;
/// Known-answer digest of [`bitrot_self_test_payload`] under `HighwayHash256S`
/// (the production default). Pinned so any platform or build where the
/// implementation drifts fails startup instead of miss-hashing shards.
const BITROT_SELF_TEST_KAT_HIGHWAY_HASH256S: [u8; 32] = [
0xb9, 0x32, 0xa2, 0xaa, 0x4a, 0xb7, 0x33, 0x6a, 0xa3, 0xca, 0x7e, 0x61, 0x9d, 0x86, 0x52, 0x14, 0x6e, 0x7f, 0xd8, 0x9e, 0xea,
0x08, 0xd9, 0x8c, 0x33, 0x85, 0x87, 0x19, 0x30, 0xd6, 0xed, 0x06,
];
/// Known-answer digest of the same payload under `HighwayHash256SLegacy`.
const BITROT_SELF_TEST_KAT_HIGHWAY_HASH256S_LEGACY: [u8; 32] = [
0x98, 0x24, 0x71, 0x4f, 0x16, 0xbb, 0x48, 0x39, 0xed, 0x68, 0xfa, 0x63, 0x5e, 0xd9, 0x07, 0x61, 0xdf, 0x0a, 0xff, 0xcf, 0x7d,
0x8c, 0xa8, 0xc7, 0xc0, 0xb6, 0x6f, 0x05, 0xdb, 0xda, 0x5a, 0x22,
];
/// FIPS 180-2 test vector: SHA-256 of the ASCII string "abc". Unlike the
/// Highway digests above this one is externally verifiable, so it guards the
/// whole `HashAlgorithm` plumbing even for readers who distrust pinned
/// self-computed constants.
const BITROT_SELF_TEST_KAT_SHA256_ABC: [u8; 32] = [
0xba, 0x78, 0x16, 0xbf, 0x8f, 0x01, 0xcf, 0xea, 0x41, 0x41, 0x40, 0xde, 0x5d, 0xae, 0x22, 0x23, 0xb0, 0x03, 0x61, 0xa3, 0x96,
0x17, 0x7a, 0x9c, 0xb4, 0x10, 0xff, 0x61, 0xf2, 0x00, 0x15, 0xad,
];
/// Deterministic self-test payload: xorshift64* from a fixed seed, so every
/// platform and every run hashes the same 4096 bytes.
fn bitrot_self_test_payload() -> [u8; BITROT_SELF_TEST_PAYLOAD_LEN] {
let mut state = 0x9E37_79B9_7F4A_7C15u64;
let mut payload = [0u8; BITROT_SELF_TEST_PAYLOAD_LEN];
for byte in payload.iter_mut() {
state ^= state >> 12;
state ^= state << 25;
state ^= state >> 27;
*byte = state.wrapping_mul(0x2545_F491_4F6C_DD1D) as u8;
}
payload
}
/// Why a bitrot self-test failed.
#[derive(Debug)]
pub enum BitrotSelfTestError {
/// A known-answer digest mismatched the pinned constant.
KnownAnswerMismatch {
algorithm: &'static str,
got: String,
want: String,
},
/// A freshly encoded shard failed `bitrot_verify`.
RoundtripVerify { algorithm: &'static str, detail: String },
/// A verified roundtrip read back different bytes than were written.
RoundtripReadback { algorithm: &'static str },
/// A deliberately tampered shard was not rejected by `bitrot_verify`.
TamperNotRejected {
algorithm: &'static str,
tampered: &'static str,
},
}
impl std::fmt::Display for BitrotSelfTestError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::KnownAnswerMismatch { algorithm, got, want } => {
write!(f, "known-answer mismatch for {algorithm}: got {got}, want {want}")
}
Self::RoundtripVerify { algorithm, detail } => write!(f, "{algorithm} roundtrip shard failed verification: {detail}"),
Self::RoundtripReadback { algorithm } => write!(f, "{algorithm} roundtrip read back different bytes"),
Self::TamperNotRejected { algorithm, tampered } => {
write!(f, "{algorithm} tampered shard ({tampered}) was not rejected")
}
}
}
}
impl std::error::Error for BitrotSelfTestError {}
fn self_test_hex(bytes: &[u8]) -> String {
rustfs_utils::hex(bytes)
}
// (kept as a named one-liner so every KAT failure site reads the same; the
// underlying formatter is the shared `rustfs_utils::hex`)
/// Compare a digest against its pinned constant. Split out so a test can drive
/// it with a wrong constant and prove the mismatch path fires.
fn bitrot_kat_check(
algorithm: &'static str,
algo: &HashAlgorithm,
payload: &[u8],
expected: &[u8; 32],
) -> Result<(), BitrotSelfTestError> {
let digest = algo.hash_encode(payload);
let digest = digest.as_ref();
if digest.len() != expected.len() || digest != expected.as_slice() {
return Err(BitrotSelfTestError::KnownAnswerMismatch {
algorithm,
got: self_test_hex(digest),
want: self_test_hex(expected),
});
}
Ok(())
}
/// Encode `payload` with `shard_size` blocks, verify it end to end, and read
/// every block back through `BitrotReader` comparing bytes.
async fn bitrot_roundtrip_check(
algorithm: &'static str,
algo: HashAlgorithm,
payload: &[u8],
shard_size: usize,
) -> Result<(), BitrotSelfTestError> {
let mut writer = BitrotWriter::new(std::io::Cursor::new(Vec::<u8>::new()), shard_size, algo.clone());
for chunk in payload.chunks(shard_size) {
writer
.write(chunk)
.await
.map_err(|err| BitrotSelfTestError::RoundtripVerify {
algorithm,
detail: format!("encode failed: {err}"),
})?;
}
let encoded = writer.into_inner().into_inner();
let on_disk = bitrot_shard_file_size(payload.len(), shard_size, algo.clone());
if encoded.len() != on_disk {
return Err(BitrotSelfTestError::RoundtripVerify {
algorithm,
detail: format!("encoded {} bytes, size formula says {on_disk}", encoded.len()),
});
}
bitrot_verify(std::io::Cursor::new(encoded.clone()), on_disk, payload.len(), algo.clone(), shard_size)
.await
.map_err(|err| BitrotSelfTestError::RoundtripVerify {
algorithm,
detail: err.to_string(),
})?;
let mut reader = BitrotReader::new(std::io::Cursor::new(encoded), shard_size, algo, false);
let mut offset = 0usize;
while offset < payload.len() {
let want = shard_size.min(payload.len() - offset);
let mut buf = vec![0u8; want];
let read = reader
.read(&mut buf)
.await
.map_err(|err| BitrotSelfTestError::RoundtripVerify {
algorithm,
detail: format!("read back failed at offset {offset}: {err}"),
})?;
if read != want || buf[..read] != payload[offset..offset + read] {
return Err(BitrotSelfTestError::RoundtripReadback { algorithm });
}
offset += read;
}
Ok(())
}
/// Flip one byte and require `bitrot_verify` to reject the result.
async fn bitrot_tamper_check(
algorithm: &'static str,
algo: HashAlgorithm,
payload: &[u8],
shard_size: usize,
tampered: &'static str,
flip_at: usize,
) -> Result<(), BitrotSelfTestError> {
let mut writer = BitrotWriter::new(std::io::Cursor::new(Vec::<u8>::new()), shard_size, algo.clone());
for chunk in payload.chunks(shard_size) {
writer.write(chunk).await.expect("self-test encode should not fail");
}
let mut corrupt = writer.into_inner().into_inner();
let flip_index = flip_at % corrupt.len();
corrupt[flip_index] ^= 0x80;
let on_disk = bitrot_shard_file_size(payload.len(), shard_size, algo.clone());
match bitrot_verify(std::io::Cursor::new(corrupt), on_disk, payload.len(), algo, shard_size).await {
// The flipped byte must be rejected as a hash mismatch specifically, not
// by any incidental read error: an in-memory cursor cannot fail reads,
// so accepting any other failure here would mask a verify path that
// errors out before it ever compares hashes.
Err(err) if err.to_string().contains("hash mismatch") => Ok(()),
Ok(()) => Err(BitrotSelfTestError::TamperNotRejected { algorithm, tampered }),
Err(err) => Err(BitrotSelfTestError::RoundtripVerify {
algorithm,
detail: format!("tampered shard rejected with an unexpected error: {err}"),
}),
}
}
/// Verify every bitrot algorithm this crate can write or verify in production:
/// both streaming Highway variants roundtrip end to end (encode → size formula
/// → `bitrot_verify` → read back) and reject a flipped byte in both the data
/// and the leading hash, while all three hashed algorithms reproduce their
/// pinned known-answer digests.
///
/// Runs in well under a millisecond on 4 KiB of data; callers may run it inline
/// at startup. Pure CPU, no allocation beyond a few KiB of scratch.
pub async fn bitrot_self_test() -> Result<(), BitrotSelfTestError> {
let payload = bitrot_self_test_payload();
// Externally verifiable vector first: it guards the HashAlgorithm plumbing
// itself, before any self-pinned constants are consulted.
let abc = HashAlgorithm::SHA256.hash_encode(b"abc");
if abc.as_ref() != BITROT_SELF_TEST_KAT_SHA256_ABC.as_slice() {
return Err(BitrotSelfTestError::KnownAnswerMismatch {
algorithm: "SHA256",
got: self_test_hex(abc.as_ref()),
want: self_test_hex(&BITROT_SELF_TEST_KAT_SHA256_ABC),
});
}
bitrot_kat_check(
"HighwayHash256S",
&HashAlgorithm::HighwayHash256S,
&payload,
&BITROT_SELF_TEST_KAT_HIGHWAY_HASH256S,
)?;
bitrot_kat_check(
"HighwayHash256SLegacy",
&HashAlgorithm::HighwayHash256SLegacy,
&payload,
&BITROT_SELF_TEST_KAT_HIGHWAY_HASH256S_LEGACY,
)?;
for (algorithm, algo) in [
("HighwayHash256S", HashAlgorithm::HighwayHash256S),
("HighwayHash256SLegacy", HashAlgorithm::HighwayHash256SLegacy),
] {
// Full blocks plus a partial tail, exactly like a real part stripe.
let tail_len = 2 * 1024 + 333;
bitrot_roundtrip_check(algorithm, algo.clone(), &payload, 1024).await?;
bitrot_roundtrip_check(algorithm, algo.clone(), &payload[..tail_len], 1024).await?;
// One flipped byte in the final data block, one in the first leading
// hash: both must fail verification.
bitrot_tamper_check(algorithm, algo.clone(), &payload, 1024, "final data byte", payload.len() - 1).await?;
bitrot_tamper_check(algorithm, algo, &payload, 1024, "leading hash byte", 0).await?;
}
Ok(())
}
#[cfg(test)]
mod tests {
use super::{
BitrotReader, BitrotWriter, BitrotWriterWrapper, CustomWriter, bitrot_kat_check, bitrot_self_test,
bitrot_self_test_payload, bitrot_shard_file_size, bitrot_verify, write_all_vectored,
BitrotReader, BitrotWriter, BitrotWriterWrapper, CustomWriter, bitrot_shard_file_size, bitrot_verify, write_all_vectored,
};
use super::{MAX_RETAINED_CHUNKS_PER_BLOCK, ShardChunkRead, ShardSource};
use bytes::Bytes;
@@ -1343,32 +1090,6 @@ mod tests {
}
}
#[test]
fn bitrot_self_test_payload_is_deterministic() {
// Two independent builds of the payload must agree byte for byte, or
// the pinned known-answer digests below would be meaningless.
assert_eq!(bitrot_self_test_payload(), bitrot_self_test_payload());
}
#[test]
fn bitrot_self_test_rejects_a_wrong_known_answer_digest() {
let payload = bitrot_self_test_payload();
let wrong = [0u8; 32];
let err = bitrot_kat_check("HighwayHash256S", &HashAlgorithm::HighwayHash256S, &payload, &wrong)
.expect_err("a zeroed digest must never match");
match err {
super::BitrotSelfTestError::KnownAnswerMismatch { algorithm, .. } => assert_eq!(algorithm, "HighwayHash256S"),
other => panic!("expected KnownAnswerMismatch, got {other:?}"),
}
}
#[tokio::test]
async fn bitrot_self_test_passes() {
bitrot_self_test()
.await
.expect("the pinned digests and roundtrip checks must all pass on this platform");
}
#[tokio::test]
async fn vectored_test_writers_cover_fallback_flush_and_shutdown_paths() {
let mut counting = VectoredCountingWriter::default();
@@ -1468,7 +1189,7 @@ mod tests {
let last = corrupt.len() - 1;
corrupt[last] ^= 0x80;
let err = bitrot_verify(
std::io::Cursor::new(corrupt),
Cursor::new(corrupt),
super::bitrot_shard_file_size(data.len(), shard_size, algo.clone()),
data.len(),
algo,
@@ -1561,7 +1282,7 @@ mod tests {
#[tokio::test]
async fn bitrot_reader_rejects_output_buffers_larger_than_shard_size() {
let mut reader = BitrotReader::new(std::io::Cursor::new(Vec::<u8>::new()), 4, HashAlgorithm::None, false);
let mut reader = BitrotReader::new(Cursor::new(Vec::<u8>::new()), 4, HashAlgorithm::None, false);
let mut out = [0u8; 5];
let err = reader
.read(&mut out)
@@ -1686,7 +1407,7 @@ mod tests {
(HashAlgorithm::HighwayHash256, true),
] {
let label = format!("{algo:?}");
let writer = std::io::Cursor::new(Vec::<u8>::new());
let writer = Cursor::new(Vec::<u8>::new());
let mut w = BitrotWriter::new(writer, shard_size, algo.clone());
w.write(&[7u8; 16]).await.unwrap();
let written = w.into_inner().into_inner();
@@ -1771,7 +1492,7 @@ mod tests {
}
async fn encode_one_block(payload: &[u8], shard_size: usize, algo: HashAlgorithm) -> Vec<u8> {
let mut w = BitrotWriter::new(std::io::Cursor::new(Vec::<u8>::new()), shard_size, algo);
let mut w = BitrotWriter::new(Cursor::new(Vec::<u8>::new()), shard_size, algo);
w.write(payload).await.unwrap();
w.into_inner().into_inner()
}
@@ -1879,7 +1600,7 @@ mod tests {
for algo in [HashAlgorithm::HighwayHash256S, HashAlgorithm::HighwayHash256SLegacy] {
for &size in &[1usize, 16, 17, 32, 40, 48] {
let payload: Vec<u8> = (0..size).map(|i| i as u8).collect();
let mut w = BitrotWriter::new(std::io::Cursor::new(Vec::<u8>::new()), shard_size, algo.clone());
let mut w = BitrotWriter::new(Cursor::new(Vec::<u8>::new()), shard_size, algo.clone());
for chunk in payload.chunks(shard_size) {
w.write(chunk).await.unwrap();
}
@@ -1953,14 +1674,14 @@ mod tests {
w.write(&data).await.expect("write shard");
let mut via_read = vec![0u8; SHARD];
let n1 = BitrotReader::new(std::io::Cursor::new(encoded.clone()), SHARD, algo.clone(), false)
let n1 = BitrotReader::new(Cursor::new(encoded.clone()), SHARD, algo.clone(), false)
.read(&mut via_read)
.await
.expect("read");
// A buffer with only capacity — no initialized bytes at all.
let mut via_append: Vec<u8> = Vec::with_capacity(SHARD);
let n2 = BitrotReader::new(std::io::Cursor::new(encoded), SHARD, algo.clone(), false)
let n2 = BitrotReader::new(Cursor::new(encoded), SHARD, algo.clone(), false)
.read_appending(&mut via_append, SHARD)
.await
.expect("read_appending");
@@ -1985,7 +1706,7 @@ mod tests {
encoded.truncate(encoded.len() - 1);
let mut out: Vec<u8> = Vec::with_capacity(SHARD);
let err = BitrotReader::new(std::io::Cursor::new(encoded), SHARD, algo.clone(), false)
let err = BitrotReader::new(Cursor::new(encoded), SHARD, algo.clone(), false)
.read_appending(&mut out, SHARD)
.await
.expect_err("a truncated shard must not succeed");
@@ -2011,7 +1732,7 @@ mod tests {
encoded[last] ^= 0xff;
let mut out: Vec<u8> = Vec::with_capacity(SHARD);
let err = BitrotReader::new(std::io::Cursor::new(encoded), SHARD, algo, false)
let err = BitrotReader::new(Cursor::new(encoded), SHARD, algo, false)
.read_appending(&mut out, SHARD)
.await
.expect_err("a corrupt shard must not verify");
@@ -2123,7 +1844,7 @@ mod tests {
"Cursor<Bytes> must be able to hand out a block, otherwise the fast path is dead code"
);
assert_eq!(mem.position(), 8, "taking a block must advance like a read of the same length");
let mut streamed = std::io::Cursor::new(encoded.clone());
let mut streamed = Cursor::new(encoded.clone());
assert!(
ShardSource::try_take_block(&mut streamed, 8).is_none(),
"a non-Bytes source must stay on the streaming path"
@@ -2151,7 +1872,7 @@ mod tests {
);
let mut via_stream: Vec<u8> = Vec::with_capacity(SHARD);
BitrotReader::new(std::io::Cursor::new(encoded), SHARD, algo, false)
BitrotReader::new(Cursor::new(encoded), SHARD, algo, false)
.read_appending(&mut via_stream, SHARD)
.await
.expect("streaming read");
-5
View File
@@ -37,11 +37,7 @@ hotpath-cpu = ["hotpath", "hotpath/hotpath-cpu"]
[dependencies]
hotpath.workspace = true
humantime.workspace = true
http.workspace = true
hyper = { workspace = true, features = ["http2", "http1", "server"] }
reqwest = { workspace = true, features = ["json"] }
rustfs-signer.workspace = true
s3s.workspace = true
jiff = { workspace = true, features = ["serde"] }
serde = { workspace = true, features = ["derive"] }
serde_json = { workspace = true, features = ["raw_value"] }
@@ -53,4 +49,3 @@ doctest = false
[dev-dependencies]
rmp-serde.workspace = true
tokio = { workspace = true, features = ["macros", "rt-multi-thread", "net"] }
-851
View File
@@ -1,851 +0,0 @@
// 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.
//! Admin API HTTP client for heal and scanner management (rustfs/backlog#1869).
//!
//! [`AdminClient`] speaks the `/rustfs/admin/v3` surface with S3 SigV4
//! request signing (the same scheme the server's admin router authenticates),
//! so `mc`-style tooling and automation can drive heal start/query/cancel and
//! read background-heal / scanner status without hand-rolling HTTP.
//!
//! Wire structs in this module mirror the server-side shapes
//! (`rustfs/src/admin/handlers/heal.rs`, `handlers/scanner.rs`,
//! `rustfs-common/src/heal_channel.rs`), following the madmin-go model where
//! the SDK owns its own copies and round-trip tests pin the encoding. Deeply
//! nested status payloads that the server composes from runtime types are
//! carried through as `serde_json::Value` and flattened maps rather than
//! duplicated field-for-field, so the client cannot silently drift on fields
//! it never interprets.
use crate::heal_commands::HealResultItem;
use http::Method;
use serde::{Deserialize, Serialize, de};
use std::time::Duration;
/// Default admin API path prefix on a RustFS endpoint.
pub const DEFAULT_ADMIN_API_PREFIX: &str = "/rustfs/admin";
/// Default SigV4 region when the server has no explicit region configured.
pub const DEFAULT_REGION: &str = "us-east-1";
/// Scan mode for a heal request, mirroring the server's numeric-or-name wire
/// encoding (`0` unknown/default, `1` normal, `2` deep).
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub enum HealScanMode {
/// Server default; behaves as [`HealScanMode::Normal`].
#[default]
Unknown,
/// Metadata-level checks only.
Normal,
/// Full bitrot verification while healing.
Deep,
}
impl HealScanMode {
fn wire_number(self) -> u8 {
match self {
Self::Unknown => 0,
Self::Normal => 1,
Self::Deep => 2,
}
}
fn from_wire_number(value: u8) -> Option<Self> {
match value {
0 => Some(Self::Unknown),
1 => Some(Self::Normal),
2 => Some(Self::Deep),
_ => None,
}
}
fn from_wire_name(value: &str) -> Option<Self> {
match value {
"unknown" => Some(Self::Unknown),
"normal" => Some(Self::Normal),
"deep" => Some(Self::Deep),
_ => None,
}
}
}
impl Serialize for HealScanMode {
fn serialize<S: serde::Serializer>(&self, serializer: S) -> Result<S::Ok, S::Error> {
serializer.serialize_u8(self.wire_number())
}
}
impl<'de> Deserialize<'de> for HealScanMode {
fn deserialize<D: serde::Deserializer<'de>>(deserializer: D) -> Result<Self, D::Error> {
struct HealScanModeVisitor;
impl de::Visitor<'_> for HealScanModeVisitor {
type Value = HealScanMode;
fn expecting(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter.write_str("a heal scan mode number or name")
}
fn visit_u64<E: de::Error>(self, value: u64) -> Result<Self::Value, E> {
u8::try_from(value)
.ok()
.and_then(HealScanMode::from_wire_number)
.ok_or_else(|| E::custom(format!("unknown heal scan mode number: {value}")))
}
fn visit_str<E: de::Error>(self, value: &str) -> Result<Self::Value, E> {
HealScanMode::from_wire_name(value).ok_or_else(|| E::custom(format!("unknown heal scan mode name: {value}")))
}
}
deserializer.deserialize_any(HealScanModeVisitor)
}
}
/// Heal options for an admin heal request (mirror of the server body type).
/// Fields default on decode: a client should tolerate a server response whose
/// settings object omits fields it never set.
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct HealOpts {
#[serde(default)]
pub recursive: bool,
#[serde(rename = "dryRun", default)]
pub dry_run: bool,
#[serde(default)]
pub remove: bool,
#[serde(default)]
pub recreate: bool,
#[serde(rename = "scanMode", default)]
pub scan_mode: HealScanMode,
#[serde(rename = "updateParity", default)]
pub update_parity: bool,
#[serde(rename = "nolock", default)]
pub no_lock: bool,
#[serde(rename = "pool", default)]
pub pool: Option<usize>,
#[serde(rename = "set", default)]
pub set: Option<usize>,
}
/// Successful heal start / path-scoped cancel response.
#[derive(Debug, Clone, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct HealStartSuccess {
pub client_token: String,
pub client_address: String,
#[serde(default)]
pub start_time: String,
}
/// Heal task status response (query, cancel-with-token, start-then-poll).
#[derive(Debug, Clone, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct HealTaskStatus {
/// `running` | `finished` | `stopped` | `notFound`.
pub summary: String,
/// Failure detail for stopped tasks; empty otherwise.
#[serde(rename = "detail", default)]
pub failure_detail: String,
#[serde(default)]
pub start_time: String,
#[serde(default)]
pub settings: HealOpts,
#[serde(default)]
pub items: Vec<HealResultItem>,
#[serde(default)]
pub truncated: bool,
/// Live progress snapshot; the exact shape is owned by the heal runtime.
#[serde(default)]
pub progress: Option<serde_json::Value>,
}
/// `POST /v3/background-heal/status` response. Known top-level fields are
/// typed; the flattened heal info and operations matrix pass through verbatim.
#[derive(Debug, Clone, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct BackgroundHealStatus {
/// `disabled` | `uninitialized` | `idle` | `active` | `degraded`.
pub state: String,
#[serde(default)]
pub heal_queue_length: u64,
#[serde(default)]
pub heal_active_tasks: u64,
#[serde(default)]
pub cluster_status_complete: bool,
#[serde(default)]
pub progress: Option<serde_json::Value>,
/// Remaining wire fields (flattened `BackgroundHealInfo` plus the
/// priority-by-source operations matrix), carried verbatim.
#[serde(flatten)]
pub extra: serde_json::Map<String, serde_json::Value>,
}
/// `GET /v3/scanner/status` response, typed at the fields operators branch
/// on; everything else passes through verbatim.
#[derive(Debug, Clone, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct ScannerStatus {
pub enabled: bool,
/// `fresh` | `stale` | `unknown`; absent when the scanner never completed
/// a cycle.
#[serde(default)]
pub freshness: Option<ScannerFreshness>,
#[serde(flatten)]
pub extra: serde_json::Map<String, serde_json::Value>,
}
/// Freshness block of the scanner status response.
#[derive(Debug, Clone, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct ScannerFreshness {
/// `fresh` | `stale` | `unknown`.
pub state: String,
}
impl ScannerStatus {
/// Convenience accessor for the freshness state string.
pub fn freshness(&self) -> &str {
self.freshness
.as_ref()
.map(|freshness| freshness.state.as_str())
.unwrap_or("unknown")
}
}
/// Everything that can go wrong in an admin client call.
#[derive(Debug)]
pub enum AdminClientError {
/// The endpoint URL could not be parsed.
InvalidEndpoint(String),
/// Request build/send failed (DNS, connect, timeout, body read).
Transport(reqwest::Error),
/// The server answered a non-2xx status.
HttpStatus { status: u16, body: String },
/// The response body did not decode into the expected shape.
Decode { message: String },
}
impl std::fmt::Display for AdminClientError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::InvalidEndpoint(message) => write!(f, "invalid admin endpoint: {message}"),
Self::Transport(err) => write!(f, "admin request transport failure: {err}"),
Self::HttpStatus { status, body } => write!(f, "admin request failed with HTTP {status}: {body}"),
Self::Decode { message } => write!(f, "admin response decode failure: {message}"),
}
}
}
impl std::error::Error for AdminClientError {}
impl From<reqwest::Error> for AdminClientError {
fn from(err: reqwest::Error) -> Self {
Self::Transport(err)
}
}
/// A signed client for a RustFS admin API.
#[derive(Debug, Clone)]
pub struct AdminClient {
endpoint: reqwest::Url,
access_key: String,
secret_key: String,
session_token: String,
region: String,
api_prefix: String,
http: reqwest::Client,
}
impl AdminClient {
/// Build a client for `endpoint` (e.g. `http://127.0.0.1:9000`) using root
/// or admin credentials. Requests are SigV4-signed with the same scheme
/// the server's admin router authenticates.
pub fn new(endpoint: &str, access_key: &str, secret_key: &str) -> Result<Self, AdminClientError> {
let url = reqwest::Url::parse(endpoint).map_err(|err| AdminClientError::InvalidEndpoint(err.to_string()))?;
if url.host_str().is_none() {
return Err(AdminClientError::InvalidEndpoint("endpoint has no host".to_string()));
}
let http = reqwest::Client::builder()
.connect_timeout(Duration::from_secs(10))
.timeout(Duration::from_secs(30))
.build()
.map_err(AdminClientError::Transport)?;
Ok(Self {
endpoint: url,
access_key: access_key.to_string(),
secret_key: secret_key.to_string(),
session_token: String::new(),
region: DEFAULT_REGION.to_string(),
api_prefix: DEFAULT_ADMIN_API_PREFIX.to_string(),
http,
})
}
/// Attach an STS session token (signed as `x-amz-security-token`).
pub fn with_session_token(mut self, session_token: impl Into<String>) -> Self {
self.session_token = session_token.into();
self
}
/// Override the SigV4 region (defaults to `us-east-1`, matching a
/// region-less RustFS deployment).
pub fn with_region(mut self, region: impl Into<String>) -> Self {
self.region = region.into();
self
}
/// Override the admin API path prefix (defaults to `/rustfs/admin`).
pub fn with_api_prefix(mut self, prefix: impl Into<String>) -> Self {
self.api_prefix = prefix.into();
self
}
/// Start a heal. `bucket` empty and `prefix` empty heals the whole
/// deployment (requires `recursive` or a `pool`/`set` pair in `opts`,
/// enforced server-side); a bucket alone heals the bucket (the server
/// forces `recursive` for bucket heals).
pub async fn heal_start(
&self,
bucket: Option<&str>,
prefix: Option<&str>,
opts: &HealOpts,
force_start: bool,
) -> Result<HealStartSuccess, AdminClientError> {
let body = serde_json::to_vec(opts).map_err(|err| AdminClientError::Decode {
message: err.to_string(),
})?;
let mut query = Vec::new();
if force_start {
query.push(("forceStart", "true".to_string()));
}
self.post_json(&heal_path(bucket, prefix), &query, body).await
}
/// Query the status of the heal identified by `client_token` (the token
/// returned by [`Self::heal_start`]) at the path it was started on.
pub async fn heal_status(
&self,
bucket: Option<&str>,
prefix: Option<&str>,
client_token: &str,
) -> Result<HealTaskStatus, AdminClientError> {
self.post_json(&heal_path(bucket, prefix), &[("clientToken", client_token.to_string())], Vec::new())
.await
}
/// Stop a heal: with a `client_token` only that task is cancelled and its
/// final status returned; without one, every heal task at the path is
/// cancelled (the server answers with a start-success-shaped receipt).
pub async fn heal_stop(
&self,
bucket: Option<&str>,
prefix: Option<&str>,
client_token: Option<&str>,
) -> Result<HealStopOutcome, AdminClientError> {
let mut query = vec![("forceStop", "true".to_string())];
if let Some(token) = client_token {
query.push(("clientToken", token.to_string()));
}
match client_token {
Some(_) => {
let status: HealTaskStatus = self.post_json(&heal_path(bucket, prefix), &query, Vec::new()).await?;
Ok(HealStopOutcome::Stopped(status))
}
None => {
let success: HealStartSuccess = self.post_json(&heal_path(bucket, prefix), &query, Vec::new()).await?;
Ok(HealStopOutcome::PathStopped(success))
}
}
}
/// Cluster-aggregated background heal status.
pub async fn background_heal_status(&self) -> Result<BackgroundHealStatus, AdminClientError> {
self.get_json("/v3/background-heal/status").await
}
/// Data scanner status (enabled state, freshness, runtime config).
pub async fn scanner_status(&self) -> Result<ScannerStatus, AdminClientError> {
self.get_json("/v3/scanner/status").await
}
/// ILM expiry worker status. The payload is owned by the expiry
/// subsystem and still evolving; returned verbatim.
pub async fn ilm_expiry_status(&self) -> Result<serde_json::Value, AdminClientError> {
self.get_json("/v3/ilm/expiry/status").await
}
/// Durable replacement-recovery status (admin v4). The payload is owned
/// by the heal runtime; returned verbatim.
pub async fn replacement_recovery_status(&self) -> Result<serde_json::Value, AdminClientError> {
self.get_json("/v4/heal/replacement-recovery").await
}
/// Signed GET returning a decoded JSON body; escape hatch for endpoints
/// this client does not wrap yet.
pub async fn get_json<T: for<'de> Deserialize<'de>>(&self, path: &str) -> Result<T, AdminClientError> {
let url = self.url_for(path, &[])?;
let request = self.sign_and_build(Method::GET, url, Vec::new(), None).await?;
self.execute(request).await
}
/// Signed POST returning a decoded JSON body.
async fn post_json<T: for<'de> Deserialize<'de>>(
&self,
path: &str,
query: &[(&str, String)],
body: Vec<u8>,
) -> Result<T, AdminClientError> {
let content_type = if body.is_empty() { None } else { Some("application/json") };
let url = self.url_for(path, query)?;
let request = self.sign_and_build(Method::POST, url, body, content_type).await?;
self.execute(request).await
}
fn url_for(&self, path: &str, query: &[(&str, String)]) -> Result<reqwest::Url, AdminClientError> {
let mut url = self
.endpoint
.join(&format!("{}{}", self.api_prefix.trim_end_matches('/'), path))
.map_err(|err| AdminClientError::InvalidEndpoint(err.to_string()))?;
if !query.is_empty() {
let mut pairs = url.query_pairs_mut();
for (key, value) in query {
pairs.append_pair(key, value);
}
}
Ok(url)
}
/// Build a SigV4-signed request via the same signer the server trusts,
/// then hand the signed headers to the HTTP client. The signature covers
/// method, path, query, and an unsigned-payload marker — the same shape
/// RustFS itself sends for peer admin calls.
async fn sign_and_build(
&self,
method: Method,
url: reqwest::Url,
body: Vec<u8>,
content_type: Option<&str>,
) -> Result<reqwest::Request, AdminClientError> {
let authority = match (url.host_str(), url.port_or_known_default()) {
(Some(host), Some(port)) => format!("{host}:{port}"),
_ => return Err(AdminClientError::InvalidEndpoint("endpoint has no authority".to_string())),
};
let mut builder = http::Request::builder()
.method(method.clone())
.uri(url.as_str())
.header(http::header::HOST, &authority)
.header("x-amz-content-sha256", rustfs_signer::constants::UNSIGNED_PAYLOAD);
if let Some(content_type) = content_type {
builder = builder.header(http::header::CONTENT_TYPE, content_type);
}
let unsigned = builder
.body(s3s::Body::empty())
.map_err(|err| AdminClientError::InvalidEndpoint(format!("build request failed: {err}")))?;
let signed = rustfs_signer::sign_v4(
unsigned,
body.len() as i64,
&self.access_key,
&self.secret_key,
&self.session_token,
&self.region,
);
let mut request = self
.http
.request(method, url)
.body(body)
.build()
.map_err(AdminClientError::Transport)?;
let headers = request.headers_mut();
for (name, value) in signed.headers().iter() {
// HOST is owned by the HTTP client; the signed value above was
// built from the same URL authority, so they always agree.
if name == http::header::HOST {
continue;
}
headers.insert(name, value.clone());
}
Ok(request)
}
async fn execute<T: for<'de> Deserialize<'de>>(&self, request: reqwest::Request) -> Result<T, AdminClientError> {
let response = self.http.execute(request).await?;
let status = response.status();
let bytes = response.bytes().await?;
if !status.is_success() {
return Err(AdminClientError::HttpStatus {
status: status.as_u16(),
body: String::from_utf8_lossy(&bytes).into_owned(),
});
}
serde_json::from_slice(&bytes).map_err(|err| AdminClientError::Decode {
message: err.to_string(),
})
}
}
/// Response of [`AdminClient::heal_stop`]: cancelling a single tokened task
/// answers with that task's status, cancelling a whole path answers with a
/// start-success-shaped receipt.
#[derive(Debug, Clone)]
pub enum HealStopOutcome {
Stopped(HealTaskStatus),
PathStopped(HealStartSuccess),
}
fn heal_path(bucket: Option<&str>, prefix: Option<&str>) -> String {
match (bucket, prefix) {
(Some(bucket), Some(prefix)) if !bucket.is_empty() && !prefix.is_empty() => {
format!("/v3/heal/{}/{}", percent_encode_path_segment(bucket), percent_encode_path_segment(prefix))
}
(Some(bucket), Some(_)) | (Some(bucket), None) if !bucket.is_empty() => {
format!("/v3/heal/{}", percent_encode_path_segment(bucket))
}
_ => "/v3/heal/".to_string(),
}
}
/// Encode a single path segment (slashes are content, not separators, inside
/// bucket/prefix path params).
fn percent_encode_path_segment(segment: &str) -> String {
let mut out = String::with_capacity(segment.len());
for byte in segment.bytes() {
match byte {
b'A'..=b'Z' | b'a'..=b'z' | b'0'..=b'9' | b'-' | b'_' | b'.' | b'~' => out.push(byte as char),
_ => out.push_str(&format!("%{byte:02X}")),
}
}
out
}
#[cfg(test)]
mod tests {
use super::{
AdminClient, AdminClientError, BackgroundHealStatus, HealOpts, HealScanMode, HealStartSuccess, HealTaskStatus,
ScannerStatus, heal_path, percent_encode_path_segment,
};
use serde_json::json;
use std::sync::{Arc, Mutex};
#[test]
fn heal_paths_cover_root_bucket_and_prefix() {
assert_eq!(heal_path(None, None), "/v3/heal/");
assert_eq!(heal_path(Some(""), Some("")), "/v3/heal/");
assert_eq!(heal_path(Some("bucket"), None), "/v3/heal/bucket");
assert_eq!(heal_path(Some("bucket"), Some("pre/fix")), "/v3/heal/bucket/pre%2Ffix");
}
#[test]
fn path_segments_percent_encode_reserved_characters() {
assert_eq!(percent_encode_path_segment("a b"), "a%20b");
assert_eq!(percent_encode_path_segment("a/b"), "a%2Fb");
assert_eq!(percent_encode_path_segment("ü"), "%C3%BC");
}
#[test]
fn heal_opts_round_trip_through_the_server_wire_shape() {
let opts = HealOpts {
recursive: true,
dry_run: false,
remove: true,
recreate: false,
scan_mode: HealScanMode::Deep,
update_parity: true,
no_lock: false,
pool: Some(1),
set: Some(2),
};
let wire = serde_json::to_value(&opts).unwrap();
assert_eq!(wire["scanMode"], json!(2), "the server body decodes scanMode as a number");
let back: HealOpts = serde_json::from_value(wire).unwrap();
assert_eq!(back.scan_mode, HealScanMode::Deep);
assert_eq!(back.pool, Some(1));
}
#[test]
fn heal_scan_mode_accepts_both_wire_encodings() {
assert_eq!(serde_json::from_value::<HealScanMode>(json!(1)).unwrap(), HealScanMode::Normal);
assert_eq!(serde_json::from_value::<HealScanMode>(json!("deep")).unwrap(), HealScanMode::Deep);
assert!(serde_json::from_value::<HealScanMode>(json!(9)).is_err());
assert!(serde_json::from_value::<HealScanMode>(json!("sideways")).is_err());
}
#[test]
fn heal_task_status_decodes_the_server_response_shape() {
let raw = json!({
"summary": "finished",
"detail": "",
"startTime": "2026-08-17T00:00:00Z",
"settings": {"recursive": false, "scanMode": 1},
"items": [{
"resultId": 1, "type": "object", "bucket": "b", "object": "o", "versionId": "", "detail": "",
"parityBlocks": 2, "dataBlocks": 2, "diskCount": 4, "setCount": 1,
"before": {"drives": []}, "after": {"drives": []}, "objectSize": 128
}],
"truncated": false
});
let status: HealTaskStatus = serde_json::from_value(raw).unwrap();
assert_eq!(status.summary, "finished");
assert_eq!(status.items.len(), 1);
assert_eq!(status.settings.scan_mode, HealScanMode::Normal);
assert!(status.progress.is_none());
}
#[test]
fn background_heal_status_types_known_fields_and_passes_the_rest_through() {
let raw = json!({
"state": "active",
"bitrotStartTime": "t",
"healQueueLength": 3,
"healActiveTasks": 1,
"healOperations": {"queueLength": 3},
"clusterStatusComplete": true
});
let status: BackgroundHealStatus = serde_json::from_value(raw).unwrap();
assert_eq!(status.state, "active");
assert_eq!(status.heal_queue_length, 3);
assert!(status.cluster_status_complete);
assert!(status.extra.contains_key("healOperations"), "unknown nested payloads must pass through");
}
#[test]
fn scanner_status_defaults_freshness_to_unknown() {
let raw = json!({"enabled": true, "freshness": {"state": "stale"}, "metrics": {}});
let status: ScannerStatus = serde_json::from_value(raw).unwrap();
assert_eq!(status.freshness(), "stale");
let bare: ScannerStatus = serde_json::from_value(json!({"enabled": false})).unwrap();
assert_eq!(bare.freshness(), "unknown");
}
#[test]
fn invalid_endpoint_is_rejected_without_io() {
let err = AdminClient::new("not a url", "ak", "sk").unwrap_err();
assert!(matches!(err, AdminClientError::InvalidEndpoint(_)));
}
#[tokio::test]
async fn signed_requests_carry_sigv4_authorization_and_correct_target() {
let server = TestServer::spawn(r#"{"clientToken":"token-1","clientAddress":"127.0.0.1:9","startTime":"t"}"#, 200).await;
let client = AdminClient::new(&format!("http://{}", server.addr), "minioadmin", "minioadmin")
.expect("client builds against the test server");
let start: HealStartSuccess = client
.heal_start(
Some("bucket"),
None,
&HealOpts {
recursive: true,
..Default::default()
},
false,
)
.await
.expect("signed heal start decodes");
assert_eq!(start.client_token, "token-1");
let request = server.recorded();
assert_eq!(request.method, "POST");
assert_eq!(request.path, "/rustfs/admin/v3/heal/bucket");
assert!(!request.query.contains("forceStart"), "absent flags must not be sent");
let auth = request.header("authorization").expect("request must be signed");
assert!(auth.starts_with("AWS4-HMAC-SHA256"), "SigV4 scheme, got: {auth}");
assert!(auth.contains("Credential=minioadmin/"), "credentials must be in the Authorization header");
assert_eq!(
request.header("x-amz-content-sha256").as_deref(),
Some("UNSIGNED-PAYLOAD"),
"the client signs the same payload marker RustFS peer calls use"
);
assert_eq!(request.header("content-type").as_deref(), Some("application/json"));
assert!(request.body.contains("\"recursive\":true"));
}
#[tokio::test]
async fn query_sends_client_token_on_the_same_path() {
let body = r#"{"summary":"running","detail":"","settings":{"recursive":false},"items":[],"truncated":false}"#;
let server = TestServer::spawn(body, 200).await;
let client = AdminClient::new(&format!("http://{}", server.addr), "ak", "sk").unwrap();
let status = client
.heal_status(Some("bucket"), None, "token-1")
.await
.expect("status decodes");
assert_eq!(status.summary, "running");
let request = server.recorded();
assert_eq!(request.path, "/rustfs/admin/v3/heal/bucket");
assert!(request.query.contains("clientToken=token-1"));
assert!(!request.query.contains("forceStop"));
}
#[tokio::test]
async fn stop_without_token_takes_the_path_cancel_branch() {
let server = TestServer::spawn(r#"{"clientToken":"path","clientAddress":"c","startTime":"t"}"#, 200).await;
let client = AdminClient::new(&format!("http://{}", server.addr), "ak", "sk").unwrap();
let outcome = client.heal_stop(Some("bucket"), None, None).await.expect("path stop decodes");
assert!(matches!(outcome, super::HealStopOutcome::PathStopped(_)));
let request = server.recorded();
assert!(request.query.contains("forceStop=true"));
assert!(!request.query.contains("clientToken"));
}
#[tokio::test]
async fn http_error_status_maps_to_a_typed_error_with_body() {
let server = TestServer::spawn(r#"{"code":"AccessDenied","message":"denied"}"#, 403).await;
let client = AdminClient::new(&format!("http://{}", server.addr), "ak", "sk").unwrap();
let err = client.scanner_status().await.unwrap_err();
match err {
AdminClientError::HttpStatus { status, body } => {
assert_eq!(status, 403);
assert!(body.contains("AccessDenied"));
}
other => panic!("expected HttpStatus, got {other:?}"),
}
}
#[tokio::test]
async fn malformed_success_body_maps_to_a_decode_error() {
let server = TestServer::spawn("not json", 200).await;
let client = AdminClient::new(&format!("http://{}", server.addr), "ak", "sk").unwrap();
assert!(matches!(client.scanner_status().await.unwrap_err(), AdminClientError::Decode { .. }));
}
/// One recorded request, parsed off the wire with the minimum needed for
/// assertions: method, path, query, headers, body.
#[derive(Debug, Clone)]
struct RecordedRequest {
method: String,
path: String,
query: String,
headers: Vec<(String, String)>,
body: String,
}
impl RecordedRequest {
fn header(&self, name: &str) -> Option<String> {
self.headers
.iter()
.find(|(key, _)| key.eq_ignore_ascii_case(name))
.map(|(_, value)| value.clone())
}
}
/// Minimal HTTP/1.1 server: one canned response per connection, every
/// request recorded behind an `Arc<Mutex>`. Deliberately dependency-free —
/// the assertions only need the raw request bytes.
struct TestServer {
addr: std::net::SocketAddr,
requests: Arc<Mutex<Vec<RecordedRequest>>>,
}
impl TestServer {
async fn spawn(response_body: &'static str, status: u16) -> Self {
use tokio::io::{AsyncReadExt, AsyncWriteExt};
let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
.await
.expect("bind ephemeral port");
let addr = listener.local_addr().expect("local addr");
let requests: Arc<Mutex<Vec<RecordedRequest>>> = Arc::new(Mutex::new(Vec::new()));
let recorded = requests.clone();
tokio::spawn(async move {
let reason = if status == 200 { "OK" } else { "Forbidden" };
let response = format!(
"HTTP/1.1 {status} {reason}\r\ncontent-type: application/json\r\ncontent-length: {}\r\nconnection: close\r\n\r\n{response_body}",
response_body.len()
);
// Each request is a fresh connection (connection: close); a
// bounded loop serves every call a test makes while letting
// the task exit instead of lingering for the whole process.
for _ in 0..16 {
let Ok((mut stream, _)) = listener.accept().await else {
break;
};
let mut buffer = Vec::with_capacity(2048);
let mut chunk = [0u8; 2048];
// Read headers plus content-length body, or stop on close.
loop {
if let Some(end) = find_header_end(&buffer) {
let content_length = extract_content_length(&buffer[..end]);
if buffer.len() >= end + content_length {
break;
}
}
let n = match stream.read(&mut chunk).await {
Ok(0) | Err(_) => break,
Ok(n) => n,
};
buffer.extend_from_slice(&chunk[..n]);
if buffer.len() > 64 * 1024 {
break;
}
}
if let Some(request) = parse_request(&buffer) {
recorded.lock().expect("recorded lock").push(request);
}
let _ = stream.write_all(response.as_bytes()).await;
let _ = stream.shutdown().await;
}
});
Self { addr, requests }
}
fn recorded(&self) -> RecordedRequest {
self.requests
.lock()
.expect("recorded lock")
.last()
.cloned()
.expect("the client call must have produced one recorded request")
}
}
fn find_header_end(buffer: &[u8]) -> Option<usize> {
buffer.windows(4).position(|window| window == b"\r\n\r\n").map(|pos| pos + 4)
}
fn extract_content_length(headers: &[u8]) -> usize {
let text = String::from_utf8_lossy(headers).to_ascii_lowercase();
text.lines()
.find_map(|line| line.strip_prefix("content-length:"))
.and_then(|value| value.trim().parse().ok())
.unwrap_or(0)
}
fn parse_request(raw: &[u8]) -> Option<RecordedRequest> {
let end = find_header_end(raw)?;
let head = String::from_utf8_lossy(&raw[..end]);
let body = String::from_utf8_lossy(&raw[end..]).into_owned();
let mut lines = head.lines();
let request_line = lines.next()?;
let mut parts = request_line.split_whitespace();
let method = parts.next()?.to_string();
let target = parts.next()?.to_string();
let (path, query) = match target.split_once('?') {
Some((path, query)) => (path.to_string(), query.to_string()),
None => (target, String::new()),
};
let headers = lines
.filter_map(|line| line.split_once(':'))
.map(|(name, value)| (name.trim().to_string(), value.trim().to_string()))
.collect();
Some(RecordedRequest {
method,
path,
query,
headers,
body,
})
}
}
-2
View File
@@ -12,7 +12,6 @@
// See the License for the specific language governing permissions and
// limitations under the License.
pub mod client;
pub mod group;
pub mod heal_commands;
pub mod health;
@@ -26,7 +25,6 @@ pub mod trace;
pub mod user;
pub mod utils;
pub use client::*;
pub use group::*;
pub use info_commands::*;
pub use policy::*;
+258 -241
View File
@@ -66,18 +66,20 @@ use rustfs_config::{
};
use rustfs_iam::error::is_err_no_such_service_account;
use rustfs_iam::federation::OIDC_VIRTUAL_PARENT_CLAIM;
use rustfs_iam::store::object::ObjectStore;
use rustfs_iam::store::{MappedPolicy, UserType, sr_wire_user_type, user_type_from_sr_wire};
use rustfs_iam::sys::{
NewServiceAccountOpts, SITE_REPLICATOR_SERVICE_ACCOUNT, UpdateServiceAccountOpts, get_claims_from_token_with_secret,
IamSys, NewServiceAccountOpts, SITE_REPLICATOR_SERVICE_ACCOUNT, UpdateServiceAccountOpts, get_claims_from_token_with_secret,
};
use rustfs_madmin::{
AddOrUpdateUserReq, BucketBandwidth, GroupAddRemove, GroupStatus, IDPSettings, InProgressMetric, InQueueMetric,
LDAPConfigSettings, LDAPSettings, OpenIDProviderSettings, PeerInfo, PeerSite, QStat, ReplProxyMetric, ReplicateAddStatus,
ReplicateEditStatus, ReplicateRemoveStatus, ResyncBucketStatus, SITE_REPL_API_VERSION, SR_IAM_ITEM_STS_ACC,
SR_IAM_ITEM_STS_ACC_LEGACY, SRBucketInfo, SRBucketMeta, SRBucketStatsSummary, SRGroupInfo, SRGroupStatsSummary, SRIAMItem,
SRIAMPolicy, SRILMExpiryStatsSummary, SRInfo, SRMetric, SRMetricsSummary, SRPeerError, SRPeerJoinReq, SRPendingOperation,
SRPolicyMapping, SRPolicyStatsSummary, SRRemoveReq, SRResyncOpStatus, SRRetryStats, SRSessionPolicy, SRSiteSummary,
SRStateEditReq, SRStateInfo, SRStatusInfo, SRSvcAccCreate, SRUserStatsSummary, SiteReplicationInfo, SyncStatus, WorkerStat,
SRIAMPolicy, SRIAMUser, SRILMExpiryStatsSummary, SRInfo, SRMetric, SRMetricsSummary, SRPeerError, SRPeerJoinReq,
SRPendingOperation, SRPolicyMapping, SRPolicyStatsSummary, SRRemoveReq, SRResyncOpStatus, SRRetryStats, SRSTSCredential,
SRSessionPolicy, SRSiteSummary, SRStateEditReq, SRStateInfo, SRStatusInfo, SRSvcAccChange, SRSvcAccCreate,
SRUserStatsSummary, SiteReplicationInfo, SyncStatus, WorkerStat,
};
use rustfs_policy::policy::{
Policy,
@@ -9278,247 +9280,16 @@ async fn apply_iam_item(item: SRIAMItem) -> S3Result<()> {
let incoming_updated_at = item.updated_at;
match item.r#type.as_str() {
"policy" => {
if let Some(policy) = item.policy {
let policy: Policy =
serde_json::from_value(policy).map_err(|e| s3_error!(InvalidRequest, "invalid policy body: {}", e))?;
iam_sys.set_policy(&item.name, policy).await.map_err(ApiError::from)?;
} else {
iam_sys.delete_policy(&item.name, true).await.map_err(ApiError::from)?;
}
Ok(())
}
"policy-mapping" => {
let Some(mapping) = item.policy_mapping else {
return Err(s3_error!(InvalidRequest, "policyMapping is required"));
};
let user_type =
user_type_from_sr_wire(mapping.user_type).ok_or_else(|| s3_error!(InvalidRequest, "invalid userType"))?;
iam_sys
.policy_db_set(&mapping.user_or_group, user_type, mapping.is_group, &mapping.policy)
.await
.map_err(ApiError::from)?;
Ok(())
}
"group-info" => {
let Some(group_info) = item.group_info else {
return Err(s3_error!(InvalidRequest, "groupInfo is required"));
};
let update = group_info.update_req;
if !group_info_requires_upsert(&update) {
iam_sys
.remove_users_from_group(&update.group, update.members)
.await
.map_err(ApiError::from)?;
return Ok(());
}
iam_sys
.add_users_to_group(&update.group, update.members)
.await
.map_err(ApiError::from)?;
iam_sys
.set_group_status(&update.group, matches!(update.status, GroupStatus::Enabled))
.await
.map_err(ApiError::from)?;
Ok(())
}
"policy" => apply_iam_policy_item(&iam_sys, &item.name, item.policy).await,
"policy-mapping" => apply_iam_policy_mapping_item(&iam_sys, item.policy_mapping).await,
"group-info" => apply_iam_group_info_item(&iam_sys, item.group_info).await,
// MinIO madmin-go sends `SRIAMItemSTSAcc = "sts-account"`. The legacy alias
// `sts-credential` (emitted by older RustFS releases) stays accepted permanently
// so mixed-version RustFS sites keep replicating STS credentials during rolling
// upgrades; it is a compatibility layer, not temporary code.
SR_IAM_ITEM_STS_ACC | SR_IAM_ITEM_STS_ACC_LEGACY => {
let Some(sts_credential) = item.sts_credential else {
return Err(s3_error!(InvalidRequest, "stsCredential is required"));
};
let Some(secret) = current_token_signing_key() else {
return Err(s3_error!(InvalidRequest, "token signing key not initialized"));
};
let claims = get_claims_from_token_with_secret(&sts_credential.session_token, &secret)
.map_err(|e| s3_error!(InvalidRequest, "invalid STS session token: {e}"))?;
let expiration = claims
.get("exp")
.and_then(claims_unix_timestamp)
.map(OffsetDateTime::from_unix_timestamp)
.transpose()
.map_err(|e| s3_error!(InvalidRequest, "invalid STS expiry: {e}"))?;
let groups = string_list_claim(&claims, "groups");
let compatibility_policy = sts_replication_compatibility_policy(&claims, &sts_credential.parent_policy_mapping);
let cred = rustfs_credentials::Credentials {
access_key: sts_credential.access_key.clone(),
secret_key: sts_credential.secret_key.clone(),
session_token: sts_credential.session_token.clone(),
expiration,
status: "on".to_string(),
parent_user: sts_credential.parent_user.clone(),
groups,
claims: Some(claims),
..Default::default()
};
iam_sys
.set_temp_user(&sts_credential.access_key, &cred, compatibility_policy)
.await
.map_err(ApiError::from)?;
Ok(())
}
"iam-user" => {
let Some(user) = item.iam_user else {
return Err(s3_error!(InvalidRequest, "iamUser is required"));
};
if let Some(local) = iam_sys.get_user(&user.access_key).await
&& is_stale_update(local.update_at.unwrap_or(OffsetDateTime::UNIX_EPOCH), incoming_updated_at)
{
return Ok(());
}
if user.is_delete_req {
iam_sys.delete_user(&user.access_key, true).await.map_err(ApiError::from)?;
} else {
let Some(user_req) = user.user_req else {
return Err(s3_error!(InvalidRequest, "userReq is required"));
};
let is_status_only_update = user_req.secret_key.is_empty() && user_req.policy.is_none();
if is_status_only_update {
iam_sys
.set_user_status(&user.access_key, user_req.status)
.await
.map_err(ApiError::from)?;
} else {
iam_sys
.create_user(&user.access_key, &user_req)
.await
.map_err(ApiError::from)?;
}
}
Ok(())
}
"service-account" => {
let Some(change) = item.svc_acc_change else {
return Err(s3_error!(InvalidRequest, "serviceAccountChange is required"));
};
let envelope = change.oidc_service_account_envelope;
if let Some(create) = change.create {
let local_updated_at = iam_sys
.get_user(&create.access_key)
.await
.map(|local| local.update_at.unwrap_or(OffsetDateTime::UNIX_EPOCH));
let replicated_policy = if create.access_key == SITE_REPLICATOR_SERVICE_ACCOUNT {
if local_updated_at.is_some_and(|local_updated_at| is_stale_update(local_updated_at, incoming_updated_at)) {
return Ok(());
}
ReplicatedServiceAccountPolicy {
policy: Some(site_replicator_service_account_policy()?),
is_envelope: false,
}
} else {
let Some(replicated_policy) = decode_service_account_replication_policy(
&create,
envelope.as_ref(),
incoming_updated_at,
local_updated_at,
)?
else {
return Ok(());
};
replicated_policy
};
match iam_sys.get_service_account(&create.access_key).await {
Ok((existing, _)) => {
if existing.parent_user != create.parent {
return Err(s3_error!(
InvalidRequest,
"service account {} already exists with a different parent user",
create.access_key
));
}
iam_sys
.update_service_account(
&create.access_key,
UpdateServiceAccountOpts {
name: replicated_policy.metadata_for_existing_account(create.name),
description: replicated_policy.metadata_for_existing_account(create.description),
session_policy: replicated_policy.for_existing_account(),
secret_key: Some(create.secret_key),
expiration: create.expiration,
status: (!create.status.is_empty()).then_some(create.status),
parent_user: None,
allow_site_replicator_account: create.access_key == SITE_REPLICATOR_SERVICE_ACCOUNT,
},
)
.await
.map_err(ApiError::from)?;
}
Err(err) if is_err_no_such_service_account(&err) => {
iam_sys
.new_service_account(
&create.parent,
Some(create.groups),
NewServiceAccountOpts {
session_policy: replicated_policy.policy,
access_key: create.access_key,
secret_key: create.secret_key,
name: (!create.name.is_empty()).then_some(create.name),
description: (!create.description.is_empty()).then_some(create.description),
expiration: create.expiration,
allow_site_replicator_account: true,
claims: Some(create.claims),
},
)
.await
.map_err(ApiError::from)?;
}
Err(err) => return Err(ApiError::from(err).into()),
}
return Ok(());
}
if let Some(update) = change.update {
if let Some(local) = iam_sys.get_user(&update.access_key).await
&& is_stale_update(local.update_at.unwrap_or(OffsetDateTime::UNIX_EPOCH), incoming_updated_at)
{
return Ok(());
}
let allow_site_replicator_account = update.access_key == SITE_REPLICATOR_SERVICE_ACCOUNT;
let session_policy = if allow_site_replicator_account {
Some(site_replicator_service_account_policy()?)
} else {
update.session_policy.as_str().and_then(|raw| serde_json::from_str(raw).ok())
};
iam_sys
.update_service_account(
&update.access_key,
UpdateServiceAccountOpts {
session_policy,
secret_key: (!update.secret_key.is_empty()).then_some(update.secret_key),
name: (!update.name.is_empty()).then_some(update.name),
description: (!update.description.is_empty()).then_some(update.description),
expiration: update.expiration,
status: (!update.status.is_empty()).then_some(update.status),
// Peers replicate credentials, never the local parent binding:
// each site resolves its own parent from its own IAM.
parent_user: None,
allow_site_replicator_account,
},
)
.await
.map_err(ApiError::from)?;
return Ok(());
}
if let Some(delete) = change.delete {
if let Some(local) = iam_sys.get_user(&delete.access_key).await
&& is_stale_update(local.update_at.unwrap_or(OffsetDateTime::UNIX_EPOCH), incoming_updated_at)
{
return Ok(());
}
iam_sys
.delete_service_account(&delete.access_key, true)
.await
.map_err(ApiError::from)?;
return Ok(());
}
Err(s3_error!(InvalidRequest, "serviceAccountChange is empty"))
}
SR_IAM_ITEM_STS_ACC | SR_IAM_ITEM_STS_ACC_LEGACY => apply_iam_sts_account_item(&iam_sys, item.sts_credential).await,
"iam-user" => apply_iam_user_item(&iam_sys, item.iam_user, incoming_updated_at).await,
"service-account" => apply_iam_service_account_item(&iam_sys, item.svc_acc_change, incoming_updated_at).await,
_ => Err(s3_error!(
NotImplemented,
"site replication IAM item type `{}` is not supported",
@@ -9527,6 +9298,252 @@ async fn apply_iam_item(item: SRIAMItem) -> S3Result<()> {
}
}
async fn apply_iam_policy_item(iam_sys: &IamSys<ObjectStore>, name: &str, policy: Option<Value>) -> S3Result<()> {
if let Some(policy) = policy {
let policy: Policy =
serde_json::from_value(policy).map_err(|e| s3_error!(InvalidRequest, "invalid policy body: {}", e))?;
iam_sys.set_policy(name, policy).await.map_err(ApiError::from)?;
} else {
iam_sys.delete_policy(name, true).await.map_err(ApiError::from)?;
}
Ok(())
}
async fn apply_iam_policy_mapping_item(iam_sys: &IamSys<ObjectStore>, policy_mapping: Option<SRPolicyMapping>) -> S3Result<()> {
let Some(mapping) = policy_mapping else {
return Err(s3_error!(InvalidRequest, "policyMapping is required"));
};
let user_type = user_type_from_sr_wire(mapping.user_type).ok_or_else(|| s3_error!(InvalidRequest, "invalid userType"))?;
iam_sys
.policy_db_set(&mapping.user_or_group, user_type, mapping.is_group, &mapping.policy)
.await
.map_err(ApiError::from)?;
Ok(())
}
async fn apply_iam_group_info_item(iam_sys: &IamSys<ObjectStore>, group_info: Option<SRGroupInfo>) -> S3Result<()> {
let Some(group_info) = group_info else {
return Err(s3_error!(InvalidRequest, "groupInfo is required"));
};
let update = group_info.update_req;
if !group_info_requires_upsert(&update) {
iam_sys
.remove_users_from_group(&update.group, update.members)
.await
.map_err(ApiError::from)?;
return Ok(());
}
iam_sys
.add_users_to_group(&update.group, update.members)
.await
.map_err(ApiError::from)?;
iam_sys
.set_group_status(&update.group, matches!(update.status, GroupStatus::Enabled))
.await
.map_err(ApiError::from)?;
Ok(())
}
async fn apply_iam_sts_account_item(iam_sys: &IamSys<ObjectStore>, sts_credential: Option<SRSTSCredential>) -> S3Result<()> {
let Some(sts_credential) = sts_credential else {
return Err(s3_error!(InvalidRequest, "stsCredential is required"));
};
let Some(secret) = current_token_signing_key() else {
return Err(s3_error!(InvalidRequest, "token signing key not initialized"));
};
let claims = get_claims_from_token_with_secret(&sts_credential.session_token, &secret)
.map_err(|e| s3_error!(InvalidRequest, "invalid STS session token: {e}"))?;
let expiration = claims
.get("exp")
.and_then(claims_unix_timestamp)
.map(OffsetDateTime::from_unix_timestamp)
.transpose()
.map_err(|e| s3_error!(InvalidRequest, "invalid STS expiry: {e}"))?;
let groups = string_list_claim(&claims, "groups");
let compatibility_policy = sts_replication_compatibility_policy(&claims, &sts_credential.parent_policy_mapping);
let cred = rustfs_credentials::Credentials {
access_key: sts_credential.access_key.clone(),
secret_key: sts_credential.secret_key.clone(),
session_token: sts_credential.session_token.clone(),
expiration,
status: "on".to_string(),
parent_user: sts_credential.parent_user.clone(),
groups,
claims: Some(claims),
..Default::default()
};
iam_sys
.set_temp_user(&sts_credential.access_key, &cred, compatibility_policy)
.await
.map_err(ApiError::from)?;
Ok(())
}
async fn apply_iam_user_item(
iam_sys: &IamSys<ObjectStore>,
iam_user: Option<SRIAMUser>,
incoming_updated_at: Option<OffsetDateTime>,
) -> S3Result<()> {
let Some(user) = iam_user else {
return Err(s3_error!(InvalidRequest, "iamUser is required"));
};
if let Some(local) = iam_sys.get_user(&user.access_key).await
&& is_stale_update(local.update_at.unwrap_or(OffsetDateTime::UNIX_EPOCH), incoming_updated_at)
{
return Ok(());
}
if user.is_delete_req {
iam_sys.delete_user(&user.access_key, true).await.map_err(ApiError::from)?;
} else {
let Some(user_req) = user.user_req else {
return Err(s3_error!(InvalidRequest, "userReq is required"));
};
let is_status_only_update = user_req.secret_key.is_empty() && user_req.policy.is_none();
if is_status_only_update {
iam_sys
.set_user_status(&user.access_key, user_req.status)
.await
.map_err(ApiError::from)?;
} else {
iam_sys
.create_user(&user.access_key, &user_req)
.await
.map_err(ApiError::from)?;
}
}
Ok(())
}
async fn apply_iam_service_account_item(
iam_sys: &IamSys<ObjectStore>,
svc_acc_change: Option<SRSvcAccChange>,
incoming_updated_at: Option<OffsetDateTime>,
) -> S3Result<()> {
let Some(change) = svc_acc_change else {
return Err(s3_error!(InvalidRequest, "serviceAccountChange is required"));
};
let envelope = change.oidc_service_account_envelope;
if let Some(create) = change.create {
let local_updated_at = iam_sys
.get_user(&create.access_key)
.await
.map(|local| local.update_at.unwrap_or(OffsetDateTime::UNIX_EPOCH));
let replicated_policy = if create.access_key == SITE_REPLICATOR_SERVICE_ACCOUNT {
if local_updated_at.is_some_and(|local_updated_at| is_stale_update(local_updated_at, incoming_updated_at)) {
return Ok(());
}
ReplicatedServiceAccountPolicy {
policy: Some(site_replicator_service_account_policy()?),
is_envelope: false,
}
} else {
let Some(replicated_policy) =
decode_service_account_replication_policy(&create, envelope.as_ref(), incoming_updated_at, local_updated_at)?
else {
return Ok(());
};
replicated_policy
};
match iam_sys.get_service_account(&create.access_key).await {
Ok((existing, _)) => {
if existing.parent_user != create.parent {
return Err(s3_error!(
InvalidRequest,
"service account {} already exists with a different parent user",
create.access_key
));
}
iam_sys
.update_service_account(
&create.access_key,
UpdateServiceAccountOpts {
name: replicated_policy.metadata_for_existing_account(create.name),
description: replicated_policy.metadata_for_existing_account(create.description),
session_policy: replicated_policy.for_existing_account(),
secret_key: Some(create.secret_key),
expiration: create.expiration,
status: (!create.status.is_empty()).then_some(create.status),
parent_user: None,
allow_site_replicator_account: create.access_key == SITE_REPLICATOR_SERVICE_ACCOUNT,
},
)
.await
.map_err(ApiError::from)?;
}
Err(err) if is_err_no_such_service_account(&err) => {
iam_sys
.new_service_account(
&create.parent,
Some(create.groups),
NewServiceAccountOpts {
session_policy: replicated_policy.policy,
access_key: create.access_key,
secret_key: create.secret_key,
name: (!create.name.is_empty()).then_some(create.name),
description: (!create.description.is_empty()).then_some(create.description),
expiration: create.expiration,
allow_site_replicator_account: true,
claims: Some(create.claims),
},
)
.await
.map_err(ApiError::from)?;
}
Err(err) => return Err(ApiError::from(err).into()),
}
return Ok(());
}
if let Some(update) = change.update {
if let Some(local) = iam_sys.get_user(&update.access_key).await
&& is_stale_update(local.update_at.unwrap_or(OffsetDateTime::UNIX_EPOCH), incoming_updated_at)
{
return Ok(());
}
let allow_site_replicator_account = update.access_key == SITE_REPLICATOR_SERVICE_ACCOUNT;
let session_policy = if allow_site_replicator_account {
Some(site_replicator_service_account_policy()?)
} else {
update.session_policy.as_str().and_then(|raw| serde_json::from_str(raw).ok())
};
iam_sys
.update_service_account(
&update.access_key,
UpdateServiceAccountOpts {
session_policy,
secret_key: (!update.secret_key.is_empty()).then_some(update.secret_key),
name: (!update.name.is_empty()).then_some(update.name),
description: (!update.description.is_empty()).then_some(update.description),
expiration: update.expiration,
status: (!update.status.is_empty()).then_some(update.status),
// Peers replicate credentials, never the local parent binding:
// each site resolves its own parent from its own IAM.
parent_user: None,
allow_site_replicator_account,
},
)
.await
.map_err(ApiError::from)?;
return Ok(());
}
if let Some(delete) = change.delete {
if let Some(local) = iam_sys.get_user(&delete.access_key).await
&& is_stale_update(local.update_at.unwrap_or(OffsetDateTime::UNIX_EPOCH), incoming_updated_at)
{
return Ok(());
}
iam_sys
.delete_service_account(&delete.access_key, true)
.await
.map_err(ApiError::from)?;
return Ok(());
}
Err(s3_error!(InvalidRequest, "serviceAccountChange is empty"))
}
fn claims_unix_timestamp(value: &Value) -> Option<i64> {
match value {
Value::Number(number) => number.as_i64(),
-29
View File
@@ -417,13 +417,6 @@ struct SystemAdminDiscovery {
struct ServerInfoResponse {
info: InfoMessage,
admin_discovery: SystemAdminDiscovery,
/// Startup bitrot algorithm self-test outcome (rustfs/backlog#1873):
/// `passed` (algorithms verified at boot), `failed` (a drifted hash
/// implementation — the process is serving with degraded integrity
/// checking unless `RUSTFS_BITROT_SELFTEST_STRICT` aborted it), or
/// `unknown` (not yet run or disabled).
#[serde(rename = "bitrotSelftest")]
bitrot_selftest: &'static str,
}
#[derive(Serialize)]
@@ -440,14 +433,6 @@ fn system_admin_discovery(usecase: &DefaultAdminUsecase) -> SystemAdminDiscovery
}
}
fn bitrot_selftest_status_str() -> &'static str {
match crate::bitrot_selftest::bitrot_selftest_passed() {
Some(true) => "passed",
Some(false) => "failed",
None => "unknown",
}
}
#[async_trait::async_trait]
impl Operation for ServerInfoHandler {
async fn call(&self, req: S3Request<Body>, _params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
@@ -479,7 +464,6 @@ impl Operation for ServerInfoHandler {
let response = ServerInfoResponse {
info,
admin_discovery: system_admin_discovery(&usecase),
bitrot_selftest: bitrot_selftest_status_str(),
};
let data = serde_json::to_vec(&response).map_err(|e| {
@@ -1551,18 +1535,6 @@ mod tests {
);
}
/// The startup bitrot self-test outcome must surface in server info as one
/// of three closed-set strings, never an internal enum or a null
/// (rustfs/backlog#1873). This test pins the string mapping; whether the
/// process-global cell holds Some(true)/Some(false)/None is owned by
/// `crate::bitrot_selftest`'s own tests.
#[test]
fn bitrot_selftest_status_str_is_a_closed_set_of_operators_strings() {
let rendered = super::bitrot_selftest_status_str();
assert!(matches!(rendered, "passed" | "failed" | "unknown"));
assert_eq!(super::bitrot_selftest_status_str(), rendered);
}
#[test]
fn server_info_response_exposes_admin_discovery_paths() {
let usecase = DefaultAdminUsecase::without_context();
@@ -1584,7 +1556,6 @@ mod tests {
pools: None,
},
admin_discovery: system_admin_discovery(&usecase),
bitrot_selftest: super::bitrot_selftest_status_str(),
};
let value = serde_json::to_value(response).expect("server info response should serialize");
-181
View File
@@ -1,181 +0,0 @@
// 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.
//! Startup bitrot algorithm self-test (rustfs/backlog#1873).
//!
//! A drifted hash implementation fails silently in production: every shard
//! reads back "corrupt", heal rewrites healthy data, and cross-platform
//! clusters disagree about which copy is good. [`run_startup_bitrot_self_test`]
//! pins the algorithms once at process start — the check itself runs in well
//! under a millisecond on 4 KiB, so it executes inline before background
//! services come up and the result is published before the server accepts
//! traffic.
//!
//! Outcome surface:
//! - one structured `bitrot_selftest` log event (`passed`/`failed`/`skipped`),
//! - the `rustfs_bitrot_selftest_status` gauge (1=passed, 0=failed, 2=skipped),
//! - [`bitrot_selftest_passed`] for admin/health surfaces,
//! - `RUSTFS_BITROT_SELFTEST_STRICT=on` turns a failure into a startup error
//! (MinIO `bitrotSelfTest` Fatal parity); the default only degrades the
//! status so a bad build cannot brick an existing fleet on upgrade.
use crate::storage_api::startup::background::{BitrotSelfTestError, bitrot_self_test};
use metrics::gauge;
use std::future::Future;
use std::io;
use std::sync::atomic::{AtomicU8, Ordering};
use std::time::Instant;
use tracing::{debug, error, info};
const LOG_COMPONENT_MAIN: &str = "main";
const LOG_SUBSYSTEM_STARTUP: &str = "startup";
const EVENT_BITROT_SELFTEST: &str = "bitrot_selftest";
const METRIC_BITROT_SELFTEST_STATUS: &str = "rustfs_bitrot_selftest_status";
/// Gauge values for [`METRIC_BITROT_SELFTEST_STATUS`].
const STATUS_PASSED: f64 = 1.0;
const STATUS_FAILED: f64 = 0.0;
const STATUS_SKIPPED: f64 = 2.0;
/// Internal cell values for [`BITROT_SELF_TEST_STATUS`].
const STATUS_CELL_UNSET: u8 = 0;
const STATUS_CELL_PASSED: u8 = 1;
const STATUS_CELL_FAILED: u8 = 2;
static BITROT_SELF_TEST_STATUS: AtomicU8 = AtomicU8::new(STATUS_CELL_UNSET);
/// Last recorded self-test outcome: `None` before the first run, then
/// `Some(true)` on a passing check and `Some(false)` on a failed one (a
/// skipped check never publishes, so it cannot read as a pass). The cell is
/// last-writer-wins rather than set-once: production runs the self-test once,
/// and last-writer-wins keeps tests that exercise both outcomes
/// order-independent.
pub fn bitrot_selftest_passed() -> Option<bool> {
match BITROT_SELF_TEST_STATUS.load(Ordering::Acquire) {
STATUS_CELL_UNSET => None,
STATUS_CELL_PASSED => Some(true),
STATUS_CELL_FAILED => Some(false),
_ => None,
}
}
/// Run the bitrot self-test and publish the outcome. In strict mode a failure
/// is returned as an error so the caller aborts startup.
pub(crate) async fn run_startup_bitrot_self_test(enabled: bool, strict: bool) -> io::Result<()> {
run_startup_bitrot_self_test_with(enabled, strict, bitrot_self_test).await
}
async fn run_startup_bitrot_self_test_with<F, Fut>(enabled: bool, strict: bool, run_check: F) -> io::Result<()>
where
F: FnOnce() -> Fut,
Fut: Future<Output = Result<(), BitrotSelfTestError>>,
{
if !enabled {
gauge!(METRIC_BITROT_SELFTEST_STATUS).set(STATUS_SKIPPED);
debug!(
target: "rustfs::main::run",
event = EVENT_BITROT_SELFTEST,
component = LOG_COMPONENT_MAIN,
subsystem = LOG_SUBSYSTEM_STARTUP,
state = "skipped",
reason = "disabled",
"Bitrot self-test skipped"
);
return Ok(());
}
let started = Instant::now();
match run_check().await {
Ok(()) => {
BITROT_SELF_TEST_STATUS.store(STATUS_CELL_PASSED, Ordering::Release);
gauge!(METRIC_BITROT_SELFTEST_STATUS).set(STATUS_PASSED);
info!(
target: "rustfs::main::run",
event = EVENT_BITROT_SELFTEST,
component = LOG_COMPONENT_MAIN,
subsystem = LOG_SUBSYSTEM_STARTUP,
state = "passed",
duration_us = started.elapsed().as_micros() as u64,
"Bitrot self-test passed"
);
}
Err(err) => {
BITROT_SELF_TEST_STATUS.store(STATUS_CELL_FAILED, Ordering::Release);
gauge!(METRIC_BITROT_SELFTEST_STATUS).set(STATUS_FAILED);
error!(
target: "rustfs::main::run",
event = EVENT_BITROT_SELFTEST,
component = LOG_COMPONENT_MAIN,
subsystem = LOG_SUBSYSTEM_STARTUP,
state = "failed",
duration_us = started.elapsed().as_micros() as u64,
error = %err,
"Bitrot self-test failed"
);
if strict {
return Err(io::Error::other(format!("bitrot self-test failed: {err}")));
}
}
}
Ok(())
}
#[cfg(test)]
mod tests {
use super::{BITROT_SELF_TEST_STATUS, STATUS_CELL_UNSET, bitrot_selftest_passed, run_startup_bitrot_self_test_with};
use crate::storage_api::startup::background::BitrotSelfTestError;
use std::future::ready;
use std::sync::atomic::Ordering;
fn failing_check() -> impl Future<Output = Result<(), BitrotSelfTestError>> {
ready(Err(BitrotSelfTestError::RoundtripReadback {
algorithm: "HighwayHash256S",
}))
}
/// All scenarios run sequentially inside one test: the status cell is
/// process-global, so parallel per-scenario tests would race the reset and
/// read each other's outcomes (the exact order-dependent flake class this
/// module exists to avoid).
#[tokio::test]
async fn startup_self_test_publishes_outcome_and_strict_gates_abort() {
BITROT_SELF_TEST_STATUS.store(STATUS_CELL_UNSET, Ordering::Release);
// Skipped: publishes nothing, never fails, never aborts.
run_startup_bitrot_self_test_with(false, true, || async { Ok(()) })
.await
.expect("a disabled self-test must not fail even in strict mode");
assert_eq!(bitrot_selftest_passed(), None, "a skipped run must leave the status unset");
// Passing: publishes Some(true), never fails.
run_startup_bitrot_self_test_with(true, false, || async { Ok(()) })
.await
.expect("a passing check must never fail startup");
assert_eq!(bitrot_selftest_passed(), Some(true), "a passing run must publish Some(true)");
// Failing, non-strict: publishes Some(false) but startup continues.
run_startup_bitrot_self_test_with(true, false, failing_check)
.await
.expect("a failed check must not abort startup in non-strict mode");
assert_eq!(bitrot_selftest_passed(), Some(false), "a failing run must publish Some(false)");
// Failing, strict: startup error carries the failure and the published
// outcome stays a failure.
let err = run_startup_bitrot_self_test_with(true, true, failing_check)
.await
.expect_err("strict mode must turn a failed check into a startup error");
assert!(err.to_string().contains("bitrot self-test failed"));
assert_eq!(bitrot_selftest_passed(), Some(false));
}
}
-1
View File
@@ -76,7 +76,6 @@ pub mod allocator_reclaim;
pub mod app;
pub mod auth;
pub mod auth_keystone;
pub(crate) mod bitrot_selftest;
pub mod capacity;
pub mod cluster_snapshot;
pub mod config;
-14
View File
@@ -33,8 +33,6 @@ pub(crate) const ENV_SCANNER_ENABLED: &str = "RUSTFS_SCANNER_ENABLED";
pub(crate) const ENV_SCANNER_ENABLED_DEPRECATED: &str = "RUSTFS_ENABLE_SCANNER";
pub(crate) const ENV_HEAL_ENABLED: &str = "RUSTFS_HEAL_ENABLED";
pub(crate) const ENV_HEAL_ENABLED_DEPRECATED: &str = "RUSTFS_ENABLE_HEAL";
pub(crate) const ENV_BITROT_SELFTEST_ENABLE: &str = "RUSTFS_BITROT_SELFTEST_ENABLE";
pub(crate) const ENV_BITROT_SELFTEST_STRICT: &str = "RUSTFS_BITROT_SELFTEST_STRICT";
static AUDIT_MODULE_ENABLED: AtomicBool = AtomicBool::new(rustfs_config::DEFAULT_AUDIT_ENABLE);
static NOTIFY_MODULE_ENABLED: AtomicBool = AtomicBool::new(rustfs_config::DEFAULT_NOTIFY_ENABLE);
@@ -49,18 +47,6 @@ pub(crate) fn heal_enabled_from_env() -> bool {
get_env_bool_with_aliases(ENV_HEAL_ENABLED, &[ENV_HEAL_ENABLED_DEPRECATED], true)
}
/// Whether the startup bitrot algorithm self-test runs, defaulting to on
/// (rustfs/backlog#1873).
pub(crate) fn bitrot_selftest_enabled_from_env() -> bool {
rustfs_utils::get_env_bool(ENV_BITROT_SELFTEST_ENABLE, true)
}
/// Whether a failed bitrot self-test aborts startup instead of only logging
/// and exposing a failed status, defaulting to off.
pub(crate) fn bitrot_selftest_strict_from_env() -> bool {
rustfs_utils::get_env_bool(ENV_BITROT_SELFTEST_STRICT, false)
}
/// Last published audit-module state.
pub fn is_audit_module_enabled() -> bool {
AUDIT_MODULE_ENABLED.load(Ordering::Relaxed)
+1 -10
View File
@@ -12,10 +12,7 @@
// See the License for the specific language governing permissions and
// limitations under the License.
use crate::bitrot_selftest::run_startup_bitrot_self_test;
use crate::module_switches::{
bitrot_selftest_enabled_from_env, bitrot_selftest_strict_from_env, heal_enabled_from_env, scanner_enabled_from_env,
};
use crate::module_switches::{heal_enabled_from_env, scanner_enabled_from_env};
use crate::storage_api::startup::background::{ECStore, set_workload_admission_snapshot_provider};
use crate::workload_admission::RustFsWorkloadAdmissionSnapshotProvider;
use rustfs_concurrency::WorkloadAdmissionSnapshotProvider;
@@ -30,12 +27,6 @@ const LOG_SUBSYSTEM_STARTUP: &str = "startup";
const EVENT_BACKGROUND_SERVICES_CONFIGURED: &str = "background_services_configured";
pub(crate) async fn init_background_service_runtime(store: Arc<ECStore>) -> Result<bool> {
// Pin the bitrot algorithms before anything can write or verify a shard:
// the check costs well under a millisecond, and in strict mode a drifted
// build must abort here rather than after it has touched data
// (rustfs/backlog#1873).
run_startup_bitrot_self_test(bitrot_selftest_enabled_from_env(), bitrot_selftest_strict_from_env()).await?;
let _ = create_ahm_services_cancel_token();
let enable_scanner = scanner_enabled_from_env();
-4
View File
@@ -569,10 +569,6 @@ pub(crate) mod ecstore_erasure {
pub(crate) use rustfs_ecstore::api::erasure::{BitrotReader, Erasure};
}
/// Startup bitrot algorithm self-test (rustfs/backlog#1873), re-exported for
/// the root facade's background-startup section.
pub(crate) use rustfs_ecstore::api::erasure::{BitrotSelfTestError, bitrot_self_test};
pub(crate) mod ecstore_storage {
#[cfg(test)]
pub(crate) use rustfs_ecstore::api::storage::init_local_disks;
+1 -3
View File
@@ -214,9 +214,7 @@ pub(crate) mod startup {
}
pub(crate) mod background {
pub(crate) use crate::storage::storage_api::{
BitrotSelfTestError, ECStore, bitrot_self_test, set_workload_admission_snapshot_provider,
};
pub(crate) use crate::storage::storage_api::{ECStore, set_workload_admission_snapshot_provider};
}
pub(crate) mod bucket_metadata {