Skip to main content

zakura_network/zakura/
legacy_gossip.rs

1//! Legacy Zebra gossip compatibility over Zakura streams.
2
3use std::{
4    collections::{HashMap, HashSet, VecDeque},
5    fmt,
6    future::Future,
7    io::Cursor,
8    pin::Pin,
9    sync::{
10        atomic::{AtomicU64, Ordering},
11        Arc, Mutex as StdMutex, OnceLock,
12    },
13    task::{Context, Poll},
14    time::Duration,
15};
16
17use thiserror::Error;
18use tokio::{
19    sync::{mpsc, oneshot, Mutex, OwnedSemaphorePermit, Semaphore},
20    time::{sleep, timeout, Instant},
21};
22use tower::{Service, ServiceExt};
23
24use zakura_chain::{
25    block::{self, Block, MAX_BLOCK_LOCATOR_LENGTH},
26    serialization::{
27        CompactSizeMessage, SerializationError, ZcashDeserialize, ZcashSerialize,
28        MAX_HEADERS_PER_MESSAGE, MAX_PROTOCOL_MESSAGE_LEN,
29    },
30    transaction::{Transaction, UnminedTx, UnminedTxId},
31};
32
33use crate::{
34    protocol::{
35        external::InventoryHash,
36        internal::{InventoryResponse, PeerSource, Request, Response},
37    },
38    BoxError, MAX_TX_INV_IN_SENT_MESSAGE,
39};
40
41use super::{
42    spawn_supervised_peer_task, BoxRunFuture, Frame, FramedSend, OrderedSendError,
43    OrderedSessionDemand, OrderedStreamOpening, OrderedStreamPolicy, Peer, RequestResponseService,
44    Service as ZakuraService, ServicePeerDirection, SinkReject, Stream, StreamMode, ZakuraConnId,
45    ZakuraPeerHandle, ZakuraPeerId, ZakuraSupervisorHandle, ZakuraTrace, FRAME_HEADER_BYTES,
46    LOCAL_MAX_CONTROL_FRAME_BYTES, ZAKURA_CAP_LEGACY_GOSSIP,
47};
48
49mod trace;
50
51use trace::{LegacyRequestError, LegacyRequestResponse, LegacyRequestStart};
52
53/// Zakura stream kind reserved for legacy gossip compatibility.
54pub const ZAKURA_STREAM_GOSSIP: u16 = 2;
55/// Zakura stream kind reserved for legacy inventory request/response compatibility.
56pub const ZAKURA_STREAM_LEGACY_REQUESTS: u16 = 3;
57/// Version of the legacy gossip compatibility stream.
58pub const LEGACY_GOSSIP_VERSION: u16 = 1;
59/// A block hash advertisement.
60pub const MSG_ADVERTISE_BLOCK: u16 = 1;
61/// A transaction-id advertisement.
62pub const MSG_ADVERTISE_TX_IDS: u16 = 2;
63/// Request block contents by hash.
64pub const MSG_REQUEST_BLOCKS_BY_HASH: u16 = 3;
65/// Request transaction contents by unmined transaction id.
66pub const MSG_REQUEST_TRANSACTIONS_BY_ID: u16 = 4;
67/// A chunk of one available block response.
68pub const MSG_RESPONSE_BLOCK: u16 = 5;
69/// A chunk of one available transaction response.
70pub const MSG_RESPONSE_TRANSACTION: u16 = 6;
71/// Missing block hashes.
72pub const MSG_RESPONSE_MISSING_BLOCKS: u16 = 7;
73/// Missing transaction ids.
74pub const MSG_RESPONSE_MISSING_TRANSACTIONS: u16 = 8;
75/// Request subsequent block hashes from a block locator.
76pub const MSG_REQUEST_FIND_BLOCKS: u16 = 9;
77/// Request subsequent block headers from a block locator.
78pub const MSG_REQUEST_FIND_HEADERS: u16 = 10;
79/// Request mempool transaction IDs.
80pub const MSG_REQUEST_MEMPOOL_TRANSACTION_IDS: u16 = 11;
81/// Ping a peer over the request/response stream.
82pub const MSG_REQUEST_PING: u16 = 12;
83/// Push an unsolicited transaction to a peer.
84pub const MSG_REQUEST_PUSH_TRANSACTION: u16 = 13;
85/// Subsequent block hashes response.
86pub const MSG_RESPONSE_BLOCK_HASHES: u16 = 14;
87/// Subsequent block headers response.
88pub const MSG_RESPONSE_BLOCK_HEADERS: u16 = 15;
89/// Mempool transaction IDs response.
90pub const MSG_RESPONSE_TRANSACTION_IDS: u16 = 16;
91/// Pong response for a ping request.
92pub const MSG_RESPONSE_PONG: u16 = 17;
93/// Nil response for fire-and-forget legacy requests.
94pub const MSG_RESPONSE_NIL: u16 = 18;
95
96const LEGACY_GOSSIP_INBOUND_QUEUE: usize = 256;
97const LEGACY_REQUEST_IN_FLIGHT_LIMIT: usize = 64;
98const LEGACY_GOSSIP_SERVICE_TIMEOUT: Duration = Duration::from_secs(30);
99const DEFAULT_FIRST_SEEN_TTL: Duration = Duration::from_secs(10 * 60);
100const DEFAULT_FIRST_SEEN_CAPACITY: usize = 50_000;
101/// A failed gossip inbound attempt that consumed at least this long is treated as
102/// "expensive" (a slow/backpressured/timed-out service), and the same inventory is
103/// placed in a short cooldown. The first-seen cache is only updated after a
104/// *successful* call, so without this an authenticated peer could replay the same
105/// valid advertisement while the inbound service is slow/erroring and make the
106/// serial gossip worker re-pay the full 30s readiness/call budget for every
107/// duplicate. Fast failures stay below this threshold and remain immediately
108/// retryable, so a transient blip does not drop the advertisement.
109const LEGACY_GOSSIP_EXPENSIVE_ATTEMPT: Duration = Duration::from_secs(1);
110/// How long an expensive failed attempt suppresses duplicate copies of the same
111/// inventory before a genuine re-advertisement may retry.
112const LEGACY_GOSSIP_DUPLICATE_COOLDOWN: Duration = Duration::from_secs(30);
113const LEGACY_REQUEST_TIMEOUT: Duration = Duration::from_secs(30);
114const SOURCE_INVENTORY_MISSING_RETRIES: usize = 8;
115const SOURCE_INVENTORY_MISSING_RETRY_DELAY: Duration = Duration::from_millis(500);
116const LEGACY_REQUEST_READY_TIMEOUT: Duration = Duration::from_secs(10);
117/// Reserve half of each connection's stream-open budget for native ordered
118/// streams, reconnects, and other request clients.
119const LEGACY_REQUEST_STREAM_RATE_DIVISOR: u32 = 2;
120/// How long the dual-stack tries the (buffered) legacy peer set for an inventory
121/// fetch before falling back to Zakura. Without this bound, a node that upgraded
122/// all its peers to Zakura (and so has no ready legacy peer) would block every
123/// fetch on the legacy peer set forever, starving the Zakura path.
124const DUAL_STACK_LEGACY_INVENTORY_TIMEOUT: Duration = Duration::from_secs(3);
125const LEGACY_RESPONSE_CHUNK_BYTES: usize = 512 * 1024;
126/// Maximum cumulative response payload bytes the inbound responder will buffer
127/// for a single legacy request before aborting.
128///
129/// A peer can name up to `MAX_TX_INV_IN_SENT_MESSAGE` block/transaction hashes
130/// on one request; without an aggregate cap `encode_response` would serialize
131/// and retain the entire multi-frame `Vec<Frame>` (worst case
132/// `MAX_TX_INV_IN_SENT_MESSAGE * MAX_PROTOCOL_MESSAGE_LEN`, tens of GiB) before
133/// the first byte is written. The outbound reader already enforces a symmetric
134/// per-response cap (`LegacyResponseBudget`); this is the responder-side mirror.
135/// Sized well above any single response the local service emits (zakurad caps
136/// `getdata` at ~1 MiB) yet far below memory-risk thresholds.
137const LEGACY_RESPONSE_MAX_AGGREGATE_BYTES: usize = 8 * MAX_PROTOCOL_MESSAGE_LEN;
138const REQUEST_ID_BYTES: usize = 8;
139const RESPONSE_CHUNK_HEADER_BYTES: usize = REQUEST_ID_BYTES + 1;
140const NO_STOP_HASH: block::Hash = block::Hash([0; 32]);
141const LEGACY_GOSSIP_SERVICE_STREAMS: [Stream; 2] = [
142    Stream {
143        kind: ZAKURA_STREAM_GOSSIP,
144        version: LEGACY_GOSSIP_VERSION,
145        frame_cap: LOCAL_MAX_CONTROL_FRAME_BYTES,
146        capability: ZAKURA_CAP_LEGACY_GOSSIP,
147        mode: StreamMode::Ordered,
148    },
149    Stream {
150        kind: ZAKURA_STREAM_LEGACY_REQUESTS,
151        version: LEGACY_GOSSIP_VERSION,
152        frame_cap: LOCAL_MAX_CONTROL_FRAME_BYTES,
153        capability: ZAKURA_CAP_LEGACY_GOSSIP,
154        mode: StreamMode::RequestResponse,
155    },
156];
157
158/// Service-declared streams for legacy gossip compatibility.
159pub(crate) fn legacy_gossip_streams() -> &'static [Stream] {
160    &LEGACY_GOSSIP_SERVICE_STREAMS
161}
162
163static FIRST_SEEN_BY_SUPERVISOR: OnceLock<std::sync::Mutex<HashMap<u64, FirstSeenCache>>> =
164    OnceLock::new();
165static NEXT_LEGACY_REQUEST_ID: AtomicU64 = AtomicU64::new(1);
166
167/// A typed legacy gossip frame carried by Zakura stream kind 2.
168#[derive(Clone, Debug, Eq, PartialEq)]
169pub enum LegacyGossipFrame {
170    /// Advertise one block hash.
171    AdvertiseBlock(block::Hash),
172    /// Advertise one or more unmined transaction IDs.
173    AdvertiseTransactionIds(Vec<UnminedTxId>),
174}
175
176impl LegacyGossipFrame {
177    /// Convert a legacy network request into a gossip frame.
178    pub fn from_request(request: Request) -> Result<Self, LegacyGossipError> {
179        match request {
180            Request::AdvertiseBlock(hash, _) | Request::AdvertiseBlockToAll(hash) => {
181                Ok(Self::AdvertiseBlock(hash))
182            }
183            Request::AdvertiseTransactionIds(ids, _) => Self::advertise_transaction_ids(ids),
184            request => Err(LegacyGossipError::UnsupportedRequest(request.command())),
185        }
186    }
187
188    /// Build a tx-id gossip frame, enforcing the outbound non-empty/cap invariant.
189    pub fn advertise_transaction_ids(
190        ids: impl IntoIterator<Item = UnminedTxId>,
191    ) -> Result<Self, LegacyGossipError> {
192        let ids = collect_outbound_tx_ids(ids)?;
193        Ok(Self::AdvertiseTransactionIds(ids))
194    }
195
196    /// Convert this typed frame to a Zakura wire frame.
197    pub fn encode_frame(&self) -> Result<Frame, LegacyGossipError> {
198        match self {
199            Self::AdvertiseBlock(hash) => {
200                let mut payload = Vec::new();
201                hash.zcash_serialize(&mut payload)?;
202                Ok(Frame {
203                    message_type: MSG_ADVERTISE_BLOCK,
204                    flags: 0,
205                    payload,
206                })
207            }
208            Self::AdvertiseTransactionIds(ids) => {
209                let ids = outbound_tx_id_slice(ids)?;
210                let mut payload = Vec::new();
211                write_tx_id_list(&mut payload, ids)?;
212                Ok(Frame {
213                    message_type: MSG_ADVERTISE_TX_IDS,
214                    flags: 0,
215                    payload,
216                })
217            }
218        }
219    }
220
221    /// Decode a Zakura frame into a typed gossip frame.
222    pub fn decode_frame(frame: Frame) -> Result<Self, LegacyGossipError> {
223        if frame.flags != 0 {
224            return Err(LegacyGossipError::UnsupportedFlags(frame.flags));
225        }
226
227        match frame.message_type {
228            MSG_ADVERTISE_BLOCK => {
229                let mut reader = Cursor::new(frame.payload.as_slice());
230                let hash = block::Hash::zcash_deserialize(&mut reader)?;
231                reject_trailing(&reader)?;
232                Ok(Self::AdvertiseBlock(hash))
233            }
234            MSG_ADVERTISE_TX_IDS => {
235                let mut reader = Cursor::new(frame.payload.as_slice());
236                let ids = read_tx_id_list(&mut reader)?;
237                if ids.is_empty() {
238                    return Err(LegacyGossipError::EmptyTransactionAdvertisement);
239                }
240                reject_trailing(&reader)?;
241                Ok(Self::AdvertiseTransactionIds(ids))
242            }
243            message_type => Err(LegacyGossipError::UnknownMessageType(message_type)),
244        }
245    }
246
247    fn into_request(self, peer_id: ZakuraPeerId) -> Request {
248        match self {
249            Self::AdvertiseBlock(hash) => {
250                Request::AdvertiseBlock(hash, Some(PeerSource::Zakura(peer_id)))
251            }
252            Self::AdvertiseTransactionIds(ids) => Request::AdvertiseTransactionIds(
253                ids.into_iter().collect(),
254                Some(PeerSource::Zakura(peer_id)),
255            ),
256        }
257    }
258}
259
260/// A typed legacy inventory request carried by Zakura stream kind 3.
261#[derive(Clone, Debug, Eq, PartialEq)]
262pub enum LegacyRequestFrame {
263    /// Request block contents by hash.
264    BlocksByHash(Vec<block::Hash>),
265    /// Request transaction contents by unmined transaction id.
266    TransactionsById(Vec<UnminedTxId>),
267    /// Request block hashes after a locator.
268    FindBlocks {
269        /// Hashes of known blocks, ordered from highest height to lowest height.
270        known_blocks: Vec<block::Hash>,
271        /// Optional stop hash.
272        stop: Option<block::Hash>,
273    },
274    /// Request block headers after a locator.
275    FindHeaders {
276        /// Hashes of known blocks, ordered from highest height to lowest height.
277        known_blocks: Vec<block::Hash>,
278        /// Optional stop hash.
279        stop: Option<block::Hash>,
280    },
281    /// Request mempool transaction IDs.
282    MempoolTransactionIds,
283    /// Ping a peer.
284    Ping,
285    /// Push an unsolicited transaction.
286    PushTransaction(UnminedTx),
287}
288
289impl LegacyRequestFrame {
290    /// Convert a legacy network request into an inventory request frame.
291    pub fn from_request(request: Request) -> Result<Self, LegacyGossipError> {
292        match request {
293            Request::BlocksByHash(hashes) | Request::BlocksByHashFrom { hashes, .. } => {
294                let hashes = truncate_to_inventory_cap(hashes)?;
295                Ok(Self::BlocksByHash(hashes))
296            }
297            Request::TransactionsById(ids) | Request::TransactionsByIdFrom { ids, .. } => {
298                let ids = truncate_to_inventory_cap(ids)?;
299                Ok(Self::TransactionsById(ids))
300            }
301            Request::FindBlocks { known_blocks, stop } => {
302                ensure_block_locator_count(known_blocks.len())?;
303                Ok(Self::FindBlocks { known_blocks, stop })
304            }
305            Request::FindHeaders { known_blocks, stop } => {
306                ensure_block_locator_count(known_blocks.len())?;
307                Ok(Self::FindHeaders { known_blocks, stop })
308            }
309            Request::MempoolTransactionIds => Ok(Self::MempoolTransactionIds),
310            Request::Ping(_) => Ok(Self::Ping),
311            Request::PushTransaction(transaction, _) => Ok(Self::PushTransaction(transaction)),
312            request => Err(LegacyGossipError::UnsupportedRequest(request.command())),
313        }
314    }
315
316    /// Convert this typed request to a Zakura wire frame.
317    pub fn encode_frame(&self) -> Result<Frame, LegacyGossipError> {
318        match self {
319            Self::BlocksByHash(hashes) => {
320                let mut payload = Vec::new();
321                write_hash_list(&mut payload, hashes)?;
322                Ok(Frame {
323                    message_type: MSG_REQUEST_BLOCKS_BY_HASH,
324                    flags: 0,
325                    payload,
326                })
327            }
328            Self::TransactionsById(ids) => {
329                let mut payload = Vec::new();
330                write_tx_id_list(&mut payload, ids)?;
331                Ok(Frame {
332                    message_type: MSG_REQUEST_TRANSACTIONS_BY_ID,
333                    flags: 0,
334                    payload,
335                })
336            }
337            Self::FindBlocks { known_blocks, stop } => {
338                let mut payload = Vec::new();
339                write_block_locator(&mut payload, known_blocks, *stop)?;
340                Ok(Frame {
341                    message_type: MSG_REQUEST_FIND_BLOCKS,
342                    flags: 0,
343                    payload,
344                })
345            }
346            Self::FindHeaders { known_blocks, stop } => {
347                let mut payload = Vec::new();
348                write_block_locator(&mut payload, known_blocks, *stop)?;
349                Ok(Frame {
350                    message_type: MSG_REQUEST_FIND_HEADERS,
351                    flags: 0,
352                    payload,
353                })
354            }
355            Self::MempoolTransactionIds => Ok(Frame {
356                message_type: MSG_REQUEST_MEMPOOL_TRANSACTION_IDS,
357                flags: 0,
358                payload: Vec::new(),
359            }),
360            Self::Ping => Ok(Frame {
361                message_type: MSG_REQUEST_PING,
362                flags: 0,
363                payload: Vec::new(),
364            }),
365            Self::PushTransaction(transaction) => Ok(Frame {
366                message_type: MSG_REQUEST_PUSH_TRANSACTION,
367                flags: 0,
368                payload: transaction.transaction().zcash_serialize_to_vec()?,
369            }),
370        }
371    }
372
373    /// Decode a Zakura frame into a typed inventory request.
374    pub fn decode_frame(frame: Frame) -> Result<Self, LegacyGossipError> {
375        if frame.flags != 0 {
376            return Err(LegacyGossipError::UnsupportedFlags(frame.flags));
377        }
378
379        match frame.message_type {
380            MSG_REQUEST_BLOCKS_BY_HASH => {
381                let mut reader = Cursor::new(frame.payload.as_slice());
382                let hashes = read_hash_list(&mut reader)?;
383                reject_trailing(&reader)?;
384                Ok(Self::BlocksByHash(hashes))
385            }
386            MSG_REQUEST_TRANSACTIONS_BY_ID => {
387                let mut reader = Cursor::new(frame.payload.as_slice());
388                let ids = read_tx_id_list(&mut reader)?;
389                reject_trailing(&reader)?;
390                Ok(Self::TransactionsById(ids))
391            }
392            MSG_REQUEST_FIND_BLOCKS => {
393                let mut reader = Cursor::new(frame.payload.as_slice());
394                let (known_blocks, stop) = read_block_locator(&mut reader)?;
395                reject_trailing(&reader)?;
396                Ok(Self::FindBlocks { known_blocks, stop })
397            }
398            MSG_REQUEST_FIND_HEADERS => {
399                let mut reader = Cursor::new(frame.payload.as_slice());
400                let (known_blocks, stop) = read_block_locator(&mut reader)?;
401                reject_trailing(&reader)?;
402                Ok(Self::FindHeaders { known_blocks, stop })
403            }
404            MSG_REQUEST_MEMPOOL_TRANSACTION_IDS => {
405                if !frame.payload.is_empty() {
406                    return Err(LegacyGossipError::TrailingBytes);
407                }
408                Ok(Self::MempoolTransactionIds)
409            }
410            MSG_REQUEST_PING => {
411                if !frame.payload.is_empty() {
412                    return Err(LegacyGossipError::TrailingBytes);
413                }
414                Ok(Self::Ping)
415            }
416            MSG_REQUEST_PUSH_TRANSACTION => {
417                let mut reader = Cursor::new(frame.payload.as_slice());
418                let transaction = Transaction::zcash_deserialize(&mut reader)?;
419                reject_trailing(&reader)?;
420                Ok(Self::PushTransaction(UnminedTx::from(transaction)))
421            }
422            message_type => Err(LegacyGossipError::UnknownMessageType(message_type)),
423        }
424    }
425
426    fn into_service_request(self, peer_id: ZakuraPeerId) -> Option<Request> {
427        match self {
428            Self::BlocksByHash(hashes) => Some(Request::BlocksByHash(hashes.into_iter().collect())),
429            Self::TransactionsById(ids) => {
430                Some(Request::TransactionsById(ids.into_iter().collect()))
431            }
432            Self::FindBlocks { known_blocks, stop } => {
433                Some(Request::FindBlocks { known_blocks, stop })
434            }
435            Self::FindHeaders { known_blocks, stop } => {
436                Some(Request::FindHeaders { known_blocks, stop })
437            }
438            Self::MempoolTransactionIds => Some(Request::MempoolTransactionIds),
439            Self::Ping => None,
440            Self::PushTransaction(transaction) => Some(Request::PushTransaction(
441                transaction,
442                Some(PeerSource::Zakura(peer_id)),
443            )),
444        }
445    }
446
447    fn kind(&self) -> LegacyRequestKind {
448        match self {
449            Self::BlocksByHash(_) => LegacyRequestKind::Blocks,
450            Self::TransactionsById(_) => LegacyRequestKind::Transactions,
451            Self::FindBlocks { .. } => LegacyRequestKind::FindBlocks,
452            Self::FindHeaders { .. } => LegacyRequestKind::FindHeaders,
453            Self::MempoolTransactionIds => LegacyRequestKind::MempoolTransactionIds,
454            Self::Ping => LegacyRequestKind::Ping,
455            Self::PushTransaction(_) => LegacyRequestKind::PushTransaction,
456        }
457    }
458}
459
460#[derive(Copy, Clone, Debug, Eq, PartialEq)]
461pub(super) enum LegacyRequestKind {
462    Blocks,
463    Transactions,
464    FindBlocks,
465    FindHeaders,
466    MempoolTransactionIds,
467    Ping,
468    PushTransaction,
469}
470
471impl LegacyRequestKind {
472    fn command(self) -> &'static str {
473        match self {
474            LegacyRequestKind::Blocks => "BlocksByHash",
475            LegacyRequestKind::Transactions => "TransactionsById",
476            LegacyRequestKind::FindBlocks => "FindBlocks",
477            LegacyRequestKind::FindHeaders => "FindHeaders",
478            LegacyRequestKind::MempoolTransactionIds => "MempoolTransactionIds",
479            LegacyRequestKind::Ping => "Ping",
480            LegacyRequestKind::PushTransaction => "PushTransaction",
481        }
482    }
483
484    fn message_type(self) -> u16 {
485        match self {
486            LegacyRequestKind::Blocks => MSG_REQUEST_BLOCKS_BY_HASH,
487            LegacyRequestKind::Transactions => MSG_REQUEST_TRANSACTIONS_BY_ID,
488            LegacyRequestKind::FindBlocks => MSG_REQUEST_FIND_BLOCKS,
489            LegacyRequestKind::FindHeaders => MSG_REQUEST_FIND_HEADERS,
490            LegacyRequestKind::MempoolTransactionIds => MSG_REQUEST_MEMPOOL_TRANSACTION_IDS,
491            LegacyRequestKind::Ping => MSG_REQUEST_PING,
492            LegacyRequestKind::PushTransaction => MSG_REQUEST_PUSH_TRANSACTION,
493        }
494    }
495}
496
497pub(super) struct LegacyResponseCodec;
498
499impl LegacyResponseCodec {
500    pub(super) fn encode_response(
501        request_id: u64,
502        response: Response,
503        max_frame_bytes: u32,
504        max_message_bytes: u32,
505    ) -> Result<Vec<Frame>, LegacyGossipError> {
506        let mut frames = Vec::new();
507        // Bound the cumulative response so a peer that requests many available
508        // blocks/transactions cannot force us to serialize and retain an
509        // unbounded `Vec<Frame>` before the first byte is written. The budget is
510        // shared across every frame of this response and aborts encoding early.
511        let mut budget = ResponseEncodeBudget::default();
512        match response {
513            Response::Blocks(blocks) => {
514                let mut missing = Vec::new();
515                for block in blocks {
516                    match block {
517                        InventoryResponse::Available((block, _)) => {
518                            push_chunked_response(
519                                &mut frames,
520                                &mut budget,
521                                MSG_RESPONSE_BLOCK,
522                                request_id,
523                                max_frame_bytes,
524                                max_message_bytes,
525                                block.zcash_serialize_to_vec()?,
526                            )?;
527                        }
528                        InventoryResponse::Missing(hash) => missing.push(hash),
529                    }
530                }
531                if !missing.is_empty() {
532                    push_response_frame(
533                        &mut frames,
534                        &mut budget,
535                        missing_blocks_frame(request_id, missing)?,
536                    )?;
537                }
538            }
539            Response::Transactions(transactions) => {
540                let mut missing = Vec::new();
541                for transaction in transactions {
542                    match transaction {
543                        InventoryResponse::Available((transaction, _)) => {
544                            push_chunked_response(
545                                &mut frames,
546                                &mut budget,
547                                MSG_RESPONSE_TRANSACTION,
548                                request_id,
549                                max_frame_bytes,
550                                max_message_bytes,
551                                transaction.transaction().zcash_serialize_to_vec()?,
552                            )?;
553                        }
554                        InventoryResponse::Missing(id) => missing.push(id),
555                    }
556                }
557                if !missing.is_empty() {
558                    push_response_frame(
559                        &mut frames,
560                        &mut budget,
561                        missing_transactions_frame(request_id, missing)?,
562                    )?;
563                }
564            }
565            Response::BlockHashes(hashes) => {
566                // FindBlocks should already be service-capped; overflowing the wire cap is a bug.
567                push_response_frame(
568                    &mut frames,
569                    &mut budget,
570                    block_hashes_frame(request_id, hashes)?,
571                )?;
572            }
573            Response::BlockHeaders(headers) => {
574                // FindHeaders is protocol-capped by MAX_HEADERS_PER_MESSAGE; reject overflow.
575                push_response_frame(
576                    &mut frames,
577                    &mut budget,
578                    block_headers_frame(request_id, headers)?,
579                )?;
580            }
581            Response::TransactionIds(ids) => {
582                // Mempools can exceed one legacy inv response, so advertise the first capped page.
583                push_response_frame(
584                    &mut frames,
585                    &mut budget,
586                    transaction_ids_frame(request_id, truncate_to_inventory_cap(ids)?)?,
587                )?;
588            }
589            Response::Pong(_) => {
590                push_response_frame(
591                    &mut frames,
592                    &mut budget,
593                    id_only_frame(MSG_RESPONSE_PONG, request_id),
594                )?;
595            }
596            Response::Nil => {
597                push_response_frame(
598                    &mut frames,
599                    &mut budget,
600                    id_only_frame(MSG_RESPONSE_NIL, request_id),
601                )?;
602            }
603            response => return Err(LegacyGossipError::UnexpectedResponse(response.command())),
604        }
605        Ok(frames)
606    }
607
608    pub(super) fn decode_response(
609        request_id: u64,
610        request_kind: LegacyRequestKind,
611        frames: Vec<Frame>,
612        requested_block_hashes: Option<&HashSet<block::Hash>>,
613    ) -> Result<Response, LegacyGossipError> {
614        let mut blocks = Vec::new();
615        let mut transactions = Vec::new();
616        let mut block_hashes = Vec::new();
617        let mut block_headers = Vec::new();
618        let mut transaction_ids = Vec::new();
619        let mut saw_pong = false;
620        let mut saw_nil = false;
621        let mut reassembler = ResponseReassembler::new(request_id);
622
623        for frame in frames {
624            if frame.flags != 0 {
625                return Err(LegacyGossipError::UnsupportedFlags(frame.flags));
626            }
627
628            match frame.message_type {
629                MSG_RESPONSE_BLOCK => {
630                    if request_kind != LegacyRequestKind::Blocks {
631                        return Err(LegacyGossipError::UnexpectedResponse("Blocks"));
632                    }
633                    if let Some(bytes) = reassembler.accept(&frame.payload)? {
634                        let block = Arc::new(Block::zcash_deserialize(&mut Cursor::new(
635                            bytes.as_slice(),
636                        ))?);
637                        // Bind the delivered block to a hash we actually requested.
638                        // Without this, a peer can substitute any other valid block
639                        // for the one requested, corrupting downstream hash/source
640                        // accounting (the response is correlated only by request id
641                        // and kind, not by hash).
642                        if let Some(requested) = requested_block_hashes {
643                            let hash = block.hash();
644                            if !requested.contains(&hash) {
645                                return Err(LegacyGossipError::UnsolicitedBlock(hash));
646                            }
647                        }
648                        blocks.push(InventoryResponse::Available((block, None)));
649                    }
650                }
651                MSG_RESPONSE_TRANSACTION => {
652                    if request_kind != LegacyRequestKind::Transactions {
653                        return Err(LegacyGossipError::UnexpectedResponse("Transactions"));
654                    }
655                    if let Some(bytes) = reassembler.accept(&frame.payload)? {
656                        let transaction =
657                            Transaction::zcash_deserialize(&mut Cursor::new(bytes.as_slice()))?;
658                        transactions.push(InventoryResponse::Available((
659                            UnminedTx::from(transaction),
660                            None,
661                        )));
662                    }
663                }
664                MSG_RESPONSE_MISSING_BLOCKS => {
665                    if request_kind != LegacyRequestKind::Blocks {
666                        return Err(LegacyGossipError::UnexpectedResponse("MissingBlocks"));
667                    }
668                    reassembler.reject_if_active()?;
669                    for hash in decode_hashes_response(request_id, frame.payload)? {
670                        // A peer may only report blocks we requested as missing.
671                        if let Some(requested) = requested_block_hashes {
672                            if !requested.contains(&hash) {
673                                return Err(LegacyGossipError::UnsolicitedBlock(hash));
674                            }
675                        }
676                        blocks.push(InventoryResponse::Missing(hash));
677                    }
678                }
679                MSG_RESPONSE_MISSING_TRANSACTIONS => {
680                    if request_kind != LegacyRequestKind::Transactions {
681                        return Err(LegacyGossipError::UnexpectedResponse("MissingTransactions"));
682                    }
683                    reassembler.reject_if_active()?;
684                    for id in decode_tx_ids_response(request_id, frame.payload)? {
685                        transactions.push(InventoryResponse::Missing(id));
686                    }
687                }
688                MSG_RESPONSE_BLOCK_HASHES => {
689                    if request_kind != LegacyRequestKind::FindBlocks {
690                        return Err(LegacyGossipError::UnexpectedResponse("BlockHashes"));
691                    }
692                    reassembler.reject_if_active()?;
693                    block_hashes.extend(decode_hashes_response(request_id, frame.payload)?);
694                }
695                MSG_RESPONSE_BLOCK_HEADERS => {
696                    if request_kind != LegacyRequestKind::FindHeaders {
697                        return Err(LegacyGossipError::UnexpectedResponse("BlockHeaders"));
698                    }
699                    reassembler.reject_if_active()?;
700                    block_headers.extend(decode_block_headers(request_id, frame.payload)?);
701                }
702                MSG_RESPONSE_TRANSACTION_IDS => {
703                    if request_kind != LegacyRequestKind::MempoolTransactionIds {
704                        return Err(LegacyGossipError::UnexpectedResponse("TransactionIds"));
705                    }
706                    reassembler.reject_if_active()?;
707                    transaction_ids.extend(decode_tx_ids_response(request_id, frame.payload)?);
708                }
709                MSG_RESPONSE_PONG => {
710                    if request_kind != LegacyRequestKind::Ping {
711                        return Err(LegacyGossipError::UnexpectedResponse("Pong"));
712                    }
713                    reassembler.reject_if_active()?;
714                    decode_id_only_response(request_id, frame.payload)?;
715                    saw_pong = true;
716                }
717                MSG_RESPONSE_NIL => {
718                    // NIL is the empty-result sentinel for chain discovery and
719                    // mempool queries: the inbound service answers an empty
720                    // FindBlocks/FindHeaders/MempoolTransactionIds (and a queued
721                    // PushTransaction) with Response::Nil, so accept it for those
722                    // kinds and let the final match below turn it into that
723                    // kind's empty response. Inventory fetches and Ping must not
724                    // receive a bare NIL — for BlocksByHash/TransactionsById it
725                    // means the peer has none of the requested items, so reject
726                    // it as unexpected and let the caller fall back to another
727                    // peer (rather than silently returning an empty fetch).
728                    match request_kind {
729                        LegacyRequestKind::FindBlocks
730                        | LegacyRequestKind::FindHeaders
731                        | LegacyRequestKind::MempoolTransactionIds
732                        | LegacyRequestKind::PushTransaction => {}
733                        LegacyRequestKind::Blocks
734                        | LegacyRequestKind::Transactions
735                        | LegacyRequestKind::Ping => {
736                            return Err(LegacyGossipError::UnexpectedResponse("Nil"));
737                        }
738                    }
739                    reassembler.reject_if_active()?;
740                    decode_id_only_response(request_id, frame.payload)?;
741                    saw_nil = true;
742                }
743                message_type => return Err(LegacyGossipError::UnknownMessageType(message_type)),
744            }
745        }
746
747        reassembler.finish()?;
748
749        match request_kind {
750            LegacyRequestKind::Blocks if blocks.is_empty() => {
751                Err(LegacyGossipError::MissingResponse(request_kind.command()))
752            }
753            LegacyRequestKind::Blocks => Ok(Response::Blocks(blocks)),
754            LegacyRequestKind::Transactions if transactions.is_empty() => {
755                Err(LegacyGossipError::MissingResponse(request_kind.command()))
756            }
757            LegacyRequestKind::Transactions => Ok(Response::Transactions(transactions)),
758            LegacyRequestKind::FindBlocks => Ok(Response::BlockHashes(block_hashes)),
759            LegacyRequestKind::FindHeaders => Ok(Response::BlockHeaders(block_headers)),
760            LegacyRequestKind::MempoolTransactionIds => {
761                Ok(Response::TransactionIds(transaction_ids))
762            }
763            LegacyRequestKind::Ping if saw_pong => Ok(Response::Pong(Duration::ZERO)),
764            LegacyRequestKind::PushTransaction if saw_nil => Ok(Response::Nil),
765            LegacyRequestKind::Ping | LegacyRequestKind::PushTransaction => {
766                Err(LegacyGossipError::MissingResponse(request_kind.command()))
767            }
768        }
769    }
770}
771
772fn collect_outbound_tx_ids(
773    ids: impl IntoIterator<Item = UnminedTxId>,
774) -> Result<Vec<UnminedTxId>, LegacyGossipError> {
775    let max_tx_inv_in_message = usize::try_from(MAX_TX_INV_IN_SENT_MESSAGE)?;
776    let ids: Vec<_> = ids.into_iter().take(max_tx_inv_in_message).collect();
777    if ids.is_empty() {
778        return Err(LegacyGossipError::EmptyTransactionAdvertisement);
779    }
780    Ok(ids)
781}
782
783fn truncate_to_inventory_cap<T>(
784    items: impl IntoIterator<Item = T>,
785) -> Result<Vec<T>, LegacyGossipError> {
786    // Legacy `inv` uses the transaction inventory cap for both tx IDs and block hashes.
787    let max = usize::try_from(MAX_TX_INV_IN_SENT_MESSAGE)?;
788    let mut collected = Vec::new();
789    let mut dropped = 0usize;
790
791    for item in items {
792        if collected.len() < max {
793            collected.push(item);
794        } else {
795            dropped = dropped.saturating_add(1);
796        }
797    }
798
799    if dropped > 0 {
800        debug!(
801            dropped,
802            max, "legacy inventory request exceeded item cap; dropping extras"
803        );
804    }
805
806    Ok(collected)
807}
808
809fn write_hash_list(out: &mut Vec<u8>, hashes: &[block::Hash]) -> Result<(), LegacyGossipError> {
810    ensure_inventory_count(hashes.len())?;
811    CompactSizeMessage::try_from(hashes.len())?.zcash_serialize(&mut *out)?;
812    for hash in hashes {
813        hash.zcash_serialize(&mut *out)?;
814    }
815    Ok(())
816}
817
818fn read_hash_list(reader: &mut Cursor<&[u8]>) -> Result<Vec<block::Hash>, LegacyGossipError> {
819    let count = bounded_inventory_count(reader)?;
820    let mut hashes = Vec::with_capacity(count);
821    for _ in 0..count {
822        hashes.push(block::Hash::zcash_deserialize(&mut *reader)?);
823    }
824    Ok(hashes)
825}
826
827fn write_block_locator(
828    out: &mut Vec<u8>,
829    known_blocks: &[block::Hash],
830    stop: Option<block::Hash>,
831) -> Result<(), LegacyGossipError> {
832    ensure_block_locator_count(known_blocks.len())?;
833    CompactSizeMessage::try_from(known_blocks.len())?.zcash_serialize(&mut *out)?;
834    for hash in known_blocks {
835        hash.zcash_serialize(&mut *out)?;
836    }
837    stop.unwrap_or(NO_STOP_HASH).zcash_serialize(&mut *out)?;
838    Ok(())
839}
840
841fn read_block_locator(
842    reader: &mut Cursor<&[u8]>,
843) -> Result<(Vec<block::Hash>, Option<block::Hash>), LegacyGossipError> {
844    let count = bounded_block_locator_count(reader)?;
845    let mut known_blocks = Vec::with_capacity(count);
846    for _ in 0..count {
847        known_blocks.push(block::Hash::zcash_deserialize(&mut *reader)?);
848    }
849    let stop_hash = block::Hash::zcash_deserialize(&mut *reader)?;
850    // A zero stop hash is the legacy locator sentinel for "no stop hash".
851    let stop = (stop_hash != NO_STOP_HASH).then_some(stop_hash);
852    Ok((known_blocks, stop))
853}
854
855fn write_tx_id_list(out: &mut Vec<u8>, ids: &[UnminedTxId]) -> Result<(), LegacyGossipError> {
856    ensure_inventory_count(ids.len())?;
857    CompactSizeMessage::try_from(ids.len())?.zcash_serialize(&mut *out)?;
858    for id in ids {
859        InventoryHash::from(id).zcash_serialize(&mut *out)?;
860    }
861    Ok(())
862}
863
864fn read_tx_id_list(reader: &mut Cursor<&[u8]>) -> Result<Vec<UnminedTxId>, LegacyGossipError> {
865    let count = bounded_inventory_count(reader)?;
866    let mut ids = Vec::with_capacity(count);
867    for _ in 0..count {
868        let inv = InventoryHash::zcash_deserialize(&mut *reader)?;
869        let Some(id) = inv.unmined_tx_id() else {
870            return Err(LegacyGossipError::NonTransactionInventory);
871        };
872        ids.push(id);
873    }
874    Ok(ids)
875}
876
877fn bounded_inventory_count(reader: &mut Cursor<&[u8]>) -> Result<usize, LegacyGossipError> {
878    let count = usize::from(CompactSizeMessage::zcash_deserialize(reader)?);
879    ensure_inventory_count(count)?;
880    Ok(count)
881}
882
883fn ensure_inventory_count(count: usize) -> Result<(), LegacyGossipError> {
884    let max = usize::try_from(MAX_TX_INV_IN_SENT_MESSAGE)?;
885    if count > max {
886        return Err(LegacyGossipError::TooManyInventoryItems(count));
887    }
888    Ok(())
889}
890
891fn bounded_block_locator_count(reader: &mut Cursor<&[u8]>) -> Result<usize, LegacyGossipError> {
892    let count = usize::from(CompactSizeMessage::zcash_deserialize(reader)?);
893    ensure_block_locator_count(count)?;
894    Ok(count)
895}
896
897fn ensure_block_locator_count(count: usize) -> Result<(), LegacyGossipError> {
898    let max = usize::try_from(MAX_BLOCK_LOCATOR_LENGTH)?;
899    if count > max {
900        return Err(LegacyGossipError::TooManyBlockLocatorHashes(count));
901    }
902    Ok(())
903}
904
905fn write_header_list(
906    out: &mut Vec<u8>,
907    headers: &[block::CountedHeader],
908) -> Result<(), LegacyGossipError> {
909    ensure_header_count(headers.len())?;
910    CompactSizeMessage::try_from(headers.len())?.zcash_serialize(&mut *out)?;
911    for header in headers {
912        header.zcash_serialize(&mut *out)?;
913    }
914    Ok(())
915}
916
917fn read_header_list(
918    reader: &mut Cursor<&[u8]>,
919) -> Result<Vec<block::CountedHeader>, LegacyGossipError> {
920    let count = bounded_header_count(reader)?;
921    let mut headers = Vec::with_capacity(count);
922    for _ in 0..count {
923        headers.push(block::CountedHeader::zcash_deserialize(&mut *reader)?);
924    }
925    Ok(headers)
926}
927
928fn bounded_header_count(reader: &mut Cursor<&[u8]>) -> Result<usize, LegacyGossipError> {
929    let count = usize::from(CompactSizeMessage::zcash_deserialize(reader)?);
930    ensure_header_count(count)?;
931    Ok(count)
932}
933
934fn ensure_header_count(count: usize) -> Result<(), LegacyGossipError> {
935    if count > MAX_HEADERS_PER_MESSAGE {
936        return Err(LegacyGossipError::TooManyHeaders(count));
937    }
938    Ok(())
939}
940
941/// Tracks the cumulative size of an encoded legacy response so the inbound
942/// responder aborts a single request's response early instead of buffering an
943/// unbounded `Vec<Frame>` before the first byte is written. The outbound reader
944/// enforces a symmetric per-response cap via `LegacyResponseBudget`.
945#[derive(Default)]
946struct ResponseEncodeBudget {
947    bytes: usize,
948}
949
950impl ResponseEncodeBudget {
951    /// Account for one buffered response frame's payload, rejecting the whole
952    /// response once cumulative payload bytes exceed the responder aggregate
953    /// budget.
954    fn account(&mut self, payload_len: usize) -> Result<(), LegacyGossipError> {
955        self.bytes = self.bytes.saturating_add(payload_len);
956        if self.bytes > LEGACY_RESPONSE_MAX_AGGREGATE_BYTES {
957            return Err(LegacyGossipError::ResponseAggregateBudget(self.bytes));
958        }
959        Ok(())
960    }
961}
962
963/// Push one fully-built response frame, charging it against the aggregate
964/// budget first so an over-budget response aborts before it is retained.
965fn push_response_frame(
966    frames: &mut Vec<Frame>,
967    budget: &mut ResponseEncodeBudget,
968    frame: Frame,
969) -> Result<(), LegacyGossipError> {
970    budget.account(frame.payload.len())?;
971    frames.push(frame);
972    Ok(())
973}
974
975fn push_chunked_response(
976    frames: &mut Vec<Frame>,
977    budget: &mut ResponseEncodeBudget,
978    message_type: u16,
979    request_id: u64,
980    max_frame_bytes: u32,
981    max_message_bytes: u32,
982    bytes: Vec<u8>,
983) -> Result<(), LegacyGossipError> {
984    if bytes.len() > MAX_PROTOCOL_MESSAGE_LEN {
985        return Err(LegacyGossipError::OversizedResponse(bytes.len()));
986    }
987
988    // Size each chunk frame against the *effective* outbound cap: the smaller of
989    // the negotiated frame payload cap (`max_frame_bytes - FRAME_HEADER_BYTES`)
990    // and the peer's negotiated `max_message_bytes`. The handshake clamps the two
991    // caps independently, so sizing against the frame cap alone produces chunk
992    // frames whose payload exceeds the peer's accepted message cap, which the
993    // peer (and our own `write_response_frame`) reject as oversize.
994    let max_payload_bytes = usize::try_from(max_frame_bytes)?
995        .saturating_sub(FRAME_HEADER_BYTES)
996        .min(usize::try_from(max_message_bytes)?);
997    let max_chunk_bytes = max_payload_bytes
998        .checked_sub(RESPONSE_CHUNK_HEADER_BYTES)
999        .ok_or(LegacyGossipError::OversizedResponse(bytes.len()))?;
1000    if max_chunk_bytes == 0 {
1001        return Err(LegacyGossipError::OversizedResponse(bytes.len()));
1002    }
1003    let chunk_bytes = LEGACY_RESPONSE_CHUNK_BYTES.min(max_chunk_bytes);
1004
1005    debug_assert!(!bytes.is_empty());
1006    let chunks = bytes.chunks(chunk_bytes).collect::<Vec<_>>();
1007
1008    for (index, chunk) in chunks.iter().enumerate() {
1009        let mut payload = Vec::with_capacity(RESPONSE_CHUNK_HEADER_BYTES + chunk.len());
1010        // The stream prelude already has this id, but each response frame repeats it so
1011        // malformed peers are rejected at the frame boundary before decoding payload bytes.
1012        payload.extend_from_slice(&request_id.to_le_bytes());
1013        payload.push(u8::from(index + 1 == chunks.len()));
1014        payload.extend_from_slice(chunk);
1015        // Charge each chunk against the shared budget so a response that
1016        // aggregates many available items aborts mid-encode rather than after
1017        // the whole `Vec<Frame>` has been materialized.
1018        budget.account(payload.len())?;
1019        frames.push(Frame {
1020            message_type,
1021            flags: 0,
1022            payload,
1023        });
1024    }
1025    Ok(())
1026}
1027
1028struct ResponseReassembler {
1029    expected_id: u64,
1030    active: bool,
1031    buffer: Vec<u8>,
1032}
1033
1034impl ResponseReassembler {
1035    fn new(expected_id: u64) -> Self {
1036        Self {
1037            expected_id,
1038            active: false,
1039            buffer: Vec::new(),
1040        }
1041    }
1042
1043    fn accept(&mut self, payload: &[u8]) -> Result<Option<Vec<u8>>, LegacyGossipError> {
1044        let (response_id, is_last, bytes) = decode_response_chunk_header(payload)?;
1045        if response_id != self.expected_id {
1046            return Err(LegacyGossipError::WrongRequestId {
1047                expected: self.expected_id,
1048                actual: response_id,
1049            });
1050        }
1051
1052        let new_len = self
1053            .buffer
1054            .len()
1055            .checked_add(bytes.len())
1056            .ok_or(LegacyGossipError::OversizedResponse(usize::MAX))?;
1057        if new_len > MAX_PROTOCOL_MESSAGE_LEN {
1058            return Err(LegacyGossipError::OversizedResponse(new_len));
1059        }
1060        self.active = true;
1061        self.buffer.extend_from_slice(bytes);
1062
1063        if is_last {
1064            self.active = false;
1065            Ok(Some(std::mem::take(&mut self.buffer)))
1066        } else {
1067            Ok(None)
1068        }
1069    }
1070
1071    fn reject_if_active(&self) -> Result<(), LegacyGossipError> {
1072        if self.active {
1073            Err(LegacyGossipError::IncompleteResponseChunk)
1074        } else {
1075            Ok(())
1076        }
1077    }
1078
1079    fn finish(self) -> Result<(), LegacyGossipError> {
1080        if self.active {
1081            Err(LegacyGossipError::IncompleteResponseChunk)
1082        } else {
1083            Ok(())
1084        }
1085    }
1086}
1087
1088fn decode_response_chunk_header(payload: &[u8]) -> Result<(u64, bool, &[u8]), LegacyGossipError> {
1089    if payload.len() < RESPONSE_CHUNK_HEADER_BYTES {
1090        return Err(LegacyGossipError::TruncatedResponse);
1091    }
1092    let (request_id, payload) = read_request_id(payload)?;
1093    Ok((request_id, payload[0] != 0, &payload[1..]))
1094}
1095
1096fn missing_blocks_frame(
1097    request_id: u64,
1098    hashes: Vec<block::Hash>,
1099) -> Result<Frame, LegacyGossipError> {
1100    let mut payload = Vec::new();
1101    // Missing responses also repeat the id, keeping all response-frame validation local.
1102    payload.extend_from_slice(&request_id.to_le_bytes());
1103    write_hash_list(&mut payload, &hashes)?;
1104    Ok(Frame {
1105        message_type: MSG_RESPONSE_MISSING_BLOCKS,
1106        flags: 0,
1107        payload,
1108    })
1109}
1110
1111fn missing_transactions_frame(
1112    request_id: u64,
1113    ids: Vec<UnminedTxId>,
1114) -> Result<Frame, LegacyGossipError> {
1115    let mut payload = Vec::new();
1116    payload.extend_from_slice(&request_id.to_le_bytes());
1117    write_tx_id_list(&mut payload, &ids)?;
1118    Ok(Frame {
1119        message_type: MSG_RESPONSE_MISSING_TRANSACTIONS,
1120        flags: 0,
1121        payload,
1122    })
1123}
1124
1125fn block_hashes_frame(
1126    request_id: u64,
1127    hashes: Vec<block::Hash>,
1128) -> Result<Frame, LegacyGossipError> {
1129    let mut payload = Vec::new();
1130    payload.extend_from_slice(&request_id.to_le_bytes());
1131    write_hash_list(&mut payload, &hashes)?;
1132    Ok(Frame {
1133        message_type: MSG_RESPONSE_BLOCK_HASHES,
1134        flags: 0,
1135        payload,
1136    })
1137}
1138
1139fn block_headers_frame(
1140    request_id: u64,
1141    headers: Vec<block::CountedHeader>,
1142) -> Result<Frame, LegacyGossipError> {
1143    let mut payload = Vec::new();
1144    payload.extend_from_slice(&request_id.to_le_bytes());
1145    write_header_list(&mut payload, &headers)?;
1146    Ok(Frame {
1147        message_type: MSG_RESPONSE_BLOCK_HEADERS,
1148        flags: 0,
1149        payload,
1150    })
1151}
1152
1153fn transaction_ids_frame(
1154    request_id: u64,
1155    ids: Vec<UnminedTxId>,
1156) -> Result<Frame, LegacyGossipError> {
1157    let mut payload = Vec::new();
1158    payload.extend_from_slice(&request_id.to_le_bytes());
1159    write_tx_id_list(&mut payload, &ids)?;
1160    Ok(Frame {
1161        message_type: MSG_RESPONSE_TRANSACTION_IDS,
1162        flags: 0,
1163        payload,
1164    })
1165}
1166
1167fn id_only_frame(message_type: u16, request_id: u64) -> Frame {
1168    Frame {
1169        message_type,
1170        flags: 0,
1171        payload: request_id.to_le_bytes().to_vec(),
1172    }
1173}
1174
1175fn read_request_id(payload: &[u8]) -> Result<(u64, &[u8]), LegacyGossipError> {
1176    if payload.len() < REQUEST_ID_BYTES {
1177        return Err(LegacyGossipError::TruncatedResponse);
1178    }
1179    let mut id_bytes = [0; REQUEST_ID_BYTES];
1180    id_bytes.copy_from_slice(&payload[..REQUEST_ID_BYTES]);
1181    Ok((u64::from_le_bytes(id_bytes), &payload[REQUEST_ID_BYTES..]))
1182}
1183
1184fn verify_response_id(request_id: u64, payload: &[u8]) -> Result<&[u8], LegacyGossipError> {
1185    let (response_id, payload) = read_request_id(payload)?;
1186    if response_id != request_id {
1187        return Err(LegacyGossipError::WrongRequestId {
1188            expected: request_id,
1189            actual: response_id,
1190        });
1191    }
1192    Ok(payload)
1193}
1194
1195fn decode_hashes_response(
1196    request_id: u64,
1197    payload: Vec<u8>,
1198) -> Result<Vec<block::Hash>, LegacyGossipError> {
1199    let payload = verify_response_id(request_id, &payload)?;
1200    let mut reader = Cursor::new(payload);
1201    let hashes = read_hash_list(&mut reader)?;
1202    reject_trailing(&reader)?;
1203    Ok(hashes)
1204}
1205
1206fn decode_tx_ids_response(
1207    request_id: u64,
1208    payload: Vec<u8>,
1209) -> Result<Vec<UnminedTxId>, LegacyGossipError> {
1210    let payload = verify_response_id(request_id, &payload)?;
1211    let mut reader = Cursor::new(payload);
1212    let ids = read_tx_id_list(&mut reader)?;
1213    reject_trailing(&reader)?;
1214    Ok(ids)
1215}
1216
1217fn decode_block_headers(
1218    request_id: u64,
1219    payload: Vec<u8>,
1220) -> Result<Vec<block::CountedHeader>, LegacyGossipError> {
1221    let payload = verify_response_id(request_id, &payload)?;
1222    let mut reader = Cursor::new(payload);
1223    let headers = read_header_list(&mut reader)?;
1224    reject_trailing(&reader)?;
1225    Ok(headers)
1226}
1227
1228fn decode_id_only_response(request_id: u64, payload: Vec<u8>) -> Result<(), LegacyGossipError> {
1229    let payload = verify_response_id(request_id, &payload)?;
1230    if !payload.is_empty() {
1231        return Err(LegacyGossipError::TrailingBytes);
1232    }
1233    Ok(())
1234}
1235
1236fn outbound_tx_id_slice(ids: &[UnminedTxId]) -> Result<&[UnminedTxId], LegacyGossipError> {
1237    let max_tx_inv_in_message = usize::try_from(MAX_TX_INV_IN_SENT_MESSAGE)?;
1238    let ids = &ids[..ids.len().min(max_tx_inv_in_message)];
1239    if ids.is_empty() {
1240        return Err(LegacyGossipError::EmptyTransactionAdvertisement);
1241    }
1242    Ok(ids)
1243}
1244
1245fn reject_trailing(reader: &Cursor<&[u8]>) -> Result<(), LegacyGossipError> {
1246    let len = u64::try_from(reader.get_ref().len())?;
1247    if reader.position() != len {
1248        return Err(LegacyGossipError::TrailingBytes);
1249    }
1250    Ok(())
1251}
1252
1253/// Broadcasts legacy gossip frames to outbound-ready Zakura peers.
1254#[derive(Clone, Debug)]
1255pub struct ZakuraGossipBroadcast {
1256    first_seen: FirstSeenCache,
1257    outbound: LegacyGossipOutbound,
1258}
1259
1260impl ZakuraGossipBroadcast {
1261    /// Create a broadcaster from a Zakura supervisor.
1262    pub fn new(supervisor: ZakuraSupervisorHandle) -> Self {
1263        let first_seen = first_seen_for_supervisor(&supervisor);
1264        let outbound = outbound_for_supervisor(&supervisor);
1265        Self {
1266            first_seen,
1267            outbound,
1268        }
1269    }
1270
1271    async fn broadcast(
1272        &self,
1273        frame: LegacyGossipFrame,
1274        exclude: Option<&ZakuraPeerId>,
1275    ) -> Result<(), BoxError> {
1276        self.outbound.remember_latest_block(&frame);
1277        self.outbound.send_to_peers(frame, exclude).await
1278    }
1279
1280    async fn record_first_seen(&self, frame: &LegacyGossipFrame) -> Option<LegacyGossipFrame> {
1281        self.first_seen.record_first_seen(frame).await
1282    }
1283
1284    async fn unseen(&self, frame: &LegacyGossipFrame) -> Option<LegacyGossipFrame> {
1285        self.first_seen.unseen(frame).await
1286    }
1287
1288    async fn mark_seen(&self, frame: &LegacyGossipFrame) {
1289        self.first_seen.record_seen(frame).await;
1290    }
1291}
1292
1293/// Ownership claim held across the transport's ordered-stream reopen gap after a
1294/// gossip session exits while its connection stays up.
1295#[derive(Clone, Copy, Debug)]
1296struct GossipGapClaim {
1297    conn_id: ZakuraConnId,
1298}
1299
1300/// Consecutive gossip sessions on one connection that ended without delivering a
1301/// single frame before the service retires the peer's gossip stream. Generous:
1302/// honest peers send frames, and transient stream errors do not repeat eight
1303/// times in a row on a healthy connection.
1304const MAX_GOSSIP_NO_FRAME_SESSIONS: u32 = 8;
1305
1306/// Per-connection counter of consecutive zero-frame session exits, so a peer
1307/// cycling stream EOFs cannot convert endless reopen churn into a permanently
1308/// claim-pinned connection.
1309#[derive(Clone, Copy, Debug)]
1310struct GossipSessionChurn {
1311    conn_id: ZakuraConnId,
1312    no_frame_exits: u32,
1313    retired: bool,
1314}
1315
1316/// Session and reopen-gap claim state, kept under one lock so a session removal
1317/// and its gap claim are always observed together.
1318#[derive(Debug, Default)]
1319struct LegacyGossipOutboundState {
1320    sessions: HashMap<ZakuraPeerId, LegacyGossipPeerSession>,
1321    gap_claims: HashMap<ZakuraPeerId, GossipGapClaim>,
1322    churn: HashMap<ZakuraPeerId, GossipSessionChurn>,
1323}
1324
1325#[derive(Clone, Debug, Default)]
1326struct LegacyGossipOutbound {
1327    state: Arc<StdMutex<LegacyGossipOutboundState>>,
1328    latest_block: Arc<StdMutex<Option<block::Hash>>>,
1329}
1330
1331impl LegacyGossipOutbound {
1332    fn insert(&self, session: LegacyGossipPeerSession) -> bool {
1333        let mut state = self
1334            .state
1335            .lock()
1336            .expect("legacy gossip outbound mutex is never poisoned");
1337        if state.sessions.get(session.peer_id()).is_some_and(|active| {
1338            active.conn_id() > session.conn_id()
1339                || (active.conn_id() == session.conn_id()
1340                    && active.session_id() >= session.session_id())
1341        }) {
1342            return false;
1343        }
1344        // A retired connection must not be resurrected by a new stream — the
1345        // remote initiates streams on inbound connections and could otherwise
1346        // churn sessions at its own pace.
1347        if state
1348            .churn
1349            .get(session.peer_id())
1350            .is_some_and(|churn| churn.conn_id == session.conn_id() && churn.retired)
1351        {
1352            return false;
1353        }
1354        state.gap_claims.remove(session.peer_id());
1355        state.sessions.insert(session.peer_id().clone(), session);
1356        true
1357    }
1358
1359    /// Retire exactly one gossip session and claim its connection across the
1360    /// transport's reopen gap.
1361    ///
1362    /// The transport reopens gossip streams on the same connection, so a delayed
1363    /// teardown from an earlier stream session must not retire a reopened
1364    /// session: removal requires an exact `(conn_id, session_id)` match, and the
1365    /// gap claim is inserted only when that removal happened.
1366    ///
1367    /// `no_frame_exit` marks a session that ended without delivering a single
1368    /// frame; enough of those in a row retires the peer's gossip stream on this
1369    /// connection instead of claiming another reopen gap.
1370    fn finish_session(
1371        &self,
1372        peer: &ZakuraPeerId,
1373        conn_id: ZakuraConnId,
1374        session_id: u64,
1375        no_frame_exit: bool,
1376    ) {
1377        let mut state = self
1378            .state
1379            .lock()
1380            .expect("legacy gossip outbound mutex is never poisoned");
1381        let matches_exact = state.sessions.get(peer).is_some_and(|session| {
1382            session.conn_id() == conn_id && session.session_id() == session_id
1383        });
1384        if !matches_exact {
1385            return;
1386        }
1387        state.sessions.remove(peer);
1388
1389        let churn = state
1390            .churn
1391            .entry(peer.clone())
1392            .and_modify(|churn| {
1393                // A newer connection starts a fresh churn record.
1394                if churn.conn_id != conn_id {
1395                    *churn = GossipSessionChurn {
1396                        conn_id,
1397                        no_frame_exits: 0,
1398                        retired: false,
1399                    };
1400                }
1401            })
1402            .or_insert(GossipSessionChurn {
1403                conn_id,
1404                no_frame_exits: 0,
1405                retired: false,
1406            });
1407        if no_frame_exit {
1408            churn.no_frame_exits = churn.no_frame_exits.saturating_add(1);
1409            if churn.no_frame_exits >= MAX_GOSSIP_NO_FRAME_SESSIONS {
1410                churn.retired = true;
1411            }
1412        } else {
1413            churn.no_frame_exits = 0;
1414        }
1415
1416        if churn.retired {
1417            // No gap claim: the retired stream will not be reopened, so holding
1418            // connection ownership across a gap would pin the connection open
1419            // for a service that has given up on it.
1420            debug!(
1421                ?peer,
1422                conn_id, "retiring Zakura legacy-gossip stream after repeated zero-frame sessions"
1423            );
1424            state.gap_claims.remove(peer);
1425        } else {
1426            state
1427                .gap_claims
1428                .insert(peer.clone(), GossipGapClaim { conn_id });
1429        }
1430    }
1431
1432    /// Whether the peer's gossip stream on this connection has been retired
1433    /// after repeated zero-frame sessions.
1434    fn is_retired(&self, peer: &ZakuraPeerId, conn_id: ZakuraConnId) -> bool {
1435        self.state
1436            .lock()
1437            .expect("legacy gossip outbound mutex is never poisoned")
1438            .churn
1439            .get(peer)
1440            .is_some_and(|churn| churn.conn_id == conn_id && churn.retired)
1441    }
1442
1443    /// Release all state for a closed connection: the active session (whatever
1444    /// its stream session id) and any reopen-gap claim.
1445    fn remove(&self, peer: &ZakuraPeerId, conn_id: ZakuraConnId) {
1446        let mut state = self
1447            .state
1448            .lock()
1449            .expect("legacy gossip outbound mutex is never poisoned");
1450        if state
1451            .sessions
1452            .get(peer)
1453            .is_some_and(|session| session.conn_id() == conn_id)
1454        {
1455            state.sessions.remove(peer);
1456        }
1457        if state
1458            .gap_claims
1459            .get(peer)
1460            .is_some_and(|claim| claim.conn_id == conn_id)
1461        {
1462            state.gap_claims.remove(peer);
1463        }
1464        if state
1465            .churn
1466            .get(peer)
1467            .is_some_and(|churn| churn.conn_id == conn_id)
1468        {
1469            state.churn.remove(peer);
1470        }
1471    }
1472
1473    fn owns_connection(&self, peer: &ZakuraPeerId, conn_id: ZakuraConnId) -> bool {
1474        let state = self
1475            .state
1476            .lock()
1477            .expect("legacy gossip outbound mutex is never poisoned");
1478        state
1479            .sessions
1480            .get(peer)
1481            .is_some_and(|session| session.conn_id() == conn_id)
1482            || state
1483                .gap_claims
1484                .get(peer)
1485                .is_some_and(|claim| claim.conn_id == conn_id)
1486    }
1487
1488    #[cfg(test)]
1489    fn contains(&self, peer: &ZakuraPeerId) -> bool {
1490        self.state
1491            .lock()
1492            .expect("legacy gossip outbound mutex is never poisoned")
1493            .sessions
1494            .contains_key(peer)
1495    }
1496
1497    fn remember_latest_block(&self, frame: &LegacyGossipFrame) {
1498        if let LegacyGossipFrame::AdvertiseBlock(hash) = frame {
1499            *self
1500                .latest_block
1501                .lock()
1502                .expect("legacy gossip latest-block mutex is never poisoned") = Some(*hash);
1503        }
1504    }
1505
1506    async fn replay_latest_block_to_peer(
1507        &self,
1508        session: LegacyGossipPeerSession,
1509    ) -> Result<(), BoxError> {
1510        let Some(hash) = *self
1511            .latest_block
1512            .lock()
1513            .expect("legacy gossip latest-block mutex is never poisoned")
1514        else {
1515            return Ok(());
1516        };
1517
1518        session.try_send_advertise_block(hash).map_err(Into::into)
1519    }
1520
1521    async fn send_to_peers(
1522        &self,
1523        frame: LegacyGossipFrame,
1524        exclude: Option<&ZakuraPeerId>,
1525    ) -> Result<(), BoxError> {
1526        let sessions: Vec<_> = {
1527            let state = self
1528                .state
1529                .lock()
1530                .expect("legacy gossip outbound mutex is never poisoned");
1531            state
1532                .sessions
1533                .iter()
1534                .filter(|(peer_id, _)| !exclude.is_some_and(|exclude| exclude == *peer_id))
1535                .map(|(_peer_id, session)| session.clone())
1536                .collect()
1537        };
1538
1539        send_to_sessions(sessions, frame)
1540    }
1541}
1542
1543#[cfg(test)]
1544fn legacy_gossip_recv_loop_panic_target() -> &'static StdMutex<Option<ZakuraPeerId>> {
1545    static TARGET: OnceLock<StdMutex<Option<ZakuraPeerId>>> = OnceLock::new();
1546    TARGET.get_or_init(Default::default)
1547}
1548
1549#[cfg(test)]
1550fn arm_legacy_gossip_recv_loop_panic(peer: ZakuraPeerId) {
1551    *legacy_gossip_recv_loop_panic_target()
1552        .lock()
1553        .expect("legacy gossip recv-loop panic target mutex is never poisoned") = Some(peer);
1554}
1555
1556#[cfg(test)]
1557fn should_panic_legacy_gossip_recv_loop(peer: &ZakuraPeerId) -> bool {
1558    let mut target = legacy_gossip_recv_loop_panic_target()
1559        .lock()
1560        .expect("legacy gossip recv-loop panic target mutex is never poisoned");
1561    if target.as_ref() == Some(peer) {
1562        *target = None;
1563        true
1564    } else {
1565        false
1566    }
1567}
1568
1569/// Typed ordered legacy gossip sender for one peer.
1570#[derive(Clone, Debug)]
1571pub struct LegacyGossipPeerSession {
1572    peer_id: ZakuraPeerId,
1573    conn_id: ZakuraConnId,
1574    session_id: u64,
1575    send: FramedSend,
1576}
1577
1578impl LegacyGossipPeerSession {
1579    fn new(
1580        peer_id: ZakuraPeerId,
1581        conn_id: ZakuraConnId,
1582        session_id: u64,
1583        send: FramedSend,
1584    ) -> Self {
1585        Self {
1586            peer_id,
1587            conn_id,
1588            session_id,
1589            send,
1590        }
1591    }
1592
1593    /// Authenticated peer identity for this legacy gossip stream.
1594    pub fn peer_id(&self) -> &ZakuraPeerId {
1595        &self.peer_id
1596    }
1597
1598    fn conn_id(&self) -> ZakuraConnId {
1599        self.conn_id
1600    }
1601
1602    fn session_id(&self) -> u64 {
1603        self.session_id
1604    }
1605
1606    /// Send an explicit block advertisement.
1607    pub fn try_send_advertise_block(&self, hash: block::Hash) -> Result<(), OrderedSendError> {
1608        self.try_send_gossip_frame(LegacyGossipFrame::AdvertiseBlock(hash))
1609    }
1610
1611    /// Send explicit transaction id advertisements.
1612    pub fn try_send_advertise_transaction_ids(
1613        &self,
1614        ids: Vec<UnminedTxId>,
1615    ) -> Result<(), OrderedSendError> {
1616        self.try_send_gossip_frame(LegacyGossipFrame::AdvertiseTransactionIds(ids))
1617    }
1618
1619    fn try_send_gossip_frame(&self, frame: LegacyGossipFrame) -> Result<(), OrderedSendError> {
1620        let frame = frame
1621            .encode_frame()
1622            .map_err(|error| OrderedSendError::Encode(Box::new(error)))?;
1623        match self.send.try_send(frame) {
1624            Ok(()) => Ok(()),
1625            Err(mpsc::error::TrySendError::Full(_frame)) => Err(OrderedSendError::Full),
1626            Err(mpsc::error::TrySendError::Closed(_frame)) => Err(OrderedSendError::Closed),
1627        }
1628    }
1629}
1630
1631static LEGACY_GOSSIP_OUTBOUND_BY_SUPERVISOR: OnceLock<
1632    std::sync::Mutex<HashMap<u64, LegacyGossipOutbound>>,
1633> = OnceLock::new();
1634
1635fn outbound_for_supervisor(supervisor: &ZakuraSupervisorHandle) -> LegacyGossipOutbound {
1636    let registry = LEGACY_GOSSIP_OUTBOUND_BY_SUPERVISOR.get_or_init(Default::default);
1637    let mut registry = registry
1638        .lock()
1639        .expect("legacy gossip outbound registry mutex is never poisoned");
1640    registry.entry(supervisor.id()).or_default().clone()
1641}
1642
1643fn first_seen_for_supervisor(supervisor: &ZakuraSupervisorHandle) -> FirstSeenCache {
1644    let registry = FIRST_SEEN_BY_SUPERVISOR.get_or_init(Default::default);
1645    let mut registry = registry
1646        .lock()
1647        .expect("legacy gossip first-seen registry mutex is never poisoned");
1648    registry
1649        .entry(supervisor.id())
1650        .or_insert_with(|| FirstSeenCache::new(DEFAULT_FIRST_SEEN_CAPACITY, DEFAULT_FIRST_SEEN_TTL))
1651        .clone()
1652}
1653
1654fn send_to_sessions(
1655    sessions: Vec<LegacyGossipPeerSession>,
1656    frame: LegacyGossipFrame,
1657) -> Result<(), BoxError> {
1658    let mut first_error = None;
1659    for session in sessions {
1660        let result = match &frame {
1661            LegacyGossipFrame::AdvertiseBlock(hash) => session.try_send_advertise_block(*hash),
1662            LegacyGossipFrame::AdvertiseTransactionIds(ids) => {
1663                session.try_send_advertise_transaction_ids(ids.clone())
1664            }
1665        };
1666
1667        if let Err(error) = result {
1668            debug!(
1669                peer = ?session.peer_id(),
1670                ?error,
1671                "failed to queue Zakura legacy gossip frame"
1672            );
1673            if first_error.is_none() {
1674                first_error = Some(error.into());
1675            }
1676        }
1677    }
1678
1679    if let Some(error) = first_error {
1680        return Err(error);
1681    }
1682    Ok(())
1683}
1684
1685/// Tower service that adapts legacy Zebra gossip requests onto Zakura streams.
1686#[derive(Clone, Debug)]
1687pub struct LegacyGossipAdapter {
1688    broadcast: ZakuraGossipBroadcast,
1689}
1690
1691impl LegacyGossipAdapter {
1692    /// Create an adapter backed by the given Zakura supervisor.
1693    pub fn new(supervisor: ZakuraSupervisorHandle) -> Self {
1694        Self {
1695            broadcast: ZakuraGossipBroadcast::new(supervisor),
1696        }
1697    }
1698}
1699
1700impl Service<Request> for LegacyGossipAdapter {
1701    type Response = Response;
1702    type Error = BoxError;
1703    type Future = Pin<Box<dyn Future<Output = Result<Response, BoxError>> + Send + 'static>>;
1704
1705    fn poll_ready(&mut self, _cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
1706        Poll::Ready(Ok(()))
1707    }
1708
1709    fn call(&mut self, request: Request) -> Self::Future {
1710        let broadcast = self.broadcast.clone();
1711        Box::pin(async move {
1712            let frame = LegacyGossipFrame::from_request(request)?;
1713            let Some(frame) = broadcast.record_first_seen(&frame).await else {
1714                return Ok(Response::Nil);
1715            };
1716            broadcast.broadcast(frame, None).await?;
1717            Ok(Response::Nil)
1718        })
1719    }
1720}
1721
1722/// Tower service that adapts legacy request/response traffic onto Zakura streams.
1723///
1724/// Legacy `Peers` discovery is intentionally not adapted here. Native peers use
1725/// the signed Zakura discovery protocol, while legacy address discovery stays on
1726/// the legacy peer set, keeping this compatibility layer narrow.
1727#[derive(Clone, Debug)]
1728pub struct LegacyRequestAdapter {
1729    client: ZakuraRequestClient,
1730}
1731
1732impl LegacyRequestAdapter {
1733    /// Create an adapter backed by the given Zakura supervisor.
1734    pub fn new(supervisor: ZakuraSupervisorHandle) -> Self {
1735        Self {
1736            client: ZakuraRequestClient::new(supervisor),
1737        }
1738    }
1739
1740    /// Create an adapter backed by the given Zakura supervisor and trace emitter.
1741    pub fn new_with_trace(supervisor: ZakuraSupervisorHandle, trace: ZakuraTrace) -> Self {
1742        Self {
1743            client: ZakuraRequestClient::new_with_trace(supervisor, trace),
1744        }
1745    }
1746
1747    fn new_with_trace_and_stream_rate(
1748        supervisor: ZakuraSupervisorHandle,
1749        trace: ZakuraTrace,
1750        stream_open_rate_per_second: u32,
1751    ) -> Self {
1752        Self {
1753            client: ZakuraRequestClient::new_with_trace_and_stream_rate(
1754                supervisor,
1755                trace,
1756                stream_open_rate_per_second,
1757            ),
1758        }
1759    }
1760
1761    #[cfg(test)]
1762    fn new_with_timeout(supervisor: ZakuraSupervisorHandle, request_timeout: Duration) -> Self {
1763        Self {
1764            client: ZakuraRequestClient::new_with_timeout(supervisor, request_timeout),
1765        }
1766    }
1767
1768    /// Request inventory, preferring the Zakura peer that advertised it when known.
1769    pub async fn request_from_source(
1770        &self,
1771        request: Request,
1772        source: Option<PeerSource>,
1773    ) -> Result<Response, BoxError> {
1774        let frame = LegacyRequestFrame::from_request(request)?;
1775        self.client.request(frame, source, false).await
1776    }
1777}
1778
1779impl Service<Request> for LegacyRequestAdapter {
1780    type Response = Response;
1781    type Error = BoxError;
1782    type Future = Pin<Box<dyn Future<Output = Result<Response, BoxError>> + Send + 'static>>;
1783
1784    fn poll_ready(&mut self, _cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
1785        Poll::Ready(Ok(()))
1786    }
1787
1788    fn call(&mut self, request: Request) -> Self::Future {
1789        let client = self.client.clone();
1790        Box::pin(async move {
1791            let source = request.inventory_source();
1792            let retry_source_missing = source.is_some();
1793            let frame = LegacyRequestFrame::from_request(request)?;
1794            client.request(frame, source, retry_source_missing).await
1795        })
1796    }
1797}
1798
1799/// Which transports a [`ZakuraDualStackService`] request fans out to.
1800#[derive(Copy, Clone)]
1801enum DualStackRoute {
1802    /// Block/transaction advertisements: fan out to both stacks, fire-and-forget.
1803    Advertise,
1804    /// Requests the Zakura adapter can serve (inventory fetches, chain-sync
1805    /// discovery, and mempool data): legacy peer set first, then Zakura fallback.
1806    ///
1807    /// This must cover every request a node needs to follow the chain, otherwise
1808    /// a node whose only peer was upgraded to Zakura sends them to its empty
1809    /// legacy peer set and stalls (its syncer can never obtain tips or fetch
1810    /// blocks). The set mirrors `LegacyRequestFrame::from_request`.
1811    LegacyFirstThenZakura,
1812    /// Anything the Zakura adapter cannot serve (e.g. peer-address discovery):
1813    /// legacy peer set only.
1814    Passthrough,
1815}
1816
1817/// Tower service that fans legacy gossip/inventory requests across both the
1818/// legacy TCP peer set and the Zakura P2P-v2 transport.
1819///
1820/// This is the production seam that turns the Zakura legacy-gossip adapters from
1821/// an isolated library into a live dual-stack: it wraps the legacy `peer_set`
1822/// returned by [`crate::init`] when `v2_p2p` is enabled, so the syncer, mempool,
1823/// and inbound downloads transparently gossip and fetch over Zakura too.
1824///
1825/// Routing (with `legacy_enabled` = `config.legacy_p2p()`):
1826/// - Advertisements fan out concurrently to the legacy peer set (when
1827///   `legacy_enabled`) and the Zakura gossip adapter. Advertise is
1828///   fire-and-forget, so per-path errors are logged and swallowed and the
1829///   composite always reports success — one stack failing never fails the local
1830///   advertisement.
1831/// - Inventory fetches hit the legacy peer set first and fall back to Zakura
1832///   when the legacy response is all-missing or errors. With `legacy_enabled`
1833///   false they route straight to Zakura.
1834/// - Every other request passes through to the legacy peer set.
1835///
1836/// The gossip and request adapters MUST be built from the same
1837/// [`ZakuraSupervisorHandle`] that backs the inbound [`LegacyGossipSink`], so the
1838/// per-supervisor first-seen cache (see [`first_seen_for_supervisor`]) dedups
1839/// locally originated gossip against echoes that peers send back to us.
1840#[derive(Clone)]
1841pub(crate) struct ZakuraDualStackService<L> {
1842    legacy: L,
1843    gossip: LegacyGossipAdapter,
1844    request: LegacyRequestAdapter,
1845    legacy_enabled: bool,
1846}
1847
1848impl<L> ZakuraDualStackService<L> {
1849    /// Wrap `legacy` with Zakura gossip/request adapters backed by `supervisor`.
1850    ///
1851    /// `supervisor` must be the endpoint's supervisor so the first-seen cache is
1852    /// shared with the inbound sink.
1853    #[cfg(test)]
1854    pub(crate) fn new(legacy: L, supervisor: ZakuraSupervisorHandle, legacy_enabled: bool) -> Self {
1855        Self {
1856            legacy,
1857            gossip: LegacyGossipAdapter::new(supervisor.clone()),
1858            request: LegacyRequestAdapter::new(supervisor),
1859            legacy_enabled,
1860        }
1861    }
1862
1863    /// Wrap `legacy` with Zakura adapters and a trace emitter.
1864    pub(crate) fn new_with_trace(
1865        legacy: L,
1866        supervisor: ZakuraSupervisorHandle,
1867        legacy_enabled: bool,
1868        trace: ZakuraTrace,
1869        stream_open_rate_per_second: u32,
1870    ) -> Self {
1871        Self {
1872            legacy,
1873            gossip: LegacyGossipAdapter::new(supervisor.clone()),
1874            request: LegacyRequestAdapter::new_with_trace_and_stream_rate(
1875                supervisor,
1876                trace,
1877                stream_open_rate_per_second,
1878            ),
1879            legacy_enabled,
1880        }
1881    }
1882}
1883
1884impl<L> Service<Request> for ZakuraDualStackService<L>
1885where
1886    L: Service<Request, Response = Response, Error = BoxError> + Clone + Send + 'static,
1887    L::Future: Send + 'static,
1888{
1889    type Response = Response;
1890    type Error = BoxError;
1891    type Future = Pin<Box<dyn Future<Output = Result<Response, BoxError>> + Send + 'static>>;
1892
1893    fn poll_ready(&mut self, cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
1894        // The route is request-dependent, so legacy readiness is awaited inside
1895        // the branches that actually use legacy. Source-aware Zakura requests
1896        // must not block behind an empty or unready legacy peer set.
1897        std::task::ready!(self.gossip.poll_ready(cx))?;
1898        std::task::ready!(self.request.poll_ready(cx))?;
1899        Poll::Ready(Ok(()))
1900    }
1901
1902    fn call(&mut self, request: Request) -> Self::Future {
1903        let route = match &request {
1904            Request::AdvertiseBlock(..)
1905            | Request::AdvertiseBlockToAll(..)
1906            | Request::AdvertiseTransactionIds(..) => DualStackRoute::Advertise,
1907            Request::BlocksByHash(..)
1908            | Request::BlocksByHashFrom { .. }
1909            | Request::TransactionsById(..)
1910            | Request::TransactionsByIdFrom { .. }
1911            | Request::FindBlocks { .. }
1912            | Request::FindHeaders { .. }
1913            | Request::MempoolTransactionIds
1914            | Request::PushTransaction(..) => DualStackRoute::LegacyFirstThenZakura,
1915            _ => DualStackRoute::Passthrough,
1916        };
1917
1918        let legacy_enabled = self.legacy_enabled;
1919        let mut gossip = self.gossip.clone();
1920        let mut request_adapter = self.request.clone();
1921        let mut legacy = self.legacy.clone();
1922
1923        Box::pin(async move {
1924            match route {
1925                DualStackRoute::Advertise => {
1926                    // Fan out concurrently so a slow/empty Zakura path can't delay
1927                    // legacy advertise, and vice versa. The gossip adapter bounds
1928                    // each peer send with its own per-peer fanout timeout.
1929                    let gossip_fut = gossip.call(request.clone());
1930                    let legacy_fut = async move {
1931                        if legacy_enabled {
1932                            Some(match legacy.ready().await {
1933                                Ok(service) => service.call(request).await,
1934                                Err(error) => Err(error),
1935                            })
1936                        } else {
1937                            None
1938                        }
1939                    };
1940                    let (legacy_res, gossip_res) = futures::join!(legacy_fut, gossip_fut);
1941                    if let Some(Err(error)) = legacy_res {
1942                        debug!(%error, "legacy advertise path failed; continuing");
1943                    }
1944                    if let Err(error) = gossip_res {
1945                        debug!(%error, "zakura advertise path failed; continuing");
1946                    }
1947                    Ok(Response::Nil)
1948                }
1949                DualStackRoute::LegacyFirstThenZakura => {
1950                    if matches!(request.inventory_source(), Some(PeerSource::Zakura(_))) {
1951                        return request_adapter.call(request).await;
1952                    }
1953                    if !legacy_enabled {
1954                        return request_adapter.call(request).await;
1955                    }
1956                    // The legacy peer set is buffered, so it queues this fetch
1957                    // even when it has no ready peer (e.g. every legacy peer was
1958                    // upgraded to Zakura). Bound the legacy attempt so we fall
1959                    // back to the Zakura path instead of blocking forever.
1960                    let legacy_attempt = timeout(DUAL_STACK_LEGACY_INVENTORY_TIMEOUT, async {
1961                        legacy.ready().await?.call(request.clone()).await
1962                    })
1963                    .await;
1964                    match legacy_attempt {
1965                        Ok(Ok(response)) if !all_inventory_missing(&response) => Ok(response),
1966                        Ok(Ok(response)) => match request_adapter.call(request).await {
1967                            Ok(zakura) if !all_inventory_missing(&zakura) => Ok(zakura),
1968                            _ => Ok(response),
1969                        },
1970                        Ok(Err(legacy_error)) => match request_adapter.call(request).await {
1971                            Ok(zakura) => Ok(zakura),
1972                            Err(_) => Err(legacy_error),
1973                        },
1974                        // Legacy timed out (no ready peer): use Zakura.
1975                        Err(_) => request_adapter.call(request).await,
1976                    }
1977                }
1978                DualStackRoute::Passthrough => legacy.ready().await?.call(request).await,
1979            }
1980        })
1981    }
1982}
1983
1984/// Outbound Zakura inventory request client.
1985#[derive(Clone, Debug)]
1986pub struct ZakuraRequestClient {
1987    supervisor: ZakuraSupervisorHandle,
1988    request_timeout: Duration,
1989    trace: ZakuraTrace,
1990    request_interval: Duration,
1991    next_request_at: Arc<Mutex<Instant>>,
1992}
1993
1994impl ZakuraRequestClient {
1995    /// Create a client from a Zakura supervisor.
1996    pub fn new(supervisor: ZakuraSupervisorHandle) -> Self {
1997        Self::new_with_trace(supervisor, ZakuraTrace::noop())
1998    }
1999
2000    /// Create a client from a Zakura supervisor and trace emitter.
2001    pub fn new_with_trace(supervisor: ZakuraSupervisorHandle, trace: ZakuraTrace) -> Self {
2002        Self::new_with_trace_and_stream_rate(
2003            supervisor,
2004            trace,
2005            super::DEFAULT_ZAKURA_STREAM_OPEN_RATE_PER_SECOND,
2006        )
2007    }
2008
2009    fn new_with_trace_and_stream_rate(
2010        supervisor: ZakuraSupervisorHandle,
2011        trace: ZakuraTrace,
2012        stream_open_rate_per_second: u32,
2013    ) -> Self {
2014        let request_rate = legacy_request_stream_rate(stream_open_rate_per_second);
2015        let interval_nanos = 1_000_000_000u64
2016            .checked_div(u64::from(request_rate))
2017            .unwrap_or(1)
2018            .max(1);
2019        Self {
2020            supervisor,
2021            request_timeout: LEGACY_REQUEST_TIMEOUT,
2022            trace,
2023            request_interval: Duration::from_nanos(interval_nanos),
2024            next_request_at: Arc::new(Mutex::new(Instant::now())),
2025        }
2026    }
2027
2028    #[cfg(test)]
2029    fn new_with_timeout(supervisor: ZakuraSupervisorHandle, request_timeout: Duration) -> Self {
2030        let mut client = Self::new(supervisor);
2031        client.request_timeout = request_timeout;
2032        client
2033    }
2034
2035    async fn wait_for_request_slot(&self) {
2036        let request_at = {
2037            let mut next_request_at = self.next_request_at.lock().await;
2038            let request_at = (*next_request_at).max(Instant::now());
2039            *next_request_at = request_at + self.request_interval;
2040            request_at
2041        };
2042        tokio::time::sleep_until(request_at).await;
2043    }
2044
2045    async fn request(
2046        &self,
2047        frame: LegacyRequestFrame,
2048        source: Option<PeerSource>,
2049        retry_source_missing: bool,
2050    ) -> Result<Response, BoxError> {
2051        let preferred = match source {
2052            Some(PeerSource::Zakura(peer_id)) => Some(peer_id),
2053            _ => None,
2054        };
2055
2056        let handles = self.ready_handles().await?;
2057        let Some(primary) = select_handle(&handles, preferred.as_ref()) else {
2058            return Err("no ready Zakura peer for legacy inventory request".into());
2059        };
2060
2061        let request_kind = frame.kind();
2062        let first_result = self
2063            .request_one(primary.clone(), frame.clone(), request_kind)
2064            .await;
2065        match first_result {
2066            Ok(response) if !all_inventory_missing(&response) => Ok(response),
2067            Ok(response) => {
2068                let Some(fallback) = select_fallback_handle(&handles, primary.peer_id()) else {
2069                    if retry_source_missing && preferred.is_some() {
2070                        return self
2071                            .retry_source_missing(primary, frame, request_kind, response)
2072                            .await;
2073                    }
2074                    return Ok(response);
2075                };
2076                self.request_one(fallback, frame, request_kind)
2077                    .await
2078                    .or(Ok(response))
2079            }
2080            Err(error) => {
2081                let Some(fallback) = select_fallback_handle(&handles, primary.peer_id()) else {
2082                    return Err(error);
2083                };
2084                self.request_one(fallback, frame, request_kind).await
2085            }
2086        }
2087    }
2088
2089    async fn retry_source_missing(
2090        &self,
2091        handle: ZakuraPeerHandle,
2092        frame: LegacyRequestFrame,
2093        request_kind: LegacyRequestKind,
2094        initial_response: Response,
2095    ) -> Result<Response, BoxError> {
2096        let mut last_response = initial_response;
2097        for _ in 0..SOURCE_INVENTORY_MISSING_RETRIES {
2098            sleep(SOURCE_INVENTORY_MISSING_RETRY_DELAY).await;
2099            match self
2100                .request_one(handle.clone(), frame.clone(), request_kind)
2101                .await
2102            {
2103                Ok(response) if !all_inventory_missing(&response) => return Ok(response),
2104                Ok(response) => last_response = response,
2105                Err(error) => return Err(error),
2106            }
2107        }
2108        Ok(last_response)
2109    }
2110
2111    async fn ready_handles(&self) -> Result<Vec<ZakuraPeerHandle>, BoxError> {
2112        let mut registered = self.supervisor.subscribe();
2113
2114        timeout(LEGACY_REQUEST_READY_TIMEOUT, async {
2115            loop {
2116                let handles = self.supervisor.outbound_peer_handles().await;
2117                if !handles.is_empty() {
2118                    return Ok(handles);
2119                }
2120
2121                if registered.changed().await.is_err() {
2122                    return Err("Zakura peer set closed before a peer was ready".into());
2123                }
2124            }
2125        })
2126        .await
2127        .map_err(|_| -> BoxError { "no ready Zakura peer for legacy inventory request".into() })?
2128    }
2129
2130    async fn request_one(
2131        &self,
2132        handle: ZakuraPeerHandle,
2133        frame: LegacyRequestFrame,
2134        request_kind: LegacyRequestKind,
2135    ) -> Result<Response, BoxError> {
2136        self.wait_for_request_slot().await;
2137        let request_id = NEXT_LEGACY_REQUEST_ID.fetch_add(1, Ordering::Relaxed);
2138        // Capture the requested block hashes (if any) before consuming the frame,
2139        // so the response can be bound to a hash we actually asked for.
2140        let requested_block_hashes: Option<HashSet<block::Hash>> = match &frame {
2141            LegacyRequestFrame::BlocksByHash(hashes) => Some(hashes.iter().copied().collect()),
2142            _ => None,
2143        };
2144        let frame = frame.encode_frame()?;
2145        self.trace.emit_event(|| {
2146            LegacyRequestStart::new(
2147                "outbound.request",
2148                Some(handle.peer_id()),
2149                request_id,
2150                request_kind,
2151                frame.message_type,
2152            )
2153        });
2154        let started_at = Instant::now();
2155        let response = match timeout(
2156            self.request_timeout,
2157            handle.request(
2158                ZAKURA_STREAM_LEGACY_REQUESTS,
2159                request_id,
2160                frame.message_type,
2161                frame.flags,
2162                frame.payload,
2163            ),
2164        )
2165        .await
2166        {
2167            Ok(Ok(response)) => response,
2168            Ok(Err(error)) => {
2169                self.trace.emit_event(|| {
2170                    LegacyRequestError::new(
2171                        "outbound.error",
2172                        Some(handle.peer_id()),
2173                        request_id,
2174                        request_kind.command(),
2175                        error.to_string(),
2176                    )
2177                });
2178                return Err(error);
2179            }
2180            Err(_) => {
2181                let error: BoxError = format!(
2182                    "Zakura legacy request timed out for peer {:?}",
2183                    handle.peer_id()
2184                )
2185                .into();
2186                self.trace.emit_event(|| {
2187                    LegacyRequestError::new(
2188                        "outbound.error",
2189                        Some(handle.peer_id()),
2190                        request_id,
2191                        request_kind.command(),
2192                        error.to_string(),
2193                    )
2194                });
2195                return Err(error);
2196            }
2197        };
2198        let mut response = match LegacyResponseCodec::decode_response(
2199            request_id,
2200            request_kind,
2201            response,
2202            requested_block_hashes.as_ref(),
2203        ) {
2204            Ok(response) => response,
2205            Err(error) => {
2206                self.trace.emit_event(|| {
2207                    LegacyRequestError::new(
2208                        "outbound.decode_error",
2209                        Some(handle.peer_id()),
2210                        request_id,
2211                        request_kind.command(),
2212                        error.to_string(),
2213                    )
2214                });
2215                return Err(error.into());
2216            }
2217        };
2218        if request_kind == LegacyRequestKind::Ping {
2219            // The responder can only acknowledge a Ping; the requester stamps the RTT.
2220            response = Response::Pong(started_at.elapsed());
2221        }
2222        self.trace.emit_event(|| {
2223            LegacyRequestResponse::new(
2224                "outbound.response",
2225                Some(handle.peer_id()),
2226                request_id,
2227                request_kind.command(),
2228                &response,
2229            )
2230        });
2231        Ok(response)
2232    }
2233}
2234
2235fn legacy_request_stream_rate(stream_open_rate_per_second: u32) -> u32 {
2236    stream_open_rate_per_second
2237        .max(1)
2238        .div_ceil(LEGACY_REQUEST_STREAM_RATE_DIVISOR)
2239}
2240
2241fn select_handle(
2242    handles: &[ZakuraPeerHandle],
2243    preferred: Option<&ZakuraPeerId>,
2244) -> Option<ZakuraPeerHandle> {
2245    // v1 controlled-network routing: source-less chain-sync requests use the
2246    // first ready Zakura peer. Production spreading/rotation belongs with the
2247    // future zakurad wiring milestone.
2248    preferred
2249        .and_then(|peer_id| {
2250            handles
2251                .iter()
2252                .find(|handle| handle.peer_id() == peer_id)
2253                .cloned()
2254        })
2255        .or_else(|| handles.first().cloned())
2256}
2257
2258fn select_fallback_handle(
2259    handles: &[ZakuraPeerHandle],
2260    primary: &ZakuraPeerId,
2261) -> Option<ZakuraPeerHandle> {
2262    handles
2263        .iter()
2264        .find(|handle| handle.peer_id() != primary)
2265        .cloned()
2266}
2267
2268fn all_inventory_missing(response: &Response) -> bool {
2269    match response {
2270        Response::Blocks(blocks) => {
2271            blocks.is_empty() || blocks.iter().all(|block| block.is_missing())
2272        }
2273        Response::Transactions(transactions) => {
2274            transactions.is_empty()
2275                || transactions
2276                    .iter()
2277                    .all(|transaction| transaction.is_missing())
2278        }
2279        Response::BlockHashes(hashes) => hashes.is_empty(),
2280        Response::BlockHeaders(headers) => headers.is_empty(),
2281        Response::TransactionIds(ids) => ids.is_empty(),
2282        _ => false,
2283    }
2284}
2285
2286fn bounded_u64(value: usize) -> u64 {
2287    u64::try_from(value).unwrap_or(u64::MAX)
2288}
2289
2290#[derive(Clone, Debug)]
2291struct LegacyGossipForwarder {
2292    broadcast: ZakuraGossipBroadcast,
2293    /// Short-lived dedup of inventory already handed to the inbound service but not
2294    /// yet confirmed `mark_seen`. Bounds duplicate work when the service is slow,
2295    /// not ready, or erroring (see [`LEGACY_GOSSIP_DUPLICATE_COOLDOWN`]).
2296    attempt_cooldown: FirstSeenCache,
2297}
2298
2299impl LegacyGossipForwarder {
2300    fn new(supervisor: ZakuraSupervisorHandle) -> Self {
2301        Self {
2302            broadcast: ZakuraGossipBroadcast::new(supervisor),
2303            attempt_cooldown: FirstSeenCache::new(
2304                DEFAULT_FIRST_SEEN_CAPACITY,
2305                LEGACY_GOSSIP_DUPLICATE_COOLDOWN,
2306            ),
2307        }
2308    }
2309
2310    async fn unseen(&self, frame: &LegacyGossipFrame) -> Option<LegacyGossipFrame> {
2311        self.broadcast.unseen(frame).await
2312    }
2313
2314    /// Return the subset of `frame`'s inventory not currently suppressed by a recent
2315    /// expensive failed attempt. Read-only: the cooldown is populated only by
2316    /// [`Self::note_failed_attempt`], so fast failures remain immediately retryable.
2317    async fn fresh_attempt(&self, frame: &LegacyGossipFrame) -> Option<LegacyGossipFrame> {
2318        self.attempt_cooldown.unseen(frame).await
2319    }
2320
2321    /// Record `frame`'s inventory in the cooldown when a failed attempt consumed at
2322    /// least [`LEGACY_GOSSIP_EXPENSIVE_ATTEMPT`], so queued duplicates skip re-paying
2323    /// a slow readiness/call. Cheap failures are left immediately retryable.
2324    async fn note_failed_attempt(&self, frame: &LegacyGossipFrame, elapsed: Duration) {
2325        if elapsed >= LEGACY_GOSSIP_EXPENSIVE_ATTEMPT {
2326            self.attempt_cooldown.record_seen(frame).await;
2327        }
2328    }
2329
2330    async fn mark_seen(&self, frame: &LegacyGossipFrame) {
2331        self.broadcast.mark_seen(frame).await;
2332    }
2333
2334    async fn forward(
2335        &self,
2336        peer_id: &ZakuraPeerId,
2337        frame: LegacyGossipFrame,
2338    ) -> Result<(), BoxError> {
2339        self.broadcast.broadcast(frame, Some(peer_id)).await
2340    }
2341}
2342
2343/// Inbound sink that adapts Zakura gossip frames into Zebra's legacy inbound service.
2344#[derive(Debug)]
2345pub struct LegacyGossipSink {
2346    inbound_tx: mpsc::Sender<LegacyInboundWork>,
2347    outbound: LegacyGossipOutbound,
2348    trace: ZakuraTrace,
2349}
2350
2351impl LegacyGossipSink {
2352    /// Spawn a bounded inbound worker around the existing legacy inbound service.
2353    pub fn spawn<Inbound>(inbound: Inbound, supervisor: ZakuraSupervisorHandle) -> Self
2354    where
2355        Inbound: Service<Request, Response = Response, Error = BoxError> + Send + Clone + 'static,
2356        Inbound::Future: Send + 'static,
2357    {
2358        Self::spawn_with_trace(inbound, supervisor, ZakuraTrace::noop())
2359    }
2360
2361    /// Spawn a bounded inbound worker with a trace emitter.
2362    pub fn spawn_with_trace<Inbound>(
2363        inbound: Inbound,
2364        supervisor: ZakuraSupervisorHandle,
2365        trace: ZakuraTrace,
2366    ) -> Self
2367    where
2368        Inbound: Service<Request, Response = Response, Error = BoxError> + Send + Clone + 'static,
2369        Inbound::Future: Send + 'static,
2370    {
2371        let (inbound_tx, inbound_rx) = mpsc::channel(LEGACY_GOSSIP_INBOUND_QUEUE);
2372        tokio::spawn(legacy_gossip_worker(
2373            inbound,
2374            inbound_rx,
2375            LegacyGossipForwarder::new(supervisor.clone()),
2376            trace.clone(),
2377        ));
2378        let outbound = outbound_for_supervisor(&supervisor);
2379        Self {
2380            inbound_tx,
2381            outbound,
2382            trace,
2383        }
2384    }
2385}
2386
2387impl LegacyGossipSink {
2388    fn enqueue_gossip_frame(
2389        inbound_tx: &mpsc::Sender<LegacyInboundWork>,
2390        peer_id: ZakuraPeerId,
2391        frame: Frame,
2392    ) -> Result<(), SinkReject> {
2393        let frame = LegacyGossipFrame::decode_frame(frame).map_err(SinkReject::protocol)?;
2394        match inbound_tx.try_send(LegacyInboundWork::Gossip(LegacyGossipInbound {
2395            peer_id,
2396            frame,
2397        })) {
2398            Ok(()) => Ok(()),
2399            Err(mpsc::error::TrySendError::Full(_)) => {
2400                debug!("legacy gossip inbound queue full: dropping frame");
2401                Ok(())
2402            }
2403            Err(mpsc::error::TrySendError::Closed(_)) => {
2404                Err(SinkReject::local("legacy gossip inbound queue closed"))
2405            }
2406        }
2407    }
2408
2409    fn deliver(
2410        &self,
2411        peer_id: ZakuraPeerId,
2412        stream_kind: u16,
2413        frame: Frame,
2414    ) -> Result<(), SinkReject> {
2415        if stream_kind != ZAKURA_STREAM_GOSSIP {
2416            return Ok(());
2417        }
2418
2419        Self::enqueue_gossip_frame(&self.inbound_tx, peer_id, frame)
2420    }
2421
2422    fn request<'a>(
2423        &'a self,
2424        peer_id: ZakuraPeerId,
2425        stream_kind: u16,
2426        request_id: u64,
2427        max_frame_bytes: u32,
2428        max_message_bytes: u32,
2429        frame: Frame,
2430    ) -> BoxRunFuture<'a, Result<Vec<Frame>, SinkReject>> {
2431        Box::pin(async move {
2432            if stream_kind != ZAKURA_STREAM_LEGACY_REQUESTS {
2433                return Err(SinkReject::protocol(
2434                    "unsupported legacy request stream kind",
2435                ));
2436            }
2437
2438            let frame = LegacyRequestFrame::decode_frame(frame).map_err(SinkReject::protocol)?;
2439            let (response_tx, response_rx) = oneshot::channel();
2440            match self
2441                .inbound_tx
2442                .try_send(LegacyInboundWork::Request(LegacyRequestInbound {
2443                    peer_id: peer_id.clone(),
2444                    request_id,
2445                    frame,
2446                    response_tx,
2447                })) {
2448                Ok(()) => {}
2449                Err(mpsc::error::TrySendError::Full(_)) => {
2450                    return Err(SinkReject::local("legacy request inbound queue full"));
2451                }
2452                Err(mpsc::error::TrySendError::Closed(_)) => {
2453                    return Err(SinkReject::local("legacy request inbound queue closed"));
2454                }
2455            }
2456
2457            let response = timeout(LEGACY_REQUEST_TIMEOUT, response_rx)
2458                .await
2459                .map_err(|_| SinkReject::local("legacy request service timed out"))?
2460                .map_err(|_| SinkReject::local("legacy request response dropped"))?
2461                .map_err(SinkReject::local)?;
2462            self.trace.emit_event(|| {
2463                LegacyRequestResponse::new(
2464                    "inbound.response",
2465                    Some(&peer_id),
2466                    request_id,
2467                    response.command(),
2468                    &response,
2469                )
2470            });
2471            LegacyResponseCodec::encode_response(
2472                request_id,
2473                response,
2474                max_frame_bytes,
2475                max_message_bytes,
2476            )
2477            .map_err(SinkReject::local)
2478        })
2479    }
2480}
2481
2482impl ZakuraService for LegacyGossipSink {
2483    fn name(&self) -> &'static str {
2484        "legacy-gossip"
2485    }
2486
2487    fn streams(&self) -> &[Stream] {
2488        legacy_gossip_streams()
2489    }
2490
2491    fn ordered_stream_policy(&self, _kind: u16) -> OrderedStreamPolicy {
2492        OrderedStreamPolicy {
2493            opening: OrderedStreamOpening::InitiatorOnly,
2494            reopen: true,
2495        }
2496    }
2497
2498    fn owns_connection_for_peer(&self, peer: &ZakuraPeerId, conn_id: ZakuraConnId) -> bool {
2499        // A reopen-gap claim counts as ownership without any demand gating: for
2500        // a non-retired peer the demand below is unconditionally `OpenNow` and
2501        // the transport re-admits as soon as the reopen backoff elapses. A
2502        // retired peer holds no session and no claim, so ownership correctly
2503        // reads false.
2504        self.outbound.owns_connection(peer, conn_id)
2505    }
2506
2507    fn ordered_session_demand(
2508        &self,
2509        conn_id: ZakuraConnId,
2510        peer: &ZakuraPeerId,
2511        _negotiated: u64,
2512        _direction: ServicePeerDirection,
2513    ) -> OrderedSessionDemand {
2514        if self.outbound.is_retired(peer, conn_id) {
2515            return OrderedSessionDemand::Retire;
2516        }
2517        OrderedSessionDemand::OpenNow
2518    }
2519
2520    fn add_peer(&self, mut peer: Peer) {
2521        let Some((session_id, mut recv, send)) =
2522            peer.take_stream_with_session_id(ZAKURA_STREAM_GOSSIP)
2523        else {
2524            return;
2525        };
2526        let outbound = self.outbound.clone();
2527        let inbound_tx = self.inbound_tx.clone();
2528        let peer_id = peer.id.clone();
2529        let conn_id = peer.conn_id;
2530        let cancel_token = peer.cancel_token();
2531        let session = LegacyGossipPeerSession::new(peer_id.clone(), conn_id, session_id, send);
2532
2533        if !outbound.insert(session.clone()) {
2534            return;
2535        }
2536        let replay_task_peer_id = peer_id.clone();
2537        let replay_panic_peer_id = replay_task_peer_id.clone();
2538        let replay_panic_outbound = outbound.clone();
2539        let replay_panic_cancel = cancel_token.clone();
2540        spawn_supervised_peer_task(
2541            replay_task_peer_id,
2542            || {},
2543            move || {
2544                replay_panic_cancel.cancel();
2545                replay_panic_outbound.finish_session(
2546                    &replay_panic_peer_id,
2547                    conn_id,
2548                    session_id,
2549                    false,
2550                );
2551            },
2552            {
2553                let outbound = outbound.clone();
2554                let session = session.clone();
2555                async move {
2556                    if let Err(error) = outbound.replay_latest_block_to_peer(session).await {
2557                        debug!(?error, "latest Zakura block gossip replay failed");
2558                    }
2559                }
2560            },
2561        );
2562
2563        let recv_task_peer_id = peer_id.clone();
2564        let recv_teardown_peer_id = recv_task_peer_id.clone();
2565        let recv_teardown_outbound = outbound.clone();
2566        let recv_panic_cancel = cancel_token.clone();
2567        spawn_supervised_peer_task(
2568            recv_task_peer_id,
2569            move || {
2570                recv_teardown_outbound.finish_session(
2571                    &recv_teardown_peer_id,
2572                    conn_id,
2573                    session_id,
2574                    false,
2575                )
2576            },
2577            move || {
2578                recv_panic_cancel.cancel();
2579            },
2580            async move {
2581                let mut received_frame = false;
2582                loop {
2583                    let frame = tokio::select! {
2584                        _ = cancel_token.cancelled() => {
2585                            outbound.finish_session(&peer_id, conn_id, session_id, false);
2586                            return;
2587                        }
2588                        frame = recv.recv() => {
2589                            let Some(frame) = frame else {
2590                                // A stream-only EOF without a single delivered
2591                                // frame counts toward the reopen-churn bound.
2592                                outbound.finish_session(
2593                                    &peer_id,
2594                                    conn_id,
2595                                    session_id,
2596                                    !received_frame,
2597                                );
2598                                return;
2599                            };
2600                            received_frame = true;
2601                            frame
2602                        }
2603                    };
2604
2605                    #[cfg(test)]
2606                    if should_panic_legacy_gossip_recv_loop(&peer_id) {
2607                        panic!("injected legacy gossip recv-loop panic after state registration");
2608                    }
2609
2610                    match Self::enqueue_gossip_frame(&inbound_tx, peer_id.clone(), frame) {
2611                        Ok(()) => {}
2612                        Err(SinkReject::Protocol(error)) => {
2613                            debug!(
2614                                ?error,
2615                                ?peer_id,
2616                                "legacy gossip stream rejected protocol-invalid frame"
2617                            );
2618                            cancel_token.cancel();
2619                            outbound.finish_session(&peer_id, conn_id, session_id, false);
2620                            return;
2621                        }
2622                        Err(SinkReject::Local(error)) => {
2623                            debug!(?error, ?peer_id, "legacy gossip inbound queue closed");
2624                            return;
2625                        }
2626                    }
2627                }
2628            },
2629        );
2630    }
2631
2632    fn remove_peer(&self, peer: &ZakuraPeerId, conn_id: ZakuraConnId) {
2633        self.outbound.remove(peer, conn_id);
2634    }
2635
2636    fn deliver_frame(
2637        &self,
2638        peer_id: ZakuraPeerId,
2639        stream_kind: u16,
2640        frame: Frame,
2641    ) -> Result<(), SinkReject> {
2642        self.deliver(peer_id, stream_kind, frame)
2643    }
2644
2645    fn as_request_response(&self) -> Option<&dyn RequestResponseService> {
2646        Some(self)
2647    }
2648}
2649
2650impl RequestResponseService for LegacyGossipSink {
2651    fn request_frame<'a>(
2652        &'a self,
2653        peer_id: ZakuraPeerId,
2654        stream_kind: u16,
2655        request_id: u64,
2656        max_frame_bytes: u32,
2657        max_message_bytes: u32,
2658        frame: Frame,
2659    ) -> BoxRunFuture<'a, Result<Vec<Frame>, SinkReject>> {
2660        self.request(
2661            peer_id,
2662            stream_kind,
2663            request_id,
2664            max_frame_bytes,
2665            max_message_bytes,
2666            frame,
2667        )
2668    }
2669}
2670
2671#[derive(Debug)]
2672enum LegacyInboundWork {
2673    Gossip(LegacyGossipInbound),
2674    Request(LegacyRequestInbound),
2675}
2676
2677#[derive(Debug)]
2678struct LegacyGossipInbound {
2679    peer_id: ZakuraPeerId,
2680    frame: LegacyGossipFrame,
2681}
2682
2683#[derive(Debug)]
2684struct LegacyRequestInbound {
2685    peer_id: ZakuraPeerId,
2686    request_id: u64,
2687    frame: LegacyRequestFrame,
2688    response_tx: oneshot::Sender<Result<Response, BoxError>>,
2689}
2690
2691async fn legacy_gossip_worker<Inbound>(
2692    mut inbound: Inbound,
2693    mut inbound_rx: mpsc::Receiver<LegacyInboundWork>,
2694    forwarder: LegacyGossipForwarder,
2695    trace: ZakuraTrace,
2696) where
2697    Inbound: Service<Request, Response = Response, Error = BoxError> + Send + Clone + 'static,
2698    Inbound::Future: Send + 'static,
2699{
2700    let request_permits = Arc::new(Semaphore::new(LEGACY_REQUEST_IN_FLIGHT_LIMIT));
2701
2702    while let Some(work) = inbound_rx.recv().await {
2703        match work {
2704            LegacyInboundWork::Gossip(gossip) => {
2705                handle_legacy_gossip(&mut inbound, gossip, &forwarder).await;
2706            }
2707            LegacyInboundWork::Request(request) => {
2708                let Ok(permit) = request_permits.clone().try_acquire_owned() else {
2709                    let _ = request.response_tx.send(Err(
2710                        "legacy request inbound concurrency limit reached".into(),
2711                    ));
2712                    continue;
2713                };
2714                tokio::spawn(handle_legacy_request(
2715                    inbound.clone(),
2716                    request,
2717                    permit,
2718                    trace.clone(),
2719                ));
2720            }
2721        }
2722    }
2723}
2724
2725async fn handle_legacy_gossip<Inbound>(
2726    inbound: &mut Inbound,
2727    gossip: LegacyGossipInbound,
2728    forwarder: &LegacyGossipForwarder,
2729) where
2730    Inbound: Service<Request, Response = Response, Error = BoxError> + Send + 'static,
2731    Inbound::Future: Send + 'static,
2732{
2733    let Some(unseen_frame) = forwarder.unseen(&gossip.frame).await else {
2734        return;
2735    };
2736    // Skip inventory whose recent service attempt was expensive but did not succeed.
2737    // The not-ready, call-error, and call-timeout paths below all return without
2738    // `mark_seen`, so without this an authenticated peer could replay the same valid
2739    // advertisement while the inbound service is slow/erroring and make the serial
2740    // worker re-pay the full 30s+30s readiness/call budget for every queued
2741    // duplicate. Fast failures are not recorded, so a transient blip stays retryable.
2742    let Some(unseen_frame) = forwarder.fresh_attempt(&unseen_frame).await else {
2743        debug!("legacy gossip duplicate suppressed after recent expensive attempt");
2744        return;
2745    };
2746    let request = unseen_frame.clone().into_request(gossip.peer_id.clone());
2747
2748    let started = Instant::now();
2749    let ready = timeout(LEGACY_GOSSIP_SERVICE_TIMEOUT, inbound.ready()).await;
2750    let Ok(Ok(service)) = ready else {
2751        debug!("legacy gossip inbound service was not ready");
2752        forwarder
2753            .note_failed_attempt(&unseen_frame, started.elapsed())
2754            .await;
2755        return;
2756    };
2757
2758    match timeout(LEGACY_GOSSIP_SERVICE_TIMEOUT, service.call(request)).await {
2759        Ok(Ok(_)) => {}
2760        Ok(Err(error)) => {
2761            debug!(?error, "legacy gossip inbound service call failed");
2762            forwarder
2763                .note_failed_attempt(&unseen_frame, started.elapsed())
2764                .await;
2765            return;
2766        }
2767        Err(_) => {
2768            debug!("legacy gossip inbound service call timed out");
2769            forwarder
2770                .note_failed_attempt(&unseen_frame, started.elapsed())
2771                .await;
2772            return;
2773        }
2774    }
2775
2776    forwarder.mark_seen(&unseen_frame).await;
2777    if let Err(error) = forwarder.forward(&gossip.peer_id, unseen_frame).await {
2778        debug!(?error, "legacy gossip forward failed");
2779    }
2780}
2781
2782async fn handle_legacy_request<Inbound>(
2783    mut inbound: Inbound,
2784    request: LegacyRequestInbound,
2785    _permit: OwnedSemaphorePermit,
2786    trace: ZakuraTrace,
2787) where
2788    Inbound: Service<Request, Response = Response, Error = BoxError> + Send + 'static,
2789    Inbound::Future: Send + 'static,
2790{
2791    let command = request.frame.to_string();
2792    let peer_id = request.peer_id.clone();
2793    let request_id = request.request_id;
2794    let request_kind = request.frame.kind();
2795    let mut response_tx = request.response_tx;
2796    let Some(legacy_request) = request.frame.into_service_request(peer_id.clone()) else {
2797        // Inbound legacy Ping is handled locally; the requester measures round-trip time.
2798        let _ = response_tx.send(Ok(Response::Pong(Duration::ZERO)));
2799        return;
2800    };
2801    trace.emit_event(|| {
2802        LegacyRequestStart::new(
2803            "inbound.request",
2804            Some(&peer_id),
2805            request_id,
2806            request_kind,
2807            request_kind.message_type(),
2808        )
2809    });
2810
2811    let service_work = async move {
2812        let ready = timeout(LEGACY_REQUEST_TIMEOUT, inbound.ready()).await;
2813        match ready {
2814            Ok(Ok(service)) => timeout(LEGACY_REQUEST_TIMEOUT, service.call(legacy_request))
2815                .await
2816                .map_err(|_| -> BoxError { format!("{command} inbound service timed out").into() })
2817                .and_then(|response| response),
2818            Ok(Err(error)) => Err(error),
2819            Err(_) => Err(format!("{command} inbound service readiness timed out").into()),
2820        }
2821    };
2822
2823    // Tie the spawned handler's lifetime to the request stream. The request-stream
2824    // side (`LegacyGossipSink::request`) waits only LEGACY_REQUEST_TIMEOUT for the
2825    // oneshot result and then drops the receiver; the peer/connection going away
2826    // drops it too. In either case `closed()` resolves, so abort the in-flight
2827    // service work and release the permit instead of letting this orphaned handler
2828    // keep one of LEGACY_REQUEST_IN_FLIGHT_LIMIT permits (and keep doing backend
2829    // work) for up to another full readiness + service-call timeout.
2830    let result = tokio::select! {
2831        biased;
2832        () = response_tx.closed() => {
2833            debug!(
2834                ?peer_id,
2835                request_id,
2836                "legacy request handler aborted: response receiver dropped before service completed"
2837            );
2838            return;
2839        }
2840        result = service_work => result,
2841    };
2842
2843    if let Err(error) = &result {
2844        trace.emit_event(|| {
2845            LegacyRequestError::new(
2846                "inbound.error",
2847                Some(&peer_id),
2848                request_id,
2849                request_kind.command(),
2850                error.to_string(),
2851            )
2852        });
2853    }
2854
2855    if response_tx.send(result).is_err() {
2856        debug!(
2857            ?peer_id,
2858            request_id, "legacy request response receiver dropped before service completed"
2859        );
2860    }
2861}
2862
2863#[derive(Clone, Debug)]
2864struct FirstSeenCache {
2865    inner: std::sync::Arc<Mutex<FirstSeenCacheInner>>,
2866}
2867
2868impl FirstSeenCache {
2869    fn new(capacity: usize, ttl: Duration) -> Self {
2870        Self {
2871            inner: std::sync::Arc::new(Mutex::new(FirstSeenCacheInner {
2872                capacity: capacity.max(1),
2873                ttl,
2874                entries: HashMap::new(),
2875                order: VecDeque::new(),
2876            })),
2877        }
2878    }
2879
2880    async fn record_unseen(
2881        &self,
2882        keys: impl IntoIterator<Item = InventoryKey>,
2883    ) -> Vec<InventoryKey> {
2884        let mut inner = self.inner.lock().await;
2885        let now = Instant::now();
2886        inner.prune(now);
2887        let mut recorded = Vec::new();
2888        for key in keys {
2889            if inner.insert(key, now) {
2890                recorded.push(key);
2891            }
2892        }
2893        recorded
2894    }
2895
2896    async fn record_first_seen(&self, frame: &LegacyGossipFrame) -> Option<LegacyGossipFrame> {
2897        let unseen = self.unseen(frame).await?;
2898        self.record_seen(&unseen).await;
2899        Some(unseen)
2900    }
2901
2902    async fn unseen(&self, frame: &LegacyGossipFrame) -> Option<LegacyGossipFrame> {
2903        match frame {
2904            LegacyGossipFrame::AdvertiseBlock(hash) => {
2905                (!self.contains(InventoryKey::Block(*hash)).await)
2906                    .then_some(LegacyGossipFrame::AdvertiseBlock(*hash))
2907            }
2908            LegacyGossipFrame::AdvertiseTransactionIds(ids) => {
2909                let unseen = self.unrecorded_transaction_ids(ids.iter().copied()).await;
2910                (!unseen.is_empty()).then_some(LegacyGossipFrame::AdvertiseTransactionIds(unseen))
2911            }
2912        }
2913    }
2914
2915    async fn record_seen(&self, frame: &LegacyGossipFrame) {
2916        match frame {
2917            LegacyGossipFrame::AdvertiseBlock(hash) => {
2918                self.record_unseen([InventoryKey::Block(*hash)]).await;
2919            }
2920            LegacyGossipFrame::AdvertiseTransactionIds(ids) => {
2921                self.record_unseen(ids.iter().copied().map(InventoryKey::Transaction))
2922                    .await;
2923            }
2924        }
2925    }
2926
2927    async fn contains(&self, key: InventoryKey) -> bool {
2928        let mut inner = self.inner.lock().await;
2929        let now = Instant::now();
2930        inner.prune(now);
2931        inner.entries.contains_key(&key)
2932    }
2933
2934    async fn unrecorded_transaction_ids(
2935        &self,
2936        ids: impl IntoIterator<Item = UnminedTxId>,
2937    ) -> Vec<UnminedTxId> {
2938        let mut inner = self.inner.lock().await;
2939        inner.prune(Instant::now());
2940        let mut unseen = Vec::new();
2941        for tx_id in ids {
2942            if !inner
2943                .entries
2944                .contains_key(&InventoryKey::Transaction(tx_id))
2945            {
2946                unseen.push(tx_id);
2947            }
2948        }
2949        unseen
2950    }
2951}
2952
2953#[derive(Debug)]
2954struct FirstSeenCacheInner {
2955    capacity: usize,
2956    ttl: Duration,
2957    entries: HashMap<InventoryKey, Instant>,
2958    order: VecDeque<InventoryKey>,
2959}
2960
2961impl FirstSeenCacheInner {
2962    fn insert(&mut self, key: InventoryKey, now: Instant) -> bool {
2963        if self.entries.contains_key(&key) {
2964            return false;
2965        }
2966        self.entries.insert(key, now);
2967        self.order.push_back(key);
2968        while self.entries.len() > self.capacity {
2969            if let Some(oldest) = self.order.pop_front() {
2970                self.entries.remove(&oldest);
2971            }
2972        }
2973        true
2974    }
2975
2976    fn prune(&mut self, now: Instant) {
2977        while let Some(key) = self.order.front().copied() {
2978            let Some(seen_at) = self.entries.get(&key).copied() else {
2979                self.order.pop_front();
2980                continue;
2981            };
2982            if now.saturating_duration_since(seen_at) < self.ttl {
2983                break;
2984            }
2985            self.order.pop_front();
2986            self.entries.remove(&key);
2987        }
2988    }
2989}
2990
2991#[derive(Copy, Clone, Debug, Eq, Hash, PartialEq)]
2992enum InventoryKey {
2993    Block(block::Hash),
2994    Transaction(UnminedTxId),
2995}
2996
2997/// Errors produced by the legacy gossip adapter.
2998#[derive(Debug, Error)]
2999pub enum LegacyGossipError {
3000    /// The request is intentionally outside this phase's gossip-only scope.
3001    #[error("unsupported legacy gossip request: {0}")]
3002    UnsupportedRequest(&'static str),
3003    /// A peer used nonzero flags before any flags are defined.
3004    #[error("unsupported legacy gossip flags: {0}")]
3005    UnsupportedFlags(u16),
3006    /// A peer sent an unknown gossip message type.
3007    #[error("unknown legacy gossip message type: {0}")]
3008    UnknownMessageType(u16),
3009    /// Transaction advertisements must not be empty.
3010    #[error("empty transaction-id advertisement")]
3011    EmptyTransactionAdvertisement,
3012    /// A peer requested or responded with too many inventory items.
3013    #[error("too many inventory items in legacy request/response: {0}")]
3014    TooManyInventoryItems(usize),
3015    /// A peer requested too many block locator hashes.
3016    #[error("too many block locator hashes in legacy request: {0}")]
3017    TooManyBlockLocatorHashes(usize),
3018    /// A peer responded with too many block headers.
3019    #[error("too many block headers in legacy response: {0}")]
3020    TooManyHeaders(usize),
3021    /// Transaction gossip payload included non-transaction inventory.
3022    #[error("transaction-id advertisement contained non-transaction inventory")]
3023    NonTransactionInventory,
3024    /// Payload had bytes after the decoded message.
3025    #[error("trailing bytes after legacy gossip payload")]
3026    TrailingBytes,
3027    /// A response frame echoed a different request id.
3028    #[error("wrong legacy request id in response: expected {expected}, got {actual}")]
3029    WrongRequestId {
3030        /// Expected request id.
3031        expected: u64,
3032        /// Actual request id in the response frame.
3033        actual: u64,
3034    },
3035    /// A chunked response ended before its final chunk.
3036    #[error("incomplete legacy response chunk")]
3037    IncompleteResponseChunk,
3038    /// A response payload was too short to contain its fixed header.
3039    #[error("truncated legacy response")]
3040    TruncatedResponse,
3041    /// A response exceeded the protocol message size.
3042    #[error("oversized legacy response: {0} bytes")]
3043    OversizedResponse(usize),
3044    /// The encoded response would exceed the responder-side aggregate budget.
3045    #[error("legacy response exceeded responder aggregate budget: {0} bytes")]
3046    ResponseAggregateBudget(usize),
3047    /// The legacy service returned an unexpected response variant.
3048    #[error("unexpected legacy response: {0}")]
3049    UnexpectedResponse(&'static str),
3050    /// A request that requires an acknowledgement received none.
3051    #[error("missing legacy response: {0}")]
3052    MissingResponse(&'static str),
3053    /// The peer returned a block we did not request, so the response is not bound
3054    /// to the requested hash. Treated as a peer fault.
3055    #[error("legacy block response contained an unrequested block: {0:?}")]
3056    UnsolicitedBlock(block::Hash),
3057    /// Zcash serialization failed.
3058    #[error(transparent)]
3059    Serialization(#[from] SerializationError),
3060    /// Local buffer serialization failed.
3061    #[error(transparent)]
3062    Io(#[from] std::io::Error),
3063    /// Integer conversion failed while checking a peer-controlled bound.
3064    #[error(transparent)]
3065    Integer(#[from] std::num::TryFromIntError),
3066}
3067
3068impl fmt::Display for LegacyGossipFrame {
3069    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
3070        match self {
3071            Self::AdvertiseBlock(_) => f.write_str("AdvertiseBlock"),
3072            Self::AdvertiseTransactionIds(ids) => {
3073                write!(f, "AdvertiseTransactionIds({})", ids.len())
3074            }
3075        }
3076    }
3077}
3078
3079impl fmt::Display for LegacyRequestFrame {
3080    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
3081        match self {
3082            Self::BlocksByHash(hashes) => write!(f, "BlocksByHash({})", hashes.len()),
3083            Self::TransactionsById(ids) => write!(f, "TransactionsById({})", ids.len()),
3084            Self::FindBlocks { known_blocks, stop } => write!(
3085                f,
3086                "FindBlocks {{ known_blocks: {}, stop: {} }}",
3087                known_blocks.len(),
3088                stop.is_some()
3089            ),
3090            Self::FindHeaders { known_blocks, stop } => write!(
3091                f,
3092                "FindHeaders {{ known_blocks: {}, stop: {} }}",
3093                known_blocks.len(),
3094                stop.is_some()
3095            ),
3096            Self::MempoolTransactionIds => f.write_str("MempoolTransactionIds"),
3097            Self::Ping => f.write_str("Ping"),
3098            Self::PushTransaction(_) => f.write_str("PushTransaction"),
3099        }
3100    }
3101}
3102
3103#[cfg(test)]
3104mod tests {
3105    use super::*;
3106    use std::{
3107        collections::HashSet,
3108        sync::{
3109            atomic::{AtomicUsize, Ordering},
3110            Arc,
3111        },
3112        task::{Context, Poll},
3113    };
3114
3115    use futures::FutureExt;
3116    use tokio::sync::mpsc::UnboundedReceiver;
3117    use tower::ServiceExt;
3118    use zakura_chain::{
3119        parameters::NetworkUpgrade,
3120        serialization::{ZcashDeserialize, ZcashSerialize},
3121        transaction::{self, LockTime, Transaction, WtxId},
3122    };
3123    use zakura_test::vectors::BLOCK_TESTNET_141042_BYTES;
3124
3125    use crate::zakura::{
3126        framed_channel,
3127        testkit::{HostilePeer, ZakuraTestNode, TEST_NET_TIMEOUT},
3128        CloseCause, Peer, ServicePeerDirection, ServiceStream, ZAKURA_CAP_LEGACY_GOSSIP,
3129    };
3130
3131    fn block_hash(byte: u8) -> block::Hash {
3132        block::Hash([byte; 32])
3133    }
3134
3135    fn legacy_tx_id(byte: u8) -> UnminedTxId {
3136        UnminedTxId::from_legacy_id(transaction::Hash([byte; 32]))
3137    }
3138
3139    fn witnessed_tx_id(byte: u8) -> UnminedTxId {
3140        UnminedTxId::from(WtxId {
3141            id: transaction::Hash([byte; 32]),
3142            auth_digest: transaction::AuthDigest::from([byte.wrapping_add(1); 32]),
3143        })
3144    }
3145
3146    fn empty_v5_transaction(byte: u8) -> Transaction {
3147        Transaction::V5 {
3148            network_upgrade: NetworkUpgrade::Nu5,
3149            lock_time: LockTime::min_lock_time_timestamp(),
3150            expiry_height: block::Height(u32::from(byte)),
3151            inputs: Vec::new(),
3152            outputs: Vec::new(),
3153            sapling_shielded_data: None,
3154            orchard_shielded_data: None,
3155        }
3156    }
3157
3158    fn encoded_count(count: usize) -> Vec<u8> {
3159        CompactSizeMessage::try_from(count)
3160            .expect("test count is within CompactSizeMessage bounds")
3161            .zcash_serialize_to_vec()
3162            .expect("compact size serializes")
3163    }
3164
3165    #[derive(Clone, Debug)]
3166    struct RequestRecorder {
3167        tx: tokio::sync::mpsc::UnboundedSender<Request>,
3168    }
3169
3170    impl Service<Request> for RequestRecorder {
3171        type Response = Response;
3172        type Error = BoxError;
3173        type Future = std::future::Ready<Result<Response, BoxError>>;
3174
3175        fn poll_ready(&mut self, _cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
3176            Poll::Ready(Ok(()))
3177        }
3178
3179        fn call(&mut self, request: Request) -> Self::Future {
3180            match self.tx.send(request) {
3181                Ok(()) => std::future::ready(Ok(Response::Nil)),
3182                Err(_) => std::future::ready(Err("request recorder receiver dropped".into())),
3183            }
3184        }
3185    }
3186
3187    #[derive(Clone, Debug)]
3188    struct FailsOnceRecorder {
3189        attempts: Arc<AtomicUsize>,
3190        tx: tokio::sync::mpsc::UnboundedSender<Request>,
3191    }
3192
3193    impl Service<Request> for FailsOnceRecorder {
3194        type Response = Response;
3195        type Error = BoxError;
3196        type Future = std::future::Ready<Result<Response, BoxError>>;
3197
3198        fn poll_ready(&mut self, _cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
3199            Poll::Ready(Ok(()))
3200        }
3201
3202        fn call(&mut self, request: Request) -> Self::Future {
3203            if self.attempts.fetch_add(1, Ordering::SeqCst) == 0 {
3204                return std::future::ready(Err("intentional transient inbound failure".into()));
3205            }
3206
3207            match self.tx.send(request) {
3208                Ok(()) => std::future::ready(Ok(Response::Nil)),
3209                Err(_) => std::future::ready(Err("request recorder receiver dropped".into())),
3210            }
3211        }
3212    }
3213
3214    #[derive(Clone, Debug)]
3215    struct InventoryResponder {
3216        transaction: UnminedTx,
3217    }
3218
3219    impl Service<Request> for InventoryResponder {
3220        type Response = Response;
3221        type Error = BoxError;
3222        type Future = std::future::Ready<Result<Response, BoxError>>;
3223
3224        fn poll_ready(&mut self, _cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
3225            Poll::Ready(Ok(()))
3226        }
3227
3228        fn call(&mut self, request: Request) -> Self::Future {
3229            let response = match request {
3230                Request::BlocksByHash(hashes) | Request::BlocksByHashFrom { hashes, .. } => {
3231                    Response::Blocks(hashes.into_iter().map(InventoryResponse::Missing).collect())
3232                }
3233                Request::TransactionsById(ids) | Request::TransactionsByIdFrom { ids, .. } => {
3234                    Response::Transactions(
3235                        ids.into_iter()
3236                            .map(|id| {
3237                                if id == self.transaction.id() {
3238                                    InventoryResponse::Available((self.transaction.clone(), None))
3239                                } else {
3240                                    InventoryResponse::Missing(id)
3241                                }
3242                            })
3243                            .collect(),
3244                    )
3245                }
3246                request => {
3247                    return std::future::ready(Err(format!(
3248                        "unexpected inventory request: {request:?}"
3249                    )
3250                    .into()));
3251                }
3252            };
3253            std::future::ready(Ok(response))
3254        }
3255    }
3256
3257    #[derive(Clone, Debug)]
3258    struct BlockInventoryResponder {
3259        block: Arc<Block>,
3260    }
3261
3262    impl Service<Request> for BlockInventoryResponder {
3263        type Response = Response;
3264        type Error = BoxError;
3265        type Future = std::future::Ready<Result<Response, BoxError>>;
3266
3267        fn poll_ready(&mut self, _cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
3268            Poll::Ready(Ok(()))
3269        }
3270
3271        fn call(&mut self, request: Request) -> Self::Future {
3272            let response = match request {
3273                Request::BlocksByHash(hashes) | Request::BlocksByHashFrom { hashes, .. } => {
3274                    Response::Blocks(
3275                        hashes
3276                            .into_iter()
3277                            .map(|hash| {
3278                                if hash == self.block.hash() {
3279                                    InventoryResponse::Available((self.block.clone(), None))
3280                                } else {
3281                                    InventoryResponse::Missing(hash)
3282                                }
3283                            })
3284                            .collect(),
3285                    )
3286                }
3287                request => {
3288                    return std::future::ready(Err(format!(
3289                        "unexpected block inventory request: {request:?}"
3290                    )
3291                    .into()));
3292                }
3293            };
3294            std::future::ready(Ok(response))
3295        }
3296    }
3297
3298    #[derive(Clone, Debug)]
3299    struct SlowRequestThenRecorder {
3300        release: tokio::sync::watch::Receiver<bool>,
3301        tx: tokio::sync::mpsc::UnboundedSender<Request>,
3302    }
3303
3304    impl Service<Request> for SlowRequestThenRecorder {
3305        type Response = Response;
3306        type Error = BoxError;
3307        type Future = Pin<Box<dyn Future<Output = Result<Response, BoxError>> + Send>>;
3308
3309        fn poll_ready(&mut self, _cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
3310            Poll::Ready(Ok(()))
3311        }
3312
3313        fn call(&mut self, request: Request) -> Self::Future {
3314            match request {
3315                Request::BlocksByHash(hashes) | Request::BlocksByHashFrom { hashes, .. } => {
3316                    let mut release = self.release.clone();
3317                    async move {
3318                        while !*release.borrow() {
3319                            if release.changed().await.is_err() {
3320                                break;
3321                            }
3322                        }
3323                        Ok(Response::Blocks(
3324                            hashes.into_iter().map(InventoryResponse::Missing).collect(),
3325                        ))
3326                    }
3327                    .boxed()
3328                }
3329                request => {
3330                    let tx = self.tx.clone();
3331                    async move {
3332                        tx.send(request)?;
3333                        Ok(Response::Nil)
3334                    }
3335                    .boxed()
3336                }
3337            }
3338        }
3339    }
3340
3341    /// Always ready, but every `call` future is pending forever. Models a slow or
3342    /// backpressured inbound service whose work outlives the request-stream timeout.
3343    #[derive(Clone, Debug)]
3344    struct NeverCompletesService;
3345
3346    impl Service<Request> for NeverCompletesService {
3347        type Response = Response;
3348        type Error = BoxError;
3349        type Future = Pin<Box<dyn Future<Output = Result<Response, BoxError>> + Send>>;
3350
3351        fn poll_ready(&mut self, _cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
3352            Poll::Ready(Ok(()))
3353        }
3354
3355        fn call(&mut self, _request: Request) -> Self::Future {
3356            std::future::pending().boxed()
3357        }
3358    }
3359
3360    /// Always ready, but every `call` future is pending forever and counts
3361    /// invocations. Models a slow/backpressured inbound service whose calls hit
3362    /// `LEGACY_GOSSIP_SERVICE_TIMEOUT`, so `handle_legacy_gossip` returns without
3363    /// `mark_seen`, exposing how many times duplicate gossip frames reach the
3364    /// expensive readiness/call path.
3365    #[derive(Clone, Debug)]
3366    struct CountingPendingService {
3367        calls: Arc<AtomicUsize>,
3368    }
3369
3370    impl Service<Request> for CountingPendingService {
3371        type Response = Response;
3372        type Error = BoxError;
3373        type Future = Pin<Box<dyn Future<Output = Result<Response, BoxError>> + Send>>;
3374
3375        fn poll_ready(&mut self, _cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
3376            Poll::Ready(Ok(()))
3377        }
3378
3379        fn call(&mut self, _request: Request) -> Self::Future {
3380            self.calls.fetch_add(1, Ordering::SeqCst);
3381            std::future::pending().boxed()
3382        }
3383    }
3384
3385    #[derive(Clone, Debug)]
3386    struct RecordingInventoryResponder {
3387        transaction: Option<UnminedTx>,
3388        tx: tokio::sync::mpsc::UnboundedSender<Request>,
3389    }
3390
3391    impl Service<Request> for RecordingInventoryResponder {
3392        type Response = Response;
3393        type Error = BoxError;
3394        type Future = std::future::Ready<Result<Response, BoxError>>;
3395
3396        fn poll_ready(&mut self, _cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
3397            Poll::Ready(Ok(()))
3398        }
3399
3400        fn call(&mut self, request: Request) -> Self::Future {
3401            if let Err(error) = self.tx.send(request.clone()) {
3402                return std::future::ready(Err(Box::new(error)));
3403            }
3404
3405            let response = match request {
3406                Request::BlocksByHash(hashes) | Request::BlocksByHashFrom { hashes, .. } => {
3407                    Response::Blocks(hashes.into_iter().map(InventoryResponse::Missing).collect())
3408                }
3409                Request::TransactionsById(ids) | Request::TransactionsByIdFrom { ids, .. } => {
3410                    Response::Transactions(
3411                        ids.into_iter()
3412                            .map(|id| match &self.transaction {
3413                                Some(transaction) if id == transaction.id() => {
3414                                    InventoryResponse::Available((transaction.clone(), None))
3415                                }
3416                                _ => InventoryResponse::Missing(id),
3417                            })
3418                            .collect(),
3419                    )
3420                }
3421                request => {
3422                    return std::future::ready(Err(format!(
3423                        "unexpected inventory request: {request:?}"
3424                    )
3425                    .into()));
3426                }
3427            };
3428            std::future::ready(Ok(response))
3429        }
3430    }
3431
3432    #[derive(Clone, Debug)]
3433    struct NormalNetworkResponder {
3434        block: Arc<Block>,
3435        transaction: UnminedTx,
3436        pushed_tx: tokio::sync::mpsc::UnboundedSender<UnminedTxId>,
3437    }
3438
3439    impl NormalNetworkResponder {
3440        fn header(&self) -> block::CountedHeader {
3441            block::CountedHeader {
3442                header: self.block.header.clone(),
3443            }
3444        }
3445    }
3446
3447    impl Service<Request> for NormalNetworkResponder {
3448        type Response = Response;
3449        type Error = BoxError;
3450        type Future = std::future::Ready<Result<Response, BoxError>>;
3451
3452        fn poll_ready(&mut self, _cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
3453            Poll::Ready(Ok(()))
3454        }
3455
3456        fn call(&mut self, request: Request) -> Self::Future {
3457            let response = match request {
3458                Request::FindBlocks { .. } => Response::BlockHashes(vec![self.block.hash()]),
3459                Request::FindHeaders { .. } => Response::BlockHeaders(vec![self.header()]),
3460                Request::MempoolTransactionIds => {
3461                    Response::TransactionIds(vec![self.transaction.id()])
3462                }
3463                Request::BlocksByHash(hashes) | Request::BlocksByHashFrom { hashes, .. } => {
3464                    Response::Blocks(
3465                        hashes
3466                            .into_iter()
3467                            .map(|hash| {
3468                                if hash == self.block.hash() {
3469                                    InventoryResponse::Available((self.block.clone(), None))
3470                                } else {
3471                                    InventoryResponse::Missing(hash)
3472                                }
3473                            })
3474                            .collect(),
3475                    )
3476                }
3477                Request::PushTransaction(transaction, _) => {
3478                    if let Err(error) = self.pushed_tx.send(transaction.id()) {
3479                        return std::future::ready(Err(Box::new(error)));
3480                    }
3481                    Response::Nil
3482                }
3483                request => {
3484                    return std::future::ready(Err(format!(
3485                        "unexpected normal-network request: {request:?}"
3486                    )
3487                    .into()));
3488                }
3489            };
3490            std::future::ready(Ok(response))
3491        }
3492    }
3493
3494    async fn legacy_node(
3495        seed: u64,
3496    ) -> Result<(ZakuraTestNode, UnboundedReceiver<Request>), BoxError> {
3497        let (tx, rx) = tokio::sync::mpsc::unbounded_channel();
3498        let node = ZakuraTestNode::builder(seed)
3499            .service_from_supervisor(move |supervisor| {
3500                Arc::new(LegacyGossipSink::spawn(RequestRecorder { tx }, supervisor))
3501            })
3502            .spawn()
3503            .await?;
3504        Ok((node, rx))
3505    }
3506
3507    async fn inventory_node(seed: u64, transaction: UnminedTx) -> Result<ZakuraTestNode, BoxError> {
3508        let node = ZakuraTestNode::builder(seed)
3509            .service_from_supervisor(move |supervisor| {
3510                Arc::new(LegacyGossipSink::spawn(
3511                    InventoryResponder { transaction },
3512                    supervisor,
3513                ))
3514            })
3515            .spawn()
3516            .await?;
3517        Ok(node)
3518    }
3519
3520    async fn block_inventory_node(
3521        seed: u64,
3522        block: Arc<Block>,
3523    ) -> Result<ZakuraTestNode, BoxError> {
3524        let node = ZakuraTestNode::builder(seed)
3525            .service_from_supervisor(move |supervisor| {
3526                Arc::new(LegacyGossipSink::spawn(
3527                    BlockInventoryResponder { block },
3528                    supervisor,
3529                ))
3530            })
3531            .spawn()
3532            .await?;
3533        Ok(node)
3534    }
3535
3536    async fn recording_inventory_node(
3537        seed: u64,
3538        transaction: Option<UnminedTx>,
3539    ) -> Result<(ZakuraTestNode, UnboundedReceiver<Request>), BoxError> {
3540        let (tx, rx) = tokio::sync::mpsc::unbounded_channel();
3541        let node = ZakuraTestNode::builder(seed)
3542            .service_from_supervisor(move |supervisor| {
3543                Arc::new(LegacyGossipSink::spawn(
3544                    RecordingInventoryResponder { transaction, tx },
3545                    supervisor,
3546                ))
3547            })
3548            .spawn()
3549            .await?;
3550        Ok((node, rx))
3551    }
3552
3553    async fn normal_network_node(
3554        seed: u64,
3555        block: Arc<Block>,
3556        transaction: UnminedTx,
3557    ) -> Result<
3558        (
3559            ZakuraTestNode,
3560            tokio::sync::mpsc::UnboundedReceiver<UnminedTxId>,
3561        ),
3562        BoxError,
3563    > {
3564        let (pushed_tx, pushed_rx) = tokio::sync::mpsc::unbounded_channel();
3565        let node = ZakuraTestNode::builder(seed)
3566            // Multi-node gossip topologies dial several loopback peers from one
3567            // node, so raise the production Zakura per-IP cap that the default
3568            // test node enforces; these tests exercise gossip routing, not the
3569            // per-IP admission gate.
3570            .max_connections_per_ip(8)
3571            .service_from_supervisor(move |supervisor| {
3572                Arc::new(LegacyGossipSink::spawn(
3573                    NormalNetworkResponder {
3574                        block,
3575                        transaction,
3576                        pushed_tx,
3577                    },
3578                    supervisor,
3579                ))
3580            })
3581            .spawn()
3582            .await?;
3583        Ok((node, pushed_rx))
3584    }
3585
3586    async fn node_peer_id(node: &ZakuraTestNode) -> Result<ZakuraPeerId, BoxError> {
3587        Ok(ZakuraPeerId::new(
3588            node.node_addr().await.node_id.as_bytes().to_vec(),
3589        )?)
3590    }
3591
3592    async fn recv_request(rx: &mut UnboundedReceiver<Request>) -> Result<Request, BoxError> {
3593        tokio::time::timeout(TEST_NET_TIMEOUT, rx.recv())
3594            .await
3595            .map_err(|_| -> BoxError { "timed out waiting for request".into() })?
3596            .ok_or_else(|| "request recorder closed".into())
3597    }
3598
3599    async fn recv_pushed_tx_id(
3600        rx: &mut tokio::sync::mpsc::UnboundedReceiver<UnminedTxId>,
3601    ) -> Result<UnminedTxId, BoxError> {
3602        tokio::time::timeout(TEST_NET_TIMEOUT, rx.recv())
3603            .await
3604            .map_err(|_| -> BoxError { "timed out waiting for pushed transaction".into() })?
3605            .ok_or_else(|| "pushed transaction recorder closed".into())
3606    }
3607
3608    fn legacy_gossip_peer(
3609        peer_id: ZakuraPeerId,
3610        cancel_token: tokio_util::sync::CancellationToken,
3611    ) -> (Peer, FramedSend) {
3612        legacy_gossip_peer_with_conn(peer_id, 0, cancel_token)
3613    }
3614
3615    fn legacy_gossip_peer_with_conn(
3616        peer_id: ZakuraPeerId,
3617        conn_id: ZakuraConnId,
3618        cancel_token: tokio_util::sync::CancellationToken,
3619    ) -> (Peer, FramedSend) {
3620        let (peer_send, service_recv) = framed_channel(8);
3621        let (service_send, _peer_recv) = framed_channel(8);
3622        let peer = Peer::new_with_conn_id_and_direction(
3623            conn_id,
3624            peer_id,
3625            None,
3626            ZAKURA_CAP_LEGACY_GOSSIP,
3627            ServicePeerDirection::Inbound,
3628            HashMap::from([(ZAKURA_STREAM_GOSSIP, (service_recv, service_send))]),
3629            cancel_token,
3630        );
3631        (peer, peer_send)
3632    }
3633
3634    fn legacy_gossip_peer_with_conn_and_session(
3635        peer_id: ZakuraPeerId,
3636        conn_id: ZakuraConnId,
3637        session_id: u64,
3638        cancel_token: tokio_util::sync::CancellationToken,
3639    ) -> (Peer, FramedSend) {
3640        let (peer_send, service_recv) = framed_channel(8);
3641        let (service_send, _peer_recv) = framed_channel(8);
3642        let stream = ServiceStream::new(
3643            session_id,
3644            LEGACY_GOSSIP_VERSION,
3645            service_recv,
3646            service_send,
3647            cancel_token.child_token(),
3648        );
3649        let peer = Peer::new_with_service_streams(
3650            conn_id,
3651            peer_id,
3652            None,
3653            ZAKURA_CAP_LEGACY_GOSSIP,
3654            ServicePeerDirection::Inbound,
3655            HashMap::from([(ZAKURA_STREAM_GOSSIP, stream)]),
3656            cancel_token,
3657            CloseCause::new(),
3658        );
3659        (peer, peer_send)
3660    }
3661
3662    async fn wait_for_legacy_gossip_panic_cleanup(
3663        outbound: &LegacyGossipOutbound,
3664        peer_id: &ZakuraPeerId,
3665        cancel_token: &tokio_util::sync::CancellationToken,
3666    ) -> Result<(), BoxError> {
3667        tokio::time::timeout(Duration::from_secs(1), async {
3668            loop {
3669                if cancel_token.is_cancelled() && !outbound.contains(peer_id) {
3670                    return;
3671                }
3672                tokio::time::sleep(Duration::from_millis(10)).await;
3673            }
3674        })
3675        .await
3676        .map_err(|_| -> BoxError { "timed out waiting for legacy gossip panic cleanup".into() })
3677    }
3678
3679    async fn wait_registered_count(node: &ZakuraTestNode, count: usize) -> Result<(), BoxError> {
3680        let result = tokio::time::timeout(TEST_NET_TIMEOUT, async {
3681            loop {
3682                let observed = node.supervisor().registered_ids().await.len();
3683                if observed == count {
3684                    return;
3685                }
3686                tokio::time::sleep(Duration::from_millis(10)).await;
3687            }
3688        })
3689        .await;
3690
3691        if result.is_err() {
3692            let observed = node.supervisor().registered_ids().await.len();
3693            return Err(format!(
3694                "timed out waiting for {count} peer registrations; observed {observed}"
3695            )
3696            .into());
3697        }
3698
3699        Ok(())
3700    }
3701
3702    /// Regression for `claude-legacy-request-orphaned-handler-permits`.
3703    ///
3704    /// The request-stream side (`LegacyGossipSink::request`) waits only
3705    /// `LEGACY_REQUEST_TIMEOUT` for the oneshot result and then drops the receiver;
3706    /// the peer disconnecting drops it too. The spawned `handle_legacy_request` task
3707    /// must not keep its in-flight permit (one of `LEGACY_REQUEST_IN_FLIGHT_LIMIT`)
3708    /// or its backend service work alive after that. Before the fix it stayed blocked
3709    /// in `service.call` for up to another full `LEGACY_REQUEST_TIMEOUT`, so an
3710    /// attacker driving slow inbound work could occupy all 64 permits.
3711    #[tokio::test(start_paused = true)]
3712    async fn legacy_request_handler_releases_permit_when_request_stream_drops_receiver() {
3713        let permits = Arc::new(Semaphore::new(LEGACY_REQUEST_IN_FLIGHT_LIMIT));
3714        let permit = permits
3715            .clone()
3716            .try_acquire_owned()
3717            .expect("a permit is available");
3718        assert_eq!(
3719            permits.available_permits(),
3720            LEGACY_REQUEST_IN_FLIGHT_LIMIT - 1,
3721            "one permit is held while the handler runs"
3722        );
3723
3724        let (response_tx, response_rx) = oneshot::channel();
3725        let request = LegacyRequestInbound {
3726            peer_id: ZakuraPeerId::new(vec![7u8; 32]).expect("valid peer id"),
3727            request_id: 1,
3728            // A non-Ping request so the handler drives the inbound service.
3729            frame: LegacyRequestFrame::MempoolTransactionIds,
3730            response_tx,
3731        };
3732
3733        let handler = tokio::spawn(handle_legacy_request(
3734            NeverCompletesService,
3735            request,
3736            permit,
3737            ZakuraTrace::noop(),
3738        ));
3739
3740        // Let the handler enter the (never-completing) service call, then model the
3741        // request-stream side giving up: it drops the oneshot receiver.
3742        tokio::task::yield_now().await;
3743        drop(response_rx);
3744
3745        // The handler must observe the dropped receiver, abort the service work, and
3746        // release the permit promptly. Without the fix this times out because the
3747        // handler stays blocked in `service.call` until `LEGACY_REQUEST_TIMEOUT`.
3748        tokio::time::timeout(TEST_NET_TIMEOUT, handler)
3749            .await
3750            .expect("handler aborts promptly after the request stream drops the receiver")
3751            .expect("handler task does not panic");
3752
3753        assert_eq!(
3754            permits.available_permits(),
3755            LEGACY_REQUEST_IN_FLIGHT_LIMIT,
3756            "the in-flight permit must be released once the handler aborts"
3757        );
3758    }
3759
3760    #[test]
3761    fn block_gossip_round_trips_and_rejects_trailing_bytes() {
3762        let frame = LegacyGossipFrame::AdvertiseBlock(block_hash(1));
3763        let mut encoded = frame.encode_frame().expect("frame encodes");
3764        assert_eq!(
3765            LegacyGossipFrame::decode_frame(encoded.clone()).expect("frame decodes"),
3766            frame
3767        );
3768
3769        encoded.payload.push(0);
3770        assert!(matches!(
3771            LegacyGossipFrame::decode_frame(encoded),
3772            Err(LegacyGossipError::TrailingBytes)
3773        ));
3774    }
3775
3776    #[tokio::test]
3777    async fn adapter_gossips_block_and_tx_ids_between_two_nodes() -> Result<(), BoxError> {
3778        let _guard = zakura_test::init();
3779        let (node_a, mut rx_a) = legacy_node(11).await?;
3780        let (node_b, mut rx_b) = legacy_node(12).await?;
3781        node_a.connect_native(&node_b, TEST_NET_TIMEOUT).await?;
3782        let a_peer_id = node_peer_id(&node_a).await?;
3783
3784        let mut adapter = LegacyGossipAdapter::new(node_a.supervisor());
3785        let block_hash = block_hash(9);
3786        adapter
3787            .ready()
3788            .await?
3789            .call(Request::AdvertiseBlockToAll(block_hash))
3790            .await?;
3791
3792        match recv_request(&mut rx_b).await? {
3793            Request::AdvertiseBlock(hash, Some(PeerSource::Zakura(peer_id))) => {
3794                assert_eq!(hash, block_hash);
3795                assert_eq!(peer_id, a_peer_id);
3796            }
3797            request => panic!("unexpected request: {request:?}"),
3798        }
3799        assert!(rx_a.try_recv().is_err());
3800
3801        let ids = HashSet::from([legacy_tx_id(10), witnessed_tx_id(11)]);
3802        adapter
3803            .ready()
3804            .await?
3805            .call(Request::AdvertiseTransactionIds(ids.clone(), None))
3806            .await?;
3807
3808        match recv_request(&mut rx_b).await? {
3809            Request::AdvertiseTransactionIds(received, Some(PeerSource::Zakura(peer_id))) => {
3810                assert_eq!(received, ids);
3811                assert_eq!(peer_id, a_peer_id);
3812            }
3813            request => panic!("unexpected request: {request:?}"),
3814        }
3815
3816        node_a.shutdown().await;
3817        node_b.shutdown().await;
3818        Ok(())
3819    }
3820
3821    #[tokio::test]
3822    async fn request_adapter_fetches_missing_block_from_advertiser() -> Result<(), BoxError> {
3823        let _guard = zakura_test::init();
3824        let transaction = UnminedTx::from(empty_v5_transaction(1));
3825        let node_a = inventory_node(61, transaction).await?;
3826        let node_b = ZakuraTestNode::builder(62).spawn().await?;
3827        node_b.connect_native(&node_a, TEST_NET_TIMEOUT).await?;
3828        let a_peer_id = node_peer_id(&node_a).await?;
3829        let hash = block_hash(90);
3830
3831        let adapter = LegacyRequestAdapter::new(node_b.supervisor());
3832        let response = adapter
3833            .request_from_source(
3834                Request::BlocksByHash(HashSet::from([hash])),
3835                Some(PeerSource::Zakura(a_peer_id)),
3836            )
3837            .await?;
3838
3839        match response {
3840            Response::Blocks(blocks) => {
3841                assert_eq!(blocks, vec![InventoryResponse::Missing(hash)]);
3842            }
3843            response => panic!("unexpected response: {response:?}"),
3844        }
3845
3846        node_a.shutdown().await;
3847        node_b.shutdown().await;
3848        Ok(())
3849    }
3850
3851    #[tokio::test]
3852    async fn request_adapter_fetches_available_multi_chunk_block() -> Result<(), BoxError> {
3853        let _guard = zakura_test::init();
3854        let block = Arc::new(Block::zcash_deserialize(
3855            BLOCK_TESTNET_141042_BYTES.as_slice(),
3856        )?);
3857        let serialized_len = block.zcash_serialize_to_vec()?.len();
3858        assert!(
3859            serialized_len > LEGACY_RESPONSE_CHUNK_BYTES,
3860            "test block must span multiple response chunks"
3861        );
3862
3863        let node_a = block_inventory_node(68, block.clone()).await?;
3864        let node_b = ZakuraTestNode::builder(69).spawn().await?;
3865        node_b.connect_native(&node_a, TEST_NET_TIMEOUT).await?;
3866        let a_peer_id = node_peer_id(&node_a).await?;
3867
3868        let adapter = LegacyRequestAdapter::new(node_b.supervisor());
3869        let response = adapter
3870            .request_from_source(
3871                Request::BlocksByHash(HashSet::from([block.hash()])),
3872                Some(PeerSource::Zakura(a_peer_id)),
3873            )
3874            .await?;
3875
3876        let Response::Blocks(blocks) = response else {
3877            panic!("unexpected response: {response:?}");
3878        };
3879        assert!(matches!(
3880            blocks.as_slice(),
3881            [InventoryResponse::Available((received, None))] if received.hash() == block.hash()
3882        ));
3883
3884        node_a.shutdown().await;
3885        node_b.shutdown().await;
3886        Ok(())
3887    }
3888
3889    #[tokio::test]
3890    async fn request_adapter_fetches_available_and_missing_transactions() -> Result<(), BoxError> {
3891        let _guard = zakura_test::init();
3892        let transaction = UnminedTx::from(empty_v5_transaction(2));
3893        let available_id = transaction.id();
3894        let missing_id = witnessed_tx_id(99);
3895        let node_a = inventory_node(63, transaction.clone()).await?;
3896        let node_b = ZakuraTestNode::builder(64).spawn().await?;
3897        node_b.connect_native(&node_a, TEST_NET_TIMEOUT).await?;
3898        let a_peer_id = node_peer_id(&node_a).await?;
3899
3900        let adapter = LegacyRequestAdapter::new(node_b.supervisor());
3901        let response = adapter
3902            .request_from_source(
3903                Request::TransactionsById(HashSet::from([available_id, missing_id])),
3904                Some(PeerSource::Zakura(a_peer_id)),
3905            )
3906            .await?;
3907
3908        let Response::Transactions(transactions) = response else {
3909            panic!("unexpected response: {response:?}");
3910        };
3911        assert!(transactions.iter().any(|response| {
3912            matches!(
3913                response,
3914                InventoryResponse::Available((tx, None)) if tx.id() == transaction.id()
3915            )
3916        }));
3917        assert!(transactions.iter().any(
3918            |response| matches!(response, InventoryResponse::Missing(id) if *id == missing_id)
3919        ));
3920
3921        node_a.shutdown().await;
3922        node_b.shutdown().await;
3923        Ok(())
3924    }
3925
3926    #[tokio::test]
3927    async fn request_adapter_drives_mock_chain_sync_progress_over_zakura() -> Result<(), BoxError> {
3928        let _guard = zakura_test::init();
3929        let block = Arc::new(Block::zcash_deserialize(
3930            BLOCK_TESTNET_141042_BYTES.as_slice(),
3931        )?);
3932        let transaction = UnminedTx::from(empty_v5_transaction(4));
3933        let (node_a, _pushed_rx) = normal_network_node(81, block.clone(), transaction).await?;
3934        let node_b = ZakuraTestNode::builder(82).spawn().await?;
3935        node_b.connect_native(&node_a, TEST_NET_TIMEOUT).await?;
3936        let a_peer_id = node_peer_id(&node_a).await?;
3937
3938        let adapter = LegacyRequestAdapter::new(node_b.supervisor());
3939        let mut local_chain = vec![block_hash(1)];
3940        let find_headers = adapter
3941            .request_from_source(
3942                Request::FindHeaders {
3943                    known_blocks: local_chain.clone(),
3944                    stop: Some(block.hash()),
3945                },
3946                Some(PeerSource::Zakura(a_peer_id.clone())),
3947            )
3948            .await?;
3949        let Response::BlockHeaders(headers) = find_headers else {
3950            panic!("unexpected FindHeaders response: {find_headers:?}");
3951        };
3952        assert!(matches!(
3953            headers.as_slice(),
3954            [header] if header.header == block.header
3955        ));
3956
3957        let find_blocks = adapter
3958            .request_from_source(
3959                Request::FindBlocks {
3960                    known_blocks: local_chain.clone(),
3961                    stop: Some(block.hash()),
3962                },
3963                Some(PeerSource::Zakura(a_peer_id.clone())),
3964            )
3965            .await?;
3966        let Response::BlockHashes(remote_hashes) = find_blocks else {
3967            panic!("unexpected FindBlocks response: {find_blocks:?}");
3968        };
3969        assert_eq!(remote_hashes, vec![block.hash()]);
3970
3971        let block_response = adapter
3972            .request_from_source(
3973                Request::BlocksByHashFrom {
3974                    hashes: remote_hashes.into_iter().collect(),
3975                    source: PeerSource::Zakura(a_peer_id),
3976                },
3977                None,
3978            )
3979            .await?;
3980        let Response::Blocks(blocks) = block_response else {
3981            panic!("unexpected block response: {block_response:?}");
3982        };
3983        for response in blocks {
3984            let received = response
3985                .available()
3986                .expect("mock sync peer serves the advertised block")
3987                .0;
3988            local_chain.push(received.hash());
3989        }
3990        assert_eq!(local_chain.last(), Some(&block.hash()));
3991        assert_eq!(local_chain.len(), 2);
3992
3993        node_a.shutdown().await;
3994        node_b.shutdown().await;
3995        Ok(())
3996    }
3997
3998    #[tokio::test]
3999    async fn controlled_three_node_network_gossips_fetches_and_advances_mock_chain(
4000    ) -> Result<(), BoxError> {
4001        let _guard = zakura_test::init();
4002        let block = Arc::new(Block::zcash_deserialize(
4003            BLOCK_TESTNET_141042_BYTES.as_slice(),
4004        )?);
4005        let transaction = UnminedTx::from(empty_v5_transaction(14));
4006        let (node_a, _pushed_rx) = normal_network_node(85, block.clone(), transaction).await?;
4007        let (node_b, mut rx_b) = legacy_node(86).await?;
4008        let (node_c, mut rx_c) = legacy_node(87).await?;
4009
4010        node_a.connect_native(&node_b, TEST_NET_TIMEOUT).await?;
4011        node_a.connect_native(&node_c, TEST_NET_TIMEOUT).await?;
4012        node_b.connect_native(&node_c, TEST_NET_TIMEOUT).await?;
4013        wait_registered_count(&node_a, 2).await?;
4014        wait_registered_count(&node_b, 2).await?;
4015        wait_registered_count(&node_c, 2).await?;
4016
4017        let a_peer_id = node_peer_id(&node_a).await?;
4018        let mut gossip = LegacyGossipAdapter::new(node_a.supervisor());
4019        gossip
4020            .ready()
4021            .await?
4022            .call(Request::AdvertiseBlockToAll(block.hash()))
4023            .await?;
4024
4025        match recv_request(&mut rx_b).await? {
4026            Request::AdvertiseBlock(hash, Some(PeerSource::Zakura(peer_id))) => {
4027                assert_eq!(hash, block.hash());
4028                assert_eq!(peer_id, a_peer_id);
4029            }
4030            request => panic!("unexpected B request: {request:?}"),
4031        }
4032        match recv_request(&mut rx_c).await? {
4033            Request::AdvertiseBlock(hash, Some(PeerSource::Zakura(_))) => {
4034                assert_eq!(hash, block.hash());
4035            }
4036            request => panic!("unexpected C request: {request:?}"),
4037        }
4038
4039        let adapter = LegacyRequestAdapter::new(node_c.supervisor());
4040        let mut local_chain = vec![block_hash(1)];
4041        let headers = adapter
4042            .request_from_source(
4043                Request::FindHeaders {
4044                    known_blocks: local_chain.clone(),
4045                    stop: Some(block.hash()),
4046                },
4047                Some(PeerSource::Zakura(a_peer_id.clone())),
4048            )
4049            .await?;
4050        assert!(matches!(
4051            headers,
4052            Response::BlockHeaders(headers) if matches!(
4053                headers.as_slice(),
4054                [header] if header.header == block.header
4055            )
4056        ));
4057
4058        let hashes = adapter
4059            .request_from_source(
4060                Request::FindBlocks {
4061                    known_blocks: local_chain.clone(),
4062                    stop: Some(block.hash()),
4063                },
4064                Some(PeerSource::Zakura(a_peer_id.clone())),
4065            )
4066            .await?;
4067        let Response::BlockHashes(hashes) = hashes else {
4068            panic!("unexpected FindBlocks response: {hashes:?}");
4069        };
4070        assert_eq!(hashes, vec![block.hash()]);
4071
4072        let blocks = adapter
4073            .request_from_source(
4074                Request::BlocksByHashFrom {
4075                    hashes: hashes.into_iter().collect(),
4076                    source: PeerSource::Zakura(a_peer_id),
4077                },
4078                None,
4079            )
4080            .await?;
4081        let Response::Blocks(blocks) = blocks else {
4082            panic!("unexpected block response: {blocks:?}");
4083        };
4084        for response in blocks {
4085            let received = response
4086                .available()
4087                .expect("mock sync peer serves the gossiped block")
4088                .0;
4089            local_chain.push(received.hash());
4090        }
4091        assert_eq!(local_chain.last(), Some(&block.hash()));
4092        assert_eq!(local_chain.len(), 2);
4093
4094        node_a.shutdown().await;
4095        node_b.shutdown().await;
4096        node_c.shutdown().await;
4097        Ok(())
4098    }
4099
4100    #[tokio::test]
4101    async fn request_adapter_supports_mempool_ping_and_push_transaction() -> Result<(), BoxError> {
4102        let _guard = zakura_test::init();
4103        let block = Arc::new(Block::zcash_deserialize(
4104            BLOCK_TESTNET_141042_BYTES.as_slice(),
4105        )?);
4106        let transaction = UnminedTx::from(empty_v5_transaction(5));
4107        let pushed_transaction = UnminedTx::from(empty_v5_transaction(6));
4108        let pushed_id = pushed_transaction.id();
4109        let (node_a, mut pushed_rx) = normal_network_node(83, block, transaction.clone()).await?;
4110        let node_b = ZakuraTestNode::builder(84).spawn().await?;
4111        node_b.connect_native(&node_a, TEST_NET_TIMEOUT).await?;
4112        let a_peer_id = node_peer_id(&node_a).await?;
4113
4114        let adapter = LegacyRequestAdapter::new(node_b.supervisor());
4115        let mempool = adapter
4116            .request_from_source(
4117                Request::MempoolTransactionIds,
4118                Some(PeerSource::Zakura(a_peer_id.clone())),
4119            )
4120            .await?;
4121        assert_eq!(mempool, Response::TransactionIds(vec![transaction.id()]));
4122
4123        let ping = adapter
4124            .request_from_source(
4125                Request::Ping(Default::default()),
4126                Some(PeerSource::Zakura(a_peer_id)),
4127            )
4128            .await?;
4129        assert!(matches!(ping, Response::Pong(_)));
4130
4131        let push_response = adapter
4132            .request_from_source(Request::PushTransaction(pushed_transaction, None), None)
4133            .await?;
4134        assert_eq!(push_response, Response::Nil);
4135        assert_eq!(recv_pushed_tx_id(&mut pushed_rx).await?, pushed_id);
4136
4137        node_a.shutdown().await;
4138        node_b.shutdown().await;
4139        Ok(())
4140    }
4141
4142    #[tokio::test]
4143    async fn slow_request_does_not_block_gossip_worker() -> Result<(), BoxError> {
4144        let (request_tx, request_rx) = oneshot::channel();
4145        let (gossip_tx, mut gossip_rx) = tokio::sync::mpsc::unbounded_channel();
4146        let (release_tx, release_rx) = tokio::sync::watch::channel(false);
4147        let (inbound_tx, inbound_rx) = mpsc::channel(4);
4148        let supervisor = ZakuraSupervisorHandle::new(1);
4149        let peer_id = ZakuraPeerId::new(vec![6; 32]).expect("test peer id is within bounds");
4150
4151        let worker = tokio::spawn(legacy_gossip_worker(
4152            SlowRequestThenRecorder {
4153                release: release_rx,
4154                tx: gossip_tx,
4155            },
4156            inbound_rx,
4157            LegacyGossipForwarder::new(supervisor),
4158            ZakuraTrace::noop(),
4159        ));
4160
4161        inbound_tx
4162            .send(LegacyInboundWork::Request(LegacyRequestInbound {
4163                peer_id: peer_id.clone(),
4164                request_id: 1,
4165                frame: LegacyRequestFrame::BlocksByHash(vec![block_hash(11)]),
4166                response_tx: request_tx,
4167            }))
4168            .await?;
4169        inbound_tx
4170            .send(LegacyInboundWork::Gossip(LegacyGossipInbound {
4171                peer_id: peer_id.clone(),
4172                frame: LegacyGossipFrame::AdvertiseBlock(block_hash(12)),
4173            }))
4174            .await?;
4175
4176        match recv_request(&mut gossip_rx).await? {
4177            Request::AdvertiseBlock(hash, Some(PeerSource::Zakura(source))) => {
4178                assert_eq!(hash, block_hash(12));
4179                assert_eq!(source, peer_id);
4180            }
4181            request => panic!("unexpected request: {request:?}"),
4182        }
4183
4184        release_tx.send(true)?;
4185        match tokio::time::timeout(Duration::from_secs(1), request_rx).await??? {
4186            Response::Blocks(blocks) => {
4187                assert_eq!(blocks, vec![InventoryResponse::Missing(block_hash(11))]);
4188            }
4189            response => panic!("unexpected request response: {response:?}"),
4190        }
4191
4192        drop(inbound_tx);
4193        worker.await?;
4194        Ok(())
4195    }
4196
4197    /// Regression: while the legacy inbound service is slow/timing out,
4198    /// `handle_legacy_gossip` returns without `mark_seen`, so an authenticated peer
4199    /// could replay the same valid advertisement and make the serial worker re-pay
4200    /// the full readiness/call timeout for every queued duplicate. After one
4201    /// expensive failed attempt the cooldown must suppress identical duplicates.
4202    ///
4203    /// Paused time auto-advances the `LEGACY_GOSSIP_SERVICE_TIMEOUT` so the first
4204    /// call's timeout fires without a real wall-clock wait.
4205    #[tokio::test(start_paused = true)]
4206    async fn duplicate_gossip_does_not_repeat_expensive_attempts_while_service_is_slow(
4207    ) -> Result<(), BoxError> {
4208        let calls = Arc::new(AtomicUsize::new(0));
4209        let (inbound_tx, inbound_rx) = mpsc::channel(16);
4210        let supervisor = ZakuraSupervisorHandle::new(1);
4211        let peer_id = ZakuraPeerId::new(vec![7; 32]).expect("test peer id is within bounds");
4212
4213        let worker = tokio::spawn(legacy_gossip_worker(
4214            CountingPendingService {
4215                calls: calls.clone(),
4216            },
4217            inbound_rx,
4218            LegacyGossipForwarder::new(supervisor),
4219            ZakuraTrace::noop(),
4220        ));
4221
4222        // The call never completes, so the permanent first-seen cache is never
4223        // updated; only the post-timeout cooldown can suppress these duplicates.
4224        let frame = LegacyGossipFrame::AdvertiseBlock(block_hash(99));
4225        for _ in 0..4 {
4226            inbound_tx
4227                .send(LegacyInboundWork::Gossip(LegacyGossipInbound {
4228                    peer_id: peer_id.clone(),
4229                    frame: frame.clone(),
4230                }))
4231                .await?;
4232        }
4233
4234        drop(inbound_tx);
4235        worker.await?;
4236
4237        assert_eq!(
4238            calls.load(Ordering::SeqCst),
4239            1,
4240            "duplicate gossip frames must not each re-pay the inbound readiness/call \
4241             timeout while the service is slow; expected one expensive attempt with \
4242             the rest suppressed by the cooldown"
4243        );
4244
4245        Ok(())
4246    }
4247
4248    #[tokio::test]
4249    async fn service_request_routes_to_advertiser_then_fallback_after_missing(
4250    ) -> Result<(), BoxError> {
4251        let _guard = zakura_test::init();
4252        let transaction = UnminedTx::from(empty_v5_transaction(3));
4253        let txid = transaction.id();
4254        let (advertiser, mut advertiser_rx) = recording_inventory_node(65, None).await?;
4255        let (fallback, mut fallback_rx) =
4256            recording_inventory_node(66, Some(transaction.clone())).await?;
4257        // The requester dials two loopback peers (advertiser + fallback), so it
4258        // raises the production Zakura per-IP cap that the default test node
4259        // enforces; this test exercises service routing, not per-IP admission.
4260        let requester = ZakuraTestNode::builder(67)
4261            .max_connections_per_ip(8)
4262            .spawn()
4263            .await?;
4264        requester
4265            .connect_native(&advertiser, TEST_NET_TIMEOUT)
4266            .await?;
4267        requester
4268            .connect_native(&fallback, TEST_NET_TIMEOUT)
4269            .await?;
4270        wait_registered_count(&requester, 2).await?;
4271        let advertiser_id = node_peer_id(&advertiser).await?;
4272
4273        let mut adapter = LegacyRequestAdapter::new(requester.supervisor());
4274        let response = adapter
4275            .ready()
4276            .await?
4277            .call(Request::TransactionsByIdFrom {
4278                ids: HashSet::from([txid]),
4279                source: PeerSource::Zakura(advertiser_id),
4280            })
4281            .await?;
4282
4283        match recv_request(&mut advertiser_rx).await? {
4284            Request::TransactionsById(ids) => assert_eq!(ids, HashSet::from([txid])),
4285            request => panic!("unexpected advertiser request: {request:?}"),
4286        }
4287        match recv_request(&mut fallback_rx).await? {
4288            Request::TransactionsById(ids) => assert_eq!(ids, HashSet::from([txid])),
4289            request => panic!("unexpected fallback request: {request:?}"),
4290        }
4291
4292        let Response::Transactions(transactions) = response else {
4293            panic!("unexpected response: {response:?}");
4294        };
4295        assert!(matches!(
4296            transactions.as_slice(),
4297            [InventoryResponse::Available((tx, None))] if tx.id() == transaction.id()
4298        ));
4299
4300        advertiser.shutdown().await;
4301        fallback.shutdown().await;
4302        requester.shutdown().await;
4303        Ok(())
4304    }
4305
4306    #[tokio::test]
4307    async fn first_seen_gossip_forwards_across_line_and_drops_duplicates() -> Result<(), BoxError> {
4308        let _guard = zakura_test::init();
4309        let (node_a, mut rx_a) = legacy_node(21).await?;
4310        let (node_b, mut rx_b) = legacy_node(22).await?;
4311        let (node_c, mut rx_c) = legacy_node(23).await?;
4312
4313        node_a.connect_native(&node_b, TEST_NET_TIMEOUT).await?;
4314        node_b.connect_native(&node_c, TEST_NET_TIMEOUT).await?;
4315
4316        wait_registered_count(&node_b, 2).await?;
4317
4318        let a_peer_id = node_peer_id(&node_a).await?;
4319        let b_peer_id = node_peer_id(&node_b).await?;
4320        let block_hash = block_hash(42);
4321        let mut adapter = LegacyGossipAdapter::new(node_a.supervisor());
4322        adapter
4323            .ready()
4324            .await?
4325            .call(Request::AdvertiseBlockToAll(block_hash))
4326            .await?;
4327
4328        match recv_request(&mut rx_b).await? {
4329            Request::AdvertiseBlock(hash, Some(PeerSource::Zakura(peer_id))) => {
4330                assert_eq!(hash, block_hash);
4331                assert_eq!(peer_id, a_peer_id);
4332            }
4333            request => panic!("unexpected B request: {request:?}"),
4334        }
4335        match recv_request(&mut rx_c).await? {
4336            Request::AdvertiseBlock(hash, Some(PeerSource::Zakura(peer_id))) => {
4337                assert_eq!(hash, block_hash);
4338                assert_eq!(peer_id, b_peer_id);
4339            }
4340            request => panic!("unexpected C request: {request:?}"),
4341        }
4342        assert!(rx_a.try_recv().is_err());
4343
4344        adapter
4345            .ready()
4346            .await?
4347            .call(Request::AdvertiseBlockToAll(block_hash))
4348            .await?;
4349        tokio::time::sleep(Duration::from_millis(100)).await;
4350        assert!(rx_c.try_recv().is_err());
4351
4352        node_a.shutdown().await;
4353        node_b.shutdown().await;
4354        node_c.shutdown().await;
4355        Ok(())
4356    }
4357
4358    #[tokio::test]
4359    async fn local_origin_echo_is_dropped_by_shared_first_seen_cache() -> Result<(), BoxError> {
4360        let _guard = zakura_test::init();
4361        let (node_a, mut rx_a) = legacy_node(51).await?;
4362        let (node_b, mut rx_b) = legacy_node(52).await?;
4363        node_a.connect_native(&node_b, TEST_NET_TIMEOUT).await?;
4364        let hostile = HostilePeer::connect_native(&node_a, 53).await?;
4365        wait_registered_count(&node_a, 2).await?;
4366
4367        let block_hash = block_hash(88);
4368        let mut adapter = LegacyGossipAdapter::new(node_a.supervisor());
4369        adapter
4370            .ready()
4371            .await?
4372            .call(Request::AdvertiseBlockToAll(block_hash))
4373            .await?;
4374        let _ = recv_request(&mut rx_b).await?;
4375
4376        let echo_payload = LegacyGossipFrame::AdvertiseBlock(block_hash)
4377            .encode_frame()?
4378            .payload;
4379        hostile
4380            .send_frame(ZAKURA_STREAM_GOSSIP, echo_payload)
4381            .await?;
4382        tokio::time::sleep(Duration::from_millis(100)).await;
4383
4384        assert!(
4385            rx_a.try_recv().is_err(),
4386            "origin must not accept its own echo"
4387        );
4388        assert!(
4389            rx_b.try_recv().is_err(),
4390            "origin must not re-forward its own echo"
4391        );
4392
4393        hostile.shutdown().await;
4394        node_a.shutdown().await;
4395        node_b.shutdown().await;
4396        Ok(())
4397    }
4398
4399    #[tokio::test]
4400    async fn saturated_peer_does_not_block_honest_fanout() -> Result<(), BoxError> {
4401        let saturated_peer = ZakuraPeerId::new(vec![1; 32]).expect("test peer id is within bounds");
4402        let honest_peer = ZakuraPeerId::new(vec![2; 32]).expect("test peer id is within bounds");
4403        let (saturated, _saturated_rx) = framed_channel(1);
4404        saturated.try_send(Frame {
4405            message_type: MSG_ADVERTISE_BLOCK,
4406            flags: 0,
4407            payload: vec![1],
4408        })?;
4409
4410        let (honest, mut honest_rx) = framed_channel(1);
4411        let block_hash = block_hash(7);
4412
4413        let result = send_to_sessions(
4414            vec![
4415                LegacyGossipPeerSession::new(saturated_peer, 0, 0, saturated),
4416                LegacyGossipPeerSession::new(honest_peer, 0, 0, honest),
4417            ],
4418            LegacyGossipFrame::AdvertiseBlock(block_hash),
4419        );
4420        let outbound = tokio::time::timeout(Duration::from_secs(1), honest_rx.recv())
4421            .await?
4422            .expect("honest peer receives fanout");
4423        assert_eq!(
4424            LegacyGossipFrame::decode_frame(outbound)?,
4425            LegacyGossipFrame::AdvertiseBlock(block_hash)
4426        );
4427
4428        assert!(
4429            result.is_err(),
4430            "saturated peer failure is reported after honest peer is attempted"
4431        );
4432        Ok(())
4433    }
4434
4435    #[tokio::test]
4436    async fn new_peer_receives_latest_block_advertised_before_gossip_stream_ready(
4437    ) -> Result<(), BoxError> {
4438        let supervisor = ZakuraSupervisorHandle::new(991);
4439        let broadcast = ZakuraGossipBroadcast::new(supervisor);
4440        let block_hash = block_hash(91);
4441
4442        broadcast
4443            .broadcast(LegacyGossipFrame::AdvertiseBlock(block_hash), None)
4444            .await?;
4445
4446        let peer_id = ZakuraPeerId::new(vec![91; 32]).expect("test peer id is within bounds");
4447        let (sender, mut receiver) = framed_channel(1);
4448        let session = LegacyGossipPeerSession::new(peer_id, 0, 0, sender);
4449        broadcast.outbound.insert(session.clone());
4450        broadcast
4451            .outbound
4452            .replay_latest_block_to_peer(session)
4453            .await?;
4454
4455        let replayed = tokio::time::timeout(Duration::from_secs(1), receiver.recv())
4456            .await?
4457            .expect("new peer receives latest block replay");
4458        assert_eq!(
4459            LegacyGossipFrame::decode_frame(replayed)?,
4460            LegacyGossipFrame::AdvertiseBlock(block_hash)
4461        );
4462        Ok(())
4463    }
4464
4465    #[tokio::test]
4466    async fn legacy_gossip_replay_panic_cancels_peer_and_removes_outbound_session(
4467    ) -> Result<(), BoxError> {
4468        let (inbound_tx, _inbound_rx) = mpsc::channel(1);
4469        let outbound = LegacyGossipOutbound::default();
4470        let sink = LegacyGossipSink {
4471            inbound_tx,
4472            outbound: outbound.clone(),
4473            trace: ZakuraTrace::noop(),
4474        };
4475        let peer_id = ZakuraPeerId::new(vec![92; 32]).expect("test peer id is within bounds");
4476        let cancel_token = tokio_util::sync::CancellationToken::new();
4477        let (peer, _peer_send) = legacy_gossip_peer(peer_id.clone(), cancel_token.clone());
4478
4479        let poisoned_latest_block = outbound.latest_block.clone();
4480        let _ = std::panic::catch_unwind(move || {
4481            let _guard = poisoned_latest_block
4482                .lock()
4483                .expect("latest-block mutex starts unpoisoned");
4484            panic!("poison latest-block mutex before replay");
4485        });
4486
4487        sink.add_peer(peer);
4488        wait_for_legacy_gossip_panic_cleanup(&outbound, &peer_id, &cancel_token).await
4489    }
4490
4491    #[tokio::test]
4492    async fn legacy_gossip_recv_loop_panic_cancels_peer_and_removes_outbound_session(
4493    ) -> Result<(), BoxError> {
4494        let (inbound_tx, _inbound_rx) = mpsc::channel(1);
4495        let outbound = LegacyGossipOutbound::default();
4496        let sink = LegacyGossipSink {
4497            inbound_tx,
4498            outbound: outbound.clone(),
4499            trace: ZakuraTrace::noop(),
4500        };
4501        let peer_id = ZakuraPeerId::new(vec![93; 32]).expect("test peer id is within bounds");
4502        let cancel_token = tokio_util::sync::CancellationToken::new();
4503        let (peer, peer_send) = legacy_gossip_peer(peer_id.clone(), cancel_token.clone());
4504
4505        arm_legacy_gossip_recv_loop_panic(peer_id.clone());
4506        sink.add_peer(peer);
4507        peer_send
4508            .send(LegacyGossipFrame::AdvertiseBlock(block_hash(93)).encode_frame()?)
4509            .await
4510            .map_err(|_| -> BoxError { "failed to send test gossip frame".into() })?;
4511
4512        wait_for_legacy_gossip_panic_cleanup(&outbound, &peer_id, &cancel_token).await
4513    }
4514
4515    #[tokio::test]
4516    async fn stale_legacy_gossip_teardown_keeps_replacement_session() {
4517        let (inbound_tx, _inbound_rx) = mpsc::channel(1);
4518        let outbound = LegacyGossipOutbound::default();
4519        let sink = LegacyGossipSink {
4520            inbound_tx,
4521            outbound: outbound.clone(),
4522            trace: ZakuraTrace::noop(),
4523        };
4524        let peer_id = ZakuraPeerId::new(vec![94; 32]).expect("test peer id is within bounds");
4525        let old_conn_id = 1;
4526        let new_conn_id = 2;
4527        let old_cancel = tokio_util::sync::CancellationToken::new();
4528        let new_cancel = tokio_util::sync::CancellationToken::new();
4529        let (old_peer, _old_peer_send) =
4530            legacy_gossip_peer_with_conn(peer_id.clone(), old_conn_id, old_cancel);
4531        let (new_peer, _new_peer_send) =
4532            legacy_gossip_peer_with_conn(peer_id.clone(), new_conn_id, new_cancel);
4533
4534        sink.add_peer(old_peer);
4535        assert!(
4536            outbound.contains(&peer_id),
4537            "old legacy gossip session is registered",
4538        );
4539        sink.add_peer(new_peer);
4540        assert!(
4541            outbound.contains(&peer_id),
4542            "replacement legacy gossip session is registered",
4543        );
4544        assert!(!sink.owns_connection_for_peer(&peer_id, old_conn_id));
4545        assert!(sink.owns_connection_for_peer(&peer_id, new_conn_id));
4546
4547        let (stale_peer, _stale_peer_send) = legacy_gossip_peer_with_conn(
4548            peer_id.clone(),
4549            old_conn_id,
4550            tokio_util::sync::CancellationToken::new(),
4551        );
4552        sink.add_peer(stale_peer);
4553
4554        sink.remove_peer(&peer_id, old_conn_id);
4555        assert!(
4556            outbound.contains(&peer_id),
4557            "stale cleanup must not remove the replacement legacy gossip session",
4558        );
4559
4560        sink.remove_peer(&peer_id, new_conn_id);
4561        assert!(
4562            !outbound.contains(&peer_id),
4563            "live cleanup removes the replacement legacy gossip session",
4564        );
4565    }
4566
4567    #[tokio::test]
4568    async fn stale_gossip_teardown_cannot_retire_a_reopened_session() {
4569        let outbound = LegacyGossipOutbound::default();
4570        let peer_id = ZakuraPeerId::new(vec![96; 32]).expect("test peer id is within bounds");
4571        let conn_id = 7;
4572        let (first_send, _first_rx) = framed_channel(1);
4573        let (second_send, _second_rx) = framed_channel(1);
4574
4575        assert!(outbound.insert(LegacyGossipPeerSession::new(
4576            peer_id.clone(),
4577            conn_id,
4578            1,
4579            first_send
4580        )));
4581        // The transport reopened the gossip stream on the same connection.
4582        assert!(outbound.insert(LegacyGossipPeerSession::new(
4583            peer_id.clone(),
4584            conn_id,
4585            2,
4586            second_send
4587        )));
4588
4589        // A delayed teardown from the first stream session must not retire the
4590        // reopened session.
4591        outbound.finish_session(&peer_id, conn_id, 1, false);
4592        assert!(
4593            outbound.contains(&peer_id),
4594            "stale teardown must not retire the reopened gossip session",
4595        );
4596        assert!(outbound.owns_connection(&peer_id, conn_id));
4597
4598        // A replayed older stream session must not replace the live one.
4599        let (stale_send, _stale_rx) = framed_channel(1);
4600        assert!(!outbound.insert(LegacyGossipPeerSession::new(
4601            peer_id.clone(),
4602            conn_id,
4603            2,
4604            stale_send
4605        )));
4606
4607        // The owning session's teardown retires it and claims the reopen gap.
4608        outbound.finish_session(&peer_id, conn_id, 2, false);
4609        assert!(!outbound.contains(&peer_id));
4610        assert!(
4611            outbound.owns_connection(&peer_id, conn_id),
4612            "reopen gap keeps the connection owned",
4613        );
4614    }
4615
4616    #[tokio::test]
4617    async fn repeated_no_frame_gossip_sessions_retire_the_service() {
4618        let (inbound_tx, _inbound_rx) = mpsc::channel(1);
4619        let outbound = LegacyGossipOutbound::default();
4620        let sink = LegacyGossipSink {
4621            inbound_tx,
4622            outbound: outbound.clone(),
4623            trace: ZakuraTrace::noop(),
4624        };
4625        let peer_id = ZakuraPeerId::new(vec![98; 32]).expect("test peer id is within bounds");
4626        let conn_id = 7;
4627
4628        for session_id in 1..=u64::from(MAX_GOSSIP_NO_FRAME_SESSIONS) {
4629            let (send, _rx) = framed_channel(1);
4630            assert!(
4631                outbound.insert(LegacyGossipPeerSession::new(
4632                    peer_id.clone(),
4633                    conn_id,
4634                    session_id,
4635                    send
4636                )),
4637                "session {session_id} is admitted before the churn bound"
4638            );
4639            outbound.finish_session(&peer_id, conn_id, session_id, true);
4640        }
4641
4642        assert!(
4643            !outbound.owns_connection(&peer_id, conn_id),
4644            "a retired gossip stream must not hold a reopen-gap claim"
4645        );
4646        assert!(matches!(
4647            sink.ordered_session_demand(conn_id, &peer_id, 0, ServicePeerDirection::Outbound),
4648            OrderedSessionDemand::Retire,
4649        ));
4650        let (send, _rx) = framed_channel(1);
4651        assert!(
4652            !outbound.insert(LegacyGossipPeerSession::new(
4653                peer_id.clone(),
4654                conn_id,
4655                u64::from(MAX_GOSSIP_NO_FRAME_SESSIONS) + 1,
4656                send
4657            )),
4658            "a retired connection must refuse new gossip sessions"
4659        );
4660
4661        // A fresh connection starts a fresh churn record.
4662        let new_conn_id = conn_id + 1;
4663        let (send, _rx) = framed_channel(1);
4664        assert!(outbound.insert(LegacyGossipPeerSession::new(
4665            peer_id.clone(),
4666            new_conn_id,
4667            1,
4668            send
4669        )));
4670    }
4671
4672    #[tokio::test]
4673    async fn a_received_frame_resets_gossip_churn() {
4674        let (inbound_tx, _inbound_rx) = mpsc::channel(1);
4675        let sink = LegacyGossipSink {
4676            inbound_tx,
4677            outbound: LegacyGossipOutbound::default(),
4678            trace: ZakuraTrace::noop(),
4679        };
4680        let outbound = sink.outbound.clone();
4681        let peer_id = ZakuraPeerId::new(vec![99; 32]).expect("test peer id is within bounds");
4682        let conn_id = 7;
4683        let mut session_id = 0;
4684        let mut run_session = |no_frame_exit: bool| {
4685            session_id += 1;
4686            let (send, _rx) = framed_channel(1);
4687            assert!(outbound.insert(LegacyGossipPeerSession::new(
4688                peer_id.clone(),
4689                conn_id,
4690                session_id,
4691                send
4692            )));
4693            outbound.finish_session(&peer_id, conn_id, session_id, no_frame_exit);
4694        };
4695
4696        for _ in 1..MAX_GOSSIP_NO_FRAME_SESSIONS {
4697            run_session(true);
4698        }
4699        // A session that delivered a frame resets the churn counter.
4700        run_session(false);
4701        for _ in 1..MAX_GOSSIP_NO_FRAME_SESSIONS {
4702            run_session(true);
4703        }
4704
4705        assert!(
4706            outbound.owns_connection(&peer_id, conn_id),
4707            "a reset churn counter must keep the reopen-gap claim"
4708        );
4709        assert!(matches!(
4710            sink.ordered_session_demand(conn_id, &peer_id, 0, ServicePeerDirection::Outbound),
4711            OrderedSessionDemand::OpenNow,
4712        ));
4713    }
4714
4715    #[tokio::test]
4716    async fn gossip_session_exit_keeps_connection_owned_across_reopen_gap() -> Result<(), BoxError>
4717    {
4718        let (inbound_tx, _inbound_rx) = mpsc::channel(1);
4719        let outbound = LegacyGossipOutbound::default();
4720        let sink = LegacyGossipSink {
4721            inbound_tx,
4722            outbound: outbound.clone(),
4723            trace: ZakuraTrace::noop(),
4724        };
4725        let peer_id = ZakuraPeerId::new(vec![97; 32]).expect("test peer id is within bounds");
4726        let conn_id = 5;
4727        let cancel_token = tokio_util::sync::CancellationToken::new();
4728        let (peer, peer_send) = legacy_gossip_peer_with_conn_and_session(
4729            peer_id.clone(),
4730            conn_id,
4731            1,
4732            cancel_token.clone(),
4733        );
4734        sink.add_peer(peer);
4735        assert!(sink.owns_connection_for_peer(&peer_id, conn_id));
4736
4737        // The remote closes only the gossip stream: the recv loop exits while
4738        // the connection stays up, and the transport backs off before reopening.
4739        drop(peer_send);
4740        tokio::time::timeout(Duration::from_secs(1), async {
4741            while outbound.contains(&peer_id) {
4742                tokio::time::sleep(Duration::from_millis(10)).await;
4743            }
4744        })
4745        .await
4746        .map_err(|_| -> BoxError { "gossip session exit did not retire the session".into() })?;
4747        assert!(
4748            sink.owns_connection_for_peer(&peer_id, conn_id),
4749            "connection must stay owned across the stream-reopen gap",
4750        );
4751
4752        // The transport reopens the stream: re-admission replaces the claim.
4753        let (reopened, _reopened_send) = legacy_gossip_peer_with_conn_and_session(
4754            peer_id.clone(),
4755            conn_id,
4756            2,
4757            cancel_token.clone(),
4758        );
4759        sink.add_peer(reopened);
4760        assert!(outbound.contains(&peer_id));
4761        assert!(sink.owns_connection_for_peer(&peer_id, conn_id));
4762
4763        // Connection close releases the session and any claim for good.
4764        sink.remove_peer(&peer_id, conn_id);
4765        assert!(!sink.owns_connection_for_peer(&peer_id, conn_id));
4766        Ok(())
4767    }
4768
4769    #[tokio::test]
4770    async fn local_legacy_gossip_exit_releases_connection_ownership() -> Result<(), BoxError> {
4771        let (inbound_tx, inbound_rx) = mpsc::channel(1);
4772        drop(inbound_rx);
4773        let outbound = LegacyGossipOutbound::default();
4774        let sink = LegacyGossipSink {
4775            inbound_tx,
4776            outbound: outbound.clone(),
4777            trace: ZakuraTrace::noop(),
4778        };
4779        let peer_id = ZakuraPeerId::new(vec![95; 32]).expect("test peer id is within bounds");
4780        let cancel_token = tokio_util::sync::CancellationToken::new();
4781        let (peer, peer_send) = legacy_gossip_peer(peer_id.clone(), cancel_token);
4782        sink.add_peer(peer);
4783        peer_send
4784            .send(LegacyGossipFrame::AdvertiseBlock(block_hash(95)).encode_frame()?)
4785            .await
4786            .map_err(|_| -> BoxError { "failed to send test gossip frame".into() })?;
4787
4788        tokio::time::timeout(Duration::from_secs(1), async {
4789            while outbound.contains(&peer_id) {
4790                tokio::time::sleep(Duration::from_millis(10)).await;
4791            }
4792        })
4793        .await
4794        .map_err(|_| -> BoxError { "local gossip exit kept stale ownership".into() })?;
4795        Ok(())
4796    }
4797
4798    #[tokio::test]
4799    async fn disconnected_peer_send_returns_error() -> Result<(), BoxError> {
4800        let peer_id = ZakuraPeerId::new(vec![9; 32]).expect("test peer id is within bounds");
4801        let (disconnected, rx) = framed_channel(1);
4802        drop(rx);
4803
4804        let error = send_to_sessions(
4805            vec![LegacyGossipPeerSession::new(peer_id, 0, 0, disconnected)],
4806            LegacyGossipFrame::AdvertiseBlock(block_hash(8)),
4807        )
4808        .expect_err("closed outbound queue reports an adapter error");
4809        assert!(
4810            error
4811                .to_string()
4812                .contains("ordered stream send queue is closed"),
4813            "unexpected error: {error}"
4814        );
4815        Ok(())
4816    }
4817
4818    #[tokio::test]
4819    async fn transient_inbound_failure_does_not_poison_first_seen_cache() -> Result<(), BoxError> {
4820        let (request_tx, mut request_rx) = tokio::sync::mpsc::unbounded_channel();
4821        let attempts = Arc::new(AtomicUsize::new(0));
4822        let (inbound_tx, inbound_rx) = mpsc::channel(4);
4823        let supervisor = ZakuraSupervisorHandle::new(1);
4824        let peer_id = ZakuraPeerId::new(vec![3; 32]).expect("test peer id is within bounds");
4825        let frame = LegacyGossipFrame::AdvertiseBlock(block_hash(5));
4826
4827        let worker = tokio::spawn(legacy_gossip_worker(
4828            FailsOnceRecorder {
4829                attempts: attempts.clone(),
4830                tx: request_tx,
4831            },
4832            inbound_rx,
4833            LegacyGossipForwarder::new(supervisor),
4834            ZakuraTrace::noop(),
4835        ));
4836
4837        inbound_tx
4838            .send(LegacyInboundWork::Gossip(LegacyGossipInbound {
4839                peer_id: peer_id.clone(),
4840                frame: frame.clone(),
4841            }))
4842            .await?;
4843        inbound_tx
4844            .send(LegacyInboundWork::Gossip(LegacyGossipInbound {
4845                peer_id: peer_id.clone(),
4846                frame: frame.clone(),
4847            }))
4848            .await?;
4849
4850        match recv_request(&mut request_rx).await? {
4851            Request::AdvertiseBlock(hash, Some(PeerSource::Zakura(source))) => {
4852                assert_eq!(hash, block_hash(5));
4853                assert_eq!(source, peer_id);
4854            }
4855            request => panic!("unexpected request: {request:?}"),
4856        }
4857        assert_eq!(attempts.load(Ordering::SeqCst), 2);
4858
4859        inbound_tx
4860            .send(LegacyInboundWork::Gossip(LegacyGossipInbound {
4861                peer_id,
4862                frame,
4863            }))
4864            .await?;
4865        tokio::time::sleep(Duration::from_millis(100)).await;
4866        assert_eq!(
4867            attempts.load(Ordering::SeqCst),
4868            2,
4869            "successful delivery marks the inventory seen"
4870        );
4871
4872        drop(inbound_tx);
4873        worker.await?;
4874        Ok(())
4875    }
4876
4877    #[test]
4878    fn inbound_queue_full_drops_without_rejecting_peer_but_malformed_rejects() {
4879        let (inbound_tx, _inbound_rx) = mpsc::channel(1);
4880        let sink = LegacyGossipSink {
4881            inbound_tx,
4882            outbound: LegacyGossipOutbound::default(),
4883            trace: ZakuraTrace::noop(),
4884        };
4885        let peer_id = ZakuraPeerId::new(vec![4; 32]).expect("test peer id is within bounds");
4886        let frame = LegacyGossipFrame::AdvertiseBlock(block_hash(6))
4887            .encode_frame()
4888            .expect("frame encodes");
4889
4890        assert!(sink
4891            .deliver(peer_id.clone(), ZAKURA_STREAM_GOSSIP, frame.clone())
4892            .is_ok());
4893        assert!(
4894            sink.deliver(peer_id.clone(), ZAKURA_STREAM_GOSSIP, frame)
4895                .is_ok(),
4896            "queue-full overload is shed without disconnecting the peer"
4897        );
4898        assert!(sink
4899            .deliver(
4900                peer_id,
4901                ZAKURA_STREAM_GOSSIP,
4902                Frame {
4903                    message_type: MSG_ADVERTISE_BLOCK,
4904                    flags: 1,
4905                    payload: Vec::new(),
4906                },
4907            )
4908            .is_err());
4909    }
4910
4911    #[test]
4912    fn tx_gossip_round_trips_legacy_and_witnessed_ids() {
4913        let frame =
4914            LegacyGossipFrame::AdvertiseTransactionIds(vec![legacy_tx_id(2), witnessed_tx_id(3)]);
4915
4916        assert_eq!(
4917            LegacyGossipFrame::decode_frame(frame.encode_frame().expect("frame encodes"))
4918                .expect("frame decodes"),
4919            frame
4920        );
4921    }
4922
4923    #[test]
4924    fn inventory_request_frames_round_trip_and_enforce_bounds() {
4925        let block_request = LegacyRequestFrame::BlocksByHash(vec![block_hash(1), block_hash(2)]);
4926        assert_eq!(
4927            LegacyRequestFrame::decode_frame(block_request.encode_frame().expect("frame encodes"))
4928                .expect("frame decodes"),
4929            block_request
4930        );
4931
4932        let tx_request =
4933            LegacyRequestFrame::TransactionsById(vec![legacy_tx_id(3), witnessed_tx_id(4)]);
4934        assert_eq!(
4935            LegacyRequestFrame::decode_frame(tx_request.encode_frame().expect("frame encodes"))
4936                .expect("frame decodes"),
4937            tx_request
4938        );
4939
4940        let find_blocks = LegacyRequestFrame::FindBlocks {
4941            known_blocks: vec![block_hash(5)],
4942            stop: Some(block_hash(6)),
4943        };
4944        assert_eq!(
4945            LegacyRequestFrame::decode_frame(find_blocks.encode_frame().expect("frame encodes"))
4946                .expect("frame decodes"),
4947            find_blocks
4948        );
4949
4950        let find_headers = LegacyRequestFrame::FindHeaders {
4951            known_blocks: vec![block_hash(7)],
4952            stop: None,
4953        };
4954        assert_eq!(
4955            LegacyRequestFrame::decode_frame(find_headers.encode_frame().expect("frame encodes"))
4956                .expect("frame decodes"),
4957            find_headers
4958        );
4959
4960        for request in [
4961            LegacyRequestFrame::MempoolTransactionIds,
4962            LegacyRequestFrame::Ping,
4963            LegacyRequestFrame::PushTransaction(UnminedTx::from(empty_v5_transaction(8))),
4964        ] {
4965            assert_eq!(
4966                LegacyRequestFrame::decode_frame(request.encode_frame().expect("frame encodes"))
4967                    .expect("frame decodes"),
4968                request
4969            );
4970        }
4971
4972        let max = usize::try_from(MAX_TX_INV_IN_SENT_MESSAGE).expect("test cap fits usize");
4973        let oversized = Frame {
4974            message_type: MSG_REQUEST_BLOCKS_BY_HASH,
4975            flags: 0,
4976            payload: encoded_count(max + 1),
4977        };
4978        assert!(matches!(
4979            LegacyRequestFrame::decode_frame(oversized),
4980            Err(LegacyGossipError::TooManyInventoryItems(count)) if count == max + 1
4981        ));
4982
4983        let max_locator =
4984            usize::try_from(MAX_BLOCK_LOCATOR_LENGTH).expect("test locator cap fits usize");
4985        let oversized_locator = Frame {
4986            message_type: MSG_REQUEST_FIND_HEADERS,
4987            flags: 0,
4988            payload: encoded_count(max_locator + 1),
4989        };
4990        assert!(matches!(
4991            LegacyRequestFrame::decode_frame(oversized_locator),
4992            Err(LegacyGossipError::TooManyBlockLocatorHashes(count)) if count == max_locator + 1
4993        ));
4994    }
4995
4996    #[test]
4997    fn response_codec_rejects_wrong_request_id_and_oversized_chunks() {
4998        let mut wrong_id_payload = Vec::new();
4999        wrong_id_payload.extend_from_slice(&2_u64.to_le_bytes());
5000        wrong_id_payload.push(1);
5001        let wrong_id = Frame {
5002            message_type: MSG_RESPONSE_BLOCK,
5003            flags: 0,
5004            payload: wrong_id_payload,
5005        };
5006        assert!(matches!(
5007            LegacyResponseCodec::decode_response(
5008                1,
5009                LegacyRequestKind::Blocks,
5010                vec![wrong_id],
5011                None
5012            ),
5013            Err(LegacyGossipError::WrongRequestId {
5014                expected: 1,
5015                actual: 2
5016            })
5017        ));
5018
5019        let mut oversized_payload = Vec::new();
5020        oversized_payload.extend_from_slice(&1_u64.to_le_bytes());
5021        oversized_payload.push(0);
5022        oversized_payload.resize(
5023            MAX_PROTOCOL_MESSAGE_LEN + RESPONSE_CHUNK_HEADER_BYTES + 1,
5024            0,
5025        );
5026        let oversized = Frame {
5027            message_type: MSG_RESPONSE_BLOCK,
5028            flags: 0,
5029            payload: oversized_payload,
5030        };
5031        assert!(matches!(
5032            LegacyResponseCodec::decode_response(
5033                1,
5034                LegacyRequestKind::Blocks,
5035                vec![oversized],
5036                None
5037            ),
5038            Err(LegacyGossipError::OversizedResponse(_))
5039        ));
5040    }
5041
5042    #[test]
5043    fn response_codec_rejects_empty_inventory_frame_sets() {
5044        for kind in [LegacyRequestKind::Blocks, LegacyRequestKind::Transactions] {
5045            assert!(matches!(
5046                LegacyResponseCodec::decode_response(1, kind, Vec::new(), None),
5047                Err(LegacyGossipError::MissingResponse(_))
5048            ));
5049        }
5050    }
5051
5052    #[tokio::test(start_paused = true)]
5053    async fn request_client_paces_stream_opens_below_the_connection_limit() {
5054        let client = ZakuraRequestClient::new_with_trace_and_stream_rate(
5055            ZakuraSupervisorHandle::new(1),
5056            ZakuraTrace::noop(),
5057            4,
5058        );
5059        let started = Instant::now();
5060
5061        client.wait_for_request_slot().await;
5062        assert_eq!(
5063            Instant::now(),
5064            started,
5065            "the first request uses the open slot"
5066        );
5067
5068        client.wait_for_request_slot().await;
5069        assert_eq!(
5070            Instant::now().duration_since(started),
5071            Duration::from_millis(500),
5072            "compatibility requests use half the configured stream-open rate"
5073        );
5074
5075        client.wait_for_request_slot().await;
5076        assert_eq!(
5077            Instant::now().duration_since(started),
5078            Duration::from_secs(1),
5079            "successive request slots stay evenly spaced"
5080        );
5081    }
5082
5083    /// The inbound service returns `Response::Nil` (a lone nil frame) for an empty
5084    /// FindBlocks/FindHeaders/MempoolTransactionIds result, so the codec must
5085    /// decode that into each kind's empty response rather than rejecting it.
5086    #[test]
5087    fn nil_response_decodes_as_each_kinds_empty_result() {
5088        let request_id = 7;
5089        let nil = || id_only_frame(MSG_RESPONSE_NIL, request_id);
5090
5091        assert_eq!(
5092            LegacyResponseCodec::decode_response(
5093                request_id,
5094                LegacyRequestKind::FindBlocks,
5095                vec![nil()],
5096                None,
5097            )
5098            .expect("nil is a valid empty FindBlocks response"),
5099            Response::BlockHashes(vec![]),
5100        );
5101        assert_eq!(
5102            LegacyResponseCodec::decode_response(
5103                request_id,
5104                LegacyRequestKind::FindHeaders,
5105                vec![nil()],
5106                None,
5107            )
5108            .expect("nil is a valid empty FindHeaders response"),
5109            Response::BlockHeaders(vec![]),
5110        );
5111        assert_eq!(
5112            LegacyResponseCodec::decode_response(
5113                request_id,
5114                LegacyRequestKind::MempoolTransactionIds,
5115                vec![nil()],
5116                None,
5117            )
5118            .expect("nil is a valid empty mempool response"),
5119            Response::TransactionIds(vec![]),
5120        );
5121        assert_eq!(
5122            LegacyResponseCodec::decode_response(
5123                request_id,
5124                LegacyRequestKind::PushTransaction,
5125                vec![nil()],
5126                None,
5127            )
5128            .expect("nil acknowledges a pushed transaction"),
5129            Response::Nil,
5130        );
5131
5132        // Inventory fetches and Ping must NOT accept a bare nil: for those it
5133        // means the peer has nothing, so the codec rejects it and the caller
5134        // falls back to another peer instead of returning an empty fetch.
5135        for kind in [
5136            LegacyRequestKind::Blocks,
5137            LegacyRequestKind::Transactions,
5138            LegacyRequestKind::Ping,
5139        ] {
5140            assert!(
5141                matches!(
5142                    LegacyResponseCodec::decode_response(request_id, kind, vec![nil()], None),
5143                    Err(LegacyGossipError::UnexpectedResponse(_)),
5144                ),
5145                "nil must be rejected for {kind:?} so the caller can fall back",
5146            );
5147        }
5148    }
5149
5150    #[test]
5151    fn response_codec_chunks_to_negotiated_frame_limit() -> Result<(), BoxError> {
5152        let block = Arc::new(Block::zcash_deserialize(
5153            BLOCK_TESTNET_141042_BYTES.as_slice(),
5154        )?);
5155        let max_frame_bytes =
5156            u32::try_from(FRAME_HEADER_BYTES + RESPONSE_CHUNK_HEADER_BYTES + 256)?;
5157        let frames = LegacyResponseCodec::encode_response(
5158            7,
5159            Response::Blocks(vec![InventoryResponse::Available((block.clone(), None))]),
5160            max_frame_bytes,
5161            max_frame_bytes,
5162        )?;
5163
5164        assert!(frames.len() > 1, "large block response must be chunked");
5165        for frame in &frames {
5166            frame.encode(max_frame_bytes)?;
5167        }
5168
5169        let response =
5170            LegacyResponseCodec::decode_response(7, LegacyRequestKind::Blocks, frames, None)?;
5171        let Response::Blocks(blocks) = response else {
5172            panic!("unexpected response: {response:?}");
5173        };
5174        assert!(matches!(
5175            blocks.as_slice(),
5176            [InventoryResponse::Available((received, None))] if received.hash() == block.hash()
5177        ));
5178
5179        Ok(())
5180    }
5181
5182    /// Regression test for `claude-outbound-write-ignores-message-cap` (legacy
5183    /// response chunking facet).
5184    ///
5185    /// The handshake clamps `max_frame_bytes` and `max_message_bytes`
5186    /// independently, so a peer can negotiate a message cap well below the frame
5187    /// cap. `push_chunked_response` sized chunk frames against the frame cap
5188    /// alone, so each chunk frame's payload could exceed the peer's accepted
5189    /// `max_message_bytes`. The request-response writer (`write_response_frame`)
5190    /// and the peer both reject such a frame as oversize, wasting the encode and
5191    /// causing avoidable disconnects/interop loss. Chunking must size against the
5192    /// effective cap `min(frame cap, message cap)` so every emitted chunk frame
5193    /// fits the negotiated message cap while still round-tripping.
5194    #[test]
5195    fn encode_response_chunks_respect_message_cap() -> Result<(), BoxError> {
5196        let block = Arc::new(Block::zcash_deserialize(
5197            BLOCK_TESTNET_141042_BYTES.as_slice(),
5198        )?);
5199
5200        // Large frame cap, small message cap: the divergence the handshake
5201        // permits. The message cap allows a 256-byte chunk payload plus the
5202        // per-chunk response header.
5203        let max_frame_bytes = u32::try_from(MAX_PROTOCOL_MESSAGE_LEN)?;
5204        let max_message_bytes = u32::try_from(RESPONSE_CHUNK_HEADER_BYTES + 256)?;
5205
5206        let frames = LegacyResponseCodec::encode_response(
5207            7,
5208            Response::Blocks(vec![InventoryResponse::Available((block.clone(), None))]),
5209            max_frame_bytes,
5210            max_message_bytes,
5211        )?;
5212
5213        assert!(
5214            frames.len() > 1,
5215            "a block larger than the negotiated message cap must be chunked, not emitted as one \
5216             over-cap frame"
5217        );
5218        for frame in &frames {
5219            assert!(
5220                frame.payload.len() <= max_message_bytes as usize,
5221                "chunk frame payload {} exceeds the negotiated max_message_bytes {}; the peer \
5222                 (and write_response_frame) would reject it as oversize",
5223                frame.payload.len(),
5224                max_message_bytes,
5225            );
5226            // Each chunk must still fit the frame cap so the transport encodes it.
5227            frame.encode(max_frame_bytes)?;
5228        }
5229
5230        // The smaller chunking must still round-trip back to the original block.
5231        let response =
5232            LegacyResponseCodec::decode_response(7, LegacyRequestKind::Blocks, frames, None)?;
5233        let Response::Blocks(blocks) = response else {
5234            panic!("unexpected response: {response:?}");
5235        };
5236        assert!(matches!(
5237            blocks.as_slice(),
5238            [InventoryResponse::Available((received, None))] if received.hash() == block.hash()
5239        ));
5240
5241        Ok(())
5242    }
5243
5244    /// A legacy block response must be bound to a hash we actually requested.
5245    ///
5246    /// Block responses are correlated only by request id and request kind, so
5247    /// without binding, a peer can substitute any other valid block for the one
5248    /// requested. The decoder must reject a delivered block whose hash is not in
5249    /// the requested set, and accept it when it is.
5250    #[test]
5251    fn decode_response_binds_blocks_to_requested_hashes() -> Result<(), BoxError> {
5252        let block = Arc::new(Block::zcash_deserialize(
5253            BLOCK_TESTNET_141042_BYTES.as_slice(),
5254        )?);
5255        let frame_cap = u32::try_from(MAX_PROTOCOL_MESSAGE_LEN)?;
5256        let frames = LegacyResponseCodec::encode_response(
5257            7,
5258            Response::Blocks(vec![InventoryResponse::Available((block.clone(), None))]),
5259            frame_cap,
5260            frame_cap,
5261        )?;
5262
5263        // A peer substitutes a real block we never asked for. Bound against a
5264        // requested-hash set that does not contain it, the codec rejects it
5265        // instead of correlating the response by request id and kind alone.
5266        let unrelated: HashSet<block::Hash> = std::iter::once(block_hash(99)).collect();
5267        assert!(
5268            matches!(
5269                LegacyResponseCodec::decode_response(
5270                    7,
5271                    LegacyRequestKind::Blocks,
5272                    frames.clone(),
5273                    Some(&unrelated),
5274                ),
5275                Err(LegacyGossipError::UnsolicitedBlock(_)),
5276            ),
5277            "a block whose hash was not requested must be rejected",
5278        );
5279
5280        // The same response is accepted when its hash is among those requested.
5281        let requested: HashSet<block::Hash> = std::iter::once(block.hash()).collect();
5282        let response = LegacyResponseCodec::decode_response(
5283            7,
5284            LegacyRequestKind::Blocks,
5285            frames,
5286            Some(&requested),
5287        )?;
5288        assert!(matches!(
5289            response,
5290            Response::Blocks(blocks)
5291                if matches!(
5292                    blocks.as_slice(),
5293                    [InventoryResponse::Available((received, None))]
5294                        if received.hash() == block.hash()
5295                )
5296        ));
5297
5298        Ok(())
5299    }
5300
5301    /// Regression test for `claude-legacy-responder-response-aggregation-unbounded`.
5302    ///
5303    /// An authenticated peer can name up to `MAX_TX_INV_IN_SENT_MESSAGE`
5304    /// block/transaction hashes on one request. Without a responder-side
5305    /// aggregate budget, `encode_response` serializes and retains the entire
5306    /// multi-frame `Vec<Frame>` for every available item before the first byte
5307    /// is written (worst case `MAX_TX_INV_IN_SENT_MESSAGE *
5308    /// MAX_PROTOCOL_MESSAGE_LEN`, tens of GiB). The outbound reader already
5309    /// enforces a symmetric `LegacyResponseBudget`; the responder must too, and
5310    /// must abort encoding early rather than buffering an over-budget response.
5311    #[test]
5312    fn encode_response_aborts_when_aggregation_exceeds_budget() -> Result<(), BoxError> {
5313        let block = Arc::new(Block::zcash_deserialize(
5314            BLOCK_TESTNET_141042_BYTES.as_slice(),
5315        )?);
5316        let block_bytes = block.zcash_serialize_to_vec()?.len();
5317        let frame_cap = u32::try_from(MAX_PROTOCOL_MESSAGE_LEN)?;
5318
5319        // A single available block is well within budget and still encodes.
5320        LegacyResponseCodec::encode_response(
5321            1,
5322            Response::Blocks(vec![InventoryResponse::Available((block.clone(), None))]),
5323            frame_cap,
5324            frame_cap,
5325        )?;
5326
5327        // Enough available blocks to overflow the cumulative byte budget must be
5328        // rejected, not fully materialized as a `Vec<Frame>`.
5329        let block_copies = LEGACY_RESPONSE_MAX_AGGREGATE_BYTES / block_bytes + 2;
5330        let many_blocks = Response::Blocks(
5331            (0..block_copies)
5332                .map(|_| InventoryResponse::Available((block.clone(), None)))
5333                .collect(),
5334        );
5335        assert!(
5336            matches!(
5337                LegacyResponseCodec::encode_response(2, many_blocks, frame_cap, frame_cap),
5338                Err(LegacyGossipError::ResponseAggregateBudget(_)),
5339            ),
5340            "an over-budget BlocksByHash response must abort encoding early",
5341        );
5342
5343        // The same cumulative budget guards the transaction responder path,
5344        // which shares `push_chunked_response`.
5345        let tx_copies = 2 * LEGACY_RESPONSE_MAX_AGGREGATE_BYTES / block_bytes + 2;
5346        let many_transactions = Response::Transactions(
5347            (0..tx_copies)
5348                .flat_map(|_| {
5349                    block
5350                        .transactions
5351                        .iter()
5352                        .map(|tx| InventoryResponse::Available((UnminedTx::from(tx.clone()), None)))
5353                })
5354                .collect(),
5355        );
5356        assert!(
5357            matches!(
5358                LegacyResponseCodec::encode_response(3, many_transactions, frame_cap, frame_cap),
5359                Err(LegacyGossipError::ResponseAggregateBudget(_)),
5360            ),
5361            "an over-budget TransactionsById response must abort encoding early",
5362        );
5363
5364        Ok(())
5365    }
5366
5367    #[test]
5368    fn response_codec_round_trips_chain_sync_and_mempool_responses() -> Result<(), BoxError> {
5369        let block = Arc::new(Block::zcash_deserialize(
5370            BLOCK_TESTNET_141042_BYTES.as_slice(),
5371        )?);
5372        let header = block::CountedHeader {
5373            header: block.header.clone(),
5374        };
5375
5376        let block_hash_response = Response::BlockHashes(vec![block.hash(), block_hash(10)]);
5377        let frames = LegacyResponseCodec::encode_response(
5378            8,
5379            block_hash_response.clone(),
5380            u32::try_from(MAX_PROTOCOL_MESSAGE_LEN)?,
5381            u32::try_from(MAX_PROTOCOL_MESSAGE_LEN)?,
5382        )?;
5383        assert_eq!(
5384            LegacyResponseCodec::decode_response(8, LegacyRequestKind::FindBlocks, frames, None)?,
5385            block_hash_response
5386        );
5387
5388        let header_response = Response::BlockHeaders(vec![header]);
5389        let frames = LegacyResponseCodec::encode_response(
5390            9,
5391            header_response.clone(),
5392            u32::try_from(MAX_PROTOCOL_MESSAGE_LEN)?,
5393            u32::try_from(MAX_PROTOCOL_MESSAGE_LEN)?,
5394        )?;
5395        assert_eq!(
5396            LegacyResponseCodec::decode_response(9, LegacyRequestKind::FindHeaders, frames, None)?,
5397            header_response
5398        );
5399
5400        let tx_ids_response = Response::TransactionIds(vec![legacy_tx_id(11), witnessed_tx_id(12)]);
5401        let frames = LegacyResponseCodec::encode_response(
5402            10,
5403            tx_ids_response.clone(),
5404            u32::try_from(MAX_PROTOCOL_MESSAGE_LEN)?,
5405            u32::try_from(MAX_PROTOCOL_MESSAGE_LEN)?,
5406        )?;
5407        assert_eq!(
5408            LegacyResponseCodec::decode_response(
5409                10,
5410                LegacyRequestKind::MempoolTransactionIds,
5411                frames,
5412                None,
5413            )?,
5414            tx_ids_response
5415        );
5416
5417        let frames = LegacyResponseCodec::encode_response(
5418            11,
5419            Response::Pong(Duration::ZERO),
5420            u32::try_from(MAX_PROTOCOL_MESSAGE_LEN)?,
5421            u32::try_from(MAX_PROTOCOL_MESSAGE_LEN)?,
5422        )?;
5423        assert!(matches!(
5424            LegacyResponseCodec::decode_response(11, LegacyRequestKind::Ping, frames, None)?,
5425            Response::Pong(_)
5426        ));
5427
5428        let frames = LegacyResponseCodec::encode_response(
5429            12,
5430            Response::Nil,
5431            u32::try_from(MAX_PROTOCOL_MESSAGE_LEN)?,
5432            u32::try_from(MAX_PROTOCOL_MESSAGE_LEN)?,
5433        )?;
5434        assert_eq!(
5435            LegacyResponseCodec::decode_response(
5436                12,
5437                LegacyRequestKind::PushTransaction,
5438                frames,
5439                None
5440            )?,
5441            Response::Nil
5442        );
5443
5444        Ok(())
5445    }
5446
5447    #[test]
5448    fn malformed_frames_are_rejected() {
5449        assert!(matches!(
5450            LegacyGossipFrame::decode_frame(Frame {
5451                message_type: MSG_ADVERTISE_BLOCK,
5452                flags: 1,
5453                payload: Vec::new(),
5454            }),
5455            Err(LegacyGossipError::UnsupportedFlags(1))
5456        ));
5457        assert!(matches!(
5458            LegacyGossipFrame::decode_frame(Frame {
5459                message_type: 99,
5460                flags: 0,
5461                payload: Vec::new(),
5462            }),
5463            Err(LegacyGossipError::UnknownMessageType(99))
5464        ));
5465        assert!(matches!(
5466            LegacyGossipFrame::decode_frame(Frame {
5467                message_type: MSG_ADVERTISE_TX_IDS,
5468                flags: 0,
5469                payload: vec![0],
5470            }),
5471            Err(LegacyGossipError::EmptyTransactionAdvertisement)
5472        ));
5473    }
5474
5475    #[test]
5476    fn tx_id_limits_are_enforced_before_allocation() {
5477        let max = usize::try_from(MAX_TX_INV_IN_SENT_MESSAGE).expect("test cap fits usize");
5478        let oversized = Frame {
5479            message_type: MSG_ADVERTISE_TX_IDS,
5480            flags: 0,
5481            payload: encoded_count(max + 1),
5482        };
5483        assert!(matches!(
5484            LegacyGossipFrame::decode_frame(oversized),
5485            Err(LegacyGossipError::TooManyInventoryItems(count)) if count == max + 1
5486        ));
5487
5488        let ids: HashSet<_> = (0..max + 3)
5489            .map(|index| {
5490                let mut bytes = [0; 32];
5491                bytes[..8].copy_from_slice(
5492                    &u64::try_from(index)
5493                        .expect("test index fits u64")
5494                        .to_le_bytes(),
5495                );
5496                UnminedTxId::from_legacy_id(transaction::Hash(bytes))
5497            })
5498            .collect();
5499        let frame = LegacyGossipFrame::from_request(Request::AdvertiseTransactionIds(ids, None))
5500            .expect("non-empty tx gossip converts");
5501        let LegacyGossipFrame::AdvertiseTransactionIds(ids) = frame else {
5502            panic!("expected tx id frame");
5503        };
5504        assert!(ids.len() <= max);
5505    }
5506
5507    #[test]
5508    fn unsupported_requests_fail_loudly() {
5509        let unsupported = [
5510            (
5511                Request::BlocksByHash(HashSet::from([block_hash(1)])),
5512                "BlocksByHash",
5513            ),
5514            (
5515                Request::TransactionsById(HashSet::from([legacy_tx_id(1)])),
5516                "TransactionsById",
5517            ),
5518            (
5519                Request::FindBlocks {
5520                    known_blocks: vec![block_hash(2)],
5521                    stop: None,
5522                },
5523                "FindBlocks",
5524            ),
5525            (
5526                Request::FindHeaders {
5527                    known_blocks: vec![block_hash(3)],
5528                    stop: None,
5529                },
5530                "FindHeaders",
5531            ),
5532            (
5533                Request::PushTransaction(
5534                    Transaction::V5 {
5535                        network_upgrade: NetworkUpgrade::Nu5,
5536                        lock_time: LockTime::min_lock_time_timestamp(),
5537                        expiry_height: block::Height(0),
5538                        inputs: Vec::new(),
5539                        outputs: Vec::new(),
5540                        sapling_shielded_data: None,
5541                        orchard_shielded_data: None,
5542                    }
5543                    .into(),
5544                    None,
5545                ),
5546                "PushTransaction",
5547            ),
5548            (Request::MempoolTransactionIds, "MempoolTransactionIds"),
5549            (Request::Peers, "Peers"),
5550        ];
5551
5552        for (request, command) in unsupported {
5553            let error =
5554                LegacyGossipFrame::from_request(request).expect_err("request is unsupported");
5555            assert!(matches!(
5556                error,
5557                LegacyGossipError::UnsupportedRequest(unsupported) if unsupported == command
5558            ));
5559        }
5560    }
5561
5562    #[test]
5563    fn legacy_peers_request_is_deferred_to_native_bootstrap_configuration() {
5564        let error =
5565            LegacyRequestFrame::from_request(Request::Peers).expect_err("Peers is deferred");
5566        assert!(matches!(
5567            error,
5568            LegacyGossipError::UnsupportedRequest("Peers")
5569        ));
5570        assert!(
5571            error
5572                .to_string()
5573                .contains("unsupported legacy gossip request"),
5574            "unexpected error: {error}"
5575        );
5576    }
5577
5578    #[tokio::test]
5579    async fn first_seen_cache_is_bounded_and_expires() {
5580        let cache = FirstSeenCache::new(1, Duration::from_millis(20));
5581        assert_eq!(
5582            cache
5583                .record_unseen([InventoryKey::Block(block_hash(1))])
5584                .await,
5585            vec![InventoryKey::Block(block_hash(1))]
5586        );
5587        assert!(cache
5588            .record_unseen([InventoryKey::Block(block_hash(1))])
5589            .await
5590            .is_empty());
5591        assert_eq!(
5592            cache
5593                .record_unseen([InventoryKey::Block(block_hash(2))])
5594                .await,
5595            vec![InventoryKey::Block(block_hash(2))]
5596        );
5597        assert_eq!(
5598            cache
5599                .record_unseen([InventoryKey::Block(block_hash(1))])
5600                .await,
5601            vec![InventoryKey::Block(block_hash(1))]
5602        );
5603
5604        let cache = FirstSeenCache::new(8, Duration::from_millis(1));
5605        assert_eq!(
5606            cache
5607                .record_unseen([InventoryKey::Block(block_hash(3))])
5608                .await,
5609            vec![InventoryKey::Block(block_hash(3))]
5610        );
5611        tokio::time::sleep(Duration::from_millis(5)).await;
5612        assert_eq!(
5613            cache
5614                .record_unseen([InventoryKey::Block(block_hash(3))])
5615                .await,
5616            vec![InventoryKey::Block(block_hash(3))]
5617        );
5618    }
5619
5620    #[tokio::test]
5621    async fn malformed_inbound_gossip_disconnects_peer() -> Result<(), BoxError> {
5622        let _guard = zakura_test::init();
5623        let (node, _rx) = legacy_node(31).await?;
5624        let hostile =
5625            HostilePeer::connect_native_with_capabilities(&node, 32, ZAKURA_CAP_LEGACY_GOSSIP)
5626                .await?;
5627        wait_registered_count(&node, 1).await?;
5628
5629        hostile.send_frame(ZAKURA_STREAM_GOSSIP, vec![1]).await?;
5630        wait_registered_count(&node, 0).await?;
5631
5632        hostile.shutdown().await;
5633        node.shutdown().await;
5634        Ok(())
5635    }
5636
5637    #[tokio::test]
5638    async fn gossip_stream_with_request_id_disconnects_peer() -> Result<(), BoxError> {
5639        let _guard = zakura_test::init();
5640        let node = ZakuraTestNode::builder(41).spawn().await?;
5641        let hostile = HostilePeer::connect_native(&node, 42).await?;
5642        wait_registered_count(&node, 1).await?;
5643
5644        let frame = LegacyGossipFrame::AdvertiseBlock(block_hash(77)).encode_frame()?;
5645        hostile
5646            .send_frame_with_request_id(ZAKURA_STREAM_GOSSIP, 99, frame)
5647            .await?;
5648        wait_registered_count(&node, 0).await?;
5649
5650        hostile.shutdown().await;
5651        node.shutdown().await;
5652        Ok(())
5653    }
5654
5655    #[tokio::test]
5656    async fn request_stream_without_request_id_disconnects_peer() -> Result<(), BoxError> {
5657        let _guard = zakura_test::init();
5658        let node = ZakuraTestNode::builder(71).spawn().await?;
5659        let hostile = HostilePeer::connect_native(&node, 72).await?;
5660        wait_registered_count(&node, 1).await?;
5661
5662        let frame = LegacyRequestFrame::BlocksByHash(vec![block_hash(1)]).encode_frame()?;
5663        hostile
5664            .send_frame(ZAKURA_STREAM_LEGACY_REQUESTS, frame.payload)
5665            .await?;
5666        wait_registered_count(&node, 0).await?;
5667
5668        hostile.shutdown().await;
5669        node.shutdown().await;
5670        Ok(())
5671    }
5672
5673    #[tokio::test]
5674    async fn malformed_outbound_response_disconnects_peer() -> Result<(), BoxError> {
5675        let _guard = zakura_test::init();
5676        let node = ZakuraTestNode::builder(73).spawn().await?;
5677        let hostile = HostilePeer::connect_native(&node, 74).await?;
5678        wait_registered_count(&node, 1).await?;
5679        let hostile_id = hostile.id()?;
5680
5681        let mut wrong_id_payload = Vec::new();
5682        wrong_id_payload.extend_from_slice(&0_u64.to_le_bytes());
5683        wrong_id_payload.push(1);
5684        let response = Frame {
5685            message_type: MSG_RESPONSE_BLOCK,
5686            flags: 0,
5687            payload: wrong_id_payload,
5688        };
5689
5690        let adapter = LegacyRequestAdapter::new(node.supervisor());
5691        let request = adapter.request_from_source(
5692            Request::BlocksByHash(HashSet::from([block_hash(1)])),
5693            Some(PeerSource::Zakura(hostile_id)),
5694        );
5695        let responder = hostile.respond_to_next_request(vec![response]);
5696        let (request_result, responder_result) = futures::join!(request, responder);
5697
5698        responder_result?;
5699        let error = request_result.expect_err("wrong request id is a request error");
5700        assert!(
5701            error
5702                .to_string()
5703                .contains("wrong legacy response request id"),
5704            "unexpected error: {error}"
5705        );
5706        wait_registered_count(&node, 0).await?;
5707
5708        hostile.shutdown().await;
5709        node.shutdown().await;
5710        Ok(())
5711    }
5712
5713    #[tokio::test]
5714    async fn outbound_request_timeout_releases_peer_for_next_request() -> Result<(), BoxError> {
5715        let _guard = zakura_test::init();
5716        let node = ZakuraTestNode::builder(79).spawn().await?;
5717        let mut hostile = HostilePeer::connect_native(&node, 80).await?;
5718        wait_registered_count(&node, 1).await?;
5719        let hostile_id = hostile.id()?;
5720
5721        let adapter =
5722            LegacyRequestAdapter::new_with_timeout(node.supervisor(), Duration::from_millis(100));
5723        let first_request = adapter.request_from_source(
5724            Request::BlocksByHash(HashSet::from([block_hash(1)])),
5725            Some(PeerSource::Zakura(hostile_id.clone())),
5726        );
5727        let hold_first = hostile.accept_next_request_without_response();
5728        let (request_result, hold_result) = futures::join!(first_request, hold_first);
5729
5730        hold_result?;
5731        let error = request_result.expect_err("silent peer should time out the request");
5732        assert!(
5733            error.to_string().contains("timed out"),
5734            "unexpected error: {error}"
5735        );
5736        wait_registered_count(&node, 1).await?;
5737
5738        let second_hash = block_hash(2);
5739        let second_request = adapter.request_from_source(
5740            Request::BlocksByHash(HashSet::from([second_hash])),
5741            Some(PeerSource::Zakura(hostile_id)),
5742        );
5743        let second_response = hostile.respond_to_next_request_with(|request_id| {
5744            vec![missing_blocks_frame(request_id, vec![second_hash])
5745                .expect("one missing block hash fits in a response frame")]
5746        });
5747        let (request_result, response_result) = futures::join!(second_request, second_response);
5748
5749        response_result?;
5750        assert_eq!(
5751            request_result?,
5752            Response::Blocks(vec![InventoryResponse::Missing(second_hash)])
5753        );
5754
5755        hostile.shutdown().await;
5756        node.shutdown().await;
5757        Ok(())
5758    }
5759
5760    #[tokio::test]
5761    async fn excessive_outbound_response_frames_disconnects_peer() -> Result<(), BoxError> {
5762        let _guard = zakura_test::init();
5763        let node = ZakuraTestNode::builder(75).spawn().await?;
5764        let hostile = HostilePeer::connect_native(&node, 76).await?;
5765        wait_registered_count(&node, 1).await?;
5766        let hostile_id = hostile.id()?;
5767
5768        let adapter = LegacyRequestAdapter::new(node.supervisor());
5769        let request = adapter.request_from_source(
5770            Request::BlocksByHash(HashSet::from([block_hash(1)])),
5771            Some(PeerSource::Zakura(hostile_id)),
5772        );
5773        let responder = hostile.respond_to_next_request_with(|request_id| {
5774            let mut frames = Vec::new();
5775            for _ in 0..10 {
5776                let mut payload = Vec::new();
5777                payload.extend_from_slice(&request_id.to_le_bytes());
5778                payload.extend_from_slice(&encoded_count(0));
5779                frames.push(Frame {
5780                    message_type: MSG_RESPONSE_MISSING_BLOCKS,
5781                    flags: 0,
5782                    payload,
5783                });
5784            }
5785            frames
5786        });
5787        let (request_result, responder_result) = futures::join!(request, responder);
5788
5789        responder_result?;
5790        let error = request_result.expect_err("excessive response frames are rejected");
5791        assert!(
5792            error
5793                .to_string()
5794                .contains("too many legacy response frames"),
5795            "unexpected error: {error}"
5796        );
5797        wait_registered_count(&node, 0).await?;
5798
5799        hostile.shutdown().await;
5800        node.shutdown().await;
5801        Ok(())
5802    }
5803
5804    #[tokio::test]
5805    async fn invalid_outbound_block_response_disconnects_peer() -> Result<(), BoxError> {
5806        let _guard = zakura_test::init();
5807        let node = ZakuraTestNode::builder(77).spawn().await?;
5808        let hostile = HostilePeer::connect_native(&node, 78).await?;
5809        wait_registered_count(&node, 1).await?;
5810        let hostile_id = hostile.id()?;
5811
5812        let adapter = LegacyRequestAdapter::new(node.supervisor());
5813        let request = adapter.request_from_source(
5814            Request::BlocksByHash(HashSet::from([block_hash(1)])),
5815            Some(PeerSource::Zakura(hostile_id)),
5816        );
5817        let responder = hostile.respond_to_next_request_with(|request_id| {
5818            let mut payload = Vec::new();
5819            payload.extend_from_slice(&request_id.to_le_bytes());
5820            payload.push(1);
5821            payload.extend_from_slice(b"not a serialized block");
5822            vec![Frame {
5823                message_type: MSG_RESPONSE_BLOCK,
5824                flags: 0,
5825                payload,
5826            }]
5827        });
5828        let (request_result, responder_result) = futures::join!(request, responder);
5829
5830        responder_result?;
5831        request_result.expect_err("invalid block bytes are rejected");
5832        wait_registered_count(&node, 0).await?;
5833
5834        hostile.shutdown().await;
5835        node.shutdown().await;
5836        Ok(())
5837    }
5838
5839    /// How the legacy stub answers inventory fetches in dual-stack tests.
5840    #[derive(Clone)]
5841    enum StubInventory {
5842        /// Return every requested id as missing (triggers Zakura fallback).
5843        Missing,
5844        /// Return the given transaction as available.
5845        Available(UnminedTx),
5846        /// Error on inventory (used to prove the legacy path is bypassed).
5847        Error,
5848    }
5849
5850    /// Stand-in for the legacy peer set inside a [`ZakuraDualStackService`]. It
5851    /// records every request it is handed and answers inventory per `inventory`.
5852    #[derive(Clone)]
5853    struct DualStackLegacyStub {
5854        tx: tokio::sync::mpsc::UnboundedSender<Request>,
5855        inventory: StubInventory,
5856    }
5857
5858    impl Service<Request> for DualStackLegacyStub {
5859        type Response = Response;
5860        type Error = BoxError;
5861        type Future = std::future::Ready<Result<Response, BoxError>>;
5862
5863        fn poll_ready(&mut self, _cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
5864            Poll::Ready(Ok(()))
5865        }
5866
5867        fn call(&mut self, request: Request) -> Self::Future {
5868            let _ = self.tx.send(request.clone());
5869            let response = match request {
5870                Request::TransactionsById(ids) | Request::TransactionsByIdFrom { ids, .. } => {
5871                    match &self.inventory {
5872                        StubInventory::Error => {
5873                            return std::future::ready(Err("legacy inventory unavailable".into()))
5874                        }
5875                        StubInventory::Missing => Response::Transactions(
5876                            ids.into_iter().map(InventoryResponse::Missing).collect(),
5877                        ),
5878                        StubInventory::Available(tx) => Response::Transactions(
5879                            ids.into_iter()
5880                                .map(|id| {
5881                                    if id == tx.id() {
5882                                        InventoryResponse::Available((tx.clone(), None))
5883                                    } else {
5884                                        InventoryResponse::Missing(id)
5885                                    }
5886                                })
5887                                .collect(),
5888                        ),
5889                    }
5890                }
5891                Request::BlocksByHash(hashes) | Request::BlocksByHashFrom { hashes, .. } => {
5892                    if let StubInventory::Error = self.inventory {
5893                        return std::future::ready(Err("legacy inventory unavailable".into()));
5894                    }
5895                    Response::Blocks(hashes.into_iter().map(InventoryResponse::Missing).collect())
5896                }
5897                _ => Response::Nil,
5898            };
5899            std::future::ready(Ok(response))
5900        }
5901    }
5902
5903    #[tokio::test]
5904    async fn dual_stack_advertise_fans_out_to_legacy_and_zakura() -> Result<(), BoxError> {
5905        let _guard = zakura_test::init();
5906        // node_a originates; node_b records gossip delivered over Zakura.
5907        let (node_a, _rx_a) = legacy_node(201).await?;
5908        let (node_b, mut rx_b) = legacy_node(202).await?;
5909        node_a.connect_native(&node_b, TEST_NET_TIMEOUT).await?;
5910        let a_peer_id = node_peer_id(&node_a).await?;
5911
5912        let (legacy_tx, mut rx_legacy) = tokio::sync::mpsc::unbounded_channel();
5913        let mut composite = ZakuraDualStackService::new(
5914            DualStackLegacyStub {
5915                tx: legacy_tx,
5916                inventory: StubInventory::Missing,
5917            },
5918            node_a.supervisor(),
5919            true,
5920        );
5921
5922        let hash = block_hash(9);
5923        composite
5924            .ready()
5925            .await?
5926            .call(Request::AdvertiseBlockToAll(hash))
5927            .await?;
5928
5929        // Legacy peer set received the advertisement verbatim.
5930        match recv_request(&mut rx_legacy).await? {
5931            Request::AdvertiseBlockToAll(received) => assert_eq!(received, hash),
5932            request => panic!("unexpected legacy request: {request:?}"),
5933        }
5934        // Zakura delivered it to node_b, attributed to node_a.
5935        match recv_request(&mut rx_b).await? {
5936            Request::AdvertiseBlock(received, Some(PeerSource::Zakura(peer_id))) => {
5937                assert_eq!(received, hash);
5938                assert_eq!(peer_id, a_peer_id);
5939            }
5940            request => panic!("unexpected zakura request: {request:?}"),
5941        }
5942
5943        node_a.shutdown().await;
5944        node_b.shutdown().await;
5945        Ok(())
5946    }
5947
5948    #[tokio::test]
5949    async fn dual_stack_inventory_falls_back_to_zakura_when_legacy_missing() -> Result<(), BoxError>
5950    {
5951        let _guard = zakura_test::init();
5952        let transaction = UnminedTx::from(empty_v5_transaction(7));
5953        let advertiser = inventory_node(203, transaction.clone()).await?;
5954        let requester = ZakuraTestNode::builder(204).spawn().await?;
5955        requester
5956            .connect_native(&advertiser, TEST_NET_TIMEOUT)
5957            .await?;
5958
5959        let (legacy_tx, _rx_legacy) = tokio::sync::mpsc::unbounded_channel();
5960        let mut composite = ZakuraDualStackService::new(
5961            DualStackLegacyStub {
5962                tx: legacy_tx,
5963                inventory: StubInventory::Missing,
5964            },
5965            requester.supervisor(),
5966            true,
5967        );
5968
5969        let response = composite
5970            .ready()
5971            .await?
5972            .call(Request::TransactionsById(HashSet::from([transaction.id()])))
5973            .await?;
5974        match response {
5975            Response::Transactions(items) => {
5976                assert_eq!(items.len(), 1);
5977                assert!(matches!(items[0], InventoryResponse::Available(_)));
5978            }
5979            other => panic!("unexpected response: {other:?}"),
5980        }
5981
5982        advertiser.shutdown().await;
5983        requester.shutdown().await;
5984        Ok(())
5985    }
5986
5987    #[tokio::test]
5988    async fn dual_stack_inventory_uses_legacy_when_available() -> Result<(), BoxError> {
5989        let _guard = zakura_test::init();
5990        let transaction = UnminedTx::from(empty_v5_transaction(8));
5991        // No Zakura advertiser is connected, so any fallback would fail; the
5992        // legacy-available short-circuit is what makes this succeed.
5993        let requester = ZakuraTestNode::builder(205).spawn().await?;
5994
5995        let (legacy_tx, _rx_legacy) = tokio::sync::mpsc::unbounded_channel();
5996        let mut composite = ZakuraDualStackService::new(
5997            DualStackLegacyStub {
5998                tx: legacy_tx,
5999                inventory: StubInventory::Available(transaction.clone()),
6000            },
6001            requester.supervisor(),
6002            true,
6003        );
6004
6005        let response = composite
6006            .ready()
6007            .await?
6008            .call(Request::TransactionsById(HashSet::from([transaction.id()])))
6009            .await?;
6010        match response {
6011            Response::Transactions(items) => {
6012                assert!(matches!(items[0], InventoryResponse::Available(_)));
6013            }
6014            other => panic!("unexpected response: {other:?}"),
6015        }
6016
6017        requester.shutdown().await;
6018        Ok(())
6019    }
6020
6021    #[tokio::test]
6022    async fn dual_stack_zakura_only_bypasses_legacy_for_inventory() -> Result<(), BoxError> {
6023        let _guard = zakura_test::init();
6024        let transaction = UnminedTx::from(empty_v5_transaction(9));
6025        let advertiser = inventory_node(206, transaction.clone()).await?;
6026        let requester = ZakuraTestNode::builder(207).spawn().await?;
6027        requester
6028            .connect_native(&advertiser, TEST_NET_TIMEOUT)
6029            .await?;
6030
6031        // legacy_enabled = false; the stub errors if it is ever consulted.
6032        let (legacy_tx, mut rx_legacy) = tokio::sync::mpsc::unbounded_channel();
6033        let mut composite = ZakuraDualStackService::new(
6034            DualStackLegacyStub {
6035                tx: legacy_tx,
6036                inventory: StubInventory::Error,
6037            },
6038            requester.supervisor(),
6039            false,
6040        );
6041
6042        let response = composite
6043            .ready()
6044            .await?
6045            .call(Request::TransactionsById(HashSet::from([transaction.id()])))
6046            .await?;
6047        match response {
6048            Response::Transactions(items) => {
6049                assert!(matches!(items[0], InventoryResponse::Available(_)));
6050            }
6051            other => panic!("unexpected response: {other:?}"),
6052        }
6053        // The legacy peer set was never consulted in Zakura-only mode.
6054        assert!(rx_legacy.try_recv().is_err());
6055
6056        advertiser.shutdown().await;
6057        requester.shutdown().await;
6058        Ok(())
6059    }
6060
6061    #[tokio::test]
6062    async fn dual_stack_zakura_only_routes_chain_sync_and_mempool_requests() -> Result<(), BoxError>
6063    {
6064        let _guard = zakura_test::init();
6065        let block = Arc::new(Block::zcash_deserialize(
6066            BLOCK_TESTNET_141042_BYTES.as_slice(),
6067        )?);
6068        let transaction = UnminedTx::from(empty_v5_transaction(33));
6069        let (responder, mut pushed_rx) =
6070            normal_network_node(220, block.clone(), transaction.clone()).await?;
6071        let requester = ZakuraTestNode::builder(221).spawn().await?;
6072        requester
6073            .connect_native(&responder, TEST_NET_TIMEOUT)
6074            .await?;
6075
6076        // legacy_enabled = false; the stub errors if the legacy peer set is ever
6077        // consulted. Chain-sync discovery and mempool data used to be routed to
6078        // the legacy peer set (`Passthrough`), so a node whose only peer is over
6079        // Zakura could never obtain tips, fetch blocks, or push transactions.
6080        let (legacy_tx, mut rx_legacy) = tokio::sync::mpsc::unbounded_channel();
6081        let mut composite = ZakuraDualStackService::new(
6082            DualStackLegacyStub {
6083                tx: legacy_tx,
6084                inventory: StubInventory::Error,
6085            },
6086            requester.supervisor(),
6087            false,
6088        );
6089
6090        let find_blocks = composite
6091            .ready()
6092            .await?
6093            .call(Request::FindBlocks {
6094                known_blocks: vec![block_hash(1)],
6095                stop: Some(block.hash()),
6096            })
6097            .await?;
6098        match find_blocks {
6099            Response::BlockHashes(hashes) => assert_eq!(hashes, vec![block.hash()]),
6100            other => panic!("unexpected FindBlocks response: {other:?}"),
6101        }
6102
6103        let find_headers = composite
6104            .ready()
6105            .await?
6106            .call(Request::FindHeaders {
6107                known_blocks: vec![block_hash(1)],
6108                stop: Some(block.hash()),
6109            })
6110            .await?;
6111        assert!(
6112            matches!(find_headers, Response::BlockHeaders(ref headers) if headers.len() == 1),
6113            "FindHeaders should be served over Zakura, got {find_headers:?}",
6114        );
6115
6116        let mempool_ids = composite
6117            .ready()
6118            .await?
6119            .call(Request::MempoolTransactionIds)
6120            .await?;
6121        assert!(
6122            matches!(mempool_ids, Response::TransactionIds(ref ids) if *ids == vec![transaction.id()]),
6123            "MempoolTransactionIds should be served over Zakura, got {mempool_ids:?}",
6124        );
6125
6126        composite
6127            .ready()
6128            .await?
6129            .call(Request::PushTransaction(transaction.clone(), None))
6130            .await?;
6131        let pushed = pushed_rx
6132            .recv()
6133            .await
6134            .expect("transaction pushed over Zakura");
6135        assert_eq!(pushed, transaction.id());
6136
6137        // None of these requests consulted the legacy peer set.
6138        assert!(
6139            rx_legacy.try_recv().is_err(),
6140            "Zakura-only mode must not consult the legacy peer set for serviceable requests",
6141        );
6142
6143        responder.shutdown().await;
6144        requester.shutdown().await;
6145        Ok(())
6146    }
6147
6148    #[tokio::test]
6149    async fn dual_stack_passthrough_request_reaches_legacy() -> Result<(), BoxError> {
6150        let _guard = zakura_test::init();
6151        let node = ZakuraTestNode::builder(208).spawn().await?;
6152        let (legacy_tx, mut rx_legacy) = tokio::sync::mpsc::unbounded_channel();
6153        let mut composite = ZakuraDualStackService::new(
6154            DualStackLegacyStub {
6155                tx: legacy_tx,
6156                inventory: StubInventory::Missing,
6157            },
6158            node.supervisor(),
6159            true,
6160        );
6161
6162        composite.ready().await?.call(Request::Peers).await?;
6163        match recv_request(&mut rx_legacy).await? {
6164            Request::Peers => {}
6165            request => panic!("unexpected passthrough request: {request:?}"),
6166        }
6167
6168        node.shutdown().await;
6169        Ok(())
6170    }
6171}