Compare commits

..

22 Commits

Author SHA1 Message Date
overtrue acc37b49c8 test(ecstore): cover final sweep cancel fence 2026-08-23 03:21:15 +08:00
overtrue 2379bfb5a3 fix(ecstore): fence final decommission sweep 2026-08-23 03:15:33 +08:00
overtrue 4527b003cc fix(ecstore): remove duplicate decommission error helper 2026-08-23 03:00:41 +08:00
overtrue c83aa7f974 Merge remote-tracking branch 'origin/main' into overtrue/fix-1911-ilm-metadata-decommission
# Conflicts:
#	crates/ecstore/src/core/pools.rs
2026-08-23 02:58:32 +08:00
overtrue b6f135f0fe fix(ecstore): drop removed decommission test import 2026-08-22 23:28:14 +08:00
overtrue b07f0a92a8 merge: sync main into ILM decommission fix 2026-08-22 22:55:14 +08:00
overtrue 34a52a1e7d fix(ecstore): restore decommission test imports 2026-08-22 19:13:21 +08:00
overtrue fe2d166cd3 fix(ecstore): distinguish active ILM target cleanup 2026-08-22 18:43:05 +08:00
overtrue c687f95260 fix(ecstore): preserve active ILM source journals 2026-08-22 18:43:05 +08:00
overtrue fe66aa509e test(ecstore): isolate durable ILM scenario stack 2026-08-22 18:43:05 +08:00
overtrue 918686bd4e test(ecstore): compile multi-source ILM recovery 2026-08-22 18:43:05 +08:00
overtrue a9cbf82943 test(ecstore): serialize multi-source ILM recovery 2026-08-22 18:43:05 +08:00
overtrue ee7a82ec1f test(ecstore): cover durable ILM recovery boundaries 2026-08-22 18:43:05 +08:00
overtrue 3229ce05e3 fix(ecstore): repair durable ILM receipt recovery 2026-08-22 18:43:04 +08:00
overtrue 7428f3138c fix(ecstore): harden durable ILM cursor receipts 2026-08-22 18:43:04 +08:00
overtrue cf0a78ce11 fix(ecstore): avoid terminal receipt shadowing 2026-08-22 18:43:04 +08:00
overtrue 7bda7f4825 fix(ecstore): re-export durable ILM checkpoint 2026-08-22 18:43:04 +08:00
overtrue 20a8c443af fix(ecstore): anchor decommission ILM receipts 2026-08-22 18:43:04 +08:00
overtrue 948ee1fc52 fix(ecstore): close ILM receipt recovery gaps 2026-08-22 18:43:04 +08:00
overtrue ecd779b6f3 fix(ecstore): track ILM recovery across decommission 2026-08-22 18:43:04 +08:00
overtrue 619fae2512 fix(ecstore): verify ILM metadata before decommission 2026-08-22 18:43:04 +08:00
overtrue 4f9d22bda7 fix(ecstore): migrate ILM metadata during decommission 2026-08-22 18:43:04 +08:00
17 changed files with 4064 additions and 438 deletions
-20
View File
@@ -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
-1
View File
@@ -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
+10
View File
@@ -78,6 +78,12 @@ test-group = 'embedded-test-ports'
filter = 'package(rustfs-ecstore) & test(manual_transition_page_checkpoint_persists_durable_job_progress)'
test-group = 'ecstore-serial-flaky'
# The durable ILM decommission regressions build isolated multi-pool stores and
# deliberately take source or target disks offline while checking fencing.
[[profile.default.overrides]]
filter = 'package(rustfs-ecstore) & (test(decommission_migrates_and_verifies_registered_durable_ilm_records) | test(decommission_durable_ilm_target_read_error_is_not_masked_by_peer_success) | test(decommission_durable_ilm_terminal_receipt_recovers_failed_source_cleanup) | test(decommission_durable_ilm_receipt_pagination_fails_closed_on_second_page) | test(decommission_durable_ilm_recovery_keeps_multiple_active_sources))'
test-group = 'ecstore-serial-flaky'
# Serialize the bucket-incarnation / lifecycle-fence tests. They drive
# init_bucket_metadata_sys and bucket_metadata_sys_of, i.e. process-global
# OnceLock state that serial_test's #[serial] cannot protect across nextest's
@@ -190,6 +196,10 @@ test-group = 'embedded-test-ports'
filter = 'package(rustfs-ecstore) & test(manual_transition_page_checkpoint_persists_durable_job_progress)'
test-group = 'ecstore-serial-flaky'
[[profile.ci.overrides]]
filter = 'package(rustfs-ecstore) & (test(decommission_migrates_and_verifies_registered_durable_ilm_records) | test(decommission_durable_ilm_target_read_error_is_not_masked_by_peer_success) | test(decommission_durable_ilm_terminal_receipt_recovers_failed_source_cleanup) | test(decommission_durable_ilm_receipt_pagination_fails_closed_on_second_page) | test(decommission_durable_ilm_recovery_keeps_multiple_active_sources))'
test-group = 'ecstore-serial-flaky'
# Serialize the bucket-incarnation / lifecycle-fence tests under the ci profile
# too (see the matching default-profile override near the top). No retries.
[[profile.ci.overrides]]
+12 -28
View File
@@ -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()
@@ -9136,6 +9136,7 @@ mod tests {
assert_eq!(loaded.report.scanned, 37);
assert_eq!(loaded.report.eligible, 11);
assert_eq!(loaded.report.enqueued, 5);
assert_eq!(loaded.cursor_revision, Some(1));
assert!(loaded.lease_expires_at_unix_nanos > 0);
let token = loaded
.report
@@ -9154,6 +9155,30 @@ mod tests {
assert_eq!(admission.lease_id, loaded.lease_id);
assert_eq!(admission.lease_expires_at_unix_nanos, loaded.lease_expires_at_unix_nanos);
let mut same_marker_report = report.clone();
same_marker_report.scanned += 1;
persist_manual_transition_page_checkpoint(
&checkpoint_options,
&same_marker_report,
Some("logs/page-end".to_string()),
Some("opaque-next-version".to_string()),
)
.await
.expect("same-marker version checkpoint should persist through the durable progress sink");
let same_marker_checkpointed = load_manual_transition_job_record(ecstore.clone(), job_id)
.await
.expect("same-marker version checkpoint should reload");
assert_eq!(same_marker_checkpointed.cursor_revision, Some(2));
let (_, version_marker) = decode_manual_transition_continuation_token(
same_marker_checkpointed
.report
.continuation_token
.as_deref()
.expect("same-marker version checkpoint should persist a cursor"),
)
.expect("same-marker version cursor should decode");
assert_eq!(version_marker.as_deref(), Some("opaque-next-version"));
create_test_bucket(&ecstore, &bucket).await;
let lifecycle_xml = format!(
r#"<?xml version="1.0" encoding="UTF-8"?>
@@ -9210,6 +9235,7 @@ mod tests {
assert_eq!(checkpointed.report.scanned, 1000);
assert_eq!(checkpointed.report.eligible, 1000);
assert_eq!(checkpointed.report.dry_run_eligible, 1000);
assert_eq!(checkpointed.cursor_revision, Some(3));
let token = checkpointed
.report
.continuation_token
File diff suppressed because it is too large Load Diff
@@ -24,6 +24,10 @@ use crate::bucket::lifecycle::bucket_lifecycle_ops::{
ManualTransitionQueueSnapshot, ManualTransitionRunOptions, ManualTransitionRunReport,
};
use crate::bucket::lifecycle::config_boundary;
use crate::bucket::lifecycle::durable_namespace::{
MANUAL_TRANSITION_JOB_NAMESPACE, MANUAL_TRANSITION_SCOPE_NAMESPACE, MANUAL_TRANSITION_TASK_NAMESPACE,
MANUAL_TRANSITION_WORKER_RESULT_NAMESPACE,
};
use crate::disk::RUSTFS_META_BUCKET;
use crate::error::{Error, Result as EcstoreResult};
use crate::object_api::ObjectOptions;
@@ -34,10 +38,10 @@ use crate::store::ECStore;
pub const MANUAL_TRANSITION_JOB_SCHEMA: &str = "rustfs-manual-transition-job-v1";
pub const MANUAL_TRANSITION_TASK_SCHEMA: &str = "rustfs-manual-transition-task-v1";
pub const MANUAL_TRANSITION_WORKER_RESULT_SCHEMA: &str = "rustfs-manual-transition-worker-result-v1";
pub const MANUAL_TRANSITION_JOB_RECORD_PREFIX: &str = "ilm/manual-transition/jobs";
pub const MANUAL_TRANSITION_SCOPE_RECORD_PREFIX: &str = "ilm/manual-transition/scopes";
pub const MANUAL_TRANSITION_TASK_PREFIX: &str = "ilm/manual-transition/tasks";
pub const MANUAL_TRANSITION_WORKER_RESULT_PREFIX: &str = "ilm/manual-transition/results";
pub const MANUAL_TRANSITION_JOB_RECORD_PREFIX: &str = MANUAL_TRANSITION_JOB_NAMESPACE.prefix;
pub const MANUAL_TRANSITION_SCOPE_RECORD_PREFIX: &str = MANUAL_TRANSITION_SCOPE_NAMESPACE.prefix;
pub const MANUAL_TRANSITION_TASK_PREFIX: &str = MANUAL_TRANSITION_TASK_NAMESPACE.prefix;
pub const MANUAL_TRANSITION_WORKER_RESULT_PREFIX: &str = MANUAL_TRANSITION_WORKER_RESULT_NAMESPACE.prefix;
pub const MAX_MANUAL_TRANSITION_JOB_RECORD_SIZE: usize = 64 * 1024;
pub const MAX_MANUAL_TRANSITION_TASK_RECORD_SIZE: usize = 16 * 1024;
pub const MAX_MANUAL_TRANSITION_WORKER_RESULT_RECORD_SIZE: usize = 8 * 1024;
@@ -195,6 +199,8 @@ pub struct ManualTransitionJobRecord {
pub updated_at_unix_nanos: i128,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub completed_at_unix_nanos: Option<i128>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub cursor_revision: Option<u64>,
pub report: ManualTransitionRunReport,
pub queue_snapshot: ManualTransitionQueueSnapshot,
#[serde(default, skip_serializing_if = "Option::is_none")]
@@ -224,6 +230,7 @@ impl ManualTransitionJobRecord {
created_at_unix_nanos: now,
updated_at_unix_nanos: now,
completed_at_unix_nanos: None,
cursor_revision: Some(0),
report: ManualTransitionRunReport {
bucket: bucket.to_string(),
prefix: options.prefix.clone(),
@@ -238,7 +245,7 @@ impl ManualTransitionJobRecord {
pub fn complete(&mut self, report: ManualTransitionRunReport, queue_snapshot: ManualTransitionQueueSnapshot) {
self.scan_completed = true;
self.report.merge_scan_report_preserving_worker(&report);
self.merge_scan_report(&report);
self.queue_snapshot = queue_snapshot;
self.error = None;
self.mark_terminal_if_worker_drained();
@@ -316,7 +323,7 @@ impl ManualTransitionJobRecord {
}
}
self.queue_snapshot = queue_snapshot;
self.updated_at_unix_nanos = OffsetDateTime::now_utc().unix_timestamp_nanos();
self.advance_updated_at();
self.mark_terminal_if_worker_drained();
}
@@ -359,14 +366,14 @@ impl ManualTransitionJobRecord {
self.report.tier_failure = scan_tier_failure.saturating_add(transition_failed);
self.report.tier_failure_by_reason = scan_tier_failure_by_reason;
self.queue_snapshot = queue_snapshot;
self.updated_at_unix_nanos = OffsetDateTime::now_utc().unix_timestamp_nanos();
self.advance_updated_at();
self.mark_terminal_if_worker_drained();
true
}
pub fn mark_cancel_requested(&mut self) {
self.cancel_requested = true;
self.updated_at_unix_nanos = OffsetDateTime::now_utc().unix_timestamp_nanos();
self.advance_updated_at();
}
pub fn claim_recovery_lease(&mut self, owner_id: impl Into<String>, queue_snapshot: ManualTransitionQueueSnapshot) {
@@ -380,7 +387,7 @@ impl ManualTransitionJobRecord {
pub fn abandon_recovery_lease(&mut self, lease_id: Uuid) {
if self.state == ManualTransitionJobState::Running && self.lease_id == lease_id {
self.lease_expires_at_unix_nanos = 0;
self.updated_at_unix_nanos = OffsetDateTime::now_utc().unix_timestamp_nanos();
self.advance_updated_at();
}
}
@@ -398,7 +405,7 @@ impl ManualTransitionJobRecord {
pub fn renew_lease(&mut self, queue_snapshot: ManualTransitionQueueSnapshot) {
let now = OffsetDateTime::now_utc().unix_timestamp_nanos();
self.updated_at_unix_nanos = now;
self.updated_at_unix_nanos = self.updated_at_unix_nanos.saturating_add(1).max(now);
self.lease_expires_at_unix_nanos = manual_transition_job_lease_expires_at(now);
self.queue_snapshot = queue_snapshot;
}
@@ -438,11 +445,18 @@ impl ManualTransitionJobRecord {
pub fn update_running_progress(&mut self, report: ManualTransitionRunReport, queue_snapshot: ManualTransitionQueueSnapshot) {
if self.state == ManualTransitionJobState::Running {
self.report.merge_scan_report_preserving_worker(&report);
self.merge_scan_report(&report);
self.renew_lease(queue_snapshot);
}
}
fn merge_scan_report(&mut self, report: &ManualTransitionRunReport) {
if self.report.continuation_token != report.continuation_token {
self.cursor_revision = Some(self.cursor_revision.unwrap_or(0).saturating_add(1));
}
self.report.merge_scan_report_preserving_worker(report);
}
pub fn mark_unknown_if_unowned(&mut self) {
if self.state == ManualTransitionJobState::Running {
self.state = ManualTransitionJobState::Unknown;
@@ -463,9 +477,13 @@ impl ManualTransitionJobRecord {
}
fn mark_updated_terminal(&mut self) {
self.advance_updated_at();
self.completed_at_unix_nanos = Some(self.updated_at_unix_nanos);
}
fn advance_updated_at(&mut self) {
let now = OffsetDateTime::now_utc().unix_timestamp_nanos();
self.updated_at_unix_nanos = now;
self.completed_at_unix_nanos = Some(now);
self.updated_at_unix_nanos = self.updated_at_unix_nanos.saturating_add(1).max(now);
}
fn mark_terminal_if_worker_drained(&mut self) {
@@ -1109,7 +1127,8 @@ pub fn manual_transition_scope_record_object_name(scope_key: &str) -> Result<Str
pub async fn save_manual_transition_job_record(api: Arc<ECStore>, job: &ManualTransitionJobRecord) -> EcstoreResult<()> {
let object = manual_transition_job_record_object_name(job.job_id).map_err(manual_transition_job_store_error)?;
let data = job.encode().map_err(manual_transition_job_store_error)?;
config_boundary::save_config(api, &object, data).await
config_boundary::save_config(api.clone(), &object, data.clone()).await?;
api.record_durable_ilm_decommission_progress(&object, &data).await
}
pub async fn load_manual_transition_job_record(api: Arc<ECStore>, job_id: Uuid) -> EcstoreResult<ManualTransitionJobRecord> {
@@ -1142,9 +1161,9 @@ pub async fn save_manual_transition_job_record_if_current(
let object = manual_transition_job_record_object_name(job.job_id).map_err(manual_transition_job_store_error)?;
let data = job.encode().map_err(manual_transition_job_store_error)?;
config_boundary::save_config_with_opts_quiet(
api,
api.clone(),
&object,
data,
data.clone(),
&ObjectOptions {
max_parity: true,
http_preconditions: Some(HTTPPreconditions {
@@ -1154,7 +1173,8 @@ pub async fn save_manual_transition_job_record_if_current(
..Default::default()
},
)
.await
.await?;
api.record_durable_ilm_decommission_progress(&object, &data).await
}
/// Applies a job-record mutation with optimistic concurrency control.
@@ -1592,9 +1612,9 @@ pub async fn save_manual_transition_scope_admission_if_absent(
let object = manual_transition_scope_record_object_name(&admission.scope_key).map_err(manual_transition_job_store_error)?;
let data = serde_json::to_vec(admission).map_err(Error::other)?;
config_boundary::save_config_with_opts(
api,
api.clone(),
&object,
data,
data.clone(),
&ObjectOptions {
max_parity: true,
http_preconditions: Some(HTTPPreconditions {
@@ -1604,7 +1624,8 @@ pub async fn save_manual_transition_scope_admission_if_absent(
..Default::default()
},
)
.await
.await?;
api.record_durable_ilm_decommission_progress(&object, &data).await
}
pub async fn load_manual_transition_scope_admission(
@@ -1642,9 +1663,9 @@ pub async fn save_manual_transition_scope_admission_if_current(
let object = manual_transition_scope_record_object_name(&admission.scope_key).map_err(manual_transition_job_store_error)?;
let data = serde_json::to_vec(admission).map_err(Error::other)?;
match config_boundary::save_config_with_opts(
api,
api.clone(),
&object,
data,
data.clone(),
&ObjectOptions {
max_parity: true,
http_preconditions: Some(HTTPPreconditions {
@@ -1660,7 +1681,8 @@ pub async fn save_manual_transition_scope_admission_if_current(
Err(Error::PreconditionFailed)
}
result => result,
}
}?;
api.record_durable_ilm_decommission_progress(&object, &data).await
}
pub async fn claim_manual_transition_scope_admission(
@@ -1953,13 +1975,15 @@ pub async fn delete_manual_transition_scope_admission_if_current(
job_id: Uuid,
lease_id: Uuid,
) -> EcstoreResult<bool> {
let etag = match load_manual_transition_scope_admission_with_etag(api.clone(), scope_key).await {
Ok((admission, etag)) if admission.job_id == job_id && admission.lease_id == lease_id => etag,
let (admission, etag) = match load_manual_transition_scope_admission_with_etag(api.clone(), scope_key).await {
Ok((admission, etag)) if admission.job_id == job_id && admission.lease_id == lease_id => (admission, etag),
Ok(_) => return Ok(false),
Err(Error::ConfigNotFound) => return Ok(true),
Err(err) => return Err(err),
};
let object = manual_transition_scope_record_object_name(scope_key).map_err(manual_transition_job_store_error)?;
let data = serde_json::to_vec(&admission).map_err(Error::other)?;
api.record_durable_ilm_decommission_terminal(&object, &data).await?;
match config_boundary::delete_config_if_match(api, &object, &etag).await {
Ok(()) | Err(Error::ConfigNotFound) => Ok(true),
Err(Error::PreconditionFailed) => Ok(false),
@@ -2577,6 +2601,25 @@ mod tests {
assert!(decoded.report.tier_failure_by_reason.is_empty());
}
#[test]
fn manual_transition_job_record_decodes_legacy_cursor_without_revision() {
let options = ManualTransitionRunOptions::default();
let record = ManualTransitionJobRecord::new(Uuid::new_v4(), "bucket", &options, TEST_OWNER);
let encoded = record.encode().expect("job record should encode");
let mut value: serde_json::Value = serde_json::from_slice(&encoded).expect("encoded job should be json");
value["job"]
.as_object_mut()
.expect("job should be object")
.remove("cursor_revision");
let record_bytes = serde_json::to_vec(&value["job"]).expect("legacy job should encode");
value["content_sha256"] = serde_json::Value::String(hex_sha256(&record_bytes, ToOwned::to_owned));
let legacy = serde_json::to_vec(&value).expect("legacy envelope should encode");
let decoded = ManualTransitionJobRecord::decode(record.job_id, &legacy).expect("legacy job should decode");
assert_eq!(decoded.cursor_revision, None);
}
#[test]
fn manual_transition_job_record_rejects_unknown_report_fields() {
let options = ManualTransitionRunOptions::default();
@@ -16,6 +16,7 @@ pub mod bucket_lifecycle_audit;
pub mod bucket_lifecycle_ops;
mod config_boundary;
pub mod core;
mod durable_namespace;
pub mod evaluator;
pub mod manual_transition_job;
mod metadata_boundary;
@@ -31,3 +32,8 @@ pub mod tier_free_version_recovery;
pub mod tier_last_day_stats;
pub mod tier_sweeper;
pub mod transition_transaction;
pub(crate) use durable_namespace::{
DurableIlmRecordCheckpoint, ILM_META_PREFIX, ValidatedDurableIlmRecord, classify_durable_ilm_record,
validate_durable_ilm_record,
};
@@ -20,6 +20,7 @@ use tokio_util::sync::CancellationToken;
use tracing::{debug, warn};
use crate::bucket::lifecycle::config_boundary;
use crate::bucket::lifecycle::durable_namespace::TIER_DELETE_JOURNAL_NAMESPACE;
use crate::bucket::lifecycle::runtime_boundary;
use crate::bucket::lifecycle::tier_sweeper::{
Jentry, TierDeleteJournalState, TierDeleteSourceIdentity,
@@ -49,7 +50,7 @@ const TIER_DELETE_JOURNAL_VERSION: u8 = 2;
const TIER_DELETE_JOURNAL_EXACT_VERSION: u8 = 3;
const TIER_DELETE_JOURNAL_STATE_VERSION: u8 = 4;
const TIER_DELETE_JOURNAL_TRANSACTION_VERSION: u8 = 5;
pub(crate) const TIER_DELETE_JOURNAL_PREFIX: &str = "ilm/tier-delete-journal/";
pub(crate) const TIER_DELETE_JOURNAL_PREFIX: &str = TIER_DELETE_JOURNAL_NAMESPACE.prefix;
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[serde(deny_unknown_fields)]
@@ -432,6 +433,21 @@ async fn process_committed_tier_delete_journal_entry(api: Arc<ECStore>, je: &Jen
)
.await?;
}
let path = tier_delete_journal_object_name(je);
let data = encode_tier_delete_journal_entry(je).map_err(std::io::Error::other)?;
let target_pool_indices = api
.record_durable_ilm_decommission_terminal_target_pools(&path, &data)
.await
.map_err(std::io::Error::other)?;
if let Some(target_pool_indices) = target_pool_indices {
for target_pool_idx in target_pool_indices {
match config_boundary::delete_config(api.pools[target_pool_idx].clone(), &path).await {
Ok(()) | Err(Error::ConfigNotFound) => {}
Err(err) => return Err(std::io::Error::other(err)),
}
}
return Ok(());
}
remove_tier_delete_journal_entry(api, je).await
}
@@ -21,6 +21,7 @@ use tracing::{debug, warn};
use uuid::Uuid;
use crate::bucket::lifecycle::config_boundary;
use crate::bucket::lifecycle::durable_namespace::TRANSITION_TRANSACTION_NAMESPACE;
use crate::bucket::lifecycle::lifecycle::TRANSITION_COMPLETE;
use crate::bucket::lifecycle::tier_sweeper::{
delete_confirmed_transition_candidate_exact_with_lease_idempotent,
@@ -42,7 +43,7 @@ const TRANSITION_TRANSACTION_RECOVERY_INTERVAL: Duration = Duration::from_secs(6
const TRANSITION_TRANSACTION_RECOVERY_TIMEOUT: Duration = Duration::from_secs(300);
pub const TRANSITION_TRANSACTION_SCHEMA: &str = "rustfs-transition-transaction-v1";
pub const TRANSITION_TRANSACTION_PREFIX: &str = "ilm/transition-transactions";
pub const TRANSITION_TRANSACTION_RECORD_PREFIX: &str = "ilm/transition-transactions/records";
pub const TRANSITION_TRANSACTION_RECORD_PREFIX: &str = TRANSITION_TRANSACTION_NAMESPACE.prefix;
pub const MAX_TRANSITION_TRANSACTION_SIZE: usize = 64 * 1024;
pub type Result<T> = std::result::Result<T, TransitionTransactionError>;
@@ -584,7 +585,8 @@ pub(crate) async fn save_transition_transaction_record(
let object =
transition_transaction_record_object_name(transaction.transaction_id).map_err(transition_transaction_store_error)?;
let data = transaction.encode().map_err(transition_transaction_store_error)?;
config_boundary::save_config(api, &object, data).await
config_boundary::save_config(api.clone(), &object, data.clone()).await?;
api.record_durable_ilm_decommission_progress(&object, &data).await
}
pub(crate) async fn load_transition_transaction_record(
@@ -596,8 +598,14 @@ pub(crate) async fn load_transition_transaction_record(
TransitionTransaction::decode(transaction_id, &data).map_err(transition_transaction_store_error)
}
pub(crate) async fn delete_transition_transaction_record(api: Arc<ECStore>, transaction_id: Uuid) -> EcstoreResult<()> {
let object = transition_transaction_record_object_name(transaction_id).map_err(transition_transaction_store_error)?;
pub(crate) async fn delete_transition_transaction_record(
api: Arc<ECStore>,
transaction: &TransitionTransaction,
) -> EcstoreResult<()> {
let object =
transition_transaction_record_object_name(transaction.transaction_id).map_err(transition_transaction_store_error)?;
let data = transaction.encode().map_err(transition_transaction_store_error)?;
api.record_durable_ilm_decommission_terminal(&object, &data).await?;
match config_boundary::delete_config(api, &object).await {
Ok(()) | Err(Error::ConfigNotFound) => Ok(()),
Err(err) => Err(err),
@@ -813,7 +821,7 @@ pub async fn finalize_missing_transition_transaction_for_operator(
if probe != TransitionOperatorProbe::Missing {
return Err(TransitionOperatorError::CandidateNotMissing(probe));
}
delete_transition_transaction_record(api, transaction_id)
delete_transition_transaction_record(api, &transaction)
.await
.map_err(TransitionOperatorError::Store)
}
@@ -849,22 +857,22 @@ pub async fn process_transition_transaction_record(
match transaction.state {
TransitionTransactionState::Uploaded => {
delete_transition_remote_candidate(api.clone(), transaction).await?;
delete_transition_transaction_record(api, transaction.transaction_id).await?;
delete_transition_transaction_record(api, transaction).await?;
Ok(TransitionTransactionRecoveryOutcome::RemoteCandidateDeleted)
}
TransitionTransactionState::CleanupPending => match local_commit_matches_transaction(api.clone(), transaction).await {
Ok(true) => {
delete_transition_transaction_record(api, transaction.transaction_id).await?;
delete_transition_transaction_record(api, transaction).await?;
Ok(TransitionTransactionRecoveryOutcome::RecordDeleted)
}
Ok(false) => {
delete_transition_remote_candidate(api.clone(), transaction).await?;
delete_transition_transaction_record(api, transaction.transaction_id).await?;
delete_transition_transaction_record(api, transaction).await?;
Ok(TransitionTransactionRecoveryOutcome::RemoteCandidateDeleted)
}
Err(err) if transition_source_is_missing(&err) => {
delete_transition_remote_candidate(api.clone(), transaction).await?;
delete_transition_transaction_record(api, transaction.transaction_id).await?;
delete_transition_transaction_record(api, transaction).await?;
Ok(TransitionTransactionRecoveryOutcome::RemoteCandidateDeleted)
}
Err(err) => Err(err),
@@ -872,7 +880,7 @@ pub async fn process_transition_transaction_record(
TransitionTransactionState::LocalCommitStarted => {
match local_commit_matches_transaction(api.clone(), transaction).await {
Ok(true) => {
delete_transition_transaction_record(api, transaction.transaction_id).await?;
delete_transition_transaction_record(api, transaction).await?;
Ok(TransitionTransactionRecoveryOutcome::RecordDeleted)
}
Ok(false) => Ok(TransitionTransactionRecoveryOutcome::Retained),
@@ -881,7 +889,7 @@ pub async fn process_transition_transaction_record(
}
}
TransitionTransactionState::AbortedNoRemote | TransitionTransactionState::Committed => {
delete_transition_transaction_record(api, transaction.transaction_id).await?;
delete_transition_transaction_record(api, transaction).await?;
Ok(TransitionTransactionRecoveryOutcome::RecordDeleted)
}
TransitionTransactionState::UploadOutcomeUnknown => recover_unknown_upload_outcome(api, transaction).await,
@@ -907,7 +915,7 @@ async fn recover_unknown_upload_outcome(
.map_err(Error::other)?
{
TransitionCandidateProbe::Missing => {
delete_transition_transaction_record(api, transaction.transaction_id).await?;
delete_transition_transaction_record(api, transaction).await?;
Ok(TransitionTransactionRecoveryOutcome::RecordDeleted)
}
TransitionCandidateProbe::UnversionedPresent => {
@@ -925,7 +933,7 @@ async fn recover_unknown_upload_outcome(
)
.await
.map_err(Error::other)?;
delete_transition_transaction_record(api, transaction.transaction_id).await?;
delete_transition_transaction_record(api, transaction).await?;
Ok(TransitionTransactionRecoveryOutcome::RemoteCandidateDeleted)
}
TransitionCandidateProbe::VersionedPresent(version_id) => {
@@ -958,7 +966,7 @@ async fn cleanup_recovered_unknown_upload_candidate(
.map_err(transition_transaction_store_error)?;
save_transition_transaction_record(api.clone(), &cleanup).await?;
delete_transition_remote_candidate(api.clone(), &cleanup).await?;
delete_transition_transaction_record(api, cleanup.transaction_id).await?;
delete_transition_transaction_record(api, &cleanup).await?;
Ok(TransitionTransactionRecoveryOutcome::RemoteCandidateDeleted)
}
+19
View File
@@ -406,6 +406,25 @@ where
Ok(data)
}
pub(crate) async fn read_config_limited_preserve_empty<S>(api: Arc<S>, file: &str, max_bytes: usize) -> Result<Vec<u8>>
where
S: EcstoreObjectIO,
{
let (data, _obj) = read_config_limited_preserve_empty_with_metadata(api, file, max_bytes).await?;
Ok(data)
}
pub(crate) async fn read_config_limited_preserve_empty_with_metadata<S>(
api: Arc<S>,
file: &str,
max_bytes: usize,
) -> Result<(Vec<u8>, ObjectInfo)>
where
S: EcstoreObjectIO,
{
read_config_with_metadata_inner(api, file, &ObjectOptions::default(), true, Some(max_bytes)).await
}
/// Read an existing config object without treating an empty payload as absent.
/// Callers that validate their own payload format need to distinguish corruption
/// from `ConfigNotFound`.
File diff suppressed because it is too large Load Diff
+7 -2
View File
@@ -37,7 +37,7 @@ use crate::bucket::lifecycle::{
transition_transaction::{
TransitionRemoteVersion, TransitionSourceIdentity, TransitionSourceVersionMode, TransitionTransaction,
TransitionTransactionInit, TransitionTransactionState, delete_transition_transaction_record,
save_transition_transaction_record,
load_transition_transaction_record, save_transition_transaction_record,
},
};
use crate::bucket::quota::reservation;
@@ -4261,7 +4261,12 @@ fn record_transition_uploaded_save_attempt(transaction: &TransitionTransaction,
async fn delete_transition_transaction_if_available(api: Option<&Arc<ECStore>>, transaction_id: Uuid) -> Result<()> {
if let Some(api) = api {
return delete_transition_transaction_record(api.clone(), transaction_id).await;
let transaction = match load_transition_transaction_record(api.clone(), transaction_id).await {
Ok(transaction) => transaction,
Err(Error::ConfigNotFound) => return Ok(()),
Err(err) => return Err(err),
};
return delete_transition_transaction_record(api.clone(), &transaction).await;
}
Ok(())
}
File diff suppressed because it is too large Load Diff
+4 -11
View File
@@ -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'
-264
View File
@@ -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())
+14 -27
View File
@@ -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,