feat(scanner): accept durable recovery intents

Add a CAS-backed scanner usage recovery intent record for async full rebuild admission. The admin reset endpoint can now persist and replay idempotent intent acceptance before returning 202, and a read-only status route exposes the durable request state.

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

Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
This commit is contained in:
houseme
2026-09-07 11:47:06 +08:00
parent 672087ec0d
commit 4ec063e55f
7 changed files with 629 additions and 7 deletions
+7 -4
View File
@@ -83,10 +83,13 @@ pub use remote_scanner::{
pub use runtime_config::{apply_scanner_runtime_config, scanner_runtime_config_status, validate_scanner_runtime_config};
pub use rustfs_scanner_metrics::last_minute;
pub use scanner::{
ScannerCycleRecoveryMarker, ScannerCycleRecoveryStatus, ScannerCycleScheduleStatus, ScannerPauseBacklogAlertReason,
ScannerPauseBacklogPhase, ScannerPauseBacklogStatus, ScannerPauseBacklogThresholds, ScannerUsageStateResetResult,
init_data_scanner, init_scanner_with_recovery, reset_scanner_cycle_recovery, reset_scanner_usage_state_for_full_rebuild,
scanner_cycle_recovery_status, scanner_cycle_schedule_status, scanner_pause_backlog_status, scanner_topology_digest,
SCANNER_RECOVERY_INTENT_ACTION_USAGE_FULL_REBUILD, ScannerCycleRecoveryMarker, ScannerCycleRecoveryStatus,
ScannerCycleScheduleStatus, ScannerPauseBacklogAlertReason, ScannerPauseBacklogPhase, ScannerPauseBacklogStatus,
ScannerPauseBacklogThresholds, ScannerRecoveryIntentAcceptResult, ScannerRecoveryIntentConflict, ScannerRecoveryIntentRecord,
ScannerRecoveryIntentRequest, ScannerUsageStateResetResult, accept_scanner_usage_recovery_intent,
get_scanner_usage_recovery_intent, init_data_scanner, init_scanner_with_recovery, reset_scanner_cycle_recovery,
reset_scanner_usage_state_for_full_rebuild, scanner_cycle_recovery_status, scanner_cycle_schedule_status,
scanner_pause_backlog_status, scanner_recovery_actor_sha256, scanner_topology_digest,
};
pub use scanner_io::{
ScannerDirtyUsageAckError, ScannerDirtyUsageBucket, ScannerDirtyUsageSnapshot, ScannerDirtyUsageState,
+5 -2
View File
@@ -3590,8 +3590,11 @@ pub use backlog::{
#[cfg(test)]
pub(crate) use cycle_state::encode_scanner_cycle_fence_for_test;
pub use cycle_state::{
ScannerCycleRecoveryMarker, ScannerCycleRecoveryStatus, ScannerUsageStateResetResult, reset_scanner_cycle_recovery,
reset_scanner_usage_state_for_full_rebuild, scanner_cycle_recovery_status,
SCANNER_RECOVERY_INTENT_ACTION_USAGE_FULL_REBUILD, ScannerCycleRecoveryMarker, ScannerCycleRecoveryStatus,
ScannerRecoveryIntentAcceptResult, ScannerRecoveryIntentConflict, ScannerRecoveryIntentRecord, ScannerRecoveryIntentRequest,
ScannerUsageStateResetResult, accept_scanner_usage_recovery_intent, get_scanner_usage_recovery_intent,
reset_scanner_cycle_recovery, reset_scanner_usage_state_for_full_rebuild, scanner_cycle_recovery_status,
scanner_recovery_actor_sha256,
};
pub(crate) use cycle_state::{
current_scanner_leader_epoch, decode_persisted_scanner_cycle_fence, load_scanner_cycle_state_for_startup,
+260
View File
@@ -34,6 +34,10 @@ const LEGACY_INCOMPLETE_USAGE_FLOOR_RECOVERY: &str = "legacy_empty_usage_floor";
const CACHE_CYCLE_AHEAD: &str = "cache_cycle_ahead";
const SCANNER_USAGE_STATE_RESET_MODE_FULL_REBUILD: &str = "full-rebuild";
const SCANNER_RECOVERY_INTENT_SCHEMA_VERSION: u16 = 1;
const SCANNER_RECOVERY_INTENT_PREFIX: &str = ".usage.v2.recovery-intents";
const MAX_SCANNER_RECOVERY_INTENT_BYTES: u64 = 16 * 1024;
pub const SCANNER_RECOVERY_INTENT_ACTION_USAGE_FULL_REBUILD: &str = "scanner-usage-full-rebuild";
#[cfg(test)]
pub(super) mod cleanup_io_fault {
@@ -146,6 +150,43 @@ pub struct ScannerUsageStateResetResult {
pub reset_paths: Vec<String>,
}
#[derive(Clone, Debug, Serialize, Deserialize, PartialEq, Eq)]
#[serde(deny_unknown_fields)]
pub struct ScannerRecoveryIntentRequest {
pub action: String,
pub mode: String,
pub idempotency_key: String,
pub actor_sha256: String,
}
#[derive(Clone, Debug, Serialize, Deserialize, PartialEq, Eq)]
#[serde(deny_unknown_fields)]
pub struct ScannerRecoveryIntentRecord {
pub schema_version: u16,
pub intent_id: String,
pub action: String,
pub mode: String,
pub state: String,
pub actor_sha256: String,
pub idempotency_key_sha256: String,
pub request_sha256: String,
pub accepted_at_unix_secs: u64,
}
#[derive(Clone, Debug, Serialize, PartialEq, Eq)]
pub struct ScannerRecoveryIntentConflict {
pub intent_id: String,
pub state: String,
}
#[derive(Clone, Debug, Serialize, PartialEq, Eq)]
#[serde(tag = "outcome", rename_all = "snake_case")]
pub enum ScannerRecoveryIntentAcceptResult {
Accepted { record: ScannerRecoveryIntentRecord },
Replayed { record: ScannerRecoveryIntentRecord },
Conflict { existing: ScannerRecoveryIntentConflict },
}
#[derive(Clone, Debug)]
pub(super) struct ScannerUsageStateResetSlot {
path: String,
@@ -360,6 +401,225 @@ fn unix_now_secs() -> u64 {
u64::try_from(Utc::now().timestamp()).unwrap_or(0)
}
fn sha256_hex(parts: &[&[u8]]) -> String {
let mut hasher = Sha256::new();
for part in parts {
hasher.update(part);
hasher.update([0]);
}
const HEX: &[u8; 16] = b"0123456789abcdef";
let digest = hasher.finalize();
let mut encoded = String::with_capacity(64);
for byte in digest {
encoded.push(char::from(HEX[usize::from(byte >> 4)]));
encoded.push(char::from(HEX[usize::from(byte & 0x0f)]));
}
encoded
}
fn is_canonical_sha256(value: &str) -> bool {
value.len() == 64
&& value
.bytes()
.all(|byte| byte.is_ascii_hexdigit() && !byte.is_ascii_uppercase())
}
fn validate_idempotency_key(key: &str) -> Result<(), ScannerError> {
if !(8..=256).contains(&key.len()) {
return Err(ScannerError::Other(
"scanner recovery intent idempotency key length is unsupported".to_string(),
));
}
if !key
.bytes()
.all(|byte| byte.is_ascii() && !byte.is_ascii_control() && !byte.is_ascii_whitespace())
{
return Err(ScannerError::Other(
"scanner recovery intent idempotency key must be printable ASCII without whitespace".to_string(),
));
}
Ok(())
}
fn validate_recovery_intent_request(request: &ScannerRecoveryIntentRequest) -> Result<(), ScannerError> {
if request.action != SCANNER_RECOVERY_INTENT_ACTION_USAGE_FULL_REBUILD {
return Err(ScannerError::Other("scanner recovery intent action is unsupported".to_string()));
}
if request.mode != SCANNER_USAGE_STATE_RESET_MODE_FULL_REBUILD {
return Err(ScannerError::Other("scanner recovery intent mode is unsupported".to_string()));
}
if !is_canonical_sha256(&request.actor_sha256) {
return Err(ScannerError::Other("scanner recovery intent actor identity is invalid".to_string()));
}
validate_idempotency_key(&request.idempotency_key)
}
fn scanner_recovery_intent_path(intent_id: &str) -> Result<String, ScannerError> {
if !is_canonical_sha256(intent_id) {
return Err(ScannerError::Other("scanner recovery intent id is invalid".to_string()));
}
Ok(format!("{SCANNER_RECOVERY_INTENT_PREFIX}/{intent_id}.json"))
}
pub fn scanner_recovery_actor_sha256(actor: &str) -> String {
sha256_hex(&[b"scanner-recovery-actor-v1", actor.as_bytes()])
}
fn scanner_recovery_intent_candidate(
request: &ScannerRecoveryIntentRequest,
) -> Result<ScannerRecoveryIntentRecord, ScannerError> {
validate_recovery_intent_request(request)?;
let intent_id = sha256_hex(&[
b"scanner-recovery-intent-v1",
request.actor_sha256.as_bytes(),
request.idempotency_key.as_bytes(),
]);
Ok(ScannerRecoveryIntentRecord {
schema_version: SCANNER_RECOVERY_INTENT_SCHEMA_VERSION,
intent_id,
action: request.action.clone(),
mode: request.mode.clone(),
state: "accepted".to_string(),
actor_sha256: request.actor_sha256.clone(),
idempotency_key_sha256: sha256_hex(&[b"scanner-recovery-idempotency-key-v1", request.idempotency_key.as_bytes()]),
request_sha256: sha256_hex(&[
b"scanner-recovery-request-v1",
request.actor_sha256.as_bytes(),
request.action.as_bytes(),
request.mode.as_bytes(),
request.idempotency_key.as_bytes(),
]),
accepted_at_unix_secs: unix_now_secs(),
})
}
fn validate_recovery_intent_record(record: &ScannerRecoveryIntentRecord) -> Result<(), ScannerError> {
if record.schema_version != SCANNER_RECOVERY_INTENT_SCHEMA_VERSION {
return Err(ScannerError::Other("scanner recovery intent schema is unsupported".to_string()));
}
if !is_canonical_sha256(&record.intent_id)
|| !is_canonical_sha256(&record.actor_sha256)
|| !is_canonical_sha256(&record.idempotency_key_sha256)
|| !is_canonical_sha256(&record.request_sha256)
{
return Err(ScannerError::Other("scanner recovery intent identity is invalid".to_string()));
}
if record.action != SCANNER_RECOVERY_INTENT_ACTION_USAGE_FULL_REBUILD
|| record.mode != SCANNER_USAGE_STATE_RESET_MODE_FULL_REBUILD
|| !matches!(record.state.as_str(), "accepted" | "running" | "completed" | "failed")
{
return Err(ScannerError::Other("scanner recovery intent state is invalid".to_string()));
}
Ok(())
}
fn decode_recovery_intent_record(data: &[u8]) -> Result<ScannerRecoveryIntentRecord, ScannerError> {
let record: ScannerRecoveryIntentRecord =
serde_json::from_slice(data).map_err(|err| ScannerError::Other(format!("scanner recovery intent is invalid: {err}")))?;
validate_recovery_intent_record(&record)?;
Ok(record)
}
fn compare_recovery_intent(
expected: &ScannerRecoveryIntentRecord,
existing: ScannerRecoveryIntentRecord,
) -> ScannerRecoveryIntentAcceptResult {
if existing.actor_sha256 == expected.actor_sha256
&& existing.action == expected.action
&& existing.mode == expected.mode
&& existing.idempotency_key_sha256 == expected.idempotency_key_sha256
&& existing.request_sha256 == expected.request_sha256
{
ScannerRecoveryIntentAcceptResult::Replayed { record: existing }
} else {
ScannerRecoveryIntentAcceptResult::Conflict {
existing: ScannerRecoveryIntentConflict {
intent_id: existing.intent_id,
state: existing.state,
},
}
}
}
async fn read_recovery_intent_record(
storeapi: Arc<impl ScannerObjectIO>,
path: &str,
) -> Result<Option<ScannerRecoveryIntentRecord>, ScannerError> {
let mut reader = match storeapi
.get_object_reader(
RUSTFS_META_BUCKET,
path,
None,
http::HeaderMap::new(),
&ScannerObjectOptions {
no_lock: true,
..Default::default()
},
)
.await
{
Ok(reader) => reader,
Err(
EcstoreError::FileNotFound
| EcstoreError::VolumeNotFound
| EcstoreError::ObjectNotFound(_, _)
| EcstoreError::BucketNotFound(_)
| EcstoreError::ConfigNotFound,
) => return Ok(None),
Err(err) => return Err(ScannerError::Other(format!("failed to read scanner recovery intent: {err}"))),
};
let max_object_size = i64::try_from(MAX_SCANNER_RECOVERY_INTENT_BYTES).unwrap_or(i64::MAX);
if reader.object_info.is_dir || reader.object_info.size < 0 || reader.object_info.size > max_object_size {
return Err(ScannerError::Other("scanner recovery intent exceeds the bounded object size".to_string()));
}
let max_len = usize::try_from(MAX_SCANNER_RECOVERY_INTENT_BYTES).unwrap_or(usize::MAX);
let mut data = Vec::new();
(&mut reader)
.take(MAX_SCANNER_RECOVERY_INTENT_BYTES.saturating_add(1))
.read_to_end(&mut data)
.await
.map_err(|err| ScannerError::Other(format!("failed to read scanner recovery intent: {err}")))?;
if data.is_empty() {
return Err(ScannerError::Other("scanner recovery intent is empty".to_string()));
}
if data.len() > max_len {
return Err(ScannerError::Other("scanner recovery intent exceeds the bounded object size".to_string()));
}
decode_recovery_intent_record(&data).map(Some)
}
pub async fn get_scanner_usage_recovery_intent(
storeapi: Arc<impl ScannerObjectIO>,
intent_id: &str,
) -> Result<Option<ScannerRecoveryIntentRecord>, ScannerError> {
let path = scanner_recovery_intent_path(intent_id)?;
read_recovery_intent_record(storeapi, &path).await
}
pub async fn accept_scanner_usage_recovery_intent(
storeapi: Arc<impl ScannerObjectIO>,
request: ScannerRecoveryIntentRequest,
) -> Result<ScannerRecoveryIntentAcceptResult, ScannerError> {
let candidate = scanner_recovery_intent_candidate(&request)?;
let path = scanner_recovery_intent_path(&candidate.intent_id)?;
let encoded = serde_json::to_vec(&candidate)
.map_err(|err| ScannerError::Other(format!("failed to encode scanner recovery intent: {err}")))?;
match save_config_with_preconditions(storeapi.clone(), &path, encoded, DataUsageCacheRevision::Missing.preconditions()).await
{
Ok(_) => Ok(ScannerRecoveryIntentAcceptResult::Accepted { record: candidate }),
Err(EcstoreError::PreconditionFailed) => {
let existing = read_recovery_intent_record(storeapi, &path).await?;
let Some(existing) = existing else {
return Err(ScannerError::Other(
"scanner recovery intent disappeared after creation conflict".to_string(),
));
};
Ok(compare_recovery_intent(&candidate, existing))
}
Err(err) => Err(ScannerError::Other(format!("failed to persist scanner recovery intent: {err}"))),
}
}
fn recovery_status(state: &str, reason: Option<&str>, retryable: bool) -> ScannerCycleRecoveryStatus {
ScannerCycleRecoveryStatus {
path: DATA_USAGE_BLOOM_NAME_PATH.clone(),
@@ -97,6 +97,15 @@ async fn run_disabled_startup(ctx: CancellationToken, store: Arc<ECStore>) {
);
}
fn recovery_intent_request(key: &str, actor: &str) -> ScannerRecoveryIntentRequest {
ScannerRecoveryIntentRequest {
action: SCANNER_RECOVERY_INTENT_ACTION_USAGE_FULL_REBUILD.to_string(),
mode: "full-rebuild".to_string(),
idempotency_key: key.to_string(),
actor_sha256: scanner_recovery_actor_sha256(actor),
}
}
async fn assert_reset_fences(store: &Arc<ECStore>) {
let data = read_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str())
.await
@@ -203,6 +212,124 @@ async fn disabled_cleanup_recovers_after_child_process_crash_boundaries() {
}
}
#[tokio::test]
#[serial]
async fn scanner_recovery_intent_accept_is_durable_and_idempotent() {
let (_dir, store) = setup_scanner_cycle_store().await;
let request = recovery_intent_request("intent-key-0001", "operator-a");
let first = accept_scanner_usage_recovery_intent(store.clone(), request.clone())
.await
.expect("first intent should persist");
let record = match first {
ScannerRecoveryIntentAcceptResult::Accepted { record } => record,
other => panic!("first request must create a durable intent: {other:?}"),
};
assert_eq!(record.state, "accepted");
assert_eq!(record.action, SCANNER_RECOVERY_INTENT_ACTION_USAGE_FULL_REBUILD);
assert_ne!(record.idempotency_key_sha256, request.idempotency_key);
let restarted = restart_scanner_cycle_store_from(&store).await;
let queried = get_scanner_usage_recovery_intent(restarted.clone(), &record.intent_id)
.await
.expect("persisted intent should read after restart")
.expect("persisted intent should exist");
assert_eq!(queried, record);
let replay = accept_scanner_usage_recovery_intent(restarted, request)
.await
.expect("lost response retry should be idempotent");
assert_eq!(replay, ScannerRecoveryIntentAcceptResult::Replayed { record });
}
#[tokio::test]
#[serial]
async fn scanner_recovery_intent_rejects_same_namespace_conflict() {
let (_dir, store) = setup_scanner_cycle_store().await;
let request = recovery_intent_request("intent-key-0002", "operator-a");
let record = match accept_scanner_usage_recovery_intent(store.clone(), request.clone())
.await
.expect("first intent should persist")
{
ScannerRecoveryIntentAcceptResult::Accepted { record } => record,
other => panic!("first request must create a durable intent: {other:?}"),
};
let mut unsupported = request.clone();
unsupported.mode = "future-mode".to_string();
let error = accept_scanner_usage_recovery_intent(store.clone(), unsupported)
.await
.expect_err("unsupported mode must fail before it can collide with a durable record");
assert!(error.to_string().contains("mode is unsupported"));
let same_key_other_actor = recovery_intent_request("intent-key-0002", "operator-b");
let accepted_other_actor = accept_scanner_usage_recovery_intent(store.clone(), same_key_other_actor)
.await
.expect("another actor owns an independent idempotency namespace");
assert!(matches!(accepted_other_actor, ScannerRecoveryIntentAcceptResult::Accepted { .. }));
let path = format!(".usage.v2.recovery-intents/{}.json", record.intent_id);
let mut existing = record.clone();
existing.request_sha256 = scanner_recovery_actor_sha256("different-request");
save_config(store.clone(), &path, serde_json::to_vec(&existing).expect("mutated record should encode"))
.await
.expect("mutate durable record");
let conflict = accept_scanner_usage_recovery_intent(store, request)
.await
.expect("same namespace collision should be reported");
assert_eq!(
conflict,
ScannerRecoveryIntentAcceptResult::Conflict {
existing: ScannerRecoveryIntentConflict {
intent_id: record.intent_id,
state: "accepted".to_string(),
},
}
);
}
#[tokio::test]
#[serial]
async fn scanner_recovery_intent_query_rejects_corrupt_or_unknown_records() {
let (_dir, store) = setup_scanner_cycle_store().await;
let request = recovery_intent_request("intent-key-0003", "operator-a");
let record = match accept_scanner_usage_recovery_intent(store.clone(), request)
.await
.expect("first intent should persist")
{
ScannerRecoveryIntentAcceptResult::Accepted { record } => record,
other => panic!("first request must create a durable intent: {other:?}"),
};
let path = format!(".usage.v2.recovery-intents/{}.json", record.intent_id);
save_config(store.clone(), &path, b"{corrupt".to_vec())
.await
.expect("corrupt durable record");
let error = get_scanner_usage_recovery_intent(store.clone(), &record.intent_id)
.await
.expect_err("corrupt intent must not decode as absent");
assert!(error.to_string().contains("scanner recovery intent is invalid"));
let unknown = get_scanner_usage_recovery_intent(store, &scanner_recovery_actor_sha256("missing"))
.await
.expect("missing intent should read as absent");
assert!(unknown.is_none());
}
#[tokio::test]
#[serial]
async fn scanner_recovery_intent_query_rejects_oversized_records() {
let (_dir, store) = setup_scanner_cycle_store().await;
let intent_id = scanner_recovery_actor_sha256("oversized-intent");
let path = format!(".usage.v2.recovery-intents/{intent_id}.json");
save_config(store.clone(), &path, vec![b'a'; 16 * 1024 + 1])
.await
.expect("oversized durable record");
let error = get_scanner_usage_recovery_intent(store, &intent_id)
.await
.expect_err("oversized intent must not be materialized");
assert!(error.to_string().contains("bounded object size"), "{error}");
}
#[tokio::test]
#[serial]
async fn disabled_cleanup_reopens_persisted_intent_without_starting_scanner() {
+194 -1
View File
@@ -62,6 +62,26 @@ struct ScannerCycleResetRequest {
#[serde(deny_unknown_fields)]
struct ScannerUsageStateResetRequest {
mode: String,
#[serde(default, rename = "async")]
async_intent: bool,
#[serde(default)]
idempotency_key: Option<String>,
}
#[derive(Debug, Serialize)]
struct ScannerRecoveryIntentResponse {
status: &'static str,
action: String,
mode: String,
intent_id: String,
state: String,
}
#[derive(Debug, Serialize)]
struct ScannerRecoveryIntentConflictResponse {
status: &'static str,
intent_id: String,
state: String,
}
#[derive(Debug, Serialize)]
@@ -242,6 +262,11 @@ pub fn register_scanner_route(r: &mut S3Router<AdminOperation>) -> std::io::Resu
format!("{ADMIN_PREFIX}/v3/scanner/usage-state/reset").as_str(),
AdminOperation(&ScannerUsageStateResetHandler {}),
)?;
r.insert(
Method::GET,
format!("{ADMIN_PREFIX}/v3/scanner/usage-state/recovery-intents/{{intent_id}}").as_str(),
AdminOperation(&ScannerUsageStateRecoveryIntentStatusHandler {}),
)?;
r.insert(
Method::GET,
format!("{ADMIN_PREFIX}/v3/ilm/expiry/status").as_str(),
@@ -269,11 +294,79 @@ async fn validate_scanner_reset_request(req: &S3Request<Body>) -> S3Result<Crede
}
fn json_response(body: Vec<u8>) -> S3Result<S3Response<(StatusCode, Body)>> {
json_response_with_status(StatusCode::OK, body)
}
fn json_response_with_status(status: StatusCode, body: Vec<u8>) -> S3Result<S3Response<(StatusCode, Body)>> {
let mut headers = HeaderMap::new();
let content_type = HeaderValue::from_str(JSON_CONTENT_TYPE)
.map_err(|err| S3Error::with_message(S3ErrorCode::InternalError, format!("invalid content type: {err}")))?;
headers.insert(CONTENT_TYPE, content_type);
Ok(S3Response::with_headers((StatusCode::OK, Body::from(body)), headers))
Ok(S3Response::with_headers((status, Body::from(body)), headers))
}
fn scanner_recovery_intent_error(err: rustfs_scanner::ScannerError) -> S3Error {
let message = err.to_string();
let code = if message.contains("action is unsupported")
|| message.contains("mode is unsupported")
|| message.contains("actor identity is invalid")
|| message.contains("intent id is invalid")
|| message.contains("idempotency key")
|| message.contains("requires")
{
S3ErrorCode::InvalidRequest
} else {
S3ErrorCode::InternalError
};
S3Error::with_message(code, message)
}
fn scanner_recovery_intent_record_response(
status_code: StatusCode,
status: &'static str,
record: rustfs_scanner::ScannerRecoveryIntentRecord,
) -> S3Result<S3Response<(StatusCode, Body)>> {
let response = ScannerRecoveryIntentResponse {
status,
action: record.action,
mode: record.mode,
intent_id: record.intent_id,
state: record.state,
};
let body = serde_json::to_vec(&response).map_err(|err| {
S3Error::with_message(
S3ErrorCode::InternalError,
format!("failed to encode scanner recovery intent response: {err}"),
)
})?;
json_response_with_status(status_code, body)
}
fn scanner_recovery_intent_accept_response(
result: rustfs_scanner::ScannerRecoveryIntentAcceptResult,
) -> S3Result<S3Response<(StatusCode, Body)>> {
match result {
rustfs_scanner::ScannerRecoveryIntentAcceptResult::Accepted { record } => {
scanner_recovery_intent_record_response(StatusCode::ACCEPTED, "accepted", record)
}
rustfs_scanner::ScannerRecoveryIntentAcceptResult::Replayed { record } => {
scanner_recovery_intent_record_response(StatusCode::ACCEPTED, "replayed", record)
}
rustfs_scanner::ScannerRecoveryIntentAcceptResult::Conflict { existing } => {
let response = ScannerRecoveryIntentConflictResponse {
status: "conflict",
intent_id: existing.intent_id,
state: existing.state,
};
let body = serde_json::to_vec(&response).map_err(|err| {
S3Error::with_message(
S3ErrorCode::InternalError,
format!("failed to encode scanner recovery intent conflict: {err}"),
)
})?;
json_response_with_status(StatusCode::CONFLICT, body)
}
}
}
pub struct ScannerStatusHandler {}
@@ -314,6 +407,8 @@ pub struct ScannerCycleStateResetHandler {}
pub struct ScannerUsageStateResetHandler {}
pub struct ScannerUsageStateRecoveryIntentStatusHandler {}
#[async_trait::async_trait]
impl Operation for ScannerCycleStateResetHandler {
async fn call(&self, mut req: S3Request<Body>, _params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
@@ -361,6 +456,24 @@ impl Operation for ScannerUsageStateResetHandler {
.ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "storage layer not initialized"))?;
let store = current_object_store_handle_for_context(Some(context.as_ref()))
.ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "storage layer not initialized"))?;
if reset.async_intent {
let idempotency_key = reset
.idempotency_key
.ok_or_else(|| S3Error::with_message(S3ErrorCode::InvalidRequest, "async reset requires idempotency_key"))?;
let request = rustfs_scanner::ScannerRecoveryIntentRequest {
action: rustfs_scanner::SCANNER_RECOVERY_INTENT_ACTION_USAGE_FULL_REBUILD.to_string(),
mode: reset.mode,
idempotency_key,
actor_sha256: rustfs_scanner::scanner_recovery_actor_sha256(&_cred.access_key),
};
let accepted = rustfs_scanner::accept_scanner_usage_recovery_intent(store, request)
.await
.map_err(scanner_recovery_intent_error)?;
return scanner_recovery_intent_accept_response(accepted);
}
if reset.idempotency_key.is_some() {
return Err(S3Error::with_message(S3ErrorCode::InvalidRequest, "idempotency_key requires async reset"));
}
let result = supervise_admin_mutation("scanner usage state reset", async move {
rustfs_scanner::scanner::reset_scanner_usage_state_for_full_rebuild(CancellationToken::new(), store)
.await
@@ -377,6 +490,27 @@ impl Operation for ScannerUsageStateResetHandler {
}
}
#[async_trait::async_trait]
impl Operation for ScannerUsageStateRecoveryIntentStatusHandler {
async fn call(&self, req: S3Request<Body>, params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
let _cred = validate_scanner_reset_request(&req).await?;
let intent_id = params.get("intent_id").unwrap_or("");
let context = app_context_from_req(&req)
.ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "storage layer not initialized"))?;
let store = current_object_store_handle_for_context(Some(context.as_ref()))
.ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "storage layer not initialized"))?;
let record = rustfs_scanner::get_scanner_usage_recovery_intent(store, intent_id)
.await
.map_err(scanner_recovery_intent_error)?
.ok_or_else(|| {
let mut err = S3Error::with_message(S3ErrorCode::NoSuchKey, "scanner recovery intent not found");
err.set_status_code(StatusCode::NOT_FOUND);
err
})?;
scanner_recovery_intent_record_response(StatusCode::OK, "found", record)
}
}
#[async_trait::async_trait]
impl Operation for IlmExpiryStatusHandler {
async fn call(&self, req: S3Request<Body>, _params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
@@ -472,6 +606,8 @@ mod tests {
let full_rebuild: ScannerUsageStateResetRequest =
serde_json::from_str(r#"{"mode":"full-rebuild"}"#).expect("full rebuild must be accepted");
assert_eq!(full_rebuild.mode, "full-rebuild");
assert!(!full_rebuild.async_intent);
assert!(full_rebuild.idempotency_key.is_none());
let cycle_mode: ScannerUsageStateResetRequest =
serde_json::from_str(r#"{"mode":"full-rescan"}"#).expect("mode validation belongs to the handler");
assert_ne!(cycle_mode.mode, "full-rebuild");
@@ -481,6 +617,63 @@ mod tests {
);
}
#[test]
fn admin_usage_reset_accepts_explicit_async_intent_contract() {
let request: ScannerUsageStateResetRequest =
serde_json::from_str(r#"{"mode":"full-rebuild","async":true,"idempotency_key":"reset-key-0001"}"#)
.expect("async reset intent contract should parse");
assert_eq!(request.mode, "full-rebuild");
assert!(request.async_intent);
assert_eq!(request.idempotency_key.as_deref(), Some("reset-key-0001"));
}
#[test]
fn scanner_recovery_intent_response_uses_accepted_status() {
let record = rustfs_scanner::ScannerRecoveryIntentRecord {
schema_version: 1,
intent_id: "0".repeat(64),
action: rustfs_scanner::SCANNER_RECOVERY_INTENT_ACTION_USAGE_FULL_REBUILD.to_string(),
mode: "full-rebuild".to_string(),
state: "accepted".to_string(),
actor_sha256: "1".repeat(64),
idempotency_key_sha256: "2".repeat(64),
request_sha256: "3".repeat(64),
accepted_at_unix_secs: 7,
};
let response =
scanner_recovery_intent_accept_response(rustfs_scanner::ScannerRecoveryIntentAcceptResult::Accepted { record })
.expect("accepted intent response");
assert_eq!(response.output.0, StatusCode::ACCEPTED);
}
#[test]
fn scanner_recovery_intent_response_reports_conflict_status() {
let response = scanner_recovery_intent_accept_response(rustfs_scanner::ScannerRecoveryIntentAcceptResult::Conflict {
existing: rustfs_scanner::ScannerRecoveryIntentConflict {
intent_id: "0".repeat(64),
state: "accepted".to_string(),
},
})
.expect("conflict intent response");
assert_eq!(response.output.0, StatusCode::CONFLICT);
}
#[test]
fn scanner_recovery_intent_error_keeps_persisted_corruption_server_side() {
let invalid_actor =
scanner_recovery_intent_error(rustfs_scanner::ScannerError::Other("actor identity is invalid".to_string()));
assert_eq!(invalid_actor.code(), &S3ErrorCode::InvalidRequest);
let corrupt_record = scanner_recovery_intent_error(rustfs_scanner::ScannerError::Other(
"scanner recovery intent is invalid: expected value".to_string(),
));
assert_eq!(corrupt_record.code(), &S3ErrorCode::InternalError);
}
#[test]
fn scanner_disabled_reason_reports_startup_env_key() {
assert_eq!(scanner_disabled_reason(true), None);
+20
View File
@@ -489,6 +489,12 @@ pub const ADMIN_ROUTE_POLICY_SPECS: &[AdminRouteSpec] = &[
CONFIG_UPDATE,
RouteRiskLevel::High,
),
admin(
HttpMethod::Get,
"/rustfs/admin/v3/scanner/usage-state/recovery-intents/{intent_id}",
CONFIG_UPDATE,
RouteRiskLevel::High,
),
admin(
HttpMethod::Get,
"/rustfs/admin/v3/ilm/expiry/status",
@@ -2181,6 +2187,20 @@ mod tests {
assert_not_action(HttpMethod::Post, "/rustfs/admin/v3/scanner/usage-state/reset", SERVER_INFO);
}
#[test]
fn route_policy_requires_config_update_for_scanner_usage_recovery_intent() {
assert_action(
HttpMethod::Get,
"/rustfs/admin/v3/scanner/usage-state/recovery-intents/{intent_id}",
CONFIG_UPDATE,
);
assert_not_action(
HttpMethod::Get,
"/rustfs/admin/v3/scanner/usage-state/recovery-intents/{intent_id}",
SERVER_INFO,
);
}
#[test]
fn route_policy_uses_tier_actions_for_transition_routes() {
assert_action(HttpMethod::Get, "/rustfs/admin/v3/ilm/recovery/records", LIST_TIER);
@@ -290,6 +290,11 @@ fn expected_admin_route_matrix() -> Vec<RouteMatrixEntry> {
admin_route(Method::GET, "/v3/scanner/status"),
admin_route(Method::POST, "/v3/scanner/cycle-state/reset"),
admin_route(Method::POST, "/v3/scanner/usage-state/reset"),
admin_route_sample(
Method::GET,
"/v3/scanner/usage-state/recovery-intents/{intent_id}",
"/v3/scanner/usage-state/recovery-intents/aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa",
),
admin_route(Method::GET, "/v3/audit/target/list"),
admin_route_sample(
Method::PUT,
@@ -945,6 +950,11 @@ fn test_register_routes_cover_representative_admin_paths() {
assert_route(&router, Method::GET, &admin_path("/v3/scanner/status"));
assert_route(&router, Method::POST, &admin_path("/v3/scanner/cycle-state/reset"));
assert_route(&router, Method::POST, &admin_path("/v3/scanner/usage-state/reset"));
assert_route(
&router,
Method::GET,
&admin_path("/v3/scanner/usage-state/recovery-intents/aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"),
);
assert_route(&router, Method::GET, &admin_path("/v3/ilm/expiry/status"));
assert_route(&router, Method::GET, &admin_path("/v3/ilm/recovery/records"));
assert_route(
@@ -1458,6 +1468,12 @@ fn test_admin_alias_paths_match_existing_admin_routes() {
(Method::GET, compat_admin_alias_path("/v3/scanner/status")),
(Method::POST, compat_admin_alias_path("/v3/scanner/cycle-state/reset")),
(Method::POST, compat_admin_alias_path("/v3/scanner/usage-state/reset")),
(
Method::GET,
compat_admin_alias_path(
"/v3/scanner/usage-state/recovery-intents/aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa",
),
),
(Method::GET, compat_admin_alias_path("/v3/ilm/expiry/status")),
(Method::PUT, compat_admin_alias_path("/v3/on-demand-migration/b")),
(Method::GET, compat_admin_alias_path("/v3/on-demand-migration/b")),