From 2b5096aa2e5647b1d0c46b7511b5e63c62727573 Mon Sep 17 00:00:00 2001 From: moudyellaz Date: Tue, 4 Aug 2026 03:06:46 +0200 Subject: [PATCH] fix(indexer): validate peer blocks before caching them --- lez/indexer/core/src/cross_zone_verifier.rs | 244 ++++++++++++++++++-- 1 file changed, 231 insertions(+), 13 deletions(-) diff --git a/lez/indexer/core/src/cross_zone_verifier.rs b/lez/indexer/core/src/cross_zone_verifier.rs index da568ab8..3ada2a73 100644 --- a/lez/indexer/core/src/cross_zone_verifier.rs +++ b/lez/indexer/core/src/cross_zone_verifier.rs @@ -12,7 +12,7 @@ use cross_zone_inbox_core::{ }; use futures::{Stream, StreamExt as _}; use lee::{GENESIS_BLOCK_ID, PublicKey}; -use log::{debug, error, info}; +use log::{debug, error, info, warn}; use logos_blockchain_core::mantle::ops::channel::ChannelId; use logos_blockchain_zone_sdk::{ CommonHttpClient, Slot, ZoneMessage, adapter::NodeHttpClient, indexer::ZoneIndexer, @@ -73,11 +73,11 @@ struct PeerChain { /// /// This, not `max(blocks.keys())`, is what the forgery test gates on: a peer /// picks its own `block_id`s, and an id that does not continue the run - /// cannot advance the run. It bounds how far this reader has read, not that - /// the chain is authentic: the link is self-asserted, since `header.hash` is - /// not recomputed on decode and the reader does not apply the pinned-key - /// check that [`crate::cross_zone_verifier::CrossZoneVerifier::rederive`] - /// applies to the block a dispatch actually names. + /// cannot advance the run. + /// + /// The link it walks means something only because [`accept_peer_block`] + /// recomputes `header.hash` and checks the pinned key before anything is + /// cached. Without that it compared two fields the peer wrote. verified_prefix: Option, } @@ -126,11 +126,38 @@ struct PeerBlocks { } impl PeerBlocks { - async fn insert(&self, zone: ZoneId, block: Block) { + /// Caches `block` unless a different block is already held at its id, and + /// says whether it was newly cached. + /// + /// First write wins. An identical re-read is a no-op, which the reader does + /// on every slot retry. A differing block at a held id is equivocation, and + /// replacing the entry is what makes it a remote halt: the prefix already + /// certified the old value, so the next dispatch naming that id re-derives + /// against the new one and reads as forged. + /// + /// A peer resetting its chain is indistinguishable from this and is refused + /// the same way. The cache is process-local, so a restart adopts it. + async fn insert(&self, zone: ZoneId, block: Block) -> bool { let mut chains = self.chains.write().await; let chain = chains.entry(zone).or_default(); + + if let Some(held) = chain.blocks.get(&block.header.block_id) { + if held.header.hash == block.header.hash { + return false; + } + error!( + "Peer zone {} equivocated at block {}: holding {}, refusing {}. Restart the indexer if this peer legitimately reset its chain.", + hex::encode(zone), + block.header.block_id, + held.header.hash, + block.header.hash + ); + return false; + } + chain.blocks.insert(block.header.block_id, block); chain.extend_prefix(); + true } /// Resolves `block_id` under a single read lock. @@ -220,6 +247,7 @@ impl CrossZoneVerifier { tokio::spawn(read_peer( ZoneIndexer::new(ChannelId::from(peer.channel_id), node), peer.channel_id, + peer_pubkeys.get(&peer.channel_id).cloned(), peers.clone(), config.consensus_info_polling_interval, )); @@ -444,6 +472,43 @@ struct PeerPass { stalled_at: Option, } +/// Whether a block read off a peer's channel may enter the cache. The channel +/// authorizes who may write, not what they may claim. +/// +/// The hash check is unconditional: `header.hash` is a field the peer wrote and +/// the prefix walk compares it against the next block's `prev_block_hash`, so +/// without recomputing it a peer can assert links it never built. The key check +/// applies only when one is pinned, mirroring the watcher; it subsumes the hash +/// check, but a peer with no pinned key still gets that one. +fn accept_peer_block( + block: &Block, + peer_zone: ZoneId, + expected_pubkey: Option<&PublicKey>, +) -> bool { + if block.recompute_hash() != block.header.hash { + warn!( + "Peer reader dropping block {} from {}: header hash {} does not match its contents", + block.header.block_id, + hex::encode(peer_zone), + block.header.hash + ); + return false; + } + + if let Some(expected) = expected_pubkey + && !block.is_signed_by(expected) + { + warn!( + "Peer reader dropping block {} from {}: not signed by the pinned block-signing key", + block.header.block_id, + hex::encode(peer_zone) + ); + return false; + } + + true +} + /// Reads a peer zone's finalized blocks from Bedrock into the shared cache. #[expect( clippy::infinite_loop, @@ -452,6 +517,7 @@ struct PeerPass { async fn read_peer( zone_indexer: ZoneIndexer, peer_zone: ZoneId, + expected_pubkey: Option, peers: PeerBlocks, poll_interval: Duration, ) { @@ -469,7 +535,15 @@ async fn read_peer( loop { match zone_indexer.next_messages(cursor).await { Ok(stream) => { - let pass = consume_peer_stream(stream, peer_zone, &peers, cursor, skip_slot).await; + let pass = consume_peer_stream( + stream, + peer_zone, + expected_pubkey.as_ref(), + &peers, + cursor, + skip_slot, + ) + .await; cursor = pass.cursor; if let Some(slot) = pass.stalled_at { let attempts = match stalled { @@ -521,6 +595,7 @@ async fn read_peer( async fn consume_peer_stream( stream: S, peer_zone: ZoneId, + expected_pubkey: Option<&PublicKey>, peers: &PeerBlocks, resume_from: Option, skip_slot: Option, @@ -543,7 +618,13 @@ where continue; }; match borsh::from_slice::(&zone_block.data) { - Ok(block) => peers.insert(peer_zone, block).await, + Ok(block) => { + // Before caching, not when a dispatch names it: an unchecked + // block steers the prefix, and by then the damage is a halt. + if accept_peer_block(&block, peer_zone, expected_pubkey) { + peers.insert(peer_zone, block).await; + } + } Err(err) if skip_slot == Some(slot) => { debug!( "Peer reader skipping undecodable block from {} at slot {slot:?}: {err}", @@ -835,7 +916,7 @@ mod tests { peer_block_msg(&chain[1], 1), ]); - let pass = consume_peer_stream(stream, PEER_ZONE, &peers, None, None).await; + let pass = consume_peer_stream(stream, PEER_ZONE, None, &peers, None, None).await; assert_eq!(pass.cursor, Some(Slot::from(1))); assert_eq!(pass.stalled_at, None); @@ -891,7 +972,7 @@ mod tests { peer_block_msg(&chain[2], 2), ]); - let pass = consume_peer_stream(stream, PEER_ZONE, &peers, None, None).await; + let pass = consume_peer_stream(stream, PEER_ZONE, None, &peers, None, None).await; assert_eq!(pass.cursor, Some(Slot::from(0))); assert_eq!(pass.stalled_at, Some(Slot::from(1))); @@ -906,7 +987,8 @@ mod tests { // One slot can carry several messages; the second one fails. let stream = stream::iter(vec![peer_block_msg(&chain[0], 7), undecodable_msg(7)]); - let pass = consume_peer_stream(stream, PEER_ZONE, &peers, Some(Slot::from(6)), None).await; + let pass = + consume_peer_stream(stream, PEER_ZONE, None, &peers, Some(Slot::from(6)), None).await; // Slot 7 is re-read whole next pass, not resumed past the failure. assert_eq!(pass.cursor, Some(Slot::from(6))); @@ -924,7 +1006,8 @@ mod tests { peer_block_msg(&chain[2], 2), ]); - let pass = consume_peer_stream(stream, PEER_ZONE, &peers, None, Some(Slot::from(1))).await; + let pass = + consume_peer_stream(stream, PEER_ZONE, None, &peers, None, Some(Slot::from(1))).await; assert_eq!(pass.cursor, Some(Slot::from(2)), "the pass drains"); assert_eq!(pass.stalled_at, None); @@ -946,6 +1029,7 @@ mod tests { consume_peer_stream( stream, PEER_ZONE, + None, &verifier.peers, None, Some(Slot::from(1)), @@ -1002,6 +1086,7 @@ mod tests { let pass = consume_peer_stream( stream::iter(vec![undecodable_msg(0)]), PEER_ZONE, + None, &verifier.peers, None, None, @@ -1016,6 +1101,7 @@ mod tests { peer_block_msg(block, u64::try_from(index).expect("test index fits in u64")) })), PEER_ZONE, + None, &verifier.peers, pass.cursor, None, @@ -1029,4 +1115,136 @@ mod tests { .await .expect("the dispatch verifies once the peer block has been read"); } + + /// The live brick from #648: a second, differing block at a used id replaced + /// the entry, so the next dispatch naming that id re-derived against the + /// wrong block and halted ingestion. One inscription, remote halt. + #[tokio::test] + async fn an_equivocating_peer_cannot_replace_a_cached_block() { + let verifier = verifier(); + let real = produce_dummy_block(PEER_BLOCK_ID, None, vec![emission(b"hi")]); + let impostor = produce_dummy_block(PEER_BLOCK_ID, None, vec![emission(b"forged")]); + assert_ne!( + real.header.hash, impostor.header.hash, + "the two blocks must differ, or this proves nothing" + ); + + assert!(verifier.peers.insert(PEER_ZONE, real.clone()).await); + assert!( + !verifier.peers.insert(PEER_ZONE, impostor).await, + "a differing block at a held id must be refused" + ); + assert_eq!( + verifier + .peers + .get(PEER_ZONE, PEER_BLOCK_ID) + .await + .unwrap() + .header + .hash, + real.header.hash, + "the block held first stays" + ); + + // And the dispatch that would have been rejected as forged still verifies. + let block = produce_dummy_block(9, None, vec![dispatch(b"hi")]); + verifier + .verify_block(&block) + .await + .expect("equivocation must not halt ingestion of an honest dispatch"); + } + + /// The reader re-reads a slot on every retry, so caching the same block + /// twice must be a quiet no-op rather than equivocation. + #[tokio::test] + async fn re_reading_the_same_block_is_not_equivocation() { + let verifier = verifier(); + let block = produce_dummy_block(PEER_BLOCK_ID, None, vec![emission(b"hi")]); + + assert!(verifier.peers.insert(PEER_ZONE, block.clone()).await); + assert!( + !verifier.peers.insert(PEER_ZONE, block.clone()).await, + "an identical re-read is already held, not newly cached" + ); + assert_eq!( + verifier + .peers + .get(PEER_ZONE, PEER_BLOCK_ID) + .await + .unwrap() + .header + .hash, + block.header.hash + ); + } + + /// Recomputing `header.hash` on the way in is what stops a peer asserting + /// links it never built. + #[tokio::test] + async fn a_block_whose_hash_does_not_match_its_contents_is_not_cached() { + let verifier = verifier(); + let mut tampered = produce_dummy_block(PEER_BLOCK_ID, None, vec![emission(b"hi")]); + tampered.header.hash = common::HashType([0xAB; 32]); + + let pass = consume_peer_stream( + stream::iter(vec![peer_block_msg(&tampered, 0)]), + PEER_ZONE, + None, + &verifier.peers, + None, + None, + ) + .await; + + assert_eq!( + pass.stalled_at, None, + "a rejected block is not a decode failure" + ); + assert!( + verifier.peers.get(PEER_ZONE, PEER_BLOCK_ID).await.is_none(), + "a block that does not hash to its own contents must not be cached" + ); + } + + /// The watcher drops unsigned peer blocks before use; the reader applied no + /// check at all, so anything on the channel entered the cache. + #[tokio::test] + async fn a_block_not_signed_by_the_pinned_key_is_not_cached() { + let verifier = verifier(); + let block = produce_dummy_block(PEER_BLOCK_ID, None, vec![emission(b"hi")]); + let wrong_key = PublicKey::try_new([42; 32]).unwrap(); + + let pass = consume_peer_stream( + stream::iter(vec![peer_block_msg(&block, 0)]), + PEER_ZONE, + Some(&wrong_key), + &verifier.peers, + None, + None, + ) + .await; + + assert_eq!(pass.stalled_at, None); + assert!( + verifier.peers.get(PEER_ZONE, PEER_BLOCK_ID).await.is_none(), + "a block not signed by the pinned key must not reach the cache" + ); + + // The same block under its real signer is cached, so the gate is the key + // and not the path. + let signer = PublicKey::new_from_private_key(&PrivateKey::try_new([37; 32]).unwrap()); + consume_peer_stream( + stream::iter(vec![peer_block_msg(&block, 1)]), + PEER_ZONE, + Some(&signer), + &verifier.peers, + None, + None, + ) + .await; + assert!( + verifier.peers.get(PEER_ZONE, PEER_BLOCK_ID).await.is_some(), + "the pinned signer's own block is cached" + ); + } }