1use bdk_kyoto::bip157;
2use bdk_kyoto::bip157::lookup_host;
3use bdk_kyoto::bip157::tokio;
4use bdk_kyoto::bip157::AddrV2;
5use bdk_kyoto::bip157::Network;
6use bdk_kyoto::bip157::Node;
7use bdk_kyoto::bip157::ServiceFlags;
8use bdk_kyoto::builder::Builder as BDKCbfBuilder;
9use bdk_kyoto::builder::BuilderExt;
10use bdk_kyoto::HashCheckpoint;
11use bdk_kyoto::Receiver;
12use bdk_kyoto::RejectReason;
13use bdk_kyoto::Requester;
14use bdk_kyoto::TrustedPeer;
15use bdk_kyoto::UnboundedReceiver;
16use bdk_kyoto::UpdateSubscriber;
17use bdk_kyoto::Warning as Warn;
18
19use std::net::{IpAddr, Ipv4Addr, Ipv6Addr, SocketAddr};
20use std::path::PathBuf;
21use std::sync::Arc;
22use std::time::Duration;
23
24use tokio::sync::Mutex;
25
26use crate::bitcoin::BlockHash;
27use crate::bitcoin::Transaction;
28use crate::bitcoin::Wtxid;
29use crate::error::CbfError;
30use crate::types::BlockId;
31use crate::types::Update;
32use crate::wallet::Wallet;
33use crate::FeeRate;
34
35const DEFAULT_CONNECTIONS: u8 = 2;
36const CWD_PATH: &str = ".";
37const TCP_HANDSHAKE_TIMEOUT: Duration = Duration::from_secs(2);
38const MESSAGE_RESPONSE_TIMEOUT: Duration = Duration::from_secs(5);
39
40#[derive(Debug, uniffi::Record)]
42pub struct CbfComponents {
43 pub client: Arc<CbfClient>,
45 pub node: Arc<CbfNode>,
47}
48
49#[derive(Debug, uniffi::Object)]
51pub struct CbfClient {
52 sender: Arc<Requester>,
53 info_rx: Mutex<Receiver<bdk_kyoto::Info>>,
54 warning_rx: Mutex<UnboundedReceiver<bdk_kyoto::Warning>>,
55 update_rx: Mutex<UpdateSubscriber<bdk_kyoto::wallets::Single>>,
56}
57
58#[derive(Debug, uniffi::Object)]
63pub struct CbfNode {
64 node: std::sync::Mutex<Option<Node>>,
65}
66
67#[uniffi::export]
68impl CbfNode {
69 pub fn run(self: Arc<Self>) {
72 let mut lock = self.node.lock().unwrap();
73 let Some(node) = lock.take() else {
74 return;
75 };
76 std::thread::spawn(|| {
77 tokio::runtime::Builder::new_multi_thread()
78 .enable_all()
79 .build()
80 .unwrap()
81 .block_on(async move {
82 let _ = node.run().await;
83 })
84 });
85 }
86}
87
88#[derive(Clone, uniffi::Object)]
102pub struct CbfBuilder {
103 connections: u8,
104 handshake_timeout: Duration,
105 response_timeout: Duration,
106 data_dir: Option<String>,
107 scan_type: ScanType,
108 socks5_proxy: Option<Socks5Proxy>,
109 peers: Vec<Peer>,
110 whitelist_only: bool,
111}
112
113#[allow(clippy::new_without_default)]
114#[uniffi::export]
115impl CbfBuilder {
116 #[uniffi::constructor]
118 pub fn new() -> Self {
119 CbfBuilder {
120 connections: DEFAULT_CONNECTIONS,
121 handshake_timeout: TCP_HANDSHAKE_TIMEOUT,
122 response_timeout: MESSAGE_RESPONSE_TIMEOUT,
123 data_dir: None,
124 scan_type: ScanType::default(),
125 socks5_proxy: None,
126 peers: Vec::new(),
127 whitelist_only: false,
128 }
129 }
130
131 pub fn connections(&self, connections: u8) -> Arc<Self> {
133 Arc::new(CbfBuilder {
134 connections,
135 ..self.clone()
136 })
137 }
138
139 pub fn data_dir(&self, data_dir: String) -> Arc<Self> {
142 Arc::new(CbfBuilder {
143 data_dir: Some(data_dir),
144 ..self.clone()
145 })
146 }
147
148 pub fn scan_type(&self, scan_type: ScanType) -> Arc<Self> {
150 Arc::new(CbfBuilder {
151 scan_type,
152 ..self.clone()
153 })
154 }
155
156 pub fn peers(&self, peers: Vec<Peer>) -> Arc<Self> {
158 Arc::new(CbfBuilder {
159 peers,
160 ..self.clone()
161 })
162 }
163
164 pub fn configure_timeout_millis(&self, handshake: u64, response: u64) -> Arc<Self> {
168 Arc::new(CbfBuilder {
169 handshake_timeout: Duration::from_millis(handshake),
170 response_timeout: Duration::from_millis(response),
171 ..self.clone()
172 })
173 }
174
175 pub fn socks5_proxy(&self, proxy: Socks5Proxy) -> Arc<Self> {
179 Arc::new(CbfBuilder {
180 socks5_proxy: Some(proxy),
181 ..self.clone()
182 })
183 }
184
185 pub fn only_configured_peers(&self) -> Arc<Self> {
188 Arc::new(CbfBuilder {
189 whitelist_only: true,
190 ..self.clone()
191 })
192 }
193
194 pub fn build(&self, wallet: &Wallet) -> CbfComponents {
196 let wallet = wallet.get_wallet();
197
198 let mut trusted_peers = Vec::new();
199 for peer in self.peers.iter() {
200 trusted_peers.push(peer.clone().into());
201 }
202
203 let scan_type = match self.scan_type.clone() {
204 ScanType::Sync => bdk_kyoto::ScanType::Sync,
205 ScanType::Recovery {
206 used_script_index,
207 checkpoint,
208 } => {
209 let network = wallet.network();
210 if !matches!(network, Network::Bitcoin) {
212 bdk_kyoto::ScanType::Recovery {
213 used_script_index,
214 checkpoint: HashCheckpoint::from_genesis(network),
215 }
216 } else {
217 match checkpoint {
218 RecoveryPoint::GenesisBlock => bdk_kyoto::ScanType::Recovery {
219 used_script_index,
220 checkpoint: HashCheckpoint::from_genesis(wallet.network()),
221 },
222 RecoveryPoint::SegwitActivation => bdk_kyoto::ScanType::Recovery {
223 used_script_index,
224 checkpoint: HashCheckpoint::segwit_activation(),
225 },
226 RecoveryPoint::TaprootActivation => bdk_kyoto::ScanType::Recovery {
227 used_script_index,
228 checkpoint: HashCheckpoint::taproot_activation(),
229 },
230 RecoveryPoint::Other { birthday } => bdk_kyoto::ScanType::Recovery {
231 used_script_index,
232 checkpoint: HashCheckpoint::new(birthday.height, birthday.hash.0),
233 },
234 }
235 }
236 }
237 };
238
239 let path_buf = self
240 .data_dir
241 .clone()
242 .map(|path| PathBuf::from(&path))
243 .unwrap_or(PathBuf::from(CWD_PATH));
244
245 let mut builder = BDKCbfBuilder::new(wallet.network())
246 .required_peers(self.connections)
247 .data_dir(path_buf)
248 .handshake_timeout(self.handshake_timeout)
249 .response_timeout(self.response_timeout)
250 .add_peers(trusted_peers);
251
252 if let Some(proxy) = &self.socks5_proxy {
253 let port = proxy.port;
254 let addr = proxy.address.inner;
255 builder = builder.socks5_proxy(SocketAddr::new(addr, port));
256 }
257
258 if self.whitelist_only {
259 builder = builder.whitelist_only();
260 }
261
262 let (client, logging, update_subscriber) = builder
263 .build_with_wallet(&wallet, scan_type)
264 .expect("networks match by definition")
265 .subscribe();
266 let (client, node) = client.managed_start();
267 let requester = client.requester();
268
269 let node = CbfNode {
270 node: std::sync::Mutex::new(Some(node)),
271 };
272
273 let client = CbfClient {
274 sender: Arc::new(requester),
275 info_rx: Mutex::new(logging.info_subscriber),
276 warning_rx: Mutex::new(logging.warning_subscriber),
277 update_rx: Mutex::new(update_subscriber),
278 };
279
280 CbfComponents {
281 client: Arc::new(client),
282 node: Arc::new(node),
283 }
284 }
285}
286
287#[uniffi::export]
288impl CbfClient {
289 pub async fn next_info(&self) -> Result<Info, CbfError> {
291 let mut info_rx = self.info_rx.lock().await;
292 info_rx
293 .recv()
294 .await
295 .map(|e| e.into())
296 .ok_or(CbfError::NodeStopped)
297 }
298
299 pub async fn next_warning(&self) -> Result<Warning, CbfError> {
301 let mut warn_rx = self.warning_rx.lock().await;
302 warn_rx
303 .recv()
304 .await
305 .map(|warn| warn.into())
306 .ok_or(CbfError::NodeStopped)
307 }
308
309 pub async fn update(&self) -> Result<Update, CbfError> {
312 let update = self
313 .update_rx
314 .lock()
315 .await
316 .update()
317 .await
318 .map_err(|_| CbfError::NodeStopped)?;
319 Ok(Update(update))
320 }
321
322 pub async fn broadcast(&self, transaction: &Transaction) -> Result<Arc<Wtxid>, CbfError> {
324 let tx: bip157::Transaction = transaction.into();
325 self.sender
326 .submit_package(tx)
327 .await
328 .map_err(From::from)
329 .map(|wtxid| Arc::new(Wtxid(wtxid)))
330 }
331
332 pub async fn min_broadcast_feerate(&self) -> Result<Arc<FeeRate>, CbfError> {
334 self.sender
335 .broadcast_min_feerate()
336 .await
337 .map_err(|_| CbfError::NodeStopped)
338 .map(|fee| Arc::new(FeeRate(fee)))
339 }
340
341 pub async fn average_fee_rate(
344 &self,
345 blockhash: Arc<BlockHash>,
346 ) -> Result<Arc<FeeRate>, CbfError> {
347 let fee_rate = self
348 .sender
349 .average_fee_rate(blockhash.0)
350 .await
351 .map_err(|_| CbfError::NodeStopped)?;
352 Ok(Arc::new(fee_rate.into()))
353 }
354
355 pub fn connect(&self, peer: Peer) -> Result<(), CbfError> {
357 self.sender
358 .add_peer(peer)
359 .map_err(|_| CbfError::NodeStopped)
360 }
361
362 pub fn lookup_host(&self, hostname: String) -> Vec<Arc<IpAddress>> {
368 let node_handle = std::thread::spawn(move || {
369 tokio::runtime::Builder::new_current_thread()
370 .enable_all()
371 .build()
372 .unwrap()
373 .block_on(lookup_host(hostname))
374 });
375 let nodes = node_handle.join().unwrap_or_default();
376 nodes
377 .into_iter()
378 .map(|ip| Arc::new(IpAddress { inner: ip }))
379 .collect()
380 }
381
382 pub async fn peer_info(&self) -> Result<Vec<Arc<IpAddress>>, CbfError> {
384 let peers = self
385 .sender
386 .peer_info()
387 .await
388 .map_err(|_| CbfError::NodeStopped)?;
389 Ok(peers
390 .into_iter()
391 .filter_map(|(ip, _)| match ip {
392 AddrV2::Ipv4(ip) => Some(IpAddr::V4(ip)),
393 AddrV2::Ipv6(ip) => Some(IpAddr::V6(ip)),
394 _ => None,
395 })
396 .map(|ip| Arc::new(IpAddress { inner: ip }))
397 .collect())
398 }
399
400 pub fn is_running(&self) -> bool {
402 self.sender.is_running()
403 }
404
405 pub fn shutdown(&self) -> Result<(), CbfError> {
407 self.sender.shutdown().map_err(From::from)
408 }
409}
410
411#[derive(Debug, uniffi::Enum)]
413pub enum Info {
414 ConnectionsMet,
416 SuccessfulHandshake,
418 Progress {
420 chain_height: u32,
422 filters_downloaded_percent: f32,
424 },
425 BlockReceived(String),
427}
428
429impl From<bdk_kyoto::Info> for Info {
430 fn from(value: bdk_kyoto::Info) -> Info {
431 match value {
432 bdk_kyoto::Info::ConnectionsMet => Info::ConnectionsMet,
433 bdk_kyoto::Info::SuccessfulHandshake => Info::SuccessfulHandshake,
434 bdk_kyoto::Info::Progress(progress) => Info::Progress {
435 filters_downloaded_percent: progress.percentage_complete(),
436 chain_height: progress.chain_height(),
437 },
438 bdk_kyoto::Info::BlockReceived(block) => Info::BlockReceived(block.to_string()),
439 }
440 }
441}
442
443#[derive(Debug, uniffi::Enum)]
445pub enum Warning {
446 NeedConnections,
448 PeerTimedOut,
450 CouldNotConnect,
452 NoCompactFilters,
454 PotentialStaleTip,
456 UnsolicitedMessage,
458 TransactionRejected {
460 wtxid: String,
461 reason: Option<String>,
462 },
463 EvaluatingFork,
465 UnexpectedSyncError { warning: String },
467 RequestFailed,
469}
470
471impl From<Warn> for Warning {
472 fn from(value: Warn) -> Warning {
473 match value {
474 Warn::NeedConnections {
475 connected: _,
476 required: _,
477 } => Warning::NeedConnections,
478 Warn::PeerTimedOut => Warning::PeerTimedOut,
479 Warn::CouldNotConnect => Warning::CouldNotConnect,
480 Warn::NoCompactFilters => Warning::NoCompactFilters,
481 Warn::PotentialStaleTip => Warning::PotentialStaleTip,
482 Warn::UnsolicitedMessage => Warning::UnsolicitedMessage,
483 Warn::TransactionRejected { payload } => {
484 let reason = payload.reason.map(|r| r.into_string());
485 Warning::TransactionRejected {
486 wtxid: payload.wtxid.to_string(),
487 reason,
488 }
489 }
490 Warn::EvaluatingFork => Warning::EvaluatingFork,
491 Warn::UnexpectedSyncError { warning } => Warning::UnexpectedSyncError { warning },
492 Warn::ChannelDropped => Warning::RequestFailed,
493 }
494 }
495}
496
497#[derive(Debug, Clone, Default, uniffi::Enum)]
500pub enum ScanType {
501 #[default]
503 Sync,
504 Recovery {
506 used_script_index: u32,
509 checkpoint: RecoveryPoint,
511 },
512}
513
514#[derive(Debug, Clone, Default, uniffi::Enum)]
515pub enum RecoveryPoint {
516 GenesisBlock,
517 #[default]
518 SegwitActivation,
519 TaprootActivation,
520 Other {
521 birthday: BlockId,
522 },
523}
524
525#[derive(Clone, uniffi::Record)]
527pub struct Peer {
528 pub address: Arc<IpAddress>,
530 pub port: Option<u16>,
533 pub v2_transport: bool,
535}
536
537#[derive(Debug, uniffi::Object)]
539#[uniffi::export(Display)]
540pub struct IpAddress {
541 inner: IpAddr,
542}
543
544impl core::fmt::Display for IpAddress {
545 fn fmt(&self, f: &mut core::fmt::Formatter) -> core::fmt::Result {
546 write!(f, "{}", self.inner)
547 }
548}
549
550#[uniffi::export]
551impl IpAddress {
552 #[uniffi::constructor]
554 pub fn from_ipv4(q1: u8, q2: u8, q3: u8, q4: u8) -> Self {
555 Self {
556 inner: IpAddr::V4(Ipv4Addr::new(q1, q2, q3, q4)),
557 }
558 }
559
560 #[allow(clippy::too_many_arguments)]
562 #[uniffi::constructor]
563 pub fn from_ipv6(a: u16, b: u16, c: u16, d: u16, e: u16, f: u16, g: u16, h: u16) -> Self {
564 Self {
565 inner: IpAddr::V6(Ipv6Addr::new(a, b, c, d, e, f, g, h)),
566 }
567 }
568}
569
570#[derive(Debug, Clone, uniffi::Record)]
573pub struct Socks5Proxy {
574 pub address: Arc<IpAddress>,
576 pub port: u16,
578}
579
580impl From<Peer> for TrustedPeer {
581 fn from(peer: Peer) -> Self {
582 let services = if peer.v2_transport {
583 let mut services = ServiceFlags::P2P_V2;
584 services.add(ServiceFlags::NETWORK);
585 services.add(ServiceFlags::COMPACT_FILTERS);
586 services
587 } else {
588 let mut services = ServiceFlags::COMPACT_FILTERS;
589 services.add(ServiceFlags::NETWORK);
590 services
591 };
592 let addr_v2 = match peer.address.inner {
593 IpAddr::V4(ipv4_addr) => AddrV2::Ipv4(ipv4_addr),
594 IpAddr::V6(ipv6_addr) => AddrV2::Ipv6(ipv6_addr),
595 };
596 TrustedPeer::new(addr_v2, peer.port, services)
597 }
598}
599
600trait DisplayExt {
601 fn into_string(self) -> String;
602}
603
604impl DisplayExt for RejectReason {
605 fn into_string(self) -> String {
606 let message = match self {
607 RejectReason::Malformed => "Message could not be decoded.",
608 RejectReason::Invalid => "Transaction was invalid for some reason.",
609 RejectReason::Obsolete => "Client version is no longer supported.",
610 RejectReason::Duplicate => "Duplicate version message received.",
611 RejectReason::NonStandard => "Transaction was nonstandard.",
612 RejectReason::Dust => "One or more outputs are below the dust threshold.",
613 RejectReason::Fee => "Transaction does not have enough fee to be mined.",
614 RejectReason::Checkpoint => "Inconsistent with compiled checkpoint.",
615 };
616 message.into()
617 }
618}
619
620#[cfg(test)]
621mod tests {
622 use super::CbfNode;
623 use std::sync::{Arc, Mutex};
624
625 #[test]
626 fn running_a_consumed_node_is_a_no_op() {
627 let node = Arc::new(CbfNode {
628 node: Mutex::new(None),
629 });
630
631 node.run();
632 }
633}