diff --git a/src/replication/audit.rs b/src/replication/audit.rs index 9b312bb3..672d17f3 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,105 @@ 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(), + "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), + 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 3701c1cc..3d2af8b6 100644 --- a/src/replication/mod.rs +++ b/src/replication/mod.rs @@ -2300,6 +2300,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 @@ -3484,6 +3486,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 @@ -5418,6 +5436,7 @@ async fn handle_replication_message( let storage_commitment_audit::Round1Work { response, content_bytes, + corrupt_keys, } = storage_commitment_audit::handle_subtree_challenge_measured_with_pointers( &challenge, &storage, @@ -5428,6 +5447,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 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 5df180f5..4d890622 100644 --- a/src/replication/storage_commitment_audit.rs +++ b/src/replication/storage_commitment_audit.rs @@ -1382,6 +1382,19 @@ 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 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. 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, } /// [`handle_subtree_challenge`], additionally reporting the read-and-hash work @@ -1418,6 +1431,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, @@ -1426,18 +1440,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, @@ -1446,6 +1463,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 { @@ -1642,6 +1660,18 @@ 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; \ + reporting it for a recheck so it can be taken out of service \ + and repaired", + hex::encode(key) + ); + corrupt_keys.push(*key); + } leaves.push(leaf); } @@ -2914,10 +2944,13 @@ mod pointer_audit_tests { use crate::replication::commitment::MAX_COMMITMENT_KEY_COUNT; use crate::replication::commitment_state::BuiltCommitment; use crate::replication::subtree::max_subtree_leaves; + 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; + use tokio_util::sync::CancellationToken; const CHALLENGE_ID: u64 = 7; @@ -2941,6 +2974,7 @@ mod pointer_audit_tests { state: Arc, peer: PeerId, peer_bytes: [u8; 32], + chunk_root: PathBuf, _dirs: (TempDir, TempDir), } @@ -2986,15 +3020,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, @@ -3010,7 +3057,6 @@ mod pointer_audit_tests { Some(&self.state), ) .await - .response } /// Round 2 as the engine serves it: with what round 1 bound for each @@ -3581,6 +3627,82 @@ mod pointer_audit_tests { ); } + /// A committed chunk whose file has rotted fails round 1 exactly as before, + /// 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; + 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: 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), + "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 { + 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(), "expected round 1 to prove some leaves"); + } + 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 3fd1f2db..dd6c4a33 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(); @@ -714,9 +766,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 +798,91 @@ impl ChunkStore { Ok(raw) } + /// 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 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. + 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. @@ -3030,7 +3167,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. @@ -4722,6 +4865,266 @@ 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 + ); + } + + /// 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 b4b087d7..110041a8 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