Skip to main content

cuprated/blockchain/
syncer.rs

1use std::{
2    future::Future,
3    sync::{
4        atomic::{AtomicU64, Ordering},
5        Arc,
6    },
7};
8
9use futures::{FutureExt, StreamExt};
10use tokio::sync::{mpsc, Notify, OwnedSemaphorePermit, Semaphore};
11use tokio_util::sync::CancellationToken;
12use tower::{Service, ServiceExt};
13use tracing::instrument;
14
15use cuprate_consensus_context::BlockchainContextService;
16use cuprate_helper::cast::usize_to_u64;
17use cuprate_p2p::{
18    block_downloader::{BlockBatch, BlockDownloaderConfig, ChainSvcRequest, ChainSvcResponse},
19    NetworkInterface, PeerSetRequest, PeerSetResponse,
20};
21use cuprate_p2p_core::{client::PeerSyncCallback, ClearNet, CoreSyncData};
22
23use super::BlockchainManagerHandle;
24use crate::monitor::FatalError;
25
26#[derive(Debug, PartialEq)]
27enum SyncStatus {
28    NoPeers,
29    BehindPeers,
30    Synced,
31}
32
33/// The syncer that makes sure we are fully synchronised with our connected peers.
34pub(crate) struct BlockchainSyncer {
35    notify_syncer: Arc<Notify>,
36    synced_tx: Option<futures::channel::oneshot::Sender<()>>,
37    target_height: Arc<AtomicU64>,
38}
39
40impl BlockchainSyncer {
41    /// Create a new [`BlockchainSyncer`] from its handle and the sender used to signal the node has synced.
42    pub(crate) fn new(
43        handle: &BlockchainSyncerHandle,
44        synced_tx: futures::channel::oneshot::Sender<()>,
45        offline: bool,
46    ) -> Self {
47        let synced_tx = if offline {
48            #[expect(clippy::let_underscore_must_use)]
49            let _ = synced_tx.send(());
50            None
51        } else {
52            Some(synced_tx)
53        };
54
55        Self {
56            notify_syncer: Arc::clone(&handle.notify_syncer),
57            synced_tx,
58            target_height: Arc::clone(&handle.target_height),
59        }
60    }
61
62    /// Run the syncer.
63    #[instrument(name = "syncer", level = "debug", skip_all)]
64    #[expect(clippy::significant_drop_tightening)]
65    #[expect(clippy::too_many_arguments)]
66    pub(crate) async fn run<CN>(
67        mut self,
68        mut context_svc: BlockchainContextService,
69        our_chain: CN,
70        mut clearnet_interface: NetworkInterface<ClearNet>,
71        incoming_block_batch_tx: mpsc::Sender<(BlockBatch, Arc<OwnedSemaphorePermit>)>,
72        stop_current_block_downloader: Arc<Notify>,
73        block_downloader_config: BlockDownloaderConfig,
74        shutdown_token: CancellationToken,
75    ) -> Result<(), FatalError>
76    where
77        CN: Service<
78                ChainSvcRequest<ClearNet>,
79                Response = ChainSvcResponse<ClearNet>,
80                Error = tower::BoxError,
81            > + Clone
82            + Send
83            + 'static,
84        CN::Future: Send + 'static,
85    {
86        tracing::info!("Starting blockchain syncer");
87        tracing::debug!("Waiting for new sync info in top sync channel");
88
89        let semaphore = Arc::new(Semaphore::new(1));
90        let mut sync_permit = Arc::new(Arc::clone(&semaphore).acquire_owned().await?);
91
92        loop {
93            tokio::select! {
94                biased;
95                () = shutdown_token.cancelled() => {
96                    tracing::info!("Blockchain syncer shut down.");
97                    return Ok(());
98                }
99                () = self.notify_syncer.notified() => {}
100            }
101
102            tracing::trace!("Checking connected peers to see if we are behind",);
103
104            match self
105                .check_sync_status(&mut context_svc, &mut clearnet_interface)
106                .await?
107            {
108                SyncStatus::BehindPeers => {}
109                SyncStatus::NoPeers => continue,
110                SyncStatus::Synced => {
111                    if let Some(synced) = self.synced_tx.take() {
112                        tracing::info!("Synchronised with the network.");
113                        #[expect(clippy::let_underscore_must_use)]
114                        let _ = synced.send(());
115                    }
116                    continue;
117                }
118            }
119
120            tracing::debug!(
121                "We are behind peers claimed cumulative difficulty, starting block downloader"
122            );
123
124            let mut block_downloader =
125                clearnet_interface.block_downloader(our_chain.clone(), block_downloader_config);
126
127            loop {
128                tokio::select! {
129                    biased;
130                    () = shutdown_token.cancelled() => break,
131                    () = stop_current_block_downloader.notified() => {
132                        tracing::info!("Received stop signal, stopping block downloader");
133
134                        drop(sync_permit);
135                        sync_permit = Arc::new(Arc::clone(&semaphore).acquire_owned().await?);
136
137                        self.notify_syncer.notify_one();
138                        break;
139                    }
140                    batch = block_downloader.stream.next() => {
141                        let Some(batch) = batch else {
142                            // Wait for all references to the permit have been dropped (which means all blocks in the queue
143                            // have been handled before checking if we are synced.
144                            drop(sync_permit);
145                            sync_permit = Arc::new(Arc::clone(&semaphore).acquire_owned().await?);
146
147                            if self.check_sync_status(&mut context_svc, &mut clearnet_interface).await? == SyncStatus::Synced {
148                                tracing::info!("Synchronised with the network.");
149                                if let Some(synced) = self.synced_tx.take() {
150                                    #[expect(clippy::let_underscore_must_use)]
151                                    let _ = synced.send(());
152                                }
153                            }
154
155                            break;
156                        };
157
158                        tracing::debug!("Got batch, len: {}", batch.blocks.len());
159                        if incoming_block_batch_tx.send((batch, Arc::clone(&sync_permit))).await.is_err() {
160                            if shutdown_token.is_cancelled() {
161                                break;
162                            }
163                            return Err("Incoming block channel closed.".into());
164                        }
165                    }
166                }
167            }
168
169            block_downloader.task.abort();
170        }
171    }
172
173    /// Checks if we are behind the connected peers.
174    async fn check_sync_status(
175        &self,
176        context_svc: &mut BlockchainContextService,
177        clearnet_interface: &mut NetworkInterface<ClearNet>,
178    ) -> Result<SyncStatus, tower::BoxError> {
179        let PeerSetResponse::MostPoWSeen {
180            cumulative_difficulty,
181            height,
182            ..
183        } = clearnet_interface
184            .peer_set()
185            .ready()
186            .await?
187            .call(PeerSetRequest::MostPoWSeen)
188            .await?
189        else {
190            unreachable!();
191        };
192
193        if cumulative_difficulty == 0 {
194            self.target_height.store(0, Ordering::Relaxed);
195            return Ok(SyncStatus::NoPeers);
196        }
197
198        if cumulative_difficulty > context_svc.blockchain_context().cumulative_difficulty {
199            self.target_height
200                .store(usize_to_u64(height), Ordering::Relaxed);
201            return Ok(SyncStatus::BehindPeers);
202        }
203
204        self.target_height.store(0, Ordering::Relaxed);
205        Ok(SyncStatus::Synced)
206    }
207}
208
209/// Handle for the `BlockchainSyncer`.
210#[derive(Clone)]
211pub struct BlockchainSyncerHandle {
212    /// The syncer notify channel, used to wake the syncer.
213    notify_syncer: Arc<Notify>,
214    /// The synced notify channel, used to wake the tasks waiting on cuprate to be synced.
215    synced: futures::future::Shared<futures::channel::oneshot::Receiver<()>>,
216    /// The target height we are syncing to, 0 if not syncing.
217    target_height: Arc<AtomicU64>,
218}
219
220impl BlockchainSyncerHandle {
221    /// Create a new handle and the sender used to signal the node has synced.
222    pub(crate) fn new() -> (Self, futures::channel::oneshot::Sender<()>) {
223        let (synced_tx, synced_rx) = futures::channel::oneshot::channel();
224
225        (
226            Self {
227                notify_syncer: Arc::new(Notify::new()),
228                synced: synced_rx.shared(),
229                target_height: Arc::new(AtomicU64::new(0)),
230            },
231            synced_tx,
232        )
233    }
234
235    /// Returns the target sync height. 0 if not syncing.
236    pub fn target_height(&self) -> u64 {
237        self.target_height.load(Ordering::Relaxed)
238    }
239
240    /// A future that resolves when cuprate has synced with the network.
241    pub fn wait_for_synced(
242        &self,
243    ) -> impl Future<Output = Result<(), futures::channel::oneshot::Canceled>> + 'static {
244        self.synced.clone()
245    }
246
247    /// Creates a [`PeerSyncCallback`] that filters and wakes the syncer.
248    pub(crate) fn callback(
249        &self,
250        context_svc: BlockchainContextService,
251        blockchain_manager: BlockchainManagerHandle,
252    ) -> PeerSyncCallback {
253        let sync_handle = self.clone();
254        let disconnect_handle = self.clone();
255
256        let on_sync = move |peer_csd: &CoreSyncData| {
257            let ctx = context_svc.blockchain_context_snapshot();
258
259            // If we are synced and the syncer hasn't yet set the node to synced, wake the syncer.
260            if peer_csd.cumulative_difficulty() == ctx.cumulative_difficulty
261                && sync_handle.synced.peek().is_none()
262            {
263                sync_handle.notify_syncer.notify_one();
264            }
265
266            // If we are behind the peer, and we aren't just one block behind with the blockchain manager handling the block, wake the syncer.
267            if peer_csd.cumulative_difficulty() > ctx.cumulative_difficulty
268                && !(peer_csd.current_height.saturating_sub(1) == ctx.chain_height as u64
269                    && blockchain_manager.is_block_being_handled(&peer_csd.top_id))
270            {
271                sync_handle.notify_syncer.notify_one();
272            }
273        };
274
275        let on_disconnect = move || disconnect_handle.notify_syncer.notify_one();
276
277        PeerSyncCallback::new(on_sync, on_disconnect)
278    }
279}