Skip to main content

devimint/
devfed.rs

1use std::ops::Deref as _;
2use std::sync::Arc;
3
4use anyhow::Result;
5use fedimint_core::runtime;
6use fedimint_core::task::jit::{JitTry, JitTryAnyhow};
7use fedimint_logging::LOG_DEVIMINT;
8use tokio::join;
9use tracing::{debug, info};
10
11use crate::LightningNode;
12use crate::external::{Bitcoind, Esplora, Lnd, NamedGateway, open_channels_between_gateways};
13use crate::federation::{Client, Federation};
14use crate::gatewayd::Gatewayd;
15use crate::recurringd::Recurringd;
16use crate::recurringdv2::Recurringdv2;
17use crate::util::{ProcessManager, supports_lnv2};
18
19async fn spawn_drop<T>(t: T)
20where
21    T: Send + 'static,
22{
23    runtime::spawn("spawn_drop", async {
24        drop(t);
25    })
26    .await
27    .expect("drop panic");
28}
29
30#[derive(Clone)]
31pub struct DevFed {
32    pub bitcoind: Bitcoind,
33    pub lnd: Lnd,
34    pub fed: Federation,
35    pub gw_lnd: Gatewayd,
36    pub gw_ldk: Gatewayd,
37    pub gw_ldk_second: Gatewayd,
38    pub esplora: Esplora,
39    pub recurringd: Recurringd,
40    pub recurringdv2: Recurringdv2,
41}
42
43impl DevFed {
44    pub async fn fast_terminate(self) {
45        let Self {
46            bitcoind,
47            lnd,
48            fed,
49            gw_lnd,
50            gw_ldk,
51            gw_ldk_second,
52            esplora,
53            recurringd,
54            recurringdv2,
55        } = self;
56
57        join!(
58            spawn_drop(gw_lnd),
59            spawn_drop(gw_ldk),
60            spawn_drop(gw_ldk_second),
61            spawn_drop(fed),
62            spawn_drop(lnd),
63            spawn_drop(esplora),
64            spawn_drop(bitcoind),
65            spawn_drop(recurringd),
66            spawn_drop(recurringdv2),
67        );
68    }
69}
70pub async fn dev_fed(process_mgr: &ProcessManager) -> Result<DevFed> {
71    DevJitFed::new(process_mgr, false, false)?
72        .to_dev_fed(process_mgr)
73        .await
74}
75
76type JitArc<T> = JitTryAnyhow<Arc<T>>;
77
78#[derive(Clone)]
79pub struct DevJitFed {
80    bitcoind: JitArc<Bitcoind>,
81    lnd: JitArc<Lnd>,
82    fed: JitArc<Federation>,
83    gw_lnd: JitArc<Gatewayd>,
84    gw_ldk: JitArc<Gatewayd>,
85    gw_ldk_second: JitArc<Gatewayd>,
86    esplora: JitArc<Esplora>,
87    recurringd: JitArc<Recurringd>,
88    recurringdv2: JitArc<Recurringdv2>,
89    start_time: std::time::SystemTime,
90    gw_lnd_registered: JitArc<()>,
91    gw_ldk_connected: JitArc<()>,
92    gw_ldk_second_connected: JitArc<()>,
93    fed_epoch_generated: JitArc<()>,
94    channel_opened: JitArc<()>,
95    lnv2_gateways_added: JitArc<()>,
96    recurringd_connected: JitArc<()>,
97
98    skip_setup: bool,
99    pre_dkg: bool,
100}
101
102impl DevJitFed {
103    pub fn new(process_mgr: &ProcessManager, skip_setup: bool, pre_dkg: bool) -> Result<DevJitFed> {
104        Self::new_with_pre_restore(process_mgr, skip_setup, pre_dkg, false)
105    }
106
107    pub fn new_with_pre_restore(
108        process_mgr: &ProcessManager,
109        skip_setup: bool,
110        pre_dkg: bool,
111        pre_restore: bool,
112    ) -> Result<DevJitFed> {
113        let fed_size = process_mgr.globals.FM_FED_SIZE;
114        let offline_nodes = process_mgr.globals.FM_OFFLINE_NODES;
115        anyhow::ensure!(
116            fed_size > 3 * offline_nodes,
117            "too many offline nodes ({offline_nodes}) to reach consensus"
118        );
119        let start_time = fedimint_core::time::now();
120
121        debug!(target: LOG_DEVIMINT, %fed_size, %offline_nodes, "Starting dev federation");
122
123        let bitcoind = JitTry::new_try({
124            let process_mgr = process_mgr.to_owned();
125            move || async move {
126                debug!(target: LOG_DEVIMINT, "Starting bitcoind...");
127                let start_time = fedimint_core::time::now();
128                let bitcoind = Bitcoind::new(&process_mgr, skip_setup).await?;
129                info!(target: LOG_DEVIMINT, elapsed_ms = %start_time.elapsed()?.as_millis(), "Started bitcoind");
130                Ok(Arc::new(bitcoind))
131            }
132        });
133        let lnd = JitTry::new_try({
134            let process_mgr = process_mgr.to_owned();
135            let bitcoind = bitcoind.clone();
136            || async move {
137                let bitcoind = bitcoind.get_try().await?.deref().clone();
138                debug!(target: LOG_DEVIMINT, "Starting lnd...");
139                let start_time = fedimint_core::time::now();
140                let lnd = Lnd::new(&process_mgr, bitcoind).await?;
141                info!(target: LOG_DEVIMINT, elapsed_ms = %start_time.elapsed()?.as_millis(), "Started lnd");
142                Ok(Arc::new(lnd))
143            }
144        });
145        let esplora = JitTryAnyhow::new_try({
146            let process_mgr = process_mgr.to_owned();
147            let bitcoind = bitcoind.clone();
148            || async move {
149                let bitcoind = bitcoind.get_try().await?.deref().clone();
150                debug!(target: LOG_DEVIMINT, "Starting esplora...");
151                let start_time = fedimint_core::time::now();
152                let esplora = Esplora::new(&process_mgr, bitcoind).await?;
153                info!(target: LOG_DEVIMINT, elapsed_ms = %start_time.elapsed()?.as_millis(), "Started esplora");
154                Ok(Arc::new(esplora))
155            }
156        });
157
158        let fed = JitTryAnyhow::new_try({
159            let process_mgr = process_mgr.to_owned();
160            let bitcoind = bitcoind.clone();
161            move || async move {
162                let bitcoind = bitcoind.get_try().await?.deref().clone();
163                debug!(target: LOG_DEVIMINT, "Starting federation...");
164                let start_time = fedimint_core::time::now();
165                let mut fed = Federation::new(
166                    &process_mgr,
167                    bitcoind,
168                    skip_setup,
169                    pre_dkg,
170                    pre_restore,
171                    0,
172                    "default".to_string(),
173                )
174                .await?;
175
176                // Create a degraded federation if there are offline nodes
177                fed.degrade_federation(&process_mgr).await?;
178
179                info!(target: LOG_DEVIMINT, elapsed_ms = %start_time.elapsed()?.as_millis(), "Started federation");
180
181                Ok(Arc::new(fed))
182            }
183        });
184
185        let gw_lnd = JitTryAnyhow::new_try({
186            let process_mgr = process_mgr.to_owned();
187            let lnd = lnd.clone();
188            || async move {
189                let lnd = lnd.get_try().await?.deref().clone();
190                debug!(target: LOG_DEVIMINT, "Starting lnd gateway...");
191                let start_time = fedimint_core::time::now();
192                let lnd_gw = Gatewayd::new(&process_mgr, LightningNode::Lnd(lnd), 0).await?;
193                info!(target: LOG_DEVIMINT, elapsed_ms = %start_time.elapsed()?.as_millis(), "Started lnd gateway");
194                Ok(Arc::new(lnd_gw))
195            }
196        });
197        let gw_lnd_registered = JitTryAnyhow::new_try({
198            let gw_lnd = gw_lnd.clone();
199            let fed = fed.clone();
200            move || async move {
201                let gw_lnd = gw_lnd.get_try().await?.deref();
202                let fed = fed.get_try().await?.deref();
203                debug!(target: LOG_DEVIMINT, "Registering lnd gateway...");
204                let start_time = fedimint_core::time::now();
205                if !skip_setup && !pre_dkg {
206                    let invite = fed.invite_code()?;
207                    gw_lnd.client().connect_fed(invite).await?;
208                }
209                info!(target: LOG_DEVIMINT, elapsed_ms = %start_time.elapsed()?.as_millis(), "Registered lnd gateway");
210                Ok(Arc::new(()))
211            }
212        });
213
214        let gw_ldk = JitTryAnyhow::new_try({
215            let process_mgr = process_mgr.to_owned();
216            let bitcoind = bitcoind.clone();
217            move || async move {
218                bitcoind.get_try().await?;
219                debug!(target: LOG_DEVIMINT, "Starting ldk gateway...");
220                let start_time = fedimint_core::time::now();
221                let ldk_gw = Gatewayd::new(
222                    &process_mgr,
223                    LightningNode::Ldk {
224                        name: "gatewayd-ldk-0".to_string(),
225                        gw_port: process_mgr.globals.FM_PORT_GW_LDK,
226                        ldk_port: process_mgr.globals.FM_PORT_LDK,
227                        metrics_port: process_mgr.globals.FM_PORT_GW_LDK_METRICS,
228                    },
229                    1,
230                )
231                .await?;
232                info!(target: LOG_DEVIMINT, elapsed_ms = %start_time.elapsed()?.as_millis(), "Started ldk gateway");
233                Ok(Arc::new(ldk_gw))
234            }
235        });
236        let gw_ldk_second = JitTryAnyhow::new_try({
237            let process_mgr = process_mgr.to_owned();
238            let bitcoind = bitcoind.clone();
239            move || async move {
240                bitcoind.get_try().await?;
241                debug!(target: LOG_DEVIMINT, "Starting ldk gateway 2...");
242                let start_time = fedimint_core::time::now();
243                let ldk_gw2 = Gatewayd::new(
244                    &process_mgr,
245                    LightningNode::Ldk {
246                        name: "gatewayd-ldk-1".to_string(),
247                        gw_port: process_mgr.globals.FM_PORT_GW_LDK2,
248                        ldk_port: process_mgr.globals.FM_PORT_LDK2,
249                        metrics_port: process_mgr.globals.FM_PORT_GW_LDK2_METRICS,
250                    },
251                    2,
252                )
253                .await?;
254                info!(target: LOG_DEVIMINT, elapsed_ms = %start_time.elapsed()?.as_millis(), "Started ldk gateway 2");
255                Ok(Arc::new(ldk_gw2))
256            }
257        });
258        let gw_ldk_connected = JitTryAnyhow::new_try({
259            let gw_ldk = gw_ldk.clone();
260            let fed = fed.clone();
261            move || async move {
262                let gw_ldk = gw_ldk.get_try().await?.deref();
263                if supports_lnv2() {
264                    let fed = fed.get_try().await?.deref();
265                    let start_time = fedimint_core::time::now();
266                    if !skip_setup && !pre_dkg {
267                        debug!(target: LOG_DEVIMINT, "Registering ldk gateway...");
268                        let invite = fed.invite_code()?;
269                        gw_ldk.client().connect_fed(invite).await?;
270                    } else {
271                        debug!(target: LOG_DEVIMINT, "Skipping registering ldk gateway");
272                    }
273                    info!(target: LOG_DEVIMINT, elapsed_ms = %start_time.elapsed()?.as_millis(), "Connected ldk gateway");
274                }
275                Ok(Arc::new(()))
276            }
277        });
278        let gw_ldk_second_connected = JitTryAnyhow::new_try({
279            let gw_ldk_second = gw_ldk_second.clone();
280            let fed = fed.clone();
281            move || async move {
282                let gw_ldk2 = gw_ldk_second.get_try().await?.deref();
283                if supports_lnv2() {
284                    let fed = fed.get_try().await?.deref();
285                    debug!(target: LOG_DEVIMINT, "Registering ldk gateway 2...");
286                    let start_time = fedimint_core::time::now();
287                    if !skip_setup && !pre_dkg {
288                        let invite = fed.invite_code()?;
289                        gw_ldk2.client().connect_fed(invite).await?;
290                    }
291                    info!(target: LOG_DEVIMINT, elapsed_ms = %start_time.elapsed()?.as_millis(), "Connected ldk gateway 2");
292                }
293                Ok(Arc::new(()))
294            }
295        });
296
297        let fed_epoch_generated = JitTryAnyhow::new_try({
298            let fed = fed.clone();
299            move || async move {
300                let fed = fed.get_try().await?.deref().clone();
301                debug!(target: LOG_DEVIMINT, "Generating federation epoch...");
302                let start_time = fedimint_core::time::now();
303                if !skip_setup && !pre_dkg {
304                    fed.mine_then_wait_blocks_sync(10).await?;
305                }
306                info!(target: LOG_DEVIMINT, elapsed_ms = %start_time.elapsed()?.as_millis(), "Generated federation epoch");
307                Ok(Arc::new(()))
308            }
309        });
310
311        let channel_opened = JitTryAnyhow::new_try({
312            let gw_lnd = gw_lnd.clone();
313            let gw_ldk = gw_ldk.clone();
314            let gw_ldk_second = gw_ldk_second.clone();
315            let bitcoind = bitcoind.clone();
316            let fed_epoch_generated = fed_epoch_generated.clone();
317            move || async move {
318                // Note: We open new channel even if starting from existing state
319                // as ports change on every start, and without this nodes will not find each
320                // other.
321                let bitcoind = bitcoind.get_try().await?.deref().clone();
322
323                // Wait for an epoch to occur since that mines blocks, which can cause opening
324                // channels to be racy
325                fed_epoch_generated.get_try().await?;
326
327                let gw_ldk_second = gw_ldk_second.get_try().await?.deref();
328                let gw_ldk = gw_ldk.get_try().await?.deref();
329                let gw_lnd = gw_lnd.get_try().await?.deref();
330                let gateways: &[NamedGateway<'_>] =
331                    &[(gw_ldk_second, "LDK2"), (gw_lnd, "LND"), (gw_ldk, "LDK")];
332                if !skip_setup && !pre_dkg {
333                    debug!(target: LOG_DEVIMINT, "Opening channels between gateways...");
334                    let start_time = fedimint_core::time::now();
335                    open_channels_between_gateways(&bitcoind, gateways).await?;
336                    info!(target: LOG_DEVIMINT, elapsed_ms = %start_time.elapsed()?.as_millis(), "Opened channels between gateways");
337                }
338
339                Ok(Arc::new(()))
340            }
341        });
342
343        let lnv2_gateways_added = JitTryAnyhow::new_try({
344            let fed = fed.clone();
345            let gw_lnd = gw_lnd.clone();
346            let gw_ldk = gw_ldk.clone();
347            let gw_ldk_second = gw_ldk_second.clone();
348            let gw_lnd_registered = gw_lnd_registered.clone();
349            let gw_ldk_connected = gw_ldk_connected.clone();
350            let gw_ldk_second_connected = gw_ldk_second_connected.clone();
351            let channel_opened = channel_opened.clone();
352            move || async move {
353                if supports_lnv2() && !skip_setup && !pre_dkg {
354                    gw_lnd_registered.get_try().await?;
355                    gw_ldk_connected.get_try().await?;
356                    gw_ldk_second_connected.get_try().await?;
357                    // The per-guardian admin calls each spin up a full
358                    // client, so run them serially and only after the
359                    // channels are open rather than competing with the
360                    // channel opens for resources.
361                    channel_opened.get_try().await?;
362                    let fed = fed.get_try().await?;
363                    debug!(target: LOG_DEVIMINT, "Adding gateways to the guardians' lnv2 gateway lists...");
364                    let start_time = fedimint_core::time::now();
365                    for gw in [
366                        gw_lnd.get_try().await?.deref(),
367                        gw_ldk.get_try().await?.deref(),
368                        gw_ldk_second.get_try().await?.deref(),
369                    ] {
370                        fed.add_lnv2_gateway(&gw.client().address()).await?;
371                    }
372                    info!(target: LOG_DEVIMINT, elapsed_ms = %start_time.elapsed()?.as_millis(), "Added gateways to the guardians' lnv2 gateway lists");
373                }
374                Ok(Arc::new(()))
375            }
376        });
377
378        let recurringd = JitTryAnyhow::new_try({
379            let process_mgr = process_mgr.to_owned();
380            move || async move {
381                debug!(target: LOG_DEVIMINT, "Starting recurringd...");
382                let start_time = fedimint_core::time::now();
383                let recurringd = Recurringd::new(&process_mgr).await?;
384                info!(target: LOG_DEVIMINT, elapsed_ms = %start_time.elapsed()?.as_millis(), "Started recurringd");
385                Ok(Arc::new(recurringd))
386            }
387        });
388
389        let recurringd_connected = JitTryAnyhow::new_try({
390            let recurringd = recurringd.clone();
391            let fed = fed.clone();
392            move || async move {
393                let recurringd = recurringd.get_try().await?.deref();
394                let fed = fed.get_try().await?.deref();
395                debug!(target: LOG_DEVIMINT, "Connecting recurringd to federation...");
396                let start_time = fedimint_core::time::now();
397                if !skip_setup && !pre_dkg {
398                    let invite_code = fed.invite_code()?;
399                    recurringd.add_federation(&invite_code).await?;
400                }
401                info!(target: LOG_DEVIMINT, elapsed_ms = %start_time.elapsed()?.as_millis(), "Connected recurringd to federation");
402                Ok(Arc::new(()))
403            }
404        });
405
406        let recurringdv2 = JitTryAnyhow::new_try({
407            let process_mgr = process_mgr.to_owned();
408            move || async move {
409                debug!(target: LOG_DEVIMINT, "Starting recurringdv2...");
410                let start_time = fedimint_core::time::now();
411                let recurringdv2 = Recurringdv2::new(&process_mgr).await?;
412                info!(target: LOG_DEVIMINT, elapsed_ms = %start_time.elapsed()?.as_millis(), "Started recurringdv2");
413                Ok(Arc::new(recurringdv2))
414            }
415        });
416
417        Ok(DevJitFed {
418            bitcoind,
419            lnd,
420            fed,
421            gw_lnd,
422            gw_ldk,
423            gw_ldk_second,
424            esplora,
425            recurringd,
426            recurringdv2,
427            start_time,
428            gw_lnd_registered,
429            gw_ldk_connected,
430            gw_ldk_second_connected,
431            fed_epoch_generated,
432            channel_opened,
433            lnv2_gateways_added,
434            recurringd_connected,
435            skip_setup,
436            pre_dkg,
437        })
438    }
439
440    pub async fn esplora(&self) -> anyhow::Result<&Esplora> {
441        Ok(self.esplora.get_try().await?.deref())
442    }
443    pub async fn lnd(&self) -> anyhow::Result<&Lnd> {
444        Ok(self.lnd.get_try().await?.deref())
445    }
446    pub async fn gw_lnd(&self) -> anyhow::Result<&Gatewayd> {
447        Ok(self.gw_lnd.get_try().await?.deref())
448    }
449    pub async fn gw_lnd_registered(&self) -> anyhow::Result<&Gatewayd> {
450        self.gw_lnd_registered.get_try().await?;
451        Ok(self.gw_lnd.get_try().await?.deref())
452    }
453    pub async fn gw_ldk(&self) -> anyhow::Result<&Gatewayd> {
454        Ok(self.gw_ldk.get_try().await?.deref())
455    }
456    pub async fn gw_ldk_second(&self) -> anyhow::Result<&Gatewayd> {
457        Ok(self.gw_ldk_second.get_try().await?.deref())
458    }
459    pub async fn gw_ldk_connected(&self) -> anyhow::Result<&Gatewayd> {
460        self.gw_ldk_connected.get_try().await?;
461        Ok(self.gw_ldk.get_try().await?.deref())
462    }
463    pub async fn gw_ldk_second_connected(&self) -> anyhow::Result<&Gatewayd> {
464        self.gw_ldk_second_connected.get_try().await?;
465        Ok(self.gw_ldk_second.get_try().await?.deref())
466    }
467    pub async fn fed(&self) -> anyhow::Result<&Federation> {
468        Ok(self.fed.get_try().await?.deref())
469    }
470    pub async fn bitcoind(&self) -> anyhow::Result<&Bitcoind> {
471        Ok(self.bitcoind.get_try().await?.deref())
472    }
473
474    pub async fn internal_client(&self) -> anyhow::Result<Client> {
475        Ok(self.fed().await?.internal_client().await?.clone())
476    }
477
478    /// Like [`Self::internal_client`] but will check and wait for a LN gateway
479    /// to be registered
480    pub async fn internal_client_gw_registered(&self) -> anyhow::Result<Client> {
481        self.fed().await?.await_gateways_registered().await?;
482        Ok(self.fed().await?.internal_client().await?.clone())
483    }
484
485    pub async fn recurringd(&self) -> anyhow::Result<&Recurringd> {
486        Ok(self.recurringd.get_try().await?.deref())
487    }
488
489    pub async fn recurringd_connected(&self) -> anyhow::Result<&Recurringd> {
490        self.recurringd_connected.get_try().await?;
491        Ok(self.recurringd.get_try().await?.deref())
492    }
493
494    pub async fn recurringdv2(&self) -> anyhow::Result<&Recurringdv2> {
495        Ok(self.recurringdv2.get_try().await?.deref())
496    }
497
498    pub async fn finalize(&self, process_mgr: &ProcessManager) -> anyhow::Result<()> {
499        let fed_size = process_mgr.globals.FM_FED_SIZE;
500        let offline_nodes = process_mgr.globals.FM_OFFLINE_NODES;
501        anyhow::ensure!(
502            fed_size > 3 * offline_nodes,
503            "too many offline nodes ({offline_nodes}) to reach consensus"
504        );
505
506        if !self.pre_dkg && !self.skip_setup {
507            let _ = self.gw_lnd_registered().await?;
508            let _ = self.internal_client_gw_registered().await?;
509        }
510        let _ = self.channel_opened.get_try().await?;
511        let _ = self.lnv2_gateways_added.get_try().await?;
512        let _ = self.gw_lnd_registered().await?;
513        let _ = self.gw_ldk_connected().await?;
514        let _ = self.gw_ldk_second_connected().await?;
515        let _ = self.lnd().await?;
516        let _ = self.esplora().await?;
517        let _ = self.recurringd_connected().await?;
518        let _ = self.recurringdv2().await?;
519        let _ = self.fed_epoch_generated.get_try().await?;
520
521        debug!(
522            target: LOG_DEVIMINT,
523            fed_size,
524            offline_nodes,
525            elapsed_ms = %self.start_time.elapsed()?.as_millis(),
526            "Dev federation ready",
527        );
528        Ok(())
529    }
530
531    pub async fn to_dev_fed(self, process_mgr: &ProcessManager) -> anyhow::Result<DevFed> {
532        self.finalize(process_mgr).await?;
533        Ok(DevFed {
534            bitcoind: self.bitcoind().await?.to_owned(),
535            lnd: self.lnd().await?.to_owned(),
536            fed: self.fed().await?.to_owned(),
537            gw_lnd: self.gw_lnd().await?.to_owned(),
538            gw_ldk: self.gw_ldk().await?.to_owned(),
539            gw_ldk_second: self.gw_ldk_second().await?.to_owned(),
540            esplora: self.esplora().await?.to_owned(),
541            recurringd: self.recurringd().await?.to_owned(),
542            recurringdv2: self.recurringdv2().await?.to_owned(),
543        })
544    }
545
546    pub async fn fast_terminate(self) {
547        let Self {
548            bitcoind,
549            lnd,
550            fed,
551            gw_lnd,
552            esplora,
553            gw_ldk,
554            gw_ldk_second,
555            recurringd,
556            recurringdv2,
557            ..
558        } = self;
559
560        join!(
561            spawn_drop(gw_lnd),
562            spawn_drop(gw_ldk),
563            spawn_drop(gw_ldk_second),
564            spawn_drop(fed),
565            spawn_drop(lnd),
566            spawn_drop(esplora),
567            spawn_drop(bitcoind),
568            spawn_drop(recurringd),
569            spawn_drop(recurringdv2),
570        );
571    }
572}