Skip to main content

solana_core/
voting_service.rs

1use {
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        // Attempt to send our vote transaction to the leaders for the next few
125        // slots. From the current slot to the forwarding slot offset
126        // (inclusive).
127        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            // Send to our own tpu vote socket if we cannot find a leader to send to
149            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}