1use 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
53pub const ZAKURA_STREAM_GOSSIP: u16 = 2;
55pub const ZAKURA_STREAM_LEGACY_REQUESTS: u16 = 3;
57pub const LEGACY_GOSSIP_VERSION: u16 = 1;
59pub const MSG_ADVERTISE_BLOCK: u16 = 1;
61pub const MSG_ADVERTISE_TX_IDS: u16 = 2;
63pub const MSG_REQUEST_BLOCKS_BY_HASH: u16 = 3;
65pub const MSG_REQUEST_TRANSACTIONS_BY_ID: u16 = 4;
67pub const MSG_RESPONSE_BLOCK: u16 = 5;
69pub const MSG_RESPONSE_TRANSACTION: u16 = 6;
71pub const MSG_RESPONSE_MISSING_BLOCKS: u16 = 7;
73pub const MSG_RESPONSE_MISSING_TRANSACTIONS: u16 = 8;
75pub const MSG_REQUEST_FIND_BLOCKS: u16 = 9;
77pub const MSG_REQUEST_FIND_HEADERS: u16 = 10;
79pub const MSG_REQUEST_MEMPOOL_TRANSACTION_IDS: u16 = 11;
81pub const MSG_REQUEST_PING: u16 = 12;
83pub const MSG_REQUEST_PUSH_TRANSACTION: u16 = 13;
85pub const MSG_RESPONSE_BLOCK_HASHES: u16 = 14;
87pub const MSG_RESPONSE_BLOCK_HEADERS: u16 = 15;
89pub const MSG_RESPONSE_TRANSACTION_IDS: u16 = 16;
91pub const MSG_RESPONSE_PONG: u16 = 17;
93pub 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;
101const LEGACY_GOSSIP_EXPENSIVE_ATTEMPT: Duration = Duration::from_secs(1);
110const 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);
117const LEGACY_REQUEST_STREAM_RATE_DIVISOR: u32 = 2;
120const DUAL_STACK_LEGACY_INVENTORY_TIMEOUT: Duration = Duration::from_secs(3);
125const LEGACY_RESPONSE_CHUNK_BYTES: usize = 512 * 1024;
126const 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
158pub(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#[derive(Clone, Debug, Eq, PartialEq)]
169pub enum LegacyGossipFrame {
170 AdvertiseBlock(block::Hash),
172 AdvertiseTransactionIds(Vec<UnminedTxId>),
174}
175
176impl LegacyGossipFrame {
177 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 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 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 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#[derive(Clone, Debug, Eq, PartialEq)]
262pub enum LegacyRequestFrame {
263 BlocksByHash(Vec<block::Hash>),
265 TransactionsById(Vec<UnminedTxId>),
267 FindBlocks {
269 known_blocks: Vec<block::Hash>,
271 stop: Option<block::Hash>,
273 },
274 FindHeaders {
276 known_blocks: Vec<block::Hash>,
278 stop: Option<block::Hash>,
280 },
281 MempoolTransactionIds,
283 Ping,
285 PushTransaction(UnminedTx),
287}
288
289impl LegacyRequestFrame {
290 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 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 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 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 push_response_frame(
568 &mut frames,
569 &mut budget,
570 block_hashes_frame(request_id, hashes)?,
571 )?;
572 }
573 Response::BlockHeaders(headers) => {
574 push_response_frame(
576 &mut frames,
577 &mut budget,
578 block_headers_frame(request_id, headers)?,
579 )?;
580 }
581 Response::TransactionIds(ids) => {
582 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 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 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 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 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 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#[derive(Default)]
946struct ResponseEncodeBudget {
947 bytes: usize,
948}
949
950impl ResponseEncodeBudget {
951 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
963fn 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 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 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 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 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#[derive(Clone, Debug)]
1255pub struct ZakuraGossipBroadcast {
1256 first_seen: FirstSeenCache,
1257 outbound: LegacyGossipOutbound,
1258}
1259
1260impl ZakuraGossipBroadcast {
1261 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#[derive(Clone, Copy, Debug)]
1296struct GossipGapClaim {
1297 conn_id: ZakuraConnId,
1298}
1299
1300const MAX_GOSSIP_NO_FRAME_SESSIONS: u32 = 8;
1305
1306#[derive(Clone, Copy, Debug)]
1310struct GossipSessionChurn {
1311 conn_id: ZakuraConnId,
1312 no_frame_exits: u32,
1313 retired: bool,
1314}
1315
1316#[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 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 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 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 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 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 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#[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 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 pub fn try_send_advertise_block(&self, hash: block::Hash) -> Result<(), OrderedSendError> {
1608 self.try_send_gossip_frame(LegacyGossipFrame::AdvertiseBlock(hash))
1609 }
1610
1611 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#[derive(Clone, Debug)]
1687pub struct LegacyGossipAdapter {
1688 broadcast: ZakuraGossipBroadcast,
1689}
1690
1691impl LegacyGossipAdapter {
1692 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#[derive(Clone, Debug)]
1728pub struct LegacyRequestAdapter {
1729 client: ZakuraRequestClient,
1730}
1731
1732impl LegacyRequestAdapter {
1733 pub fn new(supervisor: ZakuraSupervisorHandle) -> Self {
1735 Self {
1736 client: ZakuraRequestClient::new(supervisor),
1737 }
1738 }
1739
1740 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 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#[derive(Copy, Clone)]
1801enum DualStackRoute {
1802 Advertise,
1804 LegacyFirstThenZakura,
1812 Passthrough,
1815}
1816
1817#[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 #[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 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 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 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 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 Err(_) => request_adapter.call(request).await,
1976 }
1977 }
1978 DualStackRoute::Passthrough => legacy.ready().await?.call(request).await,
1979 }
1980 })
1981 }
1982}
1983
1984#[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 pub fn new(supervisor: ZakuraSupervisorHandle) -> Self {
1997 Self::new_with_trace(supervisor, ZakuraTrace::noop())
1998 }
1999
2000 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 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 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 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 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 async fn fresh_attempt(&self, frame: &LegacyGossipFrame) -> Option<LegacyGossipFrame> {
2318 self.attempt_cooldown.unseen(frame).await
2319 }
2320
2321 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#[derive(Debug)]
2345pub struct LegacyGossipSink {
2346 inbound_tx: mpsc::Sender<LegacyInboundWork>,
2347 outbound: LegacyGossipOutbound,
2348 trace: ZakuraTrace,
2349}
2350
2351impl LegacyGossipSink {
2352 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 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 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 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 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 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 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#[derive(Debug, Error)]
2999pub enum LegacyGossipError {
3000 #[error("unsupported legacy gossip request: {0}")]
3002 UnsupportedRequest(&'static str),
3003 #[error("unsupported legacy gossip flags: {0}")]
3005 UnsupportedFlags(u16),
3006 #[error("unknown legacy gossip message type: {0}")]
3008 UnknownMessageType(u16),
3009 #[error("empty transaction-id advertisement")]
3011 EmptyTransactionAdvertisement,
3012 #[error("too many inventory items in legacy request/response: {0}")]
3014 TooManyInventoryItems(usize),
3015 #[error("too many block locator hashes in legacy request: {0}")]
3017 TooManyBlockLocatorHashes(usize),
3018 #[error("too many block headers in legacy response: {0}")]
3020 TooManyHeaders(usize),
3021 #[error("transaction-id advertisement contained non-transaction inventory")]
3023 NonTransactionInventory,
3024 #[error("trailing bytes after legacy gossip payload")]
3026 TrailingBytes,
3027 #[error("wrong legacy request id in response: expected {expected}, got {actual}")]
3029 WrongRequestId {
3030 expected: u64,
3032 actual: u64,
3034 },
3035 #[error("incomplete legacy response chunk")]
3037 IncompleteResponseChunk,
3038 #[error("truncated legacy response")]
3040 TruncatedResponse,
3041 #[error("oversized legacy response: {0} bytes")]
3043 OversizedResponse(usize),
3044 #[error("legacy response exceeded responder aggregate budget: {0} bytes")]
3046 ResponseAggregateBudget(usize),
3047 #[error("unexpected legacy response: {0}")]
3049 UnexpectedResponse(&'static str),
3050 #[error("missing legacy response: {0}")]
3052 MissingResponse(&'static str),
3053 #[error("legacy block response contained an unrequested block: {0:?}")]
3056 UnsolicitedBlock(block::Hash),
3057 #[error(transparent)]
3059 Serialization(#[from] SerializationError),
3060 #[error(transparent)]
3062 Io(#[from] std::io::Error),
3063 #[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 #[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 #[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 .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 #[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 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 tokio::task::yield_now().await;
3743 drop(response_rx);
3744
3745 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 #[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 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 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 assert!(outbound.insert(LegacyGossipPeerSession::new(
4583 peer_id.clone(),
4584 conn_id,
4585 2,
4586 second_send
4587 )));
4588
4589 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 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 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 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 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 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 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 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 #[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 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 #[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 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 frame.encode(max_frame_bytes)?;
5228 }
5229
5230 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 #[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 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 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 #[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 LegacyResponseCodec::encode_response(
5321 1,
5322 Response::Blocks(vec![InventoryResponse::Available((block.clone(), None))]),
5323 frame_cap,
5324 frame_cap,
5325 )?;
5326
5327 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 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 #[derive(Clone)]
5841 enum StubInventory {
5842 Missing,
5844 Available(UnminedTx),
5846 Error,
5848 }
5849
5850 #[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 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 match recv_request(&mut rx_legacy).await? {
5931 Request::AdvertiseBlockToAll(received) => assert_eq!(received, hash),
5932 request => panic!("unexpected legacy request: {request:?}"),
5933 }
5934 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 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 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 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 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 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}