Skip to main content

solana_core/
window_service.rs

1//! `window_service` handles the data plane incoming shreds, storing them in
2//!   blockstore and retransmitting where required
3//!
4
5use {
6    crate::{
7        completed_data_sets_service::CompletedDataSetsSender,
8        repair::{
9            block_id_repair_service::{BlockIdRepairChannels, BlockIdRepairService},
10            repair_service::{
11                OutstandingShredRepairs, RepairInfo, RepairService, RepairServiceChannels,
12            },
13        },
14        result::{Error, Result},
15    },
16    agave_feature_set as feature_set,
17    crossbeam_channel::{Receiver, RecvTimeoutError, Sender, unbounded},
18    rayon::{ThreadPool, prelude::*},
19    solana_clock::Slot,
20    solana_gossip::cluster_info::ClusterInfo,
21    solana_ledger::{
22        blockstore::{Blockstore, BlockstoreInsertionMetrics, PossibleDuplicateShred},
23        blockstore_meta::BlockLocation,
24        leader_schedule_cache::LeaderScheduleCache,
25        shred::{self, ReedSolomonCache, Shred, filter::ShredRecoveryContext},
26    },
27    solana_measure::measure::Measure,
28    solana_net_utils::PinnedXdpSender,
29    solana_rayon_threadlimit::get_thread_count,
30    solana_runtime::bank_forks::{BankForks, SharableBanks},
31    solana_streamer::evicting_sender::EvictingSender,
32    std::{
33        borrow::Cow,
34        net::UdpSocket,
35        sync::{
36            Arc, RwLock,
37            atomic::{AtomicBool, AtomicUsize, Ordering},
38        },
39        thread::{self, Builder, JoinHandle},
40        time::{Duration, Instant},
41    },
42};
43
44type DuplicateSlotSender = Sender<Slot>;
45pub(crate) type DuplicateSlotReceiver = Receiver<Slot>;
46
47#[derive(Default)]
48struct WindowServiceMetrics {
49    run_insert_count: u64,
50    num_repairs: AtomicUsize,
51    num_shreds_received: usize,
52    handle_packets_elapsed_us: u64,
53    shred_receiver_elapsed_us: u64,
54    num_errors: u64,
55    num_errors_blockstore: u64,
56    num_errors_cross_beam_recv_timeout: u64,
57    num_errors_other: u64,
58    num_errors_try_crossbeam_send: u64,
59}
60
61impl WindowServiceMetrics {
62    const NAME: &str = "recv-window-insert-shreds";
63
64    fn report_metrics(&self) {
65        datapoint_info!(
66            Self::NAME,
67            (
68                "handle_packets_elapsed_us",
69                self.handle_packets_elapsed_us,
70                i64
71            ),
72            ("run_insert_count", self.run_insert_count as i64, i64),
73            ("num_repairs", self.num_repairs.load(Ordering::Relaxed), i64),
74            ("num_shreds_received", self.num_shreds_received, i64),
75            (
76                "shred_receiver_elapsed_us",
77                self.shred_receiver_elapsed_us as i64,
78                i64
79            ),
80            ("num_errors", self.num_errors, i64),
81            ("num_errors_blockstore", self.num_errors_blockstore, i64),
82            ("num_errors_other", self.num_errors_other, i64),
83            (
84                "num_errors_try_crossbeam_send",
85                self.num_errors_try_crossbeam_send,
86                i64
87            ),
88            (
89                "num_errors_cross_beam_recv_timeout",
90                self.num_errors_cross_beam_recv_timeout,
91                i64
92            ),
93        );
94    }
95
96    fn record_error(&mut self, err: &Error) {
97        self.num_errors += 1;
98        match err {
99            Error::TrySend => self.num_errors_try_crossbeam_send += 1,
100            Error::RecvTimeout(_) => self.num_errors_cross_beam_recv_timeout += 1,
101            Error::Blockstore(err) => {
102                self.num_errors_blockstore += 1;
103                error!("blockstore error: {err}");
104            }
105            _ => self.num_errors_other += 1,
106        }
107    }
108}
109
110fn run_check_duplicate(
111    cluster_info: &ClusterInfo,
112    blockstore: &Blockstore,
113    shred_receiver: &Receiver<PossibleDuplicateShred>,
114    duplicate_slots_sender: &DuplicateSlotSender,
115    bank_forks: &RwLock<BankForks>,
116) -> Result<()> {
117    let (mut root_bank, migration_status) = {
118        let bank_forks_r = bank_forks.read().unwrap();
119        (bank_forks_r.root_bank(), bank_forks_r.migration_status())
120    };
121    let mut last_updated = Instant::now();
122    let check_duplicate = |shred: PossibleDuplicateShred| -> Result<()> {
123        if last_updated.elapsed().as_nanos() > root_bank.ns_per_slot {
124            // Grabs bank forks lock once a slot
125            last_updated = Instant::now();
126            root_bank = bank_forks.read().unwrap().root_bank();
127        }
128        let shred_slot = shred.slot();
129        let validate_chained_block_id = shred::filter::check_feature_activation_from_bank(
130            &feature_set::validate_chained_block_id::id(),
131            shred_slot,
132            &root_bank,
133        );
134        let validate_chained_block_id_2 = shred::filter::check_feature_activation_from_bank(
135            &feature_set::validate_chained_block_id_2::id(),
136            shred_slot,
137            &root_bank,
138        );
139        let no_verify_chained_merkle_root = shred::filter::check_feature_activation_from_bank(
140            &feature_set::alpenglow::id(),
141            shred_slot,
142            &root_bank,
143        );
144
145        let (shred1, shred2) = match shred {
146            PossibleDuplicateShred::LastIndexConflict(shred, conflict)
147            | PossibleDuplicateShred::ErasureConflict(shred, conflict)
148            | PossibleDuplicateShred::MerkleRootConflict(shred, conflict) => (shred, conflict),
149            PossibleDuplicateShred::ChainedMerkleRootConflict(_slot) => {
150                if no_verify_chained_merkle_root {
151                    // If we're in the full alpenglow epoch, we stop validating the chained merkle root.
152                    // In Alpenglow we only use the double merkle root
153                    return Ok(());
154                }
155                if validate_chained_block_id || validate_chained_block_id_2 {
156                    // Although chained merkle roots are not necessary for agave duplicate resolution protocols,
157                    // We still need to mark the block as dead for other client teams.
158                    blockstore.set_dead_slot(shred_slot)?;
159                }
160                return Ok(());
161            }
162            PossibleDuplicateShred::FixedFECChainedMerkleRootConflict(_slot) => {
163                if no_verify_chained_merkle_root {
164                    // If we're in the full alpenglow epoch, we stop validating the chained merkle root.
165                    // In Alpenglow we only use the double merkle root
166                    return Ok(());
167                }
168                if validate_chained_block_id_2 {
169                    blockstore.set_dead_slot(shred_slot)?;
170                }
171                return Ok(());
172            }
173            PossibleDuplicateShred::Exists(shred) => {
174                // Unlike the other cases we have to wait until here to decide to handle the duplicate and store
175                // in blockstore. This is because the duplicate could have been part of the same insert batch,
176                // so we wait until the batch has been written.
177                if blockstore.has_duplicate_shreds_in_slot(shred_slot) {
178                    return Ok(()); // A duplicate is already recorded
179                }
180                let Some(existing_shred_payload) = blockstore.is_shred_duplicate(&shred) else {
181                    return Ok(()); // Not a duplicate
182                };
183                blockstore.store_duplicate_slot(
184                    shred_slot,
185                    existing_shred_payload.clone(),
186                    shred.clone().into_payload(),
187                )?;
188                (shred, shred::Payload::from(existing_shred_payload))
189            }
190        };
191
192        if migration_status.should_respond_to_ancestor_hashes_requests(shred_slot) {
193            // In alpenglow we store the duplicate block proofs in blockstore for the purposes of slashing,
194            // however we do not need to propagate the duplicate proof through gossip.
195            // We still propagate during the mixed migration epoch, to account for other nodes that are stuck
196            // and require a duplicate proof to proceed
197            cluster_info.push_duplicate_shred(&shred1, &shred2)?;
198        }
199
200        if !migration_status.is_alpenglow_enabled() {
201            // The state machine can be exited as soon as alpenglow is enabled.
202            // Notify duplicate consensus state machine. If channel is full we wait.
203            duplicate_slots_sender.send(shred_slot)?;
204        }
205
206        Ok(())
207    };
208    const RECV_TIMEOUT: Duration = Duration::from_millis(200);
209    std::iter::once(shred_receiver.recv_timeout(RECV_TIMEOUT)?)
210        .chain(shred_receiver.try_iter())
211        .try_for_each(check_duplicate)
212}
213
214#[allow(clippy::too_many_arguments)]
215fn run_insert<F>(
216    thread_pool: &ThreadPool,
217    verified_receiver: &Receiver<Vec<(shred::Payload, /*is_repaired:*/ bool, BlockLocation)>>,
218    blockstore: &Blockstore,
219    shred_recovery_context: &mut ShredRecoveryContext,
220    leader_schedule_cache: &LeaderScheduleCache,
221    handle_duplicate: F,
222    metrics: &mut BlockstoreInsertionMetrics,
223    ws_metrics: &mut WindowServiceMetrics,
224    completed_data_sets_sender: Option<&CompletedDataSetsSender>,
225) -> Result<()>
226where
227    F: Fn(PossibleDuplicateShred),
228{
229    const RECV_TIMEOUT: Duration = Duration::from_millis(200);
230    let mut shred_receiver_elapsed = Measure::start("shred_receiver_elapsed");
231    let mut shreds = verified_receiver.recv_timeout(RECV_TIMEOUT)?;
232    shreds.extend(verified_receiver.try_iter().flatten());
233    shred_receiver_elapsed.stop();
234    ws_metrics.shred_receiver_elapsed_us += shred_receiver_elapsed.as_us();
235    ws_metrics.run_insert_count += 1;
236    let handle_shred = |(shred, repair, block_location): (shred::Payload, bool, BlockLocation)| {
237        if repair {
238            ws_metrics.num_repairs.fetch_add(1, Ordering::Relaxed);
239        }
240        let shred = Shred::new_from_serialized_shred(shred).ok()?;
241        Some((Cow::Owned(shred), repair, block_location))
242    };
243    let now = Instant::now();
244    let shreds: Vec<_> = thread_pool.install(|| {
245        shreds
246            .into_par_iter()
247            .with_min_len(32)
248            .filter_map(handle_shred)
249            .collect()
250    });
251    ws_metrics.handle_packets_elapsed_us += now.elapsed().as_micros() as u64;
252    ws_metrics.num_shreds_received += shreds.len();
253    let completed_data_sets = blockstore.insert_shreds_at_location_handle_duplicate(
254        shreds,
255        Some(leader_schedule_cache),
256        false, // is_trusted
257        shred_recovery_context,
258        &handle_duplicate,
259        metrics,
260    )?;
261
262    if let Some(sender) = completed_data_sets_sender {
263        sender.try_send(completed_data_sets)?;
264    }
265
266    Ok(())
267}
268
269pub struct WindowServiceChannels {
270    pub verified_receiver: Receiver<Vec<(shred::Payload, /*is_repaired:*/ bool, BlockLocation)>>,
271    pub retransmit_sender: EvictingSender<Vec<shred::Payload>>,
272    pub completed_data_sets_sender: Option<CompletedDataSetsSender>,
273    pub duplicate_slots_sender: DuplicateSlotSender,
274    pub repair_service_channels: RepairServiceChannels,
275    pub block_id_repair_channels: BlockIdRepairChannels,
276}
277
278impl WindowServiceChannels {
279    pub fn new(
280        verified_receiver: Receiver<Vec<(shred::Payload, /*is_repaired:*/ bool, BlockLocation)>>,
281        retransmit_sender: EvictingSender<Vec<shred::Payload>>,
282        completed_data_sets_sender: Option<CompletedDataSetsSender>,
283        duplicate_slots_sender: DuplicateSlotSender,
284        repair_service_channels: RepairServiceChannels,
285        block_id_repair_channels: BlockIdRepairChannels,
286    ) -> Self {
287        Self {
288            verified_receiver,
289            retransmit_sender,
290            completed_data_sets_sender,
291            duplicate_slots_sender,
292            repair_service_channels,
293            block_id_repair_channels,
294        }
295    }
296}
297
298pub(crate) struct WindowService {
299    t_insert: JoinHandle<()>,
300    t_check_duplicate: JoinHandle<()>,
301    repair_service: RepairService,
302    block_id_repair_service: BlockIdRepairService,
303}
304
305impl WindowService {
306    #[allow(clippy::too_many_arguments)]
307    pub(crate) fn new(
308        blockstore: Arc<Blockstore>,
309        repair_socket: Arc<UdpSocket>,
310        ancestor_hashes_socket: Arc<UdpSocket>,
311        block_id_repair_socket: Arc<UdpSocket>,
312        exit: Arc<AtomicBool>,
313        repair_info: RepairInfo,
314        window_service_channels: WindowServiceChannels,
315        leader_schedule_cache: Arc<LeaderScheduleCache>,
316        shred_version: u16,
317        outstanding_repair_requests: Arc<RwLock<OutstandingShredRepairs>>,
318        repair_xdp_sender: Option<PinnedXdpSender>,
319    ) -> WindowService {
320        let cluster_info = repair_info.cluster_info.clone();
321        let bank_forks = repair_info.bank_forks.clone();
322
323        let WindowServiceChannels {
324            verified_receiver,
325            retransmit_sender,
326            completed_data_sets_sender,
327            duplicate_slots_sender,
328            repair_service_channels,
329            block_id_repair_channels,
330        } = window_service_channels;
331
332        let repair_service = RepairService::new(
333            blockstore.clone(),
334            exit.clone(),
335            repair_socket.clone(),
336            ancestor_hashes_socket,
337            repair_info.clone(),
338            outstanding_repair_requests.clone(),
339            repair_service_channels,
340            repair_xdp_sender,
341        );
342
343        let block_id_repair_service = BlockIdRepairService::new(
344            exit.clone(),
345            blockstore.clone(),
346            block_id_repair_socket,
347            repair_socket,
348            block_id_repair_channels,
349            repair_info,
350            outstanding_repair_requests,
351        );
352
353        let (duplicate_sender, duplicate_receiver) = unbounded();
354
355        let t_check_duplicate = Self::start_check_duplicate_thread(
356            cluster_info,
357            exit.clone(),
358            blockstore.clone(),
359            duplicate_receiver,
360            duplicate_slots_sender,
361            bank_forks.clone(),
362        );
363
364        let sharable_banks = bank_forks.read().unwrap().sharable_banks();
365        let t_insert = Self::start_window_insert_thread(
366            exit,
367            blockstore,
368            sharable_banks,
369            leader_schedule_cache,
370            shred_version,
371            verified_receiver,
372            duplicate_sender,
373            completed_data_sets_sender,
374            retransmit_sender,
375        );
376
377        WindowService {
378            t_insert,
379            t_check_duplicate,
380            repair_service,
381            block_id_repair_service,
382        }
383    }
384
385    fn start_check_duplicate_thread(
386        cluster_info: Arc<ClusterInfo>,
387        exit: Arc<AtomicBool>,
388        blockstore: Arc<Blockstore>,
389        duplicate_receiver: Receiver<PossibleDuplicateShred>,
390        duplicate_slots_sender: DuplicateSlotSender,
391        bank_forks: Arc<RwLock<BankForks>>,
392    ) -> JoinHandle<()> {
393        Builder::new()
394            .name("solWinCheckDup".to_string())
395            .spawn(move || {
396                while !exit.load(Ordering::Relaxed) {
397                    if let Err(e) = run_check_duplicate(
398                        &cluster_info,
399                        &blockstore,
400                        &duplicate_receiver,
401                        &duplicate_slots_sender,
402                        &bank_forks,
403                    ) && Self::should_exit_on_error(e)
404                    {
405                        break;
406                    }
407                }
408            })
409            .unwrap()
410    }
411
412    fn start_window_insert_thread(
413        exit: Arc<AtomicBool>,
414        blockstore: Arc<Blockstore>,
415        sharable_banks: SharableBanks,
416        leader_schedule_cache: Arc<LeaderScheduleCache>,
417        shred_version: u16,
418        verified_receiver: Receiver<Vec<(shred::Payload, /*is_repaired:*/ bool, BlockLocation)>>,
419        check_duplicate_sender: Sender<PossibleDuplicateShred>,
420        completed_data_sets_sender: Option<CompletedDataSetsSender>,
421        retransmit_sender: EvictingSender<Vec<shred::Payload>>,
422    ) -> JoinHandle<()> {
423        let reed_solomon_cache = ReedSolomonCache::default();
424        Builder::new()
425            .name("solWinInsert".to_string())
426            .spawn(move || {
427                let thread_pool = rayon::ThreadPoolBuilder::new()
428                    .num_threads(get_thread_count().min(8))
429                    // Use the current thread as one of the workers. This reduces overhead when the
430                    // pool is used to process a small number of shreds, since they'll be processed
431                    // directly on the current thread.
432                    .use_current_thread()
433                    .thread_name(|i| format!("solWinInsert{i:02}"))
434                    .build()
435                    .unwrap();
436                let handle_duplicate = |possible_duplicate_shred| {
437                    let _ = check_duplicate_sender.send(possible_duplicate_shred);
438                };
439
440                const METRICS_REPORTING_INTERVAL: Duration = Duration::from_secs(2);
441                let mut metrics = BlockstoreInsertionMetrics::default();
442                let mut ws_metrics = WindowServiceMetrics::default();
443                let mut last_print = Instant::now();
444                let mut shred_recovery_context = ShredRecoveryContext::new(
445                    reed_solomon_cache,
446                    retransmit_sender,
447                    sharable_banks.root(),
448                    shred_version,
449                );
450
451                while !exit.load(Ordering::Relaxed) {
452                    shred_recovery_context.maybe_update(sharable_banks.root());
453
454                    if let Err(e) = run_insert(
455                        &thread_pool,
456                        &verified_receiver,
457                        &blockstore,
458                        &mut shred_recovery_context,
459                        &leader_schedule_cache,
460                        handle_duplicate,
461                        &mut metrics,
462                        &mut ws_metrics,
463                        completed_data_sets_sender.as_ref(),
464                    ) {
465                        ws_metrics.record_error(&e);
466                        if Self::should_exit_on_error(e) {
467                            break;
468                        }
469                    }
470
471                    if last_print.elapsed() > METRICS_REPORTING_INTERVAL {
472                        metrics.report_metrics();
473                        metrics = BlockstoreInsertionMetrics::default();
474                        ws_metrics.report_metrics();
475                        ws_metrics = WindowServiceMetrics::default();
476                        last_print = Instant::now();
477                    }
478                    shred_recovery_context.maybe_submit_stats();
479                }
480            })
481            .unwrap()
482    }
483
484    fn should_exit_on_error(e: Error) -> bool {
485        match e {
486            Error::RecvTimeout(RecvTimeoutError::Disconnected) => true,
487            Error::RecvTimeout(RecvTimeoutError::Timeout) => false,
488            Error::Send => true,
489            _ => {
490                let version = solana_version::version!();
491                datapoint_error!(
492                    "error",
493                    ("thread", thread::current().name().unwrap_or("?"), String),
494                    ("message", format!("{e}"), String),
495                    ("version", version, String)
496                );
497                error!("thread {:?} error {:?}", thread::current().name(), e);
498                false
499            }
500        }
501    }
502
503    pub(crate) fn join(self) -> thread::Result<()> {
504        self.t_insert.join()?;
505        self.t_check_duplicate.join()?;
506        self.repair_service.join()?;
507        self.block_id_repair_service.join()
508    }
509}
510
511#[cfg(test)]
512mod test {
513    use {
514        super::*,
515        crossbeam_channel::bounded,
516        rand::Rng,
517        solana_entry::entry::{Entry, create_ticks},
518        solana_gossip::contact_info::ContactInfo,
519        solana_hash::Hash,
520        solana_keypair::Keypair,
521        solana_ledger::{
522            blockstore::{Blockstore, make_many_slot_entries},
523            genesis_utils::create_genesis_config,
524            get_tmp_ledger_path_auto_delete,
525            shred::{ProcessShredsStats, Shredder},
526        },
527        solana_net_utils::SocketAddrSpace,
528        solana_runtime::bank::Bank,
529        solana_signer::Signer,
530        solana_time_utils::timestamp,
531    };
532
533    fn local_entries_to_shred(
534        entries: &[Entry],
535        slot: Slot,
536        parent: Slot,
537        keypair: &Keypair,
538    ) -> Vec<Shred> {
539        let shredder = Shredder::new(slot, parent, 0, 0).unwrap();
540        let (data_shreds, _) = shredder.entries_to_merkle_shreds_for_tests(
541            keypair,
542            entries,
543            true, // is_last_in_slot
544            // chained_merkle_root
545            Hash::new_from_array(rand::rng().random()),
546            0, // next_shred_index
547            0, // next_code_index
548            &ReedSolomonCache::default(),
549            &mut ProcessShredsStats::default(),
550        );
551        data_shreds
552    }
553
554    #[test]
555    fn test_process_shred() {
556        let ledger_path = get_tmp_ledger_path_auto_delete!();
557        let blockstore = Arc::new(Blockstore::open(ledger_path.path()).unwrap());
558        let num_entries = 10;
559        let original_entries = create_ticks(num_entries, 0, Hash::default());
560        let mut shreds = local_entries_to_shred(&original_entries, 0, 0, &Keypair::new());
561        shreds.reverse();
562        blockstore
563            .insert_shreds(shreds, None, false)
564            .expect("Expect successful processing of shred");
565
566        assert_eq!(blockstore.get_slot_entries(0, 0).unwrap(), original_entries);
567    }
568
569    #[test]
570    fn test_run_check_duplicate() {
571        let ledger_path = get_tmp_ledger_path_auto_delete!();
572        let genesis_config = create_genesis_config(10_000).genesis_config;
573        let bank_forks = BankForks::new_rw_arc(Bank::new_for_tests(&genesis_config));
574        let blockstore = Arc::new(Blockstore::open(ledger_path.path()).unwrap());
575        let (sender, receiver) = bounded(1024);
576        let (duplicate_slot_sender, duplicate_slot_receiver) = bounded(1024);
577        let (shreds, _) = make_many_slot_entries(5, 5, 10);
578        blockstore
579            .insert_shreds(shreds.clone(), None, false)
580            .unwrap();
581        let duplicate_index = 0;
582        let original_shred = shreds[duplicate_index].clone();
583        let duplicate_shred = {
584            let (mut shreds, _) = make_many_slot_entries(5, 1, 10);
585            shreds.swap_remove(duplicate_index)
586        };
587        assert_eq!(duplicate_shred.slot(), shreds[0].slot());
588        let duplicate_shred_slot = duplicate_shred.slot();
589        sender
590            .send(PossibleDuplicateShred::Exists(duplicate_shred.clone()))
591            .unwrap();
592        assert!(!blockstore.has_duplicate_shreds_in_slot(duplicate_shred_slot));
593        let keypair = Keypair::new();
594        let contact_info = ContactInfo::new_localhost(&keypair.pubkey(), timestamp());
595        let cluster_info = ClusterInfo::new(
596            contact_info,
597            Arc::new(keypair),
598            SocketAddrSpace::Unspecified,
599        );
600        run_check_duplicate(
601            &cluster_info,
602            &blockstore,
603            &receiver,
604            &duplicate_slot_sender,
605            &bank_forks,
606        )
607        .unwrap();
608
609        // Make sure the correct duplicate proof was stored
610        let duplicate_proof = blockstore.get_duplicate_slot(duplicate_shred_slot).unwrap();
611        assert_eq!(duplicate_proof.shred1, *original_shred.payload());
612        assert_eq!(duplicate_proof.shred2, *duplicate_shred.payload());
613
614        // Make sure a duplicate signal was sent
615        assert_eq!(
616            duplicate_slot_receiver.try_recv().unwrap(),
617            duplicate_shred_slot
618        );
619    }
620
621    #[test]
622    fn test_store_duplicate_shreds_same_batch() {
623        let ledger_path = get_tmp_ledger_path_auto_delete!();
624        let blockstore = Arc::new(Blockstore::open(ledger_path.path()).unwrap());
625        let (duplicate_shred_sender, duplicate_shred_receiver) = bounded(1024);
626        let (duplicate_slot_sender, duplicate_slot_receiver) = bounded(1024);
627        let exit = Arc::new(AtomicBool::new(false));
628        let keypair = Keypair::new();
629        let contact_info = ContactInfo::new_localhost(&keypair.pubkey(), timestamp());
630        let cluster_info = Arc::new(ClusterInfo::new(
631            contact_info,
632            Arc::new(keypair),
633            SocketAddrSpace::Unspecified,
634        ));
635        let genesis_config = create_genesis_config(10_000).genesis_config;
636        let bank_forks = BankForks::new_rw_arc(Bank::new_for_tests(&genesis_config));
637
638        // Start duplicate thread receiving and inserting duplicates
639        let t_check_duplicate = WindowService::start_check_duplicate_thread(
640            cluster_info,
641            exit.clone(),
642            blockstore.clone(),
643            duplicate_shred_receiver,
644            duplicate_slot_sender,
645            bank_forks.clone(),
646        );
647
648        let handle_duplicate = |shred| {
649            let _ = duplicate_shred_sender.send(shred);
650        };
651        let num_trials = 100;
652        for slot in 0..num_trials {
653            let (shreds, _) = make_many_slot_entries(slot, 1, 10);
654            let duplicate_index = 0;
655            let original_shred = shreds[duplicate_index].clone();
656            let duplicate_shred = {
657                let (mut shreds, _) = make_many_slot_entries(slot, 1, 10);
658                shreds.swap_remove(duplicate_index)
659            };
660            assert_eq!(duplicate_shred.slot(), slot);
661            // Simulate storing both duplicate shreds in the same batch
662            let shreds = [&original_shred, &duplicate_shred]
663                .into_iter()
664                .map(|shred| (Cow::Borrowed(shred), /*is_repaired:*/ false));
665            let (dummy_retransmit_sender, _) = EvictingSender::new_bounded(0);
666            blockstore
667                .insert_shreds_handle_duplicate(
668                    shreds,
669                    None,
670                    false, // is_trusted
671                    &mut ShredRecoveryContext::new(
672                        ReedSolomonCache::default(),
673                        dummy_retransmit_sender,
674                        bank_forks.read().unwrap().root_bank(),
675                        0, // shred_version
676                    ),
677                    &handle_duplicate,
678                    &mut BlockstoreInsertionMetrics::default(),
679                )
680                .unwrap();
681
682            // Make sure a duplicate signal was sent
683            assert_eq!(
684                duplicate_slot_receiver
685                    .recv_timeout(Duration::from_millis(5_000))
686                    .unwrap(),
687                slot
688            );
689
690            // Make sure the correct duplicate proof was stored
691            let duplicate_proof = blockstore.get_duplicate_slot(slot).unwrap();
692            assert_eq!(duplicate_proof.shred1, *original_shred.payload());
693            assert_eq!(duplicate_proof.shred2, *duplicate_shred.payload());
694        }
695        exit.store(true, Ordering::Relaxed);
696        t_check_duplicate.join().unwrap();
697    }
698}