Skip to main content

solana_core/
tpu.rs

1//! The `tpu` module implements the Transaction Processing Unit, a
2//! multi-stage transaction processing pipeline in software.
3
4use {
5    crate::{
6        admin_rpc_post_init::{KeyUpdaterType, KeyUpdaters},
7        banking_stage::{
8            BankingControlMsg, BankingStage, BankingStageHandle,
9            transaction_scheduler::scheduler_controller::SchedulerConfig,
10        },
11        banking_trace::{Channels, TracerThread},
12        cluster_info_vote_listener::{
13            ClusterInfoVoteListener, DuplicateConfirmedSlotsSender, GossipVerifiedVoteHashSender,
14            VoteTracker,
15        },
16        fetch_stage::FetchStage,
17        forwarding_stage::{
18            ForwardAddressGetter, ForwardingClientConfig, SpawnForwardingStageResult,
19            spawn_forwarding_stage,
20        },
21        sigverify_stage::SigVerifyStage,
22        staked_nodes_updater_service::StakedNodesUpdaterService,
23        tpu_entry_notifier::TpuEntryNotifier,
24        validator::{BlockProductionMethod, GeneratorConfig},
25    },
26    agave_banking_stage_ingress_types::SchedulerPriorityFloor,
27    agave_votor::event::VotorEventSender,
28    agave_votor_messages::VerifiedVoterSlotsSender,
29    agave_xdp::transmitter::XdpSender,
30    crossbeam_channel::{Receiver, bounded, unbounded},
31    solana_clock::Slot,
32    solana_gossip::cluster_info::ClusterInfo,
33    solana_keypair::Keypair,
34    solana_ledger::{
35        blockstore::Blockstore, blockstore_processor::TransactionStatusSender,
36        entry_notifier_service::EntryNotifierSender,
37    },
38    solana_poh::{
39        poh_recorder::{PohRecorder, WorkingBankEntryOrMarker},
40        transaction_recorder::TransactionRecorder,
41    },
42    solana_pubkey::Pubkey,
43    solana_rpc::{
44        optimistically_confirmed_bank_tracker::BankNotificationSenderConfig,
45        rpc_subscriptions::RpcSubscriptions,
46    },
47    solana_runtime::{
48        bank_forks::BankForks,
49        prioritization_fee_cache::PrioritizationFeeCache,
50        vote_sender_types::{ReplayVoteReceiver, ReplayVoteSender},
51    },
52    solana_streamer::{
53        evicting_sender::EvictingSender,
54        quic::{
55            SimpleQosQuicStreamerConfig, SpawnServerResult, SwQosQuicStreamerConfig,
56            spawn_simple_qos_server, spawn_stake_weighted_qos_server,
57        },
58        quic_socket::QuicSocket,
59        streamer::StakedNodes,
60    },
61    solana_turbine::{
62        XdpSender as TurbineXdpSender,
63        broadcast_stage::{BroadcastStage, BroadcastStageType},
64    },
65    std::{
66        collections::{HashMap, HashSet},
67        net::{Ipv4Addr, UdpSocket},
68        num::NonZeroUsize,
69        path::PathBuf,
70        sync::{Arc, RwLock, atomic::AtomicBool},
71        thread::{self, JoinHandle},
72    },
73    tokio::sync::mpsc,
74    tokio_util::sync::CancellationToken,
75};
76
77pub struct TpuSockets {
78    pub vote: Vec<UdpSocket>,
79    pub broadcast: Vec<UdpSocket>,
80    pub transactions_quic: Vec<UdpSocket>,
81    pub transactions_forwards_quic: Vec<UdpSocket>,
82    pub vote_quic: Vec<UdpSocket>,
83    /// Client-side socket for the forwarding votes.
84    pub vote_forwarding_client: UdpSocket,
85}
86
87// Conservatively allow 20 TPS per validator.
88pub const MAX_VOTES_PER_SECOND: u64 = 20;
89
90/// Size of the channel between streamer and TPU sigverify stage. The values have been selected to
91/// be conservative max of obsersed on mnb during high-load events.
92const TPU_CHANNEL_SIZE: usize = 50_000;
93
94/// Size of the channel between the vote streamer and the TPU sigverify stage.
95/// Chosen based on nominal voting load for a cluster with ~2000 validators + some margin.
96pub(crate) const TPU_VOTE_CHANNEL_SIZE: usize = 4_000;
97
98/// Size of the channel between the TPU forwards streamer and the fetch stage.
99/// Mirrors `TPU_CHANNEL_SIZE`; the streamer uses `try_send`, so an over-full
100/// channel drops packets (tracked via streamer metrics) rather than blocking.
101const TPU_FORWARD_CHANNEL_SIZE: usize = 50_000;
102
103pub struct Tpu {
104    fetch_stage: FetchStage,
105    cluster_info_vote_listener: ClusterInfoVoteListener,
106    sigverify_stage: SigVerifyStage,
107    banking_stage: BankingStageHandle,
108    forwarding_stage: JoinHandle<()>,
109    broadcast_stage: BroadcastStage,
110    tpu_quic_t: thread::JoinHandle<()>,
111    tpu_forwards_quic_t: thread::JoinHandle<()>,
112    tpu_entry_notifier: Option<TpuEntryNotifier>,
113    staked_nodes_updater_service: StakedNodesUpdaterService,
114    tracer_thread_hdl: TracerThread,
115    tpu_vote_quic_t: thread::JoinHandle<()>,
116}
117
118impl Tpu {
119    #[allow(clippy::too_many_arguments)]
120    pub fn new_with_client(
121        cluster_info: &Arc<ClusterInfo>,
122        poh_recorder: &Arc<RwLock<PohRecorder>>,
123        transaction_recorder: TransactionRecorder,
124        entry_receiver: Receiver<WorkingBankEntryOrMarker>,
125        retransmit_slots_receiver: Receiver<Slot>,
126        sockets: TpuSockets,
127        subscriptions: Option<Arc<RpcSubscriptions>>,
128        transaction_status_sender: Option<TransactionStatusSender>,
129        entry_notification_sender: Option<EntryNotifierSender>,
130        blockstore: Arc<Blockstore>,
131        broadcast_type: &BroadcastStageType,
132        leader_schedule_cache: Arc<solana_ledger::leader_schedule_cache::LeaderScheduleCache>,
133        turbine_xdp_sender: Option<TurbineXdpSender>,
134        quic_xdp_sender: Option<(XdpSender, Ipv4Addr)>,
135        exit: Arc<AtomicBool>,
136        shred_version: u16,
137        vote_tracker: Arc<VoteTracker>,
138        bank_forks: Arc<RwLock<BankForks>>,
139        verified_voter_slots_sender: VerifiedVoterSlotsSender,
140        gossip_verified_vote_hash_sender: GossipVerifiedVoteHashSender,
141        replay_vote_receiver: ReplayVoteReceiver,
142        replay_vote_sender: ReplayVoteSender,
143        bank_notification_sender: Option<BankNotificationSenderConfig>,
144        duplicate_confirmed_slot_sender: DuplicateConfirmedSlotsSender,
145        tpu_forwarding_client_config: ForwardingClientConfig,
146        keypair: &Keypair,
147        log_messages_bytes_limit: Option<usize>,
148        staked_nodes: &Arc<RwLock<StakedNodes>>,
149        shared_staked_nodes_overrides: Arc<RwLock<HashMap<Pubkey, u64>>>,
150        banking_tracer_channels: Channels,
151        tracer_thread_hdl: TracerThread,
152        tpu_quic_server_config: SwQosQuicStreamerConfig,
153        tpu_fwd_quic_server_config: SwQosQuicStreamerConfig,
154        vote_quic_server_config: SimpleQosQuicStreamerConfig,
155        prioritization_fee_cache: Option<Arc<PrioritizationFeeCache>>,
156        tpu_sigverify_threads: NonZeroUsize,
157        block_production_method: BlockProductionMethod,
158        block_production_num_workers: NonZeroUsize,
159        block_production_scheduler_config: SchedulerConfig,
160        filter_keys: Arc<HashSet<Pubkey>>,
161        enable_block_production_forwarding: bool,
162        _generator_config: Option<GeneratorConfig>, /* vestigial code for replay invalidator */
163        key_notifiers: Arc<RwLock<KeyUpdaters>>,
164        banking_control_receiver: mpsc::Receiver<BankingControlMsg>,
165        scheduler_bindings: Option<(PathBuf, mpsc::Sender<BankingControlMsg>)>,
166        cancel: CancellationToken,
167        votor_event_sender: VotorEventSender,
168    ) -> Self {
169        let TpuSockets {
170            vote: tpu_vote_sockets,
171            broadcast: broadcast_sockets,
172            transactions_quic: transactions_quic_sockets,
173            transactions_forwards_quic: transactions_forwards_quic_sockets,
174            vote_quic: tpu_vote_quic_sockets,
175            vote_forwarding_client: vote_forwarding_client_socket,
176        } = sockets;
177
178        let (packet_sender, packet_receiver) = bounded(TPU_CHANNEL_SIZE);
179        let (vote_packet_sender, vote_packet_receiver) = bounded(TPU_VOTE_CHANNEL_SIZE);
180        let evicting_vote_sender =
181            EvictingSender::new(vote_packet_sender.clone(), vote_packet_receiver.clone());
182        let (forwarded_packet_sender, forwarded_packet_receiver) =
183            bounded(TPU_FORWARD_CHANNEL_SIZE);
184        let fetch_stage = FetchStage::new_with_sender(
185            tpu_vote_sockets,
186            exit.clone(),
187            &packet_sender,
188            &evicting_vote_sender,
189            forwarded_packet_receiver,
190            poh_recorder,
191            None, // coalesce
192        );
193
194        let staked_nodes_updater_service = StakedNodesUpdaterService::new(
195            exit.clone(),
196            bank_forks.clone(),
197            staked_nodes.clone(),
198            shared_staked_nodes_overrides,
199        );
200
201        let Channels {
202            non_vote_sender,
203            non_vote_receiver,
204            tpu_vote_sender,
205            tpu_vote_receiver,
206            gossip_vote_sender,
207            gossip_vote_receiver,
208        } = banking_tracer_channels;
209
210        // Streamer for Votes:
211        let quic_vote_sockets: Vec<QuicSocket> =
212            tpu_vote_quic_sockets.into_iter().map(Into::into).collect();
213        let (
214            SpawnServerResult {
215                endpoints: _,
216                thread: tpu_vote_quic_t,
217                key_updater: vote_streamer_key_updater,
218            },
219            _banlist,
220        ) = spawn_simple_qos_server(
221            "solQuicTVo",
222            "quic_streamer_tpu_vote",
223            quic_vote_sockets,
224            keypair,
225            vote_packet_sender,
226            staked_nodes.clone(),
227            vote_quic_server_config.quic_streamer_config,
228            vote_quic_server_config.qos_config,
229            cancel.clone(),
230        )
231        .unwrap();
232
233        // We check on validator startup that XDP is not mixed with multihoming, so by construction
234        // at this moment all the transactions_quic_sockets and transactions_forwards_quic_sockets
235        // have the same bind IP:PORT.
236
237        // Streamer for TPU
238        let transactions_quic_sockets =
239            into_quic_sockets(transactions_quic_sockets, quic_xdp_sender.clone());
240        let SpawnServerResult {
241            endpoints: _,
242            thread: tpu_quic_t,
243            key_updater,
244        } = spawn_stake_weighted_qos_server(
245            "solQuicTpu",
246            "quic_streamer_tpu",
247            transactions_quic_sockets,
248            keypair,
249            packet_sender,
250            staked_nodes.clone(),
251            tpu_quic_server_config.quic_streamer_config,
252            tpu_quic_server_config.qos_config,
253            cancel.clone(),
254        )
255        .unwrap();
256
257        // Streamer for TPU forward
258        let transactions_forwards_quic_sockets =
259            into_quic_sockets(transactions_forwards_quic_sockets, quic_xdp_sender);
260        let SpawnServerResult {
261            endpoints: _,
262            thread: tpu_forwards_quic_t,
263            key_updater: forwards_key_updater,
264        } = spawn_stake_weighted_qos_server(
265            "solQuicTpuFwd",
266            "quic_streamer_tpu_forwards",
267            transactions_forwards_quic_sockets,
268            keypair,
269            forwarded_packet_sender,
270            staked_nodes.clone(),
271            tpu_fwd_quic_server_config.quic_streamer_config,
272            tpu_fwd_quic_server_config.qos_config,
273            cancel,
274        )
275        .unwrap();
276
277        let (forward_stage_sender, forward_stage_receiver) = bounded(50_000);
278
279        // Shared between sigverify and scheduler. The scheduler publishes
280        // a priority floor under saturation; sigverify reads it and drops
281        // below-floor packets ahead of signature verification.
282        let scheduler_priority_floor = Arc::new(SchedulerPriorityFloor::new());
283
284        let (sigverify_stage, gossip_sigverify_handle) = SigVerifyStage::new(
285            packet_receiver,
286            vote_packet_receiver,
287            non_vote_sender,
288            tpu_vote_sender,
289            forward_stage_sender.clone(),
290            tpu_sigverify_threads,
291            enable_block_production_forwarding,
292            bank_forks.read().unwrap().sharable_banks(),
293            Some(scheduler_priority_floor.clone()),
294        );
295
296        let cluster_info_vote_listener = ClusterInfoVoteListener::new(
297            exit.clone(),
298            cluster_info.clone(),
299            gossip_sigverify_handle,
300            gossip_vote_sender,
301            vote_tracker,
302            bank_forks.clone(),
303            subscriptions,
304            verified_voter_slots_sender,
305            gossip_verified_vote_hash_sender,
306            replay_vote_receiver,
307            blockstore.clone(),
308            bank_notification_sender,
309            duplicate_confirmed_slot_sender,
310        );
311
312        let banking_stage = BankingStage::new_num_threads(
313            block_production_method,
314            poh_recorder.clone(),
315            transaction_recorder,
316            non_vote_receiver,
317            tpu_vote_receiver,
318            gossip_vote_receiver,
319            banking_control_receiver,
320            block_production_num_workers,
321            block_production_scheduler_config,
322            transaction_status_sender,
323            replay_vote_sender,
324            log_messages_bytes_limit,
325            bank_forks.clone(),
326            prioritization_fee_cache,
327            filter_keys,
328            scheduler_priority_floor,
329        );
330
331        #[cfg(unix)]
332        if let Some((path, banking_control_sender)) = scheduler_bindings {
333            super::scheduler_bindings_server::spawn(&path, banking_control_sender);
334        }
335        #[cfg(not(unix))]
336        assert!(scheduler_bindings.is_none());
337
338        let SpawnForwardingStageResult {
339            join_handle: forwarding_stage,
340            client_updater,
341        } = spawn_forwarding_stage(
342            forward_stage_receiver,
343            tpu_forwarding_client_config,
344            vote_forwarding_client_socket,
345            bank_forks.read().unwrap().sharable_banks(),
346            ForwardAddressGetter::new(cluster_info.clone(), poh_recorder.clone()),
347        );
348
349        let (entry_receiver, tpu_entry_notifier) =
350            if let Some(entry_notification_sender) = entry_notification_sender {
351                let (broadcast_entry_sender, broadcast_entry_receiver) = unbounded();
352                let tpu_entry_notifier = TpuEntryNotifier::new(
353                    entry_receiver,
354                    entry_notification_sender,
355                    broadcast_entry_sender,
356                    exit.clone(),
357                );
358                (broadcast_entry_receiver, Some(tpu_entry_notifier))
359            } else {
360                (entry_receiver, None)
361            };
362
363        let broadcast_stage = broadcast_type.new_broadcast_stage(
364            broadcast_sockets,
365            cluster_info.clone(),
366            entry_receiver,
367            retransmit_slots_receiver,
368            exit,
369            blockstore,
370            bank_forks,
371            leader_schedule_cache,
372            shred_version,
373            turbine_xdp_sender,
374            votor_event_sender,
375        );
376
377        let mut key_notifiers = key_notifiers.write().unwrap();
378        key_notifiers.add(KeyUpdaterType::Tpu, key_updater);
379        key_notifiers.add(KeyUpdaterType::TpuForwards, forwards_key_updater);
380        key_notifiers.add(KeyUpdaterType::TpuVote, vote_streamer_key_updater);
381        key_notifiers.add(KeyUpdaterType::Forward, client_updater);
382
383        Self {
384            fetch_stage,
385            cluster_info_vote_listener,
386            sigverify_stage,
387            banking_stage,
388            forwarding_stage,
389            broadcast_stage,
390            tpu_quic_t,
391            tpu_forwards_quic_t,
392            tpu_entry_notifier,
393            staked_nodes_updater_service,
394            tracer_thread_hdl,
395            tpu_vote_quic_t,
396        }
397    }
398
399    pub fn join(self) -> thread::Result<()> {
400        let results = vec![
401            self.fetch_stage.join(),
402            self.cluster_info_vote_listener.join(),
403            self.sigverify_stage.join(),
404            self.banking_stage.join(),
405            self.forwarding_stage.join(),
406            self.staked_nodes_updater_service.join(),
407            self.tpu_quic_t.join(),
408            self.tpu_forwards_quic_t.join(),
409            self.tpu_vote_quic_t.join(),
410        ];
411        let broadcast_result = self.broadcast_stage.join();
412        for result in results {
413            result?;
414        }
415        if let Some(tpu_entry_notifier) = self.tpu_entry_notifier {
416            tpu_entry_notifier.join()?;
417        }
418        let _ = broadcast_result?;
419        if let Some(tracer_thread_hdl) = self.tracer_thread_hdl
420            && let Err(tracer_result) = tracer_thread_hdl.join()?
421        {
422            error!(
423                "banking tracer thread returned error after successful thread join: \
424                 {tracer_result:?}"
425            );
426        }
427        Ok(())
428    }
429}
430
431fn into_quic_sockets(
432    sockets: impl IntoIterator<Item = UdpSocket>,
433    quic_xdp_sender: Option<(XdpSender, Ipv4Addr)>,
434) -> impl Iterator<Item = QuicSocket> {
435    sockets
436        .into_iter()
437        .map(move |socket| match &quic_xdp_sender {
438            Some((xdp_sender, fallback_src_ip)) => {
439                QuicSocket::with_xdp(socket, *fallback_src_ip, xdp_sender.clone())
440            }
441            None => QuicSocket::from(socket),
442        })
443}