From 727193d3f4093ff49dd8ed7668e20f1100b09a7e Mon Sep 17 00:00:00 2001 From: grumbach Date: Thu, 1 Oct 2026 15:12:55 +0900 Subject: [PATCH 1/2] fix(replication): repair chunks that audits find rotted Every audit reads chunk bytes through ChunkStore::get_raw, which does not check them against their address. A node whose chunk file has gone bad on disk therefore never notices from audits: it keeps committing the key and fails every subtree audit whose selected block contains it, and only a fetch of that exact key (the verifying get path) would take the file out of service and let replication repair it. In the other direction, an auditor whose own reference copy has gone bad fails honest peers in the responsible-chunk, possession and prune lanes. Responder: round 1 already hashes every leaf it reads, so a chunk leaf whose plain hash differs from its key is now reported in Round1Work::corrupt_keys. Once the reply has been sent and the admission permit released, the replication engine passes each one to the new ChunkStore::recheck_corrupt, which runs the existing verifying read: it re-reads the file under the shard write lane, removes it only if it is still wrong, drops the key from the view the next commitment is built from, and re-queues a legacy copy if there is one. The proof still carries the bytes that were read, so the auditor's verdict is unchanged. Auditor: the reference copy in the responsible-chunk, possession and prune lanes is now read with ChunkStore::get instead of get_raw, so a rotted local copy is taken out of service and the key skipped instead of failing the peer. A responsible audit in which no key could be checked against a good local copy is now idle rather than a pass, and keys_checked counts only the keys actually verified. Both sides follow the node's verify_on_read setting, like every other verifying read. --- src/replication/audit.rs | 207 +++++++++++++++++--- src/replication/mod.rs | 17 +- src/replication/possession.rs | 7 +- src/replication/pruning.rs | 43 +++- src/replication/storage_commitment_audit.rs | 115 ++++++++++- src/storage/chunk_store.rs | 80 +++++++- 6 files changed, 430 insertions(+), 39 deletions(-) diff --git a/src/replication/audit.rs b/src/replication/audit.rs index 9b312bb3..dd89c593 100644 --- a/src/replication/audit.rs +++ b/src/replication/audit.rs @@ -34,6 +34,8 @@ use crate::replication::config::REPAIR_HINT_MIN_AGE; #[cfg(test)] use crate::replication::types::{BootstrapClaimObservation, NeighborSyncState}; #[cfg(test)] +use crate::storage::file_store::CHUNKS_DIR_NAME; +#[cfg(test)] use crate::storage::ChunkStoreConfig; #[cfg(test)] use tempfile::TempDir; @@ -565,20 +567,62 @@ async fn verify_digests( .await; } - let challenged_peer_bytes = challenged_peer.as_bytes(); - let mut failed_keys = Vec::new(); + let DigestComparison { + failed_keys, + verified, + } = compare_with_local_copies(storage, nonce, challenged_peer.as_bytes(), keys, digests).await; + + if failed_keys.is_empty() { + return unfailed_audit_verdict(challenged_peer, keys.len(), verified); + } + + // Step 9: Responsibility confirmation for failed keys. + handle_classified_audit_failure( + challenged_peer, + challenge_id, + &failed_keys, + AuditFailureReason::DigestMismatch, + keys.len(), + None, + p2p_node, + config, + ) + .await +} - for (i, key) in keys.iter().enumerate() { - let received_digest = &digests[i]; +/// What comparing a peer's digests with this node's own copies found. +struct DigestComparison { + /// Keys the peer admitted missing, or whose digest did not match. + failed_keys: Vec, + /// Keys whose digest matched one recomputed from a good local copy. + verified: usize, +} +/// Compare per-key digests with ones recomputed from this node's own copies. +/// +/// The local copy is the reference the peer is judged against, so it is read +/// through the verifying path: a local file that no longer hashes to its key +/// must not fail an honest peer. Such a copy is taken out of service by the read +/// (replication repairs it) and the key is skipped, like any other key this node +/// cannot read. A skipped key is neither a failure nor a verification. +async fn compare_with_local_copies( + storage: &ChunkStore, + nonce: &[u8; 32], + challenged_peer_bytes: &[u8; 32], + keys: &[XorName], + digests: &[[u8; 32]], +) -> DigestComparison { + let mut failed_keys = Vec::new(); + let mut verified = 0usize; + + for (key, received_digest) in keys.iter().zip(digests) { // Check for absent sentinel. if *received_digest == ABSENT_KEY_DIGEST { failed_keys.push(AuditKeyFailure::absent(*key)); continue; } - // Recompute expected digest from local copy. - let local_bytes = match storage.get_raw(key).await { + let local_bytes = match storage.get(key).await { Ok(Some(bytes)) => bytes, Ok(None) => { // We should hold this key (we sampled it), but it's gone. @@ -595,34 +639,48 @@ async fn verify_digests( }; let expected = compute_audit_digest(nonce, challenged_peer_bytes, key, &local_bytes); - if *received_digest != expected { + if *received_digest == expected { + verified += 1; + } else { failed_keys.push(AuditKeyFailure::digest_mismatch(*key)); } } - if failed_keys.is_empty() { + DigestComparison { + failed_keys, + verified, + } +} + +/// The verdict on an audit in which no key failed. +/// +/// A pass needs at least one key actually checked: when every key was skipped for +/// want of a good local copy, nothing is known about the peer, so it earns no pass +/// (and no trust credit) and the audit is idle instead. +fn unfailed_audit_verdict( + challenged_peer: &PeerId, + key_count: usize, + verified: usize, +) -> AuditTickResult { + if verified == 0 { + debug!( + "Audit: no key could be checked against a local copy for {challenged_peer}; \ + not counted as a pass" + ); + return AuditTickResult::Idle; + } + if verified == key_count { + info!("Audit: peer {challenged_peer} passed (all {verified} keys verified)"); + } else { info!( - "Audit: peer {challenged_peer} passed (all {} keys verified)", - keys.len() + "Audit: peer {challenged_peer} passed ({verified} of {key_count} keys verified, \ + the rest skipped for want of a good local copy)" ); - return AuditTickResult::Passed { - challenged_peer: *challenged_peer, - keys_checked: keys.len(), - }; } - - // Step 9: Responsibility confirmation for failed keys. - handle_classified_audit_failure( - challenged_peer, - challenge_id, - &failed_keys, - AuditFailureReason::DigestMismatch, - keys.len(), - None, - p2p_node, - config, - ) - .await + AuditTickResult::Passed { + challenged_peer: *challenged_peer, + keys_checked: verified, + } } // --------------------------------------------------------------------------- @@ -1317,6 +1375,101 @@ mod tests { } } + // -- A rotted local reference copy ------------------------------------------- + + /// The auditor's own copy is the reference a peer is judged against. A local + /// copy that has rotted must not fail an honest peer, must not count as a + /// verification either, and is taken out of service so replication repairs it. + #[tokio::test] + async fn a_rotted_local_copy_neither_fails_nor_passes_the_peer() { + let temp_dir = TempDir::new().expect("temp dir"); + let storage = ChunkStore::new(ChunkStoreConfig { + root_dir: temp_dir.path().to_path_buf(), + ..ChunkStoreConfig::test_default() + }) + .await + .expect("storage"); + let nonce = [0x42; 32]; + let peer = PeerId::from_bytes([0x24; 32]); + let rot = |key: &XorName| { + let path = temp_dir + .path() + .join(CHUNKS_DIR_NAME) + .join(format!("{:02x}", key.last().copied().unwrap_or(0))) + .join(hex::encode(key)); + std::fs::write(path, b"rotted").expect("corrupt the file"); + }; + + let mut stored = Vec::new(); + for content in [ + &b"intact copy"[..], + b"first rotted copy", + b"second rotted copy", + ] { + let key = ChunkStore::compute_address(content); + storage.put(&key, content).await.expect("put"); + stored.push((key, content.to_vec())); + } + let [(good, good_bytes), (lone, lone_bytes), (paired, paired_bytes)] = stored.as_slice() + else { + panic!("three chunks stored"); + }; + rot(lone); + rot(paired); + let honest = |key: &XorName, content: &[u8]| { + compute_audit_digest(&nonce, peer.as_bytes(), key, content) + }; + + // Only a rotted reference: nothing can be checked, so no pass either way. + let nothing = compare_with_local_copies( + &storage, + &nonce, + peer.as_bytes(), + &[*lone], + &[honest(lone, lone_bytes)], + ) + .await; + assert!( + nothing.failed_keys.is_empty(), + "an honest peer is not failed" + ); + assert_eq!(nothing.verified, 0); + assert!(matches!( + unfailed_audit_verdict(&peer, 1, nothing.verified), + AuditTickResult::Idle + )); + assert!( + !storage.exists(lone).expect("exists"), + "the rotted local copy is taken out of service" + ); + + // A good and a rotted reference: the peer passes on what was checked. + let partial = compare_with_local_copies( + &storage, + &nonce, + peer.as_bytes(), + &[*good, *paired], + &[honest(good, good_bytes), honest(paired, paired_bytes)], + ) + .await; + assert!(partial.failed_keys.is_empty()); + assert_eq!(partial.verified, 1); + assert!(matches!( + unfailed_audit_verdict(&peer, 2, partial.verified), + AuditTickResult::Passed { + keys_checked: 1, + .. + } + )); + + // A wrong digest against a good reference still fails, as before. + let wrong = + compare_with_local_copies(&storage, &nonce, peer.as_bytes(), &[*good], &[[0xAB; 32]]) + .await; + assert_eq!(wrong.failed_keys.len(), 1); + assert_eq!(wrong.verified, 0); + } + // -- Scenario 55: Empty failure set means no evidence ------------------------- /// Scenario 55: Peer challenged on {K1, K2}. Both digests mismatch. diff --git a/src/replication/mod.rs b/src/replication/mod.rs index 73d91f45..d32020f8 100644 --- a/src/replication/mod.rs +++ b/src/replication/mod.rs @@ -4308,8 +4308,10 @@ impl fmt::Display for ResponderAdmissionFailure { /// RAII admission for one audit-responder task: holds the GLOBAL permit and, /// on drop, decrements the PER-PEER in-flight count. Moving this into the -/// spawned task ties both bounds to the task's exact lifetime — no manual -/// decrement to forget on an early return or panic. +/// spawned task ties both bounds to the task's lifetime — no manual decrement +/// to forget on an early return or panic. A task may drop it once its reply has +/// gone, before purely local follow-up work (round 1 does, before taking +/// rotted chunks out of service), so that work never holds back another audit. struct ResponderGuard { _permit: tokio::sync::OwnedSemaphorePermit, _peer_slot: PeerResponderSlot, @@ -5200,12 +5202,14 @@ async fn handle_replication_message( let responder_metrics = Arc::clone(&ctx.audit_responder_metrics); let subtree_round1 = ctx.subtree_round1.clone(); ctx.detached_task_tracker.spawn(async move { - let _guard = guard; // global permit + per-peer slot, held until done + // `guard` is the global permit + per-peer slot, held until the + // reply has gone. let worker_started = Instant::now(); let processing_started = Instant::now(); let storage_commitment_audit::Round1Work { response, content_bytes, + corrupt_keys, } = storage_commitment_audit::handle_subtree_challenge_measured_with_pointers( &challenge, &storage, @@ -5272,6 +5276,13 @@ async fn handle_replication_message( processing, response_send, ); + drop(guard); + // A committed chunk this proof found rotted is taken out of service + // only now, with the reply gone and the permit released, so the + // cleanup never delays the auditor or holds back another audit. + for key in &corrupt_keys { + storage.recheck_corrupt(key).await; + } }); Ok(()) } diff --git a/src/replication/possession.rs b/src/replication/possession.rs index cc552f2b..71171079 100644 --- a/src/replication/possession.rs +++ b/src/replication/possession.rs @@ -148,8 +148,11 @@ pub(crate) async fn run_possession_check( // Read our canonical copy once: the audit digest is recomputed from these // bytes for every peer (hoisted out of the per-peer loop). If we no longer // hold the chunk we cannot verify any peer's proof, and we are no longer a - // responsible checker for it — skip without penalising anyone. - let local_bytes = match storage.get_raw(&key).await { + // responsible checker for it — skip without penalising anyone. Read through + // the verifying path for the same reason: a local copy that no longer hashes + // to its key would fail every honest peer, so it is taken out of service by + // the read (replication repairs it) and the check is skipped. + let local_bytes = match storage.get(&key).await { Ok(Some(bytes)) => bytes, Ok(None) => { debug!("Possession check: checker no longer holds {key_hex}; skipping"); diff --git a/src/replication/pruning.rs b/src/replication/pruning.rs index 43cc63cf..0ae2012d 100644 --- a/src/replication/pruning.rs +++ b/src/replication/pruning.rs @@ -1969,8 +1969,12 @@ async fn local_record_digest( .map(|bytes| compute_audit_digest(nonce, peer.as_bytes(), key, &bytes)) } +/// The local copy a prune audit judges peers against. Read through the verifying +/// path: a local file that no longer hashes to its key would fail every honest +/// holder and hold the prune back indefinitely, so it is taken out of service by the +/// read instead, and the key is not audited. async fn local_record_bytes(key: &XorName, storage: &Arc) -> Option> { - match storage.get_raw(key).await { + match storage.get(key).await { Ok(Some(bytes)) => Some(bytes), Ok(None) => { debug!( @@ -2138,6 +2142,9 @@ async fn peer_is_currently_responsible( #[allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)] mod tests { use super::*; + use crate::storage::file_store::CHUNKS_DIR_NAME; + use crate::storage::ChunkStoreConfig; + use tempfile::TempDir; fn peer_id_from_byte(b: u8) -> PeerId { let mut bytes = [0u8; 32]; @@ -2159,6 +2166,40 @@ mod tests { RecordPruneCandidate { key, target_peers } } + /// A prune audit judges holders against this node's own copy. A copy that + /// no longer hashes to its key would fail every honest holder and hold the + /// prune back forever, so it is not used as a reference: the read takes it + /// out of service instead, and an intact copy is used as before. + #[tokio::test] + async fn a_rotted_local_copy_is_not_used_to_judge_holders() { + let dir = TempDir::new().expect("temp dir"); + let storage = Arc::new( + ChunkStore::new(ChunkStoreConfig { + root_dir: dir.path().to_path_buf(), + ..ChunkStoreConfig::test_default() + }) + .await + .expect("chunk store"), + ); + let good_bytes = b"intact prune candidate".to_vec(); + let bad_bytes = b"rotted prune candidate".to_vec(); + let good = ChunkStore::compute_address(&good_bytes); + let bad = ChunkStore::compute_address(&bad_bytes); + storage.put(&good, &good_bytes).await.expect("put good"); + storage.put(&bad, &bad_bytes).await.expect("put bad"); + let path = dir + .path() + .join(CHUNKS_DIR_NAME) + .join(format!("{:02x}", bad.last().copied().unwrap_or(0))) + .join(hex::encode(bad)); + std::fs::write(&path, b"rotted").expect("corrupt the file"); + + assert_eq!(local_record_bytes(&good, &storage).await, Some(good_bytes)); + assert_eq!(local_record_bytes(&bad, &storage).await, None); + assert!(!storage.exists(&bad).expect("exists")); + assert!(!path.exists(), "the rotted local copy is removed"); + } + #[test] fn prune_audit_challenges_are_batched_by_target_peer() { let peer_a = peer_id_from_byte(1); diff --git a/src/replication/storage_commitment_audit.rs b/src/replication/storage_commitment_audit.rs index d2835610..0ab00966 100644 --- a/src/replication/storage_commitment_audit.rs +++ b/src/replication/storage_commitment_audit.rs @@ -1360,6 +1360,17 @@ pub struct Round1Work { /// since the per-peer cooldown is escapable by rotating identity, the /// responder-wide work budget is the only bound that would have caught it. pub content_bytes: i64, + /// Committed chunks whose bytes, read for this proof, no longer hash to + /// their key. + /// + /// The proof still carries what was read, so the auditor's verdict on them + /// is unchanged. The caller hands each one to the chunk store's corruption + /// recheck once the reply has gone, so the node stops committing a chunk it + /// cannot prove and replication can repair it, rather than failing every + /// audit that lands on that key until a fetch of the same key happens to + /// notice. Done after the reply, not here, so the cleanup never delays the + /// answer the auditor is waiting for. + pub corrupt_keys: Vec, } /// [`handle_subtree_challenge`], additionally reporting the read-and-hash work @@ -1396,6 +1407,7 @@ pub async fn handle_subtree_challenge_measured_with_pointers( // exit reports its work by construction: a new early return cannot forget to // account for the reads that already happened. let mut content_bytes = 0i64; + let mut corrupt_keys = Vec::new(); let response = subtree_challenge_response( challenge, storage, @@ -1404,18 +1416,21 @@ pub async fn handle_subtree_challenge_measured_with_pointers( is_bootstrapping, commitment_state, &mut content_bytes, + &mut corrupt_keys, ) .await; Round1Work { response, content_bytes, + corrupt_keys, } } /// The round-1 responder proper. `content_bytes` accrues the chunk content read /// and hashed so far, and is meaningful on every return path, not just the -/// successful one. -#[allow(clippy::too_many_lines)] +/// successful one. `corrupt_keys` collects the committed chunks whose bytes did +/// not hash to their key (see [`Round1Work::corrupt_keys`]). +#[allow(clippy::too_many_lines, clippy::too_many_arguments)] async fn subtree_challenge_response( challenge: &SubtreeAuditChallenge, storage: &ChunkStore, @@ -1424,6 +1439,7 @@ async fn subtree_challenge_response( is_bootstrapping: bool, commitment_state: Option<&Arc>, content_bytes: &mut i64, + corrupt_keys: &mut Vec, ) -> SubtreeAuditResponse { if is_bootstrapping { return SubtreeAuditResponse::Bootstrapping { @@ -1620,6 +1636,17 @@ async fn subtree_challenge_response( }; } }; + // The leaf's plain hash is the content address of what was just read, so a + // committed chunk whose file no longer matches its name shows up here at no + // extra cost. See `Round1Work::corrupt_keys` for what happens to it. + if leaf.bytes_hash != *key { + warn!( + "Subtree audit: committed key {} does not hash to its address; \ + it will be taken out of service so replication can repair it", + hex::encode(key) + ); + corrupt_keys.push(*key); + } leaves.push(leaf); } @@ -2750,9 +2777,11 @@ mod pointer_audit_tests { use super::*; use crate::replication::commitment::MerkleTree; use crate::replication::commitment_state::BuiltCommitment; + use crate::storage::file_store::CHUNKS_DIR_NAME; use crate::storage::ChunkStoreConfig; use ant_protocol::pointer::{PointerTarget, PointerTargetKind}; use saorsa_pqc::api::sig::ml_dsa_65; + use std::path::PathBuf; use tempfile::TempDir; const CHALLENGE_ID: u64 = 7; @@ -2777,6 +2806,7 @@ mod pointer_audit_tests { state: Arc, peer: PeerId, peer_bytes: [u8; 32], + chunk_root: PathBuf, _dirs: (TempDir, TempDir), } @@ -2822,15 +2852,28 @@ mod pointer_audit_tests { state, peer: PeerId::from_bytes(peer_bytes), peer_bytes, + chunk_root: chunk_dir.path().to_path_buf(), _dirs: (chunk_dir, pointer_dir), } } + /// Where the file store keeps `key`'s bytes, so a test can rot them. + fn chunk_file(&self, key: &XorName) -> PathBuf { + self.chunk_root + .join(CHUNKS_DIR_NAME) + .join(format!("{:02x}", key.last().copied().unwrap_or(0))) + .join(hex::encode(key)) + } + fn committed(&self) -> Arc { self.state.current().expect("a current commitment") } async fn round1(&self, nonce: [u8; 32]) -> SubtreeAuditResponse { + self.round1_work(nonce).await.response + } + + async fn round1_work(&self, nonce: [u8; 32]) -> Round1Work { let challenge = SubtreeAuditChallenge { challenge_id: CHALLENGE_ID, nonce, @@ -2846,7 +2889,6 @@ mod pointer_audit_tests { Some(&self.state), ) .await - .response } async fn round2( @@ -3063,6 +3105,73 @@ mod pointer_audit_tests { ); } + /// A committed chunk whose file has rotted fails round 1 exactly as before, + /// but the responder reports it while answering and, once the reply has gone, + /// stops claiming it, so its next commitment leaves the chunk out and + /// replication can repair it rather than the node failing every audit that + /// lands on that key. + #[tokio::test] + async fn round_one_takes_a_rotted_committed_chunk_out_of_service() { + let responder = Responder::new(24, 0).await; + let committed = responder.committed(); + let nonce = [7u8; 32]; + let plan = subtree_plan(committed.tree(), &nonce).expect("plan"); + let victim = plan.leaf_keys.first().copied().expect("a selected chunk"); + let path = responder.chunk_file(&victim); + std::fs::write(&path, b"rotted").expect("corrupt the file"); + + let work = responder.round1_work(nonce).await; + assert_eq!( + work.corrupt_keys, + vec![victim], + "the responder reports exactly the rotted chunk" + ); + match work.response { + SubtreeAuditResponse::Proof { + commitment, proof, .. + } => assert_eq!( + evaluate_subtree_structure( + &commitment, + &proof, + &nonce, + &committed.hash(), + &responder.peer_bytes, + ), + Err(AuditFailureReason::DigestMismatch), + "the auditor's verdict on the rotted chunk is unchanged" + ), + other => panic!("expected a proof, got {other:?}"), + } + + // What the replication engine does with the report once the reply has gone. + for key in &work.corrupt_keys { + responder.storage.recheck_corrupt(key).await; + } + let keys = responder.storage.all_keys().await.expect("keys"); + assert!( + !keys.contains(&victim), + "the rotted chunk leaves the view the next commitment is built from" + ); + assert!(!path.exists(), "the rotted file is removed"); + assert_eq!(keys.len(), 23, "every intact chunk stays in service"); + } + + /// An honest round 1 over intact chunks takes nothing out of service. + #[tokio::test] + async fn round_one_over_intact_chunks_keeps_every_chunk() { + let responder = Responder::new(24, 0).await; + for seed in 0..8u8 { + assert!(responder + .round1_work([seed; 32]) + .await + .corrupt_keys + .is_empty()); + let leaves = responder.proved_leaves([seed; 32]).await; + assert!(!leaves.is_empty()); + } + assert_eq!(responder.storage.all_keys().await.expect("keys").len(), 24); + } + fn replace_records(items: &mut [SubtreeSliceItem], target: &XorName, with: &[Vec]) { for item in items { if let SubtreeSliceItem::PointerRecord { key, records } = item { diff --git a/src/storage/chunk_store.rs b/src/storage/chunk_store.rs index 042646d7..971d0ec1 100644 --- a/src/storage/chunk_store.rs +++ b/src/storage/chunk_store.rs @@ -714,9 +714,9 @@ impl ChunkStore { Ok(Some(content)) => return Ok(Some(content)), Ok(None) => {} // Same rule as `get`: while the legacy environment is there it may have the - // bytes, and this is the read that drives digest audits, possession checks - // and pruning. Answering "no digest" for a chunk the node can still produce - // is a failed audit for nothing. + // bytes, and this is the read that answers digest and subtree audits. + // Answering "no digest" for a chunk the node can still produce is a failed + // audit for nothing. Err(e) => { let Some(legacy) = fallback else { return Err(e); @@ -746,6 +746,33 @@ impl ChunkStore { Ok(raw) } + /// Take a key out of service after its raw bytes were found not to hash to it. + /// + /// [`Self::get_raw`] does not verify, and it is the read that answers audits, so a + /// rotted or torn file is never noticed there: the node goes on committing the key + /// and failing every audit that lands on it, and only a fetch of that exact key + /// would ever repair it. A caller that has hashed raw bytes anyway and found them + /// wrong passes the key here, and the verifying read does the rest exactly as it + /// does for a fetch: it re-reads the file under the write lane, removes it only if + /// it is still wrong, stops the node claiming the key so its next commitment leaves + /// it out and replication brings a good copy back, and re-queues a legacy copy if + /// there is one. + /// + /// Follows the node's `verify_on_read` setting, like every other verifying read. + pub(crate) async fn recheck_corrupt(&self, address: &XorName) { + match self.get(address).await { + Ok(Some(_)) => debug!( + "Chunk {} reads back intact or was served from the legacy environment", + hex::encode(address) + ), + Ok(None) => debug!("Chunk {} is no longer stored", hex::encode(address)), + Err(e) => debug!( + "Chunk {} taken out of service after a failed check: {e}", + hex::encode(address) + ), + } + } + /// Check whether a chunk is stored, in either backing. /// /// An in-memory lookup: no syscall, no I/O, in both phases. @@ -3022,6 +3049,7 @@ where #[allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)] mod tests { use super::*; + use crate::storage::file_store::CHUNKS_DIR_NAME; use crate::storage::migration::{now_unix, rank_closest_first, MIN_RETIRE_DELAY_HOURS}; use tempfile::TempDir; @@ -4714,6 +4742,52 @@ mod tests { ); } + #[tokio::test] + async fn recheck_corrupt_takes_a_rotted_file_out_of_service_and_spares_a_good_one() { + let dir = TempDir::new().expect("temp dir"); + let store = open(&dir).await; + let (good, good_bytes) = addressed("intact"); + let (bad, bad_bytes) = addressed("rotted"); + store.put(&good, &good_bytes).await.expect("put good"); + store.put(&bad, &bad_bytes).await.expect("put bad"); + + let path = dir + .path() + .join(CHUNKS_DIR_NAME) + .join(format!("{:02x}", bad.last().copied().unwrap_or(0))) + .join(hex::encode(bad)); + std::fs::write(&path, b"rotted").expect("corrupt the file"); + // The raw read that answers audits still hands the rotted bytes out. + assert_eq!( + store.get_raw(&bad).await.expect("raw").as_deref(), + Some(b"rotted".as_slice()) + ); + + // A false alarm must never cost a good chunk: the recheck reads it back first. + store.recheck_corrupt(&good).await; + store.recheck_corrupt(&bad).await; + + let keys = store.all_keys().await.expect("keys"); + assert!(keys.contains(&good), "an intact chunk stays in service"); + assert!( + !keys.contains(&bad), + "a rotted chunk must leave the view the next commitment is built from" + ); + assert!(!path.exists(), "the rotted file is removed"); + assert!(!store.exists(&bad).expect("exists")); + assert_eq!( + store.get(&good).await.expect("get").expect("present"), + good_bytes + ); + + // Replication's repair is an ordinary store of the right bytes. + store.put(&bad, &bad_bytes).await.expect("repair"); + assert_eq!( + store.get(&bad).await.expect("get").expect("present"), + bad_bytes + ); + } + #[tokio::test] async fn a_marker_claiming_more_than_the_file_store_holds_restarts_the_copy() { let dir = TempDir::new().expect("temp dir"); From ebfdd049e7b1d5a416405aabd0cce4ed953a526f Mon Sep 17 00:00:00 2001 From: grumbach Date: Fri, 2 Oct 2026 13:57:08 +0900 Subject: [PATCH 2/2] fix(replication): bound the cleanup of chunks audits find rotted Round 1 of the subtree audit reports the committed chunks whose bytes no longer hash to their key, and the responder took each one out of service with a verifying read after sending its reply and releasing its admission permit. That cleanup ran outside every responder limit. The verifying read ends in a quarantine that waits for the chunk's shard write lane on a blocking-pool thread, and while that lane is stalled further proofs keep reading the same rotted file: an auditor can pin a retained commitment that still contains it, and a new peer identity escapes the per-peer cooldown. Each of those proofs parked another blocking thread, with nothing capping how many. The chunk store now keeps the reported keys in a FIFO queue, deduplicated against both the queue and the key being rechecked, and capped at 1024 keys, the most a single round-1 proof can cover. Reporting never waits. A key that does not fit is dropped, and the next audit that reads it reports it again, since it stays on disk and committed until it is removed. The replication engine runs one worker that rechecks the queue a key at a time, so however many audits report rot and however long a lane stalls, the rechecks hold at most one blocking thread between them. The worker stops at the shutdown token between rechecks, and is tracked with the engine's detached storage work, so shutdown waits for a recheck in progress rather than aborting it. Keys still queued at shutdown stay on disk and committed, so after a restart the next audit that reads them reports them again. Holding the round-1 permit through the cleanup would also have bounded it. But round 1 never takes a write lane, so two proofs over rotted chunks in a stalled shard would then have stopped the node answering every subtree audit until the lane freed. The round-1 task keeps its permit until it finishes, as it did before the cleanup was added, and only queues the keys. A regression test holds a shard's real write lane while 32 reports arrive and requires at most one blocking task in flight; with each report running its own recheck, as before, it sees 33. Another checks that everything the protocol's largest subtree can report fits an empty queue. A new end-to-end test rots a node's chunk files, audits it over the live wire, and requires the audit to fail and at least one rotted chunk it read to leave service. Three assertions in this branch's tests gain failure messages for the assert_is_empty lint that Rust 1.99 added and main now passes. --- src/replication/audit.rs | 6 +- src/replication/mod.rs | 39 ++- src/replication/storage_commitment_audit.rs | 53 +-- src/storage/chunk_store.rs | 345 +++++++++++++++++++- src/storage/file_store.rs | 14 + tests/e2e/subtree_audit_testnet.rs | 74 +++++ 6 files changed, 489 insertions(+), 42 deletions(-) diff --git a/src/replication/audit.rs b/src/replication/audit.rs index dd89c593..672d17f3 100644 --- a/src/replication/audit.rs +++ b/src/replication/audit.rs @@ -1452,7 +1452,11 @@ mod tests { &[honest(good, good_bytes), honest(paired, paired_bytes)], ) .await; - assert!(partial.failed_keys.is_empty()); + assert!( + partial.failed_keys.is_empty(), + "an honest peer failed against a rotted reference: {:?}", + partial.failed_keys + ); assert_eq!(partial.verified, 1); assert!(matches!( unfailed_audit_verdict(&peer, 2, partial.verified), diff --git a/src/replication/mod.rs b/src/replication/mod.rs index d32020f8..fa079cdc 100644 --- a/src/replication/mod.rs +++ b/src/replication/mod.rs @@ -2257,6 +2257,8 @@ impl ReplicationEngine { info!("Starting replication engine"); self.start_message_handler(); + // Takes out of service the rotted chunks that round-1 proofs report. + self.start_corrupt_recheck_worker(); self.start_neighbor_sync_loop(); self.start_self_lookup_loop(); // Audit #2 (responsible-chunk): periodic tick auditing peers for the @@ -3297,6 +3299,22 @@ impl ReplicationEngine { self.task_handles.push(handle); } + /// Spawn the one task that takes chunks audits found rotted out of service. + /// + /// One, so that the rechecks hold at most one blocking thread between them however + /// many audits report rot; see [`ChunkStore::run_corrupt_rechecks`]. Tracked with + /// the detached storage work rather than the engine's loops, because shutdown aborts + /// a loop that outlasts its drain timeout and a recheck must not be dropped part-way. + /// The worker itself stops at the shutdown token, between rechecks. + fn start_corrupt_recheck_worker(&self) { + let storage = Arc::clone(&self.storage); + let shutdown = self.shutdown.clone(); + self.detached_task_tracker.spawn(async move { + storage.run_corrupt_rechecks(&shutdown).await; + debug!("Corrupt-chunk recheck worker shut down"); + }); + } + /// Periodic responsible-chunk audit loop (audit #2): every /// [`ReplicationConfig::random_audit_tick_interval`] (~10-20 min), audit one /// eligible close peer for the chunks it *should* be storing (by @@ -4308,10 +4326,8 @@ impl fmt::Display for ResponderAdmissionFailure { /// RAII admission for one audit-responder task: holds the GLOBAL permit and, /// on drop, decrements the PER-PEER in-flight count. Moving this into the -/// spawned task ties both bounds to the task's lifetime — no manual decrement -/// to forget on an early return or panic. A task may drop it once its reply has -/// gone, before purely local follow-up work (round 1 does, before taking -/// rotted chunks out of service), so that work never holds back another audit. +/// spawned task ties both bounds to the task's exact lifetime — no manual +/// decrement to forget on an early return or panic. struct ResponderGuard { _permit: tokio::sync::OwnedSemaphorePermit, _peer_slot: PeerResponderSlot, @@ -5202,8 +5218,7 @@ async fn handle_replication_message( let responder_metrics = Arc::clone(&ctx.audit_responder_metrics); let subtree_round1 = ctx.subtree_round1.clone(); ctx.detached_task_tracker.spawn(async move { - // `guard` is the global permit + per-peer slot, held until the - // reply has gone. + let _guard = guard; // global permit + per-peer slot, held until done let worker_started = Instant::now(); let processing_started = Instant::now(); let storage_commitment_audit::Round1Work { @@ -5220,6 +5235,11 @@ async fn handle_replication_message( ) .await; let processing = processing_started.elapsed(); + // Committed chunks this proof found rotted go to the chunk store's + // single recheck worker. Queueing never waits, so no number of + // audits over a rotted chunk can pile up cleanup work outside this + // task's admission. + storage.report_corrupt(&corrupt_keys); // Charge the work actually done, on EVERY outcome. // // This used to charge only the `Proof` arm, reasoning that the @@ -5276,13 +5296,6 @@ async fn handle_replication_message( processing, response_send, ); - drop(guard); - // A committed chunk this proof found rotted is taken out of service - // only now, with the reply gone and the permit released, so the - // cleanup never delays the auditor or holds back another audit. - for key in &corrupt_keys { - storage.recheck_corrupt(key).await; - } }); Ok(()) } diff --git a/src/replication/storage_commitment_audit.rs b/src/replication/storage_commitment_audit.rs index 0ab00966..495ef519 100644 --- a/src/replication/storage_commitment_audit.rs +++ b/src/replication/storage_commitment_audit.rs @@ -1364,12 +1364,14 @@ pub struct Round1Work { /// their key. /// /// The proof still carries what was read, so the auditor's verdict on them - /// is unchanged. The caller hands each one to the chunk store's corruption - /// recheck once the reply has gone, so the node stops committing a chunk it - /// cannot prove and replication can repair it, rather than failing every + /// is unchanged. The caller reports them to the chunk store + /// (`ChunkStore::report_corrupt`), whose recheck worker takes the ones its + /// bounded queue admits out of service, so the node stops committing a chunk + /// it cannot prove and replication can repair it, rather than failing every /// audit that lands on that key until a fetch of the same key happens to - /// notice. Done after the reply, not here, so the cleanup never delays the - /// answer the auditor is waiting for. + /// notice. A key the queue cannot take stays committed, so the next audit + /// that reads it reports it again. Not done here, so the cleanup never delays + /// the answer the auditor is waiting for. pub corrupt_keys: Vec, } @@ -1642,7 +1644,8 @@ async fn subtree_challenge_response( if leaf.bytes_hash != *key { warn!( "Subtree audit: committed key {} does not hash to its address; \ - it will be taken out of service so replication can repair it", + reporting it for a recheck so it can be taken out of service \ + and repaired", hex::encode(key) ); corrupt_keys.push(*key); @@ -2783,6 +2786,7 @@ mod pointer_audit_tests { use saorsa_pqc::api::sig::ml_dsa_65; use std::path::PathBuf; use tempfile::TempDir; + use tokio_util::sync::CancellationToken; const CHALLENGE_ID: u64 = 7; @@ -3106,10 +3110,10 @@ mod pointer_audit_tests { } /// A committed chunk whose file has rotted fails round 1 exactly as before, - /// but the responder reports it while answering and, once the reply has gone, - /// stops claiming it, so its next commitment leaves the chunk out and - /// replication can repair it rather than the node failing every audit that - /// lands on that key. + /// but the responder reports it, and the chunk store's recheck worker then + /// stops the node claiming it, so its next commitment leaves the chunk out + /// and replication can repair it rather than the node failing every audit + /// that lands on that key. #[tokio::test] async fn round_one_takes_a_rotted_committed_chunk_out_of_service() { let responder = Responder::new(24, 0).await; @@ -3143,10 +3147,19 @@ mod pointer_audit_tests { other => panic!("expected a proof, got {other:?}"), } - // What the replication engine does with the report once the reply has gone. - for key in &work.corrupt_keys { - responder.storage.recheck_corrupt(key).await; - } + // What the replication engine does with the report: queue it for the chunk + // store's recheck worker. The worker finishes the recheck it is on before it + // stops, so once it returns the chunk has been dealt with. + responder.storage.report_corrupt(&work.corrupt_keys); + let stop = CancellationToken::new(); + let removed = async { + let deadline = Instant::now() + Duration::from_secs(30); + while path.exists() && Instant::now() < deadline { + tokio::time::sleep(Duration::from_millis(5)).await; + } + stop.cancel(); + }; + tokio::join!(responder.storage.run_corrupt_rechecks(&stop), removed); let keys = responder.storage.all_keys().await.expect("keys"); assert!( !keys.contains(&victim), @@ -3161,13 +3174,13 @@ mod pointer_audit_tests { async fn round_one_over_intact_chunks_keeps_every_chunk() { let responder = Responder::new(24, 0).await; for seed in 0..8u8 { - assert!(responder - .round1_work([seed; 32]) - .await - .corrupt_keys - .is_empty()); + let corrupt = responder.round1_work([seed; 32]).await.corrupt_keys; + assert!( + corrupt.is_empty(), + "an intact chunk was reported rotted: {corrupt:?}" + ); let leaves = responder.proved_leaves([seed; 32]).await; - assert!(!leaves.is_empty()); + assert!(!leaves.is_empty(), "expected round 1 to prove some leaves"); } assert_eq!(responder.storage.all_keys().await.expect("keys").len(), 24); } diff --git a/src/storage/chunk_store.rs b/src/storage/chunk_store.rs index 971d0ec1..f4d5afd6 100644 --- a/src/storage/chunk_store.rs +++ b/src/storage/chunk_store.rs @@ -20,11 +20,12 @@ use crate::storage::migration::{ CopyReport, MigrationConfig, MigrationPhase, MigrationState, REQUIRED_REBUILDS_BEFORE_RETIRE, }; use crate::storage::StorageStats; -use std::collections::{BTreeMap, BTreeSet}; +use std::collections::{BTreeMap, BTreeSet, HashSet, VecDeque}; use std::io::Write; use std::path::{Path, PathBuf}; use std::sync::Arc; use std::time::Duration; +use tokio::sync::Notify; use tokio_util::sync::CancellationToken; /// Directory name of the legacy LMDB environment, under the node root. @@ -90,6 +91,17 @@ const VERIFY_LOG_EVERY: u64 = 2000; /// collapse into one. const KEY_LOCK_LANES: usize = 256; +/// Most chunks that may be waiting for a corruption recheck at once, counting the one +/// being rechecked. +/// +/// Enough for everything a single round-1 proof can report: a proof covers at most +/// `max_subtree_leaves(MAX_COMMITMENT_KEY_COUNT)` leaves, the square root of the largest +/// legal commitment rounded up to a power of two. That function is not `const`, so a test +/// checks this stays at least that large. A key reported while this is full is dropped +/// rather than waited for: it is still on disk and still committed, so the next audit +/// that reads it reports it again. +const MAX_CORRUPT_REPORTS: usize = 1024; + /// Configuration for [`ChunkStore`]. #[derive(Debug, Clone)] pub struct ChunkStoreConfig { @@ -219,6 +231,40 @@ impl Legacy { } } +/// Chunks found not to hash to their address, waiting for +/// [`ChunkStore::run_corrupt_rechecks`]. +#[derive(Default)] +struct CorruptReports { + /// Waiting to be rechecked, oldest first. + waiting: VecDeque, + /// Everything in `waiting` plus the key being rechecked now, so a chunk reported again + /// before its recheck has finished is not queued twice. + queued: HashSet, +} + +impl CorruptReports { + /// Queue each key not already queued or being rechecked, while there is room. + /// + /// Returns how many were queued, and how many were dropped for want of room. + fn add(&mut self, keys: &[XorName]) -> (usize, usize) { + let mut added = 0; + let mut dropped = 0; + for key in keys { + if self.queued.contains(key) { + continue; + } + if self.queued.len() >= MAX_CORRUPT_REPORTS { + dropped += 1; + continue; + } + self.queued.insert(*key); + self.waiting.push_back(*key); + added += 1; + } + (added, dropped) + } +} + /// Content-addressed chunk storage. pub struct ChunkStore { /// The file store. Always present, always the write target. @@ -253,6 +299,10 @@ pub struct ChunkStore { /// then the copier's write lands and resurrects it. One critical section per key, /// held across put, delete and copy, is what closes that. key_locks: Vec>, + /// Chunks [`Self::report_corrupt`] has queued for the recheck worker. + corrupt: parking_lot::Mutex, + /// Wakes the recheck worker when a report queues something. + corrupt_reported: Notify, } impl ChunkStore { @@ -371,6 +421,8 @@ impl ChunkStore { key_locks: std::iter::repeat_with(|| tokio::sync::Mutex::new(())) .take(KEY_LOCK_LANES) .collect(), + corrupt: parking_lot::Mutex::new(CorruptReports::default()), + corrupt_reported: Notify::new(), }; let (file_keys, legacy_keys) = store.split_counts(); @@ -746,20 +798,78 @@ impl ChunkStore { Ok(raw) } - /// Take a key out of service after its raw bytes were found not to hash to it. + /// Queue keys whose raw bytes were found not to hash to them, for + /// [`Self::run_corrupt_rechecks`] to take out of service. /// /// [`Self::get_raw`] does not verify, and it is the read that answers audits, so a /// rotted or torn file is never noticed there: the node goes on committing the key /// and failing every audit that lands on it, and only a fetch of that exact key /// would ever repair it. A caller that has hashed raw bytes anyway and found them - /// wrong passes the key here, and the verifying read does the rest exactly as it - /// does for a fetch: it re-reads the file under the write lane, removes it only if - /// it is still wrong, stops the node claiming the key so its next commitment leaves - /// it out and replication brings a good copy back, and re-queues a legacy copy if - /// there is one. + /// wrong reports the key here. + /// + /// Never waits. Rechecking here instead would let every audit that reads a rotted + /// chunk start a recheck of its own, outside every responder limit, and while a + /// quarantine waits on its shard's write lane each of those would hold a blocking + /// thread. A key already queued or being rechecked is not queued again, and one that + /// does not fit under [`MAX_CORRUPT_REPORTS`] is dropped. + pub(crate) fn report_corrupt(&self, keys: &[XorName]) { + let (added, dropped) = self.corrupt.lock().add(keys); + if added > 0 { + self.corrupt_reported.notify_one(); + } + if dropped > 0 { + warn!( + "{dropped} chunk(s) found rotted were not queued for a recheck, because \ + {MAX_CORRUPT_REPORTS} already are; the next audit that reads them will \ + report them again" + ); + } + } + + /// Take the keys [`Self::report_corrupt`] queued out of service, one at a time and + /// oldest first, until `stop` is cancelled. + /// + /// The node runs exactly one of these, on the replication engine, so however many + /// audits report rot, and however long a quarantine waits for its shard's write lane, + /// the rechecks hold at most one blocking thread between them. `stop` is raced only + /// against waiting for work, never against a recheck, which may be reading the legacy + /// environment and must not be dropped part-way. Keys still queued when it stops are + /// left there: each is still on disk and still committed, so after a restart the next + /// audit that reads it reports it again. + pub(crate) async fn run_corrupt_rechecks(&self, stop: &CancellationToken) { + // A worker stopped part-way through a recheck leaves that key marked as queued + // with nothing left to finish it; clear it so the key can be reported again. + { + let mut reports = self.corrupt.lock(); + reports.queued = reports.waiting.iter().copied().collect(); + } + loop { + if stop.is_cancelled() { + return; + } + let next = self.corrupt.lock().waiting.pop_front(); + let Some(key) = next else { + tokio::select! { + biased; + () = stop.cancelled() => return, + () = self.corrupt_reported.notified() => {} + } + continue; + }; + self.recheck_corrupt(&key).await; + self.corrupt.lock().queued.remove(&key); + } + } + + /// Take a key out of service if its bytes really do not hash to it. + /// + /// The verifying read does this exactly as it does for a fetch: it re-reads the file + /// under the write lane, removes it only if it is still wrong, stops the node claiming + /// the key so its next commitment leaves it out and replication brings a good copy + /// back, and re-queues a legacy copy if there is one. /// /// Follows the node's `verify_on_read` setting, like every other verifying read. - pub(crate) async fn recheck_corrupt(&self, address: &XorName) { + async fn recheck_corrupt(&self, address: &XorName) { match self.get(address).await { Ok(Some(_)) => debug!( "Chunk {} reads back intact or was served from the legacy environment", @@ -3049,8 +3159,13 @@ where #[allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)] mod tests { use super::*; + use crate::replication::commitment::MAX_COMMITMENT_KEY_COUNT; + use crate::replication::subtree::max_subtree_leaves; use crate::storage::file_store::CHUNKS_DIR_NAME; use crate::storage::migration::{now_unix, rank_closest_first, MIN_RETIRE_DELAY_HOURS}; + use std::sync::mpsc::{self, Sender}; + use std::thread::JoinHandle; + use std::time::Instant; use tempfile::TempDir; /// Everything currently legacy-only, as the set a test has "approved" for shedding. @@ -4788,6 +4903,220 @@ mod tests { ); } + /// Poll `done` until it holds, failing with `what` if it has not within 30 seconds. + async fn wait_until(what: &str, mut done: impl FnMut() -> bool) { + let deadline = Instant::now() + Duration::from_secs(30); + while !done() { + assert!(Instant::now() < deadline, "{what}"); + tokio::time::sleep(Duration::from_millis(5)).await; + } + } + + /// Hold the write lane `address` is quarantined under until the returned sender is used + /// or dropped, stalling that shard the way a hung disk would. Held from a thread of its + /// own, so nothing holds a blocking guard across an await. + fn stall_write_lane(store: &ChunkStore, address: &XorName) -> (Sender<()>, JoinHandle<()>) { + let (lanes, lane) = store.files.test_write_lane(address); + let (held_tx, held_rx) = mpsc::channel::<()>(); + let (release_tx, release_rx) = mpsc::channel::<()>(); + let holder = std::thread::spawn(move || { + let _stalled = lanes.get(lane).map(parking_lot::Mutex::lock); + held_tx.send(()).ok(); + release_rx.recv().ok(); + }); + held_rx.recv().expect("the lane is held"); + (release_tx, holder) + } + + /// A quarantine stalled on its shard's write lane holds one blocking thread, however + /// many audits report rot meanwhile, and what they reported is worked through once the + /// lane frees. + /// + /// This is the hazard as it arises: a retained commitment keeps the rotted key in + /// every subtree an auditor selects over it, the raw read that answers audits keeps + /// reading the file, and identities are cheap, so reports keep coming while the + /// quarantine waits. A recheck run by each report would park a blocking thread per + /// report, outside every responder limit. + #[tokio::test] + async fn a_stalled_quarantine_holds_one_blocking_thread_however_many_reports_arrive() { + let dir = TempDir::new().expect("temp dir"); + let store = Arc::new(open(&dir).await); + let chunk_path = |key: &XorName| { + dir.path() + .join(CHUNKS_DIR_NAME) + .join(format!("{:02x}", key.last().copied().unwrap_or(0))) + .join(hex::encode(key)) + }; + let (bad, bad_bytes) = addressed("stalled"); + let (other, other_bytes) = addressed("queued-behind-the-stall"); + let (later, later_bytes) = addressed("rotted-after-the-stall"); + for (key, bytes) in [ + (&bad, &bad_bytes), + (&other, &other_bytes), + (&later, &later_bytes), + ] { + store.put(key, bytes).await.expect("put"); + } + for key in [&bad, &other] { + std::fs::write(chunk_path(key), b"rotted").expect("rot the file"); + } + + let (release, holder) = stall_write_lane(&store, &bad); + + let stop = CancellationToken::new(); + let worker = { + let store = Arc::clone(&store); + let stop = stop.clone(); + tokio::spawn(async move { store.run_corrupt_rechecks(&stop).await }) + }; + + // The first report: the worker reads the file, proves it wrong, stops claiming it + // and parks in the quarantine behind the held lane. Waited for rather than slept + // at, so the reports below really do arrive while the quarantine is stalled. + store.report_corrupt(&[bad]); + wait_until("the recheck never reached the stalled quarantine", || { + !store.exists(&bad).expect("exists") && store.files.tasks_in_flight() > 0 + }) + .await; + + // Audits keep reporting it, and another rotted chunk, while the quarantine waits. + let reporters: Vec<_> = (0..32) + .map(|_| { + let store = Arc::clone(&store); + tokio::spawn(async move { store.report_corrupt(&[bad, other]) }) + }) + .collect(); + let mut most_in_flight = 0; + let watch_until = Instant::now() + Duration::from_millis(500); + while Instant::now() < watch_until { + most_in_flight = most_in_flight.max(store.files.tasks_in_flight()); + tokio::time::sleep(Duration::from_millis(5)).await; + } + assert!( + most_in_flight <= 1, + "{most_in_flight} rechecks held blocking threads behind one stalled write lane; \ + reports must queue for the single worker, not each start one" + ); + let waiting = reporters.iter().filter(|r| !r.is_finished()).count(); + assert_eq!( + waiting, 0, + "{waiting} reports waited on the stalled recheck" + ); + for reporter in reporters { + reporter.await.expect("report"); + } + { + let reports = store.corrupt.lock(); + assert_eq!( + reports.waiting, + VecDeque::from([other]), + "the key being rechecked must not be queued again, and the other waits its turn" + ); + assert_eq!(reports.queued.len(), 2, "one being rechecked, one waiting"); + } + + // The lane frees: both rotted chunks leave service in turn. + release.send(()).ok(); + holder.join().expect("lane holder"); + wait_until( + "the queue was not worked through after the lane freed", + || { + !chunk_path(&bad).exists() + && !chunk_path(&other).exists() + && store.corrupt.lock().queued.is_empty() + }, + ) + .await; + let keys = store.all_keys().await.expect("keys"); + assert!( + !keys.contains(&bad) && !keys.contains(&other), + "both rotted chunks must leave the view the next commitment is built from" + ); + assert!(keys.contains(&later), "an intact chunk stays in service"); + + // The worker is still there for the next report. + std::fs::write(chunk_path(&later), b"rotted").expect("rot the file"); + store.report_corrupt(&[later]); + wait_until("a report after the stall was never rechecked", || { + !chunk_path(&later).exists() + }) + .await; + + stop.cancel(); + tokio::time::timeout(Duration::from_secs(10), worker) + .await + .expect("an idle worker stops when told to") + .expect("worker"); + } + + /// Every key one maximal round-1 proof can report fits an empty queue, so none of a + /// subtree that rotted whole is dropped for want of room. + #[tokio::test] + async fn the_corrupt_queue_holds_everything_one_proof_can_report() { + let most = usize::try_from(max_subtree_leaves(MAX_COMMITMENT_KEY_COUNT)) + .expect("a leaf count fits in usize"); + let dir = TempDir::new().expect("temp dir"); + let store = open(&dir).await; + let keys: Vec = (0..most) + .map(|i| addressed(&format!("leaf-{i}")).0) + .collect(); + store.report_corrupt(&keys); + assert_eq!( + store.corrupt.lock().waiting.len(), + most, + "a maximal proof's report must fit an empty queue" + ); + } + + /// Reports are deduplicated against both the queue and the key being rechecked, and + /// what does not fit is dropped rather than waited for. + #[tokio::test] + async fn corrupt_reports_are_deduplicated_and_capped() { + let dir = TempDir::new().expect("temp dir"); + let store = open(&dir).await; + let keys: Vec = (0..MAX_CORRUPT_REPORTS + 8) + .map(|i| addressed(&format!("reported-{i}")).0) + .collect(); + + store.report_corrupt(&keys); + { + let reports = store.corrupt.lock(); + assert_eq!(reports.waiting.len(), MAX_CORRUPT_REPORTS); + assert_eq!(reports.queued.len(), MAX_CORRUPT_REPORTS); + assert_eq!(reports.waiting.front(), keys.first(), "oldest first"); + } + + // Reported again, nothing is queued twice. + store.report_corrupt(&keys); + assert_eq!(store.corrupt.lock().waiting.len(), MAX_CORRUPT_REPORTS); + + // Nor is a key the worker has taken and not yet finished, and it still counts + // toward the cap. Taken here exactly as the worker takes one. + let taken = store + .corrupt + .lock() + .waiting + .pop_front() + .expect("a waiting key"); + store.report_corrupt(&[taken]); + { + let reports = store.corrupt.lock(); + assert!( + !reports.waiting.contains(&taken), + "a key being rechecked was queued again" + ); + assert_eq!(reports.queued.len(), MAX_CORRUPT_REPORTS); + } + + // A worker stopped part-way through that recheck leaves the key marked. The next + // one to start clears the mark, so the key can be reported again. + let stopped = CancellationToken::new(); + stopped.cancel(); + store.run_corrupt_rechecks(&stopped).await; + store.report_corrupt(&[taken]); + assert_eq!(store.corrupt.lock().waiting.back(), Some(&taken)); + } + #[tokio::test] async fn a_marker_claiming_more_than_the_file_store_holds_restarts_the_copy() { let dir = TempDir::new().expect("temp dir"); diff --git a/src/storage/file_store.rs b/src/storage/file_store.rs index 3172bdda..59f72d49 100644 --- a/src/storage/file_store.rs +++ b/src/storage/file_store.rs @@ -1572,6 +1572,20 @@ impl FileStore { Arc::clone(&self.test_put_gate) } + /// Test-only handle to the write lanes, and the index of the one `address` is + /// written, repaired and quarantined under. + /// + /// Holding that lane stalls every write, repair and quarantine in the shard the way a + /// hung disk would. [`Self::test_put_gate`] cannot stand in for it: a put parks there + /// before it takes its lane. + #[cfg(test)] + pub(crate) fn test_write_lane( + &self, + address: &XorName, + ) -> (Arc>>, usize) { + (Arc::clone(&self.write_lanes), shard_index(address)) + } + /// Register a write of `address` and hand back the token that clears it. /// /// The token must be moved into the blocking closure that does the work, so the entry diff --git a/tests/e2e/subtree_audit_testnet.rs b/tests/e2e/subtree_audit_testnet.rs index ff8f86f2..42adef84 100644 --- a/tests/e2e/subtree_audit_testnet.rs +++ b/tests/e2e/subtree_audit_testnet.rs @@ -20,6 +20,7 @@ use std::time::{Duration, SystemTime}; use super::TestHarness; use ant_node::replication::audit::AuditTickResult; use ant_node::replication::{FirstAuditStats, MonetizedPinEvent, ReplicationEngine}; +use ant_node::storage::file_store::CHUNKS_DIR_NAME; use serial_test::serial; use tokio::time::sleep; @@ -182,6 +183,79 @@ async fn data_deleting_node_fails_subtree_audit() { harness.teardown().await.expect("teardown"); } +/// SELF-REPAIR: a node whose committed chunk files rotted on disk still fails the +/// audit over them, and then starts taking the rotted chunks its proof read out of +/// service, so its next commitment leaves them out and replication can bring good +/// copies back. +/// +/// Drives the shipped path end to end: the round-1 responder reports the rotted +/// leaves and the engine's recheck worker removes them. Requires at least one to +/// go, and every one gone to be unclaimed. Nothing else removes a rotted file this +/// quickly, so a responder that stopped reporting, or an engine that never started +/// the worker, leaves every file in place and fails this. +#[tokio::test] +#[serial] +async fn a_node_whose_chunks_rotted_takes_them_out_of_service_after_an_audit() { + let harness = TestHarness::setup_small().await.expect("setup"); + harness.warmup_dht().await.expect("warmup"); + + let (a_idx, b_idx) = (5, 6); + let addrs = commit_and_seed(&harness, a_idx, b_idx, 64).await; + + let a = harness.test_node(a_idx).expect("a"); + let a_store = a.ant_protocol.as_ref().expect("a protocol").storage(); + let chunk_path = |addr: &[u8; 32]| { + a.data_dir + .join(CHUNKS_DIR_NAME) + .join(format!("{:02x}", addr.last().copied().unwrap_or(0))) + .join(hex::encode(addr)) + }; + // Every file rots in place, under its own name, so the node still claims all of + // them and its raw reads still find them. + for addr in &addrs { + std::fs::write(chunk_path(addr), b"rotted").expect("rot the file"); + } + + let a_peer = *a.p2p_node.as_ref().expect("a p2p").peer_id(); + let b_engine = harness + .test_node(b_idx) + .expect("b") + .replication_engine + .as_ref() + .expect("b engine"); + let result = b_engine.audit_peer_now(&a_peer).await; + assert!( + matches!(result, AuditTickResult::Failed { .. }), + "a node serving rotted bytes must still FAIL the audit, got {result:?}" + ); + + let start = std::time::Instant::now(); + let removed = loop { + let removed: Vec<_> = addrs + .iter() + .filter(|addr| !chunk_path(addr).exists()) + .collect(); + if !removed.is_empty() || start.elapsed() >= Duration::from_secs(30) { + break removed; + } + sleep(Duration::from_millis(100)).await; + }; + assert!( + !removed.is_empty(), + "no rotted chunk the audit read was taken out of service" + ); + let claimed = a_store.all_keys().await.expect("a keys"); + for addr in &removed { + assert!( + !claimed.contains(*addr), + "{} lost its file but is still claimed", + hex::encode(addr) + ); + } + + harness.teardown().await.expect("teardown"); +} + /// SELF FILTER: a monetized-pin event targeting the node's own peer ID (the /// payment verifier emits one for every payment it verifies, because the /// node's own quote is in the payment's quote list) must be dropped by the