1use {
5 crate::{
6 admin_rpc_post_init::{KeyUpdaterType, KeyUpdaters},
7 banking_trace::BankingTracer,
8 block_creation_loop::ReplayHighestFrozen,
9 cluster_info_vote_listener::{
10 DuplicateConfirmedSlotsReceiver, GossipVerifiedVoteHashReceiver, VoteTracker,
11 },
12 cluster_slots_service::{ClusterSlotsService, cluster_slots::ClusterSlots},
13 commitment_service::AggregateCommitmentService,
14 completed_data_sets_service::CompletedDataSetsSender,
15 consensus::{Tower, tower_storage::TowerStorage},
16 cost_update_service::CostUpdateService,
17 drop_bank_service::DropBankService,
18 epoch_specs::EpochSpecs,
19 repair::{
20 block_id_repair_service::BlockIdRepairChannels,
21 repair_service::{OutstandingShredRepairs, RepairInfo, RepairServiceChannels},
22 },
23 replay_stage::{ReplayReceivers, ReplaySenders, ReplayStage, ReplayStageConfig},
24 shred_fetch_stage::{SHRED_FETCH_CHANNEL_SIZE, ShredFetchStage},
25 voting_service::VotingService,
26 warm_quic_cache_service::WarmQuicCacheService,
27 window_service::{WindowService, WindowServiceChannels},
28 },
29 agave_bls_sigverify::{
30 bls_sigverifier::{self, SigVerifierChannels, SigVerifierContext},
31 generated_cert_types::GeneratedCertTypes,
32 },
33 agave_votor::{
34 event::{LatestSwitchRequest, LeaderWindowInfo, VotorEventReceiver, VotorEventSender},
35 vote_history::VoteHistory,
36 vote_history_storage::VoteHistoryStorage,
37 voting_service::{
38 VOTOR_RATE_LIMIT_PPS, VotingService as BLSVotingService, VotingServiceOverride,
39 },
40 votor::{Votor, VotorConfig},
41 },
42 agave_votor_messages::{
43 VerifiedVoterSlotsReceiver, VerifiedVoterSlotsSender, consensus_message::Block,
44 metric_types::MAX_IN_FLIGHT_CONSENSUS_EVENTS, reward_certificate::AddVoteMessage,
45 },
46 crossbeam_channel::{Receiver, Sender, bounded, unbounded},
47 solana_client::connection_cache::ConnectionCache,
48 solana_clock::Slot,
49 solana_geyser_plugin_manager::block_metadata_notifier_interface::BlockMetadataNotifierArc,
50 solana_gossip::{
51 cluster_info::ClusterInfo, duplicate_shred_handler::DuplicateShredHandler,
52 duplicate_shred_listener::DuplicateShredListener,
53 },
54 solana_keypair::Keypair,
55 solana_ledger::{
56 blockstore::{Blockstore, MAX_COMPLETED_SLOTS_IN_CHANNEL, UpdateParentReceiver},
57 blockstore_cleanup_service::BlockstoreCleanupService,
58 blockstore_processor::TransactionStatusSender,
59 entry_notifier_service::EntryNotifierSender,
60 leader_schedule_cache::LeaderScheduleCache,
61 shred::filter::TurbineMode,
62 },
63 solana_net_utils::PinnedXdpSender,
64 solana_poh::{poh_controller::PohController, poh_recorder::PohRecorder},
65 solana_pubkey::Pubkey,
66 solana_rpc::{
67 max_slots::MaxSlots, optimistically_confirmed_bank_tracker::BankNotificationSenderConfig,
68 rpc_subscriptions::RpcSubscriptions, slot_status_notifier::SlotStatusNotifier,
69 },
70 solana_runtime::{
71 bank::MAX_ALPENGLOW_VOTE_ACCOUNTS,
72 bank_forks::BankForks,
73 bank_forks_controller::{BankForksCommandReceiver, BankForksController},
74 commitment::BlockCommitmentCache,
75 prioritization_fee_cache::PrioritizationFeeCache,
76 snapshot_controller::SnapshotController,
77 validated_block_finalization::ValidatedBlockFinalizationCert,
78 vote_sender_types::ReplayVoteSender,
79 },
80 solana_streamer::{
81 evicting_sender::EvictingSender,
82 nonblocking::simple_qos::SimpleQosConfig,
83 quic::{QuicStreamerConfig, SpawnServerResult, spawn_simple_qos_server},
84 streamer::StakedNodes,
85 },
86 solana_turbine::{XdpSender as TurbineXdpSender, retransmit_stage::RetransmitStage},
87 std::{
88 collections::HashSet,
89 net::UdpSocket,
90 num::NonZeroUsize,
91 sync::{Arc, RwLock, atomic::AtomicBool},
92 thread::{self, JoinHandle},
93 },
94 tokio_util::sync::CancellationToken,
95};
96
97const CHANNEL_SIZE_RETRANSMIT_INGRESS: usize = 16 * 1024;
104
105pub(crate) const MAX_ALPENGLOW_PACKET_NUM: usize = 10_000;
107const MAX_BLS_MESSAGES_TO_SEND: usize = 1000;
111
112pub struct Tvu {
113 fetch_stage: ShredFetchStage,
114 shred_sigverify: JoinHandle<()>,
115 retransmit_stage: RetransmitStage,
116 window_service: WindowService,
117 cluster_slots_service: ClusterSlotsService,
118 replay_stage: ReplayStage,
119 blockstore_cleanup_service: BlockstoreCleanupService,
120 cost_update_service: CostUpdateService,
121 voting_service: VotingService,
122 bls_voting_service: BLSVotingService,
123 warm_quic_cache_service: Option<WarmQuicCacheService>,
124 drop_bank_service: DropBankService,
125 duplicate_shred_listener: DuplicateShredListener,
126 bls_sigverify_threads: (JoinHandle<()>, JoinHandle<()>),
127 votor: Votor,
128 commitment_service: AggregateCommitmentService,
129}
130
131pub struct TvuSockets {
132 pub fetch: Vec<UdpSocket>,
133 pub repair: UdpSocket,
134 pub retransmit: Vec<UdpSocket>,
135 pub ancestor_hashes_requests: UdpSocket,
136 pub alpenglow: UdpSocket,
137 pub block_id_repair: UdpSocket,
138}
139
140pub struct TvuConfig {
141 pub max_ledger_shreds: Option<u64>,
142 pub shred_version: u16,
143 pub repair_validators: Option<HashSet<Pubkey>>,
145 pub repair_whitelist: Arc<RwLock<HashSet<Pubkey>>>,
147 pub wait_for_vote_to_start_leader: bool,
148 pub replay_forks_threads: NonZeroUsize,
149 pub replay_transactions_threads: NonZeroUsize,
150 pub shred_sigverify_threads: NonZeroUsize,
151 pub bls_sigverify_threads: NonZeroUsize,
152 pub turbine_xdp_sender: Option<TurbineXdpSender>,
153 pub repair_xdp_sender: Option<PinnedXdpSender>,
154}
155
156impl Default for TvuConfig {
157 fn default() -> Self {
158 Self {
159 max_ledger_shreds: None,
160 shred_version: 0,
161 repair_validators: None,
162 repair_whitelist: Arc::new(RwLock::new(HashSet::default())),
163 wait_for_vote_to_start_leader: false,
164 replay_forks_threads: NonZeroUsize::new(1).expect("1 is non-zero"),
165 replay_transactions_threads: NonZeroUsize::new(1).expect("1 is non-zero"),
166 shred_sigverify_threads: NonZeroUsize::new(1).expect("1 is non-zero"),
167 bls_sigverify_threads: NonZeroUsize::new(1).expect("1 is non-zero"),
168 turbine_xdp_sender: None,
169 repair_xdp_sender: None,
170 }
171 }
172}
173
174pub struct AlpenglowInitializationState {
176 pub leader_window_info_sender: Sender<LeaderWindowInfo>,
178 pub optimistic_parent_sender: Sender<LeaderWindowInfo>,
179 pub optimistic_parent_receiver: Receiver<LeaderWindowInfo>,
180 pub replay_highest_frozen: Arc<ReplayHighestFrozen>,
181 pub highest_parent_ready: Arc<RwLock<(Slot, Block)>>,
182 pub highest_finalized: Arc<RwLock<Option<ValidatedBlockFinalizationCert>>>,
183 pub bank_forks_controller: Arc<dyn BankForksController>,
184 pub bank_forks_controller_receiver: BankForksCommandReceiver,
185
186 pub votor_event_sender: VotorEventSender,
188 pub votor_event_receiver: VotorEventReceiver,
189
190 pub cancel: CancellationToken,
192 pub staked_nodes: Arc<RwLock<StakedNodes>>,
193 pub key_notifiers: Arc<RwLock<KeyUpdaters>>,
194
195 pub bls_connection_cache: Arc<ConnectionCache>,
197 pub voting_service_test_override: Option<VotingServiceOverride>,
198}
199
200impl Tvu {
201 #[allow(clippy::too_many_arguments)]
208 pub fn new(
209 vote_account: &Pubkey,
210 authorized_voter_keypairs: Arc<RwLock<Vec<Arc<Keypair>>>>,
211 bank_forks: Arc<RwLock<BankForks>>,
212 cluster_info: &Arc<ClusterInfo>,
213 sockets: TvuSockets,
214 blockstore: Arc<Blockstore>,
215 ledger_signal_receiver: Receiver<bool>,
216 update_parent_receiver: UpdateParentReceiver,
217 rpc_subscriptions: Option<Arc<RpcSubscriptions>>,
218 poh_recorder: &Arc<RwLock<PohRecorder>>,
219 poh_controller: PohController,
220 tower: Tower,
221 tower_storage: Arc<dyn TowerStorage>,
222 vote_history: VoteHistory,
223 vote_history_storage: Arc<dyn VoteHistoryStorage>,
224 leader_schedule_cache: &Arc<LeaderScheduleCache>,
225 exit: Arc<AtomicBool>,
226 block_commitment_cache: Arc<RwLock<BlockCommitmentCache>>,
227 turbine_mode: TurbineMode,
228 transaction_status_sender: Option<TransactionStatusSender>,
229 entry_notification_sender: Option<EntryNotifierSender>,
230 vote_tracker: Arc<VoteTracker>,
231 retransmit_slots_sender: Sender<Slot>,
232 gossip_verified_vote_hash_receiver: GossipVerifiedVoteHashReceiver,
233 verified_voter_slots_sender: VerifiedVoterSlotsSender,
234 verified_voter_slots_receiver: VerifiedVoterSlotsReceiver,
235 replay_vote_sender: ReplayVoteSender,
236 completed_data_sets_sender: Option<CompletedDataSetsSender>,
237 bank_notification_sender: Option<BankNotificationSenderConfig>,
238 duplicate_confirmed_slots_receiver: DuplicateConfirmedSlotsReceiver,
239 tvu_config: TvuConfig,
240 max_slots: &Arc<MaxSlots>,
241 block_metadata_notifier: Option<BlockMetadataNotifierArc>,
242 wait_to_vote_slot: Option<Slot>,
243 snapshot_controller: Option<Arc<SnapshotController>>,
244 log_messages_bytes_limit: Option<usize>,
245 prioritization_fee_cache: Option<Arc<PrioritizationFeeCache>>,
246 banking_tracer: Arc<BankingTracer>,
247 outstanding_repair_requests: Arc<RwLock<OutstandingShredRepairs>>,
248 cluster_slots: Arc<ClusterSlots>,
249 slot_status_notifier: Option<SlotStatusNotifier>,
250 vote_connection_cache: Arc<ConnectionCache>,
251 votor_init: AlpenglowInitializationState,
252 reward_votes_sender: Sender<AddVoteMessage>,
253 ) -> Result<Self, String> {
254 let migration_status = bank_forks.read().unwrap().migration_status();
255
256 let TvuSockets {
257 repair: repair_socket,
258 fetch: fetch_sockets,
259 retransmit: retransmit_sockets,
260 ancestor_hashes_requests: ancestor_hashes_socket,
261 alpenglow: bls_socket,
262 block_id_repair,
263 } = sockets;
264
265 let AlpenglowInitializationState {
266 leader_window_info_sender,
267 optimistic_parent_sender,
268 optimistic_parent_receiver,
269 replay_highest_frozen,
270 highest_parent_ready,
271 bank_forks_controller,
272 bank_forks_controller_receiver,
273 votor_event_sender,
274 votor_event_receiver,
275 cancel,
276 staked_nodes,
277 key_notifiers,
278 bls_connection_cache,
279 voting_service_test_override,
280 highest_finalized,
281 } = votor_init;
282
283 let (consensus_message_sender, consensus_message_receiver) =
285 bounded(MAX_ALPENGLOW_PACKET_NUM);
286 let (consensus_metrics_sender, consensus_metrics_receiver) =
287 bounded(MAX_IN_FLIGHT_CONSENSUS_EVENTS);
288 let generated_cert_types = Arc::new(GeneratedCertTypes::default());
289
290 let bls_sigverify_threads = {
291 let (bls_packet_sender, bls_packet_receiver) = bounded(MAX_ALPENGLOW_PACKET_NUM);
292
293 let (
294 SpawnServerResult {
295 endpoints: _,
296 thread: bls_streamer_t,
297 key_updater: bls_key_updater,
298 },
299 banlist,
300 ) = {
301 let quic_server_params = QuicStreamerConfig {
302 num_threads: NonZeroUsize::new(4.min(num_cpus::get())).unwrap(),
303 ..Default::default()
304 };
305 let qos_config = SimpleQosConfig {
306 max_streams_per_second: VOTOR_RATE_LIMIT_PPS,
307 max_staked_connections: MAX_ALPENGLOW_VOTE_ACCOUNTS * 2,
309 max_connections_per_peer: 2,
311 };
312 spawn_simple_qos_server(
313 "solQuicBLS",
314 "quic_streamer_bls",
315 vec![bls_socket.into()],
316 &cluster_info.keypair(),
317 bls_packet_sender,
318 staked_nodes,
319 quic_server_params,
320 qos_config,
321 cancel,
322 )
323 .unwrap()
324 };
325
326 let sharable_banks = bank_forks.read().unwrap().sharable_banks();
328 let bls_sigverifier_t = bls_sigverifier::spawn_service(
329 exit.clone(),
330 SigVerifierContext {
331 migration_status: migration_status.clone(),
332 banlist,
333 sharable_banks,
334 cluster_info: cluster_info.clone(),
335 leader_schedule: leader_schedule_cache.clone(),
336 num_threads: tvu_config.bls_sigverify_threads.get(),
337 generated_cert_types: generated_cert_types.clone(),
338 },
339 SigVerifierChannels {
340 packet_receiver: bls_packet_receiver,
341 channel_to_repair: verified_voter_slots_sender,
342 channel_to_reward: reward_votes_sender,
343 channel_to_pool: consensus_message_sender,
344 channel_to_metrics: consensus_metrics_sender.clone(),
345 },
346 );
347
348 let mut key_notifiers = key_notifiers.write().unwrap();
349 key_notifiers.add(KeyUpdaterType::Bls, bls_key_updater);
350 (bls_streamer_t, bls_sigverifier_t)
351 };
352
353 let (fetch_sender, fetch_receiver) = EvictingSender::new_bounded(SHRED_FETCH_CHANNEL_SIZE);
354
355 let repair_socket = Arc::new(repair_socket);
356 let ancestor_hashes_socket = Arc::new(ancestor_hashes_socket);
357 let block_id_repair_socket = Arc::new(block_id_repair);
358 let fetch_sockets: Vec<Arc<UdpSocket>> = fetch_sockets.into_iter().map(Arc::new).collect();
359 let fetch_stage = ShredFetchStage::new(
360 fetch_sockets,
361 repair_socket.clone(),
362 fetch_sender,
363 tvu_config.shred_version,
364 bank_forks.clone(),
365 cluster_info.clone(),
366 outstanding_repair_requests.clone(),
367 turbine_mode,
368 exit.clone(),
369 );
370
371 let (verified_sender, verified_receiver) = unbounded();
372
373 let (retransmit_sender, retransmit_receiver) =
374 EvictingSender::new_bounded(CHANNEL_SIZE_RETRANSMIT_INGRESS);
375
376 let shred_sigverify = solana_turbine::sigverify_shreds::spawn_shred_sigverify(
377 cluster_info.clone(),
378 bank_forks.clone(),
379 leader_schedule_cache.clone(),
380 fetch_receiver,
381 retransmit_sender.clone(),
382 verified_sender,
383 Arc::new({
384 let outstanding_repair_requests = outstanding_repair_requests.clone();
385 move |nonce| {
386 outstanding_repair_requests
387 .read()
388 .unwrap()
389 .fetch_metadata_for_nonce(nonce)
390 }
391 }),
392 tvu_config.shred_sigverify_threads,
393 );
394
395 let retransmit_stage = RetransmitStage::new(
396 bank_forks.clone(),
397 leader_schedule_cache.clone(),
398 cluster_info.clone(),
399 Arc::new(retransmit_sockets),
400 retransmit_receiver,
401 max_slots.clone(),
402 rpc_subscriptions.clone(),
403 slot_status_notifier.clone(),
404 tvu_config.turbine_xdp_sender,
405 votor_event_sender.clone(),
406 );
407
408 let (ancestor_duplicate_slots_sender, ancestor_duplicate_slots_receiver) = unbounded();
409 let (duplicate_slots_sender, duplicate_slots_receiver) = bounded(2048);
412 let (ancestor_hashes_replay_update_sender, ancestor_hashes_replay_update_receiver) =
413 unbounded();
414 let (dumped_slots_sender, dumped_slots_receiver) = unbounded();
415 let (popular_pruned_forks_sender, popular_pruned_forks_receiver) = unbounded();
416 let (repair_event_sender, repair_event_receiver) = bounded(100);
418
419 let (completed_slots_sender, completed_slots_receiver) =
421 bounded(MAX_COMPLETED_SLOTS_IN_CHANNEL);
422 blockstore.add_completed_slots_signal(completed_slots_sender);
423
424 let block_id_repair_channels = BlockIdRepairChannels {
425 repair_event_receiver,
426 completed_slots_receiver,
427 };
428
429 let latest_switch_request = LatestSwitchRequest::default();
431
432 let window_service = {
433 let epoch_schedule = bank_forks
434 .read()
435 .unwrap()
436 .working_bank()
437 .epoch_schedule()
438 .clone();
439 let repair_info = RepairInfo {
440 bank_forks: bank_forks.clone(),
441 epoch_schedule,
442 ancestor_duplicate_slots_sender,
443 repair_validators: tvu_config.repair_validators,
444 repair_whitelist: tvu_config.repair_whitelist,
445 cluster_info: cluster_info.clone(),
446 cluster_slots: cluster_slots.clone(),
447 };
448 let repair_service_channels = RepairServiceChannels::new(
449 verified_voter_slots_receiver,
450 dumped_slots_receiver,
451 popular_pruned_forks_sender,
452 ancestor_hashes_replay_update_receiver,
453 );
454 let window_service_channels = WindowServiceChannels::new(
455 verified_receiver,
456 retransmit_sender,
457 completed_data_sets_sender,
458 duplicate_slots_sender.clone(),
459 repair_service_channels,
460 block_id_repair_channels,
461 );
462 WindowService::new(
463 blockstore.clone(),
464 repair_socket,
465 ancestor_hashes_socket,
466 block_id_repair_socket,
467 exit.clone(),
468 repair_info,
469 window_service_channels,
470 leader_schedule_cache.clone(),
471 tvu_config.shred_version,
472 outstanding_repair_requests,
473 tvu_config.repair_xdp_sender,
474 )
475 };
476
477 let (cluster_slots_update_sender, cluster_slots_update_receiver) = unbounded();
478 let cluster_slots_service = ClusterSlotsService::new(
479 blockstore.clone(),
480 cluster_slots.clone(),
481 bank_forks.clone(),
482 cluster_info.clone(),
483 cluster_slots_update_receiver,
484 exit.clone(),
485 );
486
487 let (cost_update_sender, cost_update_receiver) = unbounded();
488 let (drop_bank_sender, drop_bank_receiver) = unbounded();
489 let (voting_sender, voting_receiver) = unbounded();
490 let (bls_sender, bls_receiver) = bounded(MAX_BLS_MESSAGES_TO_SEND);
491
492 let (lockouts_sender, votor_commitment_sender, commitment_service) =
493 AggregateCommitmentService::new(
494 exit.clone(),
495 block_commitment_cache.clone(),
496 rpc_subscriptions.clone(),
497 );
498 let (own_message_sender, own_message_receiver) = bounded(MAX_ALPENGLOW_PACKET_NUM);
499
500 let votor_config = VotorConfig {
501 exit: exit.clone(),
502 vote_account: *vote_account,
503 wait_to_vote_slot,
504 vote_history,
505 vote_history_storage: vote_history_storage.clone(),
506 generated_cert_types,
507 authorized_voter_keypairs: authorized_voter_keypairs.clone(),
508 blockstore: blockstore.clone(),
509 bank_forks: bank_forks.clone(),
510 cluster_info: cluster_info.clone(),
511 leader_schedule_cache: leader_schedule_cache.clone(),
512 consensus_metrics_sender,
513 highest_finalized: highest_finalized.clone(),
514 bank_forks_controller,
515 bls_sender: bls_sender.clone(),
516 commitment_sender: votor_commitment_sender,
517 bank_notification_sender: bank_notification_sender.clone(),
518 leader_window_info_sender,
519 highest_parent_ready,
520 event_sender: votor_event_sender.clone(),
521 latest_switch_request: latest_switch_request.clone(),
522 own_vote_sender: own_message_sender.clone(),
523 repair_event_sender,
524 event_receiver: votor_event_receiver,
525 consensus_message_receiver,
526 own_message_receiver,
527 consensus_metrics_receiver,
528 };
529 let votor = Votor::new(votor_config);
530
531 let replay_senders = ReplaySenders {
532 rpc_subscriptions,
533 slot_status_notifier,
534 transaction_status_sender,
535 entry_notification_sender,
536 bank_notification_sender,
537 ancestor_hashes_replay_update_sender,
538 retransmit_slots_sender,
539 replay_vote_sender,
540 cluster_slots_update_sender,
541 cost_update_sender,
542 voting_sender,
543 bls_sender,
544 drop_bank_sender,
545 block_metadata_notifier,
546 dumped_slots_sender,
547 votor_event_sender,
548 own_message_sender,
549 optimistic_parent_sender,
550 lockouts_sender,
551 };
552
553 let replay_receivers = ReplayReceivers {
554 ledger_signal_receiver,
555 update_parent_receiver,
556 optimistic_parent_receiver,
557 duplicate_slots_receiver,
558 ancestor_duplicate_slots_receiver,
559 duplicate_confirmed_slots_receiver,
560 gossip_verified_vote_hash_receiver,
561 popular_pruned_forks_receiver,
562 bank_forks_controller_receiver,
563 latest_switch_request,
564 };
565
566 let replay_stage_config = ReplayStageConfig {
567 vote_account: *vote_account,
568 authorized_voter_keypairs,
569 exit: exit.clone(),
570 leader_schedule_cache: leader_schedule_cache.clone(),
571 block_commitment_cache,
572 wait_for_vote_to_start_leader: tvu_config.wait_for_vote_to_start_leader,
573 tower_storage: tower_storage.clone(),
574 wait_to_vote_slot,
575 replay_forks_threads: tvu_config.replay_forks_threads,
576 replay_transactions_threads: tvu_config.replay_transactions_threads,
577 blockstore: blockstore.clone(),
578 bank_forks: bank_forks.clone(),
579 cluster_info: cluster_info.clone(),
580 poh_recorder: poh_recorder.clone(),
581 poh_controller,
582 tower,
583 vote_tracker,
584 cluster_slots,
585 log_messages_bytes_limit,
586 prioritization_fee_cache,
587 banking_tracer,
588 snapshot_controller,
589 replay_highest_frozen,
590 };
591
592 let voting_service = VotingService::new(
593 voting_receiver,
594 cluster_info.clone(),
595 poh_recorder.clone(),
596 tower_storage,
597 vote_connection_cache.clone(),
598 );
599
600 let bls_voting_service = BLSVotingService::new(
601 bls_receiver,
602 cluster_info.clone(),
603 vote_history_storage,
604 bls_connection_cache,
605 bank_forks.clone(),
606 highest_finalized,
607 voting_service_test_override,
608 );
609
610 let warm_quic_cache_service = create_cache_warmer_if_needed(
611 None,
612 vote_connection_cache,
613 cluster_info,
614 poh_recorder,
615 &exit,
616 );
617
618 let cost_update_service = CostUpdateService::new(cost_update_receiver);
619
620 let drop_bank_service = DropBankService::new(drop_bank_receiver);
621
622 let replay_stage = ReplayStage::new(replay_stage_config, replay_senders, replay_receivers)?;
623
624 let blockstore_cleanup_service = BlockstoreCleanupService::new(
625 blockstore.clone(),
626 tvu_config.max_ledger_shreds,
627 exit.clone(),
628 );
629
630 let epoch_specs: Box<dyn solana_gossip::epoch_specs::EpochSpecs> =
631 Box::new(EpochSpecs::from(bank_forks));
632
633 let duplicate_shred_listener = DuplicateShredListener::new(
634 exit,
635 cluster_info.clone(),
636 DuplicateShredHandler::new(
637 blockstore,
638 leader_schedule_cache.clone(),
639 epoch_specs,
640 duplicate_slots_sender,
641 tvu_config.shred_version,
642 ),
643 );
644
645 Ok(Tvu {
646 fetch_stage,
647 shred_sigverify,
648 retransmit_stage,
649 window_service,
650 cluster_slots_service,
651 replay_stage,
652 blockstore_cleanup_service,
653 cost_update_service,
654 voting_service,
655 bls_voting_service,
656 warm_quic_cache_service,
657 drop_bank_service,
658 duplicate_shred_listener,
659 bls_sigverify_threads,
660 votor,
661 commitment_service,
662 })
663 }
664
665 pub fn join(self) -> thread::Result<()> {
666 self.retransmit_stage.join()?;
667 self.window_service.join()?;
668 self.cluster_slots_service.join()?;
669 self.fetch_stage.join()?;
670 self.shred_sigverify.join()?;
671 self.blockstore_cleanup_service.join()?;
672 self.replay_stage.join()?;
673 self.cost_update_service.join()?;
674 self.voting_service.join()?;
675 self.bls_voting_service.join()?;
676 if let Some(warmup_service) = self.warm_quic_cache_service {
677 warmup_service.join()?;
678 }
679 self.drop_bank_service.join()?;
680 self.duplicate_shred_listener.join()?;
681 let (streamer, sigverifier) = self.bls_sigverify_threads;
682 streamer.join()?;
683 sigverifier.join()?;
684 self.votor.join()?;
685 self.commitment_service.join()?;
686 Ok(())
687 }
688}
689
690fn create_cache_warmer_if_needed(
691 connection_cache: Option<&Arc<ConnectionCache>>,
692 vote_connection_cache: Arc<ConnectionCache>,
693 cluster_info: &Arc<ClusterInfo>,
694 poh_recorder: &Arc<RwLock<PohRecorder>>,
695 exit: &Arc<AtomicBool>,
696) -> Option<WarmQuicCacheService> {
697 let tpu_connection_cache = connection_cache.filter(|cache| cache.use_quic()).cloned();
698 let vote_connection_cache = Some(vote_connection_cache).filter(|cache| cache.use_quic());
699
700 (tpu_connection_cache.is_some() || vote_connection_cache.is_some()).then(|| {
701 WarmQuicCacheService::new(
702 tpu_connection_cache,
703 vote_connection_cache,
704 cluster_info.clone(),
705 poh_recorder.clone(),
706 exit.clone(),
707 )
708 })
709}
710
711#[cfg(test)]
712pub mod tests {
713 use {
714 super::*,
715 crate::{
716 admin_rpc_post_init::KeyUpdaters, block_creation_loop::ReplayHighestFrozen,
717 consensus::tower_storage::FileTowerStorage,
718 },
719 agave_votor::{
720 event::{VotorEventReceiver, VotorEventSender},
721 vote_history::VoteHistory,
722 vote_history_storage::NullVoteHistoryStorage,
723 },
724 serial_test::serial,
725 solana_gossip::{cluster_info::ClusterInfo, node::Node},
726 solana_hash::Hash,
727 solana_keypair::Keypair,
728 solana_ledger::{
729 blockstore::BlockstoreSignals,
730 blockstore_options::BlockstoreOptions,
731 create_new_tmp_ledger,
732 genesis_utils::{GenesisConfigInfo, create_genesis_config},
733 },
734 solana_net_utils::SocketAddrSpace,
735 solana_poh::poh_recorder::create_test_recorder,
736 solana_rpc::optimistically_confirmed_bank_tracker::OptimisticallyConfirmedBank,
737 solana_runtime::{bank::Bank, bank_forks_controller::BankForksControllerHandle},
738 solana_signer::Signer,
739 solana_tpu_client::tpu_client::{DEFAULT_TPU_CONNECTION_POOL_SIZE, DEFAULT_VOTE_USE_QUIC},
740 std::{
741 sync::atomic::{AtomicU64, Ordering},
742 time::Duration,
743 },
744 };
745
746 #[test]
747 #[serial]
748 fn test_tvu_exit() {
749 agave_logger::setup();
750 let leader = Node::new_localhost();
751 let target1_keypair = Keypair::new();
752 let target1 = Node::new_localhost_with_pubkey(&target1_keypair.pubkey());
753
754 let starting_balance = 10_000;
755 let GenesisConfigInfo { genesis_config, .. } = create_genesis_config(starting_balance);
756
757 let bank_forks = BankForks::new_rw_arc(Bank::new_for_tests(&genesis_config));
758
759 let cluster_info1 = ClusterInfo::new(
761 target1.info.clone(),
762 target1_keypair.into(),
763 SocketAddrSpace::Unspecified,
764 );
765 cluster_info1.insert_info(leader.info);
766 let cref1 = Arc::new(cluster_info1);
767
768 let (blockstore_path, _) = create_new_tmp_ledger!(&genesis_config);
769 let BlockstoreSignals {
770 blockstore,
771 ledger_signal_receiver,
772 update_parent_receiver,
773 ..
774 } = Blockstore::open_with_signal(&blockstore_path, BlockstoreOptions::default())
775 .expect("Expected to successfully open ledger");
776 let blockstore = Arc::new(blockstore);
777 let bank = bank_forks.read().unwrap().working_bank();
778 let (
779 exit,
780 poh_recorder,
781 poh_controller,
782 _transaction_recorder,
783 poh_service,
784 _entry_receiver,
785 ) = create_test_recorder(bank.clone(), blockstore.clone(), None, None);
786 let vote_keypair = Keypair::new();
787 let leader_schedule_cache = Arc::new(LeaderScheduleCache::new_from_bank(&bank));
788 let block_commitment_cache = Arc::new(RwLock::new(BlockCommitmentCache::default()));
789 let (retransmit_slots_sender, _retransmit_slots_receiver) = bounded(1024);
790 let (_gossip_verified_vote_hash_sender, gossip_verified_vote_hash_receiver) = bounded(1024);
791 let (verified_voter_slots_sender, verified_voter_slots_receiver) = bounded(1024);
792 let (replay_vote_sender, _replay_vote_receiver) = bounded(1024);
793 let (_, gossip_confirmed_slots_receiver) = bounded(1024);
794 let max_complete_transaction_status_slot = Arc::new(AtomicU64::default());
795 let outstanding_repair_requests = Arc::<RwLock<OutstandingShredRepairs>>::default();
796 let cluster_slots = Arc::new(ClusterSlots::default_for_tests());
797 let connection_cache = if DEFAULT_VOTE_USE_QUIC {
798 ConnectionCache::new_quic_for_tests(
799 "connection_cache_vote_quic",
800 DEFAULT_TPU_CONNECTION_POOL_SIZE,
801 )
802 } else {
803 ConnectionCache::with_udp(
804 "connection_cache_vote_udp",
805 DEFAULT_TPU_CONNECTION_POOL_SIZE,
806 )
807 };
808 let bls_connection_cache = ConnectionCache::new_quic_for_tests(
809 "connection_cache_bls_quic",
810 DEFAULT_TPU_CONNECTION_POOL_SIZE,
811 );
812 let replay_highest_frozen = Arc::new(ReplayHighestFrozen::default());
813 let (leader_window_info_sender, _leader_window_info_receiver) = bounded(1024);
814 let (optimistic_parent_sender, optimistic_parent_receiver) = bounded(1024);
815 let highest_parent_ready = Arc::new(RwLock::new((
816 0,
817 Block {
818 slot: 0,
819 block_id: Hash::default(),
820 },
821 )));
822 let (votor_event_sender, votor_event_receiver): (VotorEventSender, VotorEventReceiver) =
823 bounded(1024);
824 let staked_nodes = Arc::new(RwLock::new(StakedNodes::default()));
825 let key_notifiers = Arc::new(RwLock::new(KeyUpdaters::default()));
826 let cancel = CancellationToken::new();
827 thread::spawn({
828 let cancel = cancel.clone();
829 let exit = exit.clone();
830 move || loop {
831 if exit.load(Ordering::Relaxed) {
832 cancel.cancel();
833 break;
834 }
835 thread::sleep(Duration::from_secs(1));
836 }
837 });
838 let (bank_forks_controller, bank_forks_controller_receiver) =
839 BankForksControllerHandle::new();
840 let bank_forks_controller = Arc::new(bank_forks_controller);
841 let (reward_votes_sender, _reward_votes_receiver) = bounded(1024);
842
843 let tvu = Tvu::new(
844 &vote_keypair.pubkey(),
845 Arc::new(RwLock::new(vec![Arc::new(vote_keypair)])),
846 bank_forks.clone(),
847 &cref1,
848 TvuSockets {
849 repair: target1.sockets.repair,
850 retransmit: target1.sockets.retransmit_sockets,
851 fetch: target1.sockets.tvu,
852 ancestor_hashes_requests: target1.sockets.ancestor_hashes_requests,
853 alpenglow: target1.sockets.alpenglow,
854 block_id_repair: target1.sockets.block_id_repair,
855 },
856 blockstore,
857 ledger_signal_receiver,
858 update_parent_receiver,
859 Some(Arc::new(RpcSubscriptions::new_for_tests(
860 exit.clone(),
861 max_complete_transaction_status_slot,
862 bank_forks.clone(),
863 block_commitment_cache.clone(),
864 OptimisticallyConfirmedBank::locked_from_bank_forks_root(&bank_forks),
865 ))),
866 &poh_recorder,
867 poh_controller,
868 Tower::default(),
869 Arc::new(FileTowerStorage::default()),
870 VoteHistory::default(),
871 Arc::new(NullVoteHistoryStorage::default()),
872 &leader_schedule_cache,
873 exit.clone(),
874 block_commitment_cache,
875 TurbineMode::default(),
876 None, None, Arc::<VoteTracker>::default(),
879 retransmit_slots_sender,
880 gossip_verified_vote_hash_receiver,
881 verified_voter_slots_sender,
882 verified_voter_slots_receiver,
883 replay_vote_sender,
884 None, None, gossip_confirmed_slots_receiver,
887 TvuConfig::default(),
888 &Arc::new(MaxSlots::default()),
889 None, None, None, None, None, BankingTracer::new_disabled(),
895 outstanding_repair_requests,
896 cluster_slots,
897 None, Arc::new(connection_cache),
899 AlpenglowInitializationState {
900 leader_window_info_sender,
901 optimistic_parent_sender,
902 optimistic_parent_receiver,
903 replay_highest_frozen,
904 highest_parent_ready,
905 votor_event_sender,
906 votor_event_receiver,
907 cancel,
908 staked_nodes,
909 key_notifiers,
910 bls_connection_cache: Arc::new(bls_connection_cache),
911 voting_service_test_override: None,
912 highest_finalized: Arc::new(RwLock::new(None)),
913 bank_forks_controller,
914 bank_forks_controller_receiver,
915 },
916 reward_votes_sender,
917 )
918 .expect("assume success");
919 exit.store(true, Ordering::Relaxed);
920 tvu.join().unwrap();
921 poh_service.join().unwrap();
922 }
923}