mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-23 04:39:04 +00:00
Compare commits
12 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 5ff5a7a06e | |||
| c58accabe7 | |||
| 4d92835d7b | |||
| 71db3327ad | |||
| 68e95d497e | |||
| 8350d73be9 | |||
| 9a2a33a7eb | |||
| 5df26a13ba | |||
| e3a362989f | |||
| 0cd2ae20e2 | |||
| 55a7fa9f03 | |||
| 9f51f37a0d |
@@ -1,20 +0,0 @@
|
||||
# Report-only calibration baseline from https://github.com/rustfs/rustfs/actions/runs/29394996173.
|
||||
# Update counts only with a linked coverage run and a reviewed explanation.
|
||||
phase = "report-only"
|
||||
allowed_drop_percentage_points = 1.0
|
||||
|
||||
[crates."crates/iam"]
|
||||
covered = 5149
|
||||
count = 8131
|
||||
|
||||
[crates."crates/kms"]
|
||||
covered = 2950
|
||||
count = 4200
|
||||
|
||||
[crates."crates/policy"]
|
||||
covered = 4636
|
||||
count = 5464
|
||||
|
||||
[crates."crates/crypto"]
|
||||
covered = 469
|
||||
count = 494
|
||||
@@ -36,7 +36,6 @@ script-tests: ## Run shell script tests
|
||||
./scripts/test_manual_transition_runbooks.sh
|
||||
./scripts/check_embedded_secrets.sh --self-test
|
||||
python3 ./scripts/check_test_wiring.py --self-test
|
||||
python3 ./scripts/check_security_coverage.py --self-test
|
||||
python3 ./scripts/check_scheduled_validation_freshness.py --self-test
|
||||
python3 ./scripts/s3-tests/test_report_compat.py
|
||||
bash -n ./scripts/validate_object_data_cache_cold_stampede.sh
|
||||
|
||||
@@ -12,12 +12,14 @@
|
||||
# See the License for the specific language governing permissions and
|
||||
# limitations under the License.
|
||||
|
||||
# Workspace line-coverage baseline and security-crate calibration
|
||||
# (backlog#1153 infra-5/infra-6).
|
||||
# Weekly workspace line-coverage baseline (backlog#1153 infra-5).
|
||||
#
|
||||
# NON-BLOCKING by design: the weekly job gives coverage a visible baseline and
|
||||
# trend, while relevant pull requests run a report-only security-crate
|
||||
# comparison. Neither job is a required check during calibration.
|
||||
# NON-BLOCKING by design: this workflow only runs on schedule and manual
|
||||
# dispatch, so it never attaches a status to a PR and must never be made a
|
||||
# required check. It exists to give coverage a visible baseline and trend
|
||||
# (per-crate table in the job summary, lcov artifact kept 90 days) — the
|
||||
# per-crate ratchet for the security-critical crates builds on it later
|
||||
# (backlog#1153 infra-6, report-only first per the ci-11 ladder).
|
||||
#
|
||||
# Measurement scope matches the PR test gate (ci.yml "Run tests"):
|
||||
# `--workspace --exclude e2e_test` with the `ci` nextest profile. Doctests are
|
||||
@@ -29,17 +31,6 @@
|
||||
name: coverage
|
||||
|
||||
on:
|
||||
pull_request:
|
||||
branches: [main]
|
||||
paths:
|
||||
- "crates/iam/**"
|
||||
- "crates/kms/**"
|
||||
- "crates/policy/**"
|
||||
- "crates/crypto/**"
|
||||
- ".config/coverage-baselines.toml"
|
||||
- "scripts/coverage_per_crate.py"
|
||||
- "scripts/check_security_coverage.py"
|
||||
- ".github/workflows/coverage.yml"
|
||||
workflow_dispatch:
|
||||
schedule:
|
||||
# 07:00 UTC Sunday — staggered clear of the other Sunday crons: ci (00:00),
|
||||
@@ -48,10 +39,6 @@ on:
|
||||
# e2e-replication-nightly (04:00) and performance-ab (06:00) lanes.
|
||||
- cron: "43 7 * * 0"
|
||||
|
||||
concurrency:
|
||||
group: ${{ github.workflow }}-${{ github.event_name }}-${{ github.event.pull_request.number || github.ref }}
|
||||
cancel-in-progress: ${{ github.event_name != 'schedule' }}
|
||||
|
||||
# Only alert-on-failure needs more than read access; it declares its own
|
||||
# job-level `issues: write`.
|
||||
permissions:
|
||||
@@ -59,13 +46,12 @@ permissions:
|
||||
|
||||
jobs:
|
||||
coverage:
|
||||
name: Workspace line coverage
|
||||
name: Workspace coverage (weekly)
|
||||
runs-on: sm-standard-4
|
||||
# The instrumented build cannot reuse the regular CI cache (different
|
||||
# RUSTFLAGS), so a cold run rebuilds the workspace before running the
|
||||
# full suite. Exact-head run 32573798257 needed 119m42s including reports
|
||||
# and artifact upload, so keep a bounded 30-minute publication margin.
|
||||
timeout-minutes: 150
|
||||
# RUSTFLAGS), so a cold week rebuilds the workspace before running the
|
||||
# full suite; give it double the test job's 60-minute budget.
|
||||
timeout-minutes: 120
|
||||
env:
|
||||
FORCE_JAVASCRIPT_ACTIONS_TO_NODE24: "true"
|
||||
# Match the PR gate's nextest semantics (ci.yml runs `--profile ci`):
|
||||
@@ -105,9 +91,7 @@ jobs:
|
||||
cargo llvm-cov report --json --output-path target/llvm-cov/coverage.json
|
||||
|
||||
- name: Write per-crate summary
|
||||
run: |
|
||||
python3 scripts/coverage_per_crate.py target/llvm-cov/coverage.json >> "$GITHUB_STEP_SUMMARY"
|
||||
python3 scripts/check_security_coverage.py target/llvm-cov/coverage.json >> "$GITHUB_STEP_SUMMARY"
|
||||
run: python3 scripts/coverage_per_crate.py target/llvm-cov/coverage.json >> "$GITHUB_STEP_SUMMARY"
|
||||
|
||||
- name: Upload coverage artifact
|
||||
if: always()
|
||||
|
||||
@@ -39,10 +39,11 @@ jobs:
|
||||
env:
|
||||
FORCE_JAVASCRIPT_ACTIONS_TO_NODE24: "true"
|
||||
steps:
|
||||
- name: Checkout repository
|
||||
- name: Checkout main branch
|
||||
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7
|
||||
with:
|
||||
persist-credentials: false
|
||||
ref: main
|
||||
|
||||
- name: Setup Rust environment
|
||||
uses: ./.github/actions/setup
|
||||
@@ -88,10 +89,11 @@ jobs:
|
||||
# either casing.
|
||||
NO_PROXY: 127.0.0.1,localhost
|
||||
steps:
|
||||
- name: Checkout repository
|
||||
- name: Checkout main branch
|
||||
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7
|
||||
with:
|
||||
persist-credentials: false
|
||||
ref: main
|
||||
|
||||
- name: Setup Rust environment
|
||||
uses: ./.github/actions/setup
|
||||
@@ -176,10 +178,11 @@ jobs:
|
||||
FORCE_JAVASCRIPT_ACTIONS_TO_NODE24: "true"
|
||||
NO_PROXY: 127.0.0.1,localhost
|
||||
steps:
|
||||
- name: Checkout repository
|
||||
- name: Checkout main branch
|
||||
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7
|
||||
with:
|
||||
persist-credentials: false
|
||||
ref: main
|
||||
|
||||
- name: Setup Rust environment
|
||||
uses: ./.github/actions/setup
|
||||
|
||||
Generated
+1
@@ -9564,6 +9564,7 @@ dependencies = [
|
||||
"serde",
|
||||
"serde_json",
|
||||
"serial_test",
|
||||
"sha2 0.11.0",
|
||||
"temp-env",
|
||||
"tempfile",
|
||||
"thiserror 2.0.20",
|
||||
|
||||
@@ -8000,10 +8000,15 @@ impl DiskAPI for LocalDisk {
|
||||
use std::io::Write as _;
|
||||
|
||||
let file_path = self.io_get_object_path(volume, path)?;
|
||||
let lock_path = file_path.with_extension("rustfs-cas.lock");
|
||||
let path = path.to_string();
|
||||
let sync_metadata = effective_durability(volume).syncs_commit_metadata();
|
||||
return Ok(tokio::task::spawn_blocking(move || {
|
||||
// A persistent directory lock bounds metadata growth. Removing
|
||||
// per-target lock files can split flock ownership across inodes.
|
||||
let lock_path = file_path
|
||||
.parent()
|
||||
.ok_or_else(|| std::io::Error::new(ErrorKind::InvalidInput, "conditional file has no parent"))?
|
||||
.join(".rustfs-cas.lock");
|
||||
let lock = std::fs::OpenOptions::new()
|
||||
.create(true)
|
||||
.truncate(false)
|
||||
@@ -8068,7 +8073,25 @@ impl DiskAPI for LocalDisk {
|
||||
.map_err(DiskError::from)??);
|
||||
}
|
||||
|
||||
#[cfg(not(unix))]
|
||||
#[cfg(windows)]
|
||||
{
|
||||
let file_path = self.io_get_object_path(volume, path)?;
|
||||
let sync_metadata = effective_durability(volume).syncs_commit_metadata();
|
||||
let publication_root = self.publication_root.clone();
|
||||
return Ok(tokio::task::spawn_blocking(move || {
|
||||
os::compare_and_update_control_file(
|
||||
&file_path,
|
||||
expected.as_deref(),
|
||||
replacement.as_deref(),
|
||||
sync_metadata,
|
||||
&publication_root,
|
||||
)
|
||||
})
|
||||
.await
|
||||
.map_err(DiskError::from)??);
|
||||
}
|
||||
|
||||
#[cfg(not(any(unix, windows)))]
|
||||
{
|
||||
let _ = (volume, path, expected, replacement);
|
||||
Err(DiskError::MethodNotAllowed)
|
||||
@@ -21823,9 +21846,9 @@ mod test {
|
||||
assert!(matches!(results[1].as_ref().unwrap_err(), DiskError::Io(_)));
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
#[cfg(any(unix, windows))]
|
||||
#[tokio::test]
|
||||
async fn conditional_file_update_never_deletes_a_new_owner() {
|
||||
async fn windows_and_unix_conditional_file_update_never_deletes_a_new_owner() {
|
||||
use tempfile::tempdir;
|
||||
|
||||
let dir = tempdir().expect("temp dir should be created");
|
||||
@@ -21856,8 +21879,18 @@ mod test {
|
||||
disk.read_all(RUSTFS_META_BUCKET, HEALING_MARKER_PATH)
|
||||
.await
|
||||
.expect("new owner marker should remain"),
|
||||
owner_b
|
||||
owner_b.clone()
|
||||
);
|
||||
assert_eq!(
|
||||
disk.compare_and_update_file(RUSTFS_META_BUCKET, HEALING_MARKER_PATH, Some(owner_b), None)
|
||||
.await
|
||||
.expect("current owner should remove marker"),
|
||||
ConditionalFileUpdate::Updated
|
||||
);
|
||||
assert!(matches!(
|
||||
disk.read_all(RUSTFS_META_BUCKET, HEALING_MARKER_PATH).await,
|
||||
Err(DiskError::FileNotFound)
|
||||
));
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
@@ -21872,7 +21905,10 @@ mod test {
|
||||
let marker_path = disk
|
||||
.get_object_path(RUSTFS_META_BUCKET, HEALING_MARKER_PATH)
|
||||
.expect("marker path should resolve");
|
||||
let lock_path = marker_path.with_extension("rustfs-cas.lock");
|
||||
let lock_path = marker_path
|
||||
.parent()
|
||||
.expect("marker path should have a parent")
|
||||
.join(".rustfs-cas.lock");
|
||||
let lock = std::fs::OpenOptions::new()
|
||||
.create(true)
|
||||
.truncate(false)
|
||||
@@ -21893,6 +21929,40 @@ mod test {
|
||||
assert!(matches!(err, DiskError::Io(ref err) if err.kind() == ErrorKind::WouldBlock));
|
||||
}
|
||||
|
||||
#[cfg(windows)]
|
||||
#[tokio::test]
|
||||
async fn windows_conditional_file_update_returns_would_block_when_marker_lock_is_contended() {
|
||||
let dir = tempfile::tempdir().expect("temp dir should be created");
|
||||
let endpoint = Endpoint::try_from(dir.path().to_str().expect("temp dir should be utf8")).expect("endpoint should parse");
|
||||
let disk = LocalDisk::new(&endpoint, false).await.expect("local disk should be created");
|
||||
ensure_test_volume(&disk, RUSTFS_META_BUCKET).await;
|
||||
let marker_path = disk
|
||||
.get_object_path(RUSTFS_META_BUCKET, HEALING_MARKER_PATH)
|
||||
.expect("marker path should resolve");
|
||||
let lock_path = marker_path
|
||||
.parent()
|
||||
.expect("marker path should have a parent")
|
||||
.join(".rustfs-cas.lock");
|
||||
let lock = std::fs::OpenOptions::new()
|
||||
.create(true)
|
||||
.truncate(false)
|
||||
.read(true)
|
||||
.write(true)
|
||||
.open(lock_path)
|
||||
.expect("marker lock should open");
|
||||
lock.try_lock().expect("marker lock should be held");
|
||||
|
||||
let err = tokio::time::timeout(
|
||||
Duration::from_secs(1),
|
||||
disk.compare_and_update_file(RUSTFS_META_BUCKET, HEALING_MARKER_PATH, None, Some(Bytes::from_static(b"owner"))),
|
||||
)
|
||||
.await
|
||||
.expect("contended conditional update must not block")
|
||||
.expect_err("contended conditional update must retry");
|
||||
|
||||
assert!(matches!(err, DiskError::Io(ref err) if err.kind() == ErrorKind::WouldBlock));
|
||||
}
|
||||
|
||||
#[cfg(target_os = "linux")]
|
||||
#[tokio::test]
|
||||
async fn replacement_io_paths_stay_under_the_mount_lease() {
|
||||
|
||||
@@ -12,6 +12,8 @@
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
#[cfg(windows)]
|
||||
use crate::disk::ConditionalFileUpdate;
|
||||
use crate::disk::error::DiskError;
|
||||
use crate::disk::error::Result;
|
||||
use crate::disk::error_conv::to_file_error;
|
||||
@@ -3458,6 +3460,89 @@ fn read_windows_relative_file(file_path: &Path, parent_guard: &ExistingBaseDirec
|
||||
Ok(Some(data))
|
||||
}
|
||||
|
||||
#[cfg(windows)]
|
||||
pub(crate) fn compare_and_update_control_file(
|
||||
file_path: &Path,
|
||||
expected: Option<&[u8]>,
|
||||
replacement: Option<&[u8]>,
|
||||
sync_metadata: bool,
|
||||
publication_root: &PublicationRoot,
|
||||
) -> io::Result<ConditionalFileUpdate> {
|
||||
use windows_sys::{
|
||||
Wdk::Storage::FileSystem::{
|
||||
FILE_NON_DIRECTORY_FILE, FILE_OPEN, FILE_OPEN_IF, FILE_OPEN_REPARSE_POINT, FILE_SYNCHRONOUS_IO_NONALERT,
|
||||
},
|
||||
Win32::Storage::FileSystem::{
|
||||
DELETE, FILE_ATTRIBUTE_NORMAL, FILE_READ_ATTRIBUTES, FILE_SHARE_READ, FILE_SHARE_WRITE, FILE_WRITE_DATA, SYNCHRONIZE,
|
||||
},
|
||||
};
|
||||
|
||||
let parent = file_path
|
||||
.parent()
|
||||
.ok_or_else(|| io::Error::new(io::ErrorKind::InvalidInput, "conditional file has no parent"))?;
|
||||
let parent_guard = lock_windows_directory_tree(parent, Some(parent), publication_root)?;
|
||||
let lock = open_windows_relative(
|
||||
parent_guard.last_handle()?,
|
||||
std::ffi::OsStr::new(".rustfs-cas.lock"),
|
||||
SYNCHRONIZE | FILE_READ_ATTRIBUTES | FILE_WRITE_DATA,
|
||||
FILE_SHARE_READ | FILE_SHARE_WRITE,
|
||||
FILE_OPEN_IF,
|
||||
FILE_NON_DIRECTORY_FILE | FILE_OPEN_REPARSE_POINT | FILE_SYNCHRONOUS_IO_NONALERT,
|
||||
FILE_ATTRIBUTE_NORMAL,
|
||||
true,
|
||||
)?;
|
||||
validate_windows_owned_file(&lock)?;
|
||||
match lock.as_file().try_lock() {
|
||||
Ok(()) => {}
|
||||
Err(std::fs::TryLockError::WouldBlock) => return Err(io::Error::from(io::ErrorKind::WouldBlock)),
|
||||
Err(std::fs::TryLockError::Error(err)) => return Err(err),
|
||||
}
|
||||
|
||||
let current = read_windows_relative_file(file_path, &parent_guard)?;
|
||||
let matches = match (¤t, expected) {
|
||||
(None, None) => true,
|
||||
(Some(current), Some(expected)) => current.as_slice() == expected,
|
||||
_ => false,
|
||||
};
|
||||
if !matches {
|
||||
return Ok(match current {
|
||||
None => ConditionalFileUpdate::Missing,
|
||||
Some(_) => ConditionalFileUpdate::Mismatch,
|
||||
});
|
||||
}
|
||||
|
||||
match replacement {
|
||||
Some(replacement) => RenameDestinationPathGuard {
|
||||
directory: parent.to_path_buf(),
|
||||
_directory_guard: parent_guard,
|
||||
}
|
||||
.write_file_for_path_access(file_path, replacement, sync_metadata, sync_metadata)?,
|
||||
None => {
|
||||
let file_name = file_path
|
||||
.file_name()
|
||||
.ok_or_else(|| io::Error::new(io::ErrorKind::InvalidInput, "conditional file must have a name"))?;
|
||||
let file = open_windows_relative(
|
||||
parent_guard.last_handle()?,
|
||||
file_name,
|
||||
DELETE | SYNCHRONIZE | FILE_READ_ATTRIBUTES,
|
||||
FILE_SHARE_READ,
|
||||
FILE_OPEN,
|
||||
FILE_NON_DIRECTORY_FILE | FILE_OPEN_REPARSE_POINT | FILE_SYNCHRONOUS_IO_NONALERT,
|
||||
0,
|
||||
true,
|
||||
)?;
|
||||
validate_windows_owned_file(&file)?;
|
||||
set_windows_file_delete_on_close(&file, true)?;
|
||||
drop(file);
|
||||
if sync_metadata {
|
||||
fsync_dir_std(parent)?;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Ok(ConditionalFileUpdate::Updated)
|
||||
}
|
||||
|
||||
#[cfg(windows)]
|
||||
fn open_windows_directory_component(
|
||||
parent: &WindowsDirectoryHandle,
|
||||
|
||||
@@ -784,24 +784,6 @@ pub(crate) fn create_deferred_bitrot_reader_with_stripe_handle(
|
||||
///
|
||||
/// # Returns
|
||||
/// A Result containing the BitrotWriterWrapper or an error
|
||||
/// Size hint handed to `DiskAPI::create_file` for a bitrot-wrapped shard.
|
||||
///
|
||||
/// A known length is grown by one checksum per shard so the on-disk file size
|
||||
/// matches what the bitrot writer emits. A negative length is the
|
||||
/// unknown-size sentinel (`HashReader::SIZE_PRESERVE_LAYER`, used by SSE and
|
||||
/// compression) and must be preserved: `RemoteDisk::create_file` forwards it
|
||||
/// in the `put_file_stream` query, and the receiver only treats `size > 0` as
|
||||
/// a fixed body length when locating the authenticated trailer. Clamping it
|
||||
/// to `0` would claim an empty body and misframe the stream. `0` stays `0`
|
||||
/// because a genuinely empty object still means an empty body.
|
||||
fn bitrot_create_file_size(length: i64, shard_size: usize, checksum_algo: &HashAlgorithm) -> i64 {
|
||||
if length <= 0 {
|
||||
return length;
|
||||
}
|
||||
let length = length as usize;
|
||||
(length.div_ceil(shard_size) * checksum_algo.size() + length) as i64
|
||||
}
|
||||
|
||||
pub async fn create_bitrot_writer(
|
||||
is_inline_buffer: bool,
|
||||
disk: Option<&DiskStore>,
|
||||
@@ -814,7 +796,12 @@ pub async fn create_bitrot_writer(
|
||||
let writer = if is_inline_buffer {
|
||||
CustomWriter::new_inline_buffer()
|
||||
} else if let Some(disk) = disk {
|
||||
let length = bitrot_create_file_size(length, shard_size, &checksum_algo);
|
||||
let length = if length > 0 {
|
||||
let length = length as usize;
|
||||
(length.div_ceil(shard_size) * checksum_algo.size() + length) as i64
|
||||
} else {
|
||||
0
|
||||
};
|
||||
|
||||
let file = disk.create_file("", volume, path, length).await?;
|
||||
#[cfg(feature = "hotpath")]
|
||||
@@ -833,25 +820,6 @@ mod tests {
|
||||
use rustfs_rio::ChunkReader;
|
||||
use std::collections::VecDeque;
|
||||
|
||||
#[test]
|
||||
fn bitrot_create_file_size_grows_known_length_by_checksums() {
|
||||
// 10 bytes over 4-byte shards = 3 shards, each followed by a 32-byte hash.
|
||||
assert_eq!(bitrot_create_file_size(10, 4, &HashAlgorithm::HighwayHash256), 10 + 3 * 32);
|
||||
assert_eq!(bitrot_create_file_size(10, 4, &HashAlgorithm::None), 10);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn bitrot_create_file_size_keeps_empty_and_unknown_distinct() {
|
||||
assert_eq!(bitrot_create_file_size(0, 4, &HashAlgorithm::HighwayHash256), 0);
|
||||
// SSE/compression streams advertise SIZE_PRESERVE_LAYER (-1); the remote
|
||||
// put_file_stream receiver relies on a non-positive size to parse the auth
|
||||
// trailer from the stream tail, so the sentinel must survive untouched.
|
||||
assert_eq!(
|
||||
bitrot_create_file_size(rustfs_rio::HashReader::SIZE_PRESERVE_LAYER, 4, &HashAlgorithm::HighwayHash256),
|
||||
rustfs_rio::HashReader::SIZE_PRESERVE_LAYER
|
||||
);
|
||||
}
|
||||
|
||||
struct TestChunkReader {
|
||||
chunks: VecDeque<Bytes>,
|
||||
}
|
||||
|
||||
@@ -91,6 +91,7 @@ metrics = { workspace = true }
|
||||
base64 = { workspace = true }
|
||||
bytes = { workspace = true }
|
||||
crc-fast = { workspace = true }
|
||||
sha2 = { workspace = true }
|
||||
|
||||
[dev-dependencies]
|
||||
serde_json = { workspace = true, features = ["raw_value"] }
|
||||
|
||||
@@ -373,6 +373,11 @@ impl ErasureSetHealer {
|
||||
set_disk_id: &str,
|
||||
buckets: &[String],
|
||||
) -> Result<(ResumeManager, CheckpointManager)> {
|
||||
if self.replacement_task_id.is_none() && CheckpointManager::is_blocked(&self.disk, task_id).await {
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: format!("Resume task {task_id} has a blocked checkpoint"),
|
||||
});
|
||||
}
|
||||
// check if resume state exists
|
||||
let has_resume_state = if self.replacement_task_id.is_some() {
|
||||
ResumeManager::has_replacement_intent(&self.disk, task_id).await
|
||||
|
||||
@@ -51,6 +51,7 @@ const RESUME_STATE_FILE: &str = "ahm_resume_state.json";
|
||||
const REPLACEMENT_INTENT_FILE: &str = "ahm_replacement_intent.json";
|
||||
const RESUME_PROGRESS_FILE: &str = "ahm_progress.json";
|
||||
pub(super) const RESUME_CHECKPOINT_FILE: &str = "ahm_checkpoint.json";
|
||||
pub(super) const RESUME_CHECKPOINT_BLOCKED_FILE: &str = "ahm_checkpoint.blocked";
|
||||
const REPLACEMENT_COMPLETION_PROOF_FILE: &str = "ahm_replacement_completion_proof.json";
|
||||
const REPLACEMENT_RECOVERY_DIR: &str = "ahm-replacement";
|
||||
const REPLACEMENT_INTENT_SEAL_FILE: &str = "ahm_replacement_intent_seal";
|
||||
|
||||
@@ -13,26 +13,31 @@
|
||||
// limitations under the License.
|
||||
|
||||
use crate::{Error, Result};
|
||||
use base64::Engine as _;
|
||||
use serde::{Deserialize, Serialize};
|
||||
use sha2::{Digest, Sha256};
|
||||
use std::collections::HashSet;
|
||||
use std::path::Path;
|
||||
use std::sync::{Arc, Mutex};
|
||||
use std::time::{SystemTime, UNIX_EPOCH};
|
||||
use tokio::sync::RwLock;
|
||||
use tokio::sync::{Mutex as AsyncMutex, RwLock};
|
||||
use tracing::{debug, warn};
|
||||
|
||||
use super::super::{BUCKET_META_PREFIX, DiskStore, HealDiskExt as _, RUSTFS_META_BUCKET};
|
||||
use super::super::storage_api::owner::{EcstoreConditionalFileUpdate, EcstoreDiskAPI, EcstoreDiskBytes};
|
||||
use super::super::{BUCKET_META_PREFIX, DiskStore, HealDiskExt, RUSTFS_META_BUCKET};
|
||||
use super::{
|
||||
LOG_COMPONENT_HEAL, LOG_SUBSYSTEM_RESUME, PersistThrottle, RESUME_CHECKPOINT_FILE, delete_resume_file, path_to_str,
|
||||
validate_resume_task_id,
|
||||
LOG_COMPONENT_HEAL, LOG_SUBSYSTEM_RESUME, PersistThrottle, RESUME_CHECKPOINT_BLOCKED_FILE, RESUME_CHECKPOINT_FILE,
|
||||
delete_resume_file, path_to_str, validate_resume_task_id,
|
||||
};
|
||||
|
||||
const EVENT_HEAL_CHECKPOINT_STATE: &str = "heal_checkpoint_state";
|
||||
const RESUME_CHECKPOINT_DIGEST_FILE: &str = "ahm_checkpoint.sha256";
|
||||
const CHECKPOINT_PER_VERSION_SCHEMA: u32 = 5;
|
||||
|
||||
/// Current on-disk schema version for `ResumeCheckpoint`. Same rationale as
|
||||
/// `CURRENT_RESUME_SCHEMA`: pre-per-version dedup identities are not comparable
|
||||
/// to the new `compose_key` identities, so a stale checkpoint is discarded.
|
||||
pub(super) const CURRENT_CHECKPOINT_SCHEMA: u32 = 5;
|
||||
pub(super) const CURRENT_CHECKPOINT_SCHEMA: u32 = 6;
|
||||
|
||||
/// resume checkpoint
|
||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||
@@ -57,6 +62,11 @@ pub struct ResumeCheckpoint {
|
||||
pub failed_objects: HashSet<String>,
|
||||
/// skipped objects
|
||||
pub skipped_objects: HashSet<String>,
|
||||
/// Integrity digest over the checkpoint with this field set to `None`.
|
||||
/// Keeping it in the checkpoint makes the payload and its authentication
|
||||
/// record one CAS generation instead of two independently-written files.
|
||||
#[serde(default)]
|
||||
pub integrity_digest: Option<String>,
|
||||
}
|
||||
|
||||
impl ResumeCheckpoint {
|
||||
@@ -70,6 +80,7 @@ impl ResumeCheckpoint {
|
||||
processed_objects: HashSet::new(),
|
||||
failed_objects: HashSet::new(),
|
||||
skipped_objects: HashSet::new(),
|
||||
integrity_digest: None,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -116,17 +127,111 @@ pub struct CheckpointManager {
|
||||
disk: DiskStore,
|
||||
checkpoint: Arc<RwLock<ResumeCheckpoint>>,
|
||||
throttle: Mutex<PersistThrottle>,
|
||||
save_lock: AsyncMutex<()>,
|
||||
last_saved: Mutex<Option<EcstoreDiskBytes>>,
|
||||
}
|
||||
|
||||
impl CheckpointManager {
|
||||
fn blocked_path(task_id: &str) -> std::path::PathBuf {
|
||||
Path::new(BUCKET_META_PREFIX).join(format!("{task_id}_{RESUME_CHECKPOINT_BLOCKED_FILE}"))
|
||||
}
|
||||
|
||||
/// Return whether a checkpoint was permanently isolated after a malformed
|
||||
/// or unsupported snapshot was observed.
|
||||
pub(crate) async fn is_blocked(disk: &DiskStore, task_id: &str) -> bool {
|
||||
if validate_resume_task_id(task_id).is_err() {
|
||||
return false;
|
||||
}
|
||||
let blocked_path = Self::blocked_path(task_id);
|
||||
let Ok(path) = path_to_str(&blocked_path) else {
|
||||
return false;
|
||||
};
|
||||
match HealDiskExt::read_all(disk.as_ref(), RUSTFS_META_BUCKET, path).await {
|
||||
Ok(_) => true,
|
||||
Err(crate::heal::DiskError::FileNotFound) => false,
|
||||
Err(_) => true,
|
||||
}
|
||||
}
|
||||
|
||||
/// Validate the checkpoint while enumerating resumable state. This reads
|
||||
/// the checkpoint once and also isolates malformed or unsupported data.
|
||||
pub(crate) async fn is_resumable(disk: &DiskStore, task_id: &str) -> Result<bool> {
|
||||
validate_resume_task_id(task_id)?;
|
||||
if Self::is_blocked(disk, task_id).await {
|
||||
return Err(Error::InvalidCheckpoint(format!("Resume task {task_id} has a blocked checkpoint")));
|
||||
}
|
||||
let file_path = Path::new(BUCKET_META_PREFIX).join(format!("{task_id}_{RESUME_CHECKPOINT_FILE}"));
|
||||
let Ok(path) = path_to_str(&file_path) else {
|
||||
return Err(Error::InvalidCheckpoint("Resume checkpoint path is not valid UTF-8".to_string()));
|
||||
};
|
||||
match HealDiskExt::read_all(disk.as_ref(), RUSTFS_META_BUCKET, path).await {
|
||||
Ok(bytes) if bytes.is_empty() => Ok(true),
|
||||
Ok(bytes) => Self::load_from_data(disk.clone(), task_id, bytes.to_vec())
|
||||
.await
|
||||
.map(|_| true),
|
||||
Err(crate::heal::DiskError::FileNotFound) => Ok(true),
|
||||
Err(error) => Err(error.into()),
|
||||
}
|
||||
}
|
||||
|
||||
async fn block_invalid_snapshot(disk: &DiskStore, task_id: &str) {
|
||||
// This marker is intentionally version-agnostic: an unsupported reader
|
||||
// must stop selector retries until an operator cleans up the snapshot.
|
||||
let blocked_path = Self::blocked_path(task_id);
|
||||
let Ok(path) = path_to_str(&blocked_path) else {
|
||||
return;
|
||||
};
|
||||
let result = EcstoreDiskAPI::compare_and_update_file(
|
||||
disk.as_ref(),
|
||||
RUSTFS_META_BUCKET,
|
||||
path,
|
||||
None,
|
||||
Some(EcstoreDiskBytes::from_static(b"blocked")),
|
||||
)
|
||||
.await;
|
||||
match result {
|
||||
Ok(EcstoreConditionalFileUpdate::Updated | EcstoreConditionalFileUpdate::Mismatch) => {}
|
||||
Ok(EcstoreConditionalFileUpdate::Missing) => warn!(
|
||||
target: "rustfs::heal::resume",
|
||||
event = EVENT_HEAL_CHECKPOINT_STATE,
|
||||
component = LOG_COMPONENT_HEAL,
|
||||
subsystem = LOG_SUBSYSTEM_RESUME,
|
||||
task_id,
|
||||
state = "blocked_marker_write_failed",
|
||||
error = "marker target disappeared",
|
||||
"Heal checkpoint could not persist its blocked marker"
|
||||
),
|
||||
Err(error) => warn!(
|
||||
target: "rustfs::heal::resume",
|
||||
event = EVENT_HEAL_CHECKPOINT_STATE,
|
||||
component = LOG_COMPONENT_HEAL,
|
||||
subsystem = LOG_SUBSYSTEM_RESUME,
|
||||
task_id,
|
||||
state = "blocked_marker_write_failed",
|
||||
error = %error,
|
||||
"Heal checkpoint could not persist its blocked marker"
|
||||
),
|
||||
}
|
||||
}
|
||||
|
||||
/// create new checkpoint manager
|
||||
pub async fn new(disk: DiskStore, task_id: String) -> Result<Self> {
|
||||
validate_resume_task_id(&task_id)?;
|
||||
let checkpoint_volume = format!("{RUSTFS_META_BUCKET}/{BUCKET_META_PREFIX}");
|
||||
if let Err(error) = EcstoreDiskAPI::make_volume(disk.as_ref(), &checkpoint_volume).await
|
||||
&& error != crate::heal::DiskError::VolumeExists
|
||||
{
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: format!("Failed to create checkpoint volume: {error}"),
|
||||
});
|
||||
}
|
||||
let checkpoint = ResumeCheckpoint::new(task_id);
|
||||
let manager = Self {
|
||||
disk,
|
||||
checkpoint: Arc::new(RwLock::new(checkpoint)),
|
||||
throttle: Mutex::new(PersistThrottle::new()),
|
||||
save_lock: AsyncMutex::new(()),
|
||||
last_saved: Mutex::new(None),
|
||||
};
|
||||
|
||||
// save initial checkpoint
|
||||
@@ -140,6 +245,7 @@ impl CheckpointManager {
|
||||
error = %e,
|
||||
"Heal checkpoint persistence failed"
|
||||
);
|
||||
return Err(e);
|
||||
}
|
||||
Ok(manager)
|
||||
}
|
||||
@@ -148,11 +254,22 @@ impl CheckpointManager {
|
||||
pub async fn load_from_disk(disk: DiskStore, task_id: &str) -> Result<Self> {
|
||||
validate_resume_task_id(task_id)?;
|
||||
let checkpoint_data = Self::read_checkpoint_file(&disk, task_id).await?;
|
||||
let mut checkpoint: ResumeCheckpoint =
|
||||
serde_json::from_slice(&checkpoint_data).map_err(|e| Error::TaskExecutionFailed {
|
||||
message: format!("Failed to deserialize checkpoint: {e}"),
|
||||
})?;
|
||||
Self::load_from_data(disk, task_id, checkpoint_data).await
|
||||
}
|
||||
|
||||
async fn load_from_data(disk: DiskStore, task_id: &str, checkpoint_data: Vec<u8>) -> Result<Self> {
|
||||
validate_resume_task_id(task_id)?;
|
||||
let mut checkpoint: ResumeCheckpoint = match serde_json::from_slice(&checkpoint_data) {
|
||||
Ok(checkpoint) => checkpoint,
|
||||
Err(error) => {
|
||||
Self::block_invalid_snapshot(&disk, task_id).await;
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: format!("Failed to deserialize checkpoint: {error}"),
|
||||
});
|
||||
}
|
||||
};
|
||||
if checkpoint.task_id != task_id {
|
||||
Self::block_invalid_snapshot(&disk, task_id).await;
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: "Resume checkpoint task id does not match filename".to_string(),
|
||||
});
|
||||
@@ -163,6 +280,7 @@ impl CheckpointManager {
|
||||
// identities. Discard the stale sets and position, then stamp the
|
||||
// current schema so the scan restarts cleanly.
|
||||
if checkpoint.schema_version > CURRENT_CHECKPOINT_SCHEMA {
|
||||
Self::block_invalid_snapshot(&disk, task_id).await;
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: format!(
|
||||
"Checkpoint schema {} is newer than supported schema {CURRENT_CHECKPOINT_SCHEMA}",
|
||||
@@ -170,7 +288,45 @@ impl CheckpointManager {
|
||||
),
|
||||
});
|
||||
}
|
||||
if checkpoint.schema_version < CURRENT_CHECKPOINT_SCHEMA {
|
||||
|
||||
let integrity_verified = if let Some(expected) = checkpoint.integrity_digest.as_deref() {
|
||||
let actual = Self::checkpoint_digest(&Self::serialize_without_digest(&checkpoint)?);
|
||||
if expected != actual {
|
||||
Self::block_invalid_snapshot(&disk, task_id).await;
|
||||
return Err(Error::InvalidCheckpoint(format!(
|
||||
"Resume checkpoint digest does not match task {task_id}"
|
||||
)));
|
||||
}
|
||||
true
|
||||
} else if checkpoint.schema_version >= CURRENT_CHECKPOINT_SCHEMA {
|
||||
Self::block_invalid_snapshot(&disk, task_id).await;
|
||||
return Err(Error::InvalidCheckpoint(format!(
|
||||
"Resume checkpoint digest is missing for task {task_id}"
|
||||
)));
|
||||
} else {
|
||||
let digest_path = Self::digest_path(task_id);
|
||||
let digest_path = path_to_str(&digest_path)?;
|
||||
match HealDiskExt::read_all(disk.as_ref(), RUSTFS_META_BUCKET, digest_path).await {
|
||||
Ok(expected) => {
|
||||
let actual = Self::checkpoint_digest(&checkpoint_data);
|
||||
if expected.as_ref() != actual.as_bytes() {
|
||||
Self::block_invalid_snapshot(&disk, task_id).await;
|
||||
return Err(Error::InvalidCheckpoint(format!(
|
||||
"Resume checkpoint digest does not match task {task_id}"
|
||||
)));
|
||||
}
|
||||
true
|
||||
}
|
||||
Err(crate::heal::DiskError::FileNotFound) => false,
|
||||
Err(error) => {
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: format!("Failed to read checkpoint digest: {error}"),
|
||||
});
|
||||
}
|
||||
}
|
||||
};
|
||||
|
||||
if checkpoint.schema_version < CHECKPOINT_PER_VERSION_SCHEMA || !integrity_verified {
|
||||
warn!(
|
||||
target: "rustfs::heal::resume",
|
||||
event = EVENT_HEAL_CHECKPOINT_STATE,
|
||||
@@ -187,13 +343,15 @@ impl CheckpointManager {
|
||||
checkpoint.skipped_objects.clear();
|
||||
checkpoint.current_bucket_index = 0;
|
||||
checkpoint.current_object_index = 0;
|
||||
checkpoint.schema_version = CURRENT_CHECKPOINT_SCHEMA;
|
||||
}
|
||||
checkpoint.schema_version = CURRENT_CHECKPOINT_SCHEMA;
|
||||
|
||||
Ok(Self {
|
||||
disk,
|
||||
checkpoint: Arc::new(RwLock::new(checkpoint)),
|
||||
throttle: Mutex::new(PersistThrottle::new()),
|
||||
save_lock: AsyncMutex::new(()),
|
||||
last_saved: Mutex::new(Some(EcstoreDiskBytes::from(checkpoint_data))),
|
||||
})
|
||||
}
|
||||
|
||||
@@ -204,7 +362,7 @@ impl CheckpointManager {
|
||||
}
|
||||
let file_path = Path::new(BUCKET_META_PREFIX).join(format!("{task_id}_{RESUME_CHECKPOINT_FILE}"));
|
||||
match path_to_str(&file_path) {
|
||||
Ok(path_str) => match disk.read_all(RUSTFS_META_BUCKET, path_str).await {
|
||||
Ok(path_str) => match HealDiskExt::read_all(disk.as_ref(), RUSTFS_META_BUCKET, path_str).await {
|
||||
Ok(data) => !data.is_empty(),
|
||||
Err(_) => false,
|
||||
},
|
||||
@@ -292,6 +450,8 @@ impl CheckpointManager {
|
||||
|
||||
let checkpoint_file = Path::new(BUCKET_META_PREFIX).join(format!("{task_id}_{RESUME_CHECKPOINT_FILE}"));
|
||||
delete_resume_file(&self.disk, &checkpoint_file).await?;
|
||||
delete_resume_file(&self.disk, &Self::digest_path(&task_id)).await?;
|
||||
delete_resume_file(&self.disk, &Self::blocked_path(&task_id)).await?;
|
||||
|
||||
debug!(
|
||||
target: "rustfs::heal::resume",
|
||||
@@ -307,21 +467,130 @@ impl CheckpointManager {
|
||||
|
||||
/// save checkpoint to disk
|
||||
async fn save_checkpoint(&self) -> Result<()> {
|
||||
let checkpoint = self.checkpoint.read().await;
|
||||
// Serialize saves and take the snapshot only after acquiring the lock:
|
||||
// a slower writer must not publish a snapshot taken before a newer one.
|
||||
let _save_guard = self.save_lock.lock().await;
|
||||
let checkpoint = self.checkpoint.read().await.clone();
|
||||
validate_resume_task_id(&checkpoint.task_id)?;
|
||||
let checkpoint_data = serde_json::to_vec(&*checkpoint).map_err(|e| Error::TaskExecutionFailed {
|
||||
message: format!("Failed to serialize checkpoint: {e}"),
|
||||
})?;
|
||||
let unsigned_checkpoint_data = Self::serialize_without_digest(&checkpoint)?;
|
||||
let digest = Self::checkpoint_digest(&unsigned_checkpoint_data);
|
||||
let mut persisted_checkpoint = checkpoint.clone();
|
||||
persisted_checkpoint.integrity_digest = Some(digest);
|
||||
let checkpoint_data =
|
||||
EcstoreDiskBytes::from(serde_json::to_vec(&persisted_checkpoint).map_err(|e| Error::TaskExecutionFailed {
|
||||
message: format!("Failed to serialize checkpoint: {e}"),
|
||||
})?);
|
||||
|
||||
let file_path = Path::new(BUCKET_META_PREFIX).join(format!("{}_{}", checkpoint.task_id, RESUME_CHECKPOINT_FILE));
|
||||
|
||||
let path_str = path_to_str(&file_path)?;
|
||||
self.disk
|
||||
.write_all(RUSTFS_META_BUCKET, path_str, checkpoint_data.into())
|
||||
let last_saved = self
|
||||
.last_saved
|
||||
.lock()
|
||||
.map_err(|_| Error::TaskExecutionFailed {
|
||||
message: "Checkpoint save state lock is poisoned; refusing to save".to_string(),
|
||||
})?
|
||||
.clone();
|
||||
let update = EcstoreDiskAPI::compare_and_update_file(
|
||||
self.disk.as_ref(),
|
||||
RUSTFS_META_BUCKET,
|
||||
path_str,
|
||||
last_saved.clone(),
|
||||
Some(checkpoint_data.clone()),
|
||||
)
|
||||
.await
|
||||
.map_err(|e| Error::TaskExecutionFailed {
|
||||
message: format!("Failed to save checkpoint: {e}"),
|
||||
})?;
|
||||
|
||||
let expected = match update {
|
||||
EcstoreConditionalFileUpdate::Updated => None,
|
||||
EcstoreConditionalFileUpdate::Missing => {
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: "Checkpoint was removed after this manager saved it; refusing to recreate it".to_string(),
|
||||
});
|
||||
}
|
||||
EcstoreConditionalFileUpdate::Mismatch => {
|
||||
// A healthy manager normally completes the CAS above without
|
||||
// another read or JSON parse. Inspect only after a mismatch so
|
||||
// corruption and future schemas cannot be overwritten blindly.
|
||||
let existing = match HealDiskExt::read_all(self.disk.as_ref(), RUSTFS_META_BUCKET, path_str).await {
|
||||
Ok(existing) => existing,
|
||||
Err(crate::heal::DiskError::FileNotFound) => {
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: "Checkpoint was removed after this manager saved it; refusing to recreate it".to_string(),
|
||||
});
|
||||
}
|
||||
Err(error) => {
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: format!("Failed to inspect checkpoint after CAS mismatch: {error}"),
|
||||
});
|
||||
}
|
||||
};
|
||||
|
||||
if existing.is_empty() && last_saved.is_none() {
|
||||
Some(existing)
|
||||
} else {
|
||||
let current: ResumeCheckpoint = match serde_json::from_slice(&existing) {
|
||||
Ok(current) => current,
|
||||
Err(error) => {
|
||||
Self::block_invalid_snapshot(&self.disk, &checkpoint.task_id).await;
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: format!("Existing checkpoint is corrupt: {error}"),
|
||||
});
|
||||
}
|
||||
};
|
||||
if current.task_id != checkpoint.task_id {
|
||||
Self::block_invalid_snapshot(&self.disk, &checkpoint.task_id).await;
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: "Existing checkpoint task id does not match filename".to_string(),
|
||||
});
|
||||
}
|
||||
if current.schema_version > CURRENT_CHECKPOINT_SCHEMA {
|
||||
Self::block_invalid_snapshot(&self.disk, &checkpoint.task_id).await;
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: format!(
|
||||
"Existing checkpoint schema {} is newer than supported schema {CURRENT_CHECKPOINT_SCHEMA}",
|
||||
current.schema_version
|
||||
),
|
||||
});
|
||||
}
|
||||
if last_saved.as_ref().is_none_or(|saved| saved.as_ref() != existing.as_ref()) {
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: "Checkpoint changed since this manager loaded it; refusing to overwrite newer progress"
|
||||
.to_string(),
|
||||
});
|
||||
}
|
||||
Some(existing)
|
||||
}
|
||||
}
|
||||
};
|
||||
|
||||
if let Some(expected) = expected {
|
||||
match EcstoreDiskAPI::compare_and_update_file(
|
||||
self.disk.as_ref(),
|
||||
RUSTFS_META_BUCKET,
|
||||
path_str,
|
||||
Some(expected),
|
||||
Some(checkpoint_data.clone()),
|
||||
)
|
||||
.await
|
||||
.map_err(|e| Error::TaskExecutionFailed {
|
||||
message: format!("Failed to save checkpoint: {e}"),
|
||||
})?;
|
||||
message: format!("Failed to save checkpoint after CAS mismatch: {e}"),
|
||||
})? {
|
||||
EcstoreConditionalFileUpdate::Updated => {}
|
||||
EcstoreConditionalFileUpdate::Missing | EcstoreConditionalFileUpdate::Mismatch => {
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: "Checkpoint changed while saving; refusing to overwrite newer progress".to_string(),
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
let mut last_saved = self.last_saved.lock().map_err(|_| Error::TaskExecutionFailed {
|
||||
message: "Checkpoint save state lock is poisoned after save".to_string(),
|
||||
})?;
|
||||
*last_saved = Some(checkpoint_data);
|
||||
|
||||
debug!(
|
||||
target: "rustfs::heal::resume",
|
||||
@@ -341,11 +610,38 @@ impl CheckpointManager {
|
||||
let file_path = Path::new(BUCKET_META_PREFIX).join(format!("{task_id}_{RESUME_CHECKPOINT_FILE}"));
|
||||
|
||||
let path_str = path_to_str(&file_path)?;
|
||||
disk.read_all(RUSTFS_META_BUCKET, path_str)
|
||||
HealDiskExt::read_all(disk.as_ref(), RUSTFS_META_BUCKET, path_str)
|
||||
.await
|
||||
.map(|bytes| bytes.to_vec())
|
||||
.map_err(|e| Error::TaskExecutionFailed {
|
||||
message: format!("Failed to read checkpoint file: {e}"),
|
||||
})
|
||||
}
|
||||
|
||||
fn serialize_without_digest(checkpoint: &ResumeCheckpoint) -> Result<Vec<u8>> {
|
||||
let mut unsigned = checkpoint.clone();
|
||||
unsigned.integrity_digest = None;
|
||||
let mut value = serde_json::to_value(&unsigned).map_err(|e| Error::TaskExecutionFailed {
|
||||
message: format!("Failed to serialize checkpoint: {e}"),
|
||||
})?;
|
||||
for field in ["processed_objects", "failed_objects", "skipped_objects"] {
|
||||
let Some(values) = value.get_mut(field).and_then(serde_json::Value::as_array_mut) else {
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: format!("Failed to canonicalize checkpoint field: {field}"),
|
||||
});
|
||||
};
|
||||
values.sort_by(|left, right| left.as_str().cmp(&right.as_str()));
|
||||
}
|
||||
serde_json::to_vec(&value).map_err(|e| Error::TaskExecutionFailed {
|
||||
message: format!("Failed to serialize checkpoint: {e}"),
|
||||
})
|
||||
}
|
||||
|
||||
fn checkpoint_digest(checkpoint_data: &[u8]) -> String {
|
||||
base64::engine::general_purpose::STANDARD.encode(Sha256::digest(checkpoint_data))
|
||||
}
|
||||
|
||||
fn digest_path(task_id: &str) -> std::path::PathBuf {
|
||||
Path::new(BUCKET_META_PREFIX).join(format!("{task_id}_{RESUME_CHECKPOINT_DIGEST_FILE}"))
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1600,6 +1600,32 @@ async fn test_checkpoint_schema_v4_discarded_on_load() {
|
||||
temp_dir.close().expect("remove schema test directory");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn downgraded_unsigned_checkpoint_resets_untrusted_progress() {
|
||||
let (temp_dir, disk) = schema_test_disk().await;
|
||||
let task_id = ResumeUtils::generate_task_id();
|
||||
let manager = CheckpointManager::new(disk.clone(), task_id.clone()).await.unwrap();
|
||||
manager.add_processed_object("victim-a".to_string()).await.unwrap();
|
||||
manager.update_position(2, 500).await.unwrap();
|
||||
let checkpoint_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_CHECKPOINT_FILE}");
|
||||
let bytes = disk.read_all(RUSTFS_META_BUCKET, &checkpoint_path).await.unwrap();
|
||||
let mut downgraded: serde_json::Value = serde_json::from_slice(&bytes).unwrap();
|
||||
downgraded["schema_version"] = serde_json::json!(CURRENT_CHECKPOINT_SCHEMA - 1);
|
||||
downgraded.as_object_mut().unwrap().remove("integrity_digest");
|
||||
downgraded["processed_objects"] = serde_json::json!(["victim-b"]);
|
||||
disk.write_all(RUSTFS_META_BUCKET, &checkpoint_path, serde_json::to_vec(&downgraded).unwrap().into())
|
||||
.await
|
||||
.expect("write downgraded checkpoint");
|
||||
|
||||
let manager = CheckpointManager::load_from_disk(disk, &task_id).await.unwrap();
|
||||
let checkpoint = manager.get_checkpoint().await;
|
||||
assert_eq!(checkpoint.schema_version, CURRENT_CHECKPOINT_SCHEMA);
|
||||
assert_eq!(checkpoint.current_bucket_index, 0);
|
||||
assert_eq!(checkpoint.current_object_index, 0);
|
||||
assert!(checkpoint.processed_objects.is_empty());
|
||||
temp_dir.close().unwrap();
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn current_normal_resume_schema_preserves_progress() {
|
||||
let (temp_dir, disk) = schema_test_disk().await;
|
||||
@@ -1675,6 +1701,369 @@ async fn future_resume_and_checkpoint_schemas_are_rejected() {
|
||||
temp_dir.close().expect("remove schema test directory");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn checkpoint_save_does_not_replace_a_non_empty_truncated_snapshot() {
|
||||
let (temp_dir, disk) = schema_test_disk().await;
|
||||
let task_id = ResumeUtils::generate_task_id();
|
||||
let manager = CheckpointManager::new(disk.clone(), task_id.clone())
|
||||
.await
|
||||
.expect("create checkpoint manager");
|
||||
let checkpoint_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_CHECKPOINT_FILE}");
|
||||
let truncated = b"{\"schema_version\":5,\"task_id\":";
|
||||
disk.write_all(RUSTFS_META_BUCKET, &checkpoint_path, truncated.as_slice().into())
|
||||
.await
|
||||
.expect("write truncated checkpoint fixture");
|
||||
|
||||
let error = manager
|
||||
.update_position(2, 7)
|
||||
.await
|
||||
.expect_err("a truncated checkpoint must fail closed during save");
|
||||
assert!(error.to_string().contains("Existing checkpoint is corrupt"));
|
||||
assert_eq!(
|
||||
disk.read_all(RUSTFS_META_BUCKET, &checkpoint_path)
|
||||
.await
|
||||
.expect("read truncated checkpoint fixture"),
|
||||
truncated.as_slice()
|
||||
);
|
||||
assert!(CheckpointManager::is_blocked(&disk, &task_id).await);
|
||||
temp_dir.close().expect("remove checkpoint save test directory");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn checkpoint_save_does_not_replace_a_future_schema_snapshot() {
|
||||
let (temp_dir, disk) = schema_test_disk().await;
|
||||
let task_id = ResumeUtils::generate_task_id();
|
||||
let manager = CheckpointManager::new(disk.clone(), task_id.clone())
|
||||
.await
|
||||
.expect("create checkpoint manager");
|
||||
let checkpoint_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_CHECKPOINT_FILE}");
|
||||
let mut future = ResumeCheckpoint::new(task_id.clone());
|
||||
future.schema_version = CURRENT_CHECKPOINT_SCHEMA + 1;
|
||||
let future_bytes = serde_json::to_vec(&future).expect("serialize future checkpoint fixture");
|
||||
disk.write_all(RUSTFS_META_BUCKET, &checkpoint_path, future_bytes.clone().into())
|
||||
.await
|
||||
.expect("write future checkpoint fixture");
|
||||
|
||||
let error = manager
|
||||
.update_position(2, 7)
|
||||
.await
|
||||
.expect_err("a future schema must fail closed during save");
|
||||
assert!(error.to_string().contains("Existing checkpoint schema"));
|
||||
assert_eq!(
|
||||
disk.read_all(RUSTFS_META_BUCKET, &checkpoint_path)
|
||||
.await
|
||||
.expect("read future checkpoint fixture"),
|
||||
future_bytes
|
||||
);
|
||||
assert!(CheckpointManager::is_blocked(&disk, &task_id).await);
|
||||
temp_dir.close().expect("remove future schema test directory");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn checkpoint_digest_rejects_same_length_progress_tampering() {
|
||||
let (temp_dir, disk) = schema_test_disk().await;
|
||||
let task_id = ResumeUtils::generate_task_id();
|
||||
let manager = CheckpointManager::new(disk.clone(), task_id.clone())
|
||||
.await
|
||||
.expect("create checkpoint manager");
|
||||
manager
|
||||
.add_processed_object("victim-a".to_string())
|
||||
.await
|
||||
.expect("persist checkpoint progress");
|
||||
manager.update_position(1, 1).await.expect("flush checkpoint progress");
|
||||
let checkpoint_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_CHECKPOINT_FILE}");
|
||||
let original = disk
|
||||
.read_all(RUSTFS_META_BUCKET, &checkpoint_path)
|
||||
.await
|
||||
.expect("read checkpoint fixture");
|
||||
let tampered = original
|
||||
.windows(b"victim-a".len())
|
||||
.position(|window| window == b"victim-a")
|
||||
.map(|index| {
|
||||
let mut bytes = original.to_vec();
|
||||
bytes[index..index + b"victim-a".len()].copy_from_slice(b"victim-b");
|
||||
bytes
|
||||
})
|
||||
.expect("checkpoint should contain the processed object");
|
||||
disk.write_all(RUSTFS_META_BUCKET, &checkpoint_path, tampered.into())
|
||||
.await
|
||||
.expect("write tampered checkpoint fixture");
|
||||
|
||||
assert!(CheckpointManager::load_from_disk(disk.clone(), &task_id).await.is_err());
|
||||
assert!(CheckpointManager::is_blocked(&disk, &task_id).await);
|
||||
temp_dir.close().expect("remove digest test directory");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn checkpoint_integrity_survives_missing_legacy_sidecar() {
|
||||
let (temp_dir, disk) = schema_test_disk().await;
|
||||
let task_id = ResumeUtils::generate_task_id();
|
||||
let manager = CheckpointManager::new(disk.clone(), task_id.clone()).await.unwrap();
|
||||
manager.update_position(2, 9).await.unwrap();
|
||||
|
||||
let digest_path = format!("{BUCKET_META_PREFIX}/{task_id}_ahm_checkpoint.sha256");
|
||||
delete_resume_file(&disk, Path::new(&digest_path)).await.unwrap();
|
||||
|
||||
let restored = CheckpointManager::load_from_disk(disk, &task_id).await.unwrap();
|
||||
let checkpoint = restored.get_checkpoint().await;
|
||||
assert_eq!(checkpoint.current_bucket_index, 2);
|
||||
assert_eq!(checkpoint.current_object_index, 9);
|
||||
assert!(checkpoint.integrity_digest.is_some());
|
||||
temp_dir.close().unwrap();
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn checkpoint_integrity_survives_multi_object_reload() {
|
||||
let (temp_dir, disk) = schema_test_disk().await;
|
||||
let task_id = ResumeUtils::generate_task_id();
|
||||
let manager = CheckpointManager::new(disk.clone(), task_id.clone()).await.unwrap();
|
||||
for index in 0..32 {
|
||||
manager.add_processed_object(format!("processed-{index}")).await.unwrap();
|
||||
manager.add_failed_object(format!("failed-{index}")).await.unwrap();
|
||||
manager.add_skipped_object(format!("skipped-{index}")).await.unwrap();
|
||||
}
|
||||
manager.update_position(2, 9).await.unwrap();
|
||||
|
||||
CheckpointManager::load_from_disk(disk, &task_id)
|
||||
.await
|
||||
.expect("a healthy multi-object checkpoint must survive reload");
|
||||
temp_dir.close().unwrap();
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn checkpoint_integrity_rejects_a_removed_embedded_digest() {
|
||||
let (temp_dir, disk) = schema_test_disk().await;
|
||||
let task_id = ResumeUtils::generate_task_id();
|
||||
let manager = CheckpointManager::new(disk.clone(), task_id.clone()).await.unwrap();
|
||||
manager.update_position(2, 9).await.unwrap();
|
||||
|
||||
let checkpoint_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_CHECKPOINT_FILE}");
|
||||
let bytes = disk
|
||||
.read_all(RUSTFS_META_BUCKET, &checkpoint_path)
|
||||
.await
|
||||
.expect("read checkpoint fixture");
|
||||
let mut value: serde_json::Value = serde_json::from_slice(&bytes).unwrap();
|
||||
value["current_object_index"] = serde_json::json!(10);
|
||||
value.as_object_mut().unwrap().remove("integrity_digest");
|
||||
disk.write_all(RUSTFS_META_BUCKET, &checkpoint_path, serde_json::to_vec(&value).unwrap().into())
|
||||
.await
|
||||
.expect("write tampered checkpoint fixture");
|
||||
|
||||
assert!(
|
||||
CheckpointManager::load_from_disk(disk.clone(), &task_id).await.is_err(),
|
||||
"a current checkpoint without its embedded digest must fail closed"
|
||||
);
|
||||
assert!(CheckpointManager::is_blocked(&disk, &task_id).await);
|
||||
temp_dir.close().unwrap();
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn new_checkpoint_manager_rebuilds_an_empty_snapshot() {
|
||||
let (temp_dir, disk) = schema_test_disk().await;
|
||||
let task_id = ResumeUtils::generate_task_id();
|
||||
let checkpoint_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_CHECKPOINT_FILE}");
|
||||
disk.write_all(RUSTFS_META_BUCKET, &checkpoint_path, EcstoreDiskBytes::new())
|
||||
.await
|
||||
.expect("write empty checkpoint fixture");
|
||||
|
||||
let manager = CheckpointManager::new(disk.clone(), task_id.clone())
|
||||
.await
|
||||
.expect("a new manager must rebuild an empty checkpoint");
|
||||
manager
|
||||
.update_position(3, 11)
|
||||
.await
|
||||
.expect("rebuilt checkpoint must remain writable");
|
||||
assert!(CheckpointManager::has_checkpoint(&disk, &task_id).await);
|
||||
temp_dir.close().expect("remove empty checkpoint test directory");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn deleted_checkpoint_is_not_recreated_by_an_old_manager() {
|
||||
let (temp_dir, disk) = schema_test_disk().await;
|
||||
let task_id = ResumeUtils::generate_task_id();
|
||||
let manager = CheckpointManager::new(disk.clone(), task_id.clone())
|
||||
.await
|
||||
.expect("create checkpoint manager");
|
||||
manager.cleanup().await.expect("delete checkpoint fixture");
|
||||
|
||||
let error = manager
|
||||
.update_position(1, 2)
|
||||
.await
|
||||
.expect_err("an old manager must not resurrect a deleted checkpoint");
|
||||
assert!(error.to_string().contains("removed after this manager saved it"));
|
||||
assert!(!CheckpointManager::has_checkpoint(&disk, &task_id).await);
|
||||
temp_dir.close().expect("remove deleted checkpoint test directory");
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
#[tokio::test]
|
||||
async fn checkpoint_cleanup_leaves_no_task_specific_lock_artifact() {
|
||||
let (temp_dir, disk) = schema_test_disk().await;
|
||||
let task_id = ResumeUtils::generate_task_id();
|
||||
let manager = CheckpointManager::new(disk.clone(), task_id.clone())
|
||||
.await
|
||||
.expect("create checkpoint manager");
|
||||
let lock_path = Path::new(BUCKET_META_PREFIX)
|
||||
.join(format!("{task_id}_{RESUME_CHECKPOINT_FILE}"))
|
||||
.with_extension("rustfs-cas.lock");
|
||||
let lock_path = temp_dir.path().join(RUSTFS_META_BUCKET).join(lock_path);
|
||||
|
||||
manager.cleanup().await.expect("delete checkpoint fixture");
|
||||
|
||||
assert!(
|
||||
!lock_path.exists(),
|
||||
"successful checkpoint cleanup must not leave a task-specific lock artifact"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn an_empty_blocked_marker_still_blocks_resume_selection() {
|
||||
let (temp_dir, disk) = schema_test_disk().await;
|
||||
let task_id = ResumeUtils::generate_task_id();
|
||||
let manager = CheckpointManager::new(disk.clone(), task_id.clone())
|
||||
.await
|
||||
.expect("create checkpoint manager");
|
||||
let blocked_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_CHECKPOINT_BLOCKED_FILE}");
|
||||
disk.write_all(RUSTFS_META_BUCKET, &blocked_path, EcstoreDiskBytes::new())
|
||||
.await
|
||||
.expect("write empty blocked marker fixture");
|
||||
|
||||
assert!(CheckpointManager::is_blocked(&disk, &task_id).await);
|
||||
assert!(CheckpointManager::is_resumable(&disk, &task_id).await.is_err());
|
||||
// Recovery requires replacing/cleaning the snapshot, then removing the
|
||||
// marker; ordinary selector retries are intentionally not an unlock path.
|
||||
manager.cleanup().await.expect("clean blocked checkpoint");
|
||||
assert!(!CheckpointManager::is_blocked(&disk, &task_id).await);
|
||||
temp_dir.close().expect("remove empty blocked marker test directory");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn resumable_selector_skips_healthy_tasks_with_blocked_markers() {
|
||||
let (temp_dir, disk) = schema_test_disk().await;
|
||||
let tasks = [
|
||||
(ResumeUtils::generate_task_id(), EcstoreDiskBytes::new()),
|
||||
(ResumeUtils::generate_task_id(), EcstoreDiskBytes::from_static(b"blocked")),
|
||||
];
|
||||
for (task_id, marker) in &tasks {
|
||||
ResumeManager::new(
|
||||
disk.clone(),
|
||||
task_id.clone(),
|
||||
"erasure_set".to_string(),
|
||||
"pool_0_set_0".to_string(),
|
||||
vec!["bucket".to_string()],
|
||||
)
|
||||
.await
|
||||
.expect("create healthy resume state");
|
||||
CheckpointManager::new(disk.clone(), task_id.clone())
|
||||
.await
|
||||
.expect("create healthy checkpoint");
|
||||
let checkpoint_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_CHECKPOINT_FILE}");
|
||||
let checkpoint_bytes = disk
|
||||
.read_all(RUSTFS_META_BUCKET, &checkpoint_path)
|
||||
.await
|
||||
.expect("read healthy checkpoint before blocking");
|
||||
let marker_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_CHECKPOINT_BLOCKED_FILE}");
|
||||
disk.write_all(RUSTFS_META_BUCKET, &marker_path, marker.clone())
|
||||
.await
|
||||
.expect("write blocked marker");
|
||||
|
||||
assert!(ResumeUtils::get_resumable_tasks(&disk).await.is_err());
|
||||
assert_eq!(
|
||||
disk.read_all(RUSTFS_META_BUCKET, &checkpoint_path)
|
||||
.await
|
||||
.expect("read healthy checkpoint after blocking"),
|
||||
checkpoint_bytes
|
||||
);
|
||||
}
|
||||
temp_dir.close().expect("remove blocked selector test directory");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn stale_checkpoint_manager_cannot_overwrite_newer_progress() {
|
||||
let (temp_dir, disk) = schema_test_disk().await;
|
||||
let task_id = ResumeUtils::generate_task_id();
|
||||
let first = CheckpointManager::new(disk.clone(), task_id.clone())
|
||||
.await
|
||||
.expect("create first checkpoint manager");
|
||||
let second = CheckpointManager::load_from_disk(disk.clone(), &task_id)
|
||||
.await
|
||||
.expect("load second checkpoint manager");
|
||||
|
||||
second
|
||||
.update_position(4, 20)
|
||||
.await
|
||||
.expect("persist newer checkpoint progress");
|
||||
let error = first
|
||||
.update_position(1, 3)
|
||||
.await
|
||||
.expect_err("stale checkpoint manager must not overwrite newer progress");
|
||||
assert!(error.to_string().contains("newer progress"));
|
||||
|
||||
let persisted = CheckpointManager::load_from_disk(disk.clone(), &task_id)
|
||||
.await
|
||||
.expect("load newer checkpoint progress")
|
||||
.get_checkpoint()
|
||||
.await;
|
||||
assert_eq!(persisted.current_bucket_index, 4);
|
||||
assert_eq!(persisted.current_object_index, 20);
|
||||
temp_dir.close().expect("remove stale manager test directory");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn resumable_selector_isolates_future_and_corrupt_checkpoints() {
|
||||
let (temp_dir, disk) = schema_test_disk().await;
|
||||
let future_task = ResumeUtils::generate_task_id();
|
||||
let corrupt_task = ResumeUtils::generate_task_id();
|
||||
for task_id in [&future_task, &corrupt_task] {
|
||||
ResumeManager::new(
|
||||
disk.clone(),
|
||||
task_id.to_string(),
|
||||
"erasure_set".to_string(),
|
||||
"pool_0_set_0".to_string(),
|
||||
vec!["bucket".to_string()],
|
||||
)
|
||||
.await
|
||||
.expect("create resumable state fixture");
|
||||
}
|
||||
|
||||
let future_path = format!("{BUCKET_META_PREFIX}/{future_task}_{RESUME_CHECKPOINT_FILE}");
|
||||
let mut future = ResumeCheckpoint::new(future_task.clone());
|
||||
future.schema_version = CURRENT_CHECKPOINT_SCHEMA + 1;
|
||||
let future_bytes = serde_json::to_vec(&future).expect("serialize future checkpoint fixture");
|
||||
disk.write_all(RUSTFS_META_BUCKET, &future_path, future_bytes.clone().into())
|
||||
.await
|
||||
.expect("write future checkpoint fixture");
|
||||
let corrupt_path = format!("{BUCKET_META_PREFIX}/{corrupt_task}_{RESUME_CHECKPOINT_FILE}");
|
||||
let corrupt_bytes = b"{truncated";
|
||||
disk.write_all(RUSTFS_META_BUCKET, &corrupt_path, corrupt_bytes.as_slice().into())
|
||||
.await
|
||||
.expect("write corrupt checkpoint fixture");
|
||||
|
||||
assert!(CheckpointManager::is_resumable(&disk, &future_task).await.is_err());
|
||||
assert!(CheckpointManager::is_resumable(&disk, &corrupt_task).await.is_err());
|
||||
assert!(ResumeUtils::get_resumable_tasks(&disk).await.is_err());
|
||||
for (task_id, path, bytes) in [
|
||||
(&future_task, future_path, future_bytes),
|
||||
(&corrupt_task, corrupt_path, corrupt_bytes.to_vec()),
|
||||
] {
|
||||
assert_eq!(
|
||||
disk.read_all(RUSTFS_META_BUCKET, &path)
|
||||
.await
|
||||
.expect("read isolated checkpoint bytes"),
|
||||
bytes
|
||||
);
|
||||
let blocked_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_CHECKPOINT_BLOCKED_FILE}");
|
||||
assert!(
|
||||
!disk
|
||||
.read_all(RUSTFS_META_BUCKET, &blocked_path)
|
||||
.await
|
||||
.expect("read checkpoint blocked marker")
|
||||
.is_empty()
|
||||
);
|
||||
}
|
||||
temp_dir.close().expect("remove selector isolation test directory");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_persist_throttle_batches_until_threshold() {
|
||||
let mut throttle = PersistThrottle::new();
|
||||
|
||||
@@ -21,7 +21,7 @@ use uuid::Uuid;
|
||||
use super::super::{BUCKET_META_PREFIX, DiskError, DiskStore, HealDiskExt as _, RUSTFS_META_BUCKET};
|
||||
use super::replacement::{ReplacementPhase, ReplacementRecoveryRecord};
|
||||
use super::{
|
||||
EVENT_HEAL_RESUME_STATE, LOG_COMPONENT_HEAL, LOG_SUBSYSTEM_RESUME, REPLACEMENT_COMPLETION_PROOF_FILE,
|
||||
CheckpointManager, EVENT_HEAL_RESUME_STATE, LOG_COMPONENT_HEAL, LOG_SUBSYSTEM_RESUME, REPLACEMENT_COMPLETION_PROOF_FILE,
|
||||
REPLACEMENT_INTENT_FILE, RESUME_STATE_FILE, ResumeManager, ResumeStateFile, is_replacement_intent, path_to_str,
|
||||
replacement_recovery_corruption_for_state_load, replacement_recovery_dir, validate_resume_task_id,
|
||||
};
|
||||
@@ -67,6 +67,7 @@ impl ResumeUtils {
|
||||
// Extract task ID from filename: {task_id}_ahm_resume_state.json
|
||||
if let Some(task_id) = entry.strip_suffix(&format!("_{RESUME_STATE_FILE}"))
|
||||
&& validate_resume_task_id(task_id).is_ok()
|
||||
&& CheckpointManager::is_resumable(disk, task_id).await?
|
||||
{
|
||||
task_ids.push(task_id.to_string());
|
||||
}
|
||||
|
||||
@@ -16,14 +16,14 @@
|
||||
//!
|
||||
//! `scripts/test/vault_ha_kms_live.sh` owns the official Vault containers and
|
||||
//! kills the active node while this test continuously decrypts through a
|
||||
//! surviving standby. KV2 and Transit must recover after the bounded circuit
|
||||
//! interval, use a bounded number of attempts, and leave the circuit and
|
||||
//! in-flight gauges at zero after a new leader is elected.
|
||||
//! surviving standby. KV2 and Transit requests must remain successful, use a
|
||||
//! bounded number of attempts, and leave the circuit and in-flight gauges at
|
||||
//! zero after a new leader is elected.
|
||||
|
||||
use std::collections::HashMap;
|
||||
use std::path::{Path, PathBuf};
|
||||
use std::sync::Arc;
|
||||
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
|
||||
use std::sync::{Arc, Mutex};
|
||||
use std::time::Duration;
|
||||
|
||||
use metrics_util::MetricKind;
|
||||
@@ -43,11 +43,6 @@ const OPERATION_ATTEMPTS: &str = "rustfs_kms_backend_operation_attempts";
|
||||
const IN_FLIGHT: &str = "rustfs_kms_backend_in_flight";
|
||||
const CIRCUIT_OPEN: &str = "rustfs_kms_backend_circuit_open";
|
||||
const MAX_ATTEMPTS: u32 = 10;
|
||||
const ATTEMPT_TIMEOUT: Duration = Duration::from_secs(2);
|
||||
const HEALTHY_PROGRESS_TIMEOUT: Duration = Duration::from_secs(20);
|
||||
// The circuit remains open for 30s after five failed attempts.
|
||||
const POST_FAILOVER_PROGRESS_TIMEOUT: Duration = Duration::from_secs(35);
|
||||
const FAILOVER_ERROR_POLL_INTERVAL: Duration = Duration::from_millis(100);
|
||||
|
||||
type MetricEntry = (
|
||||
metrics_util::CompositeKey,
|
||||
@@ -69,7 +64,7 @@ fn config(backend: KmsBackend, backend_config: BackendConfig) -> KmsConfig {
|
||||
backend,
|
||||
backend_config,
|
||||
allow_insecure_dev_defaults: true,
|
||||
timeout: ATTEMPT_TIMEOUT,
|
||||
timeout: Duration::from_secs(2),
|
||||
retry_attempts: MAX_ATTEMPTS,
|
||||
enable_cache: false,
|
||||
..KmsConfig::default()
|
||||
@@ -169,31 +164,14 @@ fn retryable_failures(snapshot: &[MetricEntry], operation: &str) -> u64 {
|
||||
.sum()
|
||||
}
|
||||
|
||||
async fn wait_for_count(
|
||||
counter: &AtomicU64,
|
||||
failure: &Mutex<Option<String>>,
|
||||
minimum: u64,
|
||||
description: &str,
|
||||
timeout: Duration,
|
||||
) {
|
||||
tokio::time::timeout(timeout, async {
|
||||
async fn wait_for_count(counter: &AtomicU64, minimum: u64, description: &str) {
|
||||
tokio::time::timeout(Duration::from_secs(20), async {
|
||||
while counter.load(Ordering::SeqCst) < minimum {
|
||||
if let Some(error) = failure.lock().expect("decrypt failure lock poisoned").as_ref() {
|
||||
panic!(
|
||||
"{description} worker failed after {} successful decrypts: {error}",
|
||||
counter.load(Ordering::SeqCst)
|
||||
);
|
||||
}
|
||||
tokio::time::sleep(Duration::from_millis(25)).await;
|
||||
}
|
||||
})
|
||||
.await
|
||||
.unwrap_or_else(|_| {
|
||||
panic!(
|
||||
"timed out after {timeout:?} waiting for {description}: completed {}, expected {minimum}",
|
||||
counter.load(Ordering::SeqCst)
|
||||
)
|
||||
});
|
||||
.unwrap_or_else(|_| panic!("timed out waiting for {description}"));
|
||||
}
|
||||
|
||||
async fn wait_for_file(path: &Path, description: &str) {
|
||||
@@ -211,8 +189,7 @@ async fn decrypt_loop<B: KmsBackendTrait + Send + Sync + 'static>(
|
||||
request: DecryptRequest,
|
||||
expected: Vec<u8>,
|
||||
completed: Arc<AtomicU64>,
|
||||
allow_failover_errors: Arc<AtomicBool>,
|
||||
failure: Arc<Mutex<Option<String>>>,
|
||||
failed: Arc<AtomicBool>,
|
||||
stop: CancellationToken,
|
||||
) {
|
||||
while !stop.is_cancelled() {
|
||||
@@ -220,18 +197,8 @@ async fn decrypt_loop<B: KmsBackendTrait + Send + Sync + 'static>(
|
||||
Ok(response) if response.plaintext == expected => {
|
||||
completed.fetch_add(1, Ordering::SeqCst);
|
||||
}
|
||||
Ok(_) => {
|
||||
*failure.lock().expect("decrypt failure lock poisoned") =
|
||||
Some("decrypt returned unexpected plaintext".to_string());
|
||||
return;
|
||||
}
|
||||
Err(rustfs_kms::KmsError::BackendError { .. } | rustfs_kms::KmsError::OperationTimedOut { .. })
|
||||
if allow_failover_errors.load(Ordering::SeqCst) =>
|
||||
{
|
||||
tokio::time::sleep(FAILOVER_ERROR_POLL_INTERVAL).await;
|
||||
}
|
||||
Err(error) => {
|
||||
*failure.lock().expect("decrypt failure lock poisoned") = Some(error.to_string());
|
||||
Ok(_) | Err(_) => {
|
||||
failed.store(true, Ordering::SeqCst);
|
||||
return;
|
||||
}
|
||||
}
|
||||
@@ -329,9 +296,7 @@ async fn exercise_failover(snapshotter: &Snapshotter) {
|
||||
);
|
||||
|
||||
let stop = CancellationToken::new();
|
||||
let allow_failover_errors = Arc::new(AtomicBool::new(false));
|
||||
let kv2_failure = Arc::new(Mutex::new(None));
|
||||
let transit_failure = Arc::new(Mutex::new(None));
|
||||
let failed = Arc::new(AtomicBool::new(false));
|
||||
let kv2_completed = Arc::new(AtomicU64::new(0));
|
||||
let transit_completed = Arc::new(AtomicU64::new(0));
|
||||
let kv2_worker = tokio::spawn(decrypt_loop(
|
||||
@@ -339,8 +304,7 @@ async fn exercise_failover(snapshotter: &Snapshotter) {
|
||||
kv2_request,
|
||||
kv2_data_key.plaintext_key,
|
||||
Arc::clone(&kv2_completed),
|
||||
Arc::clone(&allow_failover_errors),
|
||||
Arc::clone(&kv2_failure),
|
||||
Arc::clone(&failed),
|
||||
stop.clone(),
|
||||
));
|
||||
let transit_worker = tokio::spawn(decrypt_loop(
|
||||
@@ -348,21 +312,12 @@ async fn exercise_failover(snapshotter: &Snapshotter) {
|
||||
transit_request,
|
||||
transit_data_key.plaintext_key,
|
||||
Arc::clone(&transit_completed),
|
||||
Arc::clone(&allow_failover_errors),
|
||||
Arc::clone(&transit_failure),
|
||||
Arc::clone(&failed),
|
||||
stop.clone(),
|
||||
));
|
||||
|
||||
wait_for_count(&kv2_completed, &kv2_failure, 2, "two healthy KV2 decrypts", HEALTHY_PROGRESS_TIMEOUT).await;
|
||||
wait_for_count(
|
||||
&transit_completed,
|
||||
&transit_failure,
|
||||
2,
|
||||
"two healthy Transit decrypts",
|
||||
HEALTHY_PROGRESS_TIMEOUT,
|
||||
)
|
||||
.await;
|
||||
allow_failover_errors.store(true, Ordering::SeqCst);
|
||||
wait_for_count(&kv2_completed, 2, "two healthy KV2 decrypts").await;
|
||||
wait_for_count(&transit_completed, 2, "two healthy Transit decrypts").await;
|
||||
std::fs::write(&marker, b"ready").expect("publish failover readiness marker");
|
||||
|
||||
wait_for_file(&elected, "the replacement Vault leader").await;
|
||||
@@ -371,39 +326,18 @@ async fn exercise_failover(snapshotter: &Snapshotter) {
|
||||
|
||||
let kv2_after_election = kv2_completed.load(Ordering::SeqCst) + 2;
|
||||
let transit_after_election = transit_completed.load(Ordering::SeqCst) + 2;
|
||||
wait_for_count(
|
||||
&kv2_completed,
|
||||
&kv2_failure,
|
||||
kv2_after_election,
|
||||
"post-failover KV2 decrypts",
|
||||
POST_FAILOVER_PROGRESS_TIMEOUT,
|
||||
)
|
||||
.await;
|
||||
wait_for_count(
|
||||
&transit_completed,
|
||||
&transit_failure,
|
||||
transit_after_election,
|
||||
"post-failover Transit decrypts",
|
||||
POST_FAILOVER_PROGRESS_TIMEOUT,
|
||||
)
|
||||
.await;
|
||||
wait_for_count(&kv2_completed, kv2_after_election, "post-failover KV2 decrypts").await;
|
||||
wait_for_count(&transit_completed, transit_after_election, "post-failover Transit decrypts").await;
|
||||
|
||||
stop.cancel();
|
||||
kv2_worker.await.expect("KV2 decrypt worker must join");
|
||||
transit_worker.await.expect("Transit decrypt worker must join");
|
||||
assert!(
|
||||
kv2_failure.lock().expect("KV2 failure lock poisoned").is_none(),
|
||||
"no KV2 decrypt may fail or return different plaintext"
|
||||
);
|
||||
assert!(
|
||||
transit_failure.lock().expect("Transit failure lock poisoned").is_none(),
|
||||
"no Transit decrypt may fail or return different plaintext"
|
||||
);
|
||||
assert!(!failed.load(Ordering::SeqCst), "no decrypt may fail or return different plaintext");
|
||||
}
|
||||
|
||||
#[test]
|
||||
#[ignore = "requires a real three-node Vault Raft cluster; run scripts/test/vault_ha_kms_live.sh"]
|
||||
fn vault_raft_leader_failure_recovers_kv2_and_transit_decrypts() {
|
||||
fn vault_raft_leader_failure_preserves_kv2_and_transit_decrypts() {
|
||||
let recorder = DebuggingRecorder::new();
|
||||
let snapshotter = recorder.snapshotter();
|
||||
metrics::with_local_recorder(&recorder, || {
|
||||
@@ -415,6 +349,11 @@ fn vault_raft_leader_failure_recovers_kv2_and_transit_decrypts() {
|
||||
});
|
||||
let snapshot = snapshotter.snapshot().into_vec();
|
||||
|
||||
assert_eq!(
|
||||
counter_value(&snapshot, OPERATIONS_TOTAL, &[("outcome", "circuit_open")]),
|
||||
0,
|
||||
"a bounded leader election must not open the circuit"
|
||||
);
|
||||
assert_eq!(
|
||||
counter_value(&snapshot, OPERATIONS_TOTAL, &[("outcome", "budget_exhausted")]),
|
||||
0,
|
||||
|
||||
+4
-11
@@ -158,11 +158,10 @@ added by backlog#1153 infra-4.
|
||||
|
||||
## Coverage
|
||||
|
||||
Workspace line coverage is measured weekly. Pull requests that touch iam, kms,
|
||||
policy, or crypto also run a non-required, report-only comparison against
|
||||
`.config/coverage-baselines.toml`. During calibration, a regression is recorded
|
||||
in the job summary without failing the job; missing or malformed coverage
|
||||
evidence still fails closed (backlog#1153 infra-6).
|
||||
Line coverage is measured **weekly, not per-PR**, and is non-blocking: it
|
||||
exists for visibility and trend, never as a required check. Per-crate ratchets
|
||||
for the security-critical crates (iam / kms / policy / crypto) build on this
|
||||
baseline later (backlog#1153 infra-6, report-only first).
|
||||
|
||||
- **CI**: `.github/workflows/coverage.yml` runs every Sunday and on manual
|
||||
dispatch: `cargo llvm-cov nextest --workspace --exclude e2e_test` under the
|
||||
@@ -175,12 +174,6 @@ evidence still fails closed (backlog#1153 infra-6).
|
||||
plus the full suite). It prints the same per-crate table via
|
||||
`scripts/coverage_per_crate.py` and writes `target/llvm-cov/lcov.info` and
|
||||
`coverage.json`.
|
||||
- **Security-critical ratchet**: relevant pull requests compare iam / kms /
|
||||
policy / crypto line coverage with the versioned baseline. Drops greater than
|
||||
the configured one-percentage-point calibration threshold are marked
|
||||
`REGRESSION (report-only)`. The weekly summary runs the same comparison so
|
||||
calibration continues even when no relevant pull request is open. Baseline
|
||||
changes require a linked coverage run and a reviewed explanation.
|
||||
- **Trend comparison**: each run's job summary is the weekly per-crate
|
||||
snapshot — open two runs from the Actions history (workflow "coverage") and
|
||||
compare their tables. For line-level diffs, download the two runs'
|
||||
|
||||
@@ -1,264 +0,0 @@
|
||||
#!/usr/bin/env python3
|
||||
# Copyright 2024 RustFS Team
|
||||
#
|
||||
# Licensed under the Apache License, Version 2.0 (the "License");
|
||||
# you may not use this file except in compliance with the License.
|
||||
# You may obtain a copy of the License at
|
||||
#
|
||||
# http://www.apache.org/licenses/LICENSE-2.0
|
||||
#
|
||||
# Unless required by applicable law or agreed to in writing, software
|
||||
# distributed under the License is distributed on an "AS IS" BASIS,
|
||||
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
# See the License for the specific language governing permissions and
|
||||
# limitations under the License.
|
||||
|
||||
"""Compare security-critical crate line coverage with the report-only baseline."""
|
||||
|
||||
import argparse
|
||||
import json
|
||||
import math
|
||||
import os
|
||||
import sys
|
||||
import tempfile
|
||||
import tomllib
|
||||
from pathlib import Path
|
||||
|
||||
from coverage_per_crate import fmt_pct, load_coverage
|
||||
|
||||
|
||||
SECURITY_CRATES = ("crates/iam", "crates/kms", "crates/policy", "crates/crypto")
|
||||
|
||||
|
||||
def load_baselines(path: str) -> tuple[float, dict[str, tuple[int, int]]]:
|
||||
with open(path, "rb") as fh:
|
||||
config = tomllib.load(fh)
|
||||
|
||||
if config.get("phase") != "report-only":
|
||||
raise ValueError("coverage baseline phase must be report-only")
|
||||
|
||||
raw_allowed_drop = config["allowed_drop_percentage_points"]
|
||||
if isinstance(raw_allowed_drop, bool) or not isinstance(raw_allowed_drop, (int, float)):
|
||||
raise ValueError("allowed_drop_percentage_points must be a number")
|
||||
allowed_drop = float(raw_allowed_drop)
|
||||
if not math.isfinite(allowed_drop) or allowed_drop < 0:
|
||||
raise ValueError("allowed_drop_percentage_points must be finite and non-negative")
|
||||
|
||||
baselines: dict[str, tuple[int, int]] = {}
|
||||
for crate, values in config["crates"].items():
|
||||
covered = values["covered"]
|
||||
count = values["count"]
|
||||
if type(covered) is not int or type(count) is not int:
|
||||
raise ValueError(f"invalid baseline for {crate}: covered and count must be integers")
|
||||
if covered < 0 or count <= 0 or covered > count:
|
||||
raise ValueError(f"invalid baseline for {crate}: {covered}/{count}")
|
||||
baselines[crate] = (covered, count)
|
||||
missing = [crate for crate in SECURITY_CRATES if crate not in baselines]
|
||||
unexpected = sorted(set(baselines).difference(SECURITY_CRATES))
|
||||
if missing or unexpected:
|
||||
raise ValueError(f"coverage baseline crate set mismatch: missing={missing}, unexpected={unexpected}")
|
||||
return allowed_drop, baselines
|
||||
|
||||
|
||||
def compare(
|
||||
current: dict[str, list[int]],
|
||||
baselines: dict[str, tuple[int, int]],
|
||||
allowed_drop: float,
|
||||
) -> list[tuple[str, int, int, int, int, float, bool]]:
|
||||
rows = []
|
||||
for crate, (baseline_covered, baseline_count) in baselines.items():
|
||||
if crate not in current:
|
||||
raise ValueError(f"coverage report is missing {crate}")
|
||||
covered, count = current[crate]
|
||||
if type(covered) is not int or type(count) is not int:
|
||||
raise ValueError(f"invalid coverage for {crate}: covered and count must be integers")
|
||||
if covered < 0 or count <= 0 or covered > count:
|
||||
raise ValueError(f"invalid coverage for {crate}: {covered}/{count}")
|
||||
current_pct = 100.0 * covered / count
|
||||
baseline_pct = 100.0 * baseline_covered / baseline_count
|
||||
delta = current_pct - baseline_pct
|
||||
rows.append((crate, covered, count, baseline_covered, baseline_count, delta, delta < -allowed_drop))
|
||||
return rows
|
||||
|
||||
|
||||
def print_report(rows: list[tuple[str, int, int, int, int, float, bool]], allowed_drop: float) -> None:
|
||||
print("## Security-critical coverage ratchet (report-only)")
|
||||
print()
|
||||
print(f"Calibration threshold: a drop greater than {allowed_drop:.2f} percentage points is reported as a regression.")
|
||||
print()
|
||||
print("| Crate | Current | Baseline | Delta | Status |")
|
||||
print("|---|---:|---:|---:|---|")
|
||||
for crate, covered, count, baseline_covered, baseline_count, delta, regressed in rows:
|
||||
status = "REGRESSION (report-only)" if regressed else "OK"
|
||||
print(
|
||||
f"| `{crate}` | {fmt_pct(covered, count)} ({covered}/{count}) "
|
||||
f"| {fmt_pct(baseline_covered, baseline_count)} ({baseline_covered}/{baseline_count}) "
|
||||
f"| {delta:+.2f} pp | {status} |"
|
||||
)
|
||||
print()
|
||||
print("This calibration phase records regressions without failing the job; malformed or incomplete evidence still fails closed.")
|
||||
|
||||
|
||||
def self_test() -> None:
|
||||
with tempfile.TemporaryDirectory() as tmp:
|
||||
root = Path(tmp)
|
||||
coverage = root / "coverage.json"
|
||||
baseline = root / "baseline.toml"
|
||||
coverage_data = {
|
||||
"data": [
|
||||
{
|
||||
"files": [
|
||||
{
|
||||
"filename": str(root / "crates/iam/src/lib.rs"),
|
||||
"summary": {"lines": {"covered": 80, "count": 100}},
|
||||
},
|
||||
{
|
||||
"filename": str(root / "crates/kms/src/lib.rs"),
|
||||
"summary": {"lines": {"covered": 90, "count": 100}},
|
||||
},
|
||||
{
|
||||
"filename": str(root / "crates/policy/src/lib.rs"),
|
||||
"summary": {"lines": {"covered": 90, "count": 100}},
|
||||
},
|
||||
{
|
||||
"filename": str(root / "crates/crypto/src/lib.rs"),
|
||||
"summary": {"lines": {"covered": 90, "count": 100}},
|
||||
},
|
||||
],
|
||||
"totals": {"lines": {"covered": 350, "count": 400}},
|
||||
}
|
||||
]
|
||||
}
|
||||
coverage.write_text(json.dumps(coverage_data), encoding="utf-8")
|
||||
baseline_text = """phase = "report-only"
|
||||
allowed_drop_percentage_points = 1.0
|
||||
[crates."crates/iam"]
|
||||
covered = 90
|
||||
count = 100
|
||||
[crates."crates/kms"]
|
||||
covered = 85
|
||||
count = 100
|
||||
[crates."crates/policy"]
|
||||
covered = 90
|
||||
count = 100
|
||||
[crates."crates/crypto"]
|
||||
covered = 90
|
||||
count = 100
|
||||
"""
|
||||
baseline.write_text(baseline_text, encoding="utf-8")
|
||||
current, _ = load_coverage(str(coverage), str(root))
|
||||
allowed_drop, baselines = load_baselines(str(baseline))
|
||||
rows = compare(current, baselines, allowed_drop)
|
||||
assert [row[-1] for row in rows] == [True, False, False, False]
|
||||
try:
|
||||
compare({"crates/iam": current["crates/iam"]}, baselines, allowed_drop)
|
||||
except ValueError as error:
|
||||
assert str(error) == "coverage report is missing crates/kms"
|
||||
else:
|
||||
raise AssertionError("missing crate must fail closed")
|
||||
try:
|
||||
compare({**current, "crates/iam": [101, 100]}, baselines, allowed_drop)
|
||||
except ValueError as error:
|
||||
assert str(error) == "invalid coverage for crates/iam: 101/100"
|
||||
else:
|
||||
raise AssertionError("invalid coverage must fail closed")
|
||||
for invalid_threshold in ("true", '"1.0"', "nan", "inf", "-inf"):
|
||||
baseline.write_text(
|
||||
baseline_text.replace("allowed_drop_percentage_points = 1.0", f"allowed_drop_percentage_points = {invalid_threshold}"),
|
||||
encoding="utf-8",
|
||||
)
|
||||
try:
|
||||
load_baselines(str(baseline))
|
||||
except ValueError:
|
||||
pass
|
||||
else:
|
||||
raise AssertionError(f"non-finite threshold {invalid_threshold} must fail closed")
|
||||
for field, invalid_values in (
|
||||
("covered", ("true", '"90"', "90.0", "90.5")),
|
||||
("count", ("true", '"100"', "100.0", "100.5")),
|
||||
):
|
||||
for invalid_value in invalid_values:
|
||||
baseline.write_text(
|
||||
baseline_text.replace(f"{field} = {90 if field == 'covered' else 100}", f"{field} = {invalid_value}", 1),
|
||||
encoding="utf-8",
|
||||
)
|
||||
try:
|
||||
load_baselines(str(baseline))
|
||||
except ValueError:
|
||||
pass
|
||||
else:
|
||||
raise AssertionError(f"non-integer baseline {field} {invalid_value} must fail closed")
|
||||
for covered, count in (
|
||||
(True, 100),
|
||||
(80, True),
|
||||
(80.0, 100),
|
||||
(80, 100.0),
|
||||
(float("nan"), 100),
|
||||
(80, float("inf")),
|
||||
):
|
||||
try:
|
||||
compare({**current, "crates/iam": [covered, count]}, baselines, allowed_drop)
|
||||
except ValueError:
|
||||
pass
|
||||
else:
|
||||
raise AssertionError(f"invalid aggregate coverage {covered}/{count} must fail closed")
|
||||
lines = coverage_data["data"][0]["files"][0]["summary"]["lines"]
|
||||
for field, invalid_values in (
|
||||
("covered", (True, "80", 80.0, 80.5, float("nan"), float("inf"), float("-inf"))),
|
||||
("count", (True, "100", 100.0, 100.5, float("nan"), float("inf"), float("-inf"))),
|
||||
):
|
||||
original = lines[field]
|
||||
for invalid_value in invalid_values:
|
||||
lines[field] = invalid_value
|
||||
coverage.write_text(json.dumps(coverage_data), encoding="utf-8")
|
||||
try:
|
||||
load_coverage(str(coverage), str(root))
|
||||
except ValueError:
|
||||
pass
|
||||
else:
|
||||
raise AssertionError(f"invalid raw coverage {field} {invalid_value} must fail closed")
|
||||
lines[field] = original
|
||||
baseline.write_text(
|
||||
baseline_text.replace(
|
||||
'[crates."crates/crypto"]\ncovered = 90\ncount = 100\n',
|
||||
"",
|
||||
),
|
||||
encoding="utf-8",
|
||||
)
|
||||
try:
|
||||
load_baselines(str(baseline))
|
||||
except ValueError:
|
||||
pass
|
||||
else:
|
||||
raise AssertionError("missing security-crate baseline must fail closed")
|
||||
|
||||
|
||||
def main() -> int:
|
||||
parser = argparse.ArgumentParser(description=__doc__)
|
||||
parser.add_argument("coverage_json", nargs="?")
|
||||
parser.add_argument("--baseline", default=".config/coverage-baselines.toml")
|
||||
parser.add_argument("--repo-root", default=os.getcwd())
|
||||
parser.add_argument("--self-test", action="store_true")
|
||||
args = parser.parse_args()
|
||||
|
||||
if args.self_test:
|
||||
self_test()
|
||||
print("security coverage self-test passed")
|
||||
return 0
|
||||
if not args.coverage_json:
|
||||
parser.error("coverage_json is required unless --self-test is used")
|
||||
|
||||
try:
|
||||
current, _ = load_coverage(args.coverage_json, os.path.abspath(args.repo_root))
|
||||
allowed_drop, baselines = load_baselines(args.baseline)
|
||||
rows = compare(current, baselines, allowed_drop)
|
||||
except (OSError, ValueError, KeyError, IndexError, json.JSONDecodeError, tomllib.TOMLDecodeError) as error:
|
||||
print(f"error: {error}", file=sys.stderr)
|
||||
return 1
|
||||
|
||||
print_report(rows, allowed_drop)
|
||||
return 0
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
sys.exit(main())
|
||||
@@ -47,31 +47,6 @@ def fmt_pct(covered: int, count: int) -> str:
|
||||
return f"{100.0 * covered / count:.2f}%" if count else "—"
|
||||
|
||||
|
||||
def _line_counts(lines: dict[str, int], source: str) -> tuple[int, int]:
|
||||
covered = lines["covered"]
|
||||
count = lines["count"]
|
||||
if type(covered) is not int or type(count) is not int or covered < 0 or count < 0 or covered > count:
|
||||
raise ValueError(f"invalid line coverage for {source}: {covered}/{count}")
|
||||
return covered, count
|
||||
|
||||
|
||||
def load_coverage(path: str, root: str) -> tuple[dict[str, list[int]], dict[str, int]]:
|
||||
with open(path, encoding="utf-8") as fh:
|
||||
export = json.load(fh)
|
||||
|
||||
data = export["data"][0]
|
||||
files = data["files"]
|
||||
total_covered, total_count = _line_counts(data["totals"]["lines"], "totals")
|
||||
|
||||
crates: dict[str, list[int]] = {}
|
||||
for f in files:
|
||||
covered, count = _line_counts(f["summary"]["lines"], f["filename"])
|
||||
acc = crates.setdefault(crate_label(f["filename"], root), [0, 0])
|
||||
acc[0] += covered
|
||||
acc[1] += count
|
||||
return crates, {"covered": total_covered, "count": total_count}
|
||||
|
||||
|
||||
def main() -> int:
|
||||
if len(sys.argv) < 2 or len(sys.argv) > 3:
|
||||
print(__doc__.strip(), file=sys.stderr)
|
||||
@@ -79,12 +54,24 @@ def main() -> int:
|
||||
path = sys.argv[1]
|
||||
root = os.path.abspath(sys.argv[2] if len(sys.argv) == 3 else os.getcwd())
|
||||
|
||||
with open(path, encoding="utf-8") as fh:
|
||||
export = json.load(fh)
|
||||
|
||||
try:
|
||||
crates, totals = load_coverage(path, root)
|
||||
except (KeyError, IndexError, ValueError) as exc:
|
||||
data = export["data"][0]
|
||||
files = data["files"]
|
||||
totals = data["totals"]["lines"]
|
||||
except (KeyError, IndexError) as exc:
|
||||
print(f"error: unexpected llvm-cov JSON shape ({exc})", file=sys.stderr)
|
||||
return 1
|
||||
|
||||
crates: dict[str, list[int]] = {}
|
||||
for f in files:
|
||||
lines = f["summary"]["lines"]
|
||||
acc = crates.setdefault(crate_label(f["filename"], root), [0, 0])
|
||||
acc[0] += lines["covered"]
|
||||
acc[1] += lines["count"]
|
||||
|
||||
rows = sorted(
|
||||
crates.items(),
|
||||
key=lambda kv: (100.0 * kv[1][0] / kv[1][1]) if kv[1][1] else 101.0,
|
||||
|
||||
@@ -241,7 +241,7 @@ env \
|
||||
RUSTFS_TEST_VAULT_FAILOVER_MARKER="$MARKER" \
|
||||
RUSTFS_TEST_VAULT_OLD_LEADER="$OLD_LEADER" \
|
||||
cargo test -p rustfs-kms --test vault_ha_failover_live \
|
||||
vault_raft_leader_failure_recovers_kv2_and_transit_decrypts -- \
|
||||
vault_raft_leader_failure_preserves_kv2_and_transit_decrypts -- \
|
||||
--ignored --nocapture --test-threads=1 &
|
||||
TEST_PID=$!
|
||||
|
||||
|
||||
Reference in New Issue
Block a user