fix(storage): harden ODM and scanner publication

This commit is contained in:
overtrue
2026-09-05 14:44:01 +08:00
parent 123967e729
commit 5cb670d360
28 changed files with 1718 additions and 163 deletions
+6
View File
@@ -560,6 +560,12 @@ pub mod set_disk {
pub mod test_util {
pub use crate::bucket::quota::reservation::fail_next_quota_ledger_save_for_test;
pub use crate::set_disk::{MultipartCommitBarrier, MultipartCommitPause, PutObjectCommitBarrier, PutObjectCommitPause};
/// Keep a namespace commit pending until the returned owner is dropped.
#[must_use]
pub fn hold_namespace_commit(store: &crate::store::ECStore) -> impl Send + Sync {
store.ctx.begin_namespace_commit()
}
}
}
@@ -25,8 +25,8 @@
//! [`BACKFILL_SAVE_INTERVAL`], and at every page end, with an `If-Match`
//! compare-and-set so a concurrent cancel or takeover is never overwritten.
//! - The `continuation_token` only advances once every pull queued from the
//! page before it has reported back, so a crash re-lists at most one page
//! (already-present keys are then skipped, never re-pulled).
//! page before it has succeeded. After a failure it stays at that page,
//! so crash recovery cannot skip failed pulls (existing keys are skipped).
//! - The owner holds a lease of [`BACKFILL_LEASE`] renewed by every save. The
//! recovery loop ([`run_backfill_recovery_loop`]) scans the buckets this
//! node has an ODM state for every [`BACKFILL_RECOVERY_INTERVAL`] and takes
@@ -367,9 +367,8 @@ pub struct LocalBackfillObject {
pub source_etag: Option<String>,
}
/// Receiver of one queued pull's report; `None` when the pull was coalesced
/// into one already running.
pub type PullReport = Option<oneshot::Receiver<QueuedPullOutcome>>;
/// Shared report of a new or coalesced pull; absent only when not admitted.
pub type PullReport = Option<super::pull::QueuedPullReport>;
/// Everything the job needs from its bucket, so the loop can run against a
/// mock in unit tests. Production: [`BucketBackfillContext`].
@@ -1190,9 +1189,10 @@ impl Job {
}
async fn main_loop(&mut self) -> Result<(), Stop> {
let mut cursor = self.checkpoint.continuation_token.clone();
loop {
self.check_cancel()?;
let page = self.list_page().await?;
let page = self.list_page(cursor.as_deref()).await?;
for object in &page.objects {
self.check_cancel()?;
self.checkpoint.listed += 1;
@@ -1204,10 +1204,13 @@ impl Job {
self.drain_ready();
self.tick(false).await?;
}
// Only advance the cursor once every pull of this page reported
// back, so a takeover re-lists at most this page.
// A persisted cursor certifies successful work, not just listing
// progress. Keep it at the first failed page for crash recovery.
self.drain_all().await?;
self.checkpoint.continuation_token = page.next_continuation_token.clone();
cursor = page.next_continuation_token;
if self.checkpoint.failed == 0 {
self.checkpoint.continuation_token = cursor.clone();
}
self.tick(true).await?;
if !page.is_truncated {
return Ok(());
@@ -1222,7 +1225,7 @@ impl Job {
}
}
async fn list_page(&mut self) -> Result<SourcePage, Stop> {
async fn list_page(&mut self, cursor: Option<&str>) -> Result<SourcePage, Stop> {
let mut attempt = 0;
loop {
while !self.context.source_available() {
@@ -1230,7 +1233,7 @@ impl Job {
self.tick(false).await?;
}
let prefix = self.checkpoint.prefix.clone();
let token = self.checkpoint.continuation_token.clone();
let token = cursor.map(str::to_string);
match self
.context
.list_page(prefix.as_deref(), token.as_deref(), BACKFILL_LIST_PAGE_SIZE)
@@ -1304,9 +1307,10 @@ impl Job {
}
loop {
match self.context.enqueue(key) {
(EnqueueOutcome::Enqueued, report) => {
(EnqueueOutcome::Enqueued | EnqueueOutcome::Coalesced, report) => {
self.checkpoint.enqueued += 1;
if let Some(rx) = report {
let rx = report.ok_or(Stop::Unavailable)?;
{
let key = key.to_string();
self.outstanding.push(Box::pin(async move { (key, rx.await) }));
}
@@ -1321,11 +1325,6 @@ impl Job {
);
return Ok(());
}
(EnqueueOutcome::Coalesced, _) => {
// Someone else pulls it; its result is not ours to count.
self.checkpoint.enqueued += 1;
return Ok(());
}
(EnqueueOutcome::QueueFull, _) => {
// Wait, never drop: one completion frees a slot.
if self.outstanding.is_empty() {
@@ -1639,6 +1638,7 @@ mod tests {
queue_capacity: usize,
pending: Mutex<Vec<(String, oneshot::Sender<QueuedPullOutcome>)>>,
fail_keys: HashSet<String>,
coalesced: bool,
auto_complete: AtomicBool,
cancel: CancellationToken,
config_updated_at: Mutex<Option<OffsetDateTime>>,
@@ -1666,6 +1666,7 @@ mod tests {
queue_capacity: usize::MAX,
pending: Mutex::new(Vec::new()),
fail_keys: HashSet::new(),
coalesced: false,
auto_complete: AtomicBool::new(true),
cancel: CancellationToken::new(),
config_updated_at: Mutex::new(Some(ts(1_700_000_000))),
@@ -1745,7 +1746,12 @@ mod tests {
} else {
self.pending.lock().push((key.to_string(), tx));
}
(EnqueueOutcome::Enqueued, Some(rx))
let outcome = if self.coalesced {
EnqueueOutcome::Coalesced
} else {
EnqueueOutcome::Enqueued
};
(outcome, Some(futures::FutureExt::shared(rx)))
}
fn cancel_token(&self) -> CancellationToken {
@@ -1911,7 +1917,7 @@ mod tests {
#[tokio::test]
async fn failed_pulls_are_counted_hashed_and_finish_with_failures() {
let bucket = "backfill-failed";
let mut context = MockContext::new(5, 1000);
let mut context = MockContext::new(5, 2);
Arc::get_mut(&mut context)
.expect("unshared")
.fail_keys
@@ -1926,12 +1932,52 @@ mod tests {
.checkpoint;
assert_eq!(cp.state, BackfillState::CompletedWithFailures);
assert_eq!((cp.pulled, cp.failed), (4, 1));
assert_eq!(cp.continuation_token.as_deref(), Some("2"), "retain the first failed page for recovery");
assert_eq!(cp.failed_keys, vec![key_hash("k/00002")]);
let last = cp.last_error.expect("last error");
assert_eq!(last.class, "local_write");
assert_eq!(last.key_hash.as_deref(), Some(key_hash("k/00002").as_str()));
}
#[tokio::test]
async fn coalesced_pulls_block_the_checkpoint_and_report_failures() {
let bucket = "backfill-coalesced";
let mut context = MockContext::new(1, 1);
{
let ctx = Arc::get_mut(&mut context).expect("unshared");
ctx.coalesced = true;
ctx.auto_complete = AtomicBool::new(false);
ctx.fail_keys.insert("k/00000".to_string());
}
let (_dirs, store, runner) = runner_with("node-a", bucket, Arc::clone(&context)).await;
runner.start(bucket, BackfillRequest::default()).await.expect("start");
tokio::time::timeout(Duration::from_secs(10), async {
while context.pending.lock().is_empty() {
tokio::task::yield_now().await;
}
})
.await
.expect("job enqueued");
assert!(runner.is_running_locally(bucket), "coalescing is not completion");
let cp = read_checkpoint(&store, bucket)
.await
.expect("read")
.expect("checkpoint")
.checkpoint;
assert!(cp.state.is_active());
assert!(cp.continuation_token.is_none());
context.complete_pending();
runner.wait_until_idle(bucket).await;
let cp = read_checkpoint(&store, bucket)
.await
.expect("read")
.expect("checkpoint")
.checkpoint;
assert_eq!(cp.state, BackfillState::CompletedWithFailures);
assert_eq!((cp.enqueued, cp.pulled, cp.failed), (1, 0, 1));
assert_eq!(cp.failed_keys, vec![key_hash("k/00000")]);
}
#[tokio::test]
async fn listing_failure_marks_the_job_failed_with_the_error_class() {
let bucket = "backfill-list-error";
@@ -32,6 +32,9 @@ pub const LIST_THROUGH_TOKEN_VERSION: u32 = 1;
/// listing's own marker, so the decoder needs a positive signal before it
/// treats an opaque token as a merged one.
const LIST_THROUGH_TOKEN_TAG: &str = "odm-list";
// Object keys cannot contain NUL (bucket::utils::is_valid_object_prefix),
// so this framing cannot collide with a local key used as an opaque marker.
const LIST_THROUGH_TOKEN_PREFIX: &str = "\0odm-list:";
/// Pages fetched per side per request: the first page, plus at most one refill
/// when the first one was mostly consumed by the previous page. Two pages of
@@ -86,8 +89,7 @@ pub struct MergePick {
}
/// The continuation-token envelope. Opaque to clients: it is serialized as
/// JSON and then base64-encoded by the same helper that encodes a plain local
/// marker, so the wire shape is `base64(json)`.
/// framed JSON and then base64-encoded by the same helper as a local marker.
///
/// A `null` cursor with `done = false` means "list that side from the start";
/// `done = true` means the side is finished and must not be listed again.
@@ -129,7 +131,7 @@ impl ListThroughToken {
pub fn encode(&self) -> String {
// The envelope is built here from owned strings, so serialization
// cannot fail; the fallback keeps the signature infallible.
serde_json::to_string(self).unwrap_or_default()
format!("{LIST_THROUGH_TOKEN_PREFIX}{}", serde_json::to_string(self).unwrap_or_default())
}
}
@@ -153,21 +155,18 @@ pub enum ListThroughTokenError {
/// Classifies an already base64-decoded continuation token.
///
/// Only a JSON object carrying the envelope marker is read as a merged token;
/// Only a framed JSON object is read as a merged token;
/// anything else is a local marker, so a bucket that turns `list_through` off
/// keeps paginating with the tokens it handed out. A token that *is* an
/// envelope but was tampered with (unknown version, unknown field, truncated
/// JSON) is an error, never a silent fallback.
pub fn decode_continuation_token(decoded: &str) -> Result<ListThroughCursor, ListThroughTokenError> {
if !decoded.starts_with('{') {
return Ok(ListThroughCursor::Local(decoded.to_string()));
}
let Ok(value) = serde_json::from_str::<serde_json::Value>(decoded) else {
// Not JSON at all: an object key may legitimately start with '{'.
let Some(payload) = decoded.strip_prefix(LIST_THROUGH_TOKEN_PREFIX) else {
return Ok(ListThroughCursor::Local(decoded.to_string()));
};
let value = serde_json::from_str::<serde_json::Value>(payload).map_err(|_| ListThroughTokenError::Malformed)?;
if value.get("t").and_then(serde_json::Value::as_str) != Some(LIST_THROUGH_TOKEN_TAG) {
return Ok(ListThroughCursor::Local(decoded.to_string()));
return Err(ListThroughTokenError::Malformed);
}
match value.get("v").and_then(serde_json::Value::as_u64) {
Some(version) if version == u64::from(LIST_THROUGH_TOKEN_VERSION) => {}
@@ -718,14 +717,21 @@ mod tests {
assert_eq!(decode_continuation_token(&extra), Err(ListThroughTokenError::Malformed));
let truncated = &encoded[..encoded.len() - 3];
assert_eq!(decode_continuation_token(truncated), Ok(ListThroughCursor::Local(truncated.to_string())));
assert_eq!(decode_continuation_token(truncated), Err(ListThroughTokenError::Malformed));
let no_version = "{\"t\":\"odm-list\"}";
let no_version = "\0odm-list:{\"t\":\"odm-list\"}";
assert_eq!(decode_continuation_token(no_version), Err(ListThroughTokenError::Malformed));
}
#[test]
fn a_plain_local_marker_stays_local() {
for marker in [
r#"{"t":"odm-list","v":1}"#,
r#"{"t":"odm-list","v":2,"local_done":true}"#,
r#"{"t":"odm-list"}"#,
] {
assert_eq!(decode_continuation_token(marker), Ok(ListThroughCursor::Local(marker.to_string())));
}
assert_eq!(
decode_continuation_token("photos/2024/01.jpg"),
Ok(ListThroughCursor::Local("photos/2024/01.jpg".to_string()))
@@ -46,10 +46,10 @@ use super::stats::{PullFailureReason, PullPath};
use super::sys::{BucketOdmState, OnDemandMigrationSys, PullError, PullOutcome, PullSlot};
use async_trait::async_trait;
use bytes::Bytes;
use futures::{Stream, StreamExt};
use futures::{FutureExt, Stream, StreamExt, future::Shared};
use parking_lot::Mutex;
use rand::RngExt;
use std::collections::{HashMap, HashSet};
use std::collections::HashMap;
use std::fmt;
use std::io;
use std::pin::Pin;
@@ -133,6 +133,8 @@ pub enum QueuedPullOutcome {
Failed(PullError),
}
pub type QueuedPullReport = Shared<oneshot::Receiver<QueuedPullOutcome>>;
/// Result of [`PullQueue::enqueue`].
#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
pub enum EnqueueOutcome {
@@ -251,6 +253,7 @@ pub struct WriteBackRequest {
pub preserve_etag: bool,
/// `policy.emit_events`.
pub emit_events: bool,
pub respect_delete_marker: bool,
/// Source tags to copy (`policy.copy_tags`), `None` to skip.
pub tags: Option<HashMap<String, String>>,
}
@@ -266,6 +269,7 @@ impl WriteBackRequest {
pulled_at: OffsetDateTime::now_utc(),
preserve_etag: config.policy.preserve_etag,
emit_events: config.policy.emit_events,
respect_delete_marker: config.policy.respect_local_delete_marker,
tags,
}
}
@@ -830,7 +834,7 @@ pub struct PullQueue {
bucket: String,
tx: mpsc::Sender<PullJob>,
/// Keys queued or running; the job removes its key when it ends.
pending: Mutex<HashSet<String>>,
pending: Mutex<HashMap<String, QueuedPullReport>>,
capacity: usize,
cancel: CancellationToken,
stats: Arc<super::stats::OdmStats>,
@@ -869,7 +873,7 @@ impl PullQueue {
let queue = Arc::new(Self {
bucket: state.bucket().to_string(),
tx,
pending: Mutex::new(HashSet::new()),
pending: Mutex::new(HashMap::new()),
capacity,
cancel: state.cancel_token(),
stats: Arc::clone(state.stats()),
@@ -903,29 +907,24 @@ impl PullQueue {
self.enqueue_with_report(key, reason).0
}
/// [`Self::enqueue`] that also hands back the job's report channel when
/// a new job was queued (`Coalesced` pulls report to their first
/// requester only).
pub fn enqueue_with_report(
&self,
key: &str,
reason: PullReason,
) -> (EnqueueOutcome, Option<oneshot::Receiver<QueuedPullOutcome>>) {
/// [`Self::enqueue`] with a shared report, including for coalesced pulls.
pub fn enqueue_with_report(&self, key: &str, reason: PullReason) -> (EnqueueOutcome, Option<QueuedPullReport>) {
if self.cancel.is_cancelled() {
return (EnqueueOutcome::Unavailable, None);
}
let mut pending = self.pending.lock();
if pending.contains(key) {
return (EnqueueOutcome::Coalesced, None);
if let Some(report) = pending.get(key) {
return (EnqueueOutcome::Coalesced, Some(report.clone()));
}
let (report_tx, report_rx) = oneshot::channel();
let report_rx = report_rx.shared();
match self.tx.try_send(PullJob {
key: key.to_string(),
reason,
report: Some(report_tx),
}) {
Ok(()) => {
pending.insert(key.to_string());
pending.insert(key.to_string(), report_rx.clone());
(EnqueueOutcome::Enqueued, Some(report_rx))
}
Err(TrySendError::Full(_)) => {
@@ -1072,7 +1071,7 @@ impl BucketOdmState {
self: &Arc<Self>,
key: &str,
reason: PullReason,
) -> (EnqueueOutcome, Option<oneshot::Receiver<QueuedPullOutcome>>) {
) -> (EnqueueOutcome, Option<QueuedPullReport>) {
match self.pull_queue() {
Some(queue) => queue.enqueue_with_report(key, reason),
None => (EnqueueOutcome::Unavailable, None),
@@ -1399,13 +1398,21 @@ mod tests {
assert_eq!(queue.capacity(), 1024);
let mut outcomes = HashMap::new();
let mut shared_report = None;
for _ in 0..100 {
*outcomes.entry(queue.enqueue("a", PullReason::RangeGet)).or_insert(0) += 1;
let (outcome, report) = queue.enqueue_with_report("a", PullReason::RangeGet);
*outcomes.entry(outcome).or_insert(0) += 1;
shared_report = report;
}
assert_eq!(outcomes.get(&EnqueueOutcome::Enqueued), Some(&1));
assert_eq!(outcomes.get(&EnqueueOutcome::Coalesced), Some(&99));
assert_eq!(queue.pending_keys(), 1);
assert_eq!(
shared_report.expect("coalesced report").await,
Ok(QueuedPullOutcome::Stored { size: 1000 })
);
wait_until("first pull to finish", || queue.pending_keys() == 0).await;
assert_eq!(source.head_calls.load(Ordering::SeqCst), 1);
assert_eq!(source.get_calls.load(Ordering::SeqCst), 1);
@@ -1438,6 +1445,23 @@ mod tests {
assert_eq!(queue.enqueue("a", PullReason::RangeGet), EnqueueOutcome::Unavailable);
}
#[tokio::test]
async fn coalesced_enqueues_share_failure_reports() {
let sys = OnDemandMigrationSys::new();
let state = enabled_state(&sys, &config()).await;
let source = MockSource::with_object("missing", 1000, BodyKind::Bytes(body_bytes(1000)));
let queue = PullQueue::start(Arc::clone(&state), source, Arc::new(MockWriteBack::default()));
let (first, first_report) = queue.enqueue_with_report("absent", PullReason::RangeGet);
let (second, second_report) = queue.enqueue_with_report("absent", PullReason::Backfill);
assert_eq!(first, EnqueueOutcome::Enqueued);
assert_eq!(second, EnqueueOutcome::Coalesced);
let (first, second) = tokio::join!(first_report.expect("leader report"), second_report.expect("coalesced report"));
assert_eq!(first, second);
assert!(matches!(first, Ok(QueuedPullOutcome::Failed(_))));
sys.remove(BUCKET);
queue.wait_until_stopped().await;
}
#[tokio::test]
async fn queue_full_is_reported_and_cancel_drains_without_leaking_tasks() {
let sys = OnDemandMigrationSys::new();
@@ -1467,7 +1491,8 @@ mod tests {
wait_until("dispatcher to wait for a slot", || state.stats().queue_depth() == 1).await;
assert_eq!(queue.enqueue("c", PullReason::LargeObject), EnqueueOutcome::Enqueued);
assert_eq!(queue.enqueue("d", PullReason::LargeObject), EnqueueOutcome::QueueFull);
assert_eq!(queue.enqueue("c", PullReason::LargeObject), EnqueueOutcome::Coalesced);
let (coalesced, canceled_report) = queue.enqueue_with_report("c", PullReason::LargeObject);
assert_eq!(coalesced, EnqueueOutcome::Coalesced);
assert_eq!(queue.pending_keys(), 3);
assert_eq!(failures(&state).get("queue_full"), Some(&1));
assert!(!queue.is_stopped());
@@ -1477,6 +1502,12 @@ mod tests {
.await
.expect("dispatcher and in-flight job must exit after cancel");
assert!(queue.is_stopped());
assert!(
tokio::time::timeout(Duration::from_secs(5), canceled_report.expect("coalesced cancellation report"))
.await
.expect("cancellation closes the report")
.is_err()
);
assert_eq!(queue.pending_keys(), 0);
assert_eq!(state.inflight_keys(), 0);
assert_eq!(state.stats().inflight_pulls(), 0);
@@ -152,8 +152,8 @@ pub struct SourceClientSpec {
/// Wire requests one logical source call may cost. The pull pipeline and
/// the backfill job own the retry budget (`pull.rs` `PULL_MAX_RETRIES`,
/// `backfill.rs` `LIST_MAX_RETRIES`) and the breaker counts logical calls,
/// so ODM declares [`RemoteS3RetryPolicy::Disabled`] and keeps one counted
/// failure equal to one request against a struggling source.
/// so ODM declares [`RemoteS3RetryPolicy::Disabled`]. An ambiguous HEAD
/// 404 additionally probes the bucket before declaring a key absent.
pub retry: RemoteS3RetryPolicy,
/// Bytes per second the pull pipeline may consume from this source;
/// `None` means unlimited. Enforced by the consumer, not by this client.
@@ -258,7 +258,7 @@ const THROTTLE_CODES: &[&str] = &[
"TooManyRequests",
"RequestThrottled",
];
const NOT_FOUND_CODES: &[&str] = &["NoSuchKey", "NotFound", "NoSuchBucket", "NoSuchVersion"];
const NOT_FOUND_CODES: &[&str] = &["NoSuchKey"];
const ACCESS_DENIED_CODES: &[&str] = &[
"AccessDenied",
"InvalidAccessKeyId",
@@ -281,7 +281,6 @@ fn classify_status(status: u16, code: Option<&str>, message: String) -> SourceEr
}
}
match status {
404 => SourceError::NotFound,
401 | 403 => SourceError::AccessDenied,
429 | 503 => SourceError::Throttled,
500..=599 => SourceError::ServerError(status),
@@ -627,8 +626,7 @@ impl SourceClient {
}
/// `config` must come from [`SourceClientSpec::endpoint_spec`], which is
/// where the retry policy that keeps one logical call equal to one wire
/// request is declared.
/// where the policy disabling SDK-level retries is declared.
fn from_config_builder(config: aws_sdk_s3::config::Builder, endpoint: String, spec: &SourceClientSpec) -> Self {
let client = S3Client::from_conf(config.interceptor(SourceProxyMarkerInterceptor::new()).build());
Self {
@@ -749,15 +747,16 @@ impl SourceClient {
#[async_trait::async_trait]
impl SourceBackend for S3SourceBackend {
async fn head(&self, key: &str) -> Result<SourceHead, SourceError> {
let output = self
.client
.head_object()
.bucket(&self.bucket)
.key(key)
.send()
.await
.map_err(classify_sdk_error)?;
source_head_from_head_output(output)
match self.client.head_object().bucket(&self.bucket).key(key).send().await {
Ok(output) => source_head_from_head_output(output),
Err(err) if err.raw_response().is_some_and(|response| response.status().as_u16() == 404) => {
// HEAD has no error body: a missing bucket must not poison
// the per-key negative cache as though only the key was absent.
self.probe().await?;
Err(SourceError::NotFound)
}
Err(err) => Err(classify_sdk_error(err)),
}
}
/// Streams the object; `range` is passed through as an HTTP `Range`
@@ -809,8 +808,8 @@ impl SourceBackend for S3SourceBackend {
.contents
.unwrap_or_default()
.into_iter()
.filter_map(s3_source_object)
.collect();
.map(s3_source_object)
.collect::<Result<Vec<_>, _>>()?;
let common_prefixes = output
.common_prefixes
.unwrap_or_default()
@@ -849,14 +848,20 @@ impl SourceBackend for S3SourceBackend {
}
}
fn s3_source_object(object: SdkObject) -> Option<SourceObject> {
let key = object.key?;
fn s3_source_object(object: SdkObject) -> Result<SourceObject, SourceError> {
let key = object
.key
.ok_or_else(|| SourceError::Other("source listing object has no key".to_string()))?;
let size = object
.size
.and_then(|size| u64::try_from(size).ok())
.ok_or_else(|| SourceError::Other("source listing object has no valid size".to_string()))?;
let etag = normalize_etag(object.e_tag);
let is_multipart_etag = etag.as_deref().is_some_and(is_multipart_etag);
Some(SourceObject {
Ok(SourceObject {
key,
etag,
size: object.size.and_then(|size| u64::try_from(size).ok()).unwrap_or(0),
size,
last_modified: system_time(object.last_modified),
storage_class: object.storage_class.map(|class| class.as_str().to_string()),
is_multipart_etag,
@@ -1390,7 +1395,10 @@ mod tests {
#[tokio::test]
async fn source_error_classification_covers_every_class() {
let cases: Vec<(Scripted, &str, bool)> = vec![
(status(404, ""), "not_found", false),
(status(404, ""), "other", false),
(status(404, "<Error><Code>NoSuchKey</Code></Error>"), "not_found", false),
(status(404, "<Error><Code>NoSuchBucket</Code></Error>"), "other", false),
(status(404, "<Error><Code>NoSuchVersion</Code></Error>"), "other", false),
(status(403, ACCESS_DENIED_BODY), "access_denied", false),
(status(401, ""), "access_denied", false),
(status(429, ""), "throttled", true),
@@ -1413,14 +1421,35 @@ mod tests {
}
}
// HEAD carries no error body, so the classification must work from the
// status alone as well.
let (client, _) = scripted_client(&spec(None), vec![status(404, "")]).await;
let (client, requests) = scripted_client(&spec(None), vec![status(404, ""), status(200, "")]).await;
assert!(matches!(client.head_object("missing").await, Err(SourceError::NotFound)));
assert_eq!(recorded(&requests).len(), 2, "ambiguous HEAD 404 must check the bucket");
let (client, _) = scripted_client(&spec(None), vec![status(404, ""), status(404, "")]).await;
assert!(matches!(client.head_object("missing").await, Err(SourceError::Other(_))));
let (client, _) = scripted_client(&spec(None), vec![status(404, ""), status(403, "")]).await;
assert!(matches!(client.head_object("missing").await, Err(SourceError::AccessDenied)));
let (client, _) = scripted_client(&spec(None), vec![status(403, "")]).await;
assert!(matches!(client.head_object("secret").await, Err(SourceError::AccessDenied)));
}
#[test]
fn source_listing_rejects_missing_and_negative_sizes() {
for size in [None, Some(-1)] {
let object = SdkObject::builder().key("key").set_size(size).build();
assert!(matches!(s3_source_object(object), Err(SourceError::Other(_))));
}
assert!(matches!(
s3_source_object(SdkObject::builder().size(0).build()),
Err(SourceError::Other(_))
));
assert_eq!(
s3_source_object(SdkObject::builder().key("empty").size(0).build())
.expect("empty object")
.size,
0
);
}
#[tokio::test]
async fn source_client_debug_redacts_credentials() {
let (client, _) = scripted_client(&spec(Some("data/")), Vec::new()).await;
+3
View File
@@ -940,6 +940,9 @@ pub struct ObjectOptions {
pub preserve_etag: Option<String>,
pub metadata_chg: bool,
pub http_preconditions: Option<HTTPPreconditions>,
/// Internal create-only writes may also preserve an acknowledged deletion.
/// Evaluated with `http_preconditions` under the namespace commit lock.
pub preserve_delete_marker: bool,
pub delete_replication: Option<ReplicationState>,
pub delete_replication_config_snapshot: Option<Arc<DeleteReplicationConfigSnapshot>>,
+112
View File
@@ -78,6 +78,21 @@ pub(crate) struct ScannerPublicationLeaseEntry {
pub(crate) _operation_guard: OwnedRwLockReadGuard<()>,
}
pub(crate) struct NamespaceCommitGuard {
ctx: Arc<InstanceContext>,
counted: bool,
}
impl Drop for NamespaceCommitGuard {
fn drop(&mut self) {
if self.counted {
// Publish the new generation before a zero-pending publication probe.
self.ctx.advance_namespace_commit_generation();
self.ctx.namespace_commits.fetch_sub(1, Ordering::AcqRel);
}
}
}
/// Runtime state owned by a single `ECStore` instance.
///
/// This is intentionally minimal in the first migration slice; subsequent
@@ -209,9 +224,13 @@ pub struct InstanceContext {
/// Last storage-owned movement snapshot observed under the operation
/// gate. SetDisks cache writers fail closed until ECStore refreshes it.
scanner_publication_state: AtomicU8,
namespace_commits: AtomicU64,
namespace_commit_generation: AtomicU64,
/// Resolves object-encryption material at the application boundary.
object_encryption_resolver: OnceLock<Arc<dyn ObjectEncryptionResolver>>,
tier_delete_journal_recovery_stores: std::sync::Mutex<HashSet<Uuid>>,
#[cfg(test)]
suppress_tier_delete_journal_recovery: bool,
transition_transaction_recovery_stores: std::sync::Mutex<HashSet<Uuid>>,
tier_delete_journal_recovery_wakeup: tokio::sync::Notify,
}
@@ -256,8 +275,12 @@ impl InstanceContext {
data_movement_generation_exhausted: AtomicBool::new(false),
data_movement_generation_notify: Arc::new(Notify::new()),
scanner_publication_state: AtomicU8::new(SCANNER_PUBLICATION_STATE_UNKNOWN),
namespace_commits: AtomicU64::new(0),
namespace_commit_generation: AtomicU64::new(0),
object_encryption_resolver: OnceLock::new(),
tier_delete_journal_recovery_stores: std::sync::Mutex::new(HashSet::new()),
#[cfg(test)]
suppress_tier_delete_journal_recovery: false,
transition_transaction_recovery_stores: std::sync::Mutex::new(HashSet::new()),
tier_delete_journal_recovery_wakeup: tokio::sync::Notify::new(),
}
@@ -385,6 +408,36 @@ impl InstanceContext {
&& self.scanner_publication_state.load(Ordering::Acquire) == SCANNER_PUBLICATION_STATE_ALLOWED
}
pub(crate) fn begin_namespace_commit(self: &Arc<Self>) -> Arc<NamespaceCommitGuard> {
let counted = self
.namespace_commits
.fetch_update(Ordering::AcqRel, Ordering::Acquire, |count| count.checked_add(1))
.is_ok();
if counted {
self.advance_namespace_commit_generation();
} else {
self.namespace_commit_generation.store(u64::MAX, Ordering::Release);
}
Arc::new(NamespaceCommitGuard {
ctx: Arc::clone(self),
counted,
})
}
fn advance_namespace_commit_generation(&self) {
let _ = self
.namespace_commit_generation
.fetch_update(Ordering::AcqRel, Ordering::Acquire, |generation| Some(generation.saturating_add(1)));
}
pub(crate) fn namespace_commit_generation(&self) -> u64 {
self.namespace_commit_generation.load(Ordering::Acquire)
}
pub(crate) fn namespace_commits_pending(&self) -> bool {
self.namespace_commits.load(Ordering::Acquire) != 0 || self.namespace_commit_generation() == u64::MAX
}
pub(crate) fn set_scanner_publication_state(&self, blocked: bool) {
self.scanner_publication_state.store(
if blocked {
@@ -640,12 +693,21 @@ impl InstanceContext {
}
pub(crate) fn mark_tier_delete_journal_recovery_started(&self, store_id: Uuid) -> bool {
#[cfg(test)]
if self.suppress_tier_delete_journal_recovery {
return false;
}
self.tier_delete_journal_recovery_stores
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.insert(store_id)
}
#[cfg(test)]
pub(crate) fn suppress_tier_delete_journal_recovery_for_test(&mut self) {
self.suppress_tier_delete_journal_recovery = true;
}
pub(crate) fn mark_transition_transaction_recovery_started(&self, store_id: Uuid) -> bool {
self.transition_transaction_recovery_stores
.lock()
@@ -756,6 +818,50 @@ pub fn bootstrap_ctx() -> Arc<InstanceContext> {
mod tests {
use super::*;
#[test]
fn namespace_commit_guards_are_instance_local_and_count_until_last_owner() {
let first = Arc::new(InstanceContext::new());
let other = Arc::new(InstanceContext::new());
first.set_scanner_publication_state(false);
other.set_scanner_publication_state(false);
assert!(first.scanner_publication_state_allowed());
let one = first.begin_namespace_commit();
let shared_owner = Arc::clone(&one);
let two = first.begin_namespace_commit();
assert!(first.namespace_commits_pending());
assert!(first.scanner_publication_state_allowed(), "pending writes must not block scan admission");
assert_eq!(first.namespace_commit_generation(), 2);
assert!(!other.namespace_commits_pending());
assert_eq!(other.namespace_commit_generation(), 0);
assert!(other.scanner_publication_state_allowed());
drop(one);
assert_eq!(first.namespace_commit_generation(), 2);
drop(shared_owner);
assert!(first.namespace_commits_pending());
assert_eq!(first.namespace_commit_generation(), 3);
drop(two);
assert!(!first.namespace_commits_pending());
assert_eq!(first.namespace_commit_generation(), 4);
assert!(first.scanner_publication_state_allowed());
}
#[test]
fn namespace_commit_counter_exhaustion_keeps_publication_blocked() {
for (count, generation) in [(0, u64::MAX - 1), (u64::MAX, 0)] {
let ctx = Arc::new(InstanceContext::new());
ctx.set_scanner_publication_state(false);
ctx.namespace_commits.store(count, Ordering::Release);
ctx.namespace_commit_generation.store(generation, Ordering::Release);
let guard = ctx.begin_namespace_commit();
assert!(ctx.namespace_commits_pending());
assert_eq!(ctx.namespace_commit_generation(), u64::MAX);
drop(guard);
assert!(ctx.namespace_commits_pending());
assert_eq!(ctx.namespace_commit_generation(), u64::MAX);
assert_eq!(ctx.namespace_commits.load(Ordering::Acquire), count);
}
}
// The SetupType inputs must derive the exact (is_erasure,
// is_dist_erasure, is_erasure_sd) triples that the original three
// process-global erasure bools produced via update_erasure_type().
@@ -1073,6 +1179,12 @@ mod tests {
assert!(!ctx_a.mark_tier_delete_journal_recovery_started(store_a));
assert!(ctx_a.mark_tier_delete_journal_recovery_started(store_b));
assert!(ctx_b.mark_tier_delete_journal_recovery_started(store_a));
let mut manual_ctx = InstanceContext::new();
manual_ctx.suppress_tier_delete_journal_recovery_for_test();
assert!(!manual_ctx.mark_tier_delete_journal_recovery_started(store_a));
assert!(!manual_ctx.mark_tier_delete_journal_recovery_started(store_b));
assert!(ctx_b.mark_tier_delete_journal_recovery_started(store_b));
}
#[test]
@@ -3845,6 +3845,7 @@ pub(in crate::set_disk) struct RenameDataFenceOptions<'a> {
write_quorum: usize,
scanner_publication_lease_tokens: Option<&'a HashMap<String, Uuid>>,
scanner_publication_commit_scope: Option<crate::object_api::ScannerPublicationCommitScope>,
namespace_commit_guard: Option<Arc<crate::runtime::instance::NamespaceCommitGuard>>,
}
impl<'a> RenameDataFenceOptions<'a> {
@@ -3856,6 +3857,7 @@ impl<'a> RenameDataFenceOptions<'a> {
write_quorum,
scanner_publication_lease_tokens,
scanner_publication_commit_scope: None,
namespace_commit_guard: None,
}
}
@@ -3866,6 +3868,14 @@ impl<'a> RenameDataFenceOptions<'a> {
self.scanner_publication_commit_scope = scanner_publication_commit_scope;
self
}
pub(in crate::set_disk) fn with_namespace_commit_guard(
mut self,
namespace_commit_guard: Option<Arc<crate::runtime::instance::NamespaceCommitGuard>>,
) -> Self {
self.namespace_commit_guard = namespace_commit_guard;
self
}
}
#[allow(dead_code, reason = "asserted by this file's tests (backlog#1823)")]
@@ -4224,6 +4234,7 @@ impl SetDisks {
write_quorum,
scanner_publication_lease_tokens,
scanner_publication_commit_scope: _scanner_publication_commit_scope,
namespace_commit_guard,
} = fence_options;
if let Some(file_info) = disks
.iter()
@@ -4268,7 +4279,9 @@ impl SetDisks {
let dst_object = fanout_dst_object.clone();
let file_info = file_info.clone();
let successful_rename_completion_rank = successful_rename_completion_rank.clone();
let namespace_commit_guard = namespace_commit_guard.clone();
tasks.spawn(async move {
let _namespace_commit_guard = namespace_commit_guard;
let result = std::panic::AssertUnwindSafe(async move {
#[allow(clippy::let_unit_value)]
let _fanout_task_guard = Self::rename_fanout_task_guard(&dst_object);
@@ -4582,6 +4595,7 @@ impl SetDisks {
write_quorum,
scanner_publication_lease_tokens,
scanner_publication_commit_scope,
namespace_commit_guard,
} = fence_options;
if let Some(file_info) = disks
.iter()
@@ -4614,6 +4628,7 @@ impl SetDisks {
let fanout_dst_bucket = dst_bucket.clone();
let fanout_dst_object = dst_object.clone();
let fanout_publication_scope = scanner_publication_commit_scope.clone();
let fanout_namespace_commit_guard = namespace_commit_guard.clone();
// Keep one coordinator task so a cancelled caller cannot drop partially
// completed disk mutations. Per-disk futures stay ordered in `join_all`,
// preserving slot-indexed quorum and convergence accounting without a
@@ -4622,6 +4637,7 @@ impl SetDisks {
// Keep the storage-owned movement permit attached to the actual
// fan-out owner, even if the caller future is cancelled.
let _fanout_publication_scope = fanout_publication_scope;
let _namespace_commit_guard = fanout_namespace_commit_guard;
let successful_rename_completion_rank =
rustfs_io_metrics::put_stage_metrics_enabled().then(|| Arc::new(AtomicUsize::new(0)));
let futures = fanout_disks
@@ -4807,6 +4823,8 @@ impl SetDisks {
let dst_bucket = dst_bucket.clone();
let dst_object = dst_object.clone();
futures.push(tokio::spawn(async move {
#[cfg(test)]
rename_fanout_barrier::checkpoint(&dst_object, i, rename_fanout_barrier::PHASE_ROLLBACK).await;
disk.delete_version(
&dst_bucket,
&dst_object,
@@ -6565,9 +6583,9 @@ impl SetDisks {
match oi {
Ok(oi) => {
// Ordinary writes may proceed past a top-level delete marker;
// data movement must not replace an acknowledged deletion.
// data movement and guarded internal writes must preserve it.
if oi.delete_marker {
return opts.data_movement.then_some(StorageError::PreconditionFailed);
return (opts.data_movement || opts.preserve_delete_marker).then_some(StorageError::PreconditionFailed);
}
let if_none_match = http_preconditions.if_none_match_value().map(str::to_owned);
let if_match = http_preconditions.if_match_value().map(str::to_owned);
@@ -6960,6 +6978,7 @@ pub(crate) mod rename_fanout_barrier {
pub use super::rename_fanout_barrier_phase::{
CLEANUP as PHASE_CLEANUP, READ_VERSION as PHASE_READ_VERSION, RENAME as PHASE_RENAME,
};
pub const PHASE_ROLLBACK: &str = "rollback";
/// One armed barrier: the fan-out task matching `(disk_index, phase)` pauses.
struct Armed {
@@ -10534,9 +10553,35 @@ mod tests {
let mut file_infos = rename_commit_fileinfos(object, DISKS, "fresh-rollback-etag");
file_infos[3] = FileInfo::default();
SetDisks::rename_data(&disks, RUSTFS_META_TMP_BUCKET, "source", &file_infos, bucket, object, 4)
let ctx = Arc::new(crate::runtime::instance::InstanceContext::new());
ctx.set_scanner_publication_state(false);
let barrier = rename_fanout_barrier::arm(object, 0, rename_fanout_barrier::PHASE_ROLLBACK);
let rename = SetDisks::rename_data_owned_with_fence(
&disks,
(RUSTFS_META_TMP_BUCKET, "source"),
file_infos,
(bucket, object),
false,
RenameDataFenceOptions::new(4, None).with_namespace_commit_guard(Some(ctx.begin_namespace_commit())),
);
let control = async {
barrier.wait_until_paused().await;
assert!(ctx.namespace_commits_pending(), "rollback must retain namespace publication ownership");
assert!(ctx.scanner_publication_state_allowed(), "rollback must not disable namespace walks");
assert_eq!(ctx.namespace_commit_generation(), 1);
barrier.release();
};
let (result, ()) = tokio::time::timeout(BARRIER_PAUSE_GUARD, async { tokio::join!(rename, control) })
.await
.expect_err("three successful disks must fail a strict write quorum of four");
.expect("rename rollback must reach its barrier and finish after release");
assert_eq!(
result.err(),
Some(DiskError::ErasureWriteQuorum),
"three successful disks must fail a strict write quorum of four"
);
assert!(!ctx.namespace_commits_pending());
assert!(ctx.scanner_publication_state_allowed());
assert_eq!(ctx.namespace_commit_generation(), 2);
for (idx, dir) in dirs.iter().enumerate() {
let reopened = reopen_local_disk(dir).await;
+35 -19
View File
@@ -4051,6 +4051,7 @@ mod tests {
let _ = drain_global_dirty_scopes();
let rename_barrier = rename_fanout_barrier::arm(object, 0, rename_fanout_barrier::PHASE_RENAME);
let rename_tasks = rename_fanout_barrier::observe_tasks(object);
let complete_store = Arc::clone(&set_disks);
let mut complete = tokio::spawn(async move {
let mut opts = ObjectOptions::default();
@@ -4062,16 +4063,6 @@ mod tests {
tokio::time::timeout(Duration::from_secs(30), rename_barrier.wait_until_paused())
.await
.expect("multipart completion should pause one tail disk during rename");
assert!(
tokio::time::timeout(Duration::from_millis(100), &mut complete).await.is_err(),
"multipart completion must not publish success while a tail rename is still paused"
);
let initial = drain_global_dirty_scopes().into_iter().collect::<HashSet<_>>();
assert!(
initial.is_empty(),
"capacity must not be marked as committed before the full multipart rename finishes"
);
let abort_store = Arc::clone(&set_disks);
let abort = tokio::spawn(async move {
@@ -4080,21 +4071,46 @@ mod tests {
.await
});
signaling.wait_for_attempts(2).await;
assert!(!abort.is_finished(), "the in-flight completion must retain the multipart upload guard");
let retained_staging = futures::future::join_all(
disk_stores
.iter()
.map(|disk| disk.read_all(RUSTFS_META_MULTIPART_BUCKET, &staged_part)),
)
// A paused rename does not establish that the other disks reached quorum.
let retained_staging = tokio::time::timeout(Duration::from_secs(30), async {
loop {
let mut retained = 0;
for result in futures::future::join_all(
disk_stores
.iter()
.map(|disk| disk.read_all(RUSTFS_META_MULTIPART_BUCKET, &staged_part)),
)
.await
{
match result {
Ok(_) => retained += 1,
Err(DiskError::FileNotFound) => {}
Err(error) => panic!("staged rename source lookup failed: {error}"),
}
}
if retained <= 1 && rename_tasks.running() == 1 {
break retained;
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
})
.await
.into_iter()
.filter(|result| result.is_ok())
.count();
.expect("unpaused multipart renames should finish before the tail is released");
assert_eq!(
retained_staging, 1,
"only the paused tail disk should still retain the multipart rename source"
);
assert!(
tokio::time::timeout(Duration::from_millis(100), &mut complete).await.is_err(),
"multipart completion must not publish success while a tail rename is still paused"
);
let initial = drain_global_dirty_scopes().into_iter().collect::<HashSet<_>>();
assert!(
initial.is_empty(),
"capacity must not be marked as committed before the full multipart rename finishes"
);
assert!(!abort.is_finished(), "the in-flight completion must retain the multipart upload guard");
signaling.set_target(rustfs_lock::ObjectKey::new(bucket, object));
let object_attempt = signaling.attempts.load(Ordering::Acquire) + 1;
+4 -1
View File
@@ -4452,7 +4452,10 @@ impl SetDisks {
write_quorum,
commit_scanner_publication_lease_tokens.as_ref(),
)
.with_publication_scope(commit_scanner_publication_scope.clone()),
.with_publication_scope(commit_scanner_publication_scope.clone())
.with_namespace_commit_guard(
(!is_meta_bucketname(&commit_bucket)).then(|| commit_set.ctx.begin_namespace_commit()),
),
)
.await;
if let Some(scope) = commit_scanner_publication_scope.as_ref() {
+12 -2
View File
@@ -1059,6 +1059,7 @@ mod tests {
use crate::storage_api_contracts::{
bucket::{BucketOperations as _, BucketOptions, DeleteBucketOptions, MakeBucketOptions, SRBucketDeleteOp},
list::ListOperations as _,
namespace::NamespaceLocking as _,
object::{ObjectIO as _, ObjectOperations as _},
};
use crate::store::{ECStore, init_local_disks_with_instance_ctx};
@@ -1486,10 +1487,19 @@ mod tests {
.put_object(bucket, object, &mut reader, &ObjectOptions::default())
.await
.expect("object should be written");
let lock = ecstore.pools[0].disk_set[0]
.new_ns_lock(bucket, object)
.await
.expect("fixture namespace lock should be created");
drop(
lock.get_write_lock(Duration::from_secs(30))
.await
.expect("fixture rename tail should finish before checking its generation"),
);
assert_eq!(
ecstore.scanner_namespace_mutation_generation(),
generation_before_put.saturating_add(1),
"successful object creation should advance scanner namespace activity"
generation_before_put.saturating_add(3),
"successful object creation must observe the logical mutation and both fanout boundaries"
);
ecstore
.get_object_info(bucket, object, &ObjectOptions::default())
+525 -33
View File
@@ -787,6 +787,12 @@ impl ECStore {
pub fn single_pool(&self) -> bool {
self.pools.len() == 1
}
/// The set-local create-only check is atomic only when every object
/// mutation uses that same, enabled namespace lock domain.
pub fn supports_atomic_create_only_write_back(&self) -> bool {
!self.ctx.lock_manager().is_disabled() && self.pools.len() == 1 && self.pools[0].disk_set.len() == 1
}
}
#[cfg(test)]
@@ -2127,7 +2133,7 @@ mod tests {
.iter()
.map(|&drives_per_set| (1, drives_per_set))
.collect::<Vec<_>>();
build_isolated_test_store_with_layout(temp_dir, cmd_line, &pool_layouts, shutdown).await
build_isolated_test_store_with_layout(temp_dir, cmd_line, &pool_layouts, shutdown, None).await
}
async fn build_isolated_test_store_with_layout(
@@ -2135,6 +2141,7 @@ mod tests {
cmd_line: &str,
pool_layouts: &[(usize, usize)],
shutdown: CancellationToken,
instance_ctx: Option<Arc<crate::runtime::instance::InstanceContext>>,
) -> (
Arc<crate::runtime::instance::InstanceContext>,
Arc<crate::store::ECStore>,
@@ -2167,7 +2174,7 @@ mod tests {
let endpoint_pools = EndpointServerPools(pools);
crate::services::notification_sys::install_cross_pool_fence_fleet_proof_for_test();
let instance_ctx = Arc::new(crate::runtime::instance::InstanceContext::new());
let instance_ctx = instance_ctx.unwrap_or_else(|| Arc::new(crate::runtime::instance::InstanceContext::new()));
crate::store::init_local_disks_with_instance_ctx(&instance_ctx, endpoint_pools.clone())
.await
.expect("register local disks into the fresh context");
@@ -2535,6 +2542,348 @@ mod tests {
shutdown.cancel();
}
#[cfg(feature = "test-util")]
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
#[serial_test::serial(storage_class_env)]
async fn early_ack_put_tails_block_scanner_publication_until_all_renames_finish() {
use crate::storage_api_contracts::namespace::NamespaceLocking as _;
let temp_dir = tempfile::tempdir().expect("create scanner PUT tail store dir");
let (ctx, store, shutdown) =
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "scanner-put-tails", &[4])).await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(Arc::clone(&store), Vec::new()).await;
let bucket = format!("scanner-put-tails-{}", Uuid::new_v4());
store
.make_bucket(&bucket, &MakeBucketOptions::default())
.await
.expect("create scanner PUT tail bucket");
let set = &store.pools[0].disk_set[0];
let objects = [("scanner-tail-a", vec![0xA1; 273]), ("scanner-tail-b", vec![0xB2; 379])];
temp_env::async_with_vars([(crate::set_disk::ENV_RUSTFS_PUT_RENAME_EARLY_ACK_ENABLE, Some("true"))], async {
let (active, blocked, movement_generation) = store.scanner_data_movement_activity().await;
assert!(!active && !blocked);
assert!(ctx.scanner_publication_state_allowed(), "the set admission cache should start allowed");
let (old_lease, _) = store
.acquire_scanner_publication_lease(movement_generation, crate::runtime::instance::SCANNER_PUBLICATION_LEASE_TTL)
.await
.expect("publication lease should be admitted before either PUT starts");
let barriers: Vec<_> = objects
.iter()
.map(|(object, _)| {
crate::set_disk::rename_fanout_barrier::arm(object, 0, crate::set_disk::rename_fanout_barrier::PHASE_RENAME)
})
.collect();
let trackers: Vec<_> = objects
.iter()
.map(|(object, _)| crate::set_disk::rename_fanout_barrier::observe_tasks(object))
.collect();
let puts: Vec<_> = objects
.iter()
.map(|(object, body)| {
let put_store = Arc::clone(&store);
let put_bucket = bucket.clone();
let object = *object;
let body = body.clone();
tokio::spawn(async move {
let mut reader = PutObjReader::from_vec(body);
put_store
.put_object(&put_bucket, object, &mut reader, &ObjectOptions::default())
.await
})
})
.collect();
let committed = tokio::time::timeout(Duration::from_secs(30), async {
for barrier in &barriers {
barrier.wait_until_paused().await;
}
let mut committed = Vec::with_capacity(puts.len());
for put in puts {
committed.push(
put.await
.expect("early-ACK PUT task should join while its tail is paused")
.expect("root PUT should return after quorum without waiting for its tail"),
);
}
committed
})
.await
.expect("both root PUTs must quorum-ACK while their tail disks remain paused");
assert!(trackers.iter().all(|tracker| tracker.running() >= 1));
assert!(ctx.namespace_commits_pending());
assert!(
ctx.scanner_publication_state_allowed(),
"pending PUT tails must not disable scanner namespace walks"
);
let (active, blocked, observed_movement_generation) = store.scanner_data_movement_activity().await;
assert!(!active, "ordinary PUT tails are not decommission or rebalance work");
assert!(!blocked, "ordinary PUT tails must not block the movement-only scan baseline");
assert_eq!(observed_movement_generation, movement_generation);
assert!(store.scanner_data_usage_publication_blocked().await);
assert!(store.scanner_data_usage_publication_admission_guard().await.is_some());
assert!(set.scanner_data_usage_publication_admission_guard().await.is_some());
for error in [
store
.acquire_scanner_publication_lease(
movement_generation,
crate::runtime::instance::SCANNER_PUBLICATION_LEASE_TTL,
)
.await
.expect_err("a new remote publication lease must reject pending PUT tails"),
store
.validate_scanner_publication_lease(old_lease, movement_generation)
.await
.expect_err("an existing remote lease must not bypass pending PUT tails"),
store
.acquire_scanner_publication_lease_guard(old_lease)
.await
.expect_err("target-side publication admission must reject pending PUT tails"),
] {
assert!(
error.to_string().contains("blocked"),
"publication must fail because of active tails: {error}"
);
}
store.release_scanner_publication_lease(old_lease).await;
for (index, barrier) in barriers.iter().enumerate() {
let commit_generation = ctx.namespace_commit_generation();
let namespace_generation = store.scanner_namespace_mutation_generation();
barrier.release();
tokio::time::timeout(Duration::from_secs(30), async {
while trackers[index].running() != 0 || ctx.namespace_commit_generation() <= commit_generation {
tokio::task::yield_now().await;
}
if index + 1 == barriers.len() {
while ctx.namespace_commits_pending() {
tokio::task::yield_now().await;
}
}
})
.await
.expect("released tail must drain and publish its terminal namespace generation");
assert!(store.scanner_namespace_mutation_generation() > namespace_generation);
let pending = index + 1 < barriers.len();
assert_eq!(ctx.namespace_commits_pending(), pending);
assert_eq!(store.scanner_data_usage_publication_blocked().await, pending);
assert!(!store.scanner_data_movement_activity().await.1);
assert!(store.scanner_data_usage_publication_admission_guard().await.is_some());
assert!(set.scanner_data_usage_publication_admission_guard().await.is_some());
}
let (lease, generation) = store
.acquire_scanner_publication_lease(movement_generation, crate::runtime::instance::SCANNER_PUBLICATION_LEASE_TTL)
.await
.expect("remote publication lease should resume after both tails drain");
store
.validate_scanner_publication_lease(lease, generation)
.await
.expect("a resumed remote publication lease should validate");
drop(
store
.acquire_scanner_publication_lease_guard(lease)
.await
.expect("target-side publication admission should resume after both tails drain"),
);
assert!(store.release_scanner_publication_lease(lease).await);
let disks = set.disk_inventory().await;
assert_eq!(disks.len(), 4);
for ((object, body), committed) in objects.iter().zip(&committed) {
let logical_size = i64::try_from(body.len()).expect("fixture payload size should fit i64");
let etag = committed.etag.as_ref().expect("root PUT should return a committed ETag");
for (disk_index, disk) in disks.iter().enumerate() {
let file_info = disk
.as_ref()
.expect("every fixture disk should remain online")
.read_version(
"",
&bucket,
object,
"",
&crate::disk::ReadOptions {
read_data: true,
..Default::default()
},
)
.await
.unwrap_or_else(|err| panic!("disk {disk_index} should publish {object} after its tail finishes: {err}"));
assert_eq!(file_info.size, logical_size);
assert_eq!(file_info.metadata.get(http::header::ETAG.as_str()), Some(etag));
assert!(
file_info.inline_data(),
"small fixture payloads should have an inline shard on every disk"
);
let inline_data = file_info.data.as_ref().expect("every disk should retain its inline shard");
let erasure = crate::erasure::coding::Erasure::try_new_with_options(
file_info.erasure.data_blocks,
file_info.erasure.parity_blocks,
file_info.erasure.block_size,
file_info.uses_legacy_checksum,
)
.expect("persisted erasure geometry should be valid");
let shard_size =
usize::try_from(erasure.shard_file_size(logical_size)).expect("fixture shard size should fit usize");
crate::erasure::coding::bitrot_verify(
Cursor::new(inline_data.clone()),
inline_data.len(),
shard_size,
rustfs_utils::HashAlgorithm::HighwayHash256S,
erasure.shard_size(),
)
.await
.unwrap_or_else(|err| panic!("disk {disk_index} should retain a complete valid shard for {object}: {err}"));
}
let mut reader = store
.get_object_reader(&bucket, object, None, HeaderMap::new(), &ObjectOptions::default())
.await
.expect("fully drained PUT should be readable");
let mut actual = Vec::new();
reader.stream.read_to_end(&mut actual).await.expect("PUT body should drain");
assert_eq!(&actual, body);
}
let generation_before_internal_put = ctx.namespace_commit_generation();
let internal_object = "scanner-tail-regression/internal-metadata";
let internal_body = b"scanner metadata must not invalidate its own publication";
let mut internal_reader = PutObjReader::from_vec(internal_body.to_vec());
store
.put_object(RUSTFS_META_BUCKET, internal_object, &mut internal_reader, &ObjectOptions::default())
.await
.expect("internal metadata PUT should commit without scanner self-invalidation");
let internal_lock = set
.new_ns_lock(RUSTFS_META_BUCKET, internal_object)
.await
.expect("internal metadata tail lock should be available");
drop(
internal_lock
.get_write_lock(Duration::from_secs(30))
.await
.expect("internal metadata tail should drain"),
);
assert_eq!(ctx.namespace_commit_generation(), generation_before_internal_put);
assert!(!ctx.namespace_commits_pending());
assert!(store.scanner_data_usage_publication_admission_guard().await.is_some());
assert!(set.scanner_data_usage_publication_admission_guard().await.is_some());
let mut internal_reader = store
.get_object_reader(RUSTFS_META_BUCKET, internal_object, None, HeaderMap::new(), &ObjectOptions::default())
.await
.expect("internal metadata should remain readable");
let mut actual = Vec::new();
internal_reader
.stream
.read_to_end(&mut actual)
.await
.expect("internal metadata body should drain");
assert_eq!(actual, internal_body);
})
.await;
shutdown.cancel();
}
#[cfg(feature = "test-util")]
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
#[serial_test::serial(storage_class_env)]
async fn cancelled_early_ack_put_keeps_scanner_publication_blocked_until_tail_finishes() {
let temp_dir = tempfile::tempdir().expect("create cancelled scanner PUT tail store dir");
let (ctx, store, shutdown) =
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "scanner-cancelled-put-tail", &[4])).await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(Arc::clone(&store), Vec::new()).await;
let bucket = format!("scanner-cancelled-put-tail-{}", Uuid::new_v4());
let object = "scanner-cancelled-tail";
let body = vec![0xC3; 273];
store
.make_bucket(&bucket, &MakeBucketOptions::default())
.await
.expect("create cancelled scanner PUT tail bucket");
temp_env::async_with_vars([(crate::set_disk::ENV_RUSTFS_PUT_RENAME_EARLY_ACK_ENABLE, Some("true"))], async {
let tracker = crate::set_disk::rename_fanout_barrier::observe_tasks(object);
let tail =
crate::set_disk::rename_fanout_barrier::arm(object, 0, crate::set_disk::rename_fanout_barrier::PHASE_RENAME);
let quorum = crate::set_disk::PutObjectCommitBarrier::install(
&bucket,
object,
crate::set_disk::PutObjectCommitPause::AfterRenameQuorum,
);
let handoff = crate::set_disk::PutObjectCommitBarrier::install(
&bucket,
object,
crate::set_disk::PutObjectCommitPause::AfterRenameHandoff,
);
let put_store = Arc::clone(&store);
let put_bucket = bucket.clone();
let put_body = body.clone();
let put = tokio::spawn(async move {
let mut reader = PutObjReader::from_vec(put_body);
put_store
.put_object(&put_bucket, object, &mut reader, &ObjectOptions::default())
.await
});
tokio::time::timeout(Duration::from_secs(30), tail.wait_until_paused())
.await
.expect("cancelled PUT should pause one disk before rename");
quorum.wait_until_paused().await;
put.abort();
assert!(
put.await
.expect_err("caller should be cancelled after rename quorum")
.is_cancelled()
);
quorum.release();
handoff.wait_until_paused().await;
assert!(tracker.running() >= 1);
assert!(ctx.namespace_commits_pending());
assert!(!store.scanner_data_movement_activity().await.1);
assert!(store.scanner_data_usage_publication_blocked().await);
assert!(store.scanner_data_usage_publication_admission_guard().await.is_some());
assert!(
store.pools[0].disk_set[0]
.scanner_data_usage_publication_admission_guard()
.await
.is_some()
);
let generation = store.scanner_namespace_mutation_generation();
handoff.release();
tail.release();
tokio::time::timeout(Duration::from_secs(30), async {
while tracker.running() != 0 || ctx.namespace_commits_pending() {
tokio::task::yield_now().await;
}
})
.await
.expect("cancelled request's detached fanout must release scanner admission after finishing");
assert!(store.scanner_namespace_mutation_generation() > generation);
assert!(!store.scanner_data_usage_publication_blocked().await);
assert!(store.scanner_data_usage_publication_admission_guard().await.is_some());
for (disk_index, disk) in store.pools[0].disk_set[0].disk_inventory().await.iter().enumerate() {
let file_info = disk
.as_ref()
.expect("cancelled PUT fixture disk should remain online")
.read_version("", &bucket, object, "", &crate::disk::ReadOptions::default())
.await
.unwrap_or_else(|err| panic!("cancelled PUT must still publish on disk {disk_index}: {err}"));
assert_eq!(file_info.size, i64::try_from(body.len()).expect("fixture body size should fit i64"));
}
let mut reader = store
.get_object_reader(&bucket, object, None, HeaderMap::new(), &ObjectOptions::default())
.await
.expect("a cancelled caller must not discard its quorum-committed object");
let mut actual = Vec::new();
reader
.stream
.read_to_end(&mut actual)
.await
.expect("cancelled PUT body should drain");
assert_eq!(actual, body);
})
.await;
shutdown.cancel();
}
#[cfg(feature = "test-util")]
#[test]
#[serial_test::serial(storage_class_env)]
@@ -2979,6 +3328,43 @@ mod tests {
#[cfg(feature = "test-util")]
const DECOMMISSION_TEST_FAULT_STAGE_TIERED: &str = "decommission_tiered_object";
fn inject_decommission_copy_fault(faults: &AtomicUsize, attempt: usize, succeeded: bool) -> bool {
// Entry retries reset attempt, not the global post-commit fault budget.
// A real failure may consume an attempt, so preserve the final chance.
succeeded
&& faults
.fetch_update(Ordering::SeqCst, Ordering::SeqCst, |faults| {
(faults < crate::core::pools::DECOMMISSION_VERSION_COPY_ATTEMPTS.saturating_sub(1)
&& attempt < crate::core::pools::DECOMMISSION_VERSION_COPY_ATTEMPTS)
.then_some(faults.saturating_add(1))
})
.is_ok()
}
#[test]
fn decommission_copy_fault_budget_survives_entry_restarts_and_preserves_last_attempt() {
let cases: &[&[(usize, bool, bool)]] = &[
&[(1, true, true), (2, true, true), (3, true, false)],
&[(1, true, true), (2, false, false), (1, true, true), (2, true, false)],
&[(1, true, true), (2, false, false), (3, true, false)],
&[(1, true, true), (1, true, true), (1, true, false)],
&[(3, true, false), (4, true, false)],
];
for case in cases {
let faults = AtomicUsize::new(0);
let mut expected_faults = 0;
for &(attempt, succeeded, expected) in *case {
assert_eq!(
inject_decommission_copy_fault(&faults, attempt, succeeded),
expected,
"fault plan {case:?} at attempt {attempt}"
);
expected_faults += usize::from(expected);
assert_eq!(faults.load(Ordering::SeqCst), expected_faults);
}
}
}
async fn seed_decommission_source(
store: &Arc<crate::store::ECStore>,
bucket: &str,
@@ -4991,6 +5377,7 @@ mod tests {
"decommission-delete-fence",
&[(2, 4), (1, 4)],
CancellationToken::new(),
None,
))
.await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
@@ -5218,25 +5605,10 @@ mod tests {
let fault_bucket = other_bucket.clone();
let _fault_guard = crate::core::pools::DecommissionTestFaultGuard::install(Arc::new(
move |stage, bucket, object, attempt, succeeded| {
let candidate = succeeded
&& stage == DECOMMISSION_TEST_FAULT_STAGE_MIGRATE_OBJECT
let candidate = stage == DECOMMISSION_TEST_FAULT_STAGE_MIGRATE_OBJECT
&& bucket == fault_bucket.as_str()
&& object == other_object;
if !candidate {
return false;
}
// Keep the fault budget global across any
// entry-level re-list; its inner attempt counter
// restarts after SourceChanged.
ordinary_faults_for_hook
.fetch_update(Ordering::SeqCst, Ordering::SeqCst, |faults| {
let next_fault = faults.saturating_add(1);
(faults < crate::core::pools::DECOMMISSION_VERSION_COPY_ATTEMPTS.saturating_sub(1)
&& attempt == next_fault)
.then_some(next_fault)
})
.is_ok()
candidate && inject_decommission_copy_fault(&ordinary_faults_for_hook, attempt, succeeded)
},
));
@@ -5275,6 +5647,15 @@ mod tests {
changed_result.expect("SourceChanged entry retry must converge");
other_result.expect("other bucket entry must continue through ordinary copy retries");
assert_eq!(
store.pool_meta.read().await.pools[0]
.decommission
.as_ref()
.expect("decommission progress should remain available")
.items_decommission_failed,
0,
"entry completion must not hide an exhausted copy failure"
);
assert!(!rx.is_cancelled(), "entry-level SourceChanged must not cancel the shared worker token");
assert_eq!(mutation_calls.load(Ordering::SeqCst), 2, "entry must be re-listed after SourceChanged");
assert_eq!(ordinary_faults.load(Ordering::SeqCst), 2, "ordinary copy must consume the retry budget");
@@ -5903,6 +6284,7 @@ mod tests {
"reverse-decommission-fixed-target",
&[(1, 4), (1, 4)],
CancellationToken::new(),
None,
))
.await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
@@ -6324,6 +6706,7 @@ mod tests {
"multi-set-decommission-source-cleanup",
&[(2, 4)],
CancellationToken::new(),
None,
))
.await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
@@ -8834,18 +9217,17 @@ mod tests {
const MANIFEST_COUNT: usize = 10;
let temp_dir = tempfile::tempdir().expect("create fast manifest pass recovery store dir");
let (ctx, store, _shutdown) =
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "tier-delete-fast-manifest-pass", &[4])).await;
let mut instance_ctx = crate::runtime::instance::InstanceContext::new();
instance_ctx.suppress_tier_delete_journal_recovery_for_test();
let (ctx, store, shutdown) = without_storage_class_env(build_isolated_test_store_with_layout(
temp_dir.path(),
"tier-delete-fast-manifest-pass",
&[(1, 4)],
CancellationToken::new(),
Some(Arc::new(instance_ctx)),
))
.await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
let bucket = "tier-delete-fast-manifest-pass-bucket";
store
.make_bucket(bucket, &MakeBucketOptions::default())
.await
.expect("fast manifest pass bucket should be created");
let incarnation = store
.bucket_incarnation_id(bucket)
.await
.expect("fast manifest pass bucket incarnation should resolve");
let tier_name = "FAST-MANIFEST-PASS";
let backend = register_mock_tier(&ctx.tier_config_mgr(), tier_name).await;
let backend_identity = TierConfigMgr::acquire_operation_lease(&ctx.tier_config_mgr(), tier_name)
@@ -8853,9 +9235,19 @@ mod tests {
.expect("fast manifest pass tier lease should resolve")
.backend_identity();
for index in 0..MANIFEST_COUNT {
// Pagination must not depend on same-bucket lock wait deadlines.
let bucket = format!("tier-delete-fast-manifest-pass-{index}");
store
.make_bucket(&bucket, &MakeBucketOptions::default())
.await
.expect("fast manifest pass bucket should be created");
let incarnation = store
.bucket_incarnation_id(&bucket)
.await
.expect("fast manifest pass bucket incarnation should resolve");
install_aborting_dispatch_fixture(
store.clone(),
bucket,
&bucket,
incarnation,
&format!("manifest-page-{index:06}/"),
tier_name,
@@ -8886,12 +9278,78 @@ mod tests {
"one production pass must cross the default eight-manifest page limit"
);
assert_eq!(stats.manifests.scanned, MANIFEST_COUNT);
assert_eq!(stats.manifests.deleted, MANIFEST_COUNT);
assert_eq!(stats.manifests.failed, 0);
assert_eq!(stats.manifests.deleted, MANIFEST_COUNT, "full recovery result: {stats:?}");
assert_eq!(stats.manifests.failed, 0, "full recovery result: {stats:?}");
assert_eq!(manifest_marker, None);
assert_eq!(tier_delete_dispatch_manifest_count(store.clone()).await, 0);
assert_eq!(tier_delete_journal_count(store).await, 0);
assert_eq!(backend.remove_count().await, 0, "rollback recovery must not call the remote tier");
shutdown.cancel();
}
#[cfg(feature = "test-util")]
#[tokio::test]
#[serial_test::serial(storage_class_env)]
async fn tier_delete_manual_pass_retains_manifest_owned_by_startup_recovery() {
let temp_dir = tempfile::tempdir().expect("create automatic recovery ownership store dir");
let (ctx, store, shutdown) =
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "tier-delete-auto-owner", &[4])).await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
let bucket = "tier-delete-auto-owner-bucket";
store
.make_bucket(bucket, &MakeBucketOptions::default())
.await
.expect("automatic recovery bucket should be created");
let incarnation = store.bucket_incarnation_id(bucket).await.expect("bucket incarnation");
let tier_name = "AUTO-OWNER";
let backend = register_mock_tier(&ctx.tier_config_mgr(), tier_name).await;
let identity = TierConfigMgr::acquire_operation_lease(&ctx.tier_config_mgr(), tier_name)
.await
.expect("automatic recovery tier lease")
.backend_identity();
// The automatic worker must not observe a partially installed fixture.
let lifecycle_guard = store
.acquire_bucket_lifecycle_write_lock(bucket)
.await
.expect("fixture lifecycle lock");
let (manifest_name, entries) =
install_aborting_dispatch_fixture(store.clone(), bucket, incarnation, "auto-owner/", tier_name, identity, 1).await;
let journal_name = tier_delete_journal_object_name(&entries[0]);
let hook = TierDeleteDispatchRollbackTestHook::install_slow_delete(&journal_name, &journal_name);
drop(lifecycle_guard);
ctx.wake_tier_delete_journal_recovery();
tokio::time::timeout(Duration::from_secs(30), hook.wait_until_delete_paused())
.await
.expect("startup recovery should own the manifest before a manual pass");
assert!(tier_delete_dispatch_manifest_recovery_inflight_for_test(&store, &manifest_name));
let stats = recover_tier_delete_dispatch_manifests(store.clone(), 8, None)
.await
.expect("manual recovery scan");
assert_eq!(stats.scanned, 1, "{stats:?}");
assert_eq!(stats.retained, 1, "{stats:?}");
assert_eq!(stats.deleted, 0, "{stats:?}");
assert_eq!(stats.failed, 0, "{stats:?}");
assert_eq!(tier_delete_dispatch_manifest_count(store.clone()).await, 1);
assert_eq!(tier_delete_journal_count(store.clone()).await, 1);
hook.release_delete();
tokio::time::timeout(Duration::from_secs(30), async {
loop {
let manifest_gone = matches!(com::read_config(store.clone(), &manifest_name).await, Err(Error::ConfigNotFound));
if manifest_gone && !tier_delete_dispatch_manifest_recovery_inflight_for_test(&store, &manifest_name) {
break;
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
})
.await
.expect("automatic recovery should converge without a manual retry");
assert_eq!(tier_delete_dispatch_manifest_count(store.clone()).await, 0);
assert_eq!(tier_delete_journal_count(store).await, 0);
assert_eq!(backend.remove_count().await, 0, "rollback must not delete from the remote tier");
shutdown.cancel();
}
#[cfg(feature = "test-util")]
@@ -13016,6 +13474,7 @@ mod tests {
"partial-set-prefix-delete",
&[(2, 4)],
CancellationToken::new(),
None,
))
.await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
@@ -16388,6 +16847,7 @@ mod tests {
"prepared-directory-recovery",
&[(2, 4)],
shutdown,
None,
))
.await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
@@ -17036,6 +17496,38 @@ mod tests {
.expect("test thread should complete");
}
#[cfg(feature = "test-util")]
#[tokio::test]
#[serial_test::serial(storage_class_env)]
async fn odm_write_back_requires_one_set_and_enabled_namespace_locking() {
for (layout, locking, supported) in [
(&[(1, 4)][..], true, true),
(&[(1, 4), (1, 4)][..], true, false),
(&[(2, 4)][..], true, false),
(&[(1, 4)][..], false, false),
] {
temp_env::async_with_vars([("RUSTFS_LOCK_ENABLED", Some(if locking { "true" } else { "false" }))], async {
let dir = tempfile::tempdir().expect("isolated topology");
let shutdown = CancellationToken::new();
let (_ctx, store, _) = without_storage_class_env(build_isolated_test_store_with_layout(
dir.path(),
"odm-topology",
layout,
shutdown.clone(),
None,
))
.await;
assert_eq!(
store.supports_atomic_create_only_write_back(),
supported,
"layout={layout:?}, locking={locking}"
);
shutdown.cancel();
})
.await;
}
}
#[cfg(feature = "test-util")]
#[tokio::test]
#[serial_test::serial(storage_class_env)]
+8 -7
View File
@@ -848,7 +848,7 @@ impl ECStore {
}
pub fn scanner_namespace_mutation_generation(&self) -> u64 {
list_objects::scanner_namespace_mutation_generation()
list_objects::scanner_namespace_mutation_generation().saturating_add(self.ctx.namespace_commit_generation())
}
pub async fn scanner_data_movement_active(&self) -> bool {
@@ -857,7 +857,7 @@ impl ECStore {
}
/// Return the storage-owned movement state and generation as one
/// authenticated activity snapshot. The read lock is acquired before
/// authenticated activity snapshot. The read lock is acquired before
/// the state locks (cancelers, pool metadata, then rebalance metadata),
/// matching the transition writer order and preventing a terminal state
/// from being reported with the preceding generation.
@@ -886,11 +886,12 @@ impl ECStore {
/// Returns whether scanner metadata may still be hidden by a local
/// data-movement state. Terminal failed/canceled decommission entries
/// remain suspended until an operator clears or retries them, so they are
/// a publication barrier even after the worker has stopped.
/// a publication barrier even after the worker has stopped. Active PUT
/// rename fanouts also defer publication, including post-ACK tails.
pub async fn scanner_data_usage_publication_blocked(&self) -> bool {
let operation_gate = self.ctx.data_movement_operation_gate();
let _operation_guard = operation_gate.read_owned().await;
self.scanner_data_usage_publication_snapshot_blocked().await
self.scanner_data_usage_publication_snapshot_blocked().await || self.ctx.namespace_commits_pending()
}
pub async fn scanner_data_movement_pause_status(&self) -> ScannerDataMovementPauseStatus {
@@ -1070,7 +1071,7 @@ impl ECStore {
{
return Err(Error::other("scanner publication lease generation is stale"));
}
if self.scanner_data_movement_snapshot_locked().await.1 {
if self.scanner_data_movement_snapshot_locked().await.1 || self.ctx.namespace_commits_pending() {
return Err(Error::other("scanner publication lease is blocked by data movement"));
}
@@ -1109,7 +1110,7 @@ impl ECStore {
{
return Err(Error::other("scanner publication lease generation is stale"));
}
if self.scanner_data_movement_snapshot_locked().await.1 {
if self.scanner_data_movement_snapshot_locked().await.1 || self.ctx.namespace_commits_pending() {
return Err(Error::other("scanner publication lease is blocked by data movement"));
}
if !self.ctx.scanner_publication_lease_is_active(token).await {
@@ -1129,7 +1130,7 @@ impl ECStore {
if self.ctx.data_movement_generation_exhausted() || self.ctx.data_movement_operation_epoch_exhausted() {
return Err(Error::other("scanner publication lease generation is exhausted"));
}
if self.scanner_data_movement_snapshot_locked().await.1 {
if self.scanner_data_movement_snapshot_locked().await.1 || self.ctx.namespace_commits_pending() {
return Err(Error::other("scanner publication lease is blocked by data movement"));
}
let Some(lease_generation) = self.ctx.scanner_publication_lease_generation(token).await else {