solana_core/
voting_service.rs1use {
2 crate::{
3 consensus::tower_storage::{SavedTowerVersions, TowerStorage},
4 next_leader::upcoming_leader_tpu_vote_sockets,
5 },
6 crossbeam_channel::Receiver,
7 solana_client::connection_cache::ConnectionCache,
8 solana_clock::{FORWARD_TRANSACTIONS_TO_LEADER_AT_SLOT_OFFSET, Slot},
9 solana_connection_cache::client_connection::ClientConnection,
10 solana_gossip::cluster_info::ClusterInfo,
11 solana_measure::measure::Measure,
12 solana_poh::poh_recorder::PohRecorder,
13 solana_transaction::Transaction,
14 solana_transaction_error::TransportError,
15 std::{
16 net::SocketAddr,
17 sync::{Arc, RwLock},
18 thread::{self, Builder, JoinHandle},
19 },
20 thiserror::Error,
21};
22
23pub enum VoteOp {
24 PushVote {
25 tx: Transaction,
26 tower_slots: Vec<Slot>,
27 saved_tower: SavedTowerVersions,
28 },
29 RefreshVote {
30 tx: Transaction,
31 last_voted_slot: Slot,
32 },
33}
34
35impl VoteOp {
36 fn tx(&self) -> &Transaction {
37 match self {
38 VoteOp::PushVote { tx, .. } => tx,
39 VoteOp::RefreshVote { tx, .. } => tx,
40 }
41 }
42}
43
44#[derive(Debug, Error)]
45enum SendVoteError {
46 #[error(transparent)]
47 WincodeWriteError(#[from] wincode::WriteError),
48 #[error("Invalid TPU address")]
49 InvalidTpuAddress,
50 #[error(transparent)]
51 TransportError(#[from] TransportError),
52}
53
54fn send_vote_transaction(
55 cluster_info: &ClusterInfo,
56 transaction: &Transaction,
57 tpu: Option<SocketAddr>,
58 connection_cache: &Arc<ConnectionCache>,
59) -> Result<(), SendVoteError> {
60 let tpu = tpu
61 .or_else(|| {
62 cluster_info
63 .my_contact_info()
64 .tpu(connection_cache.protocol())
65 })
66 .ok_or(SendVoteError::InvalidTpuAddress)?;
67 let buf = Arc::new(wincode::serialize(transaction)?);
68 let client = connection_cache.get_connection(&tpu);
69
70 client.send_data_async(buf).map_err(|err| {
71 error!("Ran into an error when sending vote: {err:?} to {tpu:?}");
72 SendVoteError::from(err)
73 })
74}
75
76pub struct VotingService {
77 thread_hdl: JoinHandle<()>,
78}
79
80impl VotingService {
81 pub fn new(
82 vote_receiver: Receiver<VoteOp>,
83 cluster_info: Arc<ClusterInfo>,
84 poh_recorder: Arc<RwLock<PohRecorder>>,
85 tower_storage: Arc<dyn TowerStorage>,
86 connection_cache: Arc<ConnectionCache>,
87 ) -> Self {
88 let thread_hdl = Builder::new()
89 .name("solVoteService".to_string())
90 .spawn({
91 move || {
92 for vote_op in vote_receiver.iter() {
93 Self::handle_vote(
94 &cluster_info,
95 &poh_recorder,
96 tower_storage.as_ref(),
97 vote_op,
98 connection_cache.clone(),
99 );
100 }
101 }
102 })
103 .unwrap();
104 Self { thread_hdl }
105 }
106
107 pub fn handle_vote(
108 cluster_info: &ClusterInfo,
109 poh_recorder: &RwLock<PohRecorder>,
110 tower_storage: &dyn TowerStorage,
111 vote_op: VoteOp,
112 connection_cache: Arc<ConnectionCache>,
113 ) {
114 if let VoteOp::PushVote { saved_tower, .. } = &vote_op {
115 let mut measure = Measure::start("tower storage save");
116 if let Err(err) = tower_storage.store(saved_tower) {
117 error!("Unable to save tower to storage: {err:?}");
118 std::process::exit(1);
119 }
120 measure.stop();
121 trace!("{measure}");
122 }
123
124 const UPCOMING_LEADER_FANOUT_SLOTS: u64 =
128 FORWARD_TRANSACTIONS_TO_LEADER_AT_SLOT_OFFSET.saturating_add(1);
129 #[cfg(test)]
130 static_assertions::const_assert_eq!(UPCOMING_LEADER_FANOUT_SLOTS, 3);
131 let upcoming_leader_sockets = upcoming_leader_tpu_vote_sockets(
132 cluster_info,
133 poh_recorder,
134 UPCOMING_LEADER_FANOUT_SLOTS,
135 connection_cache.protocol(),
136 );
137
138 if !upcoming_leader_sockets.is_empty() {
139 for tpu_vote_socket in upcoming_leader_sockets {
140 let _ = send_vote_transaction(
141 cluster_info,
142 vote_op.tx(),
143 Some(tpu_vote_socket),
144 &connection_cache,
145 );
146 }
147 } else {
148 let _ = send_vote_transaction(cluster_info, vote_op.tx(), None, &connection_cache);
150 }
151
152 match vote_op {
153 VoteOp::PushVote {
154 tx, tower_slots, ..
155 } => {
156 cluster_info.push_vote(&tower_slots, tx);
157 }
158 VoteOp::RefreshVote {
159 tx,
160 last_voted_slot,
161 } => {
162 cluster_info.refresh_vote(tx, last_voted_slot);
163 }
164 }
165 }
166
167 pub fn join(self) -> thread::Result<()> {
168 self.thread_hdl.join()
169 }
170}