mirror of
https://github.com/rustfs/rustfs.git
synced 2026-09-08 04:58:12 +00:00
Compare commits
20 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| a8bb53218d | |||
| 5355d9f8f8 | |||
| 50d7a049ee | |||
| 3149c87cf2 | |||
| b7fa6a4615 | |||
| 6fe83f87a4 | |||
| df6981d88e | |||
| 590adad5ae | |||
| be5d14c985 | |||
| 703086d71d | |||
| a159f312f0 | |||
| 2f6f095298 | |||
| efd8ef005f | |||
| 60fa33773a | |||
| ee752b0b03 | |||
| 99c1f4418b | |||
| 7ce0ac72cf | |||
| 944e26d432 | |||
| f0b0a99260 | |||
| 33ddc10ffd |
@@ -1,5 +1,5 @@
|
||||
{
|
||||
"schema": 1,
|
||||
"schema": 2,
|
||||
"cases": {
|
||||
"background-target-restart": {
|
||||
"gate": "G14",
|
||||
@@ -28,29 +28,248 @@
|
||||
"max_objects": 65,
|
||||
"topology": {"nodes": 4, "drives_per_node": 1},
|
||||
"scope": "Target process killed during partial background rebuild, real unclean-shutdown marker, exact unversioned S3 bodies and replacement-disk shards; not power loss or EC8+4."
|
||||
},
|
||||
"ec84-target-drive-restart": {
|
||||
"gate": "G14",
|
||||
"task": "W20/W21",
|
||||
"lane": "e2e-distributed",
|
||||
"suite": "e2e_test",
|
||||
"name": "distributed::heal_test::three_node_four_drive_ec8_4_root_heal_rebuilds_replaced_drive_after_restart",
|
||||
"oracle": "ec84-target-drive-restart.json",
|
||||
"evidence": "process-restart",
|
||||
"unclean_shutdown_marker": false,
|
||||
"min_objects": 5,
|
||||
"max_objects": 5,
|
||||
"topology": {"nodes": 3, "drives_per_node": 4},
|
||||
"scope": "3-node x 4-drive single-set EC8+4, graceful target restart, preformatted replacement drive, exact unversioned S3 bodies and physical target shards; not mixed-version, multi-pool or long-window ABBA."
|
||||
}
|
||||
},
|
||||
"release_pending": {
|
||||
"G01": "W02/W04 complete root and quota authority coverage",
|
||||
"G02": "W03 bounded checkpoint progress and independent version inventory",
|
||||
"G03": "W17/W18 exact scoped ACK with durable publication and mixed peers",
|
||||
"G04": "W03/W15/W16 crash at every cache/root/floor/intent boundary",
|
||||
"G05": "W06/W07 per-object outcomes and bounded terminal retention",
|
||||
"G06": "W06/W08/W23 concurrent status, legacy clients and truncation",
|
||||
"G07": "W12/W13/W14 durable MRF responsibility at every commit boundary",
|
||||
"G08": "W12/W13/W14 MRF capacity, disk-full and replica-loss matrix",
|
||||
"G09": "W13/W18/W23 actual mixed-version reader/writer and rollback payloads",
|
||||
"G10": "W05/W09/W10/W11 bounded scheduling and pressure recovery",
|
||||
"G11": "W04/W19/W24 maintenance and complete producer coverage",
|
||||
"G12": "W02/W15/W16 both quota paths during reset and settlement",
|
||||
"G13": "W07/W14 quorum-minus-one, unknown disks, remount, Object Lock, dry-run, grace and commit tail",
|
||||
"G14": "W20/W21 same-window field evidence; 3x4 EC8+4 and multi-set/pool coverage",
|
||||
"P1": "W20 measured cold-walk share and foreground latency/throughput",
|
||||
"P2": "W20/W24 measured post-stop convergence and cold segment reuse",
|
||||
"P3": "W20 measured two-hour pressure/heal capacity and recovery window",
|
||||
"P4": "W20 measured MRF scale and replay cost with retained responsibility",
|
||||
"R-E": "W03/W05 fixed-budget real process restart through enumeration and classification",
|
||||
"R-D": "W07/W14 manager-to-event-to-ledger exact disposition, including grace",
|
||||
"R-L": "W13/W14 legacy source conflicts, migration gaps and crash-safe source retirement"
|
||||
}
|
||||
"release_lanes": {
|
||||
"single-set-restart": {
|
||||
"status": "implemented",
|
||||
"cases": ["background-target-restart", "background-target-crash"],
|
||||
"covers": ["four-node one-drive topology", "unversioned objects", "target restart/crash"]
|
||||
},
|
||||
"authority-coverage": {
|
||||
"status": "pending",
|
||||
"gates": ["G01", "G12"],
|
||||
"requires": ["root authority coverage", "quota authority coverage"]
|
||||
},
|
||||
"checkpoint-and-crash": {
|
||||
"status": "pending",
|
||||
"gates": ["G02", "G04", "R-E"],
|
||||
"requires": ["bounded checkpoint progress", "boundary crash matrix", "fixed-budget restart evidence"]
|
||||
},
|
||||
"status-and-outcome": {
|
||||
"status": "pending",
|
||||
"gates": ["G05", "G06", "R-D"],
|
||||
"requires": ["per-object outcomes", "legacy status clients", "manager/event/ledger disposition"]
|
||||
},
|
||||
"mrf-responsibility": {
|
||||
"status": "pending",
|
||||
"gates": ["G07", "G08", "P4"],
|
||||
"requires": ["durable MRF responsibility", "disk-full and replica-loss matrix", "MRF replay cost"]
|
||||
},
|
||||
"mixed-version-rollback": {
|
||||
"status": "pending",
|
||||
"gates": ["G03", "G09", "R-L"],
|
||||
"requires": ["mixed-version peers", "rollback payloads", "crash-safe source retirement"]
|
||||
},
|
||||
"scheduler-pressure": {
|
||||
"status": "pending",
|
||||
"gates": ["G10", "P1", "P2", "P3"],
|
||||
"requires": ["bounded scheduling", "foreground latency and throughput", "two-hour pressure evidence"]
|
||||
},
|
||||
"maintenance-producers": {
|
||||
"status": "pending",
|
||||
"gates": ["G11", "G13"],
|
||||
"requires": ["complete producer coverage", "quorum-minus-one and remount matrix"]
|
||||
},
|
||||
"ec8-4-multiset": {
|
||||
"status": "pending",
|
||||
"gates": ["G14"],
|
||||
"requires": ["3x4 EC8+4 topology", "multi-set coverage", "multi-pool coverage"]
|
||||
}
|
||||
},
|
||||
"release_requirements": [
|
||||
{
|
||||
"gate": "G01",
|
||||
"task": "W02/W04",
|
||||
"lane": "authority-coverage",
|
||||
"status": "pending",
|
||||
"description": "Complete root and quota authority coverage",
|
||||
"requires": ["root authority evidence", "quota authority evidence"]
|
||||
},
|
||||
{
|
||||
"gate": "G02",
|
||||
"task": "W03",
|
||||
"lane": "checkpoint-and-crash",
|
||||
"status": "pending",
|
||||
"description": "Bounded checkpoint progress and independent version inventory",
|
||||
"requires": ["bounded checkpoint oracle", "independent version inventory"]
|
||||
},
|
||||
{
|
||||
"gate": "G03",
|
||||
"task": "W17/W18",
|
||||
"lane": "mixed-version-rollback",
|
||||
"status": "pending",
|
||||
"description": "Exact scoped ACK with durable publication and mixed peers",
|
||||
"requires": ["durable scoped ACK publication", "mixed-peer evidence"],
|
||||
"evidence_fields": [
|
||||
"durable_root_publication_proof",
|
||||
"scoped_ack_request_identity",
|
||||
"participating_peer_capability_snapshot",
|
||||
"mixed_peer_ack_fallback_oracle"
|
||||
]
|
||||
},
|
||||
{
|
||||
"gate": "G04",
|
||||
"task": "W03/W15/W16",
|
||||
"lane": "checkpoint-and-crash",
|
||||
"status": "pending",
|
||||
"description": "Crash at every cache, root, floor and intent boundary",
|
||||
"requires": ["cache boundary crash evidence", "root/floor/intent crash evidence"]
|
||||
},
|
||||
{
|
||||
"gate": "G05",
|
||||
"task": "W06/W07",
|
||||
"lane": "status-and-outcome",
|
||||
"status": "pending",
|
||||
"description": "Per-object outcomes and bounded terminal retention",
|
||||
"requires": ["per-object outcome oracle", "terminal retention bounds"]
|
||||
},
|
||||
{
|
||||
"gate": "G06",
|
||||
"task": "W06/W08/W23",
|
||||
"lane": "status-and-outcome",
|
||||
"status": "pending",
|
||||
"description": "Concurrent status, legacy clients and truncation",
|
||||
"requires": ["concurrent status evidence", "legacy client compatibility", "truncation behavior"]
|
||||
},
|
||||
{
|
||||
"gate": "G07",
|
||||
"task": "W12/W13/W14",
|
||||
"lane": "mrf-responsibility",
|
||||
"status": "pending",
|
||||
"description": "Durable MRF responsibility at every commit boundary",
|
||||
"requires": ["MRF responsibility oracle", "commit-boundary crash matrix"]
|
||||
},
|
||||
{
|
||||
"gate": "G08",
|
||||
"task": "W12/W13/W14",
|
||||
"lane": "mrf-responsibility",
|
||||
"status": "pending",
|
||||
"description": "MRF capacity, disk-full and replica-loss matrix",
|
||||
"requires": ["MRF capacity evidence", "disk-full matrix", "replica-loss matrix"]
|
||||
},
|
||||
{
|
||||
"gate": "G09",
|
||||
"task": "W13/W18/W23",
|
||||
"lane": "mixed-version-rollback",
|
||||
"status": "pending",
|
||||
"description": "Actual mixed-version reader/writer and rollback payloads",
|
||||
"requires": ["mixed-version reader evidence", "mixed-version writer evidence", "rollback payload evidence"],
|
||||
"evidence_fields": [
|
||||
"mixed_version_reader_evidence",
|
||||
"mixed_version_writer_evidence",
|
||||
"rollback_payload_evidence"
|
||||
]
|
||||
},
|
||||
{
|
||||
"gate": "G10",
|
||||
"task": "W05/W09/W10/W11",
|
||||
"lane": "scheduler-pressure",
|
||||
"status": "pending",
|
||||
"description": "Bounded scheduling and pressure recovery",
|
||||
"requires": ["scheduler bound evidence", "pressure recovery evidence"]
|
||||
},
|
||||
{
|
||||
"gate": "G11",
|
||||
"task": "W04/W19/W24",
|
||||
"lane": "maintenance-producers",
|
||||
"status": "pending",
|
||||
"description": "Maintenance and complete producer coverage",
|
||||
"requires": ["maintenance producer matrix", "complete producer inventory"]
|
||||
},
|
||||
{
|
||||
"gate": "G12",
|
||||
"task": "W02/W15/W16",
|
||||
"lane": "authority-coverage",
|
||||
"status": "pending",
|
||||
"description": "Both quota paths during reset and settlement",
|
||||
"requires": ["reset quota-path evidence", "settlement quota-path evidence"]
|
||||
},
|
||||
{
|
||||
"gate": "G13",
|
||||
"task": "W07/W14",
|
||||
"lane": "maintenance-producers",
|
||||
"status": "pending",
|
||||
"description": "Quorum-minus-one, unknown disks, remount, Object Lock, dry-run, grace and commit tail",
|
||||
"requires": ["quorum-minus-one matrix", "unknown-disk/remount matrix", "Object Lock dry-run grace evidence"]
|
||||
},
|
||||
{
|
||||
"gate": "G14",
|
||||
"task": "W20/W21",
|
||||
"lane": "ec8-4-multiset",
|
||||
"status": "pending",
|
||||
"description": "Same-window field evidence with 3x4 EC8+4 and multi-set/pool coverage",
|
||||
"requires": ["same-window field evidence", "3x4 EC8+4 evidence", "multi-set evidence", "multi-pool evidence"]
|
||||
},
|
||||
{
|
||||
"gate": "P1",
|
||||
"task": "W20",
|
||||
"lane": "scheduler-pressure",
|
||||
"status": "pending",
|
||||
"description": "Measured cold-walk share and foreground latency/throughput",
|
||||
"requires": ["cold-walk share measurement", "foreground latency/throughput measurement"]
|
||||
},
|
||||
{
|
||||
"gate": "P2",
|
||||
"task": "W20/W24",
|
||||
"lane": "scheduler-pressure",
|
||||
"status": "pending",
|
||||
"description": "Measured post-stop convergence and cold segment reuse",
|
||||
"requires": ["post-stop convergence measurement", "cold segment reuse measurement"]
|
||||
},
|
||||
{
|
||||
"gate": "P3",
|
||||
"task": "W20",
|
||||
"lane": "scheduler-pressure",
|
||||
"status": "pending",
|
||||
"description": "Measured two-hour pressure/heal capacity and recovery window",
|
||||
"requires": ["two-hour pressure measurement", "heal capacity measurement", "recovery-window measurement"]
|
||||
},
|
||||
{
|
||||
"gate": "P4",
|
||||
"task": "W20",
|
||||
"lane": "mrf-responsibility",
|
||||
"status": "pending",
|
||||
"description": "Measured MRF scale and replay cost with retained responsibility",
|
||||
"requires": ["MRF scale measurement", "MRF replay-cost measurement", "retained responsibility evidence"]
|
||||
},
|
||||
{
|
||||
"gate": "R-E",
|
||||
"task": "W03/W05",
|
||||
"lane": "checkpoint-and-crash",
|
||||
"status": "pending",
|
||||
"description": "Fixed-budget real process restart through enumeration and classification",
|
||||
"requires": ["fixed-budget restart evidence", "enumeration evidence", "classification evidence"]
|
||||
},
|
||||
{
|
||||
"gate": "R-D",
|
||||
"task": "W07/W14",
|
||||
"lane": "status-and-outcome",
|
||||
"status": "pending",
|
||||
"description": "Manager-to-event-to-ledger exact disposition, including grace",
|
||||
"requires": ["manager disposition evidence", "event disposition evidence", "ledger disposition evidence", "grace handling"]
|
||||
},
|
||||
{
|
||||
"gate": "R-L",
|
||||
"task": "W13/W14",
|
||||
"lane": "mixed-version-rollback",
|
||||
"status": "pending",
|
||||
"description": "Legacy source conflicts, migration gaps and crash-safe source retirement",
|
||||
"requires": ["legacy source-conflict evidence", "migration-gap evidence", "crash-safe source retirement evidence"]
|
||||
}
|
||||
]
|
||||
}
|
||||
|
||||
@@ -22,7 +22,7 @@ name: Continuous Integration (docs only)
|
||||
on:
|
||||
pull_request:
|
||||
types: [ opened, synchronize, reopened ]
|
||||
branches: [ main ]
|
||||
branches: [ main, release ]
|
||||
paths:
|
||||
- "**.md"
|
||||
- "docs/**"
|
||||
|
||||
@@ -16,7 +16,7 @@ name: Continuous Integration
|
||||
|
||||
on:
|
||||
push:
|
||||
branches: [ main ]
|
||||
branches: [ main, release ]
|
||||
paths-ignore:
|
||||
- "**.md"
|
||||
- "docs/**"
|
||||
@@ -36,7 +36,7 @@ on:
|
||||
- "flake.lock"
|
||||
pull_request:
|
||||
types: [ opened, synchronize, reopened, closed ]
|
||||
branches: [ main ]
|
||||
branches: [ main, release ]
|
||||
# Keep this list in sync with the `paths` list in ci-docs-only.yml, which
|
||||
# reports the required "Test and Lint" check for PRs skipped here.
|
||||
paths-ignore:
|
||||
@@ -872,13 +872,14 @@ jobs:
|
||||
# Merge gate only (backlog#1149 ci-5): the never-automated user-visible
|
||||
# suites — KMS, object_lock, multipart_auth, quota, checksum, encryption,
|
||||
# security-boundary, ... — via the e2e-full nextest profile. Too heavy for
|
||||
# every PR, so it is gated to main pushes, the merge queue, and manual
|
||||
# every PR, so it is gated to main/release pushes, the merge queue, and manual
|
||||
# dispatch. protocols / the 7 cluster suites / replication / #[ignore] are
|
||||
# owned by other lanes (see .config/nextest.toml profile.e2e-full).
|
||||
if: >-
|
||||
github.event_name == 'workflow_dispatch' ||
|
||||
github.event_name == 'merge_group' ||
|
||||
(github.event_name == 'push' && github.ref == 'refs/heads/main')
|
||||
(github.event_name == 'push' &&
|
||||
(github.ref == 'refs/heads/main' || github.ref == 'refs/heads/release'))
|
||||
needs: [ build-rustfs-debug-binary ]
|
||||
runs-on: sm-standard-2
|
||||
timeout-minutes: 55
|
||||
|
||||
@@ -145,6 +145,38 @@ pub fn consume_verified_mrf_repair_events(anchors: &mut Vec<MrfDurableRepairAnch
|
||||
before.saturating_sub(anchors.len())
|
||||
}
|
||||
|
||||
/// Consume recorded verified repairs for one bucket without draining
|
||||
/// unrelated or still-unmatched proofs. If the caller crashes before
|
||||
/// persisting the retained anchor set, the proof may be replayed by repair
|
||||
/// instead of silently deleting the old responsibility.
|
||||
pub fn consume_recorded_verified_mrf_repair_events_for(bucket: &str, anchors: &mut Vec<MrfDurableRepairAnchor>) -> usize {
|
||||
let Some(registry) = MRF_VERIFIED_REPAIR_EVENTS.get() else {
|
||||
return 0;
|
||||
};
|
||||
let Ok(mut events) = registry.lock() else {
|
||||
return 0;
|
||||
};
|
||||
let before = anchors.len();
|
||||
let mut retained = std::collections::VecDeque::with_capacity(events.len());
|
||||
while let Some(event) = events.pop_front() {
|
||||
if event.bucket.as_ref() != bucket {
|
||||
retained.push_back(event);
|
||||
continue;
|
||||
}
|
||||
let mut matched = false;
|
||||
anchors.retain(|anchor| {
|
||||
let proven = anchor.is_proven_by(&event);
|
||||
matched |= proven;
|
||||
!proven
|
||||
});
|
||||
if !matched {
|
||||
retained.push_back(event);
|
||||
}
|
||||
}
|
||||
*events = retained;
|
||||
before.saturating_sub(anchors.len())
|
||||
}
|
||||
|
||||
#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
|
||||
pub struct MrfScope {
|
||||
pub pool_index: u32,
|
||||
@@ -821,6 +853,71 @@ mod tests {
|
||||
assert!(retained.is_empty());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn recorded_verified_repair_consumer_retains_unmatched_proofs() {
|
||||
let bucket = Arc::<str>::from(format!("recorded-proof-{}", Uuid::new_v4()));
|
||||
let other_bucket = Arc::<str>::from(format!("recorded-proof-other-{}", Uuid::new_v4()));
|
||||
let incarnation = Uuid::new_v4();
|
||||
let lease = MrfIngressLease::new(21);
|
||||
let retained_anchor = MrfDurableRepairAnchor {
|
||||
kind: MrfKind::PartialWrite,
|
||||
bucket: bucket.clone(),
|
||||
object: Arc::from("retained"),
|
||||
version_id: Some([7; 16]),
|
||||
scope: Some(MrfScope {
|
||||
pool_index: 1,
|
||||
set_index: 2,
|
||||
}),
|
||||
lease,
|
||||
bucket_incarnation_id: incarnation,
|
||||
};
|
||||
let waiting_anchor = MrfDurableRepairAnchor {
|
||||
object: Arc::from("waiting"),
|
||||
lease: MrfIngressLease::new(22),
|
||||
..retained_anchor.clone()
|
||||
};
|
||||
let matched_event = MrfVerifiedRepairEvent {
|
||||
kind: retained_anchor.kind,
|
||||
bucket: bucket.clone(),
|
||||
object: retained_anchor.object.clone(),
|
||||
version_id: retained_anchor.version_id,
|
||||
scope: retained_anchor.scope,
|
||||
lease: Some(retained_anchor.lease),
|
||||
bucket_incarnation_id: retained_anchor.bucket_incarnation_id,
|
||||
disposition: MrfVerifiedRepairDisposition::Repaired,
|
||||
};
|
||||
let same_bucket_unmatched = MrfVerifiedRepairEvent {
|
||||
object: Arc::from("future"),
|
||||
lease: Some(MrfIngressLease::new(23)),
|
||||
..matched_event.clone()
|
||||
};
|
||||
let other_bucket_event = MrfVerifiedRepairEvent {
|
||||
bucket: other_bucket.clone(),
|
||||
..matched_event.clone()
|
||||
};
|
||||
|
||||
note_mrf_verified_repair(matched_event);
|
||||
note_mrf_verified_repair(same_bucket_unmatched.clone());
|
||||
note_mrf_verified_repair(other_bucket_event.clone());
|
||||
|
||||
let mut anchors = vec![retained_anchor, waiting_anchor.clone()];
|
||||
assert_eq!(consume_recorded_verified_mrf_repair_events_for(&bucket, &mut anchors), 1);
|
||||
assert_eq!(anchors, vec![waiting_anchor]);
|
||||
|
||||
let remaining_bucket_events = take_mrf_verified_repair_events_for(&bucket);
|
||||
assert_eq!(
|
||||
remaining_bucket_events,
|
||||
vec![same_bucket_unmatched],
|
||||
"same-bucket proofs without a retained anchor must remain available"
|
||||
);
|
||||
let remaining_other_events = take_mrf_verified_repair_events_for(&other_bucket);
|
||||
assert_eq!(
|
||||
remaining_other_events,
|
||||
vec![other_bucket_event],
|
||||
"proofs for other buckets must not be drained by this consumer"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn try_send_delivers_and_respects_capacity() {
|
||||
let mut receiver = init_mrf_channel().expect("first initialization should succeed");
|
||||
|
||||
@@ -59,6 +59,9 @@ const POOL_META_V3_ENV: [(&str, &str); 2] = [
|
||||
|
||||
#[derive(Clone, Copy, Debug)]
|
||||
pub(crate) enum DistLayout {
|
||||
/// 3 nodes × 4 drives, one erasure pool. With `EC:4` this is the
|
||||
/// release-evidence EC8+4 geometry.
|
||||
ThreeByFourEc84,
|
||||
/// 4 nodes × 4 drives, one erasure pool spanning every endpoint.
|
||||
FourByFour,
|
||||
/// 4 nodes × 1 drive, one erasure pool (minimum 4-node 4-disk layout).
|
||||
@@ -94,6 +97,7 @@ impl DistCluster {
|
||||
|
||||
pub async fn new_stopped_with_env(layout: DistLayout, extra_env: &[(&str, &str)]) -> TestResult<Self> {
|
||||
let topology = match layout {
|
||||
DistLayout::ThreeByFourEc84 => ClusterTopology::single_pool_multidrive(3, DRIVES_PER_NODE),
|
||||
DistLayout::FourByFour => ClusterTopology::single_pool_multidrive(NODE_COUNT, DRIVES_PER_NODE),
|
||||
DistLayout::FourNodeFourDisk => ClusterTopology::single_pool(NODE_COUNT),
|
||||
DistLayout::SingleNodeFourDrive => ClusterTopology::per_node_pools(DRIVES_PER_NODE, vec![vec![0]]),
|
||||
@@ -101,7 +105,7 @@ impl DistCluster {
|
||||
let mut cluster = RustFSTestClusterEnvironment::with_topology(topology).await?;
|
||||
let pool_storage_roots = match layout {
|
||||
DistLayout::SingleNodeFourDrive => Some(configured_pool_storage_roots()?),
|
||||
DistLayout::FourByFour | DistLayout::FourNodeFourDisk => None,
|
||||
DistLayout::ThreeByFourEc84 | DistLayout::FourByFour | DistLayout::FourNodeFourDisk => None,
|
||||
};
|
||||
let mut owned_pool_dirs = Vec::new();
|
||||
if let Some(roots) = pool_storage_roots.as_deref() {
|
||||
|
||||
@@ -0,0 +1,183 @@
|
||||
// Copyright 2026 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.
|
||||
|
||||
use super::harness::{DistCluster, DistLayout, TestResult, assert_inventory, payload_for, put_object, unique_bucket, wait_until};
|
||||
use crate::chaos::{VersionShardCensus, census_object_version_on_disk, signed_admin_post};
|
||||
use crate::common::init_logging;
|
||||
use aws_sdk_s3::Client;
|
||||
use aws_sdk_s3::primitives::ByteStream;
|
||||
use std::collections::{BTreeMap, HashSet};
|
||||
use std::path::{Path, PathBuf};
|
||||
use std::time::Duration;
|
||||
|
||||
const EC84_NODE_COUNT: usize = 3;
|
||||
const EC84_DRIVES_PER_NODE: usize = 4;
|
||||
const EC84_DATA_BLOCKS: usize = 8;
|
||||
const EC84_PARITY_BLOCKS: usize = 4;
|
||||
|
||||
#[derive(Clone)]
|
||||
struct ExpectedShard {
|
||||
key: String,
|
||||
body: Vec<u8>,
|
||||
baseline: VersionShardCensus,
|
||||
}
|
||||
|
||||
fn assert_ec84_geometry(census: &VersionShardCensus, key: &str) -> TestResult {
|
||||
if census.data_blocks != Some(EC84_DATA_BLOCKS) || census.parity_blocks != Some(EC84_PARITY_BLOCKS) {
|
||||
return Err(format!("object {key} did not use EC8+4 geometry: {census:?}").into());
|
||||
}
|
||||
let erasure_index = census
|
||||
.erasure_index
|
||||
.ok_or_else(|| format!("object {key} did not record an erasure index: {census:?}"))?;
|
||||
if !(1..=EC84_DATA_BLOCKS + EC84_PARITY_BLOCKS).contains(&erasure_index) {
|
||||
return Err(format!("object {key} has out-of-range erasure index {erasure_index}: {census:?}").into());
|
||||
}
|
||||
if !census.is_complete() || census.expected_part_numbers.is_empty() {
|
||||
return Err(format!("object {key} does not have complete physical shard evidence: {census:?}").into());
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn assert_replaced_drive_empty(drive: &Path, bucket: &str, keys: &[String]) -> TestResult {
|
||||
for key in keys {
|
||||
let census = census_object_version_on_disk(drive, bucket, key, None)?;
|
||||
if census.has_xl_meta {
|
||||
return Err(format!("replacement drive unexpectedly retained {bucket}/{key}: {census:?}").into());
|
||||
}
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn put_large_inventory(client: &Client, bucket: &str) -> TestResult<Vec<ExpectedShard>> {
|
||||
let mut expected = Vec::new();
|
||||
for index in 0..4 {
|
||||
let key = format!("ec84/prefix-{}/object-{index:04}.bin", index % 2);
|
||||
let body = payload_for(&key, 10 * 1024 * 1024);
|
||||
put_object(client, bucket, &key, body.clone()).await?;
|
||||
expected.push(ExpectedShard {
|
||||
key,
|
||||
body,
|
||||
baseline: VersionShardCensus {
|
||||
version_id: None,
|
||||
has_xl_meta: false,
|
||||
data_dir: None,
|
||||
erasure_index: None,
|
||||
data_blocks: None,
|
||||
parity_blocks: None,
|
||||
expected_part_numbers: Default::default(),
|
||||
present_part_fingerprints: Default::default(),
|
||||
inline_data_fingerprint: None,
|
||||
},
|
||||
});
|
||||
}
|
||||
Ok(expected)
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn three_node_four_drive_ec8_4_root_heal_rebuilds_replaced_drive_after_restart() -> TestResult {
|
||||
init_logging();
|
||||
let mut dist = DistCluster::start_with_env(
|
||||
DistLayout::ThreeByFourEc84,
|
||||
&[
|
||||
("RUSTFS_STORAGE_CLASS_STANDARD", "EC:4"),
|
||||
("RUSTFS_HEAL_ENABLED", "true"),
|
||||
("RUSTFS_HEAL_AUTO_HEAL_ENABLE", "false"),
|
||||
("RUSTFS_HEAL_MRF_ENABLE", "false"),
|
||||
("RUSTFS_SCANNER_ENABLED", "false"),
|
||||
],
|
||||
)
|
||||
.await?;
|
||||
assert_eq!(dist.cluster.nodes.len(), EC84_NODE_COUNT);
|
||||
assert_eq!(dist.cluster.topology.drives_per_node, EC84_DRIVES_PER_NODE);
|
||||
|
||||
let bucket = unique_bucket("healec84");
|
||||
dist.create_bucket(&bucket).await?;
|
||||
let writer = dist.client(0)?;
|
||||
let mut expected = put_large_inventory(&writer, &bucket).await?;
|
||||
let replaced_node = 1;
|
||||
let replaced_drive_index = 2;
|
||||
let replaced_drive = PathBuf::from(&dist.cluster.nodes[replaced_node].data_dirs[replaced_drive_index]);
|
||||
|
||||
for item in &mut expected {
|
||||
item.baseline = census_object_version_on_disk(&replaced_drive, &bucket, &item.key, None)?;
|
||||
assert_ec84_geometry(&item.baseline, &item.key)?;
|
||||
}
|
||||
|
||||
let format_path = replaced_drive.join(".rustfs.sys").join("format.json");
|
||||
let format_json = std::fs::read(&format_path)?;
|
||||
dist.cluster.stop_node_gracefully(replaced_node).await?;
|
||||
let retired_drive = PathBuf::from(format!("{}.retired", replaced_drive.display()));
|
||||
std::fs::rename(&replaced_drive, &retired_drive)?;
|
||||
std::fs::create_dir_all(format_path.parent().ok_or("replacement format path has no parent")?)?;
|
||||
std::fs::write(&format_path, format_json)?;
|
||||
assert_replaced_drive_empty(
|
||||
&replaced_drive,
|
||||
&bucket,
|
||||
&expected.iter().map(|item| item.key.clone()).collect::<Vec<_>>(),
|
||||
)?;
|
||||
|
||||
let outage_key = "ec84/written-while-node-restarting.bin";
|
||||
let outage_body = payload_for(outage_key, 10 * 1024 * 1024);
|
||||
writer
|
||||
.put_object()
|
||||
.bucket(&bucket)
|
||||
.key(outage_key)
|
||||
.body(ByteStream::from(outage_body.clone()))
|
||||
.send()
|
||||
.await?;
|
||||
|
||||
dist.cluster.start_node(replaced_node).await?;
|
||||
let heal_body =
|
||||
r#"{"recursive":true,"dryRun":false,"remove":false,"recreate":true,"scanMode":2,"updateParity":false,"nolock":false}"#;
|
||||
let heal_url = format!("{}/rustfs/admin/v3/heal/{bucket}?forceStart=true", dist.cluster.nodes[0].url);
|
||||
signed_admin_post(&heal_url, Some(heal_body), &dist.cluster.access_key, &dist.cluster.secret_key).await?;
|
||||
|
||||
wait_until(
|
||||
Duration::from_secs(120),
|
||||
|| async {
|
||||
for item in &expected {
|
||||
let current = census_object_version_on_disk(&replaced_drive, &bucket, &item.key, None)?;
|
||||
if !current.matches_manifest(&item.baseline) {
|
||||
return Ok(false);
|
||||
}
|
||||
}
|
||||
let outage = census_object_version_on_disk(&replaced_drive, &bucket, outage_key, None)?;
|
||||
Ok(outage.is_complete()
|
||||
&& outage.data_blocks == Some(EC84_DATA_BLOCKS)
|
||||
&& outage.parity_blocks == Some(EC84_PARITY_BLOCKS))
|
||||
},
|
||||
"EC8+4 replacement drive rebuilt baseline and outage shards",
|
||||
)
|
||||
.await?;
|
||||
|
||||
let inventory = expected
|
||||
.iter()
|
||||
.map(|item| (item.key.clone(), item.body.clone()))
|
||||
.chain(std::iter::once((outage_key.to_string(), outage_body.clone())))
|
||||
.collect::<BTreeMap<_, _>>();
|
||||
let expected_keys = inventory.keys().cloned().collect::<HashSet<_>>();
|
||||
for node_index in 0..dist.cluster.nodes.len() {
|
||||
let client = dist.client(node_index)?;
|
||||
assert_inventory(&client, &bucket, &inventory).await?;
|
||||
let listing = client.list_objects_v2().bucket(&bucket).send().await?;
|
||||
let observed = listing
|
||||
.contents()
|
||||
.iter()
|
||||
.filter_map(|object| object.key().map(str::to_owned))
|
||||
.collect::<HashSet<_>>();
|
||||
assert_eq!(observed, expected_keys, "node {node_index} listing diverged after EC8+4 heal");
|
||||
}
|
||||
|
||||
Ok(())
|
||||
}
|
||||
@@ -25,6 +25,7 @@ mod data_integrity_movement_test;
|
||||
mod expand_decommission_rebalance_test;
|
||||
mod extra_test;
|
||||
mod harness;
|
||||
mod heal_test;
|
||||
mod object_lock_test;
|
||||
mod observability_test;
|
||||
mod replication_quota_test;
|
||||
|
||||
@@ -8220,7 +8220,7 @@ mod tests {
|
||||
}
|
||||
mod multipart_transport_tests {
|
||||
use super::super::super::replication_filemeta_boundary::ObjectPartInfo;
|
||||
use super::super::super::replication_storage_boundary::ObjectIO as _;
|
||||
use super::super::super::replication_storage_boundary::{ObjectIO as _, ReadPlan};
|
||||
use super::*;
|
||||
use bytes::Bytes;
|
||||
use http_body_util::{BodyExt, Full};
|
||||
@@ -8264,8 +8264,7 @@ mod tests {
|
||||
} else {
|
||||
self.full_reads.fetch_add(1, Ordering::Relaxed);
|
||||
}
|
||||
let plan =
|
||||
crate::object_api::ReadPlan::build_for_request(range, &self.info, opts, &HeaderMap::new(), None).await?;
|
||||
let plan = ReadPlan::build_for_request(range, &self.info, opts, &HeaderMap::new(), None).await?;
|
||||
let start = plan.storage_offset();
|
||||
let end = start + usize::try_from(plan.storage_length()).expect("nonnegative storage length");
|
||||
return plan.into_object_reader(Box::new(std::io::Cursor::new(stored.slice(start..end))), &self.info);
|
||||
|
||||
@@ -22,7 +22,7 @@ pub(crate) use crate::object_api::{
|
||||
GetObjectReader, ObjectInfo, ObjectOptions, PutObjReader, ReplicationStatusWritebackCondition, ReplicationStatusWritebackMode,
|
||||
};
|
||||
#[cfg(test)]
|
||||
pub(crate) use crate::object_api::{NamespaceLockFence, NamespaceLockSignalTestFence};
|
||||
pub(crate) use crate::object_api::{NamespaceLockFence, NamespaceLockSignalTestFence, ReadPlan};
|
||||
pub(crate) use crate::storage_api_contracts::list::{
|
||||
ListOperations, StorageListObjectVersionsInfo, StorageListObjectsV2Info, StorageObjectInfoOrErr, StorageWalkOptions,
|
||||
};
|
||||
|
||||
@@ -285,6 +285,20 @@ fn peer_replay_state(audience: &str) -> PeerReplayState {
|
||||
.unwrap_or_default()
|
||||
}
|
||||
|
||||
pub(crate) fn clear_peer_replay_state_for_addr(addr: &str) -> std::io::Result<()> {
|
||||
let uri = addr
|
||||
.parse::<Uri>()
|
||||
.map_err(|_| std::io::Error::other("Invalid gRPC peer URI"))?;
|
||||
let audience = uri
|
||||
.authority()
|
||||
.map(|authority| normalize_tonic_rpc_audience(authority.as_str()))
|
||||
.ok_or_else(|| std::io::Error::other("Missing gRPC peer authority"))??;
|
||||
if let Ok(mut states) = PEER_REPLAY_STATES.lock() {
|
||||
states.remove(&audience);
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn apply_peer_replay_response(
|
||||
audience: String,
|
||||
sent_state: PeerReplayState,
|
||||
@@ -619,6 +633,13 @@ mod tests {
|
||||
.remove(audience);
|
||||
}
|
||||
|
||||
fn set_peer_capability(audience: &str, state: PeerReplayState) {
|
||||
PEER_REPLAY_STATES
|
||||
.lock()
|
||||
.expect("peer capability cache lock must not be poisoned")
|
||||
.insert(audience.to_string(), state);
|
||||
}
|
||||
|
||||
fn rolling_mutation_request(method: &'static str) -> tonic::Request<()> {
|
||||
let mut request = tonic::Request::new(rustfs_protos::proto_gen::node_service::GenerallyLockRequest {
|
||||
args: "canonical mutation request".to_string(),
|
||||
@@ -1090,6 +1111,23 @@ mod tests {
|
||||
clear_peer_capability(audience);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn clear_peer_replay_state_for_addr_removes_normalized_audience() {
|
||||
let audience = "clear-peer-replay-state-test:9000";
|
||||
let boot_epoch = Uuid::new_v4();
|
||||
set_peer_capability(
|
||||
audience,
|
||||
PeerReplayState {
|
||||
boot_epoch: Some(boot_epoch),
|
||||
cache_capability: Some(PeerReplayCapability::Capable { boot_epoch }),
|
||||
},
|
||||
);
|
||||
|
||||
clear_peer_replay_state_for_addr("http://clear-peer-replay-state-test:9000").expect("peer URI should clear replay state");
|
||||
|
||||
assert_eq!(peer_replay_state(audience), PeerReplayState::default());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn interceptor_snapshot_prevents_delayed_legacy_response_from_revoking_capability() {
|
||||
ensure_test_rpc_secret();
|
||||
|
||||
@@ -13,8 +13,9 @@
|
||||
// limitations under the License.
|
||||
|
||||
use crate::cluster::rpc::client::{
|
||||
AuthenticatedChannel, TonicInterceptor, embedded_tonic_status, gen_tonic_signature_interceptor, heal_control_time_out_client,
|
||||
is_network_like_status, message_has_network_needle, node_service_time_out_client, tier_mutation_control_time_out_client,
|
||||
AuthenticatedChannel, TonicInterceptor, clear_peer_replay_state_for_addr, embedded_tonic_status,
|
||||
gen_tonic_signature_interceptor, heal_control_time_out_client, is_network_like_status, message_has_network_needle,
|
||||
node_service_time_out_client, tier_mutation_control_time_out_client,
|
||||
};
|
||||
use crate::cluster::rpc::{set_tonic_canonical_body_digest, set_tonic_mutation_body_digest, verify_tonic_rpc_response_proof};
|
||||
use crate::error::{Error, Result};
|
||||
@@ -544,6 +545,16 @@ fn validate_heal_control_response_proof(canonical_response: &[u8], proof: &[u8])
|
||||
.map_err(|_| Error::other("peer returned an invalid heal control response proof"))
|
||||
}
|
||||
|
||||
fn heal_control_auth_may_need_replay_scope_refresh(err: &Error) -> bool {
|
||||
matches!(
|
||||
err,
|
||||
Error::Io(io_err)
|
||||
if embedded_tonic_status(io_err).is_some_and(|status| {
|
||||
status.code() == tonic::Code::Unauthenticated && status.message() == "No valid auth token"
|
||||
})
|
||||
)
|
||||
}
|
||||
|
||||
fn decode_remote_version_state_capability(expected_member: &str, result: &[u8]) -> Result<Uuid> {
|
||||
let (topology_member, process_epoch) = rustfs_protos::decode_remote_version_state_capability(result).map_err(Error::other)?;
|
||||
if topology_member != expected_member {
|
||||
@@ -1720,45 +1731,72 @@ impl PeerRestClient {
|
||||
return Err(Error::other("heal control command exceeds size limit"));
|
||||
}
|
||||
let capability_probe = rustfs_protos::is_heal_control_capability_probe(&command);
|
||||
self.finalize_result(
|
||||
async {
|
||||
let mut client = self
|
||||
.get_heal_control_client()
|
||||
.await?
|
||||
.max_encoding_message_size(rustfs_protos::HEAL_CONTROL_RPC_MAX_MESSAGE_SIZE)
|
||||
.max_decoding_message_size(rustfs_protos::HEAL_CONTROL_RPC_MAX_MESSAGE_SIZE);
|
||||
let canonical_body = rustfs_protos::canonical_heal_control_request_body(version, &topology_fingerprint, &command)
|
||||
.map_err(|_| Error::other("heal control request length cannot be represented"))?;
|
||||
let mut request = Request::new(HealControlRequest {
|
||||
version,
|
||||
topology_fingerprint: topology_fingerprint.clone(),
|
||||
command: command.clone().into(),
|
||||
});
|
||||
request.set_timeout(rustfs_protos::heal_control_execution_timeout());
|
||||
set_tonic_canonical_body_digest(&mut request, &canonical_body)?;
|
||||
let response = client.heal_control(request).await?.into_inner();
|
||||
if !response.success {
|
||||
return Err(Error::other(
|
||||
response
|
||||
.error_info
|
||||
.unwrap_or_else(|| "peer heal control failed without an error".to_string()),
|
||||
));
|
||||
}
|
||||
if !capability_probe {
|
||||
let canonical_response = rustfs_protos::canonical_heal_control_response_body(
|
||||
version,
|
||||
&topology_fingerprint,
|
||||
&command,
|
||||
&response.result,
|
||||
)
|
||||
let result = self
|
||||
.heal_control_once(version, &topology_fingerprint, &command, capability_probe)
|
||||
.await;
|
||||
if result
|
||||
.as_ref()
|
||||
.err()
|
||||
.is_some_and(heal_control_auth_may_need_replay_scope_refresh)
|
||||
{
|
||||
self.prepare_heal_control_auth_retry().await;
|
||||
return self
|
||||
.finalize_result(
|
||||
self.heal_control_once(version, &topology_fingerprint, &command, capability_probe)
|
||||
.await,
|
||||
)
|
||||
.await;
|
||||
}
|
||||
self.finalize_result(result).await
|
||||
}
|
||||
|
||||
async fn prepare_heal_control_auth_retry(&self) {
|
||||
if let Err(err) = clear_peer_replay_state_for_addr(&self.grid_host) {
|
||||
debug!(
|
||||
peer = %self.grid_host,
|
||||
error = %err,
|
||||
"could not clear heal control replay state before retry"
|
||||
);
|
||||
}
|
||||
self.evict_connection().await;
|
||||
}
|
||||
|
||||
async fn heal_control_once(
|
||||
&self,
|
||||
version: u32,
|
||||
topology_fingerprint: &str,
|
||||
command: &[u8],
|
||||
capability_probe: bool,
|
||||
) -> Result<Vec<u8>> {
|
||||
let mut client = self
|
||||
.get_heal_control_client()
|
||||
.await?
|
||||
.max_encoding_message_size(rustfs_protos::HEAL_CONTROL_RPC_MAX_MESSAGE_SIZE)
|
||||
.max_decoding_message_size(rustfs_protos::HEAL_CONTROL_RPC_MAX_MESSAGE_SIZE);
|
||||
let canonical_body = rustfs_protos::canonical_heal_control_request_body(version, topology_fingerprint, command)
|
||||
.map_err(|_| Error::other("heal control request length cannot be represented"))?;
|
||||
let mut request = Request::new(HealControlRequest {
|
||||
version,
|
||||
topology_fingerprint: topology_fingerprint.to_string(),
|
||||
command: command.to_vec().into(),
|
||||
});
|
||||
request.set_timeout(rustfs_protos::heal_control_execution_timeout());
|
||||
set_tonic_canonical_body_digest(&mut request, &canonical_body)?;
|
||||
let response = client.heal_control(request).await?.into_inner();
|
||||
if !response.success {
|
||||
return Err(Error::other(
|
||||
response
|
||||
.error_info
|
||||
.unwrap_or_else(|| "peer heal control failed without an error".to_string()),
|
||||
));
|
||||
}
|
||||
if !capability_probe {
|
||||
let canonical_response =
|
||||
rustfs_protos::canonical_heal_control_response_body(version, topology_fingerprint, command, &response.result)
|
||||
.map_err(|_| Error::other("heal control response length cannot be represented"))?;
|
||||
validate_heal_control_response_proof(&canonical_response, &response.response_proof)?;
|
||||
}
|
||||
Ok(response.result.to_vec())
|
||||
}
|
||||
.await,
|
||||
)
|
||||
.await
|
||||
validate_heal_control_response_proof(&canonical_response, &response.response_proof)?;
|
||||
}
|
||||
Ok(response.result.to_vec())
|
||||
}
|
||||
|
||||
/// Confirms that a peer supports the current heal-control coordination
|
||||
@@ -3728,6 +3766,22 @@ mod tests {
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn heal_control_auth_retry_is_limited_to_transport_auth_rejection() {
|
||||
assert!(heal_control_auth_may_need_replay_scope_refresh(&Error::from(
|
||||
tonic::Status::unauthenticated("No valid auth token")
|
||||
)));
|
||||
assert!(!heal_control_auth_may_need_replay_scope_refresh(&Error::from(
|
||||
tonic::Status::permission_denied("bad signature")
|
||||
)));
|
||||
assert!(!heal_control_auth_may_need_replay_scope_refresh(&Error::from(
|
||||
tonic::Status::unauthenticated("application rejected heal control")
|
||||
)));
|
||||
assert!(!heal_control_auth_may_need_replay_scope_refresh(&Error::other(
|
||||
"Io error: code: 'Unauthenticated', message: \"No valid auth token\""
|
||||
)));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn peer_rest_client_network_classifier_keeps_slow_peers_online() {
|
||||
// The per-RPC channel deadline (RUSTFS_INTERNODE_RPC_TIMEOUT, 30s)
|
||||
|
||||
@@ -506,6 +506,25 @@ fn active_heal_for_dedup_key(active_heals: &HashMap<String, Arc<HealTask>>, key:
|
||||
.map(|(task_id, task)| (task_id.clone(), task.heal_type.clone()))
|
||||
}
|
||||
|
||||
fn request_matches_task(request: &HealRequest, task: &HealTask) -> bool {
|
||||
request.heal_type == task.heal_type
|
||||
&& request.options == task.options
|
||||
&& request.priority == task.priority
|
||||
&& request.source == task.source
|
||||
&& request.retry_attempts == task.retry_attempts
|
||||
&& request.heal_endpoints == task.heal_endpoints
|
||||
}
|
||||
|
||||
fn request_matches_request(request: &HealRequest, existing: &HealRequest) -> bool {
|
||||
request.heal_type == existing.heal_type
|
||||
&& request.options == existing.options
|
||||
&& request.priority == existing.priority
|
||||
&& request.source == existing.source
|
||||
&& request.force_start == existing.force_start
|
||||
&& request.retry_attempts == existing.retry_attempts
|
||||
&& request.heal_endpoints == existing.heal_endpoints
|
||||
}
|
||||
|
||||
fn retrying_heal_for_dedup_key(retrying_heals: &HashMap<String, RetryingHeal>, key: &str) -> Option<(String, HealType)> {
|
||||
retrying_heals
|
||||
.iter()
|
||||
@@ -1613,6 +1632,64 @@ impl HealManager {
|
||||
pause_duplicate_admission_after_active_lock(&request.id).await;
|
||||
let mut queue = self.heal_queue.lock().await;
|
||||
let retrying_heals = self.retrying_heals.lock().await;
|
||||
|
||||
let request_id_admission = active_heals
|
||||
.get(&request.id)
|
||||
.map(|task| (request_matches_task(&request, task), "active"))
|
||||
.or_else(|| {
|
||||
queue
|
||||
.requests()
|
||||
.find(|queued| queued.id == request.id)
|
||||
.map(|queued| (request_matches_request(&request, queued), "queued"))
|
||||
})
|
||||
.or_else(|| {
|
||||
retrying_heals
|
||||
.get(&request.id)
|
||||
.map(|retrying| (request_matches_request(&request, &retrying.request), "retrying"))
|
||||
});
|
||||
if let Some((matches_existing, duplicate_state)) = request_id_admission {
|
||||
let admission = if matches_existing {
|
||||
HealAdmissionResult::Accepted
|
||||
} else {
|
||||
HealAdmissionResult::Dropped(HealAdmissionDropReason::AlreadyRunning)
|
||||
};
|
||||
if matches!(admission, HealAdmissionResult::Accepted | HealAdmissionResult::Merged)
|
||||
&& let Some(target) = mrf_notice_target
|
||||
{
|
||||
let mut targets = lock_mrf_repair_notice_targets(&self.mrf_repair_notice_targets);
|
||||
Self::insert_mrf_repair_notice_target(&mut targets, &request.id, target);
|
||||
}
|
||||
drop(retrying_heals);
|
||||
drop(queue);
|
||||
drop(active_heals);
|
||||
let lock_phase = lock_phase_start.elapsed();
|
||||
Self::record_admission_metric(request.source, admission, "duplicate");
|
||||
self.record_admission_observation(HealAdmissionObservation {
|
||||
source,
|
||||
result: admission,
|
||||
context: "duplicate",
|
||||
force_start,
|
||||
displaced: false,
|
||||
start_duration: admission_start.elapsed(),
|
||||
lock_phase,
|
||||
});
|
||||
debug!(
|
||||
target: "rustfs::heal::manager",
|
||||
event = EVENT_HEAL_QUEUE_ADMISSION,
|
||||
component = LOG_COMPONENT_HEAL,
|
||||
subsystem = LOG_SUBSYSTEM_MANAGER,
|
||||
request_id = %request.id,
|
||||
duplicate_state,
|
||||
result = admission.result_label(),
|
||||
reason = admission.reason_label(),
|
||||
"Heal queue admission reused an existing request id"
|
||||
);
|
||||
return Ok(HealAdmissionReceipt {
|
||||
result: admission,
|
||||
task_id: request.id,
|
||||
});
|
||||
}
|
||||
|
||||
let duplicate = (!request.force_start).then(|| {
|
||||
active_heal_for_dedup_key(&active_heals, &dedup_key)
|
||||
.map(|(task_id, _)| (task_id, "active"))
|
||||
|
||||
@@ -4537,6 +4537,48 @@ async fn test_force_start_marks_dedup_key_for_future_duplicates() {
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn same_request_id_replay_reuses_existing_task_without_force_start_duplication() {
|
||||
let storage: Arc<dyn HealStorageAPI> = Arc::new(MockStorage);
|
||||
let manager = HealManager::new(storage, None);
|
||||
|
||||
let mut original = admin_prefix_request("bucket", "logs/");
|
||||
original.force_start = true;
|
||||
let original_id = original.id.clone();
|
||||
let accepted = manager
|
||||
.submit_heal_request_with_receipt(original.clone())
|
||||
.await
|
||||
.expect("original forceStart request should queue");
|
||||
assert_eq!(accepted.result, HealAdmissionResult::Accepted);
|
||||
assert_eq!(accepted.task_id, original_id);
|
||||
|
||||
let replayed = manager
|
||||
.submit_heal_request_with_receipt(original.clone())
|
||||
.await
|
||||
.expect("same request id and payload should reuse the existing task");
|
||||
assert_eq!(replayed.result, HealAdmissionResult::Accepted);
|
||||
assert_eq!(replayed.task_id, original_id);
|
||||
assert_eq!(
|
||||
manager.get_queue_length().await,
|
||||
1,
|
||||
"exact forceStart replay must not create a second queued task"
|
||||
);
|
||||
|
||||
let mut changed = original;
|
||||
changed.options.remove_corrupted = true;
|
||||
let changed = manager
|
||||
.submit_heal_request_with_receipt(changed)
|
||||
.await
|
||||
.expect("same request id with a changed payload should fail closed");
|
||||
assert_eq!(changed.result, HealAdmissionResult::Dropped(HealAdmissionDropReason::AlreadyRunning));
|
||||
assert_eq!(changed.task_id, original_id);
|
||||
assert_eq!(
|
||||
manager.get_queue_length().await,
|
||||
1,
|
||||
"same-id conflict must not displace or duplicate the original task"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_running_heal_set_counts_groups_set_scoped_tasks() {
|
||||
let storage: Arc<dyn HealStorageAPI> = Arc::new(MockStorage);
|
||||
|
||||
@@ -527,6 +527,9 @@ struct MrfRuntime {
|
||||
/// True while a journal snapshot exists on disk that may still be needed
|
||||
/// for replay or cleanup.
|
||||
journal_on_disk: bool,
|
||||
/// True after replay admitted records into the manager and the startup
|
||||
/// journal must stay until a durable replay proof is persisted.
|
||||
retain_replay_journal: bool,
|
||||
/// Earliest instant a full-admission retry may proceed.
|
||||
backoff_until: Option<tokio::time::Instant>,
|
||||
}
|
||||
@@ -673,10 +676,11 @@ pub async fn replay_journal_once(manager: &Arc<HealManager>) -> usize {
|
||||
struct ReplayOutcome {
|
||||
replayed: usize,
|
||||
journal_on_disk: bool,
|
||||
retain_journal_for_replay: bool,
|
||||
}
|
||||
|
||||
fn replay_must_retain_journal(rearm_incomplete: bool, pending_depth: usize) -> bool {
|
||||
rearm_incomplete || pending_depth > 0
|
||||
fn replay_must_retain_journal(rearm_incomplete: bool, pending_depth: usize, accepted_or_merged: bool) -> bool {
|
||||
rearm_incomplete || pending_depth > 0 || accepted_or_merged
|
||||
}
|
||||
|
||||
/// Shared replay core: read + decode + re-arm, then drain what fits. The
|
||||
@@ -698,6 +702,7 @@ async fn replay_into(
|
||||
return ReplayOutcome {
|
||||
replayed: 0,
|
||||
journal_on_disk: false,
|
||||
retain_journal_for_replay: false,
|
||||
};
|
||||
}
|
||||
},
|
||||
@@ -722,6 +727,7 @@ async fn replay_into(
|
||||
// prefix.
|
||||
queue.raise_limits_for_replay(intents.len(), replay_bytes);
|
||||
let mut rearm_incomplete = false;
|
||||
let mut accepted_or_merged = false;
|
||||
for intent in intents {
|
||||
let result = queue.try_push_typed(intent.clone());
|
||||
match result {
|
||||
@@ -748,7 +754,9 @@ async fn replay_into(
|
||||
break;
|
||||
}
|
||||
match submit_mrf_heal_request(manager, &intent).await {
|
||||
Ok(HealAdmissionResult::Accepted) | Ok(HealAdmissionResult::Merged) => {}
|
||||
Ok(HealAdmissionResult::Accepted) | Ok(HealAdmissionResult::Merged) => {
|
||||
accepted_or_merged = true;
|
||||
}
|
||||
Ok(HealAdmissionResult::Full) | Ok(HealAdmissionResult::Dropped(HealAdmissionDropReason::QueueFull)) => {
|
||||
intent.attempts = intent.attempts.saturating_add(1);
|
||||
if intent.attempts < MRF_MAX_ATTEMPTS {
|
||||
@@ -779,7 +787,8 @@ async fn replay_into(
|
||||
}
|
||||
}
|
||||
}
|
||||
let journal_on_disk = if replay_must_retain_journal(rearm_incomplete, queue.depth()) {
|
||||
let retain_journal_for_replay = replay_must_retain_journal(rearm_incomplete, queue.depth(), accepted_or_merged);
|
||||
let journal_on_disk = if retain_journal_for_replay {
|
||||
true
|
||||
} else {
|
||||
!delete_journals().await
|
||||
@@ -787,6 +796,7 @@ async fn replay_into(
|
||||
ReplayOutcome {
|
||||
replayed,
|
||||
journal_on_disk,
|
||||
retain_journal_for_replay,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -800,6 +810,7 @@ async fn run_mrf_consumer(manager: Arc<HealManager>, mut receiver: mpsc::Receive
|
||||
new_since_flush: 0,
|
||||
dirty: false,
|
||||
journal_on_disk: false,
|
||||
retain_replay_journal: false,
|
||||
backoff_until: None,
|
||||
};
|
||||
|
||||
@@ -807,6 +818,7 @@ async fn run_mrf_consumer(manager: Arc<HealManager>, mut receiver: mpsc::Receive
|
||||
// on disk whenever any replayed intent still needs a successor snapshot.
|
||||
let replay = replay_into(&manager, &mut runtime.queue, &mut runtime.backoff_until).await;
|
||||
runtime.journal_on_disk = replay.journal_on_disk;
|
||||
runtime.retain_replay_journal = replay.retain_journal_for_replay;
|
||||
// Anything still pending (e.g. the manager was full and backoff armed)
|
||||
// must be re-persisted by the next flush before replay can delete the
|
||||
// startup anchor.
|
||||
@@ -854,6 +866,7 @@ async fn run_mrf_consumer(manager: Arc<HealManager>, mut receiver: mpsc::Receive
|
||||
runtime.dirty,
|
||||
runtime.queue.depth(),
|
||||
runtime.journal_on_disk,
|
||||
runtime.retain_replay_journal,
|
||||
) {
|
||||
TickAction::Flush => {
|
||||
runtime.flush().await;
|
||||
@@ -897,12 +910,12 @@ enum TickAction {
|
||||
Idle,
|
||||
}
|
||||
|
||||
fn tick_action(dirty: bool, depth: usize, journal_on_disk: bool) -> TickAction {
|
||||
fn tick_action(dirty: bool, depth: usize, journal_on_disk: bool, retain_replay_journal: bool) -> TickAction {
|
||||
if dirty {
|
||||
TickAction::Flush
|
||||
} else if depth > 0 {
|
||||
TickAction::Retry
|
||||
} else if journal_on_disk {
|
||||
} else if journal_on_disk && !retain_replay_journal {
|
||||
TickAction::DeleteJournal
|
||||
} else {
|
||||
TickAction::Idle
|
||||
@@ -934,34 +947,39 @@ mod tests {
|
||||
|
||||
// Dirty dominates: a changed pending set flushes even when idle
|
||||
// otherwise.
|
||||
assert!(matches!(tick_action(true, 0, false), Flush));
|
||||
assert!(matches!(tick_action(true, 3, true), Flush));
|
||||
assert!(matches!(tick_action(true, 0, false, false), Flush));
|
||||
assert!(matches!(tick_action(true, 3, true, false), Flush));
|
||||
|
||||
// Clean backlog: no rewrite, but keep draining so an expired
|
||||
// admission backoff retries on time.
|
||||
assert!(matches!(tick_action(false, 1, false), Retry));
|
||||
assert!(matches!(tick_action(false, 2, true), Retry));
|
||||
assert!(matches!(tick_action(false, 1, false, false), Retry));
|
||||
assert!(matches!(tick_action(false, 2, true, false), Retry));
|
||||
|
||||
// Quiescent with a stale journal file on disk: remove it.
|
||||
assert!(matches!(tick_action(false, 0, true), DeleteJournal));
|
||||
assert!(matches!(tick_action(false, 0, true, false), DeleteJournal));
|
||||
assert!(matches!(tick_action(false, 0, true, true), Idle));
|
||||
|
||||
// Fully quiescent: nothing to do.
|
||||
assert!(matches!(tick_action(false, 0, false), Idle));
|
||||
assert!(matches!(tick_action(false, 0, false, false), Idle));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn replay_cleanup_retains_journal_for_unarmed_or_refused_records() {
|
||||
assert!(
|
||||
replay_must_retain_journal(true, 0),
|
||||
replay_must_retain_journal(true, 0, false),
|
||||
"a rejected replay record still needs its disk anchor"
|
||||
);
|
||||
assert!(
|
||||
replay_must_retain_journal(false, 1),
|
||||
replay_must_retain_journal(false, 1, false),
|
||||
"a Full admission retry must keep the startup journal until the next snapshot"
|
||||
);
|
||||
assert!(
|
||||
!replay_must_retain_journal(false, 0),
|
||||
"only a fully consumed replay snapshot may be deleted"
|
||||
!replay_must_retain_journal(false, 0, false),
|
||||
"a fully consumed replay snapshot without accepts may be deleted"
|
||||
);
|
||||
assert!(
|
||||
replay_must_retain_journal(false, 0, true),
|
||||
"accepted or merged replay records still need a durable successor"
|
||||
);
|
||||
}
|
||||
|
||||
|
||||
@@ -27,7 +27,9 @@
|
||||
|
||||
use super::{MRF_JOURNAL_PATH, MRF_SCOPED_JOURNAL_PATH, decode_journal};
|
||||
use crate::heal::RUSTFS_META_BUCKET;
|
||||
use crate::heal::storage_api::owner::{EcstoreDiskAPI, EcstoreDiskError, EcstoreDiskStore};
|
||||
use crate::heal::storage_api::owner::{
|
||||
EcstoreConditionalFileUpdate, EcstoreDiskAPI, EcstoreDiskBytes, EcstoreDiskError, EcstoreDiskStore,
|
||||
};
|
||||
use sha2::{Digest, Sha256};
|
||||
use std::collections::HashMap;
|
||||
use tokio::io::AsyncReadExt;
|
||||
@@ -54,6 +56,8 @@ pub enum SnapshotError {
|
||||
TooLarge,
|
||||
#[error("MRF checkpoint replicas disagree at the same sequence")]
|
||||
Conflict,
|
||||
#[error("MRF checkpoint has no writable replica")]
|
||||
NoWritableReplica,
|
||||
#[error("MRF checkpoint storage is unavailable")]
|
||||
Disk(#[source] EcstoreDiskError),
|
||||
#[error("MRF checkpoint body could not be read")]
|
||||
@@ -121,6 +125,7 @@ impl Manifest {
|
||||
pub struct CommittedSnapshot {
|
||||
manifest: Manifest,
|
||||
payload: Vec<u8>,
|
||||
slot: usize,
|
||||
}
|
||||
|
||||
#[derive(Default)]
|
||||
@@ -141,13 +146,18 @@ impl CommittedSnapshot {
|
||||
self.manifest.owner
|
||||
}
|
||||
|
||||
/// Slot that supplied this committed checkpoint.
|
||||
pub fn slot(&self) -> usize {
|
||||
self.slot
|
||||
}
|
||||
|
||||
/// Complete, checksum-validated record bytes. Inspection does not consume
|
||||
/// these records or acknowledge completion to any producer.
|
||||
pub fn payload(&self) -> &[u8] {
|
||||
&self.payload
|
||||
}
|
||||
|
||||
fn decode(manifest: &[u8], payload: Vec<u8>, limit: usize) -> Result<Self, SnapshotError> {
|
||||
fn decode(slot: usize, manifest: &[u8], payload: Vec<u8>, limit: usize) -> Result<Self, SnapshotError> {
|
||||
let manifest = Manifest::decode(manifest, limit)?;
|
||||
let checksum: [u8; 32] = Sha256::digest(&payload).into();
|
||||
if payload.len() != manifest.payload_len || checksum != manifest.payload_digest {
|
||||
@@ -156,10 +166,20 @@ impl CommittedSnapshot {
|
||||
if decode_journal(&payload).1 != 0 {
|
||||
return Err(SnapshotError::Corrupt);
|
||||
}
|
||||
Ok(Self { manifest, payload })
|
||||
Ok(Self { manifest, payload, slot })
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||
pub struct SnapshotPublication {
|
||||
pub owner: Uuid,
|
||||
pub sequence: u64,
|
||||
pub slot: usize,
|
||||
pub payload_len: usize,
|
||||
pub payload_replicas: usize,
|
||||
pub manifest_replicas: usize,
|
||||
}
|
||||
|
||||
#[derive(Debug)]
|
||||
pub enum RecoverySnapshot {
|
||||
/// An intact legacy snapshot, without a comparable commit sequence.
|
||||
@@ -230,7 +250,7 @@ async fn read_committed_with_stats(
|
||||
let mut damaged = None;
|
||||
let mut identities = HashMap::new();
|
||||
for disk in disks {
|
||||
for (manifest_path, payload_path) in MANIFEST_PATHS.into_iter().zip(PAYLOAD_PATHS) {
|
||||
for (slot, (manifest_path, payload_path)) in MANIFEST_PATHS.into_iter().zip(PAYLOAD_PATHS).enumerate() {
|
||||
let candidate = async {
|
||||
let Some(manifest) = read_bounded_with_stats(disk, manifest_path, MANIFEST_LEN, stats.as_deref_mut()).await?
|
||||
else {
|
||||
@@ -240,7 +260,7 @@ async fn read_committed_with_stats(
|
||||
let payload = read_bounded_with_stats(disk, payload_path, header.payload_len, stats.as_deref_mut())
|
||||
.await?
|
||||
.ok_or(SnapshotError::Corrupt)?;
|
||||
CommittedSnapshot::decode(&manifest, payload, limit).map(Some)
|
||||
CommittedSnapshot::decode(slot, &manifest, payload, limit).map(Some)
|
||||
}
|
||||
.await;
|
||||
match candidate {
|
||||
@@ -272,6 +292,125 @@ async fn read_committed_with_stats(
|
||||
}
|
||||
}
|
||||
|
||||
async fn cas_replace(
|
||||
disk: &EcstoreDiskStore,
|
||||
path: &str,
|
||||
replacement: &[u8],
|
||||
limit: usize,
|
||||
) -> Result<EcstoreConditionalFileUpdate, SnapshotError> {
|
||||
let expected = read_bounded(disk, path, limit).await?.map(EcstoreDiskBytes::from);
|
||||
cas_replace_expected(disk, path, expected, replacement).await
|
||||
}
|
||||
|
||||
async fn cas_replace_expected(
|
||||
disk: &EcstoreDiskStore,
|
||||
path: &str,
|
||||
expected: Option<EcstoreDiskBytes>,
|
||||
replacement: &[u8],
|
||||
) -> Result<EcstoreConditionalFileUpdate, SnapshotError> {
|
||||
EcstoreDiskAPI::compare_and_update_file(
|
||||
disk.as_ref(),
|
||||
RUSTFS_META_BUCKET,
|
||||
path,
|
||||
expected,
|
||||
Some(EcstoreDiskBytes::copy_from_slice(replacement)),
|
||||
)
|
||||
.await
|
||||
.map_err(SnapshotError::Disk)
|
||||
}
|
||||
|
||||
fn validate_reusable_manifest_slot(existing: Option<&[u8]>, sequence: u64, payload_limit: usize) -> Result<(), SnapshotError> {
|
||||
let Some(existing) = existing else {
|
||||
return Ok(());
|
||||
};
|
||||
let manifest = Manifest::decode(existing, payload_limit)?;
|
||||
if manifest.sequence >= sequence {
|
||||
return Err(SnapshotError::Conflict);
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Publish a committed checkpoint into the inactive slot.
|
||||
///
|
||||
/// The writer is a narrow production primitive for the ownership-aware MRF
|
||||
/// handoff: it validates the whole journal payload, preserves the previous
|
||||
/// committed slot, and publishes the manifest only after the successor payload
|
||||
/// reaches the same disk. It does not delete legacy journals, tombstone older
|
||||
/// anchors, or activate the live consumer.
|
||||
pub async fn publish_committed_snapshot(
|
||||
disks: &[EcstoreDiskStore],
|
||||
owner: Uuid,
|
||||
sequence: u64,
|
||||
payload: &[u8],
|
||||
limit: usize,
|
||||
) -> Result<SnapshotPublication, SnapshotError> {
|
||||
if disks.is_empty() {
|
||||
return Err(SnapshotError::NoWritableReplica);
|
||||
}
|
||||
if owner.is_nil() || sequence == 0 || sequence == u64::MAX {
|
||||
return Err(SnapshotError::Corrupt);
|
||||
}
|
||||
if payload.len() > limit || decode_journal(payload).1 != 0 {
|
||||
return Err(SnapshotError::Corrupt);
|
||||
}
|
||||
let current = read_committed(disks, limit).await?;
|
||||
if current.as_ref().is_some_and(|snapshot| snapshot.sequence() >= sequence) {
|
||||
return Err(SnapshotError::Conflict);
|
||||
}
|
||||
let slot = current.as_ref().map_or(0, |snapshot| 1usize.saturating_sub(snapshot.slot()));
|
||||
let manifest = Manifest::encode(owner, sequence, payload)?;
|
||||
let mut payload_replicas = 0usize;
|
||||
let mut manifest_replicas = 0usize;
|
||||
let mut first_error = None;
|
||||
for disk in disks {
|
||||
let expected_manifest = match read_bounded(disk, MANIFEST_PATHS[slot], MANIFEST_LEN).await {
|
||||
Ok(expected) => expected,
|
||||
Err(error) => {
|
||||
if first_error.is_none() {
|
||||
first_error = Some(error);
|
||||
}
|
||||
continue;
|
||||
}
|
||||
};
|
||||
if let Err(error) = validate_reusable_manifest_slot(expected_manifest.as_deref(), sequence, limit) {
|
||||
if first_error.is_none() {
|
||||
first_error = Some(error);
|
||||
}
|
||||
continue;
|
||||
}
|
||||
match cas_replace(disk, PAYLOAD_PATHS[slot], payload, limit).await {
|
||||
Ok(EcstoreConditionalFileUpdate::Updated) => payload_replicas += 1,
|
||||
Ok(EcstoreConditionalFileUpdate::Missing | EcstoreConditionalFileUpdate::Mismatch) => continue,
|
||||
Err(error) => {
|
||||
if first_error.is_none() {
|
||||
first_error = Some(error);
|
||||
}
|
||||
continue;
|
||||
}
|
||||
}
|
||||
match cas_replace_expected(disk, MANIFEST_PATHS[slot], expected_manifest.map(EcstoreDiskBytes::from), &manifest).await {
|
||||
Ok(EcstoreConditionalFileUpdate::Updated) => manifest_replicas += 1,
|
||||
Ok(EcstoreConditionalFileUpdate::Missing | EcstoreConditionalFileUpdate::Mismatch) => {}
|
||||
Err(error) => {
|
||||
if first_error.is_none() {
|
||||
first_error = Some(error);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
if manifest_replicas == 0 {
|
||||
return Err(first_error.unwrap_or(SnapshotError::NoWritableReplica));
|
||||
}
|
||||
Ok(SnapshotPublication {
|
||||
owner,
|
||||
sequence,
|
||||
slot,
|
||||
payload_len: payload.len(),
|
||||
payload_replicas,
|
||||
manifest_replicas,
|
||||
})
|
||||
}
|
||||
|
||||
async fn read_legacy(disks: &[EcstoreDiskStore], path: &str, limit: usize) -> Result<Option<Vec<u8>>, SnapshotError> {
|
||||
let mut selected = None;
|
||||
let mut incomplete: Option<Vec<u8>> = None;
|
||||
@@ -337,7 +476,6 @@ async fn read_recovery_snapshot(disks: &[EcstoreDiskStore], limit: usize) -> Res
|
||||
mod tests {
|
||||
use super::*;
|
||||
use crate::heal::mrf_queue::encode_intent;
|
||||
use crate::heal::storage_api::owner::{EcstoreConditionalFileUpdate, EcstoreDiskBytes};
|
||||
use crate::heal::{DiskOption, Endpoint, new_disk};
|
||||
use rustfs_common::mrf_channel::{MrfIntent, MrfKind, MrfScope};
|
||||
use std::sync::Arc;
|
||||
@@ -359,6 +497,24 @@ mod tests {
|
||||
bytes
|
||||
}
|
||||
|
||||
fn many_record_payload(records: usize) -> Vec<u8> {
|
||||
let mut bytes = Vec::new();
|
||||
for index in 0..records {
|
||||
let intent = MrfIntent {
|
||||
bucket: Arc::from("b"),
|
||||
object: Arc::from(format!("o-{index:06}")),
|
||||
version_id: None,
|
||||
kind: MrfKind::PartialWrite,
|
||||
scope: None,
|
||||
lease: None,
|
||||
enqueued_at_ms: 1234,
|
||||
attempts: 0,
|
||||
};
|
||||
assert!(encode_intent(&intent, &mut bytes), "large fixture record must encode");
|
||||
}
|
||||
bytes
|
||||
}
|
||||
|
||||
fn manifest(owner: Uuid, sequence: u64, payload: &[u8]) -> Vec<u8> {
|
||||
let mut bytes = Vec::with_capacity(MANIFEST_LEN);
|
||||
bytes.extend_from_slice(MAGIC);
|
||||
@@ -417,7 +573,7 @@ mod tests {
|
||||
fn manifest_validates_identity_sequence_length_and_digest() {
|
||||
let bytes = payload("object");
|
||||
let owner = Uuid::new_v4();
|
||||
assert!(CommittedSnapshot::decode(&manifest(owner, 1, &bytes), bytes.clone(), bytes.len()).is_ok());
|
||||
assert!(CommittedSnapshot::decode(0, &manifest(owner, 1, &bytes), bytes.clone(), bytes.len()).is_ok());
|
||||
for (owner, sequence) in [(Uuid::nil(), 1), (owner, 0), (owner, u64::MAX)] {
|
||||
assert!(matches!(
|
||||
Manifest::decode(&manifest(owner, sequence, &bytes), bytes.len()),
|
||||
@@ -442,12 +598,12 @@ mod tests {
|
||||
let owner = Uuid::new_v4();
|
||||
let header = manifest(owner, 1, &bytes);
|
||||
assert!(matches!(
|
||||
CommittedSnapshot::decode(&header, bytes[..bytes.len() - 1].to_vec(), bytes.len()),
|
||||
CommittedSnapshot::decode(0, &header, bytes[..bytes.len() - 1].to_vec(), bytes.len()),
|
||||
Err(SnapshotError::Corrupt)
|
||||
));
|
||||
let invalid = b"not an MRF record".to_vec();
|
||||
assert!(matches!(
|
||||
CommittedSnapshot::decode(&manifest(owner, 2, &invalid), invalid, bytes.len()),
|
||||
CommittedSnapshot::decode(0, &manifest(owner, 2, &invalid), invalid, bytes.len()),
|
||||
Err(SnapshotError::Corrupt)
|
||||
));
|
||||
}
|
||||
@@ -585,6 +741,174 @@ mod tests {
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn committed_snapshot_writer_uses_inactive_slot_and_reports_replicas() {
|
||||
let root = TempDir::new().expect("test directory");
|
||||
let first = disk(&root, "first").await;
|
||||
let second = disk(&root, "second").await;
|
||||
let owner = Uuid::new_v4();
|
||||
let old = payload("old");
|
||||
let next = payload("next");
|
||||
commit(&first, 0, owner, 1, &old).await;
|
||||
commit(&second, 0, owner, 1, &old).await;
|
||||
|
||||
let publication = publish_committed_snapshot(&[first.clone(), second.clone()], owner, 2, &next, 4096)
|
||||
.await
|
||||
.expect("publish successor");
|
||||
|
||||
assert_eq!(publication.slot, 1);
|
||||
assert_eq!(publication.payload_len, next.len());
|
||||
assert_eq!(publication.payload_replicas, 2);
|
||||
assert_eq!(publication.manifest_replicas, 2);
|
||||
let recovered = read_committed(&[first.clone(), second.clone()], 4096)
|
||||
.await
|
||||
.expect("read committed")
|
||||
.expect("successor committed");
|
||||
assert_eq!(recovered.sequence(), 2);
|
||||
assert_eq!(recovered.slot(), 1);
|
||||
assert_eq!(recovered.payload(), next.as_slice());
|
||||
for disk in [&first, &second] {
|
||||
assert_eq!(
|
||||
EcstoreDiskAPI::read_all(disk.as_ref(), RUSTFS_META_BUCKET, PAYLOAD_PATHS[0])
|
||||
.await
|
||||
.expect("old payload retained")
|
||||
.as_ref(),
|
||||
old.as_slice()
|
||||
);
|
||||
assert_eq!(
|
||||
EcstoreDiskAPI::read_all(disk.as_ref(), RUSTFS_META_BUCKET, MANIFEST_PATHS[0])
|
||||
.await
|
||||
.expect("old manifest retained")
|
||||
.as_ref(),
|
||||
manifest(owner, 1, &old).as_slice()
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn committed_snapshot_writer_manifest_failure_preserves_previous_anchor() {
|
||||
let root = TempDir::new().expect("test directory");
|
||||
let store = disk(&root, "disk").await;
|
||||
let owner = Uuid::new_v4();
|
||||
let old = payload("old");
|
||||
let next = payload("next");
|
||||
commit(&store, 0, owner, 1, &old).await;
|
||||
std::fs::create_dir(root.path().join("disk").join(RUSTFS_META_BUCKET).join(MANIFEST_PATHS[1]))
|
||||
.expect("manifest path blocks successor commit");
|
||||
|
||||
let result = publish_committed_snapshot(std::slice::from_ref(&store), owner, 2, &next, 4096).await;
|
||||
|
||||
assert!(
|
||||
matches!(
|
||||
result,
|
||||
Err(SnapshotError::Disk(_) | SnapshotError::Read(_) | SnapshotError::NoWritableReplica)
|
||||
),
|
||||
"manifest failure must be visible: {result:?}"
|
||||
);
|
||||
let reopened = disk(&root, "disk").await;
|
||||
let recovered = read_committed(std::slice::from_ref(&reopened), 4096)
|
||||
.await
|
||||
.expect("read previous committed snapshot")
|
||||
.expect("old anchor remains committed");
|
||||
assert_eq!(recovered.sequence(), 1);
|
||||
assert_eq!(recovered.slot(), 0);
|
||||
assert_eq!(recovered.payload(), old.as_slice());
|
||||
assert_eq!(
|
||||
EcstoreDiskAPI::read_all(reopened.as_ref(), RUSTFS_META_BUCKET, PAYLOAD_PATHS[0])
|
||||
.await
|
||||
.expect("old payload retained")
|
||||
.as_ref(),
|
||||
old.as_slice()
|
||||
);
|
||||
assert_eq!(
|
||||
EcstoreDiskAPI::read_all(reopened.as_ref(), RUSTFS_META_BUCKET, MANIFEST_PATHS[0])
|
||||
.await
|
||||
.expect("old manifest retained")
|
||||
.as_ref(),
|
||||
manifest(owner, 1, &old).as_slice()
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn committed_snapshot_writer_does_not_overwrite_damaged_inactive_manifest() {
|
||||
let root = TempDir::new().expect("test directory");
|
||||
let store = disk(&root, "disk").await;
|
||||
let owner = Uuid::new_v4();
|
||||
let old = payload("old");
|
||||
let next = payload("next");
|
||||
let damaged = b"damaged successor manifest".to_vec();
|
||||
commit(&store, 0, owner, 1, &old).await;
|
||||
EcstoreDiskAPI::write_all(
|
||||
store.as_ref(),
|
||||
RUSTFS_META_BUCKET,
|
||||
MANIFEST_PATHS[1],
|
||||
EcstoreDiskBytes::copy_from_slice(&damaged),
|
||||
)
|
||||
.await
|
||||
.expect("damaged inactive manifest fixture");
|
||||
|
||||
let result = publish_committed_snapshot(std::slice::from_ref(&store), owner, 2, &next, 4096).await;
|
||||
|
||||
assert!(
|
||||
matches!(result, Err(SnapshotError::Corrupt)),
|
||||
"damaged manifest must fail closed: {result:?}"
|
||||
);
|
||||
let reopened = disk(&root, "disk").await;
|
||||
let recovered = read_committed(std::slice::from_ref(&reopened), 4096)
|
||||
.await
|
||||
.expect("read previous committed snapshot")
|
||||
.expect("old anchor remains committed");
|
||||
assert_eq!(recovered.sequence(), 1);
|
||||
assert_eq!(recovered.payload(), old.as_slice());
|
||||
assert_eq!(
|
||||
EcstoreDiskAPI::read_all(reopened.as_ref(), RUSTFS_META_BUCKET, MANIFEST_PATHS[1])
|
||||
.await
|
||||
.expect("damaged manifest retained")
|
||||
.as_ref(),
|
||||
damaged.as_slice()
|
||||
);
|
||||
assert!(
|
||||
matches!(
|
||||
EcstoreDiskAPI::read_all(reopened.as_ref(), RUSTFS_META_BUCKET, PAYLOAD_PATHS[1]).await,
|
||||
Err(EcstoreDiskError::FileNotFound | EcstoreDiskError::VolumeNotFound)
|
||||
),
|
||||
"successor payload must not be written before manifest slot is reusable"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn committed_snapshot_writer_publishes_100k_records_with_bounded_readback() {
|
||||
let root = TempDir::new().expect("test directory");
|
||||
let first = disk(&root, "first").await;
|
||||
let second = disk(&root, "second").await;
|
||||
let owner = Uuid::new_v4();
|
||||
let records = 100_000usize;
|
||||
let bytes = many_record_payload(records);
|
||||
let limit = rustfs_config::DEFAULT_HEAL_MRF_JOURNAL_MAX_BYTES;
|
||||
assert!(bytes.len() < limit, "100k compact MRF records must fit the configured journal limit");
|
||||
|
||||
let publication = publish_committed_snapshot(&[first.clone(), second.clone()], owner, 1, &bytes, limit)
|
||||
.await
|
||||
.expect("publish 100k-record successor");
|
||||
|
||||
assert_eq!(publication.payload_replicas, 2);
|
||||
assert_eq!(publication.manifest_replicas, 2);
|
||||
assert_eq!(publication.payload_len, bytes.len());
|
||||
let mut stats = SnapshotReadStats::default();
|
||||
let recovered = read_committed_with_stats(&[first, second], limit, Some(&mut stats))
|
||||
.await
|
||||
.expect("read committed large snapshot")
|
||||
.expect("large snapshot committed");
|
||||
assert_eq!(recovered.sequence(), 1);
|
||||
assert_eq!(recovered.payload().len(), bytes.len());
|
||||
let (decoded, truncated) = decode_journal(recovered.payload());
|
||||
assert_eq!(truncated, 0);
|
||||
assert_eq!(decoded.len(), records);
|
||||
assert_eq!(stats.file_reads, 4, "two manifest and two payload files should be read");
|
||||
assert_eq!(stats.bytes_read, (MANIFEST_LEN * 2) + (bytes.len() * 2));
|
||||
assert_eq!(stats.peak_file_bytes, bytes.len().max(MANIFEST_LEN));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn stale_manifest_cas_cannot_replace_committed_anchor() {
|
||||
let root = TempDir::new().expect("test directory");
|
||||
@@ -609,6 +933,78 @@ mod tests {
|
||||
assert_eq!(recovered.manifest.sequence, 1);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn manifest_cas_publication_transitions_from_legacy_without_losing_anchor() {
|
||||
let root = TempDir::new().expect("test directory");
|
||||
let store = disk(&root, "disk").await;
|
||||
let owner = Uuid::new_v4();
|
||||
let legacy = payload("legacy");
|
||||
let committed = payload("committed");
|
||||
let successor = payload("successor");
|
||||
|
||||
EcstoreDiskAPI::write_all(store.as_ref(), RUSTFS_META_BUCKET, MRF_SCOPED_JOURNAL_PATH, legacy.clone().into())
|
||||
.await
|
||||
.expect("legacy fixture");
|
||||
assert!(
|
||||
matches!(read_recovery_snapshot(std::slice::from_ref(&store), 4096).await.expect("legacy read"), Some(RecoverySnapshot::Legacy(data)) if data == legacy),
|
||||
"complete legacy journal remains the fallback before committed publication"
|
||||
);
|
||||
|
||||
install(&store, PAYLOAD_PATHS[0], &committed).await;
|
||||
assert!(
|
||||
matches!(read_recovery_snapshot(std::slice::from_ref(&store), 4096).await.expect("payload-only read"), Some(RecoverySnapshot::Legacy(data)) if data == legacy),
|
||||
"payload-only successor is not a committed snapshot"
|
||||
);
|
||||
|
||||
install(&store, MANIFEST_PATHS[0], &manifest(owner, 1, &committed)).await;
|
||||
let recovered = read_recovery_snapshot(std::slice::from_ref(&store), 4096)
|
||||
.await
|
||||
.expect("committed read")
|
||||
.expect("committed snapshot");
|
||||
assert!(
|
||||
matches!(recovered, RecoverySnapshot::Committed(snapshot) if snapshot.sequence() == 1 && snapshot.payload() == committed),
|
||||
"manifest CAS completion promotes the committed snapshot above legacy"
|
||||
);
|
||||
|
||||
let stale_manifest = manifest(owner, 2, &successor);
|
||||
let result = EcstoreDiskAPI::compare_and_update_file(
|
||||
store.as_ref(),
|
||||
RUSTFS_META_BUCKET,
|
||||
MANIFEST_PATHS[0],
|
||||
None,
|
||||
Some(EcstoreDiskBytes::copy_from_slice(&stale_manifest)),
|
||||
)
|
||||
.await
|
||||
.expect("stale CAS call");
|
||||
assert_eq!(result, EcstoreConditionalFileUpdate::Mismatch);
|
||||
install(&store, PAYLOAD_PATHS[1], &successor).await;
|
||||
install(&store, MANIFEST_PATHS[1], &stale_manifest[..20]).await;
|
||||
|
||||
let reopened = disk(&root, "disk").await;
|
||||
let recovered = read_recovery_snapshot(std::slice::from_ref(&reopened), 4096)
|
||||
.await
|
||||
.expect("committed anchor after stale successor")
|
||||
.expect("committed snapshot");
|
||||
assert!(
|
||||
matches!(recovered, RecoverySnapshot::Committed(snapshot) if snapshot.sequence() == 1 && snapshot.payload() == committed),
|
||||
"failed or torn successor publication must not fall back to legacy"
|
||||
);
|
||||
assert_eq!(
|
||||
EcstoreDiskAPI::read_all(reopened.as_ref(), RUSTFS_META_BUCKET, MANIFEST_PATHS[0])
|
||||
.await
|
||||
.expect("old manifest retained")
|
||||
.as_ref(),
|
||||
manifest(owner, 1, &committed)
|
||||
);
|
||||
assert_eq!(
|
||||
EcstoreDiskAPI::read_all(reopened.as_ref(), RUSTFS_META_BUCKET, MRF_SCOPED_JOURNAL_PATH)
|
||||
.await
|
||||
.expect("legacy bytes retained")
|
||||
.as_ref(),
|
||||
legacy
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn legacy_import_requires_complete_consistent_replicas() {
|
||||
let root = TempDir::new().expect("test directory");
|
||||
|
||||
@@ -106,7 +106,7 @@ const EVENT_HEAL_ERASURE_SET_STAGE: &str = "heal_erasure_set_stage";
|
||||
const EVENT_HEAL_ERASURE_SET_RESULT: &str = "heal_erasure_set_result";
|
||||
|
||||
/// Heal type
|
||||
#[derive(Debug, Clone)]
|
||||
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||
pub enum HealType {
|
||||
/// Cluster heal
|
||||
Cluster,
|
||||
@@ -209,7 +209,7 @@ impl HealPriority {
|
||||
}
|
||||
|
||||
/// Heal options
|
||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
|
||||
pub struct HealOptions {
|
||||
/// Scan mode
|
||||
pub scan_mode: HealScanMode,
|
||||
|
||||
@@ -263,6 +263,14 @@ pub struct ScannerUsageRecoveryIntentResponse {
|
||||
pub mode: String,
|
||||
pub intent_id: String,
|
||||
pub state: String,
|
||||
#[serde(default)]
|
||||
pub actor_sha256: Option<String>,
|
||||
#[serde(default)]
|
||||
pub idempotency_key_sha256: Option<String>,
|
||||
#[serde(default)]
|
||||
pub request_sha256: Option<String>,
|
||||
#[serde(default)]
|
||||
pub accepted_at_unix_secs: Option<u64>,
|
||||
#[serde(flatten)]
|
||||
pub extra: serde_json::Map<String, serde_json::Value>,
|
||||
}
|
||||
@@ -869,6 +877,10 @@ mod tests {
|
||||
"mode": "full-rebuild",
|
||||
"intent_id": "0123456789abcdef0123456789abcdef0123456789abcdef0123456789abcdef",
|
||||
"state": "accepted",
|
||||
"actor_sha256": "1111111111111111111111111111111111111111111111111111111111111111",
|
||||
"idempotency_key_sha256": "2222222222222222222222222222222222222222222222222222222222222222",
|
||||
"request_sha256": "3333333333333333333333333333333333333333333333333333333333333333",
|
||||
"accepted_at_unix_secs": 7,
|
||||
"future": {"worker": "pending"}
|
||||
}))
|
||||
.unwrap();
|
||||
@@ -877,7 +889,33 @@ mod tests {
|
||||
assert_eq!(intent.mode, "full-rebuild");
|
||||
assert_eq!(intent.intent_id, "0123456789abcdef0123456789abcdef0123456789abcdef0123456789abcdef");
|
||||
assert_eq!(intent.state, "accepted");
|
||||
assert_eq!(
|
||||
intent.actor_sha256.as_deref(),
|
||||
Some("1111111111111111111111111111111111111111111111111111111111111111")
|
||||
);
|
||||
assert_eq!(
|
||||
intent.idempotency_key_sha256.as_deref(),
|
||||
Some("2222222222222222222222222222222222222222222222222222222222222222")
|
||||
);
|
||||
assert_eq!(
|
||||
intent.request_sha256.as_deref(),
|
||||
Some("3333333333333333333333333333333333333333333333333333333333333333")
|
||||
);
|
||||
assert_eq!(intent.accepted_at_unix_secs, Some(7));
|
||||
assert_eq!(intent.extra["future"]["worker"], "pending");
|
||||
|
||||
let legacy_intent: ScannerUsageRecoveryIntentResponse = serde_json::from_value(json!({
|
||||
"status": "accepted",
|
||||
"action": "usage-full-rebuild",
|
||||
"mode": "full-rebuild",
|
||||
"intent_id": "0123456789abcdef0123456789abcdef0123456789abcdef0123456789abcdef",
|
||||
"state": "accepted"
|
||||
}))
|
||||
.unwrap();
|
||||
assert!(legacy_intent.actor_sha256.is_none());
|
||||
assert!(legacy_intent.idempotency_key_sha256.is_none());
|
||||
assert!(legacy_intent.request_sha256.is_none());
|
||||
assert!(legacy_intent.accepted_at_unix_secs.is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
@@ -971,7 +1009,7 @@ mod tests {
|
||||
#[tokio::test]
|
||||
async fn scanner_usage_async_reset_posts_explicit_intent_contract() {
|
||||
let server = TestServer::spawn(
|
||||
r#"{"status":"accepted","action":"usage-full-rebuild","mode":"full-rebuild","intent_id":"aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","state":"accepted"}"#,
|
||||
r#"{"status":"accepted","action":"usage-full-rebuild","mode":"full-rebuild","intent_id":"aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","state":"accepted","actor_sha256":"1111111111111111111111111111111111111111111111111111111111111111","idempotency_key_sha256":"2222222222222222222222222222222222222222222222222222222222222222","request_sha256":"3333333333333333333333333333333333333333333333333333333333333333","accepted_at_unix_secs":7}"#,
|
||||
202,
|
||||
)
|
||||
.await;
|
||||
@@ -985,6 +1023,7 @@ mod tests {
|
||||
assert_eq!(accepted.status, "accepted");
|
||||
assert_eq!(accepted.mode, "full-rebuild");
|
||||
assert_eq!(accepted.state, "accepted");
|
||||
assert_eq!(accepted.accepted_at_unix_secs, Some(7));
|
||||
let request = server.recorded();
|
||||
assert_eq!(request.method, "POST");
|
||||
assert_eq!(request.path, "/rustfs/admin/v3/scanner/usage-state/reset");
|
||||
|
||||
@@ -97,9 +97,11 @@ pub use scanner::{
|
||||
pub use scanner_io::{
|
||||
ScannerDirtyUsageAckError, ScannerDirtyUsageBucket, ScannerDirtyUsageSnapshot, ScannerDirtyUsageState,
|
||||
acknowledge_dirty_usage_generation, acknowledge_scoped_dirty_usage, clear_dirty_usage_bucket, record_dirty_usage_bucket,
|
||||
record_dirty_usage_object, record_scanner_maintenance_change, scanner_activity_epoch, scanner_dirty_usage_snapshot,
|
||||
scanner_dirty_usage_state, scanner_maintenance_generation,
|
||||
record_dirty_usage_bucket_from_producer, record_dirty_usage_object, record_dirty_usage_object_from_producer,
|
||||
record_scanner_maintenance_change, scanner_activity_epoch, scanner_dirty_usage_snapshot, scanner_dirty_usage_state,
|
||||
scanner_maintenance_generation,
|
||||
};
|
||||
pub use segment_invalidation::SegmentInvalidationProducerIdentity;
|
||||
pub use sleeper::{DynamicSleeper, SCANNER_IDLE_MODE, SCANNER_SLEEPER};
|
||||
use std::sync::atomic::{AtomicU64, Ordering};
|
||||
pub use storage_api::ScannerReplicationConfig as ReplicationConfig;
|
||||
|
||||
@@ -3,7 +3,8 @@
|
||||
use super::*;
|
||||
use crate::segment_invalidation::{
|
||||
MAX_SEGMENT_INVALIDATION_BYTES, MAX_SEGMENT_INVALIDATION_ENTRIES, SegmentInvalidationDomain, SegmentInvalidationEnvelope,
|
||||
SegmentInvalidationError, SegmentInvalidationProducer, SegmentInvalidationProof, admit_segment_invalidation,
|
||||
SegmentInvalidationError, SegmentInvalidationProducer, SegmentInvalidationProducerIdentity, SegmentInvalidationProof,
|
||||
admit_segment_invalidation, complete_segment_invalidation_producers,
|
||||
};
|
||||
use std::collections::BTreeSet;
|
||||
|
||||
@@ -11,7 +12,8 @@ const MAX_WALK_SAMPLES: usize = 32;
|
||||
const MAX_WALK_BYTES: usize = 1024;
|
||||
|
||||
fn segment_producers() -> BTreeSet<SegmentInvalidationProducer> {
|
||||
SegmentInvalidationProducer::REQUIRED.into_iter().collect()
|
||||
complete_segment_invalidation_producers(SegmentInvalidationProducerIdentity::REQUIRED_PRODUCTION)
|
||||
.expect("fixture should enumerate the complete production producer matrix")
|
||||
}
|
||||
|
||||
fn segment_envelope() -> SegmentInvalidationEnvelope {
|
||||
|
||||
@@ -39,7 +39,7 @@ use s3s::dto::{
|
||||
BucketLifecycleConfiguration, ObjectLockConfiguration, ObjectLockEnabled, ReplicationConfiguration, VersioningConfiguration,
|
||||
};
|
||||
use sha2::{Digest as _, Sha256};
|
||||
use std::collections::{HashMap, HashSet};
|
||||
use std::collections::{BTreeSet, HashMap, HashSet};
|
||||
use std::future::Future;
|
||||
use std::path::Path;
|
||||
use std::pin::Pin;
|
||||
@@ -1224,8 +1224,9 @@ pub(crate) use cache::{
|
||||
pub use dirty_usage::{
|
||||
ScannerDirtyUsageAckError, ScannerDirtyUsageBucket, ScannerDirtyUsageSnapshot, ScannerDirtyUsageState,
|
||||
acknowledge_dirty_usage_generation, acknowledge_scoped_dirty_usage, clear_dirty_usage_bucket, record_dirty_usage_bucket,
|
||||
record_dirty_usage_object, record_scanner_maintenance_change, scanner_activity_epoch, scanner_dirty_usage_snapshot,
|
||||
scanner_dirty_usage_state, scanner_maintenance_generation,
|
||||
record_dirty_usage_bucket_from_producer, record_dirty_usage_object, record_dirty_usage_object_from_producer,
|
||||
record_scanner_maintenance_change, scanner_activity_epoch, scanner_dirty_usage_snapshot, scanner_dirty_usage_state,
|
||||
scanner_maintenance_generation,
|
||||
};
|
||||
#[cfg(test)]
|
||||
pub(crate) use dirty_usage::{clear_dirty_usage_buckets_for_tests, dirty_usage_buckets_for_tests};
|
||||
|
||||
@@ -22,6 +22,11 @@ pub(super) static DIRTY_USAGE_BUCKETS: LazyLock<StdMutex<DirtyUsageBuckets>> = L
|
||||
// matching scope.
|
||||
pub(super) static DIRTY_USAGE_BUCKET_SCOPES: LazyLock<StdMutex<DirtyUsageBucketScopes>> =
|
||||
LazyLock::new(|| StdMutex::new(HashMap::new()));
|
||||
// Non-authoritative process-local producer coverage. Any future segment reuse
|
||||
// activation must bind this to the exact generation window and durable proof.
|
||||
pub(super) static DIRTY_USAGE_PRODUCER_IDENTITIES: LazyLock<
|
||||
StdMutex<BTreeSet<crate::segment_invalidation::SegmentInvalidationProducerIdentity>>,
|
||||
> = LazyLock::new(|| StdMutex::new(BTreeSet::new()));
|
||||
pub(super) static DIRTY_USAGE_BUCKET_NOTIFY: LazyLock<Notify> = LazyLock::new(Notify::new);
|
||||
pub(super) static SCANNER_ACTIVITY_EPOCH: LazyLock<String> = LazyLock::new(|| format!("{:032x}", rand::random::<u128>()));
|
||||
pub(super) static SCANNER_MAINTENANCE_GENERATION: AtomicU64 = AtomicU64::new(0);
|
||||
@@ -153,6 +158,7 @@ fn apply_scoped_dirty_usage_ack(
|
||||
#[cfg(test)]
|
||||
mod scoped_dirty_usage_tests {
|
||||
use super::*;
|
||||
use crate::segment_invalidation::SegmentInvalidationProducerIdentity;
|
||||
|
||||
#[test]
|
||||
fn scoped_dirty_usage_preserves_uncovered_newer_and_replayed_generations() {
|
||||
@@ -213,6 +219,30 @@ mod scoped_dirty_usage_tests {
|
||||
assert_eq!(scopes, original_scopes);
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn dirty_usage_tracks_known_segment_producer_identities_without_authorizing_unknown_sources() {
|
||||
clear_dirty_usage_buckets_for_tests();
|
||||
record_dirty_usage_object_from_producer("photos", "hot/object", SegmentInvalidationProducerIdentity::PutObject);
|
||||
record_dirty_usage_object_from_producer("photos", "archive/object", SegmentInvalidationProducerIdentity::DeleteObject);
|
||||
record_dirty_usage_bucket_from_producer("photos", SegmentInvalidationProducerIdentity::Unknown);
|
||||
|
||||
assert_eq!(
|
||||
dirty_usage_producer_identities_for_tests(),
|
||||
BTreeSet::from([
|
||||
SegmentInvalidationProducerIdentity::PutObject,
|
||||
SegmentInvalidationProducerIdentity::DeleteObject
|
||||
])
|
||||
);
|
||||
assert_eq!(
|
||||
dirty_usage_bucket_scopes_for_tests().get("photos"),
|
||||
Some(&DirtyUsageBucketScope::WholeBucket),
|
||||
"an unknown producer keeps the bucket dirty but must not count as producer coverage"
|
||||
);
|
||||
|
||||
clear_dirty_usage_buckets_for_tests();
|
||||
assert!(dirty_usage_producer_identities_for_tests().is_empty());
|
||||
}
|
||||
}
|
||||
|
||||
pub(super) fn dirty_usage_buckets() -> MutexGuard<'static, DirtyUsageBuckets> {
|
||||
@@ -225,6 +255,13 @@ fn dirty_usage_bucket_scopes() -> MutexGuard<'static, DirtyUsageBucketScopes> {
|
||||
.unwrap_or_else(|poisoned| poisoned.into_inner())
|
||||
}
|
||||
|
||||
fn dirty_usage_producer_identities()
|
||||
-> MutexGuard<'static, BTreeSet<crate::segment_invalidation::SegmentInvalidationProducerIdentity>> {
|
||||
DIRTY_USAGE_PRODUCER_IDENTITIES
|
||||
.lock()
|
||||
.unwrap_or_else(|poisoned| poisoned.into_inner())
|
||||
}
|
||||
|
||||
pub(super) fn usize_to_u64_saturated(value: usize) -> u64 {
|
||||
u64::try_from(value).unwrap_or(u64::MAX)
|
||||
}
|
||||
@@ -240,6 +277,22 @@ pub fn record_dirty_usage_bucket(bucket: &str) {
|
||||
return;
|
||||
}
|
||||
|
||||
record_dirty_usage_bucket_inner(bucket);
|
||||
}
|
||||
|
||||
pub fn record_dirty_usage_bucket_from_producer(
|
||||
bucket: &str,
|
||||
producer: crate::segment_invalidation::SegmentInvalidationProducerIdentity,
|
||||
) {
|
||||
if bucket.is_empty() {
|
||||
return;
|
||||
}
|
||||
|
||||
record_segment_invalidation_producer_identity(producer);
|
||||
record_dirty_usage_bucket_inner(bucket);
|
||||
}
|
||||
|
||||
fn record_dirty_usage_bucket_inner(bucket: &str) {
|
||||
let pending_buckets = {
|
||||
let mut dirty_buckets = dirty_usage_buckets();
|
||||
let mut dirty_scopes = dirty_usage_bucket_scopes();
|
||||
@@ -263,6 +316,23 @@ pub fn record_dirty_usage_bucket(bucket: &str) {
|
||||
/// local: after restart or any unverified distributed path the scanner falls
|
||||
/// back to its ordinary bucket scan.
|
||||
pub fn record_dirty_usage_object(bucket: &str, object: &str) {
|
||||
record_dirty_usage_object_inner(bucket, object);
|
||||
}
|
||||
|
||||
pub fn record_dirty_usage_object_from_producer(
|
||||
bucket: &str,
|
||||
object: &str,
|
||||
producer: crate::segment_invalidation::SegmentInvalidationProducerIdentity,
|
||||
) {
|
||||
if bucket.is_empty() {
|
||||
return;
|
||||
}
|
||||
|
||||
record_segment_invalidation_producer_identity(producer);
|
||||
record_dirty_usage_object_inner(bucket, object);
|
||||
}
|
||||
|
||||
fn record_dirty_usage_object_inner(bucket: &str, object: &str) {
|
||||
let Some(top_level_entry) = dirty_usage_top_level_entry(object) else {
|
||||
record_dirty_usage_bucket(bucket);
|
||||
return;
|
||||
@@ -296,6 +366,17 @@ pub fn record_dirty_usage_object(bucket: &str, object: &str) {
|
||||
DIRTY_USAGE_BUCKET_NOTIFY.notify_one();
|
||||
}
|
||||
|
||||
fn record_segment_invalidation_producer_identity(producer: crate::segment_invalidation::SegmentInvalidationProducerIdentity) {
|
||||
if producer.producer().is_some() {
|
||||
dirty_usage_producer_identities().insert(producer);
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
fn dirty_usage_producer_identities_for_tests() -> BTreeSet<crate::segment_invalidation::SegmentInvalidationProducerIdentity> {
|
||||
dirty_usage_producer_identities().clone()
|
||||
}
|
||||
|
||||
fn dirty_usage_top_level_entry(object: &str) -> Option<String> {
|
||||
let (top_level_entry, _) = object.split_once('/').unwrap_or((object, ""));
|
||||
(!top_level_entry.is_empty()
|
||||
@@ -577,6 +658,7 @@ pub(super) fn dirty_usage_bucket_count() -> usize {
|
||||
pub(crate) fn clear_dirty_usage_buckets_for_tests() {
|
||||
dirty_usage_buckets().clear();
|
||||
dirty_usage_bucket_scopes().clear();
|
||||
dirty_usage_producer_identities().clear();
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
|
||||
@@ -16,23 +16,29 @@ use super::*;
|
||||
use crate::data_usage_define::{DATA_USAGE_OBJ_NAME_PATH, read_config_with_revision};
|
||||
|
||||
async fn create_cohort_bucket(store: &ECStore, bucket: &str) {
|
||||
create_cohort_bucket_objects(store, bucket, 1).await;
|
||||
}
|
||||
|
||||
async fn create_cohort_bucket_objects(store: &ECStore, bucket: &str, objects: usize) {
|
||||
store
|
||||
.make_bucket(bucket, &MakeBucketOptions::default())
|
||||
.await
|
||||
.expect("fixture bucket");
|
||||
for set in store.all_set_disks() {
|
||||
let mut reader = ScannerPutObjReader::from_vec(b"cohort".to_vec());
|
||||
set.put_object(
|
||||
bucket,
|
||||
"initial",
|
||||
&mut reader,
|
||||
&ScannerObjectOptions {
|
||||
no_lock: true,
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await
|
||||
.expect("fixture object and all rename tails should persist");
|
||||
for index in 0..objects {
|
||||
let mut reader = ScannerPutObjReader::from_vec(b"cohort".to_vec());
|
||||
set.put_object(
|
||||
bucket,
|
||||
&format!("object-{index:04}"),
|
||||
&mut reader,
|
||||
&ScannerObjectOptions {
|
||||
no_lock: true,
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await
|
||||
.expect("fixture object and all rename tails should persist");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -168,6 +174,92 @@ async fn service_cohort_production_dispatch_services_waiters_across_sources() {
|
||||
clear_dirty_usage_buckets_for_tests();
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn service_cohort_flat_bucket_budget_does_not_publish_unscanned_small_bucket() {
|
||||
let (_dir, store) = setup_two_pool_scanner_store().await;
|
||||
clear_dirty_usage_buckets_for_tests();
|
||||
let flat = format!("a-flat-{}", Uuid::new_v4().simple());
|
||||
let small = format!("z-small-{}", Uuid::new_v4().simple());
|
||||
create_cohort_bucket_objects(&store, &flat, 6).await;
|
||||
create_cohort_bucket(&store, &small).await;
|
||||
let cohort = Arc::new(StdMutex::new(ScannerServiceCohort::default()));
|
||||
|
||||
let expected_flat = store
|
||||
.all_set_disks()
|
||||
.iter()
|
||||
.map(|set| (DataUsageCacheSource::new(set.pool_index, set.set_index), flat.clone()))
|
||||
.collect::<HashSet<_>>();
|
||||
let expected_small = store
|
||||
.all_set_disks()
|
||||
.iter()
|
||||
.map(|set| (DataUsageCacheSource::new(set.pool_index, set.set_index), small.clone()))
|
||||
.collect::<HashSet<_>>();
|
||||
|
||||
let ctx = CancellationToken::new();
|
||||
let budget = ScannerCycleBudget::new_with_progress_tracking(
|
||||
&ctx,
|
||||
ScannerCycleBudgetConfig {
|
||||
max_objects: Some(1),
|
||||
..Default::default()
|
||||
},
|
||||
);
|
||||
let (result, usage) = run_cohort_cycle(&store, cohort.clone(), 1, budget.clone()).await;
|
||||
assert_eq!(result.status, ScannerCycleStatus::Incomplete);
|
||||
assert!(usage.is_none(), "wide-bucket budget exhaustion must not publish a partial aggregate");
|
||||
assert!(budget.budget_elapsed());
|
||||
let first_round_admitted = cohort
|
||||
.lock()
|
||||
.expect("cohort lock")
|
||||
.admitted_members()
|
||||
.into_iter()
|
||||
.collect::<HashSet<_>>();
|
||||
assert!(
|
||||
!first_round_admitted.is_disjoint(&expected_flat),
|
||||
"the first fixed budget round should exercise the wide flat bucket"
|
||||
);
|
||||
assert!(
|
||||
first_round_admitted.is_disjoint(&expected_small),
|
||||
"a small bucket not yet reached by the real scanner must not be marked admitted"
|
||||
);
|
||||
|
||||
for cycle in 2..=4 {
|
||||
let ctx = CancellationToken::new();
|
||||
let (result, usage) = run_cohort_cycle(
|
||||
&store,
|
||||
cohort.clone(),
|
||||
cycle,
|
||||
ScannerCycleBudget::new_with_progress_tracking(
|
||||
&ctx,
|
||||
ScannerCycleBudgetConfig {
|
||||
max_objects: Some(1),
|
||||
..Default::default()
|
||||
},
|
||||
),
|
||||
)
|
||||
.await;
|
||||
assert_eq!(result.status, ScannerCycleStatus::Incomplete);
|
||||
assert!(usage.is_none(), "mixed partial coverage still cannot publish the set root");
|
||||
}
|
||||
let admitted_after_budgeted_rounds = cohort
|
||||
.lock()
|
||||
.expect("cohort lock")
|
||||
.admitted_members()
|
||||
.into_iter()
|
||||
.collect::<HashSet<_>>();
|
||||
assert!(
|
||||
expected_small.is_subset(&admitted_after_budgeted_rounds),
|
||||
"tracked small buckets must receive real execution opportunities within their fixed service-round bound"
|
||||
);
|
||||
|
||||
let ctx = CancellationToken::new();
|
||||
let (result, usage) =
|
||||
run_cohort_cycle(&store, cohort, 5, ScannerCycleBudget::new(&ctx, ScannerCycleBudgetConfig::default())).await;
|
||||
assert_eq!(result.status, ScannerCycleStatus::Complete);
|
||||
assert_eq!(usage.expect("final complete aggregate").objects_total_count, 14);
|
||||
clear_dirty_usage_buckets_for_tests();
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn service_cohort_fresh_complete_aggregate_preserves_reordered_sources() {
|
||||
|
||||
@@ -25,6 +25,7 @@ pub enum SegmentInvalidationError {
|
||||
ByteLimit,
|
||||
InvalidProof,
|
||||
InvalidKey,
|
||||
UnknownProducer,
|
||||
}
|
||||
|
||||
#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord)]
|
||||
@@ -50,6 +51,74 @@ impl SegmentInvalidationProducer {
|
||||
];
|
||||
}
|
||||
|
||||
#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord)]
|
||||
pub enum SegmentInvalidationProducerIdentity {
|
||||
PutObject,
|
||||
DeleteObject,
|
||||
DeleteMarker,
|
||||
CompleteMultipartUpload,
|
||||
Replication,
|
||||
TierTransition,
|
||||
TierExpiration,
|
||||
DirectoryObject,
|
||||
Unknown,
|
||||
TestFixture,
|
||||
}
|
||||
|
||||
impl SegmentInvalidationProducerIdentity {
|
||||
pub const REQUIRED_PRODUCTION: [Self; 8] = [
|
||||
Self::PutObject,
|
||||
Self::DeleteObject,
|
||||
Self::DeleteMarker,
|
||||
Self::CompleteMultipartUpload,
|
||||
Self::Replication,
|
||||
Self::TierTransition,
|
||||
Self::TierExpiration,
|
||||
Self::DirectoryObject,
|
||||
];
|
||||
|
||||
pub fn producer(self) -> Option<SegmentInvalidationProducer> {
|
||||
match self {
|
||||
Self::PutObject => Some(SegmentInvalidationProducer::Put),
|
||||
Self::DeleteObject => Some(SegmentInvalidationProducer::Delete),
|
||||
Self::DeleteMarker => Some(SegmentInvalidationProducer::DeleteMarker),
|
||||
Self::CompleteMultipartUpload => Some(SegmentInvalidationProducer::Multipart),
|
||||
Self::Replication => Some(SegmentInvalidationProducer::Replication),
|
||||
Self::TierTransition | Self::TierExpiration => Some(SegmentInvalidationProducer::Tier),
|
||||
Self::DirectoryObject => Some(SegmentInvalidationProducer::DirectoryObject),
|
||||
Self::Unknown | Self::TestFixture => None,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
pub fn complete_segment_invalidation_producers<I>(
|
||||
identities: I,
|
||||
) -> Result<BTreeSet<SegmentInvalidationProducer>, SegmentInvalidationError>
|
||||
where
|
||||
I: IntoIterator<Item = SegmentInvalidationProducerIdentity>,
|
||||
{
|
||||
let mut covered_identities = BTreeSet::new();
|
||||
let mut producers = BTreeSet::new();
|
||||
for identity in identities {
|
||||
let Some(producer) = identity.producer() else {
|
||||
return Err(SegmentInvalidationError::UnknownProducer);
|
||||
};
|
||||
covered_identities.insert(identity);
|
||||
producers.insert(producer);
|
||||
}
|
||||
if SegmentInvalidationProducerIdentity::REQUIRED_PRODUCTION
|
||||
.iter()
|
||||
.all(|identity| covered_identities.contains(identity))
|
||||
&& SegmentInvalidationProducer::REQUIRED
|
||||
.iter()
|
||||
.all(|producer| producers.contains(producer))
|
||||
{
|
||||
Ok(producers)
|
||||
} else {
|
||||
Err(SegmentInvalidationError::InvalidProof)
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
|
||||
pub enum SegmentInvalidationDomain {
|
||||
LocalSingleSet,
|
||||
@@ -170,7 +239,8 @@ mod tests {
|
||||
use super::*;
|
||||
|
||||
fn producers() -> BTreeSet<SegmentInvalidationProducer> {
|
||||
SegmentInvalidationProducer::REQUIRED.into_iter().collect()
|
||||
complete_segment_invalidation_producers(SegmentInvalidationProducerIdentity::REQUIRED_PRODUCTION)
|
||||
.expect("production producer matrix should be complete")
|
||||
}
|
||||
|
||||
fn envelope() -> SegmentInvalidationEnvelope {
|
||||
@@ -320,6 +390,63 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn segment_invalidation_producer_identities_must_be_known_and_complete() {
|
||||
assert_eq!(
|
||||
complete_segment_invalidation_producers(SegmentInvalidationProducerIdentity::REQUIRED_PRODUCTION),
|
||||
Ok(SegmentInvalidationProducer::REQUIRED.into_iter().collect())
|
||||
);
|
||||
assert_eq!(
|
||||
complete_segment_invalidation_producers([
|
||||
SegmentInvalidationProducerIdentity::PutObject,
|
||||
SegmentInvalidationProducerIdentity::DeleteObject,
|
||||
SegmentInvalidationProducerIdentity::DeleteMarker,
|
||||
SegmentInvalidationProducerIdentity::CompleteMultipartUpload,
|
||||
SegmentInvalidationProducerIdentity::Replication,
|
||||
SegmentInvalidationProducerIdentity::TierTransition,
|
||||
SegmentInvalidationProducerIdentity::DirectoryObject,
|
||||
SegmentInvalidationProducerIdentity::Unknown,
|
||||
]),
|
||||
Err(SegmentInvalidationError::UnknownProducer)
|
||||
);
|
||||
assert_eq!(
|
||||
complete_segment_invalidation_producers([
|
||||
SegmentInvalidationProducerIdentity::PutObject,
|
||||
SegmentInvalidationProducerIdentity::DeleteObject,
|
||||
SegmentInvalidationProducerIdentity::DeleteMarker,
|
||||
SegmentInvalidationProducerIdentity::CompleteMultipartUpload,
|
||||
SegmentInvalidationProducerIdentity::Replication,
|
||||
SegmentInvalidationProducerIdentity::TierTransition,
|
||||
SegmentInvalidationProducerIdentity::DirectoryObject,
|
||||
SegmentInvalidationProducerIdentity::TestFixture,
|
||||
]),
|
||||
Err(SegmentInvalidationError::UnknownProducer)
|
||||
);
|
||||
assert_eq!(
|
||||
complete_segment_invalidation_producers([
|
||||
SegmentInvalidationProducerIdentity::PutObject,
|
||||
SegmentInvalidationProducerIdentity::DeleteObject,
|
||||
SegmentInvalidationProducerIdentity::DeleteMarker,
|
||||
SegmentInvalidationProducerIdentity::Replication,
|
||||
SegmentInvalidationProducerIdentity::TierTransition,
|
||||
SegmentInvalidationProducerIdentity::DirectoryObject,
|
||||
]),
|
||||
Err(SegmentInvalidationError::InvalidProof)
|
||||
);
|
||||
assert_eq!(
|
||||
complete_segment_invalidation_producers([
|
||||
SegmentInvalidationProducerIdentity::PutObject,
|
||||
SegmentInvalidationProducerIdentity::DeleteObject,
|
||||
SegmentInvalidationProducerIdentity::DeleteMarker,
|
||||
SegmentInvalidationProducerIdentity::CompleteMultipartUpload,
|
||||
SegmentInvalidationProducerIdentity::Replication,
|
||||
SegmentInvalidationProducerIdentity::TierTransition,
|
||||
SegmentInvalidationProducerIdentity::DirectoryObject,
|
||||
]),
|
||||
Err(SegmentInvalidationError::InvalidProof)
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn segment_invalidation_entries_are_bounded_and_key_checked() {
|
||||
let envelope = envelope();
|
||||
|
||||
@@ -66,6 +66,7 @@ The manifest has the following JSON contract (all fields are required):
|
||||
| `rounds`, `duration_seconds`, `min_free_bytes` | 3..10 groups, 900..86400 seconds for measured runs, and the independently estimated free-space reservation in bytes. Synthetic runs may use 1 second. |
|
||||
| `baseline`, `candidate` | Each contains executable `binary`, full 40-character `revision`, and verified `sha256`. The runner rehashes binaries before every leg. |
|
||||
| `fixed` | `config_sha256`, `dataset_sha256`, `release_flags`, `durability`, `disk_type`, `cache_state`, `load_command`, `resource_isolation`, `topology` (`EC8+4`), and positive `offered_load_ops`. Hashes use 64 lowercase hexadecimal characters. |
|
||||
| `release_evidence` | Required for `measured` runs. It binds the 3x4 EC8+4 topology, multi-pool/multi-set coverage, per-node metrics endpoints, same-window distributed sampling, process restart and crash-restart fault modes, mixed-version reader/writer/rollback participation, and allocation/flamegraph/RSS/save-frequency profile artifact requirements. Synthetic runs do not need this field and still cannot approve release evidence. |
|
||||
| `oracles` | A map with all five scenario names. Each value contains positive integer `objects`, `versions`, `bytes`, and `sha256` of the independently prepared canonical object/version/content manifest. |
|
||||
| `expected_healed_objects` | A map with all five scenario names and independently seeded repair counts. Running-heal and MRF-replay require a positive count. |
|
||||
|
||||
@@ -75,6 +76,14 @@ object/version/content result. Fix the foreground arrival rate (offered load),
|
||||
cache preparation procedure, configuration, and hardware across every leg.
|
||||
Do not include credentials in the manifest, adapter output, or saved commands;
|
||||
the collector reads `RUSTFS_ACCESS_KEY` and `RUSTFS_SECRET_KEY` from its environment.
|
||||
The adapter must echo the measured run's `release_evidence` object in every
|
||||
measurement response. A mismatch fails the cell because it means the deployment,
|
||||
mixed-version set, crash mode, or profiler contract no longer matches the
|
||||
operator-reviewed manifest. This echo is provenance binding only; it does not
|
||||
replace the independent correctness oracle, distributed metrics samples, profile
|
||||
artifacts, or ABBA comparison thresholds. The summary tool revalidates the same
|
||||
manifest contract before it can print a measured PASS result, so hand-built or
|
||||
trimmed reports without this provenance fail closed.
|
||||
|
||||
#### Deployment Adapter Contract
|
||||
|
||||
@@ -204,6 +213,14 @@ The command prints only `PASS scanner_heal_perf ...` for measured passing ABBA
|
||||
evidence, otherwise `FAIL scanner_heal_perf ...`. The JSON and Markdown outputs
|
||||
carry the key p99/throughput/P1/P2/cache-cost fields and artifact provenance
|
||||
hashes; raw per-cell logs remain in the original artifact tree for audit.
|
||||
Failed or interrupted ABBA reports that contain only `status`, `performance`,
|
||||
`completed_cells`, and `error` also summarize as `FAIL`; they do not become
|
||||
performance evidence, and a missing comparison matrix is accepted only for a
|
||||
non-passing report.
|
||||
Measured passing reports must also retain the W10/W11 foreground-pressure,
|
||||
heal-lock-wait, and heal-attempt-cost fields emitted by the ABBA evaluator. If
|
||||
those fields are removed, empty, malformed, or length-mismatched, the quiet
|
||||
summary fails closed instead of treating the report as performance evidence.
|
||||
|
||||
They cover the complete 120-cell schedule, data isolation, missing builds and
|
||||
oracles, zero samples/requests, swallowed request errors, offered-load drift,
|
||||
|
||||
+59
-12
@@ -148,11 +148,21 @@ and exact S3 content.
|
||||
This case is a **four-node, one-drive-per-node process-restart test**. It is not
|
||||
power-loss validation, a 3x4 EC8+4 experiment, an all-version inventory, or proof
|
||||
of scanner enumeration, exact MRF disposition, legacy migration, or rollback.
|
||||
The registry keeps all G01-G14/P1-P4 and R-E/R-D/R-L release requirements pending
|
||||
until their actual feature-specific oracles and required topologies exist.
|
||||
Missing cases cannot be supplied by synthetic W20 results. W20's bounded JSON
|
||||
and file-hash helpers are reused; its ABBA performance contracts remain in
|
||||
The schema 2 registry separates the implemented single-set restart lane from
|
||||
structured release lanes for authority coverage, checkpoint/crash, status and
|
||||
outcome, MRF responsibility, mixed-version rollback, scheduler pressure,
|
||||
maintenance producers, and EC8+4 multi-set coverage. All G01-G14/P1-P4 and
|
||||
R-E/R-D/R-L release requirements stay `pending` until their actual
|
||||
feature-specific oracles, measurements and required topologies exist. Missing
|
||||
cases cannot be supplied by synthetic W20 results. W20's bounded JSON and
|
||||
file-hash helpers are reused; its ABBA performance contracts remain in
|
||||
`docs/operations/scanner-benchmark-runbook.md`.
|
||||
Measured ABBA manifests must also carry the runbook's `release_evidence`
|
||||
contract. The runner rejects reports that cannot bind the exact 3x4 EC8+4
|
||||
topology, multi-pool/multi-set shape, distributed same-window metrics endpoints,
|
||||
restart/crash modes, mixed-version reader/writer/rollback participation, and
|
||||
allocation/flamegraph/RSS/save-frequency profile artifact plan. Synthetic runs
|
||||
and manifests missing that contract remain harness-only evidence.
|
||||
|
||||
### Recording One Case
|
||||
|
||||
@@ -232,14 +242,51 @@ For automation, `--check-scanner-heal-release "$RUN_DIR"` emits one compact
|
||||
JSON decision and exits nonzero while blocked. `verified_cases` contains only
|
||||
cases that pass the complete receipt, build provenance, nextest/JUnit and real
|
||||
oracle checks; `rejected_cases` names registered cases that do not, and
|
||||
`pending_gates` names the unimplemented release requirements. Approval requires
|
||||
every registered case to verify, `pending_gates` to be empty, and a future
|
||||
registry schema capable of representing the complete release matrix. Schema 1
|
||||
is deliberately marked `release_schema_capable: false`: it models only the
|
||||
single-version, unversioned-object restart/crash cases and cannot represent
|
||||
mixed-version, rollback, EC8+4 or performance evidence. A focused run,
|
||||
synthetic harness, compile-only result, skipped/retried test, ordinary CI
|
||||
success, or removal of pending text therefore cannot become a release approval.
|
||||
`pending_gates` names the unimplemented release requirements and
|
||||
`pending_lanes` names the structured release lanes that still need real
|
||||
evidence. Schema 1 is deliberately marked `release_schema_capable: false`
|
||||
because it models only the single-version, unversioned-object restart/crash
|
||||
cases. Schema 2 can describe the wider release matrix, but approval still
|
||||
requires every registered case to verify and every required gate to leave
|
||||
`pending` only after a future checker can bind it to real feature-specific
|
||||
evidence. The current checker hard-rejects missing structured requirements and
|
||||
pending gates mapped to an implemented lane, so clearing pending text cannot
|
||||
become approval. A focused run, synthetic harness, compile-only result,
|
||||
skipped/retried test, ordinary CI success, or unregistered mixed-version,
|
||||
rollback, EC8+4 or performance claim therefore cannot become a release approval.
|
||||
For high-risk rollback gates, `evidence_fields` records the specific proof
|
||||
fields that a future real-evidence checker must bind before a pending gate can
|
||||
move out of the blocked set. G03 keeps scoped ACK tied to durable root
|
||||
publication, ACK request identity, participating peer capability snapshots, and
|
||||
mixed-peer fallback oracles; G09 keeps mixed-version reader, writer, and rollback
|
||||
payload evidence explicit. These fields are part of the release contract, not
|
||||
evidence by themselves.
|
||||
|
||||
When the real release lanes have produced their dedicated artifacts, validate
|
||||
the complete hard-gate bundle with:
|
||||
|
||||
```bash
|
||||
scripts/python_bin.sh scripts/check_test_wiring.py \
|
||||
--check-scanner-heal-release-bundle /path/to/release-evidence.json
|
||||
```
|
||||
|
||||
The bundle checker is intentionally stricter than the case checker. It requires
|
||||
schema 2 registry metadata, `evidence: measured`, the current checkout revision,
|
||||
all G01-G14/P1-P4/R-E/R-D/R-L gates, per-gate `status: pass`, lane identity,
|
||||
relative artifact paths, matching SHA256 hashes, and non-empty summaries. It
|
||||
also binds the hard evidence shape for the release claims: mixed-version gates
|
||||
must name at least two participating versions, crash/durable replay gates must
|
||||
include crash-boundary evidence, G14 must record EC8+4 with at least three nodes
|
||||
and four drives per node plus multi-set and multi-pool evidence, performance
|
||||
gates need measured durations, P3's pressure run needs at least two hours, and
|
||||
P1 needs a symbolized profile summary with resolved samples. Missing, synthetic,
|
||||
stale, tampered, undersized, or topology-mismatched evidence returns a compact
|
||||
blocked or invalid JSON result and a nonzero exit.
|
||||
|
||||
This command validates the evidence package; it does not create evidence. A
|
||||
handwritten JSON file, a synthetic harness pass, a single focused case, or a
|
||||
local unit fixture still cannot satisfy the distributed, mixed-version,
|
||||
crash-restart, durable MRF replay, EC8+4, ABBA, or profiling gates.
|
||||
|
||||
Run parser/receipt regressions with
|
||||
`scripts/python_bin.sh scripts/check_test_wiring.py --self-test`. Those fixtures
|
||||
|
||||
@@ -75,6 +75,10 @@ struct ScannerRecoveryIntentResponse {
|
||||
mode: String,
|
||||
intent_id: String,
|
||||
state: String,
|
||||
actor_sha256: String,
|
||||
idempotency_key_sha256: String,
|
||||
request_sha256: String,
|
||||
accepted_at_unix_secs: u64,
|
||||
}
|
||||
|
||||
#[derive(Debug, Serialize)]
|
||||
@@ -332,6 +336,10 @@ fn scanner_recovery_intent_record_response(
|
||||
mode: record.mode,
|
||||
intent_id: record.intent_id,
|
||||
state: record.state,
|
||||
actor_sha256: record.actor_sha256,
|
||||
idempotency_key_sha256: record.idempotency_key_sha256,
|
||||
request_sha256: record.request_sha256,
|
||||
accepted_at_unix_secs: record.accepted_at_unix_secs,
|
||||
};
|
||||
let body = serde_json::to_vec(&response).map_err(|err| {
|
||||
S3Error::with_message(
|
||||
@@ -371,10 +379,7 @@ fn scanner_recovery_intent_accept_response(
|
||||
|
||||
fn scanner_recovery_intent_executor_id(result: &rustfs_scanner::ScannerRecoveryIntentAcceptResult) -> Option<String> {
|
||||
match result {
|
||||
rustfs_scanner::ScannerRecoveryIntentAcceptResult::Accepted { record }
|
||||
| rustfs_scanner::ScannerRecoveryIntentAcceptResult::Replayed { record }
|
||||
if matches!(record.state.as_str(), "accepted" | "running") =>
|
||||
{
|
||||
rustfs_scanner::ScannerRecoveryIntentAcceptResult::Accepted { record } if record.state == "accepted" => {
|
||||
Some(record.intent_id.clone())
|
||||
}
|
||||
_ => None,
|
||||
@@ -673,7 +678,7 @@ mod tests {
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn scanner_recovery_intent_executor_only_starts_non_terminal_work() {
|
||||
fn scanner_recovery_intent_executor_starts_only_newly_accepted_work() {
|
||||
let mut record = rustfs_scanner::ScannerRecoveryIntentRecord {
|
||||
schema_version: 1,
|
||||
intent_id: "0".repeat(64),
|
||||
@@ -699,7 +704,16 @@ mod tests {
|
||||
record: record.clone(),
|
||||
})
|
||||
.as_deref(),
|
||||
Some(record.intent_id.as_str())
|
||||
None,
|
||||
"a lost-response retry must not start a duplicate executor"
|
||||
);
|
||||
record.state = "accepted".to_string();
|
||||
assert!(
|
||||
scanner_recovery_intent_executor_id(&rustfs_scanner::ScannerRecoveryIntentAcceptResult::Replayed {
|
||||
record: record.clone(),
|
||||
})
|
||||
.is_none(),
|
||||
"replayed accepted records remain durable for startup/control recovery instead of duplicating work"
|
||||
);
|
||||
record.state = "completed".to_string();
|
||||
assert!(
|
||||
|
||||
@@ -814,7 +814,11 @@ impl DefaultObjectUsecase {
|
||||
}
|
||||
}
|
||||
|
||||
rustfs_scanner::record_dirty_usage_object(&bucket, &key);
|
||||
rustfs_scanner::record_dirty_usage_object_from_producer(
|
||||
&bucket,
|
||||
&key,
|
||||
rustfs_scanner::SegmentInvalidationProducerIdentity::PutObject,
|
||||
);
|
||||
Ok::<_, S3Error>((oi, dest_versioned))
|
||||
}
|
||||
});
|
||||
|
||||
@@ -1175,7 +1175,12 @@ impl DefaultObjectUsecase {
|
||||
let manager = get_capacity_manager();
|
||||
manager.record_write_operation().await;
|
||||
let _ = helper.complete(&result);
|
||||
rustfs_scanner::record_dirty_usage_object(&bucket, &key);
|
||||
let producer = if delete_marker && version_id_clone.is_none() {
|
||||
rustfs_scanner::SegmentInvalidationProducerIdentity::DeleteMarker
|
||||
} else {
|
||||
rustfs_scanner::SegmentInvalidationProducerIdentity::DeleteObject
|
||||
};
|
||||
rustfs_scanner::record_dirty_usage_object_from_producer(&bucket, &key, producer);
|
||||
result
|
||||
}
|
||||
}
|
||||
|
||||
@@ -689,7 +689,11 @@ impl DefaultObjectUsecase {
|
||||
schedule_object_replication(obj_info.clone(), store, completion_replication_decision).await;
|
||||
}
|
||||
|
||||
rustfs_scanner::record_dirty_usage_object(&bucket, &key);
|
||||
rustfs_scanner::record_dirty_usage_object_from_producer(
|
||||
&bucket,
|
||||
&key,
|
||||
rustfs_scanner::SegmentInvalidationProducerIdentity::CompleteMultipartUpload,
|
||||
);
|
||||
Ok::<_, ApiError>(obj_info)
|
||||
}
|
||||
});
|
||||
|
||||
@@ -2054,7 +2054,11 @@ impl DefaultObjectUsecase {
|
||||
schedule_object_replication(obj_info.clone(), store, dsc).await;
|
||||
}
|
||||
|
||||
rustfs_scanner::record_dirty_usage_object(&bucket, &key);
|
||||
rustfs_scanner::record_dirty_usage_object_from_producer(
|
||||
&bucket,
|
||||
&key,
|
||||
rustfs_scanner::SegmentInvalidationProducerIdentity::PutObject,
|
||||
);
|
||||
rustfs_io_metrics::record_put_object_stage_duration_from("app_post_store_bookkeeping", post_store_stage_start);
|
||||
|
||||
let capacity_update_stage_start = put_stage_metrics_enabled.then(Instant::now);
|
||||
|
||||
@@ -3475,6 +3475,9 @@ mod tests {
|
||||
1,
|
||||
"post-admission response loss must leave exactly one canonical task"
|
||||
);
|
||||
if let Some(cache) = super::HEAL_CONTROL_REPLAY_CACHE.get() {
|
||||
cache.lock().await.clear();
|
||||
}
|
||||
|
||||
let mut retry = connect_faulty_heal_control_client(
|
||||
Arc::clone(&manager),
|
||||
|
||||
+470
-27
@@ -27,6 +27,47 @@ from scanner_abba import MAX_JSON_BYTES, digest, number, read_json, require, sha
|
||||
|
||||
|
||||
ROOT = Path(__file__).resolve().parents[1]
|
||||
SCANNER_HEAL_REGISTRY_SCHEMA_MAX = 2
|
||||
SCANNER_HEAL_RELEASE_REQUIRED_GATES = (
|
||||
"G01", "G02", "G03", "G04", "G05", "G06", "G07", "G08", "G09", "G10", "G11", "G12", "G13", "G14",
|
||||
"P1", "P2", "P3", "P4", "R-E", "R-D", "R-L",
|
||||
)
|
||||
SCANNER_HEAL_RELEASE_REQUIRED_EVIDENCE_FIELDS = {
|
||||
"G03": (
|
||||
"durable_root_publication_proof",
|
||||
"scoped_ack_request_identity",
|
||||
"participating_peer_capability_snapshot",
|
||||
"mixed_peer_ack_fallback_oracle",
|
||||
),
|
||||
"G09": (
|
||||
"mixed_version_reader_evidence",
|
||||
"mixed_version_writer_evidence",
|
||||
"rollback_payload_evidence",
|
||||
),
|
||||
}
|
||||
SCANNER_HEAL_RELEASE_BUNDLE_REQUIRED_EVIDENCE_FIELDS = {
|
||||
"G01": ("root_authority_evidence", "quota_authority_evidence"),
|
||||
"G02": ("bounded_checkpoint_oracle", "independent_version_inventory"),
|
||||
"G03": SCANNER_HEAL_RELEASE_REQUIRED_EVIDENCE_FIELDS["G03"],
|
||||
"G04": ("cache_boundary_crash_evidence", "root_floor_intent_crash_evidence"),
|
||||
"G05": ("per_object_outcome_oracle", "terminal_retention_bounds"),
|
||||
"G06": ("concurrent_status_evidence", "legacy_client_compatibility", "truncation_behavior"),
|
||||
"G07": ("mrf_responsibility_oracle", "commit_boundary_crash_matrix"),
|
||||
"G08": ("mrf_capacity_evidence", "disk_full_matrix", "replica_loss_matrix"),
|
||||
"G09": SCANNER_HEAL_RELEASE_REQUIRED_EVIDENCE_FIELDS["G09"],
|
||||
"G10": ("scheduler_bound_evidence", "pressure_recovery_evidence"),
|
||||
"G11": ("maintenance_producer_matrix", "complete_producer_inventory"),
|
||||
"G12": ("reset_quota_path_evidence", "settlement_quota_path_evidence"),
|
||||
"G13": ("quorum_minus_one_matrix", "unknown_disk_remount_matrix", "object_lock_dry_run_grace_evidence"),
|
||||
"G14": ("same_window_field_evidence", "ec8_4_evidence", "multi_set_evidence", "multi_pool_evidence"),
|
||||
"P1": ("cold_walk_share_measurement", "foreground_latency_throughput_measurement", "profile_evidence"),
|
||||
"P2": ("post_stop_convergence_measurement", "cold_segment_reuse_measurement"),
|
||||
"P3": ("two_hour_pressure_measurement", "heal_capacity_measurement", "recovery_window_measurement"),
|
||||
"P4": ("mrf_scale_measurement", "mrf_replay_cost_measurement", "retained_responsibility_evidence"),
|
||||
"R-E": ("fixed_budget_restart_evidence", "enumeration_evidence", "classification_evidence"),
|
||||
"R-D": ("manager_disposition_evidence", "event_disposition_evidence", "ledger_disposition_evidence", "grace_handling"),
|
||||
"R-L": ("legacy_source_conflict_evidence", "migration_gap_evidence", "crash_safe_source_retirement_evidence"),
|
||||
}
|
||||
SCHEDULED_ALERT_WORKFLOWS = tuple(
|
||||
item["workflow"]
|
||||
for item in json.loads((ROOT / ".github/scheduled-validations.json").read_text())
|
||||
@@ -886,9 +927,13 @@ def evidence_integer(value: object, name: str, minimum: int, maximum: int) -> in
|
||||
return value
|
||||
|
||||
|
||||
def scanner_heal_registry_schema(registry: dict[str, object]) -> int:
|
||||
return evidence_integer(registry.get("schema"), "registry schema", 1, SCANNER_HEAL_REGISTRY_SCHEMA_MAX)
|
||||
|
||||
|
||||
def scanner_heal_oracle_names(root: Path) -> tuple[str, ...]:
|
||||
registry = read_json(root / ".config/scanner-heal-required-tests.json")
|
||||
evidence_integer(registry.get("schema"), "registry schema", 1, 1)
|
||||
scanner_heal_registry_schema(registry)
|
||||
cases = registry.get("cases")
|
||||
require(isinstance(cases, dict) and cases, "invalid scanner/heal registry")
|
||||
names = set()
|
||||
@@ -906,6 +951,86 @@ def scanner_heal_oracle_names(root: Path) -> tuple[str, ...]:
|
||||
return tuple(sorted(names))
|
||||
|
||||
|
||||
def scanner_heal_release_requirements(registry: dict[str, object]) -> tuple[dict[str, dict[str, object]], bool, list[str]]:
|
||||
schema = scanner_heal_registry_schema(registry)
|
||||
if schema == 1:
|
||||
pending = registry.get("release_pending")
|
||||
require(isinstance(pending, dict), "invalid scanner/heal release requirements")
|
||||
requirements = {}
|
||||
for gate, reason in pending.items():
|
||||
require(isinstance(gate, str) and re.fullmatch(r"[A-Z][A-Z0-9-]*", gate) is not None,
|
||||
"invalid scanner/heal release gate")
|
||||
require(isinstance(reason, str) and reason.strip(), f"missing release requirement for {gate}")
|
||||
requirements[gate] = {"gate": gate, "status": "pending", "lane": "schema-1-pending",
|
||||
"description": reason, "requires": [reason]}
|
||||
return requirements, False, ["schema-1-pending"] if requirements else []
|
||||
|
||||
lanes = registry.get("release_lanes")
|
||||
cases = registry.get("cases")
|
||||
require(isinstance(cases, dict) and cases, "invalid scanner/heal registry")
|
||||
require(isinstance(lanes, dict) and lanes, "invalid scanner/heal release lanes")
|
||||
lane_statuses = {}
|
||||
lane_gates = {}
|
||||
for lane_id, lane in lanes.items():
|
||||
require(isinstance(lane_id, str) and re.fullmatch(r"[a-z0-9-]+", lane_id) is not None,
|
||||
"invalid scanner/heal release lane")
|
||||
require(isinstance(lane, dict), f"invalid release lane {lane_id}")
|
||||
status = lane.get("status")
|
||||
require(status in ("implemented", "pending"), f"invalid release lane status for {lane_id}")
|
||||
lane_statuses[lane_id] = status
|
||||
if status == "implemented":
|
||||
lane_cases = lane.get("cases")
|
||||
require(isinstance(lane_cases, list) and lane_cases and all(isinstance(case, str) and case for case in lane_cases),
|
||||
f"implemented release lane {lane_id} has no cases")
|
||||
require(all(case in cases for case in lane_cases), f"implemented release lane {lane_id} has unknown cases")
|
||||
else:
|
||||
gates = lane.get("gates")
|
||||
require(isinstance(gates, list) and gates and all(isinstance(gate, str) and gate for gate in gates),
|
||||
f"pending release lane {lane_id} has no gates")
|
||||
lane_gates[lane_id] = set(gates)
|
||||
|
||||
raw_requirements = registry.get("release_requirements")
|
||||
require(isinstance(raw_requirements, list) and raw_requirements, "invalid scanner/heal release requirements")
|
||||
requirements: dict[str, dict[str, object]] = {}
|
||||
for item in raw_requirements:
|
||||
require(isinstance(item, dict), "invalid scanner/heal release requirement")
|
||||
gate = item.get("gate")
|
||||
require(isinstance(gate, str) and re.fullmatch(r"[A-Z][A-Z0-9-]*", gate) is not None,
|
||||
"invalid scanner/heal release gate")
|
||||
require(gate not in requirements, f"duplicate scanner/heal release gate {gate}")
|
||||
status = item.get("status")
|
||||
require(status == "pending", f"release gate {gate} must stay pending until real evidence is registered")
|
||||
lane = item.get("lane")
|
||||
require(isinstance(lane, str) and lane in lane_statuses, f"unknown release lane for {gate}")
|
||||
require(lane_statuses[lane] == "pending", f"pending release gate {gate} mapped to non-pending lane {lane}")
|
||||
description = item.get("description")
|
||||
require(isinstance(description, str) and description.strip(), f"missing release requirement for {gate}")
|
||||
requires = item.get("requires")
|
||||
require(isinstance(requires, list) and requires and
|
||||
all(isinstance(requirement, str) and requirement.strip() for requirement in requires),
|
||||
f"missing concrete evidence requirements for {gate}")
|
||||
evidence_fields = item.get("evidence_fields", [])
|
||||
require(isinstance(evidence_fields, list) and
|
||||
all(isinstance(field, str) and re.fullmatch(r"[a-z0-9][a-z0-9_]*", field) is not None
|
||||
for field in evidence_fields),
|
||||
f"invalid evidence fields for release gate {gate}")
|
||||
required_fields = set(SCANNER_HEAL_RELEASE_REQUIRED_EVIDENCE_FIELDS.get(gate, ()))
|
||||
missing_fields = sorted(required_fields - set(evidence_fields))
|
||||
require(not missing_fields,
|
||||
f"release gate {gate} missing required evidence fields: {', '.join(missing_fields)}")
|
||||
requirements[gate] = item
|
||||
|
||||
missing = sorted(set(SCANNER_HEAL_RELEASE_REQUIRED_GATES) - set(requirements))
|
||||
require(not missing, f"missing scanner/heal release requirements: {', '.join(missing)}")
|
||||
for lane, gates in lane_gates.items():
|
||||
unknown = sorted(gates - set(requirements))
|
||||
require(not unknown, f"pending release lane {lane} has unknown gates: {', '.join(unknown)}")
|
||||
mapped = {gate for gate, requirement in requirements.items() if requirement["lane"] == lane}
|
||||
require(gates == mapped, f"pending release lane {lane} gates do not match release requirements")
|
||||
pending_lanes = sorted(lane for lane, status in lane_statuses.items() if status == "pending")
|
||||
return requirements, True, pending_lanes
|
||||
|
||||
|
||||
def begin_scanner_heal_receipt(root: Path, directory: Path, binary: Path, test_binary: Path) -> None:
|
||||
"""Record an existing build; this command never builds or runs a test."""
|
||||
require(not directory.exists(), "scanner/heal run directory must be new")
|
||||
@@ -969,7 +1094,7 @@ def check_scanner_heal_evidence(root: Path, directory: Path, case_id: str) -> li
|
||||
"""Validate one actual case, or fail the release while required lanes are pending."""
|
||||
try:
|
||||
registry = read_json(root / ".config/scanner-heal-required-tests.json")
|
||||
evidence_integer(registry.get("schema"), "registry schema", 1, 1)
|
||||
scanner_heal_registry_schema(registry)
|
||||
require(registry.get("cases"), "invalid scanner/heal registry")
|
||||
selected = registry["cases"] if case_id == "release" else {case_id: registry["cases"][case_id]}
|
||||
run = read_json(directory / "run.json")
|
||||
@@ -1088,7 +1213,10 @@ def check_scanner_heal_evidence(root: Path, directory: Path, case_id: str) -> li
|
||||
require(all(keys == sorted(obj["key"] for obj in objects) for keys in node_listings),
|
||||
"S3 listing differs from object oracle")
|
||||
if case_id == "release":
|
||||
errors.extend(f"pending {gate}: {reason}" for gate, reason in registry["release_pending"].items())
|
||||
release_requirements, _, _ = scanner_heal_release_requirements(registry)
|
||||
errors.extend(
|
||||
f"pending {gate}: {requirement['description']}" for gate, requirement in release_requirements.items()
|
||||
)
|
||||
return errors
|
||||
except (OSError, KeyError, TypeError, ValueError, ET.ParseError) as error:
|
||||
return [f"scanner/heal evidence rejected: {error}"]
|
||||
@@ -1097,15 +1225,10 @@ def check_scanner_heal_evidence(root: Path, directory: Path, case_id: str) -> li
|
||||
def scanner_heal_release_status(root: Path, directory: Path) -> dict[str, object]:
|
||||
"""Return a compact release decision without weakening case validation."""
|
||||
registry = read_json(root / ".config/scanner-heal-required-tests.json")
|
||||
evidence_integer(registry.get("schema"), "registry schema", 1, 1)
|
||||
scanner_heal_registry_schema(registry)
|
||||
cases = registry.get("cases")
|
||||
require(isinstance(cases, dict) and cases, "invalid scanner/heal registry")
|
||||
pending = registry.get("release_pending")
|
||||
require(isinstance(pending, dict), "invalid scanner/heal release requirements")
|
||||
for gate, reason in pending.items():
|
||||
require(isinstance(gate, str) and re.fullmatch(r"[A-Z][A-Z0-9-]*", gate) is not None,
|
||||
"invalid scanner/heal release gate")
|
||||
require(isinstance(reason, str) and reason.strip(), f"missing release requirement for {gate}")
|
||||
release_requirements, release_schema_capable, pending_lanes = scanner_heal_release_requirements(registry)
|
||||
|
||||
verified_cases = []
|
||||
rejected_cases = []
|
||||
@@ -1119,10 +1242,128 @@ def scanner_heal_release_status(root: Path, directory: Path) -> dict[str, object
|
||||
"schema": 1,
|
||||
"decision": "blocked",
|
||||
"release_approved": False,
|
||||
"release_schema_capable": False,
|
||||
"release_schema_capable": release_schema_capable,
|
||||
"verified_cases": verified_cases,
|
||||
"rejected_cases": rejected_cases,
|
||||
"pending_gates": sorted(pending),
|
||||
"pending_gates": sorted(release_requirements),
|
||||
"pending_lanes": pending_lanes,
|
||||
}
|
||||
|
||||
|
||||
def release_bundle_artifact_path(bundle_path: Path, raw_path: object, gate: str, field: str) -> Path:
|
||||
require(isinstance(raw_path, str) and raw_path.strip(), f"{gate}.{field} missing artifact")
|
||||
path = Path(raw_path)
|
||||
require(not path.is_absolute() and ".." not in path.parts, f"{gate}.{field} artifact path escapes bundle directory")
|
||||
resolved = (bundle_path.parent / path).resolve()
|
||||
require(resolved.is_relative_to(bundle_path.parent.resolve()), f"{gate}.{field} artifact path escapes bundle directory")
|
||||
require(resolved.is_file(), f"{gate}.{field} artifact is missing")
|
||||
return resolved
|
||||
|
||||
|
||||
def validate_release_bundle_artifact(bundle_path: Path, gate: str, field: str, evidence: dict[str, object]) -> None:
|
||||
require(evidence.get("evidence_type") == "measured", f"{gate}.{field} must be measured evidence")
|
||||
artifact = release_bundle_artifact_path(bundle_path, evidence.get("artifact"), gate, field)
|
||||
require(sha(evidence.get("sha256")) and digest(artifact) == evidence["sha256"], f"{gate}.{field} artifact hash mismatch")
|
||||
summary = evidence.get("summary")
|
||||
require(isinstance(summary, str) and summary.strip(), f"{gate}.{field} missing human summary")
|
||||
if gate.startswith("P"):
|
||||
duration = evidence_integer(evidence.get("duration_seconds"), f"{gate}.{field}.duration_seconds", 1, 86400)
|
||||
require(duration >= 900, f"{gate}.{field} requires at least 900 seconds")
|
||||
if gate == "P3" and field == "two_hour_pressure_measurement":
|
||||
require(duration >= 7200, f"{gate}.{field} requires at least two hours")
|
||||
elif "duration_seconds" in evidence:
|
||||
evidence_integer(evidence.get("duration_seconds"), f"{gate}.{field}.duration_seconds", 1, 86400)
|
||||
if gate in ("G03", "G09", "R-L"):
|
||||
versions = evidence.get("versions")
|
||||
require(isinstance(versions, list) and
|
||||
len({version for version in versions if isinstance(version, str) and version.strip()}) >= 2,
|
||||
f"{gate}.{field} requires mixed-version evidence")
|
||||
if gate in ("G04", "G07", "R-E", "R-L"):
|
||||
crash_points = evidence.get("crash_points")
|
||||
require(isinstance(crash_points, list) and crash_points,
|
||||
f"{gate}.{field} requires crash-boundary evidence")
|
||||
if gate == "G14":
|
||||
if field == "ec8_4_evidence":
|
||||
topology = evidence.get("topology")
|
||||
require(isinstance(topology, dict), "G14.ec8_4_evidence missing topology")
|
||||
require(topology.get("erasure") == "EC8+4", "G14.ec8_4_evidence must record EC8+4")
|
||||
require(evidence_integer(topology.get("nodes"), "G14 topology nodes", 3, 64) >= 3,
|
||||
"G14.ec8_4_evidence requires at least three nodes")
|
||||
require(evidence_integer(topology.get("drives_per_node"), "G14 topology drives", 4, 64) >= 4,
|
||||
"G14.ec8_4_evidence requires at least four drives per node")
|
||||
if field == "multi_set_evidence":
|
||||
evidence_integer(evidence.get("sets"), "G14 multi_set_evidence.sets", 2, 1024)
|
||||
if field == "multi_pool_evidence":
|
||||
evidence_integer(evidence.get("pools"), "G14 multi_pool_evidence.pools", 2, 1024)
|
||||
if field == "profile_evidence":
|
||||
evidence_integer(evidence.get("resolved_samples"), f"{gate}.{field}.resolved_samples", 1, 2**63 - 1)
|
||||
|
||||
|
||||
def scanner_heal_release_bundle_status(root: Path, bundle_path: Path) -> dict[str, object]:
|
||||
"""Validate a complete hard-gate evidence bundle without accepting synthetic claims."""
|
||||
registry = read_json(root / ".config/scanner-heal-required-tests.json")
|
||||
requirements, release_schema_capable, pending_lanes = scanner_heal_release_requirements(registry)
|
||||
require(release_schema_capable, "scanner/heal release bundle requires schema 2 registry")
|
||||
bundle_path = bundle_path.resolve()
|
||||
bundle = read_json(bundle_path)
|
||||
require(bundle.get("schema") == 1, "unsupported scanner/heal release evidence bundle schema")
|
||||
require(bundle.get("evidence") == "measured", "scanner/heal release evidence bundle must be measured")
|
||||
revision = subprocess.check_output(["git", "rev-parse", "HEAD"], cwd=root, text=True).strip()
|
||||
require(bundle.get("source_revision") == revision, "scanner/heal release evidence source revision mismatch")
|
||||
raw_gates = bundle.get("gates")
|
||||
require(isinstance(raw_gates, dict), "scanner/heal release evidence bundle missing gates")
|
||||
|
||||
verified: list[str] = []
|
||||
rejected: dict[str, list[str]] = {}
|
||||
for gate, requirement in requirements.items():
|
||||
gate_errors: list[str] = []
|
||||
gate_evidence = raw_gates.get(gate)
|
||||
if not isinstance(gate_evidence, dict):
|
||||
rejected[gate] = ["missing gate evidence"]
|
||||
continue
|
||||
if gate_evidence.get("status") != "pass":
|
||||
gate_errors.append("gate status must be pass")
|
||||
if gate_evidence.get("lane") != requirement["lane"]:
|
||||
gate_errors.append("gate lane mismatch")
|
||||
if gate_evidence.get("evidence_type") != "measured":
|
||||
gate_errors.append("gate evidence type must be measured")
|
||||
fields = gate_evidence.get("evidence_fields")
|
||||
if not isinstance(fields, dict):
|
||||
gate_errors.append("missing gate evidence fields")
|
||||
fields = {}
|
||||
required_fields = tuple(SCANNER_HEAL_RELEASE_BUNDLE_REQUIRED_EVIDENCE_FIELDS[gate])
|
||||
missing_fields = [field for field in required_fields if field not in fields]
|
||||
if missing_fields:
|
||||
gate_errors.append(f"missing required fields: {', '.join(missing_fields)}")
|
||||
for field in required_fields:
|
||||
if field not in fields:
|
||||
continue
|
||||
evidence = fields[field]
|
||||
if not isinstance(evidence, dict):
|
||||
gate_errors.append(f"{field} must be an object")
|
||||
continue
|
||||
try:
|
||||
validate_release_bundle_artifact(bundle_path, gate, field, evidence)
|
||||
except (OSError, KeyError, TypeError, ValueError) as error:
|
||||
gate_errors.append(str(error))
|
||||
if gate_errors:
|
||||
rejected[gate] = gate_errors
|
||||
else:
|
||||
verified.append(gate)
|
||||
|
||||
unknown = sorted(set(raw_gates) - set(requirements))
|
||||
if unknown:
|
||||
rejected["unknown"] = [f"unknown gates: {', '.join(unknown)}"]
|
||||
approved = not rejected and sorted(verified) == sorted(requirements)
|
||||
return {
|
||||
"schema": 1,
|
||||
"decision": "approved" if approved else "blocked",
|
||||
"release_approved": approved,
|
||||
"release_schema_capable": release_schema_capable,
|
||||
"verified_gates": sorted(verified),
|
||||
"rejected_gates": rejected,
|
||||
"pending_gates": [] if approved else sorted(gate for gate in requirements if gate not in verified),
|
||||
"pending_lanes": [] if approved else pending_lanes,
|
||||
}
|
||||
|
||||
|
||||
@@ -1310,16 +1551,26 @@ class SelfTests(unittest.TestCase):
|
||||
)
|
||||
+ "</testsuite></testsuites>"
|
||||
)
|
||||
physical = {"has_xl_meta": True, "version_id": None, "data_dir": "data-generation",
|
||||
"erasure_index": 1, "data_blocks": 2, "parity_blocks": 2, "expected_part_numbers": [1],
|
||||
"present_part_fingerprints": {"1": {"size": 12, "sha256": "c" * 64}},
|
||||
"inline_data_fingerprint": None}
|
||||
obj = {"key": "object", "version_id": None, "expected_bytes": 16, "actual_bytes": 16,
|
||||
"expected_sha256": "d" * 64, "actual_sha256": "d" * 64,
|
||||
"expected_physical": physical, "physical": physical}
|
||||
objects = [dict(obj, key=f"object-{index}") for index in range(9)]
|
||||
objects[-1] = dict(objects[-1], expected_physical=None)
|
||||
def oracle_objects(requirement: dict[str, object]) -> list[dict[str, object]]:
|
||||
topology = requirement["topology"]
|
||||
total_blocks = topology["nodes"] * topology["drives_per_node"]
|
||||
parity_blocks = 4 if total_blocks == 12 else total_blocks // 2
|
||||
data_blocks = total_blocks - parity_blocks
|
||||
physical = {"has_xl_meta": True, "version_id": None, "data_dir": "data-generation",
|
||||
"erasure_index": 1, "data_blocks": data_blocks, "parity_blocks": parity_blocks,
|
||||
"expected_part_numbers": [1],
|
||||
"present_part_fingerprints": {"1": {"size": 12, "sha256": "c" * 64}},
|
||||
"inline_data_fingerprint": None}
|
||||
obj = {"key": "object", "version_id": None, "expected_bytes": 16, "actual_bytes": 16,
|
||||
"expected_sha256": "d" * 64, "actual_sha256": "d" * 64,
|
||||
"expected_physical": physical, "physical": physical}
|
||||
count = requirement.get("min_objects", 9)
|
||||
objects = [dict(obj, key=f"object-{index}") for index in range(count)]
|
||||
objects[-1] = dict(objects[-1], expected_physical=None)
|
||||
return objects
|
||||
|
||||
for case_id, requirement in requirements.items():
|
||||
objects = oracle_objects(requirement)
|
||||
write_json(run_dir / requirement["oracle"], {
|
||||
"schema": 1, "evidence": requirement["evidence"], "case": case_id,
|
||||
"run_id": "a" * 32, "source_revision": "b" * 40,
|
||||
@@ -1328,11 +1579,114 @@ class SelfTests(unittest.TestCase):
|
||||
"binary_sha256": build["sha256"], "test_binary_sha256": build["sha256"],
|
||||
"topology": requirement["topology"], "pid_before": 10, "pid_after": 11,
|
||||
"unclean_shutdown_marker": requirement["unclean_shutdown_marker"],
|
||||
"objects": objects, "node_listings": [[item["key"] for item in objects]] * 4,
|
||||
"objects": objects, "node_listings": [[item["key"] for item in objects]] * requirement["topology"]["nodes"],
|
||||
})
|
||||
finish_scanner_heal_receipt(run_dir, 0, root)
|
||||
return root, run_dir
|
||||
|
||||
def scanner_heal_release_bundle_fixture(self, directory: Path) -> tuple[Path, Path]:
|
||||
"""Parser fixtures only; the bundle is not runtime evidence."""
|
||||
root, _ = self.scanner_heal_fixture(directory)
|
||||
bundle_dir = directory / "bundle"
|
||||
artifact_dir = bundle_dir / "artifacts"
|
||||
artifact_dir.mkdir(parents=True)
|
||||
registry = read_json(root / ".config/scanner-heal-required-tests.json")
|
||||
requirements, _, _ = scanner_heal_release_requirements(registry)
|
||||
gates = {}
|
||||
for gate, requirement in requirements.items():
|
||||
fields = {}
|
||||
for field in SCANNER_HEAL_RELEASE_BUNDLE_REQUIRED_EVIDENCE_FIELDS[gate]:
|
||||
artifact = artifact_dir / f"{gate}-{field}.json"
|
||||
write_json(artifact, {"gate": gate, "field": field, "fixture": True})
|
||||
evidence = {
|
||||
"artifact": artifact.relative_to(bundle_dir).as_posix(),
|
||||
"sha256": digest(artifact),
|
||||
"evidence_type": "measured",
|
||||
"summary": f"parser fixture for {gate}.{field}",
|
||||
}
|
||||
if gate.startswith("P"):
|
||||
evidence["duration_seconds"] = 900
|
||||
if gate == "P3" and field == "two_hour_pressure_measurement":
|
||||
evidence["duration_seconds"] = 7200
|
||||
if gate in ("G03", "G09", "R-L"):
|
||||
evidence["versions"] = ["previous", "candidate"]
|
||||
if gate in ("G04", "G07", "R-E", "R-L"):
|
||||
evidence["crash_points"] = ["before-commit"]
|
||||
if gate == "G14" and field == "ec8_4_evidence":
|
||||
evidence["topology"] = {"erasure": "EC8+4", "nodes": 3, "drives_per_node": 4}
|
||||
if gate == "G14" and field == "multi_set_evidence":
|
||||
evidence["sets"] = 2
|
||||
if gate == "G14" and field == "multi_pool_evidence":
|
||||
evidence["pools"] = 2
|
||||
if field == "profile_evidence":
|
||||
evidence["resolved_samples"] = 1
|
||||
fields[field] = evidence
|
||||
gates[gate] = {
|
||||
"status": "pass",
|
||||
"lane": requirement["lane"],
|
||||
"evidence_type": "measured",
|
||||
"evidence_fields": fields,
|
||||
}
|
||||
bundle = bundle_dir / "release-evidence.json"
|
||||
write_json(bundle, {"schema": 1, "evidence": "measured", "source_revision": "b" * 40, "gates": gates})
|
||||
return root, bundle
|
||||
|
||||
def test_scanner_heal_release_bundle_accepts_complete_measured_evidence(self) -> None:
|
||||
with tempfile.TemporaryDirectory() as tmp:
|
||||
root, bundle = self.scanner_heal_release_bundle_fixture(Path(tmp))
|
||||
with mock.patch("subprocess.check_output", return_value="b" * 40):
|
||||
status = scanner_heal_release_bundle_status(root, bundle)
|
||||
self.assertEqual(status["decision"], "approved")
|
||||
self.assertTrue(status["release_approved"])
|
||||
self.assertEqual(len(status["verified_gates"]), len(SCANNER_HEAL_RELEASE_REQUIRED_GATES))
|
||||
self.assertEqual(status["pending_gates"], [])
|
||||
self.assertEqual(status["pending_lanes"], [])
|
||||
|
||||
def test_scanner_heal_release_bundle_rejects_synthetic_or_missing_fields(self) -> None:
|
||||
for fault in ("synthetic", "missing-field", "hash"):
|
||||
with self.subTest(fault=fault), tempfile.TemporaryDirectory() as tmp:
|
||||
root, bundle = self.scanner_heal_release_bundle_fixture(Path(tmp))
|
||||
data = read_json(bundle)
|
||||
if fault == "synthetic":
|
||||
data["evidence"] = "synthetic"
|
||||
elif fault == "missing-field":
|
||||
del data["gates"]["G09"]["evidence_fields"]["rollback_payload_evidence"]
|
||||
else:
|
||||
artifact = bundle.parent / data["gates"]["G01"]["evidence_fields"]["root_authority_evidence"]["artifact"]
|
||||
artifact.write_text(artifact.read_text(encoding="utf-8") + "\n", encoding="utf-8")
|
||||
write_json(bundle, data)
|
||||
|
||||
with mock.patch("subprocess.check_output", return_value="b" * 40):
|
||||
if fault == "synthetic":
|
||||
with self.assertRaisesRegex(ValueError, "must be measured"):
|
||||
scanner_heal_release_bundle_status(root, bundle)
|
||||
else:
|
||||
status = scanner_heal_release_bundle_status(root, bundle)
|
||||
if fault != "synthetic":
|
||||
self.assertEqual(status["decision"], "blocked")
|
||||
self.assertFalse(status["release_approved"])
|
||||
self.assertTrue(status["rejected_gates"])
|
||||
|
||||
def test_scanner_heal_release_bundle_enforces_topology_duration_profile_and_versions(self) -> None:
|
||||
for fault, gate, field, mutation, expected in (
|
||||
("topology", "G14", "ec8_4_evidence", lambda item: item.update({"topology": {"erasure": "EC4+2", "nodes": 2, "drives_per_node": 3}}), "EC8+4"),
|
||||
("missing-duration", "P1", "cold_walk_share_measurement", lambda item: item.pop("duration_seconds"), "duration_seconds"),
|
||||
("duration", "P3", "two_hour_pressure_measurement", lambda item: item.update({"duration_seconds": 7199}), "two hours"),
|
||||
("profile", "P1", "profile_evidence", lambda item: item.pop("resolved_samples"), "resolved_samples"),
|
||||
("versions", "G09", "mixed_version_reader_evidence", lambda item: item.update({"versions": [1, 2]}), "mixed-version"),
|
||||
):
|
||||
with self.subTest(fault=fault), tempfile.TemporaryDirectory() as tmp:
|
||||
root, bundle = self.scanner_heal_release_bundle_fixture(Path(tmp))
|
||||
data = read_json(bundle)
|
||||
mutation(data["gates"][gate]["evidence_fields"][field])
|
||||
write_json(bundle, data)
|
||||
|
||||
with mock.patch("subprocess.check_output", return_value="b" * 40):
|
||||
status = scanner_heal_release_bundle_status(root, bundle)
|
||||
self.assertEqual(status["decision"], "blocked")
|
||||
self.assertFalse(status["release_approved"])
|
||||
self.assertTrue(any(expected in error for error in status["rejected_gates"][gate]))
|
||||
|
||||
def test_scanner_heal_case_does_not_approve_pending_release(self) -> None:
|
||||
with tempfile.TemporaryDirectory() as tmp:
|
||||
root, run_dir = self.scanner_heal_fixture(Path(tmp))
|
||||
@@ -1346,14 +1700,21 @@ class SelfTests(unittest.TestCase):
|
||||
status = scanner_heal_release_status(root, run_dir)
|
||||
self.assertEqual(status["decision"], "blocked")
|
||||
self.assertFalse(status["release_approved"])
|
||||
self.assertTrue(status["release_schema_capable"])
|
||||
self.assertEqual(status["rejected_cases"], [])
|
||||
self.assertEqual(len(status["pending_gates"]), 21)
|
||||
self.assertIn("mixed-version-rollback", status["pending_lanes"])
|
||||
self.assertIn("ec8-4-multiset", status["pending_lanes"])
|
||||
self.assertIn("scheduler-pressure", status["pending_lanes"])
|
||||
|
||||
def test_scanner_heal_case_only_schema_cannot_approve_release(self) -> None:
|
||||
with tempfile.TemporaryDirectory() as tmp:
|
||||
root, run_dir = self.scanner_heal_fixture(Path(tmp))
|
||||
registry = read_json(root / ".config/scanner-heal-required-tests.json")
|
||||
registry["schema"] = 1
|
||||
registry["release_pending"] = {}
|
||||
registry.pop("release_lanes")
|
||||
registry.pop("release_requirements")
|
||||
write_json(root / ".config/scanner-heal-required-tests.json", registry)
|
||||
|
||||
status = scanner_heal_release_status(root, run_dir)
|
||||
@@ -1363,11 +1724,78 @@ class SelfTests(unittest.TestCase):
|
||||
self.assertEqual(status["rejected_cases"], [])
|
||||
self.assertEqual(status["pending_gates"], [])
|
||||
|
||||
def test_scanner_heal_release_requirements_cannot_be_cleared(self) -> None:
|
||||
with tempfile.TemporaryDirectory() as tmp:
|
||||
root, run_dir = self.scanner_heal_fixture(Path(tmp))
|
||||
registry = read_json(root / ".config/scanner-heal-required-tests.json")
|
||||
registry["release_requirements"] = []
|
||||
write_json(root / ".config/scanner-heal-required-tests.json", registry)
|
||||
|
||||
with self.assertRaisesRegex(ValueError, "invalid scanner/heal release requirements"):
|
||||
scanner_heal_release_status(root, run_dir)
|
||||
|
||||
def test_scanner_heal_release_matrix_lanes_remain_pending_without_real_evidence(self) -> None:
|
||||
with tempfile.TemporaryDirectory() as tmp:
|
||||
root, run_dir = self.scanner_heal_fixture(Path(tmp))
|
||||
|
||||
errors = check_scanner_heal_evidence(root, run_dir, "release")
|
||||
for gate, text in (
|
||||
("G03", "Exact scoped ACK"),
|
||||
("G09", "mixed-version reader/writer"),
|
||||
("G14", "3x4 EC8+4"),
|
||||
("P3", "two-hour pressure/heal capacity"),
|
||||
("R-L", "Legacy source conflicts"),
|
||||
):
|
||||
self.assertTrue(any(error.startswith(f"pending {gate}:") and text in error for error in errors), gate)
|
||||
|
||||
status = scanner_heal_release_status(root, run_dir)
|
||||
self.assertFalse(status["release_approved"])
|
||||
self.assertIn("mixed-version-rollback", status["pending_lanes"])
|
||||
self.assertIn("ec8-4-multiset", status["pending_lanes"])
|
||||
self.assertIn("scheduler-pressure", status["pending_lanes"])
|
||||
requirements, _, _ = scanner_heal_release_requirements(read_json(root / ".config/scanner-heal-required-tests.json"))
|
||||
self.assertIn("durable_root_publication_proof", requirements["G03"]["evidence_fields"])
|
||||
self.assertIn("mixed_version_writer_evidence", requirements["G09"]["evidence_fields"])
|
||||
|
||||
def test_scanner_heal_required_evidence_fields_cannot_be_removed(self) -> None:
|
||||
with tempfile.TemporaryDirectory() as tmp:
|
||||
root, run_dir = self.scanner_heal_fixture(Path(tmp))
|
||||
registry = read_json(root / ".config/scanner-heal-required-tests.json")
|
||||
for gate in ("G03", "G09"):
|
||||
for requirement in registry["release_requirements"]:
|
||||
if requirement["gate"] == gate:
|
||||
requirement["evidence_fields"] = []
|
||||
break
|
||||
write_json(root / ".config/scanner-heal-required-tests.json", registry)
|
||||
|
||||
with self.subTest(gate=gate):
|
||||
with self.assertRaisesRegex(ValueError, f"release gate {gate} missing required evidence fields"):
|
||||
scanner_heal_release_status(root, run_dir)
|
||||
|
||||
registry = read_json(root / ".config/scanner-heal-required-tests.json")
|
||||
for requirement in registry["release_requirements"]:
|
||||
if requirement["gate"] == gate:
|
||||
requirement["evidence_fields"] = list(SCANNER_HEAL_RELEASE_REQUIRED_EVIDENCE_FIELDS[gate])
|
||||
break
|
||||
|
||||
def test_scanner_heal_pending_gate_cannot_map_to_implemented_lane(self) -> None:
|
||||
with tempfile.TemporaryDirectory() as tmp:
|
||||
root, run_dir = self.scanner_heal_fixture(Path(tmp))
|
||||
registry = read_json(root / ".config/scanner-heal-required-tests.json")
|
||||
registry["release_requirements"][0]["lane"] = "single-set-restart"
|
||||
write_json(root / ".config/scanner-heal-required-tests.json", registry)
|
||||
|
||||
with self.assertRaisesRegex(ValueError, "mapped to non-pending lane"):
|
||||
scanner_heal_release_status(root, run_dir)
|
||||
|
||||
def test_scanner_heal_release_status_rejects_synthetic_case(self) -> None:
|
||||
with tempfile.TemporaryDirectory() as tmp:
|
||||
root, run_dir = self.scanner_heal_fixture(Path(tmp))
|
||||
registry = read_json(root / ".config/scanner-heal-required-tests.json")
|
||||
registry["schema"] = 1
|
||||
registry["release_pending"] = {}
|
||||
registry.pop("release_lanes")
|
||||
registry.pop("release_requirements")
|
||||
write_json(root / ".config/scanner-heal-required-tests.json", registry)
|
||||
path = run_dir / "background-target-crash.json"
|
||||
oracle = read_json(path)
|
||||
@@ -1379,6 +1807,7 @@ class SelfTests(unittest.TestCase):
|
||||
status = scanner_heal_release_status(root, run_dir)
|
||||
self.assertEqual(status["decision"], "blocked")
|
||||
self.assertFalse(status["release_approved"])
|
||||
self.assertFalse(status["release_schema_capable"])
|
||||
self.assertEqual(status["rejected_cases"], ["background-target-crash"])
|
||||
self.assertEqual(status["pending_gates"], [])
|
||||
|
||||
@@ -1386,9 +1815,14 @@ class SelfTests(unittest.TestCase):
|
||||
with tempfile.TemporaryDirectory() as tmp:
|
||||
root, run_dir = self.scanner_heal_fixture(Path(tmp))
|
||||
registry = read_json(root / ".config/scanner-heal-required-tests.json")
|
||||
registry["schema"] = 1
|
||||
registry["release_pending"] = {}
|
||||
registry.pop("release_lanes")
|
||||
registry.pop("release_requirements")
|
||||
write_json(root / ".config/scanner-heal-required-tests.json", registry)
|
||||
(run_dir / "background-target-crash.json").unlink()
|
||||
for case_id, requirement in registry["cases"].items():
|
||||
if case_id != "background-target-restart":
|
||||
(run_dir / requirement["oracle"]).unlink()
|
||||
(run_dir / "execution.json").unlink()
|
||||
finish_scanner_heal_receipt(run_dir, 0, root)
|
||||
|
||||
@@ -1396,7 +1830,7 @@ class SelfTests(unittest.TestCase):
|
||||
self.assertEqual(status["decision"], "blocked")
|
||||
self.assertFalse(status["release_approved"])
|
||||
self.assertEqual(status["verified_cases"], ["background-target-restart"])
|
||||
self.assertEqual(status["rejected_cases"], ["background-target-crash"])
|
||||
self.assertEqual(status["rejected_cases"], sorted(set(registry["cases"]) - {"background-target-restart"}))
|
||||
|
||||
def test_scanner_heal_finish_collects_oracles_from_registry(self) -> None:
|
||||
with tempfile.TemporaryDirectory() as tmp:
|
||||
@@ -2246,7 +2680,7 @@ def main() -> int:
|
||||
suite = unittest.defaultTestLoader.loadTestsFromTestCase(SelfTests)
|
||||
return 0 if unittest.TextTestRunner(verbosity=2).run(suite).wasSuccessful() else 1
|
||||
if sys.argv[1:2] in (["--begin-scanner-heal"], ["--finish-scanner-heal"], ["--check-scanner-heal"],
|
||||
["--check-scanner-heal-release"]):
|
||||
["--check-scanner-heal-release"], ["--check-scanner-heal-release-bundle"]):
|
||||
try:
|
||||
if len(sys.argv) == 5 and sys.argv[1] == "--begin-scanner-heal":
|
||||
begin_scanner_heal_receipt(ROOT, Path(sys.argv[2]), Path(sys.argv[3]), Path(sys.argv[4]))
|
||||
@@ -2270,7 +2704,16 @@ def main() -> int:
|
||||
return 2
|
||||
print(json.dumps(status, sort_keys=True, separators=(",", ":")))
|
||||
return 0 if status["release_approved"] else 1
|
||||
raise ValueError("expected --begin-scanner-heal DIR BINARY TEST_BINARY, --finish-scanner-heal DIR EXIT, --check-scanner-heal DIR CASE|release, or --check-scanner-heal-release DIR")
|
||||
if len(sys.argv) == 3 and sys.argv[1] == "--check-scanner-heal-release-bundle":
|
||||
try:
|
||||
status = scanner_heal_release_bundle_status(ROOT, Path(sys.argv[2]))
|
||||
except (OSError, KeyError, TypeError, ValueError, ET.ParseError) as error:
|
||||
print(json.dumps({"schema": 1, "decision": "invalid", "release_approved": False,
|
||||
"error": str(error)}, sort_keys=True, separators=(",", ":")))
|
||||
return 2
|
||||
print(json.dumps(status, sort_keys=True, separators=(",", ":")))
|
||||
return 0 if status["release_approved"] else 1
|
||||
raise ValueError("expected --begin-scanner-heal DIR BINARY TEST_BINARY, --finish-scanner-heal DIR EXIT, --check-scanner-heal DIR CASE|release, --check-scanner-heal-release DIR, or --check-scanner-heal-release-bundle FILE")
|
||||
except (OSError, KeyError, TypeError, ValueError, subprocess.SubprocessError) as error:
|
||||
print(f"ERROR: {error}", file=sys.stderr)
|
||||
return 1
|
||||
@@ -2297,7 +2740,7 @@ def main() -> int:
|
||||
if sys.argv[1:]:
|
||||
print(
|
||||
"usage: check_test_wiring.py [--self-test | --check-core LISTING | --check-profile PROFILE LISTING | "
|
||||
"--update-profile PROFILE LISTING PLATFORM]",
|
||||
"--update-profile PROFILE LISTING PLATFORM | --check-scanner-heal-release-bundle FILE]",
|
||||
file=sys.stderr,
|
||||
)
|
||||
return 2
|
||||
|
||||
@@ -103,8 +103,8 @@ import sys
|
||||
status = json.loads(pathlib.Path(sys.argv[1]).read_text())
|
||||
if status.get("decision") != "blocked" or status.get("release_approved") is not False:
|
||||
raise SystemExit("release status did not record a blocked decision")
|
||||
if status.get("release_schema_capable") is not False:
|
||||
raise SystemExit("case-only evidence schema unexpectedly became release-capable")
|
||||
if not status.get("pending_gates"):
|
||||
raise SystemExit("release status did not retain pending gates")
|
||||
PY
|
||||
}
|
||||
|
||||
|
||||
@@ -29,6 +29,16 @@ METRICS = (
|
||||
)
|
||||
REPEATABILITY_LIMIT = Decimal("0.05")
|
||||
P2_WORK_MULTIPLE_LIMIT = Decimal("1.2")
|
||||
RELEASE_PROFILE_ARTIFACTS = (
|
||||
"allocation-profile",
|
||||
"flamegraph",
|
||||
"rss-samples",
|
||||
"save-frequency",
|
||||
)
|
||||
RELEASE_FAULT_MODES = (
|
||||
"process-restart",
|
||||
"process-crash-restart",
|
||||
)
|
||||
|
||||
|
||||
def require(condition, message):
|
||||
@@ -140,6 +150,101 @@ def validate_manifest(manifest):
|
||||
number(manifest["expected_healed_objects"].get(scenario), f"{scenario} expected repairs")
|
||||
if scenario in ("running-heal", "mrf-replay"):
|
||||
require(manifest["expected_healed_objects"][scenario] > 0, f"{scenario} requires repairs")
|
||||
validate_release_evidence_manifest(manifest)
|
||||
|
||||
|
||||
def release_evidence_integer(value, name, minimum=1, maximum=1024):
|
||||
require(type(value) is int and minimum <= value <= maximum, f"invalid release_evidence.{name}")
|
||||
return value
|
||||
|
||||
|
||||
def release_evidence_string(value, name):
|
||||
require(isinstance(value, str) and value.strip(), f"missing release_evidence.{name}")
|
||||
return value
|
||||
|
||||
|
||||
def release_evidence_bool(value, name):
|
||||
require(type(value) is bool, f"invalid release_evidence.{name}")
|
||||
return value
|
||||
|
||||
|
||||
def release_evidence_true(value, name):
|
||||
release_evidence_bool(value, name)
|
||||
require(value is True, f"missing release_evidence.{name}")
|
||||
|
||||
|
||||
def validate_release_evidence_manifest(manifest):
|
||||
if manifest["evidence"] != "measured":
|
||||
return
|
||||
|
||||
evidence = manifest.get("release_evidence")
|
||||
require(isinstance(evidence, dict), "missing release_evidence for measured ABBA")
|
||||
|
||||
topology = evidence.get("topology")
|
||||
require(isinstance(topology, dict), "missing release_evidence.topology")
|
||||
nodes = release_evidence_integer(topology.get("nodes"), "topology.nodes", 3, 64)
|
||||
drives = release_evidence_integer(topology.get("drives_per_node"), "topology.drives_per_node", 1, 64)
|
||||
set_size = release_evidence_integer(topology.get("erasure_set_size"), "topology.erasure_set_size", 12, 12)
|
||||
data = release_evidence_integer(topology.get("erasure_data_blocks"), "topology.erasure_data_blocks", 8, 8)
|
||||
parity = release_evidence_integer(topology.get("erasure_parity_blocks"), "topology.erasure_parity_blocks", 4, 4)
|
||||
require(data + parity == set_size, "release_evidence.topology must be EC8+4")
|
||||
require(nodes * drives >= set_size, "release_evidence.topology cannot host one EC8+4 set")
|
||||
pools = release_evidence_integer(topology.get("pools"), "topology.pools", 1)
|
||||
sets_total = release_evidence_integer(topology.get("sets_total"), "topology.sets_total", 1)
|
||||
sampled_pools = release_evidence_integer(topology.get("sampled_pools"), "topology.sampled_pools", 2)
|
||||
sampled_sets = release_evidence_integer(topology.get("sampled_sets"), "topology.sampled_sets", 2)
|
||||
require(sampled_pools <= pools, "release_evidence.topology sampled pools exceed total pools")
|
||||
require(sampled_sets <= sets_total, "release_evidence.topology sampled sets exceed total sets")
|
||||
|
||||
distributed = evidence.get("distributed")
|
||||
require(isinstance(distributed, dict), "missing release_evidence.distributed")
|
||||
endpoints = distributed.get("metrics_endpoints")
|
||||
require(isinstance(endpoints, list) and len(endpoints) >= nodes, "missing release_evidence.distributed.metrics_endpoints")
|
||||
require(
|
||||
all(isinstance(endpoint, str) and endpoint.strip() for endpoint in endpoints)
|
||||
and len(set(endpoints)) == len(endpoints),
|
||||
"invalid release_evidence.distributed.metrics_endpoints",
|
||||
)
|
||||
release_evidence_string(distributed.get("failure_domain"), "distributed.failure_domain")
|
||||
release_evidence_true(distributed.get("same_window_sampling"), "distributed.same_window_sampling")
|
||||
|
||||
crash = evidence.get("crash_restart")
|
||||
require(isinstance(crash, dict), "missing release_evidence.crash_restart")
|
||||
fault_modes = crash.get("fault_modes")
|
||||
require(
|
||||
isinstance(fault_modes, list)
|
||||
and all(mode in fault_modes for mode in RELEASE_FAULT_MODES)
|
||||
and all(isinstance(mode, str) and mode.strip() for mode in fault_modes),
|
||||
"missing release_evidence.crash_restart.fault_modes",
|
||||
)
|
||||
release_evidence_true(crash.get("unclean_shutdown_marker"), "crash_restart.unclean_shutdown_marker")
|
||||
|
||||
mixed = evidence.get("mixed_version")
|
||||
require(isinstance(mixed, dict), "missing release_evidence.mixed_version")
|
||||
revisions = mixed.get("participating_revisions")
|
||||
require(
|
||||
isinstance(revisions, list)
|
||||
and len(set(revisions)) >= 2
|
||||
and all(isinstance(revision, str) and len(revision) == 40 and all(c in "0123456789abcdef" for c in revision)
|
||||
for revision in revisions),
|
||||
"invalid release_evidence.mixed_version.participating_revisions",
|
||||
)
|
||||
for revision in (manifest["baseline"]["revision"], manifest["candidate"]["revision"]):
|
||||
require(revision in revisions, "release_evidence.mixed_version omits tested build revision")
|
||||
for key in ("reader", "writer", "rollback_payload"):
|
||||
require(mixed.get(key) is True, f"missing release_evidence.mixed_version.{key}")
|
||||
|
||||
profile = evidence.get("profile")
|
||||
require(isinstance(profile, dict), "missing release_evidence.profile")
|
||||
artifacts = profile.get("required_artifacts")
|
||||
require(
|
||||
isinstance(artifacts, list)
|
||||
and all(item in artifacts for item in RELEASE_PROFILE_ARTIFACTS)
|
||||
and all(isinstance(item, str) and item.strip() for item in artifacts),
|
||||
"missing release_evidence.profile.required_artifacts",
|
||||
)
|
||||
for key in ("collector_config_sha256", "profiler_config_sha256"):
|
||||
require(sha(profile.get(key)), f"invalid release_evidence.profile.{key}")
|
||||
|
||||
|
||||
class OwnedCommand:
|
||||
@@ -254,6 +359,8 @@ def validate_result(result, request, expected):
|
||||
require(result.get("build") == request["build"], "deployed build provenance mismatch")
|
||||
require(result.get("data_dir") == request["data_dir"], "adapter data isolation mismatch")
|
||||
require(result.get("background") == request["background"], "background mode mismatch")
|
||||
if request["evidence"] == "measured":
|
||||
require(result.get("release_evidence") == request["release_evidence"], "release evidence provenance mismatch")
|
||||
require(type(result.get("sample_count")) is int and 1 <= result["sample_count"] <= 3600,
|
||||
"sample_count must be 1..3600")
|
||||
number(result.get("elapsed_seconds"), "elapsed_seconds", request["duration_seconds"])
|
||||
@@ -445,6 +552,9 @@ def collect_live(prepared, request, request_path, adapter):
|
||||
connection = prepared["collector"]
|
||||
require(set(connection) == {"alias", "endpoint", "metrics_endpoints"}, "invalid collector connection")
|
||||
require(all(isinstance(value, str) and value for value in connection.values()), "missing collector endpoint")
|
||||
expected_metrics_endpoints = None
|
||||
if request.get("evidence") == "measured":
|
||||
expected_metrics_endpoints = request["release_evidence"]["distributed"]["metrics_endpoints"]
|
||||
output = request_path.parent / "telemetry"
|
||||
args = ["bash", str(collector), "--alias", connection["alias"], "--endpoint", connection["endpoint"],
|
||||
"--metrics-endpoints", connection["metrics_endpoints"], "--deployment", "distributed",
|
||||
@@ -470,6 +580,9 @@ def collect_live(prepared, request, request_path, adapter):
|
||||
require(isinstance(status.get("healOperations"), dict) and status["healOperations"], "invalid heal status response")
|
||||
metrics = list((output / "metrics").glob("admin-metrics.*.ndjson"))
|
||||
endpoints = [endpoint for endpoint in connection["metrics_endpoints"].split(",") if endpoint]
|
||||
if expected_metrics_endpoints is not None:
|
||||
require(endpoints == expected_metrics_endpoints,
|
||||
"collector metrics endpoints do not match release evidence")
|
||||
require(metrics and len(metrics) == len(endpoints) * len(samples), "missing distributed metrics samples")
|
||||
for sample in metrics:
|
||||
# The collector requests n=1, so each file contains one final JSON record.
|
||||
@@ -518,6 +631,8 @@ def run(manifest, adapter, output, data_root):
|
||||
"duration_seconds": manifest["duration_seconds"], "data_dir": str(data_dir),
|
||||
"expected_healed_objects": manifest["expected_healed_objects"][scenario],
|
||||
"expected_oracle": manifest["oracles"][scenario]}
|
||||
if manifest["evidence"] == "measured":
|
||||
request["release_evidence"] = manifest["release_evidence"]
|
||||
require(digest(Path(request["build"]["binary"])) == request["build"]["sha256"], "binary changed during run")
|
||||
require(digest(adapter) == manifest["adapter_sha256"], "adapter changed during run")
|
||||
require(shutil.disk_usage(data_root).free >= manifest["min_free_bytes"], "insufficient free disk space")
|
||||
|
||||
@@ -12,6 +12,8 @@ from pathlib import Path
|
||||
import sys
|
||||
from typing import Any
|
||||
|
||||
from scanner_abba import LEGS, SCENARIOS, validate_release_evidence_manifest
|
||||
|
||||
MAX_JSON_BYTES = 1024 * 1024
|
||||
CACHE_COST_PREFIX = "CACHE_COST "
|
||||
PASS_STATES = {"pass"}
|
||||
@@ -77,12 +79,89 @@ def max_decimal(values: list[Decimal | None]) -> Decimal | None:
|
||||
return max(present)
|
||||
|
||||
|
||||
def require_metric_series(value: Any, name: str, minimum: Decimal | None = None,
|
||||
maximum: Decimal | None = None) -> list[Decimal | None]:
|
||||
require(isinstance(value, list) and value, f"missing performance evidence field: {name}")
|
||||
parsed = [maybe_number(item, name) for item in value]
|
||||
for item in parsed:
|
||||
if item is None:
|
||||
continue
|
||||
if minimum is not None:
|
||||
require(item >= minimum, f"{name} below minimum")
|
||||
if maximum is not None:
|
||||
require(item <= maximum, f"{name} above maximum")
|
||||
return parsed
|
||||
|
||||
|
||||
def require_measured_comparison_evidence(comparison: dict[str, Any], index: int) -> None:
|
||||
w10_w11 = comparison.get("w10_w11")
|
||||
require(isinstance(w10_w11, dict), f"comparison {index} missing W10/W11 evidence")
|
||||
pressure = require_metric_series(
|
||||
w10_w11.get("foreground_pressure_high_sample_ratios"),
|
||||
f"comparison {index} foreground_pressure_high_sample_ratios",
|
||||
Decimal("0"),
|
||||
Decimal("1"),
|
||||
)
|
||||
lock_wait = require_metric_series(
|
||||
w10_w11.get("heal_lock_wait_p99_ms"),
|
||||
f"comparison {index} heal_lock_wait_p99_ms",
|
||||
Decimal("0"),
|
||||
)
|
||||
attempt_cost = require_metric_series(
|
||||
w10_w11.get("attempt_cost_per_healed_object"),
|
||||
f"comparison {index} attempt_cost_per_healed_object",
|
||||
Decimal("0"),
|
||||
)
|
||||
require(len(pressure) == len(lock_wait) == len(attempt_cost),
|
||||
f"comparison {index} W10/W11 evidence length mismatch")
|
||||
candidate_attempt_cost = maybe_number(
|
||||
w10_w11.get("candidate_attempt_cost_per_healed_object"),
|
||||
f"comparison {index} candidate_attempt_cost_per_healed_object",
|
||||
)
|
||||
require(candidate_attempt_cost is None or candidate_attempt_cost >= 0,
|
||||
f"comparison {index} candidate attempt cost below minimum")
|
||||
|
||||
|
||||
def require_complete_abba_matrix(manifest: dict[str, Any], report: dict[str, Any], comparisons: list[dict[str, Any]]) -> None:
|
||||
require(report.get("evidence") == manifest.get("evidence"), "manifest/report evidence mismatch")
|
||||
rounds = manifest.get("rounds")
|
||||
require(type(rounds) is int and 3 <= rounds <= 10, "invalid manifest.rounds")
|
||||
expected_cells = len(SCENARIOS) * 2 * rounds * len(LEGS)
|
||||
require(
|
||||
report.get("cells") == expected_cells,
|
||||
f"ABBA matrix cell count mismatch: expected {expected_cells}, got {report.get('cells')}",
|
||||
)
|
||||
expected_keys = {
|
||||
(scenario, comparison, round_id)
|
||||
for scenario in SCENARIOS
|
||||
for comparison in ("build", "background")
|
||||
for round_id in range(1, rounds + 1)
|
||||
}
|
||||
observed_keys = []
|
||||
for index, comparison in enumerate(comparisons):
|
||||
key = (comparison.get("scenario"), comparison.get("comparison"), comparison.get("round"))
|
||||
require(key in expected_keys, f"comparison {index} is outside the ABBA matrix")
|
||||
require(comparison.get("status") in PASS_STATES, f"comparison {index} did not pass")
|
||||
observed_keys.append(key)
|
||||
observed_set = set(observed_keys)
|
||||
require(len(observed_keys) == len(observed_set), "duplicate ABBA matrix comparison")
|
||||
missing = sorted(expected_keys - observed_set)
|
||||
require(not missing, f"missing ABBA matrix comparison: {missing[0] if missing else ''}")
|
||||
|
||||
|
||||
def summarize_abba(abba_dir: Path) -> dict[str, Any]:
|
||||
manifest_path = abba_dir / "manifest.json"
|
||||
report_path = abba_dir / "report.json"
|
||||
manifest = read_json(manifest_path)
|
||||
report = read_json(report_path)
|
||||
report_state = report.get("status")
|
||||
performance_state = report.get("performance")
|
||||
require(isinstance(report_state, str) and report_state, "report.status missing")
|
||||
require(isinstance(performance_state, str) and performance_state, "report.performance missing")
|
||||
comparisons = report.get("comparisons")
|
||||
if comparisons is None:
|
||||
require(report_state not in PASS_STATES, "passing report requires comparisons")
|
||||
comparisons = []
|
||||
require(isinstance(comparisons, list), "report.comparisons must be a list")
|
||||
|
||||
counts = Counter()
|
||||
@@ -115,15 +194,18 @@ def summarize_abba(abba_dir: Path) -> dict[str, Any]:
|
||||
if isinstance(p2, list):
|
||||
p2_values.extend(maybe_number(value, "p2_post_stop_work_multiple") for value in p2)
|
||||
|
||||
report_state = report.get("status")
|
||||
performance_state = report.get("performance")
|
||||
require(isinstance(report_state, str) and report_state, "report.status missing")
|
||||
require(isinstance(performance_state, str) and performance_state, "report.performance missing")
|
||||
measured = report.get("evidence") == "measured"
|
||||
passed = report_state in PASS_STATES and performance_state in PASS_STATES and measured
|
||||
if passed:
|
||||
require_complete_abba_matrix(manifest, report, comparisons)
|
||||
validate_release_evidence_manifest({**manifest, "evidence": "measured"})
|
||||
for index, comparison in enumerate(comparisons):
|
||||
require_measured_comparison_evidence(comparison, index)
|
||||
gate_state = "pass" if passed else "fail"
|
||||
if report_state == "synthetic_validated":
|
||||
reason = "synthetic evidence validates the harness only; measured performance remains pending"
|
||||
elif report_state in FAIL_STATES and isinstance(report.get("error"), str) and report["error"]:
|
||||
reason = f"ABBA report status is {report_state}: {report['error']}"
|
||||
elif report_state not in PASS_STATES:
|
||||
reason = f"ABBA report status is {report_state}"
|
||||
elif performance_state not in PASS_STATES:
|
||||
@@ -142,6 +224,8 @@ def summarize_abba(abba_dir: Path) -> dict[str, Any]:
|
||||
"performance": performance_state,
|
||||
"evidence": report.get("evidence"),
|
||||
"cells": report.get("cells", 0),
|
||||
"completed_cells": report.get("completed_cells"),
|
||||
"error": report.get("error"),
|
||||
"comparisons_total": len(comparisons),
|
||||
"comparison_status_counts": dict(sorted(counts.items())),
|
||||
"worst_p99_regression": None if not p99_regressions else float(max(p99_regressions)),
|
||||
@@ -164,6 +248,7 @@ def summarize_abba(abba_dir: Path) -> dict[str, Any]:
|
||||
"durability": fixed.get("durability"),
|
||||
"topology": fixed.get("topology"),
|
||||
"offered_load_ops": fixed.get("offered_load_ops"),
|
||||
"release_evidence": manifest.get("release_evidence"),
|
||||
},
|
||||
}
|
||||
|
||||
@@ -247,6 +332,10 @@ def markdown(summary: dict[str, Any]) -> str:
|
||||
f"- worst_throughput_loss: {pct(throughput)}",
|
||||
f"- p2_worst_post_stop_work_multiple: {ratio(p2)}",
|
||||
]
|
||||
if abba.get("completed_cells") is not None:
|
||||
lines.append(f"- completed_cells: {abba['completed_cells']}")
|
||||
if abba.get("error"):
|
||||
lines.append(f"- error: {abba['error']}")
|
||||
if summary.get("cache_cost") is not None:
|
||||
cache = summary["cache_cost"]
|
||||
lines.extend([
|
||||
|
||||
@@ -49,6 +49,8 @@ def fake_adapter():
|
||||
if fault == "measure-exit":
|
||||
return 42
|
||||
result = {key: request[key] for key in ("evidence", "fixed", "build", "data_dir", "background")}
|
||||
if request["evidence"] == "measured":
|
||||
result["release_evidence"] = copy.deepcopy(request["release_evidence"])
|
||||
result.update({"sample_count": 10, "elapsed_seconds": request["duration_seconds"],
|
||||
"metrics": dict.fromkeys(harness.METRICS, 10)})
|
||||
baseline = request["comparison"] == "build" and request["leg"].startswith("A")
|
||||
@@ -163,6 +165,45 @@ class ScannerAbbaTest(unittest.TestCase):
|
||||
build = {"binary": str(self.binary), "sha256": harness.digest(self.binary), "revision": "a" * 40}
|
||||
self.manifest.update(baseline=build.copy(), candidate=build.copy())
|
||||
|
||||
def measured_manifest(self):
|
||||
manifest = copy.deepcopy(self.manifest)
|
||||
manifest.update(evidence="measured", duration_seconds=900)
|
||||
manifest["candidate"]["revision"] = "b" * 40
|
||||
manifest["release_evidence"] = {
|
||||
"topology": {
|
||||
"nodes": 3,
|
||||
"drives_per_node": 4,
|
||||
"pools": 2,
|
||||
"sets_total": 2,
|
||||
"sampled_pools": 2,
|
||||
"sampled_sets": 2,
|
||||
"erasure_set_size": 12,
|
||||
"erasure_data_blocks": 8,
|
||||
"erasure_parity_blocks": 4,
|
||||
},
|
||||
"distributed": {
|
||||
"metrics_endpoints": ["https://node-1:9000", "https://node-2:9000", "https://node-3:9000"],
|
||||
"failure_domain": "three-node-localhost-lab",
|
||||
"same_window_sampling": True,
|
||||
},
|
||||
"crash_restart": {
|
||||
"fault_modes": ["process-restart", "process-crash-restart"],
|
||||
"unclean_shutdown_marker": True,
|
||||
},
|
||||
"mixed_version": {
|
||||
"participating_revisions": ["a" * 40, "b" * 40],
|
||||
"reader": True,
|
||||
"writer": True,
|
||||
"rollback_payload": True,
|
||||
},
|
||||
"profile": {
|
||||
"required_artifacts": ["allocation-profile", "flamegraph", "rss-samples", "save-frequency"],
|
||||
"collector_config_sha256": "4" * 64,
|
||||
"profiler_config_sha256": "5" * 64,
|
||||
},
|
||||
}
|
||||
return manifest
|
||||
|
||||
def run_harness(self, fault=""):
|
||||
with patch.dict(os.environ, {"SCANNER_ABBA_TEST_FAULT": fault}), contextlib.redirect_stdout(io.StringIO()):
|
||||
return harness.run(copy.deepcopy(self.manifest), self.adapter, self.root / "out", self.root / "data")
|
||||
@@ -440,6 +481,34 @@ class ScannerAbbaTest(unittest.TestCase):
|
||||
|
||||
process.finish.assert_called_once_with(terminate=True)
|
||||
|
||||
def test_live_collector_binds_release_evidence_metrics_endpoints(self):
|
||||
telemetry = self.root / "telemetry"
|
||||
for name in ("status", "heal", "metrics"):
|
||||
(telemetry / name).mkdir(parents=True)
|
||||
(telemetry / "scanner-summary.csv").write_text("timestamp\n")
|
||||
for index in range(16):
|
||||
harness.write_json(telemetry / f"status/scanner-status.{index}.json", {"metrics": {"objects": 10}})
|
||||
for node in ("node-a", "node-b"):
|
||||
harness.write_json(telemetry / f"heal/background-heal-status.{node}.{index}.json",
|
||||
{"healOperations": {"queueLength": 0}})
|
||||
harness.write_json(telemetry / f"metrics/admin-metrics.{node}.{index}.ndjson",
|
||||
{"errors": [], "final": True,
|
||||
"by_host": {f"{node}:9000": {"scanner": {"objects": 10}}}})
|
||||
prepared = {"collector": {"alias": "test", "endpoint": "http://node-a:9000",
|
||||
"metrics_endpoints": "http://node-a:9000,http://node-b:9000"}}
|
||||
request = {
|
||||
"duration_seconds": 900,
|
||||
"evidence": "measured",
|
||||
"release_evidence": self.measured_manifest()["release_evidence"],
|
||||
}
|
||||
process = Mock(pid=123, wait=Mock(return_value=0))
|
||||
with patch.object(harness, "OwnedCommand", return_value=process), \
|
||||
patch.object(harness, "invoke", return_value={"sample_count": 10}), \
|
||||
patch.object(harness.time, "monotonic", side_effect=(0, 900)):
|
||||
with self.assertRaisesRegex(ValueError, "collector metrics endpoints"):
|
||||
harness.collect_live(prepared, request, self.root / "request.json", self.adapter)
|
||||
process.finish.assert_called_once_with(terminate=True)
|
||||
|
||||
def test_unstable_p1_work_control_is_inconclusive(self):
|
||||
with patch.object(harness, "SCENARIOS", ("cold-hot",)):
|
||||
self.assertEqual(self.run_harness("unstable-p1-control"), 3)
|
||||
@@ -463,6 +532,81 @@ class ScannerAbbaTest(unittest.TestCase):
|
||||
with self.assertRaisesRegex(ValueError, "rounds"):
|
||||
harness.validate_manifest(self.manifest)
|
||||
|
||||
def test_measured_manifest_requires_release_evidence_contract(self):
|
||||
harness.validate_manifest(self.measured_manifest())
|
||||
faults = {
|
||||
"missing root": lambda manifest: manifest.pop("release_evidence"),
|
||||
"single-set": lambda manifest: manifest["release_evidence"]["topology"].update(sets_total=1),
|
||||
"unsampled-set": lambda manifest: manifest["release_evidence"]["topology"].update(sampled_sets=1),
|
||||
"wrong geometry": lambda manifest: manifest["release_evidence"]["topology"].update(erasure_set_size=11),
|
||||
"duplicate endpoint": lambda manifest: manifest["release_evidence"]["distributed"].update(
|
||||
metrics_endpoints=["https://node-1:9000", "https://node-1:9000", "https://node-3:9000"],
|
||||
),
|
||||
"split sampling": lambda manifest: manifest["release_evidence"]["distributed"].update(
|
||||
same_window_sampling=False,
|
||||
),
|
||||
"missing crash": lambda manifest: manifest["release_evidence"]["crash_restart"].update(
|
||||
fault_modes=["process-restart"],
|
||||
),
|
||||
"clean crash marker": lambda manifest: manifest["release_evidence"]["crash_restart"].update(
|
||||
unclean_shutdown_marker=False,
|
||||
),
|
||||
"mixed version false": lambda manifest: manifest["release_evidence"]["mixed_version"].update(writer=False),
|
||||
"missing candidate": lambda manifest: manifest["release_evidence"]["mixed_version"].update(
|
||||
participating_revisions=["a" * 40, "c" * 40],
|
||||
),
|
||||
"missing profile": lambda manifest: manifest["release_evidence"]["profile"].update(
|
||||
required_artifacts=["allocation-profile", "flamegraph", "rss-samples"],
|
||||
),
|
||||
"bad profile hash": lambda manifest: manifest["release_evidence"]["profile"].update(
|
||||
profiler_config_sha256="not-a-sha",
|
||||
),
|
||||
}
|
||||
for name, mutate in faults.items():
|
||||
with self.subTest(fault=name):
|
||||
manifest = self.measured_manifest()
|
||||
mutate(manifest)
|
||||
with self.assertRaisesRegex(ValueError, "release_evidence"):
|
||||
harness.validate_manifest(manifest)
|
||||
|
||||
def test_measured_result_must_echo_release_evidence(self):
|
||||
manifest = self.measured_manifest()
|
||||
request = {
|
||||
"schema": 1,
|
||||
"scenario": "cold-hot",
|
||||
"comparison": "build",
|
||||
"round": 1,
|
||||
"leg": "B1",
|
||||
"background": "on",
|
||||
"build": manifest["candidate"],
|
||||
"evidence": manifest["evidence"],
|
||||
"fixed": manifest["fixed"],
|
||||
"release_evidence": manifest["release_evidence"],
|
||||
"duration_seconds": manifest["duration_seconds"],
|
||||
"data_dir": str(self.root / "data"),
|
||||
"expected_healed_objects": manifest["expected_healed_objects"]["cold-hot"],
|
||||
}
|
||||
metrics = dict.fromkeys(harness.METRICS, 10)
|
||||
metrics.update(p99_ms=10, throughput_ops=100, errors=0, requests=100,
|
||||
walk_objects=100, cold_walk_objects=20, healed_objects=10)
|
||||
result = {
|
||||
"evidence": request["evidence"],
|
||||
"fixed": request["fixed"],
|
||||
"build": request["build"],
|
||||
"data_dir": request["data_dir"],
|
||||
"background": request["background"],
|
||||
"release_evidence": request["release_evidence"],
|
||||
"sample_count": 10,
|
||||
"elapsed_seconds": request["duration_seconds"],
|
||||
"metrics": metrics,
|
||||
"oracle": manifest["oracles"]["cold-hot"],
|
||||
}
|
||||
harness.validate_result(result, request, manifest["oracles"]["cold-hot"])
|
||||
result["release_evidence"] = copy.deepcopy(result["release_evidence"])
|
||||
result["release_evidence"]["profile"]["required_artifacts"].remove("flamegraph")
|
||||
with self.assertRaisesRegex(ValueError, "release evidence provenance mismatch"):
|
||||
harness.validate_result(result, request, manifest["oracles"]["cold-hot"])
|
||||
|
||||
def test_existing_data_preserved(self):
|
||||
(self.root / "data").mkdir()
|
||||
marker = self.root / "data/keep"
|
||||
|
||||
@@ -2,6 +2,7 @@
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import copy
|
||||
import contextlib
|
||||
import hashlib
|
||||
import io
|
||||
@@ -33,6 +34,9 @@ class ScannerHealPerfSummaryTest(unittest.TestCase):
|
||||
self.abba = self.root / "abba"
|
||||
self.abba.mkdir()
|
||||
self.manifest = {
|
||||
"schema": 1,
|
||||
"evidence": "measured",
|
||||
"rounds": 3,
|
||||
"fixed": {
|
||||
"config_sha256": "1" * 64,
|
||||
"dataset_sha256": "2" * 64,
|
||||
@@ -45,6 +49,39 @@ class ScannerHealPerfSummaryTest(unittest.TestCase):
|
||||
"candidate": {"revision": "b" * 40, "sha256": "4" * 64},
|
||||
"adapter_sha256": "5" * 64,
|
||||
"collector_sha256": "6" * 64,
|
||||
"release_evidence": {
|
||||
"topology": {
|
||||
"nodes": 3,
|
||||
"drives_per_node": 4,
|
||||
"pools": 2,
|
||||
"sets_total": 2,
|
||||
"sampled_pools": 2,
|
||||
"sampled_sets": 2,
|
||||
"erasure_set_size": 12,
|
||||
"erasure_data_blocks": 8,
|
||||
"erasure_parity_blocks": 4,
|
||||
},
|
||||
"distributed": {
|
||||
"metrics_endpoints": ["https://node-1:9000", "https://node-2:9000", "https://node-3:9000"],
|
||||
"failure_domain": "three-node-localhost-lab",
|
||||
"same_window_sampling": True,
|
||||
},
|
||||
"crash_restart": {
|
||||
"fault_modes": ["process-restart", "process-crash-restart"],
|
||||
"unclean_shutdown_marker": True,
|
||||
},
|
||||
"mixed_version": {
|
||||
"participating_revisions": ["a" * 40, "b" * 40],
|
||||
"reader": True,
|
||||
"writer": True,
|
||||
"rollback_payload": True,
|
||||
},
|
||||
"profile": {
|
||||
"required_artifacts": ["allocation-profile", "flamegraph", "rss-samples", "save-frequency"],
|
||||
"collector_config_sha256": "7" * 64,
|
||||
"profiler_config_sha256": "8" * 64,
|
||||
},
|
||||
},
|
||||
}
|
||||
self.comparison = {
|
||||
"scenario": "cold-hot",
|
||||
@@ -55,16 +92,32 @@ class ScannerHealPerfSummaryTest(unittest.TestCase):
|
||||
"throughput_change": -0.01,
|
||||
"p1": {"required_reduction": 0.8, "observed_reduction": 0.82, "repeatability_drift": 0.01},
|
||||
"p2_post_stop_work_multiples": [None, 1.1, 1.0, None],
|
||||
"w10_w11": {
|
||||
"foreground_pressure_high_sample_ratios": [0.0, 0.25, 0.25, 0.0],
|
||||
"heal_lock_wait_p99_ms": [12.0, 8.0, 9.0, 13.0],
|
||||
"attempt_cost_per_healed_object": [None, 1.2, 1.3, None],
|
||||
"candidate_attempt_cost_per_healed_object": 1.3,
|
||||
},
|
||||
}
|
||||
self.report = {
|
||||
"status": "pass",
|
||||
"performance": "pass",
|
||||
"evidence": "measured",
|
||||
"cells": 120,
|
||||
"comparisons": [self.comparison],
|
||||
"comparisons": self.full_comparisons(),
|
||||
}
|
||||
self.write_inputs()
|
||||
|
||||
def full_comparisons(self):
|
||||
comparisons = []
|
||||
for scenario in summary.SCENARIOS:
|
||||
for comparison in ("build", "background"):
|
||||
for round_id in range(1, 4):
|
||||
row = copy.deepcopy(self.comparison)
|
||||
row.update(scenario=scenario, comparison=comparison, round=round_id)
|
||||
comparisons.append(row)
|
||||
return comparisons
|
||||
|
||||
def write_inputs(self):
|
||||
(self.abba / "manifest.json").write_text(json.dumps(self.manifest), encoding="utf-8")
|
||||
(self.abba / "report.json").write_text(json.dumps(self.report), encoding="utf-8")
|
||||
@@ -112,6 +165,107 @@ class ScannerHealPerfSummaryTest(unittest.TestCase):
|
||||
self.assertEqual(result["verdict"], "FAIL")
|
||||
self.assertIn("synthetic evidence", result["reason"])
|
||||
|
||||
def test_failed_abba_report_without_comparisons_writes_fail_closed_summary(self):
|
||||
self.report = {
|
||||
"status": "failed",
|
||||
"performance": "pending",
|
||||
"completed_cells": 7,
|
||||
"error": "collector failed",
|
||||
}
|
||||
self.write_inputs()
|
||||
args = type("Args", (), {
|
||||
"abba_dir": self.abba,
|
||||
"cache_cost_log": None,
|
||||
"require_cache_cost": False,
|
||||
"json_out": None,
|
||||
"markdown_out": None,
|
||||
})
|
||||
result = summary.build_summary(args)
|
||||
self.assertEqual(result["verdict"], "FAIL")
|
||||
self.assertEqual(result["abba"]["completed_cells"], 7)
|
||||
self.assertEqual(result["abba"]["comparisons_total"], 0)
|
||||
self.assertIn("collector failed", result["reason"])
|
||||
self.assertIn("- completed_cells: 7", summary.markdown(result))
|
||||
self.assertIn("- error: collector failed", summary.markdown(result))
|
||||
|
||||
def test_passing_abba_report_requires_comparisons(self):
|
||||
del self.report["comparisons"]
|
||||
self.write_inputs()
|
||||
args = type("Args", (), {
|
||||
"abba_dir": self.abba,
|
||||
"cache_cost_log": None,
|
||||
"require_cache_cost": False,
|
||||
"json_out": None,
|
||||
"markdown_out": None,
|
||||
})
|
||||
with self.assertRaisesRegex(ValueError, "passing report requires comparisons"):
|
||||
summary.build_summary(args)
|
||||
|
||||
def test_passing_abba_report_requires_complete_matrix(self):
|
||||
cases = {
|
||||
"trimmed": lambda: self.report["comparisons"].pop(),
|
||||
"duplicate": lambda: self.report["comparisons"].__setitem__(1, copy.deepcopy(self.report["comparisons"][0])),
|
||||
"bad cells": lambda: self.report.update(cells=119),
|
||||
"bad evidence": lambda: self.manifest.update(evidence="synthetic"),
|
||||
"outside": lambda: self.report["comparisons"][0].update(round=99),
|
||||
"failed comparison": lambda: self.report["comparisons"][0].update(status="inconclusive"),
|
||||
}
|
||||
for name, mutate in cases.items():
|
||||
with self.subTest(fault=name):
|
||||
self.setUp()
|
||||
mutate()
|
||||
self.write_inputs()
|
||||
args = type("Args", (), {
|
||||
"abba_dir": self.abba,
|
||||
"cache_cost_log": None,
|
||||
"require_cache_cost": False,
|
||||
"json_out": None,
|
||||
"markdown_out": None,
|
||||
})
|
||||
with self.assertRaisesRegex(ValueError, "ABBA matrix|manifest/report evidence|comparison"):
|
||||
summary.build_summary(args)
|
||||
|
||||
def test_passing_abba_report_requires_w10_w11_evidence(self):
|
||||
for fault in ("missing", "pressure", "lock", "attempt", "length", "range"):
|
||||
with self.subTest(fault=fault):
|
||||
self.setUp()
|
||||
target = self.report["comparisons"][0]
|
||||
if fault == "missing":
|
||||
del target["w10_w11"]
|
||||
elif fault == "pressure":
|
||||
del target["w10_w11"]["foreground_pressure_high_sample_ratios"]
|
||||
elif fault == "lock":
|
||||
del target["w10_w11"]["heal_lock_wait_p99_ms"]
|
||||
elif fault == "attempt":
|
||||
del target["w10_w11"]["attempt_cost_per_healed_object"]
|
||||
elif fault == "length":
|
||||
target["w10_w11"]["attempt_cost_per_healed_object"] = [None]
|
||||
else:
|
||||
target["w10_w11"]["foreground_pressure_high_sample_ratios"] = [1.5, 0.0, 0.0, 0.0]
|
||||
self.write_inputs()
|
||||
args = type("Args", (), {
|
||||
"abba_dir": self.abba,
|
||||
"cache_cost_log": None,
|
||||
"require_cache_cost": False,
|
||||
"json_out": None,
|
||||
"markdown_out": None,
|
||||
})
|
||||
with self.assertRaisesRegex(ValueError, "W10/W11|performance evidence|length mismatch|above maximum"):
|
||||
summary.build_summary(args)
|
||||
|
||||
def test_passing_measured_report_requires_release_evidence_manifest(self):
|
||||
del self.manifest["release_evidence"]
|
||||
self.write_inputs()
|
||||
args = type("Args", (), {
|
||||
"abba_dir": self.abba,
|
||||
"cache_cost_log": None,
|
||||
"require_cache_cost": False,
|
||||
"json_out": None,
|
||||
"markdown_out": None,
|
||||
})
|
||||
with self.assertRaisesRegex(ValueError, "release_evidence"):
|
||||
summary.build_summary(args)
|
||||
|
||||
def test_requires_cache_profile_when_requested(self):
|
||||
args = type("Args", (), {
|
||||
"abba_dir": self.abba,
|
||||
|
||||
Reference in New Issue
Block a user