Skip to main content

solana_core/
tvu.rs

1//! The `tvu` module implements the Transaction Validation Unit, a multi-stage transaction
2//! validation pipeline in software.
3
4use {
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
97/// Sets the upper bound on the number of batches stored in the retransmit
98/// stage ingress channel.
99/// Allows for a max of 16k batches of up to 64 packets each
100/// (PACKETS_PER_BATCH).
101/// This translates to about 1 GB of RAM for packet storage in the worst case.
102/// In reality this means about 200K shreds since most batches are not full.
103const CHANNEL_SIZE_RETRANSMIT_INGRESS: usize = 16 * 1024;
104
105/// The maximum number of alpenglow packets that can be processed in a single batch
106pub(crate) const MAX_ALPENGLOW_PACKET_NUM: usize = 10_000;
107/// The maximum number of distinct bls messages that can be sent in a single batch.
108/// This is overprovisioned to account for standstill scenarios, where a large amount
109/// of votes / certificate need to be refreshed.
110const 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    // Validators from which repairs are requested
144    pub repair_validators: Option<HashSet<Pubkey>>,
145    // Validators which should be given priority when serving repairs
146    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
174/// Shared state from validator necessary to instantiate votor and related services
175pub struct AlpenglowInitializationState {
176    // Shared with block creation loop
177    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    // Main communication channel
187    pub votor_event_sender: VotorEventSender,
188    pub votor_event_receiver: VotorEventReceiver,
189
190    // For BLS streamer setup
191    pub cancel: CancellationToken,
192    pub staked_nodes: Arc<RwLock<StakedNodes>>,
193    pub key_notifiers: Arc<RwLock<KeyUpdaters>>,
194
195    // For BLS voting service
196    pub bls_connection_cache: Arc<ConnectionCache>,
197    pub voting_service_test_override: Option<VotingServiceOverride>,
198}
199
200impl Tvu {
201    /// This service receives messages from a leader in the network and processes the transactions
202    /// on the bank state.
203    /// # Arguments
204    /// * `cluster_info` - The cluster_info state.
205    /// * `sockets` - fetch, repair, and retransmit sockets
206    /// * `blockstore` - the ledger itself
207    #[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        // streamer and sigverify for A2A BLS messages
284        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                    // Cap by # of active validators (some overhead for epoch boundaries)
308                    max_staked_connections: MAX_ALPENGLOW_VOTE_ACCOUNTS * 2,
309                    // Two staked connection per validator to account for hotspares
310                    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            // sigverifier
327            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        // This channel is used by both gossip and window service. Gossip will not fill this
410        // beyond 50%, while window_service will perform a blocking send.
411        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        // Create repair event channel for BlockIdRepairService
417        let (repair_event_sender, repair_event_receiver) = bounded(100);
418
419        // Create completed slots channel for BlockIdRepairService
420        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        // Shared latest switch-bank request from Votor to ReplayStage.
430        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        //start cluster_info1
760        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, // transaction_status_sender
877            None, // entry_notification_sender
878            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, // completed_data_sets_sender
885            None, // bank_notification_sender
886            gossip_confirmed_slots_receiver,
887            TvuConfig::default(),
888            &Arc::new(MaxSlots::default()),
889            None, // block_metadata_notifier
890            None, // wait_to_vote_slot
891            None, // snapshot_controller
892            None, // log_messages_bytes_limit
893            None, // prioritization_fee_cache
894            BankingTracer::new_disabled(),
895            outstanding_repair_requests,
896            cluster_slots,
897            None, // slot_status_notifier
898            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}