bdkffi/
kyoto.rs

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/// Receive a [`CbfClient`] and [`CbfNode`].
41#[derive(Debug, uniffi::Record)]
42pub struct CbfComponents {
43    /// Publish events to the node, like broadcasting transactions or adding scripts.
44    pub client: Arc<CbfClient>,
45    /// The node to run and fetch transactions for a [`Wallet`].
46    pub node: Arc<CbfNode>,
47}
48
49/// A [`CbfClient`] handles wallet updates from a [`CbfNode`].
50#[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/// A [`CbfNode`] gathers transactions for a [`Wallet`].
59/// To receive [`Update`] for [`Wallet`], refer to the
60/// [`CbfClient`]. The [`CbfNode`] will run until instructed
61/// to stop.
62#[derive(Debug, uniffi::Object)]
63pub struct CbfNode {
64    node: std::sync::Mutex<Option<Node>>,
65}
66
67#[uniffi::export]
68impl CbfNode {
69    /// Start the node on a detached OS thread and immediately return.
70    /// Subsequent calls have no effect.
71    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/// Build a BIP 157/158 light client to fetch transactions for a `Wallet`.
89///
90/// Options:
91/// * List of `Peer`: Bitcoin full-nodes for the light client to connect to. May be empty.
92/// * `connections`: The number of connections for the light client to maintain.
93/// * `scan_type`: Sync, recover, or start a new wallet. For more information see [`ScanType`].
94/// * `data_dir`: Optional directory to store block headers and peers.
95///
96/// A note on recovering wallets. Developers should allow users to provide an
97/// approximate recovery height and an estimated number of transactions for the
98/// wallet. When determining how many scripts to check filters for, the `Wallet`
99/// `lookahead` value will be used. To ensure all transactions are recovered, the
100/// `lookahead` should be roughly the number of transactions in the wallet history.
101#[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    /// Start a new [`CbfBuilder`]
117    #[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    /// The number of connections for the light client to maintain. Default is two.
132    pub fn connections(&self, connections: u8) -> Arc<Self> {
133        Arc::new(CbfBuilder {
134            connections,
135            ..self.clone()
136        })
137    }
138
139    /// Directory to store block headers and peers. If none is provided, the current
140    /// working directory will be used.
141    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    /// Select between syncing, recovering, or scanning for new wallets.
149    pub fn scan_type(&self, scan_type: ScanType) -> Arc<Self> {
150        Arc::new(CbfBuilder {
151            scan_type,
152            ..self.clone()
153        })
154    }
155
156    /// Bitcoin full-nodes to attempt a connection with.
157    pub fn peers(&self, peers: Vec<Peer>) -> Arc<Self> {
158        Arc::new(CbfBuilder {
159            peers,
160            ..self.clone()
161        })
162    }
163
164    /// Configure the time in milliseconds that a node has to:
165    /// 1. Respond to the initial connection
166    /// 2. Respond to a request
167    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    /// Configure connections to be established through a `Socks5 proxy. The vast majority of the
176    /// time, the connection is to a local Tor daemon, which is typically exposed at
177    /// `127.0.0.1:9050`.
178    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    /// When set, only connect to peers configured at build time. This will skip DNS and will not
186    /// use gossiped nodes.
187    pub fn only_configured_peers(&self) -> Arc<Self> {
188        Arc::new(CbfBuilder {
189            whitelist_only: true,
190            ..self.clone()
191        })
192    }
193
194    /// Construct a [`CbfComponents`] for a [`Wallet`].
195    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                // Any other network has taproot and segwit baked in since the genesis block.
211                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    /// Return the next available info message from a node. If none is returned, the node has stopped.
290    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    /// Return the next available warning message from a node. If none is returned, the node has stopped.
300    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    /// Return an [`Update`]. This is method returns once the node syncs to the rest of
310    /// the network or a new block has been gossiped.
311    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    /// Broadcast a transaction to the network, erroring if the node has stopped running.
323    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    /// The minimum fee rate required to broadcast a transcation to all connected peers.
333    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    /// Fetch the average fee rate for a block by requesting it from a peer. Not recommend for
342    /// resource-limited devices.
343    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    /// Add another [`Peer`] to attempt a connection with.
356    pub fn connect(&self, peer: Peer) -> Result<(), CbfError> {
357        self.sender
358            .add_peer(peer)
359            .map_err(|_| CbfError::NodeStopped)
360    }
361
362    /// Query a Bitcoin DNS seeder using the configured resolver.
363    ///
364    /// This is **not** a generic DNS implementation. Host names are prefixed with a `x849` to filter
365    /// for compact block filter nodes from the seeder. For example `dns.myseeder.com` will be queried
366    /// as `x849.dns.myseeder.com`. This has no guarantee to return any `IpAddr`.
367    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    /// Get the list of current connections.
383    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    /// Check if the node is still running in the background.
401    pub fn is_running(&self) -> bool {
402        self.sender.is_running()
403    }
404
405    /// Stop the [`CbfNode`]. Errors if the node is already stopped.
406    pub fn shutdown(&self) -> Result<(), CbfError> {
407        self.sender.shutdown().map_err(From::from)
408    }
409}
410
411/// A log message from the node.
412#[derive(Debug, uniffi::Enum)]
413pub enum Info {
414    /// All the required connections have been met. This is subject to change.
415    ConnectionsMet,
416    /// The node was able to successfully connect to a remote peer.
417    SuccessfulHandshake,
418    /// A percentage value of filters that have been scanned.
419    Progress {
420        /// The height of the local block chain.
421        chain_height: u32,
422        /// The percent of filters downloaded.
423        filters_downloaded_percent: f32,
424    },
425    /// A relevant block was downloaded from a peer.
426    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/// Warnings a node may issue while running.
444#[derive(Debug, uniffi::Enum)]
445pub enum Warning {
446    /// The node is looking for connections to peers.
447    NeedConnections,
448    /// A connection to a peer timed out.
449    PeerTimedOut,
450    /// The node was unable to connect to a peer in the database.
451    CouldNotConnect,
452    /// A connection was maintained, but the peer does not signal for compact block filers.
453    NoCompactFilters,
454    /// The node has been waiting for new inv and will find new peers to avoid block withholding.
455    PotentialStaleTip,
456    /// A peer sent us a peer-to-peer message the node did not request.
457    UnsolicitedMessage,
458    /// A transaction got rejected, likely for being an insufficient fee or non-standard transaction.
459    TransactionRejected {
460        wtxid: String,
461        reason: Option<String>,
462    },
463    /// The peer sent us a potential fork.
464    EvaluatingFork,
465    /// An unexpected error occurred processing a peer-to-peer message.
466    UnexpectedSyncError { warning: String },
467    /// The node failed to respond to a message sent from the client.
468    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/// Sync a wallet from the last known block hash or recover a wallet from a specified recovery
498/// point.
499#[derive(Debug, Clone, Default, uniffi::Enum)]
500pub enum ScanType {
501    /// Sync an existing wallet from the last stored chain checkpoint.
502    #[default]
503    Sync,
504    /// Recover an existing wallet by scanning from the specified height.
505    Recovery {
506        /// The estimated number of scripts the user has revealed for the wallet being recovered.
507        /// If unknown, a conservative estimate, say 1,000, could be used.
508        used_script_index: u32,
509        /// A relevant starting point or soft fork to start the sync.
510        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/// A peer to connect to over the Bitcoin peer-to-peer network.
526#[derive(Clone, uniffi::Record)]
527pub struct Peer {
528    /// The IP address to reach the node.
529    pub address: Arc<IpAddress>,
530    /// The port to reach the node. If none is provided, the default
531    /// port for the selected network will be used.
532    pub port: Option<u16>,
533    /// Does the remote node offer encrypted peer-to-peer connection.
534    pub v2_transport: bool,
535}
536
537/// An IP address to connect to over TCP.
538#[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    /// Build an IPv4 address.
553    #[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    /// Build an IPv6 address.
561    #[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/// A proxy to route network traffic, most likely through a Tor daemon. Normally this proxy is
571/// exposed at 127.0.0.1:9050.
572#[derive(Debug, Clone, uniffi::Record)]
573pub struct Socks5Proxy {
574    /// The IP address, likely `127.0.0.1`
575    pub address: Arc<IpAddress>,
576    /// The listening port, likely `9050`
577    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}