Merge branch 'main' into overtrue/fix/lifecycle-rule-validation

This commit is contained in:
Zhengchao An
2026-09-05 16:50:27 +08:00
committed by GitHub
20 changed files with 2358 additions and 247 deletions
@@ -42,6 +42,7 @@ jobs:
- name: Check latest scheduled runs
env:
GH_TOKEN: ${{ secrets.GITHUB_TOKEN }}
RUSTFS_DEFAULT_BRANCH: ${{ github.event.repository.default_branch }}
run: |
set +e
python3 scripts/check_scheduled_validation_freshness.py \
+60 -12
View File
@@ -45,6 +45,11 @@ use tracing::{debug, error, info, warn};
use super::{DiskError, Endpoint, HealDiskExt as _, local_disk_map_read};
const KEEP_HEAL_TASK_STATUS_DURATION: Duration = Duration::from_secs(10 * 60);
// Each cache includes alias tokens in its count and byte budget. Eviction
// removes every token sharing a snapshot; neither cache retains repair state.
const MAX_COMPLETED_HEAL_TOKENS: usize = 1024;
const MAX_COMPLETED_HEAL_BYTES: usize = 64 * 1024 * 1024;
const MAX_COMPLETED_HEAL_RESULT_BYTES: usize = 1024 * 1024;
const DISPLACED_HEAL_REASON: &str = "reason=displaced; retry_hint=submit_again";
const LOG_COMPONENT_HEAL: &str = "heal";
const LOG_SUBSYSTEM_DISK_SCANNER: &str = "disk_scanner";
@@ -180,6 +185,8 @@ fn record_displaced_terminal(
request: &HealRequest,
) -> Arc<CompletedHealStatus> {
let terminal = Arc::new(CompletedHealStatus {
progress: None,
retained_bytes: std::sync::OnceLock::new(),
heal_type: request.heal_type.clone(),
status: HealTaskStatus::Failed {
error: format!("heal task displaced by a higher-priority request ({DISPLACED_HEAL_REASON})"),
@@ -193,6 +200,7 @@ fn record_displaced_terminal(
let mut terminals = lock_displaced_terminals(registry);
prune_completed_heal_statuses(&mut terminals);
terminals.insert(request.id.clone(), Arc::clone(&terminal));
prune_completed_heal_statuses(&mut terminals);
terminal
}
@@ -209,9 +217,15 @@ async fn remove_displaced_task_aliases(
.collect::<Vec<_>>();
let mut displaced_terminals = lock_displaced_terminals(terminals);
prune_completed_heal_statuses(&mut displaced_terminals);
for alias_id in alias_ids {
displaced_terminals.insert(alias_id, Arc::clone(terminal));
if displaced_terminals
.get(task_id)
.is_some_and(|current| Arc::ptr_eq(current, terminal))
{
for alias_id in alias_ids {
displaced_terminals.insert(alias_id, Arc::clone(terminal));
}
}
prune_completed_heal_statuses(&mut displaced_terminals);
aliases.retain(|alias_id, alias| alias_id != task_id && alias.task_id != task_id);
}
@@ -222,6 +236,36 @@ async fn remove_task_aliases_for_task(registry: &Arc<Mutex<HashMap<String, HealT
.retain(|alias_id, alias| alias_id != task_id && alias.task_id != task_id);
}
// Callers hold active ownership until publication. Lock order is active ->
// retrying (when needed) -> aliases -> completed; queries release aliases
// before looking up active state. Publishing aliases before removing their
// mapping keeps both an already-resolved token and a new lookup valid.
async fn publish_completed_heal(
completed_heals: &Mutex<HashMap<String, Arc<CompletedHealStatus>>>,
task_aliases: &Mutex<HashMap<String, HealTaskAlias>>,
task_id: &str,
completed: CompletedHealStatus,
terminal: bool,
) {
let completed = Arc::new(completed);
completed.retained_bytes();
let mut aliases = task_aliases.lock().await;
let mut retained = completed_heals.lock().await;
if let Some(previous) = retained.get(task_id).cloned() {
for entry in retained.values_mut().filter(|entry| Arc::ptr_eq(entry, &previous)) {
*entry = Arc::clone(&completed);
}
}
retained.insert(task_id.to_owned(), Arc::clone(&completed));
if terminal {
for (alias_id, _) in aliases.iter().filter(|(_, alias)| alias.task_id == task_id) {
retained.insert(alias_id.clone(), Arc::clone(&completed));
}
aliases.retain(|alias_id, alias| alias_id != task_id && alias.task_id != task_id);
}
prune_completed_heal_statuses(&mut retained);
}
#[derive(Debug, Clone)]
pub struct HealTaskReport {
pub status: HealTaskStatus,
@@ -268,7 +312,7 @@ fn completed_task_report(completed: &CompletedHealStatus, since: Option<u64>) ->
let result_items = match since {
None => completed.seqed_items.iter().map(|(_, item)| item.clone()).collect(),
Some(cursor) => {
if cursor + 1 < completed.min_seq {
if cursor.saturating_add(1) < completed.min_seq {
lagged = true;
}
completed
@@ -283,7 +327,7 @@ fn completed_task_report(completed: &CompletedHealStatus, since: Option<u64>) ->
status: completed.status.clone(),
result_items,
result_items_truncated: completed.result_items_truncated || lagged,
progress: None,
progress: completed.progress.clone(),
next_seq: completed.next_seq,
min_seq: completed.min_seq,
}
@@ -1847,14 +1891,14 @@ impl HealManager {
pub async fn get_task_progress(&self, task_id: &str) -> Result<HealProgress> {
let canonical_task_id = self.canonical_task_id(task_id).await;
let active_heals = self.active_heals.lock().await;
if let Some(task) = active_heals.get(&canonical_task_id) {
Ok(task.get_progress().await)
} else {
Err(Error::TaskNotFound {
task_id: task_id.to_string(),
})
}
let progress = match self.lookup_task_state(&canonical_task_id, None).await {
TaskStateLookup::Active(task) => Some(task.get_progress().await),
TaskStateLookup::Completed(completed) => completed.progress.clone(),
_ => None,
};
progress.ok_or_else(|| Error::TaskNotFound {
task_id: task_id.to_string(),
})
}
/// Cancel task
@@ -1864,6 +1908,8 @@ impl HealManager {
let mut active_heals = self.active_heals.lock().await;
if let Some(task) = active_heals.get(&canonical_task_id) {
task.cancel().await?;
let completed = CompletedHealStatus::snapshot(task, HealTaskStatus::Cancelled).await;
publish_completed_heal(&self.completed_heals, &self.task_aliases, &canonical_task_id, completed, true).await;
active_heals.remove(&canonical_task_id);
publish_active_heal_count(&active_heals);
info!(
@@ -1940,6 +1986,8 @@ impl HealManager {
for task_id in &task_ids {
if let Some(task) = active_heals.get(task_id) {
task.cancel().await?;
let completed = CompletedHealStatus::snapshot(task, HealTaskStatus::Cancelled).await;
publish_completed_heal(&self.completed_heals, &self.task_aliases, task_id, completed, true).await;
}
active_heals.remove(task_id);
cancelled += 1;
+129
View File
@@ -82,6 +82,8 @@ pub(super) enum QueuePushOutcome {
pub(super) struct CompletedHealStatus {
pub(super) heal_type: HealType,
pub(super) status: HealTaskStatus,
pub(super) progress: Option<HealProgress>,
pub(super) retained_bytes: std::sync::OnceLock<usize>,
pub(super) result_items_truncated: bool,
pub(super) completed_at: SystemTime,
/// Sequence-stamped retained window, archived with the completion so
@@ -92,6 +94,133 @@ pub(super) struct CompletedHealStatus {
pub(super) min_seq: u64,
}
impl CompletedHealStatus {
// Account for owned capacities, including nested drive arrays. Aliases
// conservatively charge the shared allocation again, keeping both token
// count and retained payload bounded without a second ownership index.
pub(super) fn retained_bytes(&self) -> usize {
*self.retained_bytes.get_or_init(|| self.measure_retained_bytes())
}
fn measure_retained_bytes(&self) -> usize {
let mut bytes = size_of::<Self>();
let mut add = |amount: usize| bytes = bytes.saturating_add(amount);
match &self.heal_type {
HealType::Cluster => {}
HealType::Bucket { bucket } => add(bucket.capacity()),
HealType::Object {
bucket,
object,
version_id,
}
| HealType::ECDecode {
bucket,
object,
version_id,
} => {
add(bucket.capacity());
add(object.capacity());
add(version_id.as_ref().map_or(0, String::capacity));
}
HealType::Prefix { bucket, prefix } => {
add(bucket.capacity());
add(prefix.capacity());
}
HealType::Metadata { bucket, object } => {
add(bucket.capacity());
add(object.capacity());
}
HealType::ErasureSet { buckets, set_disk_id } => {
add(buckets.capacity().saturating_mul(size_of::<String>()));
for bucket in buckets {
add(bucket.capacity());
}
add(set_disk_id.capacity());
}
}
if let HealTaskStatus::Failed { error } | HealTaskStatus::Retrying { error, .. } = &self.status {
add(error.capacity());
}
add(self
.progress
.as_ref()
.and_then(|progress| progress.current_object.as_ref())
.map_or(0, String::capacity));
add(self.seqed_items.capacity().saturating_mul(size_of::<(u64, HealResultItem)>()));
for (_, item) in &self.seqed_items {
add(Self::result_item_heap_bytes(item));
}
bytes
}
fn result_item_heap_bytes(item: &HealResultItem) -> usize {
let mut bytes = 0usize;
let mut add = |amount: usize| bytes = bytes.saturating_add(amount);
for value in [
&item.heal_item_type,
&item.bucket,
&item.object,
&item.version_id,
&item.detail,
] {
add(value.capacity());
}
for infos in [&item.before, &item.after] {
add(infos
.drives
.capacity()
.saturating_mul(size_of::<rustfs_madmin::heal_commands::HealDriveInfo>()));
for drive in &infos.drives {
add(drive.uuid.capacity());
add(drive.endpoint.capacity());
add(drive.state.capacity());
}
}
bytes
}
pub(super) fn bound_result_window(&mut self) {
let mut bytes = 0usize;
let retained = self
.seqed_items
.iter()
.rev()
.take_while(|(_, item)| {
bytes = bytes
.saturating_add(size_of::<(u64, HealResultItem)>())
.saturating_add(Self::result_item_heap_bytes(item));
bytes <= MAX_COMPLETED_HEAL_RESULT_BYTES
})
.count();
let truncated = retained < self.seqed_items.len();
if truncated {
self.seqed_items.drain(..self.seqed_items.len() - retained);
self.seqed_items.shrink_to_fit();
self.min_seq = self.seqed_items.first().map_or(self.next_seq, |(seq, _)| *seq);
self.result_items_truncated = true;
self.retained_bytes.take();
}
}
pub(super) async fn snapshot(task: &HealTask, status: HealTaskStatus) -> Self {
let seqed_items = task.get_seqed_result_items().await;
let (next_seq, min_seq) = task.result_seq_cursors();
let mut snapshot = Self {
heal_type: task.heal_type.clone(),
status,
progress: Some(task.get_progress().await),
retained_bytes: std::sync::OnceLock::new(),
result_items_truncated: task.result_items_truncated(),
completed_at: SystemTime::now(),
seqed_items,
next_seq,
min_seq,
};
snapshot.bound_result_window();
snapshot
}
}
#[derive(Debug, Clone)]
pub(super) struct HealTaskAlias {
pub(super) task_id: String,
+75 -39
View File
@@ -264,7 +264,7 @@ impl HealManager {
error: error.clone(),
retry_attempt: request.retry_attempts,
});
let retry_request_for_queue = retry_request;
let mut retry_request_for_queue = retry_request;
let retry_cancel_token = retry_request_for_queue.as_ref().map(|_| CancellationToken::new());
if retry_request_for_queue.is_none() {
replacement_recovery_anchors_clone
@@ -272,7 +272,35 @@ impl HealManager {
.unwrap_or_else(|poisoned| poisoned.into_inner())
.remove(&task_id);
}
let mut completed_status = match retry_request_for_status {
Some(status) => status,
None => task.get_status().await,
};
let mut completed_status_entry = CompletedHealStatus::snapshot(&task, completed_status.clone()).await;
let completed_progress = task.get_progress().await;
#[cfg(test)]
tests::pause_completed_retention_before_publish(&task_id, &completed_status).await;
let mut active_heals_guard = active_heals_clone.lock().await;
let owns_completion = active_heals_guard.contains_key(&task_id);
let cancelled_completion = if owns_completion {
false
} else {
// Cancellation can win while a finished worker waits
// for active ownership. It must not resurrect a retry
// or replace an acknowledged cancellation with success.
retry_request_for_queue = None;
completed_heals_clone
.lock()
.await
.get(&task_id)
.is_some_and(|completed| completed.status == HealTaskStatus::Cancelled)
};
if cancelled_completion {
completed_status = HealTaskStatus::Cancelled;
completed_status_entry.status = HealTaskStatus::Cancelled;
}
let terminal_completion = !matches!(completed_status, HealTaskStatus::Retrying { .. });
let successful_completion = matches!(completed_status, HealTaskStatus::Completed);
// Keep retry ownership continuous: status snapshots acquire
// these locks in the same active -> retrying order.
let mut retrying_heals_guard = if let (Some((request, _, error)), Some(cancel_token)) =
@@ -295,6 +323,16 @@ impl HealManager {
} else {
None
};
if owns_completion || cancelled_completion {
publish_completed_heal(
&completed_heals_clone,
&task_aliases_clone,
&task_id,
completed_status_entry,
terminal_completion,
)
.await;
}
let completed_task = active_heals_guard.remove(&task_id);
if let Some(completed_task) = completed_task.as_ref() {
publish_active_heal_count(&active_heals_guard);
@@ -304,33 +342,10 @@ impl HealManager {
drop(retrying_heals_guard.take());
drop(active_heals_guard);
if let Some(completed_task) = completed_task {
let completed_status = if let Some(status) = retry_request_for_status {
status
} else {
completed_task.get_status().await
};
let terminal_completion = !matches!(completed_status, HealTaskStatus::Retrying { .. });
let successful_completion = matches!(completed_status, HealTaskStatus::Completed);
let completed_progress = completed_task.get_progress().await;
// Single snapshot of the retained window: the task is
// finished and already off the active map, so there is
// no concurrent writer to race with.
let seqed_items = completed_task.get_seqed_result_items().await;
let (next_seq, min_seq) = completed_task.result_seq_cursors();
let completed_status_entry = CompletedHealStatus {
heal_type: completed_task.heal_type.clone(),
status: completed_status.clone(),
result_items_truncated: completed_task.result_items_truncated(),
completed_at: SystemTime::now(),
seqed_items,
next_seq,
min_seq,
};
let mut completed_heals_guard = completed_heals_clone.lock().await;
prune_completed_heal_statuses(&mut completed_heals_guard);
completed_heals_guard.insert(task_id.clone(), Arc::new(completed_status_entry));
drop(completed_heals_guard);
#[cfg(test)]
tests::pause_completed_retention_handoff(&task_id).await;
if completed_task.is_some() {
// update statistics
let mut stats = statistics_clone.write().await;
match completed_status {
@@ -352,10 +367,6 @@ impl HealManager {
} else {
release_mrf_repair_notice_targets(notice_targets);
}
task_aliases_clone
.lock()
.await
.retain(|alias_id, alias| alias_id != &task_id && alias.task_id != task_id);
}
}
@@ -718,17 +729,42 @@ pub(super) fn heal_request_set_key_for_task(task: &HealTask) -> Option<String> {
}
pub(super) fn prune_completed_heal_statuses(completed_heals: &mut HashMap<String, Arc<CompletedHealStatus>>) {
let Ok(now) = SystemTime::now().duration_since(SystemTime::UNIX_EPOCH) else {
return;
};
prune_completed_heal_statuses_at(completed_heals, SystemTime::now());
}
pub(super) fn prune_completed_heal_statuses_at(completed_heals: &mut HashMap<String, Arc<CompletedHealStatus>>, now: SystemTime) {
completed_heals.retain(|_, completed| {
completed
.completed_at
.duration_since(SystemTime::UNIX_EPOCH)
.map(|completed_at| now.saturating_sub(completed_at) <= KEEP_HEAL_TASK_STATUS_DURATION)
now.duration_since(completed.completed_at)
.map(|age| age <= KEEP_HEAL_TASK_STATUS_DURATION)
.unwrap_or(false)
});
let entry_bytes = |key: &String, value: &Arc<CompletedHealStatus>| {
key.capacity()
.saturating_add(size_of::<(String, Arc<CompletedHealStatus>)>())
.saturating_add(value.retained_bytes())
};
let mut bytes = completed_heals
.iter()
.fold(0usize, |total, (key, value)| total.saturating_add(entry_bytes(key, value)));
while completed_heals.len() > MAX_COMPLETED_HEAL_TOKENS || bytes > MAX_COMPLETED_HEAL_BYTES {
let Some(oldest) = completed_heals
.iter()
.min_by(|(left_id, left), (right_id, right)| {
left.completed_at.cmp(&right.completed_at).then_with(|| left_id.cmp(right_id))
})
.map(|(_, value)| Arc::clone(value))
else {
break;
};
completed_heals.retain(|key, value| {
if Arc::ptr_eq(value, &oldest) {
bytes = bytes.saturating_sub(entry_bytes(key, value));
false
} else {
true
}
});
}
}
pub(super) fn can_schedule_request(
+348 -3
View File
@@ -101,6 +101,326 @@ async fn process_manager_queue_once(manager: &HealManager) {
struct MockStorage;
fn completed_retention_fixture(completed_at: SystemTime) -> CompletedHealStatus {
CompletedHealStatus {
heal_type: HealType::Cluster,
status: HealTaskStatus::Completed,
progress: Some(HealProgress {
objects_scanned: 9,
objects_healed: 8,
objects_failed: 1,
..Default::default()
}),
retained_bytes: std::sync::OnceLock::new(),
result_items_truncated: false,
completed_at,
seqed_items: vec![(3, HealResultItem::default()), (4, HealResultItem::default())],
next_seq: 5,
min_seq: 3,
}
}
#[test]
fn completed_retention_cursor_boundaries_preserve_progress() {
let completed = completed_retention_fixture(SystemTime::now());
for (cursor, count, lagged) in [
(0, 2, true),
(1, 2, true),
(2, 2, false),
(3, 1, false),
(4, 0, false),
(5, 0, false),
(u64::MAX, 0, false),
] {
let report = completed_task_report(&completed, Some(cursor));
assert_eq!(report.result_items.len(), count, "cursor={cursor}");
assert_eq!(report.result_items_truncated, lagged, "cursor={cursor}");
assert_eq!(report.progress, completed.progress);
assert_eq!((report.next_seq, report.min_seq), (5, 3));
}
assert_eq!(completed_task_report(&completed, None).result_items.len(), 2);
}
#[tokio::test]
async fn completed_retention_displaced_alias_does_not_resurrect_evicted_snapshot() {
let manager = HealManager::new(Arc::new(MockStorage), None);
let request = HealRequest::bucket("bucket".to_string());
manager.insert_task_alias("alias", &request.id).await;
let terminal = record_displaced_terminal(&manager.displaced_terminals, &request);
lock_displaced_terminals(&manager.displaced_terminals).remove(&request.id);
remove_displaced_task_aliases(&manager.task_aliases, &manager.displaced_terminals, &request.id, &terminal).await;
for token in [&request.id, &"alias".to_string()] {
assert!(matches!(manager.get_task_report(token).await, Err(Error::TaskNotFound { .. })));
}
assert!(manager.task_aliases.lock().await.is_empty());
assert!(lock_displaced_terminals(&manager.displaced_terminals).is_empty());
}
#[test]
fn completed_retention_count_ttl_and_alias_eviction_are_bounded() {
let now = SystemTime::now();
let mut entries = HashMap::new();
let oldest = Arc::new(completed_retention_fixture(now - KEEP_HEAL_TASK_STATUS_DURATION));
entries.insert("oldest".to_string(), Arc::clone(&oldest));
entries.insert("oldest-alias".to_string(), Arc::clone(&oldest));
for index in 2..MAX_COMPLETED_HEAL_TOKENS {
entries.insert(format!("task-{index}"), Arc::new(completed_retention_fixture(now)));
}
prune_completed_heal_statuses_at(&mut entries, now);
assert_eq!(entries.len(), MAX_COMPLETED_HEAL_TOKENS);
entries.insert("cap-plus-one".to_string(), Arc::new(completed_retention_fixture(now)));
prune_completed_heal_statuses_at(&mut entries, now);
assert_eq!(entries.len(), MAX_COMPLETED_HEAL_TOKENS - 1);
assert!(!entries.contains_key("oldest"));
assert!(!entries.contains_key("oldest-alias"));
entries.clear();
entries.insert("ttl-boundary".to_string(), oldest);
entries.insert(
"expired".to_string(),
Arc::new(completed_retention_fixture(
now - KEEP_HEAL_TASK_STATUS_DURATION - Duration::from_nanos(1),
)),
);
entries.insert("future".to_string(), Arc::new(completed_retention_fixture(now + Duration::from_nanos(1))));
prune_completed_heal_statuses_at(&mut entries, now);
assert_eq!(entries.len(), 1);
assert!(entries.contains_key("ttl-boundary"));
prune_completed_heal_statuses_at(&mut entries, now + Duration::from_nanos(1));
assert!(entries.is_empty());
}
#[test]
fn completed_retention_total_byte_cap_and_cap_plus_one() {
let now = SystemTime::now();
let key = "large".to_string();
let mut entry = completed_retention_fixture(now);
let base_bytes = entry.retained_bytes() + key.capacity() + size_of::<(String, Arc<CompletedHealStatus>)>();
entry.retained_bytes.take();
entry.status = HealTaskStatus::Failed {
error: "x".repeat(MAX_COMPLETED_HEAL_BYTES - base_bytes),
};
assert_eq!(
entry.retained_bytes() + key.capacity() + size_of::<(String, Arc<CompletedHealStatus>)>(),
MAX_COMPLETED_HEAL_BYTES
);
let mut entries = HashMap::from([(key, Arc::new(entry))]);
prune_completed_heal_statuses_at(&mut entries, now);
assert_eq!(entries.len(), 1, "exact byte cap remains retained");
let mut over = Arc::try_unwrap(entries.remove("large").expect("entry retained")).expect("entry not shared");
over.retained_bytes.take();
if let HealTaskStatus::Failed { error } = &mut over.status {
*error = "x".repeat(error.len() + 1);
}
entries.insert("large".to_string(), Arc::new(over));
prune_completed_heal_statuses_at(&mut entries, now);
assert!(entries.is_empty(), "oversized metadata cannot escape total byte bound");
}
#[tokio::test]
async fn completed_retention_large_window_keeps_cursors_and_progress() {
let task = HealTask::from_request(HealRequest::bucket("bucket".to_string()), Arc::new(MockStorage));
let mut snapshot = completed_retention_fixture(SystemTime::now());
snapshot.seqed_items[0].1.detail = "x".repeat(MAX_COMPLETED_HEAL_RESULT_BYTES);
snapshot.bound_result_window();
assert_eq!(snapshot.seqed_items.len(), 1);
assert_eq!((snapshot.min_seq, snapshot.next_seq), (4, 5));
assert!(snapshot.result_items_truncated);
assert!(snapshot.retained_bytes() < MAX_COMPLETED_HEAL_RESULT_BYTES);
let report = completed_task_report(&snapshot, Some(0));
assert_eq!(report.progress.expect("progress retained").objects_scanned, 9);
assert!(report.result_items_truncated);
let active_max = task.get_result_items_since(Some(u64::MAX)).await;
assert!(active_max.items.is_empty());
assert!(!active_max.lagged);
}
#[test]
fn completed_retention_result_byte_cap_and_cap_plus_one() {
for extra in [0, 1] {
let mut snapshot = completed_retention_fixture(SystemTime::now());
snapshot.seqed_items = vec![(
4,
HealResultItem {
detail: "x".repeat(MAX_COMPLETED_HEAL_RESULT_BYTES - size_of::<(u64, HealResultItem)>() + extra),
..Default::default()
},
)];
snapshot.min_seq = 4;
snapshot.bound_result_window();
assert_eq!(snapshot.seqed_items.len(), 1 - extra);
assert_eq!(snapshot.result_items_truncated, extra == 1);
assert_eq!(snapshot.min_seq, if extra == 0 { 4 } else { 5 });
assert_eq!(snapshot.next_seq, 5);
assert_eq!(snapshot.progress.as_ref().expect("progress retained").objects_scanned, 9);
}
}
#[derive(Default)]
struct CompletedRetentionHook {
started: Notify,
execute: Notify,
handoff: Notify,
finish: Notify,
pause_before_publish: bool,
before_publish: Notify,
publish: Notify,
prepared_status: Mutex<Option<HealTaskStatus>>,
}
static COMPLETED_RETENTION_HOOKS: LazyLock<Mutex<HashMap<String, Arc<CompletedRetentionHook>>>> =
LazyLock::new(|| Mutex::new(HashMap::new()));
pub(super) async fn pause_completed_retention_handoff(task_id: &str) {
let hook = COMPLETED_RETENTION_HOOKS.lock().await.get(task_id).cloned();
if let Some(hook) = hook {
hook.handoff.notify_one();
hook.finish.notified().await;
}
}
pub(super) async fn pause_completed_retention_before_publish(task_id: &str, status: &HealTaskStatus) {
let hook = COMPLETED_RETENTION_HOOKS.lock().await.get(task_id).cloned();
if let Some(hook) = hook.filter(|hook| hook.pause_before_publish) {
*hook.prepared_status.lock().await = Some(status.clone());
hook.before_publish.notify_one();
hook.publish.notified().await;
}
}
#[tokio::test]
async fn completed_retention_cancel_wins_over_a_prepared_retry_snapshot() {
let bucket = "completed-retention-retry-cancel";
let manager = HealManager::new(Arc::new(MockStorage), None);
let request = HealRequest::object(bucket.to_string(), "object".to_string(), None);
let task_id = request.id.clone();
let duplicate = HealRequest::object(bucket.to_string(), "object".to_string(), None);
let alias = duplicate.id.clone();
let hook = Arc::new(CompletedRetentionHook {
pause_before_publish: true,
..Default::default()
});
{
let mut hooks = COMPLETED_RETENTION_HOOKS.lock().await;
hooks.insert(bucket.to_string(), Arc::clone(&hook));
hooks.insert(task_id.clone(), Arc::clone(&hook));
}
manager.submit_heal_request(request).await.expect("admit original");
manager.submit_heal_request(duplicate).await.expect("admit alias");
process_manager_queue_once(&manager).await;
tokio::time::timeout(Duration::from_secs(5), hook.started.notified())
.await
.expect("scheduler starts");
let task = manager.active_heals.lock().await.get(&task_id).cloned().expect("active task");
task.progress.write().await.update_object_progress(1, 1, 0, 0, 4096);
hook.execute.notify_one();
tokio::time::timeout(Duration::from_secs(5), hook.before_publish.notified())
.await
.expect("retry snapshot prepared");
manager.cancel_task(&alias).await.expect("cancel wins active ownership");
assert!(matches!(*hook.prepared_status.lock().await, Some(HealTaskStatus::Retrying { .. })));
hook.publish.notify_one();
tokio::time::timeout(Duration::from_secs(5), hook.handoff.notified())
.await
.expect("scheduler finishes handoff");
for token in [&task_id, &alias] {
let report = manager.get_task_report(token).await.expect("cancelled token retained");
assert_eq!(report.status, HealTaskStatus::Cancelled);
assert_eq!(report.progress.expect("frozen progress").objects_scanned, 1);
}
assert!(!manager.retrying_heals.lock().await.contains_key(&task_id));
assert!(!manager.heal_queue.lock().await.contains_request_id(&task_id));
hook.finish.notify_one();
COMPLETED_RETENTION_HOOKS
.lock()
.await
.retain(|key, _| key != bucket && key != &task_id);
}
#[tokio::test]
async fn completed_retention_scheduler_preserves_progress_aliases_and_atomic_handoff() {
for outcome in ["success", "failed", "cancelled"] {
let bucket = format!("completed-retention-{outcome}");
let hook = Arc::new(CompletedRetentionHook::default());
let manager = Arc::new(HealManager::new(Arc::new(MockStorage), None));
let request = HealRequest::object(bucket.clone(), "object".to_string(), None);
let task_id = request.id.clone();
let duplicate = HealRequest::object(bucket.clone(), "object".to_string(), None);
let alias = duplicate.id.clone();
{
let mut hooks = COMPLETED_RETENTION_HOOKS.lock().await;
hooks.insert(bucket.clone(), Arc::clone(&hook));
hooks.insert(task_id.clone(), Arc::clone(&hook));
}
manager.submit_heal_request(request).await.expect("admit original");
manager.submit_heal_request(duplicate).await.expect("admit alias");
process_manager_queue_once(&manager).await;
tokio::time::timeout(Duration::from_secs(5), hook.started.notified())
.await
.expect("scheduler reaches storage");
let task = manager
.active_heals
.lock()
.await
.get(&task_id)
.cloned()
.expect("task is active");
task.progress.write().await.update_object_progress(1, 1, 0, 0, 4096);
let before = manager.get_task_report(&alias).await.expect("alias resolves active progress");
assert_eq!(before.progress.as_ref().expect("active progress").objects_scanned, 1);
let poll_manager = Arc::clone(&manager);
let poll_alias = alias.clone();
let stop = CancellationToken::new();
let poll_stop = stop.clone();
let polling = tokio::spawn(async move {
while !poll_stop.is_cancelled() {
let report = poll_manager
.get_task_report(&poll_alias)
.await
.expect("handoff must never return NotFound");
assert!(report.progress.expect("progress never disappears").objects_scanned >= 1);
tokio::task::yield_now().await;
}
});
if outcome == "cancelled" {
manager.cancel_task(&alias).await.expect("cancel active task by alias");
} else {
hook.execute.notify_one();
}
tokio::time::timeout(Duration::from_secs(5), hook.handoff.notified())
.await
.expect("scheduler archives terminal");
assert!(!manager.active_heals.lock().await.contains_key(&task_id));
let expected = task.get_progress().await;
for token in [&task_id, &alias] {
assert_eq!(manager.get_task_progress(token).await.expect("terminal progress query"), expected);
let report = manager
.get_task_report_for_path_since(&format!("{bucket}/object"), token, Some(u64::MAX))
.await
.expect("terminal token remains queryable at handoff");
assert_eq!(report.progress.as_ref(), Some(&expected));
assert!(report.result_items.is_empty());
match outcome {
"success" => assert_eq!(report.status, HealTaskStatus::Completed),
"failed" => assert!(matches!(report.status, HealTaskStatus::Failed { .. })),
_ => assert_eq!(report.status, HealTaskStatus::Cancelled),
}
}
let retained = manager.completed_heals.lock().await;
assert!(Arc::ptr_eq(&retained[&task_id], &retained[&alias]));
drop(retained);
stop.cancel();
polling.await.expect("concurrent polling succeeds");
// Archived progress must not alias a mutable live progress object.
task.progress.write().await.objects_scanned = 999;
assert_eq!(manager.get_task_report(&alias).await.expect("frozen report").progress, Some(expected));
hook.finish.notify_one();
COMPLETED_RETENTION_HOOKS
.lock()
.await
.retain(|key, _| key != &bucket && key != &task_id);
}
}
#[async_trait::async_trait]
impl HealStorageAPI for MockStorage {
async fn get_object_meta(&self, _bucket: &str, _object: &str) -> Result<Option<HealObjectInfo>> {
@@ -123,6 +443,12 @@ impl HealStorageAPI for MockStorage {
}
async fn object_exists(&self, bucket: &str, _object: &str) -> Result<bool> {
let hook = COMPLETED_RETENTION_HOOKS.lock().await.get(bucket).cloned();
if let Some(hook) = hook {
hook.started.notify_one();
hook.execute.notified().await;
return Ok(true);
}
Ok(bucket == "retry-transition")
}
@@ -133,13 +459,18 @@ impl HealStorageAPI for MockStorage {
_version_id: Option<&str>,
_opts: &HealOpts,
) -> Result<(HealResultItem, Option<Error>)> {
if bucket == "completed-retention-failed" {
return Err(Error::TaskExecutionFailed {
message: "retention fixture failure".to_string(),
});
}
if let Some(hook) = manager_recovery_test_hook() {
*hook
.heal_object_calls
.lock()
.expect("manager recovery object call lock should not poison") += 1;
}
if bucket == "retry-transition" {
if matches!(bucket, "retry-transition" | "completed-retention-retry-cancel") {
return Ok((
HealResultItem::default(),
Some(Error::Storage(EcstoreError::InsufficientReadQuorum(
@@ -1145,7 +1476,13 @@ async fn test_active_duplicate_token_can_query_and_cancel_original_task() {
.expect("duplicate token should cancel merged active task");
assert!(manager.active_heals.lock().await.get(&active_task_id).is_none());
assert!(matches!(manager.get_task_status(&active_task_id).await, Err(Error::TaskNotFound { .. })));
assert_eq!(
manager
.get_task_status(&active_task_id)
.await
.expect("cancelled task remains queryable"),
HealTaskStatus::Cancelled
);
}
#[tokio::test]
@@ -1638,6 +1975,8 @@ async fn insert_retrying_request(manager: &HealManager, request: HealRequest) ->
manager.completed_heals.lock().await.insert(
task_id,
Arc::new(CompletedHealStatus {
progress: None,
retained_bytes: std::sync::OnceLock::new(),
heal_type: request.heal_type,
status: HealTaskStatus::Retrying {
error: "Lock acquisition timeout".to_string(),
@@ -2053,7 +2392,7 @@ async fn admin_force_start_cancels_overlapping_active_task_first() {
"the overlapping admin task must be cancelled (removed from the active table) before the new one starts"
);
assert!(
matches!(manager.get_task_status(&old_id).await, Err(Error::TaskNotFound { .. })),
matches!(manager.get_task_status(&old_id).await, Ok(HealTaskStatus::Cancelled)),
"a cancelled task must no longer resolve as an active heal"
);
}
@@ -2360,6 +2699,8 @@ async fn test_retrying_completion_outranks_the_queue_for_the_same_id() {
manager.completed_heals.lock().await.insert(
task_id.clone(),
Arc::new(CompletedHealStatus {
progress: None,
retained_bytes: std::sync::OnceLock::new(),
heal_type: request.heal_type.clone(),
status: HealTaskStatus::Retrying {
error: "transient disk failure".to_string(),
@@ -2395,6 +2736,8 @@ async fn test_get_task_status_reads_recent_completed_status() {
manager.completed_heals.lock().await.insert(
"completed-token".to_string(),
Arc::new(CompletedHealStatus {
progress: None,
retained_bytes: std::sync::OnceLock::new(),
heal_type: HealType::Bucket {
bucket: "bucket".to_string(),
},
@@ -2424,6 +2767,8 @@ async fn test_get_task_report_for_path_reads_completed_items() {
manager.completed_heals.lock().await.insert(
"completed-token".to_string(),
Arc::new(CompletedHealStatus {
progress: None,
retained_bytes: std::sync::OnceLock::new(),
heal_type: HealType::Object {
bucket: "bucket".to_string(),
object: "object".to_string(),
+1 -1
View File
@@ -999,7 +999,7 @@ impl HealTask {
let items = match since {
None => result_items.iter().map(|(_, item)| item.clone()).collect::<Vec<_>>(),
Some(cursor) => {
if cursor + 1 < min_seq {
if cursor.saturating_add(1) < min_seq {
lagged = true;
}
result_items
+72
View File
@@ -182,6 +182,9 @@ pub struct BackgroundHealStatus {
pub heal_active_tasks: u64,
#[serde(default)]
pub cluster_status_complete: bool,
/// Missing on older servers; absent coverage or counts mean unknown.
#[serde(default)]
pub coverage: Option<BackgroundHealCoverage>,
#[serde(default)]
pub progress: Option<serde_json::Value>,
/// Remaining wire fields (flattened `BackgroundHealInfo` plus the
@@ -190,6 +193,22 @@ pub struct BackgroundHealStatus {
pub extra: serde_json::Map<String, serde_json::Value>,
}
/// Node coverage of a background heal status snapshot. Counters describe only
/// nodes with usable snapshots; unknown peers may still be running heal work.
#[derive(Debug, Clone, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct BackgroundHealCoverage {
#[serde(default)]
pub expected: Option<usize>,
#[serde(default)]
pub responded: Option<usize>,
#[serde(default)]
pub unknown: Option<usize>,
/// Stable reason codes; unknown future codes are preserved verbatim.
#[serde(default)]
pub reasons: Vec<String>,
}
/// `GET /v3/scanner/status` response, typed at the fields operators branch
/// on; everything else passes through verbatim.
#[derive(Debug, Clone, Deserialize)]
@@ -630,9 +649,39 @@ mod tests {
assert_eq!(status.state, "active");
assert_eq!(status.heal_queue_length, 3);
assert!(status.cluster_status_complete);
assert!(status.coverage.is_none(), "legacy payloads have unknown coverage");
assert!(status.extra.contains_key("healOperations"), "unknown nested payloads must pass through");
}
#[test]
fn background_heal_status_missing_coverage_fields_remain_unknown() {
for raw in [json!({"state": "degraded"}), json!({"state": "degraded", "coverage": {}})] {
let status: BackgroundHealStatus = serde_json::from_value(raw).expect("partial legacy payload decodes");
assert!(!status.cluster_status_complete);
if let Some(coverage) = status.coverage {
assert_eq!(coverage.expected, None);
assert_eq!(coverage.responded, None);
assert_eq!(coverage.unknown, None);
}
}
}
#[test]
fn background_heal_status_preserves_future_fields_and_reasons() {
let raw = json!({
"state": "degraded", "clusterStatusComplete": false,
"coverage": {"expected": 3, "responded": 1, "unknown": 2, "reasons": ["future_reason"], "futureCoverage": true},
"futureStatus": {"value": 7}
});
let status: BackgroundHealStatus = serde_json::from_value(raw).expect("future additive fields decode");
assert_eq!(status.extra["futureStatus"]["value"], 7);
let coverage = status.coverage.expect("coverage supplied");
assert_eq!(coverage.expected, Some(3));
assert_eq!(coverage.responded, Some(1));
assert_eq!(coverage.unknown, Some(2));
assert_eq!(coverage.reasons, ["future_reason"]);
}
#[test]
fn scanner_status_defaults_freshness_to_unknown() {
let raw = json!({"enabled": true, "freshness": {"state": "stale"}, "metrics": {}});
@@ -721,6 +770,7 @@ mod tests {
let status = client.background_heal_status().await.expect("status decodes");
assert_eq!(status.state, "idle");
assert!(status.coverage.is_none(), "older HTTP responses retain unknown coverage");
let request = server.recorded();
// The server registers this route POST-only; a GET here answers 405.
assert_eq!(request.method, "POST");
@@ -728,6 +778,28 @@ mod tests {
assert_eq!(request.query, "");
}
#[tokio::test]
async fn background_heal_status_decodes_partial_coverage_over_http() {
let body = r#"{"state":"degraded","healQueueLength":0,"healActiveTasks":0,"clusterStatusComplete":false,"coverage":{"expected":3,"responded":1,"unknown":2,"reasons":["notification_system_unavailable"]},"futureStatus":true}"#;
let server = TestServer::spawn(body, 200).await;
let client = AdminClient::new(&format!("http://{}", server.addr), "ak", "sk").expect("client builds");
let status = client
.background_heal_status()
.await
.expect("partial status is a successful response");
assert_eq!(status.state, "degraded");
assert!(!status.cluster_status_complete);
assert_eq!(status.extra["futureStatus"], true);
let coverage = status.coverage.expect("partial coverage supplied");
assert_eq!(coverage.expected, Some(3));
assert_eq!(coverage.responded, Some(1));
assert_eq!(coverage.unknown, Some(2));
assert_eq!(coverage.reasons, ["notification_system_unavailable"]);
let request = server.recorded();
assert_eq!(request.method, "POST");
assert_eq!(request.query, "", "reading status must not send heal control parameters");
}
#[tokio::test]
async fn http_error_status_maps_to_a_typed_error_with_body() {
let server = TestServer::spawn(r#"{"code":"AccessDenied","message":"denied"}"#, 403).await;
+15 -8
View File
@@ -1703,14 +1703,18 @@ where
let (sender, receiver) = mpsc::channel::<DataUsageInfo>(1);
let done_cycle = Metrics::time(Metric::ScanCycle);
let scan_result = crate::scanner_io::nsscanner_with_storage_status(
let scan_result = crate::scanner_io::nsscanner_with_storage_status_scoped(
storeapi.as_ref(),
cycle_budget.token(),
cycle_budget.clone(),
sender,
cycle_info.current,
leader_epoch,
scan_mode,
crate::scanner_io::ScannerCycleRequest {
ctx: cycle_budget.token(),
budget: cycle_budget.clone(),
updates: sender,
want_cycle: cycle_info.current,
leader_epoch,
scan_mode,
scan_scope: crate::scanner_io::ScannerBucketScanScope::default(),
persisted_usage_baseline: usage_persist_baseline.data.clone(),
},
)
.await;
let publication_defer_reason = match &scan_result {
@@ -3424,10 +3428,13 @@ use cycle_state::*;
use leadership::*;
use usage_store::*;
#[cfg(test)]
pub(crate) use activity::scanner_activity_snapshot_digest;
pub use activity::scanner_topology_digest;
pub(crate) use activity::{
ScannerActivitySnapshot, ScannerDirtyUsageAcknowledgement, probe_scanner_activity, scanner_activity_allows_usage_publication,
scanner_activity_publication_lease_targets, scanner_activity_snapshot_digest, scanner_dirty_usage_acknowledgements,
scanner_activity_dirty_usage_state_for_host, scanner_activity_publication_lease_targets, scanner_activity_structural_digest,
scanner_dirty_usage_acknowledgements,
};
pub(crate) use activity::{ScannerCycleOutcome, scanner_cycle_outcome_with_pending_maintenance};
pub use backlog::{
+41
View File
@@ -902,6 +902,7 @@ where
observation
}
#[cfg(test)]
pub(crate) fn scanner_activity_snapshot_digest(snapshot: &ScannerActivitySnapshot) -> [u8; 32] {
let mut hasher = Sha256::new();
hasher.update(u64::try_from(snapshot.len()).unwrap_or(u64::MAX).to_be_bytes());
@@ -925,6 +926,30 @@ pub(crate) fn scanner_activity_snapshot_digest(snapshot: &ScannerActivitySnapsho
hasher.finalize().into()
}
/// Hash the activity inputs that make an existing scanner cache unsafe to
/// reuse. Regular namespace writes and dirty-usage generations are omitted:
/// their affected buckets are tracked separately and may be refreshed from a
/// complete authoritative cache baseline.
pub(crate) fn scanner_activity_structural_digest(snapshot: &ScannerActivitySnapshot) -> [u8; 32] {
let mut hasher = Sha256::new();
hasher.update(u64::try_from(snapshot.len()).unwrap_or(u64::MAX).to_be_bytes());
for (host, activity) in snapshot {
let host = host.as_bytes();
let instance_id = activity.instance_id.as_bytes();
hasher.update(u64::try_from(host.len()).unwrap_or(u64::MAX).to_be_bytes());
hasher.update(host);
hasher.update(u64::try_from(instance_id.len()).unwrap_or(u64::MAX).to_be_bytes());
hasher.update(instance_id);
hasher.update(activity.maintenance_generation.to_be_bytes());
hasher.update(activity.protocol_version.to_be_bytes());
hasher.update(activity.topology_digest);
hasher.update([u8::from(activity.data_movement_active)]);
hasher.update(activity.movement_generation.to_be_bytes());
hasher.update([u8::from(activity.publication_blocked)]);
}
hasher.finalize().into()
}
pub(crate) fn scanner_activity_allows_usage_publication(snapshot: &ScannerActivitySnapshot) -> bool {
!snapshot.is_empty()
&& snapshot.values().all(|activity| {
@@ -955,6 +980,22 @@ pub(crate) fn scanner_dirty_usage_acknowledgements(snapshot: &ScannerActivitySna
.collect()
}
pub(crate) fn scanner_activity_dirty_usage_state_for_host<'a>(
snapshot: &'a ScannerActivitySnapshot,
host: &str,
) -> Option<(&'a str, u64, bool)> {
snapshot
.get(host)
.filter(|_| host != LOCAL_SCANNER_ACTIVITY_NODE)
.map(|activity| {
(
activity.instance_id.as_str(),
activity.dirty_usage_generation,
activity.dirty_usage_pending,
)
})
}
pub fn scanner_topology_digest(storeapi: &ECStore) -> [u8; 32] {
let endpoint_pools = storeapi.endpoints();
let mut hasher = Sha256::new();
+280 -51
View File
@@ -379,6 +379,12 @@ pub(super) fn decode_recovery_marker_for_reset(
if !matches!(marker_revision, DataUsageCacheRevision::Etag(_)) {
return Err(ScannerError::Other("cycle recovery marker has no object revision".to_string()));
}
if let Ok(value) = serde_json::from_slice::<serde_json::Value>(data)
&& let Some(state) = value.get("state")
&& !matches!(state.as_str(), Some("blocked" | "cleanup-pending"))
{
return Err(ScannerError::Other("cycle recovery marker state is unsupported".to_string()));
}
let compat = serde_json::from_slice::<ScannerCycleRecoveryMarkerCompat>(data).ok();
let _schema_version = compat.as_ref().and_then(|marker| marker.schema_version);
let primary_revision = compat
@@ -406,7 +412,10 @@ pub(super) fn decode_recovery_marker_for_reset(
};
let state = match compat.as_ref().and_then(|marker| marker.state.as_deref()) {
Some("cleanup-pending") => "cleanup-pending",
_ => "blocked",
Some("blocked") | None => "blocked",
Some(_) => {
return Err(ScannerError::Other("cycle recovery marker state is unsupported".to_string()));
}
};
let now = unix_now_secs();
Ok(ScannerCycleRecoveryMarker {
@@ -721,17 +730,19 @@ async fn mark_cycle_recovery_cleanup_pending(
mut marker: ScannerCycleRecoveryMarker,
marker_revision: &DataUsageCacheRevision,
expected_epoch: u64,
owns_reset: &(impl Fn() -> bool + Sync),
) -> Result<(ScannerCycleRecoveryMarker, DataUsageCacheRevision), ScannerError> {
marker.state = "cleanup-pending".to_string();
marker.last_attempt_at_unix_secs = unix_now_secs();
let bytes = serde_json::to_vec(&marker)
.map_err(|err| ScannerError::Other(format!("failed to encode cycle recovery marker: {err}")))?;
let info = save_config_with_publication_admission_for_epoch(
let info = save_reset_config(
storeapi.clone(),
DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(),
bytes,
marker_revision.preconditions(),
expected_epoch,
owns_reset,
)
.await
.map_err(|err| ScannerError::Other(format!("failed to mark cycle recovery cleanup pending: {err}")))?;
@@ -933,6 +944,7 @@ pub async fn reset_scanner_cycle_recovery(ctx: CancellationToken, storeapi: Arc<
.get_write_lock_quiet(Duration::from_secs(5))
.await
.map_err(|err| ScannerError::Other(format!("scanner leader lock is busy: {err}")))?;
let owns_reset = || !guard.is_lock_lost() && !ctx.is_cancelled();
if guard.is_lock_lost() {
return Err(ScannerError::Other("scanner leader lock was lost before recovery reset".to_string()));
@@ -952,7 +964,27 @@ pub async fn reset_scanner_cycle_recovery(ctx: CancellationToken, storeapi: Arc<
}
Err(err) => return Err(ScannerError::Other(format!("failed to read cycle recovery marker: {err}"))),
};
let marker_data = marker_data.ok_or_else(|| ScannerError::Other("scanner cycle recovery marker is absent".to_string()))?;
let Some(marker_data) = marker_data else {
// A delete may commit before its reply is lost. Confirm both durable
// fences before treating a retry without its marker as completed.
let (cycle, epoch, revision) = read_cycle_state_for_usage_reset(storeapi.clone()).await?;
let floor = persisted_usage_floor(storeapi.clone()).await?;
if !matches!(revision, DataUsageCacheRevision::Etag(_))
|| epoch < floor.leader_epoch
|| cycle.next < floor.next_cycle
|| !owns_reset()
|| scanner_publication_admission_for_epoch(storeapi.clone(), reset_epoch)
.await
.is_none()
{
return Err(ScannerError::Other(
"scanner cycle recovery marker is absent without a completed reset fence".to_string(),
));
}
set_scanner_cycle_recovery_status(recovery_status("healthy", None, false));
super::notify_scanner_cycle_recovery_wake();
return Ok(());
};
let (marker, force_full_rescan) = match serde_json::from_slice::<ScannerCycleRecoveryMarker>(&marker_data) {
Ok(marker) if validate_recovery_marker(&marker).is_ok() => (marker, false),
_ => (decode_recovery_marker_for_reset(&marker_data, &marker_revision)?, true),
@@ -1026,8 +1058,10 @@ pub async fn reset_scanner_cycle_recovery(ctx: CancellationToken, storeapi: Arc<
}
};
if let Some((primary_cycle, primary_epoch)) = primary_state {
verify_cycle_reset_intent(storeapi.clone(), &marker_revision, &owns_reset).await?;
let (cleanup_marker, cleanup_marker_revision) =
mark_cycle_recovery_cleanup_pending(storeapi.clone(), marker.clone(), &marker_revision, reset_epoch).await?;
mark_cycle_recovery_cleanup_pending(storeapi.clone(), marker.clone(), &marker_revision, reset_epoch, &owns_reset)
.await?;
set_scanner_cycle_recovery_status(recovery_status_from_marker(&cleanup_marker, "cleanup-pending"));
let usage_floor = persisted_usage_floor(storeapi.clone()).await?;
let fence_epoch = primary_epoch
@@ -1047,12 +1081,14 @@ pub async fn reset_scanner_cycle_recovery(ctx: CancellationToken, storeapi: Arc<
"preserved scanner cycle state exceeds the bounded object size".to_string(),
));
}
let preserved_info = save_config_with_publication_admission_for_epoch(
verify_cycle_reset_intent(storeapi.clone(), &cleanup_marker_revision, &owns_reset).await?;
let preserved_info = save_reset_config(
storeapi.clone(),
DATA_USAGE_BLOOM_NAME_PATH.as_str(),
preserved_data,
primary_revision.preconditions(),
reset_epoch,
&owns_reset,
)
.await
.map_err(|err| {
@@ -1072,9 +1108,17 @@ pub async fn reset_scanner_cycle_recovery(ctx: CancellationToken, storeapi: Arc<
"scanner leader lock was lost after fencing newer cycle state".to_string(),
));
}
fence_scanner_usage_epoch_with_expected_epoch(&ctx, storeapi.clone(), fence_epoch, Some(reset_epoch), false)
.await
.map_err(|err| ScannerError::Other(format!("failed to fence preserved scanner usage epoch: {err}")))?;
verify_cycle_reset_intent(storeapi.clone(), &cleanup_marker_revision, &owns_reset).await?;
fence_scanner_usage_epoch_with_expected_epoch(
&ctx,
storeapi.clone(),
fence_epoch,
Some(reset_epoch),
false,
&owns_reset,
)
.await
.map_err(|err| ScannerError::Other(format!("failed to fence preserved scanner usage epoch: {err}")))?;
if guard.is_lock_lost() {
return Err(ScannerError::Other(
"scanner leader lock was lost after fencing newer cycle state".to_string(),
@@ -1088,7 +1132,8 @@ pub async fn reset_scanner_cycle_recovery(ctx: CancellationToken, storeapi: Arc<
"scanner cycle state changed before recovery marker cleanup".to_string(),
));
}
delete_config_with_publication_admission_for_epoch(
verify_cycle_reset_intent(storeapi.clone(), &cleanup_marker_revision, &owns_reset).await?;
delete_reset_config(
storeapi.clone(),
RUSTFS_META_BUCKET,
DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(),
@@ -1100,6 +1145,7 @@ pub async fn reset_scanner_cycle_recovery(ctx: CancellationToken, storeapi: Arc<
..Default::default()
},
reset_epoch,
&owns_reset,
)
.await
.map_err(|err| {
@@ -1149,17 +1195,20 @@ pub async fn reset_scanner_cycle_recovery(ctx: CancellationToken, storeapi: Arc<
// Persist the cleanup-pending phase before rewriting the primary. If the
// process dies after the rewrite, startup still sees a durable fence and
// cannot mistake the partially completed reset for a healthy state.
verify_cycle_reset_intent(storeapi.clone(), &marker_revision, &owns_reset).await?;
let (marker, marker_revision) = if marker.state == "cleanup-pending" {
(marker, marker_revision)
} else {
mark_cycle_recovery_cleanup_pending(storeapi.clone(), marker, &marker_revision, reset_epoch).await?
mark_cycle_recovery_cleanup_pending(storeapi.clone(), marker, &marker_revision, reset_epoch, &owns_reset).await?
};
let rebuilt_info = save_config_with_publication_admission_for_epoch(
verify_cycle_reset_intent(storeapi.clone(), &marker_revision, &owns_reset).await?;
let rebuilt_info = save_reset_config(
storeapi.clone(),
DATA_USAGE_BLOOM_NAME_PATH.as_str(),
data,
primary_revision.preconditions(),
reset_epoch,
&owns_reset,
)
.await
.map_err(|err| {
@@ -1178,8 +1227,10 @@ pub async fn reset_scanner_cycle_recovery(ctx: CancellationToken, storeapi: Arc<
"scanner leader lock was lost after rebuilding cycle state".to_string(),
));
}
verify_cycle_reset_intent(storeapi.clone(), &marker_revision, &owns_reset).await?;
if let Err(err) =
fence_scanner_usage_epoch_with_expected_epoch(&ctx, storeapi.clone(), leader_epoch, Some(reset_epoch), false).await
fence_scanner_usage_epoch_with_expected_epoch(&ctx, storeapi.clone(), leader_epoch, Some(reset_epoch), false, &owns_reset)
.await
{
set_scanner_cycle_recovery_status(ScannerCycleRecoveryStatus {
path: DATA_USAGE_BLOOM_NAME_PATH.clone(),
@@ -1249,7 +1300,8 @@ pub async fn reset_scanner_cycle_recovery(ctx: CancellationToken, storeapi: Arc<
));
}
if let Err(err) = delete_config_with_publication_admission_for_epoch(
verify_cycle_reset_intent(storeapi.clone(), &marker_revision, &owns_reset).await?;
if let Err(err) = delete_reset_config(
storeapi.clone(),
RUSTFS_META_BUCKET,
DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(),
@@ -1261,6 +1313,7 @@ pub async fn reset_scanner_cycle_recovery(ctx: CancellationToken, storeapi: Arc<
..Default::default()
},
reset_epoch,
&owns_reset,
)
.await
{
@@ -1310,6 +1363,57 @@ pub async fn reset_scanner_cycle_recovery(ctx: CancellationToken, storeapi: Arc<
Ok(())
}
async fn verify_cycle_reset_intent(
storeapi: Arc<impl ScannerObjectIO>,
expected_revision: &DataUsageCacheRevision,
owns_reset: &(impl Fn() -> bool + Sync),
) -> Result<(), ScannerError> {
let revision = read_config_revision(storeapi, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str())
.await
.map_err(|err| ScannerError::Other(format!("failed to verify scanner cycle reset intent: {err}")))?;
if &revision != expected_revision {
return Err(ScannerError::Other("scanner cycle reset intent changed".to_string()));
}
if !owns_reset() {
return Err(ScannerError::Other("scanner cycle reset ownership was lost".to_string()));
}
Ok(())
}
async fn save_reset_config(
storeapi: Arc<impl ScannerObjectIO + ScannerConfigObjectDelete>,
path: &str,
data: Vec<u8>,
preconditions: crate::HTTPPreconditions,
expected_epoch: u64,
owns_reset: &(impl Fn() -> bool + Sync),
) -> Result<crate::ScannerObjectInfo, EcstoreError> {
let Some(_admission) = scanner_publication_admission_for_epoch(storeapi.clone(), expected_epoch).await else {
return Err(EcstoreError::other(SCANNER_PUBLICATION_EPOCH_CHANGED));
};
if !owns_reset() {
return Err(EcstoreError::other("scanner reset ownership was lost before write"));
}
save_config_with_preconditions(storeapi, path, data, preconditions).await
}
async fn delete_reset_config(
storeapi: Arc<impl ScannerObjectIO + ScannerConfigObjectDelete>,
bucket: &str,
path: &str,
options: ScannerObjectOptions,
expected_epoch: u64,
owns_reset: &(impl Fn() -> bool + Sync),
) -> Result<crate::ScannerObjectInfo, EcstoreError> {
let Some(_admission) = scanner_publication_admission_for_epoch(storeapi.clone(), expected_epoch).await else {
return Err(EcstoreError::other(SCANNER_PUBLICATION_EPOCH_CHANGED));
};
if !owns_reset() {
return Err(EcstoreError::other("scanner reset ownership was lost before delete"));
}
storeapi.delete_config_object(bucket, path, options).await
}
fn scanner_usage_state_reset_paths() -> Vec<String> {
vec![
DATA_USAGE_OBJ_NAME_PATH.as_str().to_string(),
@@ -1333,8 +1437,14 @@ pub(super) async fn read_usage_state_reset_slots(
Ok(slots)
}
fn usage_state_reset_floor(slots: &[ScannerUsageStateResetSlot]) -> Result<PersistedUsageFloor, ScannerError> {
let mut floor = PersistedUsageFloor::default();
enum ScannerUsageResetFloor {
Missing,
Trusted(PersistedUsageFloor),
Corrupt,
}
fn usage_state_reset_floor(slots: &[ScannerUsageStateResetSlot]) -> Result<ScannerUsageResetFloor, ScannerError> {
let mut floor = None;
for slot in slots {
let Some(data) = slot.data.as_deref() else {
continue;
@@ -1342,9 +1452,21 @@ fn usage_state_reset_floor(slots: &[ScannerUsageStateResetSlot]) -> Result<Persi
let Ok(usage) = serde_json::from_slice::<DataUsageInfo>(data) else {
continue;
};
update_persisted_usage_floor(&mut floor, &usage, &slot.path)?;
if !data_usage_info_has_persisted_baseline_identity(&usage)
&& !(slot.path == DATA_USAGE_OBJ_NAME_PATH.as_str() && data_usage_info_is_bootstrap_pending(&usage))
&& legacy_incomplete_usage_fence(data, &usage)
.and_then(|fence| fence.claimable_epoch())
.is_none()
{
continue;
}
update_persisted_usage_floor(floor.get_or_insert_with(PersistedUsageFloor::default), &usage, &slot.path)?;
}
Ok(floor)
Ok(match floor {
Some(floor) => ScannerUsageResetFloor::Trusted(floor),
None if slots.iter().any(|slot| slot.data.is_some()) => ScannerUsageResetFloor::Corrupt,
None => ScannerUsageResetFloor::Missing,
})
}
async fn read_cycle_state_for_usage_reset(
@@ -1401,11 +1523,12 @@ async fn delete_usage_state_reset_slot(
storeapi: Arc<impl ScannerObjectIO + ScannerConfigObjectDelete>,
slot: &ScannerUsageStateResetSlot,
expected_epoch: u64,
owns_reset: &(impl Fn() -> bool + Sync),
) -> Result<bool, ScannerError> {
if matches!(slot.revision, DataUsageCacheRevision::Missing) {
return Ok(false);
}
let delete_result = delete_config_with_publication_admission_for_epoch(
let delete_result = delete_reset_config(
storeapi.clone(),
RUSTFS_META_BUCKET,
&slot.path,
@@ -1415,6 +1538,7 @@ async fn delete_usage_state_reset_slot(
..Default::default()
},
expected_epoch,
owns_reset,
)
.await;
match delete_result {
@@ -1486,21 +1610,24 @@ pub(super) async fn publish_scanner_usage_bootstrap_primary(
expected_publication_epoch: u64,
leader_epoch: Option<u64>,
context: ScannerUsageBootstrapPublishContext,
owns_publication: impl Fn() -> bool + Sync,
) -> Result<(), ScannerError> {
async fn inner(
storeapi: Arc<impl ScannerObjectIO + ScannerConfigObjectDelete>,
expected_revision: &DataUsageCacheRevision,
expected_publication_epoch: u64,
leader_epoch: Option<u64>,
owns_publication: &(impl Fn() -> bool + Sync),
) -> Result<(), ScannerUsageBootstrapPublishError> {
let marker = scanner_usage_bootstrap_marker(std::time::SystemTime::now(), leader_epoch);
let data = serde_json::to_vec(&marker).map_err(ScannerUsageBootstrapPublishError::Encode)?;
let save_result = save_config_with_publication_admission_for_epoch(
let save_result = save_reset_config(
storeapi.clone(),
DATA_USAGE_OBJ_NAME_PATH.as_str(),
data.clone(),
expected_revision.preconditions(),
expected_publication_epoch,
owns_publication,
)
.await;
if save_result
@@ -1524,7 +1651,7 @@ pub(super) async fn publish_scanner_usage_bootstrap_primary(
})
}
inner(storeapi, expected_revision, expected_publication_epoch, leader_epoch)
inner(storeapi, expected_revision, expected_publication_epoch, leader_epoch, &owns_publication)
.await
.map_err(|err| err.into_scanner_error(context))
}
@@ -1534,32 +1661,108 @@ pub(super) async fn reset_scanner_usage_state_slots_for_full_rebuild(
slots: &[ScannerUsageStateResetSlot],
expected_epoch: u64,
leader_epoch: u64,
owns_reset: impl Fn() -> bool + Sync,
) -> Result<Vec<String>, ScannerError> {
let mut reset_paths = Vec::new();
let primary = slots
.iter()
.find(|slot| slot.path == DATA_USAGE_OBJ_NAME_PATH.as_str())
.ok_or_else(|| ScannerError::Other("scanner usage reset primary slot was not inspected".to_string()))?;
publish_scanner_usage_bootstrap_primary(
storeapi.clone(),
&primary.revision,
expected_epoch,
Some(leader_epoch),
ScannerUsageBootstrapPublishContext::Reset,
)
.await?;
if !owns_reset() {
return Err(ScannerError::Other("scanner usage reset ownership was lost".to_string()));
}
let resume_epoch = usage_state_reset_resume_epoch(slots)?;
match resume_epoch {
Some(epoch) if epoch == leader_epoch => {}
Some(_) => return Err(ScannerError::Other("scanner usage reset bootstrap epoch changed".to_string())),
None => {
publish_scanner_usage_bootstrap_primary(
storeapi.clone(),
&primary.revision,
expected_epoch,
Some(leader_epoch),
ScannerUsageBootstrapPublishContext::Reset,
&owns_reset,
)
.await?;
}
}
let (data, intent_revision) = read_config_with_revision(storeapi.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str())
.await
.map_err(|err| ScannerError::Other(format!("failed to inspect scanner usage reset intent: {err}")))?;
data.as_deref()
.and_then(|data| serde_json::from_slice::<DataUsageInfo>(data).ok())
.filter(|usage| data_usage_info_is_bootstrap_pending(usage) && usage.scanner_epoch == Some(leader_epoch))
.ok_or_else(|| ScannerError::Other("scanner usage reset intent changed before cleanup".to_string()))?;
if !matches!(intent_revision, DataUsageCacheRevision::Etag(_))
|| (resume_epoch.is_some() && intent_revision != primary.revision)
{
return Err(ScannerError::Other("scanner usage reset intent revision changed".to_string()));
}
reset_paths.push(DATA_USAGE_OBJ_NAME_PATH.as_str().to_string());
for slot in slots.iter().filter(|slot| slot.path != DATA_USAGE_OBJ_NAME_PATH.as_str()) {
if delete_usage_state_reset_slot(storeapi.clone(), slot, expected_epoch).await? {
if let Some(usage) = slot
.data
.as_deref()
.and_then(|data| serde_json::from_slice::<DataUsageInfo>(data).ok())
&& usage_epoch(&usage) >= leader_epoch
{
return Err(ScannerError::Other(format!(
"scanner usage reset slot is not older than its intent: {}",
slot.path
)));
}
let revision = read_config_revision(storeapi.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str())
.await
.map_err(|err| ScannerError::Other(format!("failed to verify scanner usage reset intent: {err}")))?;
if revision != intent_revision {
return Err(ScannerError::Other("scanner usage reset intent changed during cleanup".to_string()));
}
if !owns_reset() {
return Err(ScannerError::Other("scanner usage reset ownership was lost".to_string()));
}
if delete_usage_state_reset_slot(storeapi.clone(), slot, expected_epoch, &owns_reset).await? {
reset_paths.push(slot.path.clone());
}
}
let revision = read_config_revision(storeapi.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str())
.await
.map_err(|err| ScannerError::Other(format!("failed to confirm scanner usage reset intent: {err}")))?;
if revision != intent_revision || !owns_reset() {
return Err(ScannerError::Other(
"scanner usage reset intent or ownership changed before completion".to_string(),
));
}
invalidate_admin_data_usage_snapshot_cache().await;
invalidate_data_usage_snapshot_cache().await;
Ok(reset_paths)
}
fn usage_state_reset_resume_epoch(slots: &[ScannerUsageStateResetSlot]) -> Result<Option<u64>, ScannerError> {
let primary = slots.iter().find(|slot| slot.path == DATA_USAGE_OBJ_NAME_PATH.as_str());
let usage = primary
.and_then(|slot| slot.data.as_deref())
.and_then(|data| serde_json::from_slice::<DataUsageInfo>(data).ok());
match usage {
Some(usage) if usage.usage_snapshot_bootstrap_pending => {
if !data_usage_info_is_bootstrap_pending(&usage) {
return Err(ScannerError::Other("scanner usage reset bootstrap is invalid".to_string()));
}
if usage.scanner_epoch.is_none() {
// Initial bootstrap has no reset owner yet.
return Ok(None);
}
usage
.scanner_epoch
.filter(|epoch| *epoch > 0 && *epoch < u64::MAX)
.map(Some)
.ok_or_else(|| ScannerError::Other("scanner usage reset bootstrap has no valid epoch".to_string()))
}
_ => Ok(None),
}
}
pub async fn reset_scanner_usage_state_for_full_rebuild(
ctx: CancellationToken,
storeapi: Arc<ECStore>,
@@ -1584,12 +1787,31 @@ pub async fn reset_scanner_usage_state_for_full_rebuild(
};
let (cycle, cycle_epoch, cycle_revision) = read_cycle_state_for_usage_reset(storeapi.clone()).await?;
let slots = read_usage_state_reset_slots(storeapi.clone()).await?;
let usage_floor = usage_state_reset_floor(&slots)?;
let leader_epoch = cycle_epoch
.max(usage_floor.leader_epoch)
.checked_add(1)
.filter(|epoch| *epoch < u64::MAX)
.ok_or_else(|| ScannerError::Other("scanner leader epoch is exhausted".to_string()))?;
let usage_floor = match usage_state_reset_floor(&slots)? {
ScannerUsageResetFloor::Trusted(floor) => floor,
ScannerUsageResetFloor::Corrupt if matches!(cycle_revision, DataUsageCacheRevision::Missing) => {
return Err(ScannerError::Other("scanner usage reset has no trusted cycle or usage floor".to_string()));
}
ScannerUsageResetFloor::Missing | ScannerUsageResetFloor::Corrupt => PersistedUsageFloor {
next_cycle: cycle.next,
leader_epoch: cycle_epoch,
},
};
let resume_epoch = usage_state_reset_resume_epoch(&slots)?;
let leader_epoch = if let Some(epoch) = resume_epoch {
if epoch != cycle_epoch || usage_floor.leader_epoch > epoch || usage_floor.next_cycle > cycle.next {
return Err(ScannerError::Other(
"scanner usage reset bootstrap conflicts with the persisted cycle fence".to_string(),
));
}
epoch
} else {
cycle_epoch
.max(usage_floor.leader_epoch)
.checked_add(1)
.filter(|epoch| *epoch < u64::MAX)
.ok_or_else(|| ScannerError::Other("scanner leader epoch is exhausted".to_string()))?
};
let rebuilt_cycle = CurrentCycle {
next: cycle.next.max(usage_floor.next_cycle),
..Default::default()
@@ -1602,21 +1824,24 @@ pub async fn reset_scanner_usage_state_for_full_rebuild(
"scanner leader lock was lost before fencing usage reset cycle state".to_string(),
));
}
save_config_with_publication_admission_for_epoch(
storeapi.clone(),
DATA_USAGE_BLOOM_NAME_PATH.as_str(),
cycle_data,
cycle_revision.preconditions(),
reset_epoch,
)
.await
.map_err(|err| {
if scanner_publication_epoch_changed(&err) {
ScannerError::Other("scanner usage reset deferred by a movement epoch change".to_string())
} else {
ScannerError::Other(format!("failed to fence scanner cycle state for usage reset: {err}"))
}
})?;
if resume_epoch.is_none() {
save_reset_config(
storeapi.clone(),
DATA_USAGE_BLOOM_NAME_PATH.as_str(),
cycle_data,
cycle_revision.preconditions(),
reset_epoch,
&|| !guard.is_lock_lost() && !ctx.is_cancelled(),
)
.await
.map_err(|err| {
if scanner_publication_epoch_changed(&err) {
ScannerError::Other("scanner usage reset deferred by a movement epoch change".to_string())
} else {
ScannerError::Other(format!("failed to fence scanner cycle state for usage reset: {err}"))
}
})?;
}
if guard.is_lock_lost() {
return Err(ScannerError::Other(
@@ -1624,7 +1849,10 @@ pub async fn reset_scanner_usage_state_for_full_rebuild(
));
}
let reset_paths =
reset_scanner_usage_state_slots_for_full_rebuild(storeapi.clone(), &slots, reset_epoch, leader_epoch).await?;
reset_scanner_usage_state_slots_for_full_rebuild(storeapi.clone(), &slots, reset_epoch, leader_epoch, || {
!guard.is_lock_lost() && !ctx.is_cancelled()
})
.await?;
if guard.is_lock_lost() {
return Err(ScannerError::Other(
"scanner leader lock was lost after publishing usage reset marker".to_string(),
@@ -2135,6 +2363,7 @@ async fn recover_legacy_incomplete_usage_floor(
expected_publication_epoch,
Some(primary.epoch),
ScannerUsageBootstrapPublishContext::Recovery,
|| true,
)
.await?;
warn!(
+7 -1
View File
@@ -191,6 +191,7 @@ pub(super) async fn initialize_usage_baseline_bootstrap(
expected_epoch,
None,
ScannerUsageBootstrapPublishContext::Initial,
|| true,
)
.await
}
@@ -201,9 +202,10 @@ pub(super) async fn fence_scanner_usage_epoch_with_expected_epoch(
claimed_epoch: u64,
expected_publication_epoch: Option<u64>,
allow_bootstrap_pending: bool,
owns_fence: impl Fn() -> bool,
) -> Result<(), ScannerError> {
for retry in 0..=SCANNER_PERSIST_CAS_RETRIES {
if ctx.is_cancelled() {
if ctx.is_cancelled() || !owns_fence() {
return Err(ScannerError::Other("scanner leadership was cancelled before usage fencing".to_string()));
}
@@ -264,6 +266,9 @@ pub(super) async fn fence_scanner_usage_epoch_with_expected_epoch(
"scanner usage epoch fence changed while preparing its conditional write".to_string(),
));
};
if ctx.is_cancelled() || !owns_fence() {
return Err(ScannerError::Other("scanner leadership was lost before usage fencing".to_string()));
}
save_config_with_preconditions(storeapi.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str(), data, revision.preconditions())
.await
};
@@ -319,6 +324,7 @@ pub(super) async fn complete_scanner_leadership_claim(
claimed_epoch,
expected_publication_epoch,
allow_bootstrap_pending,
|| true,
)
.await
{
+415 -3
View File
@@ -634,6 +634,8 @@ struct MemoryConfigStore {
cancel_after_successful_puts: Mutex<HashMap<String, (usize, CancellationToken)>>,
replace_after_successful_puts: Mutex<HashMap<String, (usize, Vec<u8>)>>,
error_after_commit_deletes: Mutex<HashSet<String>>,
cancel_after_deletes: Mutex<HashMap<String, CancellationToken>>,
pause_next_publication_admission: Mutex<Option<(Arc<tokio::sync::Notify>, Arc<tokio::sync::Notify>)>>,
put_counts: Mutex<HashMap<String, usize>>,
publication_admission_blocked: AtomicBool,
block_publication_after_admissions: AtomicUsize,
@@ -4081,6 +4083,9 @@ impl crate::ScannerConfigObjectDelete for MemoryConfigStore {
revisions.remove(&key);
drop(revisions);
drop(objects);
if let Some(token) = self.cancel_after_deletes.lock().await.remove(&key) {
token.cancel();
}
if self.error_after_commit_deletes.lock().await.remove(&key) {
return Err(EcstoreError::other("injected delete error after commit"));
}
@@ -4088,6 +4093,11 @@ impl crate::ScannerConfigObjectDelete for MemoryConfigStore {
}
async fn scanner_data_usage_publication_admission(&self) -> Option<crate::ScannerDataUsagePublicationAdmission> {
let pause = self.pause_next_publication_admission.lock().await.take();
if let Some((entered, resume)) = pause {
entered.notify_one();
resume.notified().await;
}
if self.publication_admission_blocked.load(Ordering::Acquire) {
return None;
}
@@ -4589,7 +4599,7 @@ async fn scanner_legacy_usage_backup_survives_fencing_and_restart_after_real_met
.expect("publication must also read the intact backup");
assert_eq!(baseline.data.as_deref(), Some(data.as_slice()));
assert_eq!(baseline.revision, DataUsageCacheRevision::Missing);
fence_scanner_usage_epoch_with_expected_epoch(&CancellationToken::new(), store.clone(), 7, None, false)
fence_scanner_usage_epoch_with_expected_epoch(&CancellationToken::new(), store.clone(), 7, None, false, || true)
.await
.expect("legacy backup must be fenced into v2");
let fenced = read_config(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str())
@@ -4818,6 +4828,28 @@ async fn scanner_usage_state_reset_publishes_fenced_bootstrap_marker() {
);
}
let cycle_before_retry = read_config_with_revision(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str())
.await
.expect("cycle should remain before retry");
let marker_before_retry = read_config_with_revision(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str())
.await
.expect("bootstrap should remain before retry");
let retry = reset_scanner_usage_state_for_full_rebuild(CancellationToken::new(), store.clone())
.await
.expect("completed cleanup should be reentrant");
assert_eq!(retry.leader_epoch, result.leader_epoch);
assert_eq!(
read_config_with_revision(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str())
.await
.expect("cycle should remain"),
cycle_before_retry
);
assert_eq!(
read_config_with_revision(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str())
.await
.expect("bootstrap should remain"),
marker_before_retry
);
let (floor, state) = persisted_usage_floor_for_startup(store, false)
.await
.expect("reset marker should be resumable");
@@ -4973,7 +5005,7 @@ async fn scanner_usage_state_reset_slots_reject_primary_aba() {
store.objects.lock().await.insert(key.clone(), b"newer-json".to_vec());
store.revisions.lock().await.insert(key, 2);
let err = reset_scanner_usage_state_slots_for_full_rebuild(store, &slots, 0, 3)
let err = reset_scanner_usage_state_slots_for_full_rebuild(store, &slots, 0, 3, || true)
.await
.expect_err("stale primary revision must not be overwritten");
assert!(
@@ -4983,6 +5015,348 @@ async fn scanner_usage_state_reset_slots_reject_primary_aba() {
);
}
#[tokio::test]
async fn scanner_usage_state_reset_resumes_every_cleanup_boundary_without_rewriting_intent() {
for completed in 0..=4 {
let store = Arc::new(MemoryConfigStore::default());
let primary_path = DATA_USAGE_OBJ_NAME_PATH.as_str();
let cleanup_paths = [
format!("{primary_path}.bkp"),
LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str().to_string(),
format!("{}.bkp", LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str()),
DATA_USAGE_OBSERVED_OBJ_NAME_PATH.as_str().to_string(),
];
for path in std::iter::once(primary_path).chain(cleanup_paths.iter().map(String::as_str)) {
let mut usage = complete_usage_with_bucket_count(Some(std::time::SystemTime::UNIX_EPOCH), 0);
usage.scanner_epoch = Some(1);
save_config(store.clone(), path, serde_json::to_vec(&usage).expect("fixture should encode"))
.await
.expect("fixture should persist");
}
// These objects belong to other owners, even when reset cleanup resumes.
for path in ["buckets/quota-reservations/ledger", "buckets/example/incarnation"] {
save_config(store.clone(), path, b"retain".to_vec())
.await
.expect("unrelated state should persist");
}
let slots = read_usage_state_reset_slots(store.clone()).await.expect("slots should load");
let cancelled = CancellationToken::new();
if completed == 0 {
store
.cancel_after_successful_puts
.lock()
.await
.insert(memory_config_key(RUSTFS_META_BUCKET, primary_path), (2, cancelled.clone()));
} else {
store
.cancel_after_deletes
.lock()
.await
.insert(memory_config_key(RUSTFS_META_BUCKET, &cleanup_paths[completed - 1]), cancelled.clone());
}
let err = reset_scanner_usage_state_slots_for_full_rebuild(store.clone(), &slots, 0, 3, || !cancelled.is_cancelled())
.await
.expect_err("interruption should stop cleanup");
assert!(err.to_string().contains("ownership"), "boundary {completed}: {err}");
for (index, path) in cleanup_paths.iter().enumerate() {
assert_eq!(
store
.objects
.lock()
.await
.contains_key(&memory_config_key(RUSTFS_META_BUCKET, path)),
index >= completed,
"boundary {completed}, slot {index}"
);
}
let intent = read_config_with_revision(store.clone(), primary_path)
.await
.expect("intent should persist");
let slots = read_usage_state_reset_slots(store.clone())
.await
.expect("restart should reload slots");
reset_scanner_usage_state_slots_for_full_rebuild(store.clone(), &slots, 0, 3, || true)
.await
.expect("restart should complete the same intent");
assert_eq!(
read_config_with_revision(store.clone(), primary_path)
.await
.expect("intent should remain"),
intent
);
assert_eq!(store.put_counts.lock().await[&memory_config_key(RUSTFS_META_BUCKET, primary_path)], 2);
for path in cleanup_paths {
assert!(
!store
.objects
.lock()
.await
.contains_key(&memory_config_key(RUSTFS_META_BUCKET, &path))
);
}
for path in ["buckets/quota-reservations/ledger", "buckets/example/incarnation"] {
assert_eq!(read_config(store.clone(), path).await.expect("unrelated state should remain"), b"retain");
}
}
}
#[tokio::test]
async fn scanner_usage_state_reset_stops_usage_fence_after_owner_loss() {
let store = Arc::new(MemoryConfigStore::default());
let mut usage = complete_usage_with_bucket_count(Some(std::time::SystemTime::UNIX_EPOCH), 0);
usage.scanner_epoch = Some(1);
let bytes = serde_json::to_vec(&usage).expect("baseline should encode");
save_config(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str(), bytes.clone())
.await
.expect("baseline should persist");
let checks = AtomicUsize::new(0);
let err = fence_scanner_usage_epoch_with_expected_epoch(&CancellationToken::new(), store.clone(), 3, Some(0), false, || {
checks.fetch_add(1, Ordering::SeqCst) == 0
})
.await
.expect_err("ownership lost during reads must prevent the write");
assert!(err.to_string().contains("leadership was lost"), "{err}");
assert_eq!(
read_config(store, DATA_USAGE_OBJ_NAME_PATH.as_str())
.await
.expect("baseline should remain"),
bytes
);
}
#[tokio::test]
async fn scanner_usage_state_reset_cancels_during_publication_admission() {
for resuming in [false, true] {
let store = Arc::new(MemoryConfigStore::default());
let usage = if resuming {
scanner_usage_bootstrap_marker(std::time::SystemTime::UNIX_EPOCH, Some(3))
} else {
complete_usage_with_bucket_count(Some(std::time::SystemTime::UNIX_EPOCH), 0)
};
save_config(
store.clone(),
DATA_USAGE_OBJ_NAME_PATH.as_str(),
serde_json::to_vec(&usage).expect("primary should encode"),
)
.await
.expect("primary should persist");
save_config(store.clone(), LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str(), b"corrupt".to_vec())
.await
.expect("cleanup target should persist");
let slots = read_usage_state_reset_slots(store.clone()).await.expect("slots should load");
let before = store.objects.lock().await.clone();
let revisions_before = store.revisions.lock().await.clone();
let entered = Arc::new(tokio::sync::Notify::new());
let resume = Arc::new(tokio::sync::Notify::new());
*store.pause_next_publication_admission.lock().await = Some((entered.clone(), resume.clone()));
let cancelled = CancellationToken::new();
let (result, ()) = tokio::join!(
reset_scanner_usage_state_slots_for_full_rebuild(store.clone(), &slots, 0, 3, || !cancelled.is_cancelled()),
async {
entered.notified().await;
cancelled.cancel();
resume.notify_one();
}
);
let err = result.expect_err("losing ownership during admission must prevent mutation");
assert!(err.to_string().contains("ownership was lost"), "resuming={resuming}: {err}");
assert_eq!(*store.objects.lock().await, before);
assert_eq!(*store.revisions.lock().await, revisions_before);
}
}
#[tokio::test]
#[serial]
async fn scanner_usage_state_reset_rejects_corruption_without_a_trusted_floor() {
let (_temp_dir, store) = setup_scanner_cycle_store_with_usage_baseline(false).await;
save_config(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str(), b"{corrupt".to_vec())
.await
.expect("corrupt primary should persist");
let before = read_config_with_revision(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str())
.await
.expect("evidence should load");
let err = reset_scanner_usage_state_for_full_rebuild(CancellationToken::new(), store.clone())
.await
.expect_err("corruption must not become a zero floor");
assert!(err.to_string().contains("no trusted cycle or usage floor"), "{err}");
assert_eq!(
read_config_with_revision(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str())
.await
.expect("evidence should remain"),
before
);
assert!(matches!(
read_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str()).await,
Err(EcstoreError::ConfigNotFound)
));
let mut backup = complete_usage_with_bucket_count(Some(std::time::SystemTime::UNIX_EPOCH), 0);
backup.scanner_epoch = Some(7);
backup.scanner_cycle = Some(40);
save_config(
store.clone(),
&format!("{}.bkp", DATA_USAGE_OBJ_NAME_PATH.as_str()),
serde_json::to_vec(&backup).expect("backup should encode"),
)
.await
.expect("valid backup should persist");
let result = reset_scanner_usage_state_for_full_rebuild(CancellationToken::new(), store)
.await
.expect("valid backup should supply the recovery floor");
assert_eq!(result.leader_epoch, 8);
assert_eq!(result.next_cycle, 41);
}
#[tokio::test]
async fn scanner_usage_state_reset_rejects_replaced_intent_and_newer_cleanup_slot() {
let store = Arc::new(MemoryConfigStore::default());
let marker = scanner_usage_bootstrap_marker(std::time::SystemTime::UNIX_EPOCH, Some(3));
let bytes = serde_json::to_vec(&marker).expect("marker should encode");
save_config(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str(), bytes.clone())
.await
.expect("intent should persist");
let slots = read_usage_state_reset_slots(store.clone()).await.expect("slots should load");
save_config(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str(), bytes)
.await
.expect("another intent should persist");
let err = reset_scanner_usage_state_slots_for_full_rebuild(store.clone(), &slots, 0, 3, || true)
.await
.expect_err("same epoch cannot replace an intent revision");
assert!(err.to_string().contains("intent revision changed"), "{err}");
let mut newer = complete_usage_with_bucket_count(Some(std::time::SystemTime::UNIX_EPOCH), 0);
newer.scanner_epoch = Some(3);
let path = format!("{}.bkp", DATA_USAGE_OBJ_NAME_PATH.as_str());
let bytes = serde_json::to_vec(&newer).expect("newer snapshot should encode");
save_config(store.clone(), &path, bytes.clone())
.await
.expect("newer snapshot should persist");
let slots = read_usage_state_reset_slots(store.clone())
.await
.expect("slots should reload");
let err = reset_scanner_usage_state_slots_for_full_rebuild(store.clone(), &slots, 0, 3, || true)
.await
.expect_err("cleanup cannot delete same-epoch progress");
assert!(err.to_string().contains("not older than its intent"), "{err}");
assert_eq!(read_config(store, &path).await.expect("newer snapshot should remain"), bytes);
}
#[tokio::test]
#[serial]
async fn scanner_usage_state_reset_rejects_decodable_untrusted_floor() {
let (_temp_dir, store) = setup_scanner_cycle_store_with_usage_baseline(false).await;
let invalid_identity = DataUsageInfo {
usage_snapshot_complete: true,
buckets_count: 1,
last_update: Some(std::time::SystemTime::UNIX_EPOCH),
..Default::default()
};
for usage in [DataUsageInfo::default(), invalid_identity] {
save_config(
store.clone(),
DATA_USAGE_OBJ_NAME_PATH.as_str(),
serde_json::to_vec(&usage).expect("fixture should encode"),
)
.await
.expect("untrusted primary should persist");
let before = read_config_with_revision(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str())
.await
.expect("primary should load");
let err = reset_scanner_usage_state_for_full_rebuild(CancellationToken::new(), store.clone())
.await
.expect_err("valid JSON alone cannot prove a usage floor");
assert!(err.to_string().contains("no trusted cycle or usage floor"), "{err}");
assert_eq!(
read_config_with_revision(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str())
.await
.expect("evidence should remain"),
before
);
assert!(matches!(
read_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str()).await,
Err(EcstoreError::ConfigNotFound)
));
}
}
#[test]
fn full_rescan_reset_rejects_unknown_marker_phase_even_with_invalid_compat_fields() {
for state in [serde_json::json!("rewrite-v2"), serde_json::json!(7), serde_json::Value::Null] {
let marker = serde_json::json!({"state": state, "retry_count": "future-type", "schema_version": 99});
let err = super::cycle_state::decode_recovery_marker_for_reset(
&serde_json::to_vec(&marker).expect("future marker should encode"),
&DataUsageCacheRevision::Etag("intent-1".to_string()),
)
.expect_err("unknown persistent phases must remain fenced");
assert!(err.to_string().contains("state is unsupported"), "{err}");
}
}
#[tokio::test]
#[serial]
async fn full_rescan_reset_preserves_unknown_phase_and_retries_completed_cleanup() {
let (_temp_dir, store) = setup_scanner_cycle_store().await;
save_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str(), b"corrupt".to_vec())
.await
.expect("corrupt primary should persist");
save_config(
store.clone(),
DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(),
br#"{"state":"future-rewrite"}"#.to_vec(),
)
.await
.expect("future marker should persist");
let primary_before = read_config_with_revision(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str())
.await
.expect("primary should load");
let marker_before = read_config_with_revision(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str())
.await
.expect("marker should load");
let err = reset_scanner_cycle_recovery(CancellationToken::new(), store.clone())
.await
.expect_err("unknown phase must block explicit reset");
assert!(err.to_string().contains("state is unsupported"), "{err}");
assert_eq!(
read_config_with_revision(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str())
.await
.expect("primary should remain"),
primary_before
);
assert_eq!(
read_config_with_revision(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str())
.await
.expect("marker should remain"),
marker_before
);
save_config(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), b"{malformed".to_vec())
.await
.expect("recoverable marker should persist");
reset_scanner_cycle_recovery(CancellationToken::new(), store.clone())
.await
.expect("reset should complete");
let primary = read_config_with_revision(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str())
.await
.expect("rebuilt primary should load");
let usage = read_config_with_revision(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str())
.await
.expect("fenced usage should load");
reset_scanner_cycle_recovery(CancellationToken::new(), store.clone())
.await
.expect("retry after marker deletion should complete");
assert_eq!(
read_config_with_revision(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str())
.await
.expect("rebuilt primary should remain"),
primary
);
assert_eq!(
read_config_with_revision(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str())
.await
.expect("fenced usage should remain"),
usage
);
}
#[tokio::test]
async fn scanner_usage_state_reset_slots_defer_when_publication_epoch_moves() {
let store = Arc::new(MemoryConfigStore::default());
@@ -4993,7 +5367,7 @@ async fn scanner_usage_state_reset_slots_defer_when_publication_epoch_moves() {
.expect("usage reset slots should be inspected");
store.publication_admission_blocked.store(true, Ordering::Release);
let err = reset_scanner_usage_state_slots_for_full_rebuild(store, &slots, 0, 3)
let err = reset_scanner_usage_state_slots_for_full_rebuild(store, &slots, 0, 3, || true)
.await
.expect_err("movement admission loss must defer reset");
assert!(
@@ -8169,6 +8543,44 @@ fn scanner_activity_snapshot_digest_fences_dirty_usage_state() {
assert_ne!(scanner_activity_snapshot_digest(&clean), scanner_activity_snapshot_digest(&pending));
}
#[test]
fn scanner_activity_structural_digest_ignores_regular_bucket_writes() {
let baseline = BTreeMap::from([("node-2".to_string(), scanner_node_activity("epoch-a", 7, 3))]);
let mut written = baseline.clone();
let activity = written.get_mut("node-2").expect("node should exist");
activity.namespace_generation = 8;
activity.dirty_usage_generation = 6;
activity.dirty_usage_pending = true;
assert_ne!(scanner_activity_snapshot_digest(&baseline), scanner_activity_snapshot_digest(&written));
assert_eq!(
scanner_activity_structural_digest(&baseline),
scanner_activity_structural_digest(&written),
"bucket writes are refreshed through the dirty-bucket scope rather than invalidating every cache"
);
}
#[test]
fn scanner_activity_structural_digest_fences_restart_and_maintenance() {
let baseline = BTreeMap::from([("node-2".to_string(), scanner_node_activity("epoch-a", 7, 3))]);
let mut restarted = baseline.clone();
restarted.get_mut("node-2").expect("node should exist").instance_id = "epoch-b".to_string();
let mut maintained = baseline.clone();
maintained
.get_mut("node-2")
.expect("node should exist")
.maintenance_generation = 4;
assert_ne!(
scanner_activity_structural_digest(&baseline),
scanner_activity_structural_digest(&restarted)
);
assert_ne!(
scanner_activity_structural_digest(&baseline),
scanner_activity_structural_digest(&maintained)
);
}
#[test]
fn scanner_dirty_usage_acknowledgements_exclude_local_and_clean_nodes() {
let snapshot = BTreeMap::from([
+118 -1
View File
@@ -21,6 +21,7 @@ use crate::{
DataUsageCacheSource, DataUsageEntry, DataUsageEntryInfo, DataUsageInfo, DataUsageScanPlanDigest, DataUsageSnapshotSetState,
ScannerError, SizeSummary, TierStats,
};
use bytes::Bytes;
use futures::future::join_all;
use metrics::counter;
use rand::seq::SliceRandom as _;
@@ -54,6 +55,7 @@ use tokio_util::task::AbortOnDropHandle;
use tracing::{debug, error, warn};
use crate::ScannerObjectInfo as ObjectInfo;
use crate::storage_api::EcstoreScannerPeerDirtyUsageSnapshot;
use crate::storage_api::ScannerStorage;
use crate::storage_api::scan::NamespaceLocking as _;
use crate::storage_api::scanner_io::{BucketInfo, BucketOptions};
@@ -111,6 +113,121 @@ pub(crate) struct ScannerBucketScanScope {
baseline_scan_plan_digest: Option<DataUsageScanPlanDigest>,
}
impl ScannerBucketScanScope {
fn is_default(&self) -> bool {
self.selected_buckets.is_none() && self.baseline_scan_plan_digest.is_none()
}
fn from_dirty_buckets(selected_buckets: HashSet<String>, baseline_scan_plan_digest: DataUsageScanPlanDigest) -> Self {
Self {
selected_buckets: Some(Arc::new(selected_buckets)),
baseline_scan_plan_digest: Some(baseline_scan_plan_digest),
}
}
}
#[derive(Clone, Copy)]
pub(super) struct ScannerCacheBaselineProof<'a> {
pub(super) data: Option<&'a Bytes>,
pub(super) expected_sources: &'a HashSet<DataUsageCacheSource>,
pub(super) leader_epoch: u64,
pub(super) want_cycle: u64,
pub(super) scan_plan_digest: DataUsageScanPlanDigest,
}
#[derive(Clone, Debug, PartialEq, Eq)]
struct ScannerPeerDirtyUsageExpectation {
instance_id: String,
generation: u64,
pending: bool,
}
fn verified_remote_dirty_usage_buckets(
expected_peers: &HashMap<String, ScannerPeerDirtyUsageExpectation>,
peer_snapshots: Vec<(String, EcstoreScannerPeerDirtyUsageSnapshot)>,
) -> Option<HashSet<String>> {
if expected_peers.is_empty() || peer_snapshots.len() != expected_peers.len() {
return None;
}
let mut received_peers = HashSet::with_capacity(peer_snapshots.len());
let mut dirty_buckets = HashSet::new();
for (host, snapshot) in peer_snapshots {
let expected = expected_peers.get(&host)?;
if !received_peers.insert(host)
|| snapshot.instance_id != expected.instance_id
|| snapshot.generation != expected.generation
|| snapshot.generation == u64::MAX
|| snapshot.protocol_version != crate::SCANNER_DIRTY_USAGE_SNAPSHOT_PROTOCOL_VERSION
|| !snapshot.complete
|| snapshot.pending_bucket_count != u64::try_from(snapshot.buckets.len()).unwrap_or(u64::MAX)
|| (expected.pending && snapshot.pending_bucket_count == 0)
{
return None;
}
dirty_buckets.extend(snapshot.buckets.into_keys());
}
(received_peers.len() == expected_peers.len()).then_some(dirty_buckets)
}
fn complete_scanner_cache_baseline_plan_digest(proof: ScannerCacheBaselineProof<'_>) -> Option<DataUsageScanPlanDigest> {
let data = proof.data?;
let baseline = serde_json::from_slice::<DataUsageInfo>(data).ok()?;
if !baseline.is_complete_bucket_usage_snapshot()
|| baseline.usage_snapshot_partial
|| baseline.usage_snapshot_converged != Some(true)
|| baseline.scanner_epoch != Some(proof.leader_epoch)
|| baseline.usage_snapshot_set_states.len() != proof.expected_sources.len()
{
return None;
}
let mut states = HashSet::with_capacity(baseline.usage_snapshot_set_states.len());
for state in &baseline.usage_snapshot_set_states {
let source = DataUsageCacheSource::new(usize::try_from(state.pool_index).ok()?, usize::try_from(state.set_index).ok()?);
if !proof.expected_sources.contains(&source)
|| !states.insert(source)
|| !state.complete
|| state.tombstone
|| state.scanner_epoch != Some(proof.leader_epoch)
|| state.scanner_cycle.is_none_or(|cycle| cycle > proof.want_cycle)
|| state.scan_plan_digest != Some(proof.scan_plan_digest.0)
{
return None;
}
}
(states == *proof.expected_sources).then_some(proof.scan_plan_digest)
}
fn scoped_scan_scope_from_dirty_buckets(
requested_scope: ScannerBucketScanScope,
dirty_buckets: HashSet<String>,
dirty_snapshot_complete: bool,
all_buckets: &[BucketInfo],
baseline_proof: ScannerCacheBaselineProof<'_>,
) -> ScannerBucketScanScope {
if !requested_scope.is_default() || !dirty_snapshot_complete {
return requested_scope;
}
let current_buckets = all_buckets.iter().map(|bucket| bucket.name.as_str()).collect::<HashSet<_>>();
let selected_buckets = dirty_buckets
.into_iter()
.filter(|bucket| current_buckets.contains(bucket.as_str()))
.collect::<HashSet<_>>();
if selected_buckets.is_empty() {
return requested_scope;
}
let Some(baseline_scan_plan_digest) = complete_scanner_cache_baseline_plan_digest(baseline_proof) else {
return requested_scope;
};
ScannerBucketScanScope::from_dirty_buckets(selected_buckets, baseline_scan_plan_digest)
}
pub(crate) fn is_scanner_metadata_corrupt_error(err: &StorageError) -> bool {
matches!(err, StorageError::Io(io) if io.to_string().starts_with(SCANNER_METADATA_CORRUPT_ERROR))
}
@@ -749,7 +866,7 @@ mod io_cache;
mod io_cycle;
#[cfg(test)]
use io_cache::{ScannerSetCacheGeneration, prepare_scoped_set_scan};
pub(crate) use io_cycle::nsscanner_with_storage_status;
pub(crate) use io_cycle::{ScannerCycleRequest, nsscanner_with_storage_status_scoped};
mod io_disk;
#[cfg(test)]
mod publish_gate_tests;
+18
View File
@@ -282,9 +282,26 @@ pub(super) fn completed_data_usage_info(
.iter()
.map(|(bucket, usage)| (bucket.clone(), usage.size))
.collect();
let mut usage_snapshot_set_states = results
.iter()
.map(|result| {
let source = result.info.source?;
Some(DataUsageSnapshotSetState {
pool_index: u64::try_from(source.pool_index).ok()?,
set_index: u64::try_from(source.set_index).ok()?,
scanner_cycle: Some(result.info.next_cycle),
scanner_epoch: Some(result.info.leader_epoch),
scan_plan_digest: Some(result.info.scan_plan_digest?.0),
complete: true,
tombstone: false,
})
})
.collect::<Option<Vec<_>>>()?;
usage_snapshot_set_states.sort_by_key(|state| (state.pool_index, state.set_index));
let data_usage_info = DataUsageInfo {
last_update: Some(merged_last_update),
scanner_cycle: Some(results.first()?.info.next_cycle),
scanner_epoch: Some(results.first()?.info.leader_epoch),
objects_total_count: u64::try_from(total.objects).ok()?,
versions_total_count: u64::try_from(total.versions).ok()?,
delete_markers_total_count: u64::try_from(total.delete_markers).ok()?,
@@ -295,6 +312,7 @@ pub(super) fn completed_data_usage_info(
bucket_sizes,
buckets_usage,
usage_snapshot_complete: true,
usage_snapshot_set_states,
..Default::default()
};
Some((data_usage_info, merged_last_update))
+94 -1
View File
@@ -71,6 +71,7 @@ where
leader_epoch,
scan_mode,
scan_scope: ScannerBucketScanScope::default(),
persisted_usage_baseline: None,
};
nsscanner_with_storage_status_scoped(store, request).await
}
@@ -83,6 +84,79 @@ pub(crate) struct ScannerCycleRequest {
pub(crate) leader_epoch: u64,
pub(crate) scan_mode: HealScanMode,
pub(crate) scan_scope: ScannerBucketScanScope,
pub(crate) persisted_usage_baseline: Option<Bytes>,
}
struct ScannerBucketScopeResolution<'a> {
requested_scope: ScannerBucketScanScope,
baseline_proof: ScannerCacheBaselineProof<'a>,
activity_before: &'a crate::scanner::ScannerActivitySnapshot,
dirty_usage_snapshot: &'a DirtyUsageSnapshot,
all_buckets: &'a [BucketInfo],
}
async fn resolve_scanner_bucket_scan_scope<S>(
store: &S,
distributed: bool,
resolution: ScannerBucketScopeResolution<'_>,
) -> ScannerBucketScanScope
where
S: ScannerStorage,
{
if !resolution.requested_scope.is_default()
|| !resolution.dirty_usage_snapshot.covers_all_pending
|| resolution.dirty_usage_snapshot.generation == u64::MAX
|| resolution.dirty_usage_snapshot.buckets.len() > crate::SCANNER_DIRTY_USAGE_SNAPSHOT_MAX_ENTRIES
{
return resolution.requested_scope;
}
let mut dirty_buckets = resolution
.dirty_usage_snapshot
.buckets
.keys()
.cloned()
.collect::<HashSet<_>>();
if distributed {
let Some(notification_system) = store.scanner_notification_system() else {
return resolution.requested_scope;
};
let Ok(peer_snapshots) = notification_system.scanner_dirty_usage_snapshots().await else {
return resolution.requested_scope;
};
let mut expected_peers = HashMap::new();
for (host, lease_instance_id, _) in crate::scanner::scanner_activity_publication_lease_targets(resolution.activity_before)
{
let Some((activity_instance_id, generation, pending)) =
crate::scanner::scanner_activity_dirty_usage_state_for_host(resolution.activity_before, &host)
else {
return resolution.requested_scope;
};
if activity_instance_id != lease_instance_id || expected_peers.contains_key(&host) {
return resolution.requested_scope;
}
expected_peers.insert(
host,
ScannerPeerDirtyUsageExpectation {
instance_id: activity_instance_id.to_string(),
generation,
pending,
},
);
}
let Some(remote_dirty_buckets) = verified_remote_dirty_usage_buckets(&expected_peers, peer_snapshots) else {
return resolution.requested_scope;
};
dirty_buckets.extend(remote_dirty_buckets);
}
scoped_scan_scope_from_dirty_buckets(
resolution.requested_scope,
dirty_buckets,
true,
resolution.all_buckets,
resolution.baseline_proof,
)
}
pub(crate) async fn nsscanner_with_storage_status_scoped<S>(store: &S, request: ScannerCycleRequest) -> Result<ScannerCycleResult>
@@ -97,6 +171,7 @@ where
leader_epoch,
scan_mode,
scan_scope,
persisted_usage_baseline,
} = request;
let child_token = ctx.child_token();
let _tier_cycle_guard = begin_tier_registry_cycle(want_cycle, leader_epoch);
@@ -186,8 +261,26 @@ where
}
bucket_plan_complete &= buckets_by_source.keys().copied().collect::<HashSet<_>>() == *expected_sources;
let scan_plan_digest =
scanner_bucket_plan_digest(&all_buckets, crate::scanner::scanner_activity_snapshot_digest(&activity_before));
scanner_bucket_plan_digest(&all_buckets, crate::scanner::scanner_activity_structural_digest(&activity_before));
let dirty_usage_snapshot = Arc::new(snapshot_dirty_usage_buckets(&all_buckets, dirty_generation_before_bucket_list));
let scan_scope = resolve_scanner_bucket_scan_scope(
store,
distributed,
ScannerBucketScopeResolution {
requested_scope: scan_scope,
baseline_proof: ScannerCacheBaselineProof {
data: persisted_usage_baseline.as_ref(),
expected_sources: &expected_sources,
leader_epoch,
want_cycle,
scan_plan_digest,
},
activity_before: &activity_before,
dirty_usage_snapshot: &dirty_usage_snapshot,
all_buckets: &all_buckets,
},
)
.await;
let cache_cycle_floor = Arc::new(AtomicU64::new(want_cycle));
let tier_registry = runtime_tier_registry_for_cycle(want_cycle, leader_epoch).await;
let tier_registry_generation = tier_registry.generation;
@@ -655,9 +655,31 @@ fn completed_data_usage_info_requires_every_set_before_publish() {
.expect("all completed sets should produce a publishable data usage snapshot");
assert_eq!(last_update, SystemTime::UNIX_EPOCH + Duration::from_secs(20));
assert_eq!(data_usage_info.scanner_cycle, Some(0));
assert_eq!(data_usage_info.scanner_epoch, Some(0));
assert_eq!(data_usage_info.objects_total_count, 3);
assert_eq!(data_usage_info.buckets_usage.len(), 3);
assert!(data_usage_info.usage_snapshot_complete);
assert_eq!(
data_usage_info
.usage_snapshot_set_states
.iter()
.map(|state| {
(
state.pool_index,
state.set_index,
state.scanner_cycle,
state.scanner_epoch,
state.scan_plan_digest,
state.complete,
state.tombstone,
)
})
.collect::<Vec<_>>(),
vec![
(0, 0, Some(0), Some(0), Some(TEST_PLAN_DIGEST.0), true, false),
(1, 0, Some(0), Some(0), Some(TEST_PLAN_DIGEST.0), true, false),
]
);
assert_eq!(
data_usage_info
.buckets_usage
+190
View File
@@ -17,6 +17,7 @@ use super::io_disk::tier_stats_template;
use super::*;
use crate::scanner_budget::ScannerCycleBudgetConfig;
use crate::scanner_folder::ScannerItem;
use crate::storage_api::EcstoreScannerPeerDirtyUsageSnapshot;
use crate::storage_api::owner::{
EcstorePoolDecommissionInfo, EcstoreRebalStatus, EcstoreRebalanceInfo, EcstoreRebalanceMeta, EcstoreRebalanceStats,
};
@@ -796,6 +797,195 @@ fn complete_set_usage_cache(buckets: &[(&str, usize)], scan_plan_digest: DataUsa
cache
}
fn complete_usage_baseline(
source: DataUsageCacheSource,
scan_plan_digest: DataUsageScanPlanDigest,
scanner_cycle: u64,
scanner_epoch: u64,
) -> bytes::Bytes {
let baseline = DataUsageInfo {
last_update: Some(SystemTime::UNIX_EPOCH + Duration::from_secs(10)),
scanner_cycle: Some(scanner_cycle),
scanner_epoch: Some(scanner_epoch),
buckets_count: 1,
buckets_usage: HashMap::from([("photos".to_string(), Default::default())]),
usage_snapshot_complete: true,
usage_snapshot_converged: Some(true),
usage_snapshot_set_states: vec![DataUsageSnapshotSetState {
pool_index: u64::try_from(source.pool_index).expect("test pool index should fit"),
set_index: u64::try_from(source.set_index).expect("test set index should fit"),
scanner_cycle: Some(scanner_cycle),
scanner_epoch: Some(scanner_epoch),
scan_plan_digest: Some(scan_plan_digest.0),
complete: true,
tombstone: false,
}],
..Default::default()
};
bytes::Bytes::from(serde_json::to_vec(&baseline).expect("test baseline should encode"))
}
#[test]
fn scoped_scan_requires_a_converged_complete_baseline_with_exact_set_provenance() {
let source = DataUsageCacheSource::new(1, 2);
let expected_sources = HashSet::from([source]);
let scan_plan_digest = DataUsageScanPlanDigest([9; 32]);
let baseline = complete_usage_baseline(source, scan_plan_digest, 7, 11);
assert_eq!(
complete_scanner_cache_baseline_plan_digest(ScannerCacheBaselineProof {
data: Some(&baseline),
expected_sources: &expected_sources,
leader_epoch: 11,
want_cycle: 8,
scan_plan_digest,
}),
Some(scan_plan_digest)
);
let mut incomplete = serde_json::from_slice::<DataUsageInfo>(&baseline).expect("test baseline should decode");
incomplete.usage_snapshot_converged = Some(false);
let incomplete = bytes::Bytes::from(serde_json::to_vec(&incomplete).expect("test baseline should encode"));
assert_eq!(
complete_scanner_cache_baseline_plan_digest(ScannerCacheBaselineProof {
data: Some(&incomplete),
expected_sources: &expected_sources,
leader_epoch: 11,
want_cycle: 8,
scan_plan_digest,
}),
None
);
let mut wrong_provenance = serde_json::from_slice::<DataUsageInfo>(&baseline).expect("test baseline should decode");
wrong_provenance.usage_snapshot_set_states[0].scan_plan_digest = Some([8; 32]);
let wrong_provenance = bytes::Bytes::from(serde_json::to_vec(&wrong_provenance).expect("test baseline should encode"));
assert_eq!(
complete_scanner_cache_baseline_plan_digest(ScannerCacheBaselineProof {
data: Some(&wrong_provenance),
expected_sources: &expected_sources,
leader_epoch: 11,
want_cycle: 8,
scan_plan_digest,
}),
None
);
}
#[test]
fn scoped_scan_selects_only_current_dirty_buckets_after_baseline_validation() {
let source = DataUsageCacheSource::new(1, 2);
let expected_sources = HashSet::from([source]);
let baseline_scan_plan_digest = DataUsageScanPlanDigest([4; 32]);
let current_scan_plan_digest = DataUsageScanPlanDigest([5; 32]);
let baseline = complete_usage_baseline(source, current_scan_plan_digest, 7, 11);
let scope = scoped_scan_scope_from_dirty_buckets(
ScannerBucketScanScope::default(),
HashSet::from(["photos".to_string(), "deleted".to_string()]),
true,
&[bucket_info("photos")],
ScannerCacheBaselineProof {
data: Some(&baseline),
expected_sources: &expected_sources,
leader_epoch: 11,
want_cycle: 8,
scan_plan_digest: current_scan_plan_digest,
},
);
assert_eq!(scope.baseline_scan_plan_digest, Some(current_scan_plan_digest));
assert_eq!(
scope
.selected_buckets
.as_deref()
.expect("validated scope should select a bucket"),
&HashSet::from(["photos".to_string()])
);
assert_ne!(scope.baseline_scan_plan_digest, Some(baseline_scan_plan_digest));
}
fn peer_dirty_usage_snapshot(
instance_id: &str,
generation: u64,
complete: bool,
buckets: &[(&str, u64)],
) -> EcstoreScannerPeerDirtyUsageSnapshot {
EcstoreScannerPeerDirtyUsageSnapshot {
instance_id: instance_id.to_string(),
generation,
pending_bucket_count: u64::try_from(buckets.len()).expect("test bucket count should fit"),
protocol_version: crate::SCANNER_DIRTY_USAGE_SNAPSHOT_PROTOCOL_VERSION,
complete,
buckets: buckets
.iter()
.map(|(bucket, generation)| ((*bucket).to_string(), *generation))
.collect(),
}
}
#[test]
fn verified_remote_dirty_usage_buckets_merges_only_complete_current_snapshots() {
let expected_peers = HashMap::from([
(
"node-a:9000".to_string(),
ScannerPeerDirtyUsageExpectation {
instance_id: "instance-a".to_string(),
generation: 7,
pending: true,
},
),
(
"node-b:9000".to_string(),
ScannerPeerDirtyUsageExpectation {
instance_id: "instance-b".to_string(),
generation: 3,
pending: false,
},
),
]);
assert_eq!(
verified_remote_dirty_usage_buckets(
&expected_peers,
vec![
(
"node-a:9000".to_string(),
peer_dirty_usage_snapshot("instance-a", 7, true, &[("photos", 7)]),
),
(
"node-b:9000".to_string(),
peer_dirty_usage_snapshot("instance-b", 3, true, &[("archive", 3)]),
),
],
),
Some(HashSet::from(["photos".to_string(), "archive".to_string()]))
);
}
#[test]
fn verified_remote_dirty_usage_buckets_rejects_incomplete_or_stale_peer_state() {
let expected_peers = HashMap::from([(
"node-a:9000".to_string(),
ScannerPeerDirtyUsageExpectation {
instance_id: "instance-a".to_string(),
generation: 7,
pending: true,
},
)]);
for snapshot in [
peer_dirty_usage_snapshot("instance-a", 7, false, &[("photos", 7)]),
peer_dirty_usage_snapshot("instance-a", 6, true, &[("photos", 6)]),
peer_dirty_usage_snapshot("instance-b", 7, true, &[("photos", 7)]),
peer_dirty_usage_snapshot("instance-a", 7, true, &[]),
] {
assert!(
verified_remote_dirty_usage_buckets(&expected_peers, vec![("node-a:9000".to_string(), snapshot)]).is_none(),
"incomplete, stale, mismatched, or empty pending peer state must fall back to a full scan"
);
}
}
#[test]
fn scoped_set_scan_preserves_unselected_usage_and_drops_deleted_buckets() {
let baseline_digest = DataUsageScanPlanDigest([1; 32]);
+3 -1
View File
@@ -103,7 +103,9 @@ pub(crate) use rustfs_ecstore::api::rebalance::{
RebalStatus as EcstoreRebalStatus, RebalanceInfo as EcstoreRebalanceInfo, RebalanceMeta as EcstoreRebalanceMeta,
RebalanceStats as EcstoreRebalanceStats,
};
pub(crate) use rustfs_ecstore::api::rpc::ScannerBucketListing as EcstoreScannerBucketListing;
pub(crate) use rustfs_ecstore::api::rpc::{
ScannerBucketListing as EcstoreScannerBucketListing, ScannerPeerDirtyUsageSnapshot as EcstoreScannerPeerDirtyUsageSnapshot,
};
#[cfg(test)]
pub(crate) use rustfs_ecstore::api::runtime::InstanceContext as EcstoreInstanceContext;
pub(crate) use rustfs_ecstore::api::runtime::{
+177 -20
View File
@@ -39,7 +39,7 @@ use rustfs_utils::path::path_join;
use s3s::header::{CONTENT_LENGTH, CONTENT_TYPE};
use s3s::{Body, S3Request, S3Response, S3Result, s3_error};
use serde::{Deserialize, Serialize};
use std::collections::{BTreeMap, HashSet};
use std::collections::{BTreeMap, BTreeSet, HashSet};
use std::future::Future;
use std::path::PathBuf;
use std::sync::Arc;
@@ -261,6 +261,7 @@ struct BackgroundHealStatus<'a> {
heal_active_tasks: u64,
heal_operations: rustfs_heal::HealOperationsSnapshot,
cluster_status_complete: bool,
coverage: &'a BackgroundHealCoverage,
#[serde(skip_serializing_if = "Option::is_none")]
progress: Option<BackgroundHealProgress>,
}
@@ -300,6 +301,23 @@ fn background_heal_runtime_state(
type BackgroundHealProgress = rustfs_heal::HealProgress;
#[derive(Debug, Serialize)]
struct BackgroundHealCoverage {
expected: usize,
responded: usize,
unknown: usize,
reasons: BTreeSet<BackgroundHealCoverageReason>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize)]
#[serde(rename_all = "snake_case")]
enum BackgroundHealCoverageReason {
NotificationSystemUnavailable,
PeerTopologyIncomplete,
PeerStatusUnsupported,
PeerStatusUnavailable,
}
#[derive(Debug)]
struct ClusterHealStatusSnapshot {
info: BackgroundHealInfo,
@@ -307,6 +325,7 @@ struct ClusterHealStatusSnapshot {
operations: rustfs_heal::HealOperationsSnapshot,
progress: Option<BackgroundHealProgress>,
complete: bool,
coverage: BackgroundHealCoverage,
}
fn add_priority_counts(total: &mut rustfs_heal::HealPriorityCounts, next: rustfs_heal::HealPriorityCounts) {
@@ -338,6 +357,7 @@ fn add_operations(total: &mut rustfs_heal::HealOperationsSnapshot, next: rustfs_
}
fn aggregate_cluster_heal_status(snapshots: Vec<NodeHealStatusSnapshot>) -> ClusterHealStatusSnapshot {
let responded = snapshots.len();
let mut info = BackgroundHealInfo::default();
let mut operations = rustfs_heal::HealOperationsSnapshot::default();
let mut progress = Vec::new();
@@ -379,6 +399,12 @@ fn aggregate_cluster_heal_status(snapshots: Vec<NodeHealStatusSnapshot>) -> Clus
operations,
progress,
complete: true,
coverage: BackgroundHealCoverage {
expected: responded,
responded,
unknown: 0,
reasons: BTreeSet::new(),
},
}
}
@@ -413,12 +439,14 @@ fn merge_peer_heal_statuses(
mut snapshots: Vec<NodeHealStatusSnapshot>,
peer_statuses: Vec<Result<Option<NodeHealStatusSnapshot>, String>>,
expected_nodes: usize,
topology_complete: bool,
coverage_reason: Option<BackgroundHealCoverageReason>,
) -> S3Result<ClusterHealStatusSnapshot> {
let mut reasons: BTreeSet<_> = coverage_reason.into_iter().collect();
for peer_status in peer_statuses {
match peer_status {
Ok(Some(snapshot)) => snapshots.push(snapshot),
Ok(None) => {
reasons.insert(BackgroundHealCoverageReason::PeerStatusUnsupported);
warn!(
event = EVENT_ADMIN_REQUEST_FAILED,
component = LOG_COMPONENT_ADMIN_API,
@@ -430,6 +458,7 @@ fn merge_peer_heal_statuses(
);
}
Err(err) => {
reasons.insert(BackgroundHealCoverageReason::PeerStatusUnavailable);
warn!(
event = EVENT_ADMIN_REQUEST_FAILED,
component = LOG_COMPONENT_ADMIN_API,
@@ -452,9 +481,12 @@ fn merge_peer_heal_statuses(
// so during a reconfiguration the count can equal `expected_nodes` while
// the topology is known-incomplete. Counting alone would report a
// definitive answer precisely when the membership itself is in doubt.
let complete = topology_complete && snapshots.len() == expected_nodes;
let complete = reasons.is_empty() && snapshots.len() == expected_nodes;
let mut status = aggregate_cluster_heal_status(snapshots);
status.complete = complete;
status.coverage.expected = expected_nodes;
status.coverage.unknown = expected_nodes.saturating_sub(status.coverage.responded);
status.coverage.reasons = reasons;
// A partial answer must never be mistakable for a definitive verdict: an
// unreachable peer might be mid-heal, so reporting the reachable nodes'
// "idle" (or disabled/uninitialized) as the cluster state would falsely
@@ -497,7 +529,12 @@ async fn read_cluster_heal_status(
return Ok(aggregate_cluster_heal_status(snapshots));
}
let Some(notification_system) = notification_system else {
return Err(cluster_heal_status_unavailable("notification_system_unavailable"));
return merge_peer_heal_statuses(
snapshots,
Vec::new(),
expected_nodes,
Some(BackgroundHealCoverageReason::NotificationSystemUnavailable),
);
};
// An incomplete peer topology (a down member's client slot, a rolling
// upgrade) previously failed the whole endpoint here, before any peer was
@@ -540,7 +577,12 @@ async fn read_cluster_heal_status(
}))
.await;
merge_peer_heal_statuses(snapshots, peer_statuses, expected_nodes, topology_complete)
merge_peer_heal_statuses(
snapshots,
peer_statuses,
expected_nodes,
(!topology_complete).then_some(BackgroundHealCoverageReason::PeerTopologyIncomplete),
)
}
async fn query_peer_replacement_recovery_status<E>(
@@ -1164,6 +1206,7 @@ fn encode_background_heal_status(
heal_operations: rustfs_heal::HealOperationsSnapshot,
progress: Option<BackgroundHealProgress>,
cluster_status_complete: bool,
coverage: &BackgroundHealCoverage,
) -> S3Result<Vec<u8>> {
let status = BackgroundHealStatus {
info,
@@ -1172,6 +1215,7 @@ fn encode_background_heal_status(
heal_active_tasks: heal_operations.active_tasks,
heal_operations,
cluster_status_complete,
coverage,
progress,
};
serde_json::to_vec(&status).map_err(|e| {
@@ -1461,6 +1505,7 @@ impl Operation for BackgroundHealStatusHandler {
cluster_status.operations,
cluster_status.progress,
cluster_status.complete,
&cluster_status.coverage,
)?;
info!(
event = EVENT_ADMIN_RESPONSE_EMITTED,
@@ -1515,13 +1560,14 @@ impl Operation for ReplacementRecoveryStatusHandler {
mod tests {
use super::extract_heal_init_params;
use super::{
BackgroundHealProgress, HealInitParams, HealResp, HealRuntimeState, aggregate_cluster_heal_status,
aggregate_replacement_recovery_cluster_status, background_heal_runtime_state, build_heal_channel_request,
build_replacement_recovery_status_response, encode_background_heal_status, encode_heal_control_path,
encode_heal_start_success, encode_heal_task_status, execute_after_heal_control_capability, heal_channel_response_items,
heal_channel_response_progress, heal_channel_response_summary, heal_control_response_id, json_response,
map_heal_response, merge_peer_heal_statuses, peer_topology_complete, query_peer_heal_status,
query_peer_replacement_recovery_status, reject_heal_admission, validate_heal_request_mode, validate_heal_target,
BackgroundHealCoverage, BackgroundHealCoverageReason, BackgroundHealProgress, HealInitParams, HealResp, HealRuntimeState,
aggregate_cluster_heal_status, aggregate_replacement_recovery_cluster_status, background_heal_runtime_state,
build_heal_channel_request, build_replacement_recovery_status_response, encode_background_heal_status,
encode_heal_control_path, encode_heal_start_success, encode_heal_task_status, execute_after_heal_control_capability,
heal_channel_response_items, heal_channel_response_progress, heal_channel_response_summary, heal_control_response_id,
json_response, map_heal_response, merge_peer_heal_statuses, peer_topology_complete, query_peer_heal_status,
query_peer_replacement_recovery_status, read_cluster_heal_status, reject_heal_admission, validate_heal_request_mode,
validate_heal_target,
};
use crate::storage::rpc::node_service::heal::{
NodeHealProgress, NodeHealStatusSnapshot, NodeReplacementRecoveryStatusSnapshot, encode_node_replacement_recovery_status,
@@ -2175,7 +2221,13 @@ mod tests {
..Default::default()
};
let encoded = encode_background_heal_status(&info, HealRuntimeState::Active, operations, None, true)
let coverage = BackgroundHealCoverage {
expected: 1,
responded: 1,
unknown: 0,
reasons: Default::default(),
};
let encoded = encode_background_heal_status(&info, HealRuntimeState::Active, operations, None, true, &coverage)
.expect("background heal info should serialize");
let json: serde_json::Value = serde_json::from_slice(&encoded).expect("json should deserialize");
@@ -2225,6 +2277,12 @@ mod tests {
rustfs_heal::HealOperationsSnapshot::default(),
Some(progress),
true,
&BackgroundHealCoverage {
expected: 1,
responded: 1,
unknown: 0,
reasons: Default::default(),
},
)
.expect("background heal info should serialize");
let json: serde_json::Value = serde_json::from_slice(&encoded).expect("json should deserialize");
@@ -2398,6 +2456,84 @@ mod tests {
assert!(peer_topology_complete(1, 0, 0, 1, 0));
}
#[tokio::test]
async fn test_background_heal_status_without_notification_preserves_local_snapshot() {
let initialized = rustfs_heal::heal_runtime_initialized();
let info = BackgroundHealInfo {
bitrot_start_cycle: 37,
current_scan_mode: HealScanMode::Deep,
..Default::default()
};
for expected in [1, 3] {
let status = tokio::time::timeout(Duration::from_secs(1), read_cluster_heal_status(info.clone(), None, expected))
.await
.expect("local status must not wait for remote peers")
.expect("missing notification must retain the local snapshot");
assert_eq!(status.info.bitrot_start_cycle, 37);
assert_eq!(status.info.current_scan_mode, HealScanMode::Deep);
assert_eq!(status.complete, expected == 1);
assert_eq!(status.coverage.expected, expected);
assert_eq!(status.coverage.responded, 1);
assert_eq!(status.coverage.unknown, expected - 1);
if expected == 1 {
assert!(status.coverage.reasons.is_empty());
} else {
assert!(matches!(status.state, HealRuntimeState::Degraded | HealRuntimeState::Active));
assert_eq!(
status.coverage.reasons,
[BackgroundHealCoverageReason::NotificationSystemUnavailable].into()
);
}
let encoded = encode_background_heal_status(
&status.info,
status.state,
status.operations,
status.progress,
status.complete,
&status.coverage,
)
.expect("fallback status must encode");
let decoded: rustfs_madmin::client::BackgroundHealStatus =
serde_json::from_slice(&encoded).expect("the actual madmin client must decode the server response");
assert_eq!(decoded.cluster_status_complete, expected == 1);
let coverage = decoded.coverage.expect("new server supplies coverage");
assert_eq!(coverage.expected, Some(expected));
assert_eq!(coverage.responded, Some(1));
assert_eq!(coverage.unknown, Some(expected - 1));
if expected > 1 {
assert_eq!(coverage.reasons, ["notification_system_unavailable"]);
}
}
assert_eq!(
rustfs_heal::heal_runtime_initialized(),
initialized,
"reading status must not initialize heal"
);
}
#[test]
fn test_background_heal_status_coverage_reasons_are_bounded() {
let local = NodeHealStatusSnapshot::for_test(true, true, BackgroundHealInfo::default(), Default::default(), None);
let peers = (0..100)
.map(|index| {
if index % 2 == 0 {
Ok(None)
} else {
Err("peer unavailable".to_owned())
}
})
.collect();
let status = merge_peer_heal_statuses(vec![local], peers, 101, None).expect("local status remains available");
assert_eq!(status.coverage.responded, 1);
assert_eq!(status.coverage.unknown, 100);
assert_eq!(status.coverage.reasons.len(), 2);
let encoded = serde_json::to_vec(&status.coverage).expect("coverage encodes");
assert!(encoded.len() < 256, "coverage must not grow with peer failures");
let decoded: rustfs_madmin::client::BackgroundHealCoverage =
serde_json::from_slice(&encoded).expect("client coverage decodes");
assert_eq!(decoded.reasons, ["peer_status_unsupported", "peer_status_unavailable"]);
}
#[test]
fn test_peer_status_merge_degrades_explicitly_and_never_claims_idle() {
let local = || {
@@ -2413,15 +2549,20 @@ mod tests {
// but the safety property of the previous fail-closed behaviour is
// preserved: the partial answer is labelled Degraded, never Idle, so
// unknown peer work cannot be mistaken for "nothing is running".
let partial = merge_peer_heal_statuses(vec![local()], vec![Err("peer timeout".to_string())], 2, true)
let partial = merge_peer_heal_statuses(vec![local()], vec![Err("peer timeout".to_string())], 2, None)
.expect("an unreachable peer degrades the answer instead of destroying it");
assert!(!partial.complete);
assert_eq!(partial.state, HealRuntimeState::Degraded);
assert_eq!(partial.coverage.expected, 2);
assert_eq!(partial.coverage.responded, 1);
assert_eq!(partial.coverage.unknown, 1);
assert_eq!(partial.coverage.reasons, [BackgroundHealCoverageReason::PeerStatusUnavailable].into());
let older_peer = merge_peer_heal_statuses(vec![local()], vec![Ok(None)], 2, true)
let older_peer = merge_peer_heal_statuses(vec![local()], vec![Ok(None)], 2, None)
.expect("an older peer degrades the answer instead of destroying it");
assert!(!older_peer.complete);
assert_eq!(older_peer.state, HealRuntimeState::Degraded);
assert_eq!(older_peer.coverage.reasons, [BackgroundHealCoverageReason::PeerStatusUnsupported].into());
let known_active = NodeHealStatusSnapshot::for_test(
true,
@@ -2433,12 +2574,12 @@ mod tests {
},
None,
);
let partial_active = merge_peer_heal_statuses(vec![known_active], vec![Ok(None)], 2, true)
let partial_active = merge_peer_heal_statuses(vec![known_active], vec![Ok(None)], 2, None)
.expect("known active work may be reported as an explicit partial status");
assert!(!partial_active.complete);
assert_eq!(partial_active.state, HealRuntimeState::Active);
merge_peer_heal_statuses(Vec::new(), vec![Err("peer timeout".to_string())], 2, true)
merge_peer_heal_statuses(Vec::new(), vec![Err("peer timeout".to_string())], 2, None)
.expect_err("no snapshot at all still fails closed");
}
@@ -2459,12 +2600,22 @@ mod tests {
None,
)
};
let full_count_incomplete_topology = merge_peer_heal_statuses(vec![snapshot()], vec![Ok(Some(snapshot()))], 2, false)
.expect("incomplete topology degrades the answer instead of destroying it");
let full_count_incomplete_topology = merge_peer_heal_statuses(
vec![snapshot()],
vec![Ok(Some(snapshot()))],
2,
Some(BackgroundHealCoverageReason::PeerTopologyIncomplete),
)
.expect("incomplete topology degrades the answer instead of destroying it");
assert!(!full_count_incomplete_topology.complete);
assert_eq!(full_count_incomplete_topology.state, HealRuntimeState::Degraded);
assert_eq!(full_count_incomplete_topology.coverage.unknown, 0);
assert_eq!(
full_count_incomplete_topology.coverage.reasons,
[BackgroundHealCoverageReason::PeerTopologyIncomplete].into()
);
let full_count_complete_topology = merge_peer_heal_statuses(vec![snapshot()], vec![Ok(Some(snapshot()))], 2, true)
let full_count_complete_topology = merge_peer_heal_statuses(vec![snapshot()], vec![Ok(Some(snapshot()))], 2, None)
.expect("complete topology and full count is a definitive answer");
assert!(full_count_complete_topology.complete);
assert_eq!(full_count_complete_topology.state, HealRuntimeState::Idle);
@@ -2478,6 +2629,12 @@ mod tests {
rustfs_heal::HealOperationsSnapshot::default(),
None,
false,
&BackgroundHealCoverage {
expected: 2,
responded: 1,
unknown: 1,
reasons: [BackgroundHealCoverageReason::PeerStatusUnavailable].into(),
},
)
.expect("degraded status must serialize");
let json: serde_json::Value = serde_json::from_slice(&encoded).expect("valid json");
+292 -106
View File
@@ -1,10 +1,11 @@
#!/usr/bin/env python3
"""Fail when a critical scheduled validation has not started recently."""
"""Require recent scheduled attempts and completed successes on the default branch."""
from __future__ import annotations
import argparse
from datetime import datetime, timedelta, timezone
import io
import json
import os
from pathlib import Path
@@ -13,7 +14,7 @@ import sys
import tempfile
import unittest
from unittest import mock
from urllib.parse import quote, urlencode
from urllib.parse import parse_qs, quote, urlencode, urlsplit
from urllib.request import Request, urlopen
@@ -75,15 +76,8 @@ def stale_reason(
run: dict[str, object] | None,
now: datetime,
max_age_hours: int,
never_ran_grace_until: datetime | None = None,
) -> str | None:
if run is None:
# The grace deadline only covers a workflow whose first scheduled slot
# has not arrived yet (for example a monthly cron enabled mid-month).
# A recorded-but-old run proves the schedule used to fire and stopped,
# so the grace never masks that case.
if never_ran_grace_until is not None and now <= never_ran_grace_until:
return None
return "no scheduled run has been recorded"
created_at = parse_timestamp(run.get("created_at"))
age = now - created_at
@@ -93,14 +87,23 @@ def stale_reason(
def fetch_latest_scheduled_run(
repository: str, workflow: str, token: str, api_url: str
repository: str,
workflow: str,
token: str,
api_url: str,
default_branch: str,
successful: bool = False,
) -> dict[str, object] | None:
owner, repo = repository.split("/", 1)
workflow_name = Path(workflow).name
query = {"event": "schedule", "branch": default_branch, "per_page": 1}
if successful:
# Filter on the server: the last success may be beyond a page of failures.
query["status"] = "success"
endpoint = (
f"{api_url.rstrip('/')}/repos/{quote(owner, safe='')}/{quote(repo, safe='')}"
f"/actions/workflows/{quote(workflow_name, safe='')}/runs?"
+ urlencode({"event": "schedule", "per_page": 1})
+ urlencode(query)
)
request = Request(
endpoint,
@@ -110,57 +113,104 @@ def fetch_latest_scheduled_run(
"X-GitHub-Api-Version": "2022-11-28",
},
)
with urlopen(request, timeout=30) as response:
# Two requests per manifest entry must fit the watchdog's ten-minute job.
with urlopen(request, timeout=15) as response:
payload = json.load(response)
runs = payload.get("workflow_runs")
runs = payload.get("workflow_runs") if isinstance(payload, dict) else None
if not isinstance(runs, list):
raise ValueError(f"GitHub returned no workflow_runs list for {workflow}")
total_count = payload.get("total_count")
if not isinstance(total_count, int) or isinstance(total_count, bool) or total_count < len(runs):
raise ValueError(f"GitHub returned an invalid run count for {workflow}")
if not runs:
if total_count:
raise ValueError(f"GitHub returned an empty first page with recorded runs for {workflow}")
return None
if not isinstance(runs[0], dict):
run = runs[0]
if not isinstance(run, dict):
raise ValueError(f"GitHub returned an invalid workflow run for {workflow}")
return runs[0]
if run.get("event") != "schedule" or run.get("head_branch") != default_branch:
raise ValueError(f"GitHub returned a run outside the scheduled default-branch query for {workflow}")
if not isinstance(run.get("status"), str) or not run["status"]:
raise ValueError(f"GitHub returned no run status for {workflow}")
conclusion = run.get("conclusion")
if (conclusion is not None and not isinstance(conclusion, str)) or (
run["status"] == "completed" and not conclusion
):
raise ValueError(f"GitHub returned an invalid run conclusion for {workflow}")
if successful and (run["status"] != "completed" or conclusion != "success"):
raise ValueError(f"GitHub returned a run without a completed success for {workflow}")
parse_timestamp(run.get("created_at"))
if not isinstance(run.get("html_url"), str) or not run["html_url"]:
raise ValueError(f"GitHub returned no run URL for {workflow}")
return run
def write_report(path: Path, failures: list[tuple[str, int, str, str]]) -> None:
lines = ["## Scheduled validation freshness"]
if not failures:
lines.append("")
lines.append("All critical scheduled validations have a recent scheduled run.")
else:
lines.extend(
[
"",
"The following critical validations are stale or could not be inspected:",
"",
"| Workflow | Limit | Result | Last run |",
"| --- | ---: | --- | --- |",
]
)
for workflow, max_age_hours, reason, run_url in failures:
link = f"[open]({run_url})" if run_url else ""
lines.append(f"| `{workflow}` | {max_age_hours}h | {reason} | {link} |")
def describe_run(run: dict[str, object] | None) -> str:
if run is None:
return "No recorded run"
outcome = run["status"]
if run.get("conclusion"):
outcome = f"{outcome}/{run['conclusion']}"
return f"[{outcome}]({run['html_url']}) — created {run['created_at']}"
def write_report(path: Path, rows: list[tuple[str, int, str, str, str]], default_branch: str) -> None:
lines = [
"## Scheduled validation freshness",
"",
f"Default branch: `{default_branch}`. Ages use scheduled-run creation time; rerunning an old commit does not refresh its evidence.",
"Attempt outcomes are shown independently of successful-run freshness.",
"Success is the GitHub workflow run conclusion; suite completeness remains the responsibility of each workflow.",
"",
"| Workflow | Limit | Freshness | Last attempt | Last completed success |",
"| --- | ---: | --- | --- | --- |",
]
for workflow, max_age_hours, result, attempt, success in rows:
cells = [f"`{workflow}`", f"{max_age_hours}h", result, attempt, success]
lines.append("| " + " | ".join(cell.replace("|", "\\|").replace("\n", " ") for cell in cells) + " |")
path.write_text("\n".join(lines) + "\n")
def check_freshness(
config: Path, report: Path, repository: str, token: str, api_url: str
config: Path, report: Path, repository: str, token: str, api_url: str, default_branch: str
) -> int:
now = datetime.now(timezone.utc)
failures: list[tuple[str, int, str, str]] = []
rows: list[tuple[str, int, str, str, str]] = []
failed = False
for workflow, max_age_hours, never_ran_grace_until in load_validations(config):
try:
run = fetch_latest_scheduled_run(repository, workflow, token, api_url)
reason = stale_reason(run, now, max_age_hours, never_ran_grace_until)
if reason is not None:
run_url = str(run.get("html_url", "")) if run else ""
failures.append((workflow, max_age_hours, reason, run_url))
except Exception as error:
failures.append(
(workflow, max_age_hours, f"inspection failed: {error}", "")
)
write_report(report, failures)
return 1 if failures else 0
runs: dict[str, dict[str, object] | None] = {}
reasons: list[str] = []
for label, successful in (("Last attempt", False), ("Last completed success", True)):
try:
runs[label] = fetch_latest_scheduled_run(
repository, workflow, token, api_url, default_branch, successful
)
except Exception as error:
reasons.append(f"{label}: inspection failed: {error}")
# A failed inspection or any recorded attempt ends first-run grace.
initial_grace = (
len(runs) == 2
and all(run is None for run in runs.values())
and never_ran_grace_until is not None
and now <= never_ran_grace_until
)
if not initial_grace:
for label, run in runs.items():
reason = stale_reason(run, now, max_age_hours)
if reason is not None:
reasons.append(f"{label}: {reason}")
failed |= bool(reasons)
result = "; ".join(reasons) if reasons else "Fresh"
if initial_grace:
result = f"Initial grace until {never_ran_grace_until.isoformat()}"
evidence = [
describe_run(runs[label]) if label in runs else "Inspection failed"
for label in ("Last attempt", "Last completed success")
]
rows.append((workflow, max_age_hours, result, *evidence))
write_report(report, rows, default_branch)
return 1 if failed else 0
class SelfTests(unittest.TestCase):
@@ -173,15 +223,6 @@ class SelfTests(unittest.TestCase):
self.assertIsNotNone(stale_reason(past_limit, self.NOW, 36))
self.assertIsNotNone(stale_reason(None, self.NOW, 36))
def test_never_ran_grace_only_covers_missing_runs(self) -> None:
future_grace = self.NOW + timedelta(hours=1)
past_grace = self.NOW - timedelta(seconds=1)
self.assertIsNone(stale_reason(None, self.NOW, 36, future_grace))
self.assertIsNone(stale_reason(None, self.NOW, 36, self.NOW))
self.assertIsNotNone(stale_reason(None, self.NOW, 36, past_grace))
stale_run = {"created_at": "2026-08-20T23:59:59Z"}
self.assertIsNotNone(stale_reason(stale_run, self.NOW, 36, future_grace))
def test_config_rejects_duplicate_and_invalid_entries(self) -> None:
with tempfile.TemporaryDirectory() as tmp:
path = Path(tmp) / "validations.json"
@@ -238,63 +279,205 @@ class SelfTests(unittest.TestCase):
],
)
def test_check_reports_missing_runs(self) -> None:
@staticmethod
def run_fixture(**overrides: object) -> dict[str, object]:
return {
"status": "completed",
"conclusion": "success",
"event": "schedule",
"head_branch": "release/current",
"created_at": "2026-08-22T00:00:00Z",
"html_url": "https://github.test/rustfs/rustfs/actions/runs/1",
**overrides,
}
def check_payloads(
self, payloads: list[object], *, grace: str | None = None, workflows: int = 1
) -> tuple[int, str, list]:
with tempfile.TemporaryDirectory() as tmp:
root = Path(tmp)
config = root / "validations.json"
report = root / "report.md"
config.write_text(
json.dumps(
[
{"workflow": ".github/workflows/ci.yml", "max_age_hours": 36},
{"workflow": ".github/workflows/fuzz.yml", "max_age_hours": 36},
{"workflow": ".github/workflows/mint.yml", "max_age_hours": 36},
]
)
)
with mock.patch(
__name__ + ".fetch_latest_scheduled_run",
side_effect=[
{"created_at": "2999-01-01T00:00:00Z"},
None,
RuntimeError("API unavailable"),
],
entries = [
{"workflow": f".github/workflows/check-{index}.yml", "max_age_hours": 36}
for index in range(workflows)
]
if grace is not None:
entries[0]["never_ran_grace_until"] = grace
config.write_text(json.dumps(entries))
responses = []
for payload in payloads:
if isinstance(payload, dict) and isinstance(payload.get("workflow_runs"), list):
payload = {"total_count": len(payload["workflow_runs"]), **payload}
responses.append(payload if isinstance(payload, Exception) else io.StringIO(json.dumps(payload)))
with (
mock.patch(__name__ + ".urlopen", side_effect=responses) as request,
mock.patch(__name__ + ".datetime", wraps=datetime) as clock,
):
self.assertEqual(
check_freshness(
config,
report,
"rustfs/rustfs",
"token",
"https://api.github.test",
),
1,
clock.now.return_value = self.NOW
status = check_freshness(
config, report, "rustfs/rustfs", "test-token",
"https://api.github.test", "release/current",
)
contents = report.read_text()
self.assertIn(".github/workflows/fuzz.yml", contents)
self.assertIn("inspection failed: API unavailable", contents)
self.assertNotIn(".github/workflows/ci.yml`", contents)
return status, report.read_text(), request.call_args_list
config.write_text(
json.dumps(
[{"workflow": ".github/workflows/ci.yml", "max_age_hours": 36}]
def test_requests_filter_schedule_default_branch_and_success_on_server(self) -> None:
attempt = self.run_fixture(status="in_progress", conclusion=None)
success = self.run_fixture(html_url="https://github.test/rustfs/rustfs/actions/runs/2")
status, report, calls = self.check_payloads([
{"workflow_runs": [attempt], "total_count": 1001},
{"workflow_runs": [success], "total_count": 1},
])
self.assertEqual(status, 0)
self.assertEqual(len(calls), 2)
for call, successful in zip(calls, (False, True)):
request = call.args[0]
url = urlsplit(request.full_url)
self.assertEqual(url.path, "/repos/rustfs/rustfs/actions/workflows/check-0.yml/runs")
expected = {"event": ["schedule"], "branch": ["release/current"], "per_page": ["1"]}
if successful:
expected["status"] = ["success"]
self.assertEqual(parse_qs(url.query), expected)
self.assertEqual(request.get_header("Authorization"), "Bearer test-token")
self.assertEqual(call.kwargs, {"timeout": 15})
self.assertIn("[in_progress]", report)
self.assertIn(str(attempt["html_url"]), report)
self.assertIn(str(success["html_url"]), report)
def test_cancelled_attempt_cannot_refresh_expired_success(self) -> None:
attempt = self.run_fixture(conclusion="cancelled")
success = self.run_fixture(
created_at="2026-08-20T23:59:59Z", updated_at="2026-08-22T11:59:59Z",
html_url="https://github.test/rustfs/rustfs/actions/runs/2",
)
status, report, _ = self.check_payloads([
{"workflow_runs": [attempt]}, {"workflow_runs": [success]},
])
self.assertEqual(status, 1)
self.assertIn("Last completed success: last scheduled run is", report)
self.assertIn("[completed/cancelled]", report)
for run in (attempt, success):
self.assertIn(str(run["html_url"]), report)
self.assertIn(str(run["created_at"]), report)
def test_attempt_outcome_does_not_replace_recent_success(self) -> None:
success = self.run_fixture(created_at="2026-08-21T00:00:00Z")
for state, conclusion in (
("completed", "failure"), ("completed", "cancelled"),
("completed", "timed_out"), ("completed", "success"),
("queued", None), ("in_progress", None),
):
with self.subTest(state=state, conclusion=conclusion):
status, report, _ = self.check_payloads([
{"workflow_runs": [self.run_fixture(status=state, conclusion=conclusion)]},
{"workflow_runs": [success]},
])
self.assertEqual(status, 0)
self.assertIn(f"[{state}" + (f"/{conclusion}" if conclusion else "") + "]", report)
self.assertIn("Fresh", report)
self.assertNotIn("All critical scheduled validations", report)
def test_grace_requires_two_successful_queries_with_no_history(self) -> None:
for attempt, success, grace, expected in (
(None, None, "2026-08-22T12:00:00Z", 0),
(None, None, "2026-08-22T11:59:59Z", 1),
(self.run_fixture(conclusion="failure"), None, "2026-08-23T00:00:00Z", 1),
(self.run_fixture(status="queued", conclusion=None), None, "2026-08-23T00:00:00Z", 1),
(None, self.run_fixture(), "2026-08-23T00:00:00Z", 1),
):
with self.subTest(attempt=attempt, success=success, grace=grace):
status, report, _ = self.check_payloads([
{"workflow_runs": [] if attempt is None else [attempt]},
{"workflow_runs": [] if success is None else [success]},
], grace=grace)
self.assertEqual(status, expected)
self.assertEqual("Initial grace until" in report, expected == 0)
def test_api_failures_preserve_other_evidence_and_never_enter_grace(self) -> None:
good = {"workflow_runs": [self.run_fixture()]}
for first, second in (
(RuntimeError("API unavailable"), good),
(good, RuntimeError("API unavailable")),
(RuntimeError("API unavailable"), {"workflow_runs": []}),
):
with self.subTest(first=first, second=second):
status, report, calls = self.check_payloads(
[first, second], grace="2026-08-23T00:00:00Z"
)
)
with mock.patch(
__name__ + ".fetch_latest_scheduled_run",
return_value={"created_at": "2999-01-01T00:00:00Z"},
):
self.assertEqual(
check_freshness(
config,
report,
"rustfs/rustfs",
"token",
"https://api.github.test",
),
0,
)
self.assertIn("All critical scheduled validations", report.read_text())
self.assertEqual(status, 1)
self.assertEqual(len(calls), 2)
self.assertIn("inspection failed: API unavailable", report)
self.assertNotIn("Initial grace until", report)
if first is good or second is good:
self.assertIn(str(self.run_fixture()["html_url"]), report)
def test_invalid_api_evidence_fails_closed(self) -> None:
malformed = [
[], {}, {"workflow_runs": {}}, {"workflow_runs": [None]},
{"workflow_runs": [], "total_count": 1},
{"workflow_runs": [], "total_count": -1},
{"workflow_runs": [], "total_count": None},
{"workflow_runs": [], "total_count": True},
*({"workflow_runs": [self.run_fixture(**override)]} for override in (
{"event": "workflow_dispatch"}, {"head_branch": "other"},
{"created_at": "invalid"}, {"created_at": "2026-08-22T00:00:00"},
{"status": None}, {"conclusion": None}, {"conclusion": 1},
{"html_url": ""},
)),
]
for payload in malformed:
for index, label in enumerate(("Last attempt", "Last completed success")):
with self.subTest(payload=payload, label=label):
payloads = [{"workflow_runs": [self.run_fixture()]} for _ in range(2)]
payloads[index] = payload
status, report, _ = self.check_payloads(payloads, grace="2026-08-23T00:00:00Z")
self.assertEqual(status, 1)
self.assertIn(f"{label}: inspection failed", report)
self.assertNotIn("Initial grace until", report)
self.assertIn(str(self.run_fixture()["html_url"]), report)
for state, conclusion in (("in_progress", "success"), ("completed", "failure"), ("completed", "skipped")):
with self.subTest(state=state, conclusion=conclusion):
status, report, _ = self.check_payloads([
{"workflow_runs": [self.run_fixture()]},
{"workflow_runs": [self.run_fixture(status=state, conclusion=conclusion)]},
])
self.assertEqual(status, 1)
self.assertIn("without a completed success", report)
def test_report_retains_every_workflow(self) -> None:
status, report, calls = self.check_payloads([
{"workflow_runs": [self.run_fixture()]}, {"workflow_runs": [self.run_fixture()]},
{"workflow_runs": []}, {"workflow_runs": []},
RuntimeError("API unavailable"), {"workflow_runs": [self.run_fixture()]},
], workflows=3)
self.assertEqual(status, 1)
self.assertEqual(len(calls), 6)
for index in range(3):
self.assertEqual(report.count(f"`.github/workflows/check-{index}.yml`"), 1)
self.assertIn("No recorded run", report)
self.assertIn("Inspection failed", report)
def test_cli_requires_the_repository_default_branch(self) -> None:
from check_test_wiring import yaml_block
workflow = (ROOT / ".github/workflows/scheduled-validation-freshness.yml").read_text().splitlines()
job = yaml_block(workflow, "check-freshness", 2)
self.assertIsNotNone(job)
start = job.index(" - name: Check latest scheduled runs")
end = next((index for index in range(start + 1, len(job)) if job[index].startswith(" - ")), len(job))
environment = yaml_block(job[start:end], "env", 8)
self.assertIsNotNone(environment)
self.assertIn(" RUSTFS_DEFAULT_BRANCH: ${{ github.event.repository.default_branch }}", environment)
with (
mock.patch.dict(os.environ, {"GITHUB_REPOSITORY": "rustfs/rustfs", "GH_TOKEN": "test-token"}, clear=True),
mock.patch.object(sys, "argv", ["checker", "--report", "unused.md"]),
mock.patch("sys.stderr", new=io.StringIO()) as stderr,
self.assertRaises(SystemExit) as error,
):
main()
self.assertEqual(error.exception.code, 2)
self.assertIn("RUSTFS_DEFAULT_BRANCH", stderr.getvalue())
def main() -> int:
@@ -318,11 +501,14 @@ def main() -> int:
repository = os.environ.get("GITHUB_REPOSITORY", "")
token = os.environ.get("GH_TOKEN", "")
api_url = os.environ.get("GITHUB_API_URL", "https://api.github.com")
default_branch = os.environ.get("RUSTFS_DEFAULT_BRANCH", "")
if not re.fullmatch(r"[^/\s]+/[^/\s]+", repository):
parser.error("GITHUB_REPOSITORY must be owner/repository")
if not token:
parser.error("GH_TOKEN is required")
return check_freshness(args.config, args.report, repository, token, api_url)
if not default_branch or any(character.isspace() for character in default_branch):
parser.error("RUSTFS_DEFAULT_BRANCH is required and must name the repository default branch")
return check_freshness(args.config, args.report, repository, token, api_url, default_branch)
if __name__ == "__main__":