fix(storage): harden ODM and scanner publication (#7187)

* fix(storage): harden ODM and scanner publication

* fix(app): simplify absent SSE configuration matching

* test(heal): settle PUT rename tails before disk-wipe fixtures

* fix(ecstore): remove duplicate local rename implementation

Keep the canonical commit module after concurrent storage changes merged.
The control-write and rollback changes are already present there.

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

* fix(ci): satisfy new clippy lints

* style(scanner): order merged test imports

* fix(scanner): invalidate bucket work after namespace completion

* fix(scanner): fence cached snapshots by scan execution

---------

Co-authored-by: houseme <housemecn@gmail.com>
Co-authored-by: heihutu <heihutu@gmail.com>
Co-authored-by: zhi22915 <qiuzgang@gmail.com>
This commit is contained in:
Zhengchao An
2026-09-05 21:47:12 +08:00
committed by GitHub
parent 447f3c704b
commit e2a921bc16
35 changed files with 2252 additions and 274 deletions
+3
View File
@@ -551,6 +551,9 @@ mod tests {
#[test]
fn a_plain_local_token_is_passed_through_and_a_tampered_one_is_rejected() {
let json_key = r#"{"t":"odm-list","v":1,"local_done":true}"#;
assert!(decode_list_cursor(Some(json_key)).expect("valid local key").is_none());
assert!(matches!(local_cursor(Some(json_key), None), LocalListCursor::Token(Some(local)) if local == json_key));
assert!(
decode_list_cursor(Some("photos/a.jpg"))
.expect("plain markers decode")
+83 -5
View File
@@ -4615,6 +4615,7 @@ fn odm_inline_client_body(primary: TeePrimary) -> StreamingBlob {
async fn odm_get_passthrough<S: OdmGetSource>(
state: &Arc<BucketOdmState>,
source: &S,
headers: &HeaderMap,
key: &str,
range: Option<&HTTPRangeSpec>,
backfill: Option<PullReason>,
@@ -4623,6 +4624,9 @@ async fn odm_get_passthrough<S: OdmGetSource>(
Ok(get) => get,
Err(err) => return OdmGetReply::Error(odm_get_source_failure(state, &err)),
};
if let Err(err) = odm_check_source_preconditions(headers, &get.head) {
return OdmGetReply::Error(err);
}
let content_length = match odm_content_length(get.head.size) {
Ok(length) => length,
Err(err) => {
@@ -4648,6 +4652,7 @@ async fn odm_get_passthrough<S: OdmGetSource>(
async fn odm_get_inline<S: OdmGetSource>(
state: &Arc<BucketOdmState>,
source: &S,
headers: &HeaderMap,
key: &str,
leader: PullLeader,
request_context: Option<request_context::RequestContext>,
@@ -4676,6 +4681,12 @@ async fn odm_get_inline<S: OdmGetSource>(
body,
content_range,
} = get;
// HEAD and GET can observe different source versions. Validate the
// representation whose body will actually be returned and persisted.
if let Err(err) = odm_check_source_preconditions(headers, &head) {
leader.complete(Err(PullError::canceled("source GET did not satisfy request preconditions")));
return OdmGetReply::Error(err);
}
// The object outgrew the inline budget between HEAD and GET: followers
// stream through on their own and the background pull stores it.
if head.size > policy.inline_max_bytes {
@@ -4758,19 +4769,19 @@ pub(super) async fn odm_get_from_source<S: OdmGetSource>(
let policy = &state.config().policy;
if let Some(range) = range {
let backfill = (policy.range_get == RangeGetPolicy::ServeAndBackfill).then_some(PullReason::RangeGet);
return odm_get_passthrough(state, source, key, Some(range), backfill).await;
return odm_get_passthrough(state, source, headers, key, Some(range), backfill).await;
}
if head.size > policy.inline_max_bytes {
return odm_get_passthrough(state, source, key, None, Some(PullReason::LargeObject)).await;
return odm_get_passthrough(state, source, headers, key, None, Some(PullReason::LargeObject)).await;
}
let slot = match state.acquire_pull_slot(key).await {
Ok(slot) => slot,
// The bucket state was torn down under this request: serve it
// without queueing anything on the old state.
Err(_) => return odm_get_passthrough(state, source, key, None, None).await,
Err(_) => return odm_get_passthrough(state, source, headers, key, None, None).await,
};
match slot {
PullSlot::Leader(leader) => odm_get_inline(state, source, key, leader, request_context).await,
PullSlot::Leader(leader) => odm_get_inline(state, source, headers, key, leader, request_context).await,
PullSlot::Follower(follower) => {
let first_byte = Duration::from_millis(policy.source_timeout.first_byte_ms);
match tokio::time::timeout(first_byte, follower.wait()).await {
@@ -4778,7 +4789,7 @@ pub(super) async fn odm_get_from_source<S: OdmGetSource>(
stats.record_request(OdmOp::Get, OdmOutcome::SourceHit);
OdmGetReply::RetryLocal
}
Ok(Err(_)) | Err(_) => odm_get_passthrough(state, source, key, None, None).await,
Ok(Err(_)) | Err(_) => odm_get_passthrough(state, source, headers, key, None, None).await,
}
}
}
@@ -5296,6 +5307,73 @@ mod on_demand_migration_tests {
assert!(rt.write_back.puts().is_empty());
}
#[tokio::test]
async fn odm_get_rechecks_conditions_against_the_get_representation() {
for inline_max_bytes in [0, 1024] {
for range in [
None,
Some(HTTPRangeSpec {
is_suffix_length: false,
start: 0,
end: 2,
}),
] {
let rt = runtime(
"changed-source",
PolicyConfig {
inline_max_bytes,
..Default::default()
},
)
.await;
let state = rt.state("changed-source");
let before = source_head(b"before");
let after = source_head(b"after!");
let source = ScriptedSource::new(vec![Ok(before.clone())], vec![Ok((after, b"after!".to_vec(), None))]);
let mut headers = HeaderMap::new();
headers.insert(
http::header::IF_MATCH,
HeaderValue::from_str(&format!("\"{}\"", before.etag.expect("etag"))).expect("header"),
);
let error = failed(odm_get_from_source(&state, &source, &headers, KEY, range.as_ref(), None).await);
assert_eq!(error.code(), &S3ErrorCode::PreconditionFailed);
assert_eq!(source.get_calls(), 1);
assert_eq!(state.inflight_keys(), 0);
assert!(rt.write_back.puts().is_empty(), "a failed condition must not start write-back");
}
}
}
#[tokio::test]
async fn odm_get_missing_validators_cannot_bypass_a_condition() {
for inline_max_bytes in [0, 1024] {
let rt = runtime(
"missing-validator",
PolicyConfig {
inline_max_bytes,
..Default::default()
},
)
.await;
let state = rt.state("missing-validator");
let before = source_head(b"before");
let after = SourceHead {
size: 6,
..Default::default()
};
let source = ScriptedSource::new(vec![Ok(before.clone())], vec![Ok((after, b"after!".to_vec(), None))]);
let mut headers = HeaderMap::new();
headers.insert(
http::header::IF_MATCH,
HeaderValue::from_str(&format!("\"{}\"", before.etag.expect("etag"))).expect("header"),
);
let error = failed(odm_get_from_source(&state, &source, &headers, KEY, None, None).await);
assert_eq!(error.status_code(), Some(StatusCode::FAILED_DEPENDENCY));
assert_eq!(error.message(), Some("missing_source_validator"));
assert!(rt.write_back.puts().is_empty());
}
}
#[tokio::test]
async fn odm_get_source_not_found_is_404_and_negative_cached() {
let rt = runtime("n", PolicyConfig::default()).await;
+17 -2
View File
@@ -57,6 +57,9 @@ pub(crate) struct InternalPutContext {
pub(crate) expected_md5_hex: Option<String>,
/// ETag to store instead of the computed one.
pub(crate) preserve_etag: Option<String>,
/// Reject an existing current object under the storage commit lock.
pub(crate) if_absent: bool,
pub(crate) preserve_delete_marker: bool,
pub(crate) content_headers: HashMap<String, String>,
pub(crate) user_metadata: HashMap<String, String>,
pub(crate) tags: Option<String>,
@@ -240,6 +243,8 @@ impl DefaultObjectUsecase {
size,
expected_md5_hex,
preserve_etag,
if_absent,
preserve_delete_marker,
content_headers,
user_metadata,
tags,
@@ -252,7 +257,10 @@ impl DefaultObjectUsecase {
};
let size = i64::try_from(size).map_err(|_| ApiError::invalid_request("internal put size exceeds the supported range"))?;
let headers = internal_put_headers(&content_headers)?;
let mut headers = internal_put_headers(&content_headers)?;
if if_absent {
headers.insert(http::header::IF_NONE_MATCH, HeaderValue::from_static("*"));
}
validate_internal_write_target(&key, &bucket, &headers).await?;
remove_source_replication_bookkeeping(&mut internal_metadata);
@@ -287,6 +295,7 @@ impl DefaultObjectUsecase {
origin: PutObjectOrigin::Internal {
principal_id,
emit_events,
preserve_delete_marker,
},
};
let committed = self
@@ -527,10 +536,14 @@ impl DefaultObjectUsecase {
.map_err(api_error_from_s3)?;
let store = self.object_store().ok_or_else(not_initialized)?;
let headers = HeaderMap::new();
let mut headers = HeaderMap::new();
if ctx.if_absent {
headers.insert(http::header::IF_NONE_MATCH, HeaderValue::from_static("*"));
}
let mut opts =
get_complete_multipart_upload_opts_with_replication_authorization(&headers, false).map_err(ApiError::from)?;
opts.preserve_etag = ctx.preserve_etag.clone();
opts.preserve_delete_marker = ctx.preserve_delete_marker;
let versioned = BucketVersioningSys::prefix_enabled(&bucket, &key).await;
opts.versioned = versioned;
opts.version_suspended = BucketVersioningSys::prefix_suspended(&bucket, &key).await;
@@ -747,6 +760,8 @@ mod tests {
size: Some(body.len() as u64),
expected_md5_hex: Some(md5_hex(body)),
preserve_etag: None,
if_absent: false,
preserve_delete_marker: false,
content_headers: HashMap::from([
("Content-Type".to_string(), "text/plain".to_string()),
("Cache-Control".to_string(), "max-age=60".to_string()),
@@ -66,6 +66,15 @@ impl OnDemandMigrationWriteBack {
.object_store()
.ok_or_else(|| WriteBackError::Local("object store is not initialized".to_string()))
}
fn require_atomic_write_back(&self) -> Result<(), WriteBackError> {
if !self.store()?.supports_atomic_create_only_write_back() {
return Err(WriteBackError::Unsupported(
"write-back requires namespace locking and exactly one pool with one erasure set".to_string(),
));
}
Ok(())
}
}
fn rfc3339(time: OffsetDateTime) -> String {
@@ -161,6 +170,8 @@ pub(super) async fn write_back_context(request: &WriteBackRequest, single_part:
size: Some(head.size),
expected_md5_hex: single_part.then(|| expected_md5_hex(head)).flatten(),
preserve_etag,
if_absent: true,
preserve_delete_marker: request.respect_delete_marker,
content_headers: content_headers(head),
user_metadata: head.user_metadata.clone(),
tags: request.tags.as_ref().and_then(encode_tags),
@@ -207,6 +218,7 @@ impl OdmWriteBack for OnDemandMigrationWriteBack {
}
async fn put_object(&self, request: &WriteBackRequest, body: WriteBackBody) -> Result<WriteBackOutcome, WriteBackError> {
self.require_atomic_write_back()?;
let ctx = write_back_context(request, true).await;
self.usecase()
.internal_put_object(ctx, body)
@@ -216,6 +228,7 @@ impl OdmWriteBack for OnDemandMigrationWriteBack {
}
async fn create_multipart_upload(&self, request: &WriteBackRequest) -> Result<String, WriteBackError> {
self.require_atomic_write_back()?;
let ctx = write_back_context(request, false).await;
self.usecase()
.internal_create_multipart_upload(&ctx)
@@ -249,6 +262,7 @@ impl OdmWriteBack for OnDemandMigrationWriteBack {
upload_id: &str,
parts: Vec<WriteBackPart>,
) -> Result<WriteBackOutcome, WriteBackError> {
self.require_atomic_write_back()?;
let ctx = write_back_context(request, false).await;
let parts = parts
.into_iter()
@@ -334,6 +348,7 @@ mod tests {
pulled_at: OffsetDateTime::from_unix_timestamp(1_756_800_000).expect("valid timestamp"),
preserve_etag: true,
emit_events: true,
respect_delete_marker: true,
tags: Some(HashMap::from([
("team".to_string(), "storage".to_string()),
("env".to_string(), "prod".to_string()),
@@ -486,6 +501,33 @@ mod tests {
assert!(!local.delete_marker);
}
#[tokio::test]
#[serial_test::serial]
async fn write_back_rejects_unsupported_topology_before_any_mutation() {
let (_dir, _paths, store) = crate::app::gating_test_env::isolated_multi_pool_ecstore().await;
crate::app::runtime_sources::install_test_app_context(Arc::clone(&store)).await;
let bucket = "odm-unsupported";
store
.make_bucket(bucket, &MakeBucketOptions::default())
.await
.expect("bucket");
let write_back = OnDemandMigrationWriteBack::new();
let req = request(bucket, "key", source_head(b"source"));
assert!(matches!(
write_back.put_object(&req, body_stream(b"source")).await,
Err(WriteBackError::Unsupported(_))
));
assert!(matches!(
write_back.create_multipart_upload(&req).await,
Err(WriteBackError::Unsupported(_))
));
assert!(matches!(
write_back.complete_multipart_upload(&req, "no-session", Vec::new()).await,
Err(WriteBackError::Unsupported(_))
));
assert_nothing_left(&store, bucket, "key").await;
}
#[tokio::test]
#[serial_test::serial]
async fn write_back_integrity_failure_leaves_nothing_behind() {
@@ -503,6 +545,133 @@ mod tests {
assert_nothing_left(&store, &bucket, "wrong.bin").await;
}
#[tokio::test]
#[serial_test::serial]
async fn write_back_commit_does_not_overwrite_a_concurrent_client_put() {
use crate::app::storage_api::test::set_disk::{PutObjectCommitBarrier, PutObjectCommitPause};
for versioned in [false, true] {
let (store, bucket) = write_back_test_bucket("odm-wb-race", versioned).await;
let source = b"old source bytes";
let client = b"new client bytes";
let req = request(&bucket, "race", source_head(source));
let client_req = request(&bucket, "race", source_head(client));
let mut client_ctx = write_back_context(&client_req, true).await;
client_ctx.if_absent = false;
let client_after = PutObjectCommitBarrier::install(&bucket, "race", PutObjectCommitPause::AfterNamespace);
let client_put = tokio::spawn(async move {
DefaultObjectUsecase::from_global()
.internal_put_object(client_ctx, body_stream(client))
.await
});
client_after.wait_until_paused().await;
let source_before = PutObjectCommitBarrier::install(&bucket, "race", PutObjectCommitPause::BeforeNamespace);
let write_back = OnDemandMigrationWriteBack::new();
let (result, ()) = tokio::join!(write_back.put_object(&req, body_stream(source)), async {
source_before.wait_until_paused().await;
drop(source_before);
drop(client_after);
});
let committed = client_put.await.expect("client task").expect("ordinary client write wins");
assert!(
matches!(result, Err(WriteBackError::Local(ref error)) if error.contains("PreconditionFailed")),
"{result:?}"
);
let stored = stored_object(&store, &bucket, "race").await;
assert_eq!(stored.etag, committed.etag);
assert_eq!(stored.version_id, committed.version_id);
assert_eq!(committed.version_id.is_some(), versioned);
assert_eq!(raw_object_bytes(&store, &bucket, "race").await, client);
}
}
#[tokio::test]
#[serial_test::serial]
async fn write_back_multipart_completion_preserves_a_client_put_after_staging() {
let (store, bucket) = write_back_test_bucket("odm-mpu-race", false).await;
let write_back = OnDemandMigrationWriteBack::new();
let req = request(&bucket, "race", source_head(b"source"));
let upload_id = write_back.create_multipart_upload(&req).await.expect("create");
let part = write_back
.upload_part(&req, &upload_id, 1, 6, body_stream(b"source"))
.await
.expect("stage");
let mut client_ctx = write_back_context(&request(&bucket, "race", source_head(b"client")), true).await;
client_ctx.if_absent = false;
let committed = DefaultObjectUsecase::from_global()
.internal_put_object(client_ctx, body_stream(b"client"))
.await
.expect("client put after staging");
let result = write_back.complete_multipart_upload(&req, &upload_id, vec![part]).await;
assert!(
matches!(result, Err(WriteBackError::Local(ref error)) if error.contains("PreconditionFailed")),
"{result:?}"
);
write_back
.abort_multipart_upload(&bucket, "race", &upload_id)
.await
.expect("abort rejected upload");
let stored = stored_object(&store, &bucket, "race").await;
assert_eq!(stored.etag, committed.etag);
assert_eq!(stored.version_id, committed.version_id);
assert_eq!(raw_object_bytes(&store, &bucket, "race").await, b"client");
}
#[tokio::test]
#[serial_test::serial]
async fn write_back_preserves_delete_markers_unless_policy_allows_revival() {
for multipart in [false, true] {
let (store, bucket) = write_back_test_bucket("odm-wb-tombstone", true).await;
let write_back = OnDemandMigrationWriteBack::new();
let mut req = request(&bucket, "deleted", source_head(b"source"));
let staged = if multipart {
let id = write_back.create_multipart_upload(&req).await.expect("create");
let part = write_back
.upload_part(&req, &id, 1, 6, body_stream(b"source"))
.await
.expect("part");
Some((id, part))
} else {
None
};
store
.delete_object(
&bucket,
"deleted",
ObjectOptions {
versioned: true,
..Default::default()
},
)
.await
.expect("delete marker");
let marker = stored_object(&store, &bucket, "deleted").await;
assert!(marker.delete_marker);
let rejected = if let Some((id, part)) = staged {
let result = write_back.complete_multipart_upload(&req, &id, vec![part]).await;
write_back
.abort_multipart_upload(&bucket, "deleted", &id)
.await
.expect("abort");
result
} else {
write_back.put_object(&req, body_stream(b"source")).await
};
assert!(
matches!(rejected, Err(WriteBackError::Local(ref error)) if error.contains("PreconditionFailed")),
"{rejected:?}"
);
let retained = stored_object(&store, &bucket, "deleted").await;
assert!(retained.delete_marker);
assert_eq!(retained.version_id, marker.version_id);
req.respect_delete_marker = false;
write_back
.put_object(&req, body_stream(b"source"))
.await
.expect("explicit revival policy");
assert!(!stored_object(&store, &bucket, "deleted").await.delete_marker);
}
}
#[tokio::test]
#[serial_test::serial]
async fn write_back_truncated_stream_leaves_nothing_behind() {
+12 -1
View File
@@ -949,7 +949,11 @@ pub(super) enum PutObjectOrigin<'a> {
/// request and no credential: managed-SSE authorization treats the write
/// as internal, and the creation event, when requested, names
/// `principal_id` instead of an access key.
Internal { principal_id: &'static str, emit_events: bool },
Internal {
principal_id: &'static str,
emit_events: bool,
preserve_delete_marker: bool,
},
}
impl PutObjectOrigin<'_> {
@@ -1603,6 +1607,12 @@ impl DefaultObjectUsecase {
if let Some(etag) = preserve_etag {
opts.preserve_etag = Some(etag);
}
if let PutObjectOrigin::Internal {
preserve_delete_marker, ..
} = &origin
{
opts.preserve_delete_marker = *preserve_delete_marker;
}
if let Some(quota_check) = quota_check.as_ref() {
apply_quota_admission(&mut opts, quota_check)?;
}
@@ -1769,6 +1779,7 @@ impl DefaultObjectUsecase {
PutObjectOrigin::Internal {
principal_id,
emit_events,
..
} => {
let principal_id = *principal_id;
let request_context = request_context::RequestContext::fallback();
+69 -2
View File
@@ -1014,11 +1014,44 @@ pub(crate) fn mark_on_demand_migration_list_local_only(headers: &mut HeaderMap)
/// forwarded to the source: a 304/412 answered by the source would be
/// indistinguishable from a source failure.
pub(crate) fn odm_check_source_preconditions(headers: &HeaderMap, head: &SourceHead) -> S3Result<()> {
let if_match = headers
.get(http::header::IF_MATCH)
.and_then(|value| value.to_str().ok())
.map(str::trim);
let if_none_match = headers
.get(http::header::IF_NONE_MATCH)
.and_then(|value| value.to_str().ok())
.map(str::trim);
let needs_etag = if_match.is_some_and(|value| value != "*") || if_none_match.is_some_and(|value| value != "*");
let needs_mtime = (!headers.contains_key(http::header::IF_MATCH) && headers.contains_key(http::header::IF_UNMODIFIED_SINCE))
|| (!headers.contains_key(http::header::IF_NONE_MATCH) && headers.contains_key(http::header::IF_MODIFIED_SINCE));
if (needs_etag && head.etag.is_none()) || (needs_mtime && head.last_modified.is_none()) {
return Err(odm_source_unavailable_error("missing_source_validator"));
}
let info = ObjectInfo {
etag: head.etag.clone(),
mod_time: head.last_modified.map(OffsetDateTime::from),
..Default::default()
};
// A successful source read establishes wildcard existence, but the
// remaining conditions must still run in their ordinary precedence.
if head.etag.is_none() && (if_match == Some("*") || if_none_match == Some("*")) {
let mut remaining = headers.clone();
if if_match == Some("*") {
remaining.remove(http::header::IF_MATCH);
remaining.remove(http::header::IF_UNMODIFIED_SINCE);
}
if if_none_match == Some("*") {
remaining.remove(http::header::IF_NONE_MATCH);
remaining.remove(http::header::IF_MODIFIED_SINCE);
}
check_preconditions(&remaining, &info)?;
return if if_none_match == Some("*") {
Err(S3Error::new(S3ErrorCode::NotModified))
} else {
Ok(())
};
}
check_preconditions(headers, &info)
}
@@ -2075,8 +2108,42 @@ mod on_demand_migration_tests {
.expect_err("modified since an earlier date is 412");
assert_eq!(err.code(), &S3ErrorCode::PreconditionFailed);
// A source without validators cannot fail a precondition.
let bare = SourceHead::default();
assert!(odm_check_source_preconditions(&headers_with(http::header::IF_MATCH, "\"other\""), &bare).is_ok());
let err =
odm_check_source_preconditions(&headers_with(http::header::IF_MATCH, "\"other\""), &bare).expect_err("missing ETag");
assert_eq!(err.status_code(), Some(http::StatusCode::FAILED_DEPENDENCY));
assert!(odm_check_source_preconditions(&headers_with(http::header::IF_MATCH, "*"), &bare).is_ok());
let err =
odm_check_source_preconditions(&headers_with(http::header::IF_NONE_MATCH, "*"), &bare).expect_err("source exists");
assert_eq!(err.code(), &S3ErrorCode::NotModified);
let dated = SourceHead {
last_modified: head.last_modified,
..Default::default()
};
assert!(odm_check_source_preconditions(&headers_with(http::header::IF_MATCH, "*"), &dated).is_ok());
let mut combined = headers_with(http::header::IF_NONE_MATCH, "*");
combined.insert(http::header::IF_MATCH, HeaderValue::from_static("\"other\""));
assert_eq!(
odm_check_source_preconditions(&combined, &dated)
.expect_err("specific ETag unavailable")
.status_code(),
Some(http::StatusCode::FAILED_DEPENDENCY)
);
combined.remove(http::header::IF_MATCH);
combined.insert(
http::header::IF_UNMODIFIED_SINCE,
HeaderValue::from_static("Wed, 21 Oct 2015 07:28:00 GMT"),
);
assert_eq!(
odm_check_source_preconditions(&combined, &dated)
.expect_err("unmodified-since fails before none-match")
.code(),
&S3ErrorCode::PreconditionFailed
);
for header in [http::header::IF_MODIFIED_SINCE, http::header::IF_UNMODIFIED_SINCE] {
let err = odm_check_source_preconditions(&headers_with(header, "Wed, 21 Oct 2015 07:28:00 GMT"), &bare)
.expect_err("missing timestamp");
assert_eq!(err.status_code(), Some(http::StatusCode::FAILED_DEPENDENCY));
}
}
}
+86 -1
View File
@@ -2131,9 +2131,9 @@ impl Node for NodeService {
)
.map_err(|err| Status::failed_precondition(err.to_string()))?;
}
let namespace_generation = store.scanner_namespace_mutation_generation();
let topology_digest = rustfs_scanner::scanner_topology_digest(store.as_ref());
let (data_movement_active, publication_blocked, movement_generation) = store.scanner_data_movement_activity().await;
let namespace_generation = store.scanner_namespace_mutation_generation();
let mut response = match request_protocol {
SCANNER_ACTIVITY_LEGACY_PROTOCOL_VERSION | SCANNER_ACTIVITY_PREVIOUS_PROTOCOL_VERSION => {
previous_scanner_activity_response(namespace_generation, topology_digest, data_movement_active)
@@ -6264,6 +6264,91 @@ mod tests {
assert_eq!(unavailable.code(), tonic::Code::Unavailable);
}
#[tokio::test]
async fn scanner_activity_samples_namespace_generation_after_waiting_for_movement_state() {
use crate::storage::storage_api::{ObjectOptions, PutObjReader, contract::object::ObjectIO as _};
let _ = rustfs_credentials::set_global_rpc_secret("scanner-activity-generation-test-secret".to_string());
let _ = rustfs_credentials::init_global_action_credentials(
Some("TESTROOTACCESSKEY".to_string()),
Some("TESTROOTSECRET123".to_string()),
);
let temp_dir = tempfile::tempdir().expect("scanner activity RPC test directory");
let env = rustfs_test_utils::TestECStoreEnv::builder()
.base_dir(temp_dir.path())
.build()
.await;
ObjectStore::new(Arc::clone(&env.ecstore))
.save_iam_config(serde_json::json!({"version": 1}), format!("{}/format.json", *IAM_CONFIG_PREFIX))
.await
.expect("seed IAM format");
let iam = rustfs_iam::build_iam_sys(Arc::clone(&env.ecstore))
.await
.expect("build isolated IAM");
let context = Arc::new(crate::runtime_sources::AppContext::with_default_interfaces(
Arc::clone(&env.ecstore),
iam,
Arc::new(KmsServiceManager::new()),
));
let service = make_server_for_context(Some(context));
let bucket = "scanner-activity-generation";
env.make_bucket(bucket, false).await;
let generation_before = env.ecstore.scanner_namespace_mutation_generation();
let mut request = Request::new(ScannerActivityRequest {
challenge: vec![7; 16].into(),
protocol_version: rustfs_scanner::SCANNER_ACTIVITY_PROTOCOL_VERSION,
acknowledge_instance_id: String::new(),
acknowledge_dirty_usage_generation: 0,
});
let canonical = rustfs_protos::canonical_scanner_activity_request_body(request.get_ref())
.expect("scanner activity request should encode");
set_tonic_canonical_body_digest(&mut request, &canonical).expect("digest metadata should encode");
mark_v2_authenticated(&mut request);
let pool_meta = env.ecstore.pool_meta.write().await;
drop(
env.ecstore
.decommission_cancelers
.try_write()
.expect("movement snapshot should not hold the cancelers before the RPC"),
);
let mut activity = Box::pin(tokio::task::unconstrained(service.scanner_activity(request)));
assert!(futures::poll!(activity.as_mut()).is_pending());
assert!(
env.ecstore.decommission_cancelers.try_write().is_err(),
"the RPC must hold the cancelers read guard while waiting for pool metadata"
);
// Select the existing set directly: ECStore pool selection reads the lock held by this test.
let mut reader = PutObjReader::from_vec(b"namespace changed during activity probe".to_vec());
tokio::time::timeout(
Duration::from_secs(30),
env.ecstore.pools[0].disk_set[0].put_object(
bucket,
"object",
&mut reader,
&ObjectOptions {
no_lock: true,
..Default::default()
},
),
)
.await
.expect("the namespace mutation must not wait for the RPC's pool lock")
.expect("the namespace mutation must complete while the RPC waits");
let generation_after = env.ecstore.scanner_namespace_mutation_generation();
assert!(generation_after > generation_before);
drop(pool_meta);
let response = tokio::time::timeout(Duration::from_secs(30), activity)
.await
.expect("scanner activity RPC should resume after the pool lock is released")
.expect("authenticated scanner activity RPC should succeed")
.into_inner();
assert_eq!(response.namespace_generation, generation_after);
assert_eq!(response.publication_blocked, Some(false));
}
#[tokio::test]
async fn test_scanner_dirty_usage_snapshot_requires_body_bound_auth_and_signs_a_consistent_view() {
let _ = rustfs_credentials::set_global_rpc_secret("scanner-dirty-usage-snapshot-test-secret".to_string());