Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
211 changes: 184 additions & 27 deletions src/replication/audit.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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<AuditKeyFailure>,
/// 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.
Expand All @@ -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,
}
}

// ---------------------------------------------------------------------------
Expand Down Expand Up @@ -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.
Expand Down
24 changes: 24 additions & 0 deletions src/replication/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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,
Expand All @@ -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
Expand Down
7 changes: 5 additions & 2 deletions src/replication/possession.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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");
Expand Down
43 changes: 42 additions & 1 deletion src/replication/pruning.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<ChunkStore>) -> Option<Vec<u8>> {
match storage.get_raw(key).await {
match storage.get(key).await {
Ok(Some(bytes)) => Some(bytes),
Ok(None) => {
debug!(
Expand Down Expand Up @@ -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];
Expand All @@ -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);
Expand Down
Loading
Loading