cuprated/blockchain/
syncer.rs1use 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
33pub(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 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 #[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 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 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#[derive(Clone)]
211pub struct BlockchainSyncerHandle {
212 notify_syncer: Arc<Notify>,
214 synced: futures::future::Shared<futures::channel::oneshot::Receiver<()>>,
216 target_height: Arc<AtomicU64>,
218}
219
220impl BlockchainSyncerHandle {
221 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 pub fn target_height(&self) -> u64 {
237 self.target_height.load(Ordering::Relaxed)
238 }
239
240 pub fn wait_for_synced(
242 &self,
243 ) -> impl Future<Output = Result<(), futures::channel::oneshot::Canceled>> + 'static {
244 self.synced.clone()
245 }
246
247 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 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 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}