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 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 let bitcoind = bitcoind.get_try().await?.deref().clone();
322
323 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 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 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}