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 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#[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 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}