Skip to main content

solana_runtime/
bank_forks_controller.rs

1use {
2    crate::{bank::Bank, bank_forks::BankForks, installed_scheduler_pool::BankWithScheduler},
3    agave_votor_messages::consensus_message::Block,
4    crossbeam_channel::{Receiver, RecvTimeoutError, Sender, bounded},
5    log::warn,
6    solana_clock::Slot,
7    solana_metrics::datapoint_info,
8    std::{
9        fmt,
10        sync::{Arc, Mutex},
11        time::{Duration, Instant},
12    },
13    thiserror::Error,
14};
15
16const CHANNEL_SIZE: usize = 16;
17
18#[derive(Debug, Error)]
19pub enum BankForksControllerError {
20    #[error("bank forks controller is disconnected")]
21    Disconnected,
22    #[error("bank to insert for slot {0} was stale, failed to insert")]
23    UnableToInsertStaleBank(Slot),
24}
25
26pub enum BankForksCommand {
27    InsertBank {
28        bank: Box<Bank>,
29        response_sender: Sender<Option<BankWithScheduler>>,
30    },
31    ClearBank {
32        slot: Slot,
33        response_sender: Sender<()>,
34    },
35}
36
37#[derive(Clone, Copy, Debug)]
38pub struct SetRootCommand {
39    pub new_root: Block,
40}
41
42impl SetRootCommand {
43    /// Whether the requested root still identifies a frozen bank newer than the applied root.
44    pub fn matches_frozen_bank(&self, bank_forks: &BankForks) -> bool {
45        self.new_root.slot > bank_forks.root()
46            && bank_forks.get(self.new_root.slot).is_some_and(|bank| {
47                bank.is_frozen() && bank.block_id() == Some(self.new_root.block_id)
48            })
49    }
50}
51
52impl BankForksCommand {
53    fn metric_slot(&self) -> Slot {
54        match self {
55            Self::InsertBank { bank, .. } => bank.slot(),
56            Self::ClearBank { slot, .. } => *slot,
57        }
58    }
59}
60
61impl fmt::Display for BankForksCommand {
62    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
63        match self {
64            Self::InsertBank { .. } => write!(f, "insert_bank"),
65            Self::ClearBank { .. } => write!(f, "clear_bank"),
66        }
67    }
68}
69
70pub trait BankForksController: Send + Sync {
71    fn insert_bank(&self, bank: Bank) -> Result<BankWithScheduler, BankForksControllerError>;
72
73    fn enqueue_set_root(&self, new_root: Block);
74
75    fn clear_bank(&self, slot: Slot) -> Result<(), BankForksControllerError>;
76}
77
78/// Handle used by non-replay threads to serialize BankForks writes onto ReplayStage.
79#[derive(Clone)]
80pub struct BankForksControllerHandle {
81    sender: Sender<BankForksCommand>,
82    pending_set_root: Arc<Mutex<Option<SetRootCommand>>>,
83    set_root_signal_sender: Sender<()>,
84}
85
86impl BankForksControllerHandle {
87    pub fn new() -> (Self, BankForksCommandReceiver) {
88        let (sender, receiver) = bounded(CHANNEL_SIZE);
89        let (set_root_signal_sender, set_root_signal_receiver) = bounded(1);
90        let pending_set_root = Arc::new(Mutex::new(None));
91        (
92            Self {
93                sender,
94                pending_set_root: pending_set_root.clone(),
95                set_root_signal_sender,
96            },
97            BankForksCommandReceiver {
98                receiver,
99                pending_set_root,
100                set_root_signal_receiver,
101            },
102        )
103    }
104
105    fn send_command<T>(
106        &self,
107        command: BankForksCommand,
108        response_receiver: Receiver<T>,
109    ) -> Result<T, BankForksControllerError> {
110        let command_name = command.to_string();
111        let slot = command.metric_slot();
112        let queue_len_before_send = self.sender.len();
113        let total_start = Instant::now();
114        let send_start = Instant::now();
115        if self.sender.send(command).is_err() {
116            return Err(BankForksControllerError::Disconnected);
117        }
118        let send_us = send_start.elapsed().as_micros() as i64;
119
120        let response_wait_start = Instant::now();
121        let response = loop {
122            match response_receiver.recv_timeout(Duration::from_millis(100)) {
123                Ok(response) => break response,
124                Err(RecvTimeoutError::Disconnected) => {
125                    return Err(BankForksControllerError::Disconnected);
126                }
127                Err(RecvTimeoutError::Timeout) => (),
128            }
129            warn!(
130                "Replay is stuck, waiting for {}ms no response to {command_name} for {slot}",
131                response_wait_start.elapsed().as_millis()
132            );
133        };
134
135        let response_wait_us = response_wait_start.elapsed().as_micros() as i64;
136        let total_us = total_start.elapsed().as_micros() as i64;
137        datapoint_info!(
138            "bank_forks_controller-command",
139            ("command", command_name, String),
140            ("slot", slot as i64, i64),
141            ("queue_len_before_send", queue_len_before_send as i64, i64),
142            ("send_us", send_us, i64),
143            ("response_wait_us", response_wait_us, i64),
144            ("total_us", total_us, i64),
145        );
146
147        Ok(response)
148    }
149}
150
151impl BankForksController for BankForksControllerHandle {
152    fn insert_bank(&self, bank: Bank) -> Result<BankWithScheduler, BankForksControllerError> {
153        let slot = bank.slot();
154        let (response_sender, response_receiver) = bounded(1);
155        let bank = self.send_command(
156            BankForksCommand::InsertBank {
157                bank: Box::new(bank),
158                response_sender,
159            },
160            response_receiver,
161        )?;
162        bank.ok_or(BankForksControllerError::UnableToInsertStaleBank(slot))
163    }
164
165    fn enqueue_set_root(&self, new_root: Block) {
166        let total_start = Instant::now();
167        let command = SetRootCommand { new_root };
168
169        {
170            let mut pending_set_root = self.pending_set_root.lock().unwrap();
171            // Replay only needs to process the highest pending root.
172            if pending_set_root
173                .as_ref()
174                .is_none_or(|pending| command.new_root.slot > pending.new_root.slot)
175            {
176                *pending_set_root = Some(command);
177            }
178        }
179        let _ = self.set_root_signal_sender.try_send(());
180
181        let total_us = total_start.elapsed().as_micros() as i64;
182
183        datapoint_info!(
184            "bank_forks_controller-command",
185            ("command", "set_root", String),
186            ("slot", new_root.slot as i64, i64),
187            ("total_us", total_us, i64),
188        );
189    }
190
191    fn clear_bank(&self, slot: Slot) -> Result<(), BankForksControllerError> {
192        let (response_sender, response_receiver) = bounded(1);
193        self.send_command(
194            BankForksCommand::ClearBank {
195                slot,
196                response_sender,
197            },
198            response_receiver,
199        )
200    }
201}
202
203pub struct BankForksCommandReceiver {
204    receiver: Receiver<BankForksCommand>,
205    pending_set_root: Arc<Mutex<Option<SetRootCommand>>>,
206    set_root_signal_receiver: Receiver<()>,
207}
208
209impl BankForksCommandReceiver {
210    pub fn receiver(&self) -> &Receiver<BankForksCommand> {
211        &self.receiver
212    }
213
214    pub fn set_root_signal_receiver(&self) -> &Receiver<()> {
215        &self.set_root_signal_receiver
216    }
217
218    pub fn take_set_root_command(&self) -> Option<SetRootCommand> {
219        self.pending_set_root.lock().unwrap().take()
220    }
221}
222
223#[cfg(test)]
224mod tests {
225    use {
226        super::*,
227        crate::{bank::SlotLeader, bank_forks::BankForks, genesis_utils::create_genesis_config},
228        solana_hash::Hash,
229        std::{thread, time::Duration},
230    };
231
232    #[test]
233    fn test_bank_forks_controller_keeps_highest_pending_set_root() {
234        let (controller, receiver) = BankForksControllerHandle::new();
235
236        let block_id_5 = Hash::new_unique();
237        controller.enqueue_set_root(Block {
238            slot: 5,
239            block_id: block_id_5,
240        });
241        controller.enqueue_set_root(Block {
242            slot: 3,
243            block_id: Hash::new_unique(),
244        });
245        let command = receiver.take_set_root_command().unwrap();
246        assert_eq!(command.new_root.slot, 5);
247        assert_eq!(command.new_root.block_id, block_id_5);
248        assert!(receiver.take_set_root_command().is_none());
249
250        controller.enqueue_set_root(Block {
251            slot: 3,
252            block_id: Hash::new_unique(),
253        });
254        controller.enqueue_set_root(Block {
255            slot: 5,
256            block_id: block_id_5,
257        });
258        assert_eq!(receiver.take_set_root_command().unwrap().new_root.slot, 5);
259    }
260
261    #[test]
262    fn test_bank_forks_controller_signals_pending_set_root() {
263        let (controller, receiver) = BankForksControllerHandle::new();
264
265        controller.enqueue_set_root(Block {
266            slot: 1,
267            block_id: Hash::new_unique(),
268        });
269        receiver
270            .set_root_signal_receiver()
271            .recv_timeout(Duration::from_secs(1))
272            .unwrap();
273        assert_eq!(receiver.take_set_root_command().unwrap().new_root.slot, 1);
274
275        controller.enqueue_set_root(Block {
276            slot: 2,
277            block_id: Hash::new_unique(),
278        });
279        controller.enqueue_set_root(Block {
280            slot: 3,
281            block_id: Hash::new_unique(),
282        });
283        receiver
284            .set_root_signal_receiver()
285            .recv_timeout(Duration::from_secs(1))
286            .unwrap();
287        assert!(receiver.set_root_signal_receiver().try_recv().is_err());
288        assert_eq!(receiver.take_set_root_command().unwrap().new_root.slot, 3);
289    }
290
291    #[test]
292    fn test_set_root_command_matches_frozen_bank() {
293        let genesis = create_genesis_config(10_000);
294        let bank_forks = BankForks::new_rw_arc(Bank::new_for_tests(&genesis.genesis_config));
295        let parent_bank = bank_forks.read().unwrap().root_bank();
296        let bank = Bank::new_from_parent(parent_bank, SlotLeader::default(), 1);
297        let block_id = Hash::new_unique();
298        bank.set_block_id(Some(block_id));
299        let bank = bank_forks
300            .write()
301            .unwrap()
302            .insert(bank)
303            .clone_without_scheduler();
304        let command = SetRootCommand {
305            new_root: Block { slot: 1, block_id },
306        };
307
308        assert!(!command.matches_frozen_bank(&bank_forks.read().unwrap()));
309        bank.freeze();
310        assert!(command.matches_frozen_bank(&bank_forks.read().unwrap()));
311
312        let mismatched_command = SetRootCommand {
313            new_root: Block {
314                block_id: Hash::new_unique(),
315                ..command.new_root
316            },
317        };
318        assert!(!mismatched_command.matches_frozen_bank(&bank_forks.read().unwrap()));
319
320        let missing_command = SetRootCommand {
321            new_root: Block {
322                slot: 2,
323                block_id: Hash::new_unique(),
324            },
325        };
326        assert!(!missing_command.matches_frozen_bank(&bank_forks.read().unwrap()));
327
328        bank_forks.write().unwrap().set_root(1, None, None);
329        assert!(!command.matches_frozen_bank(&bank_forks.read().unwrap()));
330    }
331
332    #[test]
333    fn test_bank_forks_controller_insert_and_set_root() {
334        let genesis = create_genesis_config(10_000);
335        let bank_forks = BankForks::new_rw_arc(Bank::new_for_tests(&genesis.genesis_config));
336        let (controller, receiver) = BankForksControllerHandle::new();
337        let replay_bank_forks = bank_forks.clone();
338        let (root_sender, root_receiver) = bounded(1);
339        let replay_thread = thread::spawn(move || {
340            loop {
341                if let Some(command) = receiver.take_set_root_command() {
342                    let new_root = command.new_root.slot;
343                    {
344                        let mut bank_forks = replay_bank_forks.write().unwrap();
345                        bank_forks.set_root(new_root, None, None);
346                    }
347                    root_sender.send(new_root).unwrap();
348                }
349                let command = match receiver.receiver().recv_timeout(Duration::from_millis(10)) {
350                    Ok(command) => command,
351                    Err(RecvTimeoutError::Timeout) => continue,
352                    Err(RecvTimeoutError::Disconnected) => break,
353                };
354                match command {
355                    BankForksCommand::InsertBank {
356                        bank,
357                        response_sender,
358                    } => {
359                        let bank = {
360                            let mut bank_forks = replay_bank_forks.write().unwrap();
361                            bank_forks.insert(*bank)
362                        };
363                        response_sender.send(Some(bank)).unwrap();
364                    }
365                    BankForksCommand::ClearBank {
366                        slot,
367                        response_sender,
368                    } => {
369                        let bank_to_clear =
370                            replay_bank_forks.read().unwrap().get_with_scheduler(slot);
371                        if let Some(bank) = bank_to_clear {
372                            let _ = bank.wait_for_completed_scheduler();
373                        }
374
375                        {
376                            let mut bank_forks = replay_bank_forks.write().unwrap();
377                            bank_forks.clear_bank(slot, false);
378                        }
379                        response_sender.send(()).unwrap();
380                    }
381                }
382            }
383        });
384
385        let parent_bank = bank_forks.read().unwrap().root_bank();
386        let bank = Bank::new_from_parent(parent_bank, SlotLeader::default(), 1);
387        let block_id = Hash::new_unique();
388        bank.set_block_id(Some(block_id));
389        bank.freeze();
390        let inserted_bank = controller.insert_bank(bank).unwrap();
391        assert_eq!(inserted_bank.slot(), 1);
392        assert!(bank_forks.read().unwrap().get(1).is_some());
393
394        controller.enqueue_set_root(Block { slot: 1, block_id });
395        assert_eq!(root_receiver.recv().unwrap(), 1);
396        assert_eq!(bank_forks.read().unwrap().root(), 1);
397
398        let parent_bank = bank_forks.read().unwrap().root_bank();
399        let bank = Bank::new_from_parent(parent_bank, SlotLeader::default(), 2);
400        bank.freeze();
401        let inserted_bank = controller.insert_bank(bank).unwrap();
402        assert_eq!(inserted_bank.slot(), 2);
403        assert!(bank_forks.read().unwrap().get(2).is_some());
404
405        controller.clear_bank(2).unwrap();
406        assert!(bank_forks.read().unwrap().get(2).is_none());
407
408        drop(controller);
409        replay_thread.join().unwrap();
410    }
411}