Compare commits

..

3 Commits

Author SHA1 Message Date
houseme 984c705713 docs(ecstore): fix bitrot comment typo (#6168)
Co-authored-by: heihutu <heihutu@gmail.com>
2026-08-17 08:24:54 +00:00
houseme 23b17c2d5a feat(madmin): add a SigV4-signed admin client for heal and scanner APIs (HS-05) (#6166)
feat(madmin): add a SigV4-signed admin client for heal and scanner APIs

The madmin crate held only wire types; automation and mc-style tooling
had no way to drive the heal/scanner admin surface without hand-rolled
HTTP. Add `AdminClient`, which signs with the same rustfs-signer path
the server authenticates (UNSIGNED-PAYLOAD marker, matching RustFS peer
admin calls) and wraps:

- heal_start / heal_status / heal_stop over POST /rustfs/admin/v3/heal/
  (bucket/prefix path params percent-encoded per segment; stop models
  the server's two cancel branches: token-scoped task status vs
  path-scoped start-success receipt);
- background_heal_status, scanner_status (freshness typed), plus
  ilm_expiry_status / replacement_recovery_status passthroughs;
- a public get_json escape hatch for endpoints not wrapped yet.

Wire types follow the madmin-go model (SDK-owned mirrors pinned by
round-trip tests): HealOpts with serde defaults so partial settings
objects decode, HealScanMode accepting both the numeric and name
encodings, and status structs that type the fields operators branch on
while flattening unknown nested payloads verbatim so server additions
cannot break the client. Errors map to a closed AdminClientError enum
(InvalidEndpoint / Transport / HttpStatus with body / Decode).

Tests cover wire round-trips, path building, both stop branches, error
mapping, and — via a dependency-free raw-TCP test server — that signed
requests carry a SigV4 Authorization header, the right method/path/
query, and the expected JSON body.

Closes rustfs/backlog#1869 (first increment; single-sourcing the wire
structs server-side and an embedded-server e2e roundtrip are noted as
follow-ups there).

Co-authored-by: heihutu <heihutu@gmail.com>
2026-08-17 15:04:49 +08:00
houseme 89e2513205 feat(ecstore): pin bitrot algorithms with a startup self-test (HS-11) (#6165)
feat(ecstore): pin bitrot algorithms with a startup self-test

A drifted HighwayHash implementation fails silently: every shard reads
back corrupt, heal rewrites healthy data, and cross-platform clusters
disagree about which copy is good. Mirror MinIO's bitrotSelfTest by
verifying, once at process start:

- known-answer digests for HighwayHash256S / HighwayHash256SLegacy over
  a deterministic 4096-byte xorshift64* payload, plus the externally
  verifiable FIPS SHA-256 "abc" vector guarding the HashAlgorithm
  plumbing itself;
- an end-to-end roundtrip per streaming variant (encode -> size formula
  -> bitrot_verify -> BitrotReader read-back), over full blocks and a
  partial tail;
- tamper detection: one flipped byte in the final data block and one in
  the leading hash must both be rejected as a hash mismatch, not by an
  incidental read error.

The check costs microseconds and runs inline in
init_background_service_runtime before any shard can be written or
verified. Outcome surfaces as one structured bitrot_selftest log event,
the rustfs_bitrot_selftest_status gauge (1=passed / 0=failed / 2=skipped),
a bitrotSelftest field on the admin server-info response, and
RUSTFS_BITROT_SELFTEST_STRICT=on turns a failure into a startup error
(MinIO Fatal parity; the default only degrades the status so a bad build
cannot brick an existing fleet on upgrade).

Closes rustfs/backlog#1873 (HS-11).

Co-authored-by: heihutu <heihutu@gmail.com>
2026-08-17 15:04:23 +08:00
17 changed files with 2534 additions and 1575 deletions
Generated
+5
View File
@@ -9825,14 +9825,19 @@ 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,6 +40,7 @@ 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, BitrotWriter, BitrotWriterWrapper, CustomWriter, Erasure, ErasureConstructionError, ReedSolomonEncoder,
calc_shard_size, calc_shard_size_legacy,
BitrotReader, BitrotSelfTestError, BitrotWriter, BitrotWriterWrapper, CustomWriter, Erasure, ErasureConstructionError,
ReedSolomonEncoder, bitrot_self_test, calc_shard_size, calc_shard_size_legacy,
};
}
@@ -667,368 +667,6 @@ 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 {
@@ -1583,12 +1221,71 @@ impl<S: ReplicationStorage> ReplicationPool<S> {
let storage = self.storage.clone();
let handle = tokio::spawn(async move {
let Some(recovery_guard) = acquire_mrf_recovery_guard(&storage).await else {
return;
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(entries) = load_mrf_recovery_entries(&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;
}
};
set_durable_mrf_backlog_snapshot(durable_mrf_backlog_summary_from_entries(&entries));
@@ -1597,8 +1294,187 @@ impl<S: ReplicationStorage> ReplicationPool<S> {
let mut retry_entries = Vec::new();
for entry in entries.iter() {
let Some(admission) = replay_mrf_entry(entry, &storage, &mut retry_entries).await else {
continue;
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
}
}
}
};
if admission == ReplicationQueueAdmission::Missed {
@@ -1608,7 +1484,29 @@ impl<S: ReplicationStorage> ReplicationPool<S> {
}
}
let retained = resolve_retained_mrf_entries(&storage, &recovery_guard, &entries, &retry_entries).await;
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_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
+291 -12
View File
@@ -820,10 +820,263 @@ 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_shard_file_size, bitrot_verify, write_all_vectored,
BitrotReader, BitrotWriter, BitrotWriterWrapper, CustomWriter, bitrot_kat_check, bitrot_self_test,
bitrot_self_test_payload, bitrot_shard_file_size, bitrot_verify, write_all_vectored,
};
use super::{MAX_RETAINED_CHUNKS_PER_BLOCK, ShardChunkRead, ShardSource};
use bytes::Bytes;
@@ -1090,6 +1343,32 @@ 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();
@@ -1189,7 +1468,7 @@ mod tests {
let last = corrupt.len() - 1;
corrupt[last] ^= 0x80;
let err = bitrot_verify(
Cursor::new(corrupt),
std::io::Cursor::new(corrupt),
super::bitrot_shard_file_size(data.len(), shard_size, algo.clone()),
data.len(),
algo,
@@ -1282,7 +1561,7 @@ mod tests {
#[tokio::test]
async fn bitrot_reader_rejects_output_buffers_larger_than_shard_size() {
let mut reader = BitrotReader::new(Cursor::new(Vec::<u8>::new()), 4, HashAlgorithm::None, false);
let mut reader = BitrotReader::new(std::io::Cursor::new(Vec::<u8>::new()), 4, HashAlgorithm::None, false);
let mut out = [0u8; 5];
let err = reader
.read(&mut out)
@@ -1407,7 +1686,7 @@ mod tests {
(HashAlgorithm::HighwayHash256, true),
] {
let label = format!("{algo:?}");
let writer = Cursor::new(Vec::<u8>::new());
let writer = std::io::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();
@@ -1492,7 +1771,7 @@ mod tests {
}
async fn encode_one_block(payload: &[u8], shard_size: usize, algo: HashAlgorithm) -> Vec<u8> {
let mut w = BitrotWriter::new(Cursor::new(Vec::<u8>::new()), shard_size, algo);
let mut w = BitrotWriter::new(std::io::Cursor::new(Vec::<u8>::new()), shard_size, algo);
w.write(payload).await.unwrap();
w.into_inner().into_inner()
}
@@ -1600,7 +1879,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(Cursor::new(Vec::<u8>::new()), shard_size, algo.clone());
let mut w = BitrotWriter::new(std::io::Cursor::new(Vec::<u8>::new()), shard_size, algo.clone());
for chunk in payload.chunks(shard_size) {
w.write(chunk).await.unwrap();
}
@@ -1674,14 +1953,14 @@ mod tests {
w.write(&data).await.expect("write shard");
let mut via_read = vec![0u8; SHARD];
let n1 = BitrotReader::new(Cursor::new(encoded.clone()), SHARD, algo.clone(), false)
let n1 = BitrotReader::new(std::io::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(Cursor::new(encoded), SHARD, algo.clone(), false)
let n2 = BitrotReader::new(std::io::Cursor::new(encoded), SHARD, algo.clone(), false)
.read_appending(&mut via_append, SHARD)
.await
.expect("read_appending");
@@ -1706,7 +1985,7 @@ mod tests {
encoded.truncate(encoded.len() - 1);
let mut out: Vec<u8> = Vec::with_capacity(SHARD);
let err = BitrotReader::new(Cursor::new(encoded), SHARD, algo.clone(), false)
let err = BitrotReader::new(std::io::Cursor::new(encoded), SHARD, algo.clone(), false)
.read_appending(&mut out, SHARD)
.await
.expect_err("a truncated shard must not succeed");
@@ -1732,7 +2011,7 @@ mod tests {
encoded[last] ^= 0xff;
let mut out: Vec<u8> = Vec::with_capacity(SHARD);
let err = BitrotReader::new(Cursor::new(encoded), SHARD, algo, false)
let err = BitrotReader::new(std::io::Cursor::new(encoded), SHARD, algo, false)
.read_appending(&mut out, SHARD)
.await
.expect_err("a corrupt shard must not verify");
@@ -1844,7 +2123,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 = Cursor::new(encoded.clone());
let mut streamed = std::io::Cursor::new(encoded.clone());
assert!(
ShardSource::try_take_block(&mut streamed, 8).is_none(),
"a non-Bytes source must stay on the streaming path"
@@ -1872,7 +2151,7 @@ mod tests {
);
let mut via_stream: Vec<u8> = Vec::with_capacity(SHARD);
BitrotReader::new(Cursor::new(encoded), SHARD, algo, false)
BitrotReader::new(std::io::Cursor::new(encoded), SHARD, algo, false)
.read_appending(&mut via_stream, SHARD)
.await
.expect("streaming read");
+5
View File
@@ -37,7 +37,11 @@ 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"] }
@@ -49,3 +53,4 @@ doctest = false
[dev-dependencies]
rmp-serde.workspace = true
tokio = { workspace = true, features = ["macros", "rt-multi-thread", "net"] }
+851
View File
@@ -0,0 +1,851 @@
// 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,6 +12,7 @@
// 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;
@@ -25,6 +26,7 @@ pub mod trace;
pub mod user;
pub mod utils;
pub use client::*;
pub use group::*;
pub use info_commands::*;
pub use policy::*;
+241 -258
View File
@@ -66,20 +66,18 @@ 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::{
IamSys, NewServiceAccountOpts, SITE_REPLICATOR_SERVICE_ACCOUNT, UpdateServiceAccountOpts, get_claims_from_token_with_secret,
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, 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,
SRIAMPolicy, SRILMExpiryStatsSummary, SRInfo, SRMetric, SRMetricsSummary, SRPeerError, SRPeerJoinReq, SRPendingOperation,
SRPolicyMapping, SRPolicyStatsSummary, SRRemoveReq, SRResyncOpStatus, SRRetryStats, SRSessionPolicy, SRSiteSummary,
SRStateEditReq, SRStateInfo, SRStatusInfo, SRSvcAccCreate, SRUserStatsSummary, SiteReplicationInfo, SyncStatus, WorkerStat,
};
use rustfs_policy::policy::{
Policy,
@@ -9280,16 +9278,247 @@ async fn apply_iam_item(item: SRIAMItem) -> S3Result<()> {
let incoming_updated_at = item.updated_at;
match item.r#type.as_str() {
"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,
"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(())
}
// 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 => 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,
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"))
}
_ => Err(s3_error!(
NotImplemented,
"site replication IAM item type `{}` is not supported",
@@ -9298,252 +9527,6 @@ 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,6 +417,13 @@ 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)]
@@ -433,6 +440,14 @@ 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)>> {
@@ -464,6 +479,7 @@ 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| {
@@ -1535,6 +1551,18 @@ 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();
@@ -1556,6 +1584,7 @@ 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
@@ -0,0 +1,181 @@
// 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,6 +76,7 @@ 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,6 +33,8 @@ 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);
@@ -47,6 +49,18 @@ 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)
+10 -1
View File
@@ -12,7 +12,10 @@
// See the License for the specific language governing permissions and
// limitations under the License.
use crate::module_switches::{heal_enabled_from_env, scanner_enabled_from_env};
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::storage_api::startup::background::{ECStore, set_workload_admission_snapshot_provider};
use crate::workload_admission::RustFsWorkloadAdmissionSnapshotProvider;
use rustfs_concurrency::WorkloadAdmissionSnapshotProvider;
@@ -27,6 +30,12 @@ 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,6 +569,10 @@ 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;
+3 -1
View File
@@ -214,7 +214,9 @@ pub(crate) mod startup {
}
pub(crate) mod background {
pub(crate) use crate::storage::storage_api::{ECStore, set_workload_admission_snapshot_provider};
pub(crate) use crate::storage::storage_api::{
BitrotSelfTestError, ECStore, bitrot_self_test, set_workload_admission_snapshot_provider,
};
}
pub(crate) mod bucket_metadata {