1use {
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 pub vote_forwarding_client: UdpSocket,
85}
86
87pub const MAX_VOTES_PER_SECOND: u64 = 20;
89
90const TPU_CHANNEL_SIZE: usize = 50_000;
93
94pub(crate) const TPU_VOTE_CHANNEL_SIZE: usize = 4_000;
97
98const 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>, 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, );
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 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 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 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 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}