diff --git a/Cargo.lock b/Cargo.lock index 9d95793..ab39c20 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2400,6 +2400,27 @@ dependencies = [ "tracing", ] +[[package]] +name = "fedimint-lnv2-common" +version = "0.12.0" +source = "git+https://github.com/fedimint/fedimint?tag=v0.12.0#2ead392ba22141bbe5898cdb724ec3dbc55387a9" +dependencies = [ + "anyhow", + "async-trait", + "bitcoin 0.32.102", + "fedimint-connectors", + "fedimint-core 0.12.0", + "fedimint-ln-common", + "fedimint-tpe", + "group", + "lightning-invoice 0.33.3", + "rand 0.8.8", + "reqwest 0.12.28", + "serde", + "serde_json", + "thiserror 2.0.20", +] + [[package]] name = "fedimint-logging" version = "0.4.4" @@ -2519,6 +2540,21 @@ dependencies = [ "zeroize", ] +[[package]] +name = "fedimint-tpe" +version = "0.12.0" +source = "git+https://github.com/fedimint/fedimint?tag=v0.12.0#2ead392ba22141bbe5898cdb724ec3dbc55387a9" +dependencies = [ + "bitcoin_hashes 0.14.101", + "bls12_381", + "fedimint-core 0.12.0", + "group", + "rand 0.8.8", + "rand_chacha 0.3.1", + "serde", + "serde-big-array", +] + [[package]] name = "fedimint-util-error" version = "0.12.0" @@ -2613,6 +2649,7 @@ dependencies = [ "fedimint-connectors", "fedimint-core 0.12.0", "fedimint-ln-common", + "fedimint-lnv2-common", "fedimint-meta-client", "fedimint-mint-common", "fedimint-wallet-common", diff --git a/Cargo.toml b/Cargo.toml index be39e51..c9ddbce 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -11,6 +11,7 @@ fedimint-api-client = { version = "0.12.0", git = "https://github.com/fedimint/f fedimint-connectors = { version = "0.12.0", git = "https://github.com/fedimint/fedimint", tag = "v0.12.0" } fedimint-core = { version = "0.12.0", git = "https://github.com/fedimint/fedimint", tag = "v0.12.0" } fedimint-ln-common = { version = "0.12.0", git = "https://github.com/fedimint/fedimint", tag = "v0.12.0" } +fedimint-lnv2-common = { version = "0.12.0", git = "https://github.com/fedimint/fedimint", tag = "v0.12.0" } fedimint-meta-client = { version = "0.12.0", git = "https://github.com/fedimint/fedimint", tag = "v0.12.0" } fedimint-mint-common = { version = "0.12.0", git = "https://github.com/fedimint/fedimint", tag = "v0.12.0" } fedimint-wallet-common = { version = "0.12.0", git = "https://github.com/fedimint/fedimint", tag = "v0.12.0" } diff --git a/fmo_api_types/src/lib.rs b/fmo_api_types/src/lib.rs index 416d080..f392914 100644 --- a/fmo_api_types/src/lib.rs +++ b/fmo_api_types/src/lib.rs @@ -111,6 +111,65 @@ pub struct GatewayInfo { /// `7d` #[serde(skip_serializing_if = "Option::is_none")] pub metrics_window: Option, + /// Lightning module protocols this gateway is registered under. A gateway + /// appears under both only if its identity was verified to match. + #[serde(skip_serializing_if = "Option::is_none")] + pub protocols: Option>, + /// LNv2 registration details, present if the gateway is registered with + /// the federation's LNv2 module + #[serde(skip_serializing_if = "Option::is_none")] + pub lnv2: Option, +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)] +#[serde(rename_all = "lowercase")] +pub enum GatewayProtocol { + Lnv1, + Lnv2, +} + +/// Result of asking an LNv2 gateway for its routing info. `Ok` only means the +/// gateway advertises routing info for this federation, not that payments +/// through it succeed. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum Lnv2RoutingStatus { + /// The gateway returned routing info for this federation + Ok, + /// The gateway answered but does not serve this federation + NotServing, + /// The gateway could not be reached or returned an invalid response + Unreachable, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct Lnv2GatewayInfo { + /// Gateway API URL as listed in the federation's LNv2 registry + pub api_endpoint: String, + pub routing_status: Lnv2RoutingStatus, + #[serde(skip_serializing_if = "Option::is_none")] + pub routing_checked_at: Option>, + /// LN node public key (hex-encoded), from the gateway's routing info + #[serde(skip_serializing_if = "Option::is_none")] + pub lightning_public_key: Option, + /// Key the gateway uses to claim and refund LNv2 contracts (hex-encoded) + #[serde(skip_serializing_if = "Option::is_none")] + pub module_public_key: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub send_fee_minimum: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub send_fee_default: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub receive_fee: Option, + /// Registry presence of the LNv2 registration over the requested window + #[serde(skip_serializing_if = "Option::is_none")] + pub uptime_window: Option, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct GatewayFee { + pub base_msat: u64, + pub parts_per_million: u64, } #[derive(Debug, Clone, Serialize, Deserialize)] diff --git a/fmo_frontend_react/src/types/api.ts b/fmo_frontend_react/src/types/api.ts index 25f0391..6b4de23 100644 --- a/fmo_frontend_react/src/types/api.ts +++ b/fmo_frontend_react/src/types/api.ts @@ -37,6 +37,31 @@ export interface GatewayInfo { activity_window?: GatewayActivityMetrics; uptime_window?: GatewayUptimeMetrics; metrics_window?: GatewayWindow; + protocols?: GatewayProtocol[]; + lnv2?: Lnv2GatewayInfo; +} + +export type GatewayProtocol = 'lnv1' | 'lnv2'; + +// 'ok' only means the gateway advertises routing info for this federation, +// not that payments through it succeed. +export type Lnv2RoutingStatus = 'ok' | 'not_serving' | 'unreachable'; + +export interface Lnv2GatewayInfo { + api_endpoint: string; + routing_status: Lnv2RoutingStatus; + routing_checked_at?: string; + lightning_public_key?: string; + module_public_key?: string; + send_fee_minimum?: GatewayFee; + send_fee_default?: GatewayFee; + receive_fee?: GatewayFee; + uptime_window?: GatewayUptimeMetrics; +} + +export interface GatewayFee { + base_msat: number; + parts_per_million: number; } export interface GatewayActivityMetrics { diff --git a/fmo_server/Cargo.toml b/fmo_server/Cargo.toml index c5a1f79..45da9f4 100644 --- a/fmo_server/Cargo.toml +++ b/fmo_server/Cargo.toml @@ -10,6 +10,7 @@ fedimint-api-client = { workspace = true } fedimint-connectors = { workspace = true } fedimint-core = { workspace = true } fedimint-ln-common = { workspace = true } +fedimint-lnv2-common = { workspace = true } fedimint-meta-client = { workspace = true } fedimint-mint-common = { workspace = true } fedimint-wallet-common = { workspace = true } diff --git a/fmo_server/schema/v11.sql b/fmo_server/schema/v11.sql new file mode 100644 index 0000000..da485a3 --- /dev/null +++ b/fmo_server/schema/v11.sql @@ -0,0 +1,53 @@ +-- LNv2 gateway discovery: track which Lightning module protocol a gateway is registered under +BEGIN; + +INSERT INTO schema_version (version) +VALUES (11); + +-- Existing rows all came from the LNv1 (`ln`) module registry. LNv2 rows use the normalized +-- gateway API URL as `gateway_id`, since the LNv2 registry only lists URLs. +ALTER TABLE gateways + ADD COLUMN IF NOT EXISTS protocol TEXT NOT NULL DEFAULT 'lnv1'; + +ALTER TABLE gateways DROP CONSTRAINT gateways_pkey; +ALTER TABLE gateways ADD PRIMARY KEY (federation_id, protocol, gateway_id); + +-- LNv2 identity comes from probing the gateway, which may fail +ALTER TABLE gateways ALTER COLUMN node_pub_key DROP NOT NULL; +ALTER TABLE gateways ALTER COLUMN lightning_alias DROP NOT NULL; + +-- LNv2 only: result of the most recent `routing_info` probe +-- 'ok' - the gateway returned routing info for this federation +-- 'not_serving' - the gateway answered but does not serve this federation +-- 'unreachable' - the probe failed or timed out +ALTER TABLE gateways + ADD COLUMN IF NOT EXISTS module_public_key TEXT, + ADD COLUMN IF NOT EXISTS routing_status TEXT, + ADD COLUMN IF NOT EXISTS routing_checked_at TIMESTAMPTZ; + +ALTER TABLE gateway_poll_snapshots + ADD COLUMN IF NOT EXISTS protocol TEXT NOT NULL DEFAULT 'lnv1'; + +ALTER TABLE gateway_poll_snapshots DROP CONSTRAINT gateway_poll_snapshots_pkey; +ALTER TABLE gateway_poll_snapshots ADD PRIMARY KEY (federation_id, protocol, gateway_id, poll_time); + +-- `is_seen` remains registry presence; `routing_ok` is LNv2 routing readiness (NULL for LNv1) +ALTER TABLE gateway_poll_snapshots + ADD COLUMN IF NOT EXISTS routing_ok BOOLEAN; + +-- Every LNv2 module key ever observed for a gateway, so that contracts can still be +-- attributed after a gateway rotates its key or leaves the registry. +CREATE TABLE IF NOT EXISTS lnv2_gateway_keys ( + federation_id BYTEA NOT NULL REFERENCES federations (federation_id), + gateway_id TEXT NOT NULL, + module_public_key TEXT NOT NULL, + lightning_public_key TEXT NOT NULL, + first_seen TIMESTAMPTZ NOT NULL, + last_seen TIMESTAMPTZ NOT NULL, + PRIMARY KEY (federation_id, gateway_id, module_public_key) +); + +CREATE INDEX IF NOT EXISTS lnv2_gateway_keys_module_public_key + ON lnv2_gateway_keys (federation_id, module_public_key); + +COMMIT; diff --git a/fmo_server/src/config/mod.rs b/fmo_server/src/config/mod.rs index 7d72d83..dbc7a3f 100644 --- a/fmo_server/src/config/mod.rs +++ b/fmo_server/src/config/mod.rs @@ -1,13 +1,14 @@ use std::collections::HashMap; use std::sync::Arc; -use axum::extract::{Path, State}; +use axum::extract::{Path, Query, State}; use axum::routing::get; use axum::{Json, Router}; use fedimint_core::config::{FederationId, JsonClientConfig}; use fedimint_core::invite_code::InviteCode; use fmo_api_types::GatewayInfo; use reqwest::Method; +use serde::Deserialize; use tower_http::cors::{Any, CorsLayer}; use tracing::warn; @@ -59,15 +60,33 @@ pub async fn fetch_federation_config( .into()) } +#[derive(Debug, Deserialize)] +pub struct FetchFederationGatewaysParams { + /// `all` to include LNv2 gateways. Omitted, only LNv1 gateways are + /// returned, exactly as before LNv2 support existed. + protocols: Option, +} + pub async fn fetch_federation_gateways( Path(invite): Path, + Query(params): Query, ) -> Result>> { + let include_lnv2 = match params.protocols.as_deref() { + None => false, + Some("all") => true, + Some(invalid) => { + return Err( + anyhow::anyhow!("Invalid protocols '{invalid}'. Supported values: all").into(), + ) + } + }; + let connectors = fedimint_connectors::ConnectorRegistry::build_from_client_env()? .bind() .await?; let (config, _api) = fedimint_api_client::download_from_invite_code(&connectors, &invite).await?; - let gateways = fetch_gateways_for_config(&config).await?; + let gateways = fetch_gateways_for_config(&config, include_lnv2).await?; Ok(gateways.into()) } diff --git a/fmo_server/src/federation/gateways.rs b/fmo_server/src/federation/gateways.rs index 03e989f..78b458b 100644 --- a/fmo_server/src/federation/gateways.rs +++ b/fmo_server/src/federation/gateways.rs @@ -1,23 +1,34 @@ -use std::collections::HashMap; +use std::collections::{BTreeMap, HashMap}; use std::time::Duration; -use anyhow::{bail, Context}; +use anyhow::bail; use axum::extract::{Path, Query, State}; use axum::Json; use chrono::{DateTime, Utc}; +use deadpool_postgres::Transaction; use fedimint_api_client::api::{DynGlobalApi, FederationApiExt}; use fedimint_core::config::{ClientConfig, FederationId}; -use fedimint_core::core::ModuleInstanceId; +use fedimint_core::core::{ModuleInstanceId, ModuleKind}; use fedimint_core::encoding::Encodable; use fedimint_core::module::ApiRequestErased; +use fedimint_core::util::SafeUrl; +use fedimint_core::PeerId; +use fedimint_ln_common::client::GatewayApi; use fedimint_ln_common::federation_endpoint_constants::LIST_GATEWAYS_ENDPOINT; use fedimint_ln_common::LightningGatewayAnnouncement; +use fedimint_lnv2_common::endpoint_constants::GATEWAYS_ENDPOINT; +use fedimint_lnv2_common::gateway_api::{ + GatewayConnection, PaymentFee, RealGatewayConnection, RoutingInfo, +}; use fmo_api_types::{ - GatewayActivityMetrics, GatewayInfo, GatewayUptimeMetrics, GatewayUptimeTrendPoint, + GatewayActivityMetrics, GatewayFee, GatewayInfo, GatewayProtocol, GatewayUptimeMetrics, + GatewayUptimeTrendPoint, Lnv2GatewayInfo, Lnv2RoutingStatus, }; use futures::future::join_all; +use futures::StreamExt; +use serde::de::DeserializeOwned; use serde::Deserialize; -use tracing::{info, warn}; +use tracing::{debug, info, warn}; use crate::federation::observer::FederationObserver; use crate::util::query; @@ -25,6 +36,8 @@ use crate::util::query; const GATEWAY_POLL_INTERVAL_MINUTES: u64 = 5; const GATEWAY_SNAPSHOT_RETENTION_DAYS: i64 = 90; const GATEWAY_PRUNE_INTERVAL_HOURS: i64 = 6; +const LNV2_PROBE_TIMEOUT: Duration = Duration::from_secs(5); +const LNV2_PROBE_CONCURRENCY: usize = 8; #[derive(Debug, Clone, Copy, Eq, PartialEq)] enum GatewayMetricsWindow { @@ -75,9 +88,363 @@ pub(super) struct GetFederationGatewaysParams { window: Option, } +/// Instance ids of the Lightning modules a federation runs. Either, both or +/// neither may be present. +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] +struct LnInstances { + lnv1: Option, + lnv2: Option, +} + +impl LnInstances { + fn from_config(config: &ClientConfig) -> Self { + Self::from_kinds( + config + .modules + .iter() + .map(|(&instance_id, module)| (instance_id, &module.kind)), + ) + } + + fn from_kinds<'a>( + modules: impl IntoIterator, + ) -> Self { + let mut instances = Self::default(); + for (instance_id, kind) in modules { + if *kind == fedimint_ln_common::KIND { + instances.lnv1.get_or_insert(instance_id); + } else if *kind == fedimint_lnv2_common::KIND { + instances.lnv2.get_or_insert(instance_id); + } + } + instances + } + + fn is_empty(&self) -> bool { + self.lnv1.is_none() && self.lnv2.is_none() + } +} + +/// An LNv2 registry entry together with what the gateway itself told us. +#[derive(Debug, Clone)] +struct Lnv2Gateway { + /// Normalized API URL, see [`normalize_gateway_url`] + gateway_id: String, + /// API URL exactly as listed in the registry + api_endpoint: SafeUrl, + probe: Lnv2Probe, +} + +#[derive(Debug, Clone)] +enum Lnv2Probe { + Ok(Box), + NotServing, + Unreachable, +} + +impl Lnv2Probe { + fn status(&self) -> Lnv2RoutingStatus { + match self { + Self::Ok(_) => Lnv2RoutingStatus::Ok, + Self::NotServing => Lnv2RoutingStatus::NotServing, + Self::Unreachable => Lnv2RoutingStatus::Unreachable, + } + } + + fn routing_info(&self) -> Option<&RoutingInfo> { + match self { + Self::Ok(info) => Some(info), + Self::NotServing | Self::Unreachable => None, + } + } +} + +/// Identity key for a gateway API URL. +/// +/// Gateways serve their API both at the root and under `/v1`, and guardians +/// enter LNv2 registry URLs by hand, so a trailing slash or `/v1` must not make +/// the same gateway look like two. Scheme and host are already lowercased by +/// URL parsing. +pub(crate) fn normalize_gateway_url(url: &str) -> String { + let url = url.trim_end_matches('/'); + url.strip_suffix("/v1") + .unwrap_or(url) + .trim_end_matches('/') + .to_owned() +} + +fn routing_status_str(status: Lnv2RoutingStatus) -> &'static str { + match status { + Lnv2RoutingStatus::Ok => "ok", + Lnv2RoutingStatus::NotServing => "not_serving", + Lnv2RoutingStatus::Unreachable => "unreachable", + } +} + +fn parse_routing_status(status: &str) -> Option { + match status { + "ok" => Some(Lnv2RoutingStatus::Ok), + "not_serving" => Some(Lnv2RoutingStatus::NotServing), + "unreachable" => Some(Lnv2RoutingStatus::Unreachable), + _ => None, + } +} + +fn gateway_fee(fee: &PaymentFee) -> GatewayFee { + GatewayFee { + base_msat: fee.base.msats, + parts_per_million: fee.parts_per_million, + } +} + +/// Query one module endpoint on every peer and concatenate the answers. Each +/// guardian keeps its own gateway registry, so the union is what matters. +/// Returns `None` if no peer answered, i.e. registry contents are unknown. +async fn fetch_registry( + api: &DynGlobalApi, + instance_id: ModuleInstanceId, + endpoint: &str, + peer_ids: &[PeerId], +) -> Option> { + let peer_results = join_all(peer_ids.iter().copied().map(|peer_id| async move { + let result: anyhow::Result> = api + .with_module(instance_id) + .request_single_peer(endpoint.to_owned(), ApiRequestErased::default(), peer_id) + .await + .map_err(anyhow::Error::from); + (peer_id, result) + })) + .await; + + let mut any_success = false; + let mut entries = Vec::new(); + for (peer_id, result) in peer_results { + match result { + Ok(peer_entries) => { + any_success = true; + entries.extend(peer_entries); + } + Err(e) => { + warn!("Failed to fetch {endpoint} from peer {peer_id}: {e:?}"); + } + } + } + any_success.then_some(entries) +} + +async fn fetch_lnv1_registry( + api: &DynGlobalApi, + instance_id: ModuleInstanceId, + peer_ids: &[PeerId], +) -> Option> { + let announcements: Vec = + fetch_registry(api, instance_id, LIST_GATEWAYS_ENDPOINT, peer_ids).await?; + + let mut merged: HashMap = HashMap::new(); + for gw in announcements { + merged.entry(gw.info.gateway_id.to_string()).or_insert(gw); + } + Some(merged.into_values().collect()) +} + +async fn fetch_lnv2_registry( + api: &DynGlobalApi, + instance_id: ModuleInstanceId, + peer_ids: &[PeerId], +) -> Option> { + let urls: Vec = fetch_registry(api, instance_id, GATEWAYS_ENDPOINT, peer_ids).await?; + + let mut merged: BTreeMap = BTreeMap::new(); + for url in urls { + merged + .entry(normalize_gateway_url(&url.to_string())) + .or_insert(url); + } + Some(merged.into_values().collect()) +} + +/// Ask each LNv2 gateway for its routing info for this federation. This is +/// the only way to learn an LNv2 gateway's keys, since the registry lists +/// URLs only. +async fn probe_lnv2_gateways( + gateway_conn: &RealGatewayConnection, + federation_id: FederationId, + urls: Vec, +) -> Vec { + futures::stream::iter(urls) + .map(|url| async move { + let probe = match tokio::time::timeout( + LNV2_PROBE_TIMEOUT, + gateway_conn.routing_info(url.clone(), &federation_id), + ) + .await + { + Ok(Ok(Some(routing_info))) => Lnv2Probe::Ok(Box::new(routing_info)), + Ok(Ok(None)) => Lnv2Probe::NotServing, + Ok(Err(e)) => { + debug!("LNv2 gateway {url} routing info request failed: {e:?}"); + Lnv2Probe::Unreachable + } + Err(_) => { + debug!("LNv2 gateway {url} routing info request timed out"); + Lnv2Probe::Unreachable + } + }; + Lnv2Gateway { + gateway_id: normalize_gateway_url(&url.to_string()), + api_endpoint: url, + probe, + } + }) + .buffer_unordered(LNV2_PROBE_CONCURRENCY) + .collect() + .await +} + +fn lnv1_gateway_info(gw: LightningGatewayAnnouncement) -> anyhow::Result { + let raw = serde_json::to_value(&gw)?; + Ok(GatewayInfo { + gateway_id: gw.info.gateway_id.to_string(), + node_pub_key: gw.info.node_pub_key.to_string(), + lightning_alias: gw.info.lightning_alias, + api_endpoint: gw.info.api.to_string(), + vetted: gw.vetted, + raw: Some(raw), + first_seen: None, + last_seen: None, + activity_7d: None, + activity_window: None, + uptime_window: None, + metrics_window: None, + protocols: None, + lnv2: None, + }) +} + +fn lnv2_raw(gw: &Lnv2Gateway) -> serde_json::Value { + serde_json::json!({ + "api_endpoint": gw.api_endpoint.to_string(), + "routing_info": gw.probe.routing_info(), + }) +} + +/// Builds the API view of an LNv2-only registration. `gateway_id` is the +/// normalized URL because LNv2 has no other stable identifier; LNv2 has no +/// notion of vetting, so `vetted` is always false. +#[allow(clippy::too_many_arguments)] +fn lnv2_gateway_info( + gateway_id: String, + api_endpoint: String, + routing_status: Lnv2RoutingStatus, + routing_checked_at: Option>, + lightning_public_key: Option, + lightning_alias: Option, + module_public_key: Option, + routing_info: Option<&RoutingInfo>, + raw: serde_json::Value, +) -> GatewayInfo { + GatewayInfo { + gateway_id, + node_pub_key: lightning_public_key.clone().unwrap_or_default(), + lightning_alias: lightning_alias.unwrap_or_default(), + api_endpoint: api_endpoint.clone(), + vetted: false, + raw: Some(raw), + first_seen: None, + last_seen: None, + activity_7d: None, + activity_window: None, + uptime_window: None, + metrics_window: None, + protocols: Some(vec![GatewayProtocol::Lnv2]), + lnv2: Some(Lnv2GatewayInfo { + api_endpoint, + routing_status, + routing_checked_at, + lightning_public_key, + module_public_key, + send_fee_minimum: routing_info.map(|info| gateway_fee(&info.send_fee_minimum)), + send_fee_default: routing_info.map(|info| gateway_fee(&info.send_fee_default)), + receive_fee: routing_info.map(|info| gateway_fee(&info.receive_fee)), + uptime_window: None, + }), + } +} + +fn lnv2_gateway_info_from_probe(gw: &Lnv2Gateway, checked_at: DateTime) -> GatewayInfo { + let routing_info = gw.probe.routing_info(); + lnv2_gateway_info( + gw.gateway_id.clone(), + gw.api_endpoint.to_string(), + gw.probe.status(), + Some(checked_at), + routing_info.map(|info| info.lightning_public_key.to_string()), + routing_info.and_then(|info| info.lightning_alias.clone()), + routing_info.map(|info| info.module_public_key.to_string()), + routing_info, + lnv2_raw(gw), + ) +} + +/// Merges LNv2 registrations into LNv1 gateways where identity is verified: +/// the LNv1 node key must equal the LNv2 lightning key **and** the normalized +/// API URLs must match. Anything else, including an ambiguous match against +/// more than one LNv1 gateway, stays a separate row. LNv1 rows must already +/// carry `protocols: Some([Lnv1])`. +pub(crate) fn merge_protocols( + mut lnv1: Vec, + lnv2: Vec, +) -> Vec { + let mut unmerged = Vec::new(); + + for v2 in lnv2 { + let candidates: Vec = match &v2.lnv2 { + Some(Lnv2GatewayInfo { + lightning_public_key: Some(lightning_key), + api_endpoint, + .. + }) => { + let v2_url = normalize_gateway_url(api_endpoint); + lnv1.iter() + .enumerate() + .filter(|(_, v1)| { + v1.lnv2.is_none() + && v1.node_pub_key == *lightning_key + && normalize_gateway_url(&v1.api_endpoint) == v2_url + }) + .map(|(idx, _)| idx) + .collect() + } + _ => vec![], + }; + + if let [idx] = candidates[..] { + let v1 = &mut lnv1[idx]; + v1.protocols = Some(vec![GatewayProtocol::Lnv1, GatewayProtocol::Lnv2]); + v1.lnv2 = v2.lnv2; + } else { + unmerged.push(v2); + } + } + + lnv1.extend(unmerged); + lnv1 +} + +/// Live registry lookup used by the stable `/config/:invite/gateways` API. +/// +/// With `include_lnv2 == false` the output is exactly what this API returned +/// before LNv2 support: LNv1 gateways only, without the `protocols`/`lnv2` +/// fields. A federation without a Lightning module yields an empty list. pub(crate) async fn fetch_gateways_for_config( config: &ClientConfig, + include_lnv2: bool, ) -> anyhow::Result> { + let instances = LnInstances::from_config(config); + if instances.is_empty() { + return Ok(vec![]); + } + let connectors = fedimint_connectors::ConnectorRegistry::build_from_client_env()? .bind() .await?; @@ -87,76 +454,50 @@ pub(crate) async fn fetch_gateways_for_config( .iter() .map(|(&peer_id, peer_url)| (peer_id, peer_url.url.clone())) .collect(); - let api = DynGlobalApi::new(connectors, peers, None)?; + let api = DynGlobalApi::new(connectors.clone(), peers, None)?; + let peer_ids: Vec = config.global.api_endpoints.keys().copied().collect(); - let ln_instance_id = config - .modules - .iter() - .find_map(|(&instance_id, module)| (module.kind.as_str() == "ln").then_some(instance_id)) - .context("No LN module found in federation config")?; + let mut lnv1 = match instances.lnv1 { + Some(instance_id) => fetch_lnv1_registry(&api, instance_id, &peer_ids) + .await + .unwrap_or_default() + .into_iter() + .map(lnv1_gateway_info) + .collect::>>()?, + None => vec![], + }; - let peer_ids: Vec = - config.global.api_endpoints.keys().copied().collect(); - let mut merged: HashMap = HashMap::new(); - let peer_results = join_all(peer_ids.iter().copied().map(|peer_id| { - let api = api.clone(); - async move { - let result: anyhow::Result> = api - .with_module(ln_instance_id) - .request_single_peer( - LIST_GATEWAYS_ENDPOINT.to_owned(), - ApiRequestErased::default(), - peer_id, - ) - .await - .map_err(anyhow::Error::from) - .and_then(|v| serde_json::from_value(v).map_err(anyhow::Error::from)); - (peer_id, result) - } - })) - .await; + let Some(lnv2_instance_id) = instances.lnv2.filter(|_| include_lnv2) else { + return Ok(lnv1); + }; - for (peer_id, result) in peer_results { - match result { - Ok(gateways) => { - for gw in gateways { - merged.entry(gw.info.gateway_id.to_string()).or_insert(gw); - } - } - Err(e) => { - warn!( - "Failed to fetch live gateways from peer {}: {:?}", - peer_id, e - ); - } - } + for gw in &mut lnv1 { + gw.protocols = Some(vec![GatewayProtocol::Lnv1]); } - merged - .into_values() - .map(|gw| { - let raw = serde_json::to_value(&gw)?; - Ok(GatewayInfo { - gateway_id: gw.info.gateway_id.to_string(), - node_pub_key: gw.info.node_pub_key.to_string(), - lightning_alias: gw.info.lightning_alias, - api_endpoint: gw.info.api.to_string(), - vetted: gw.vetted, - raw: Some(raw), - first_seen: None, - last_seen: None, - activity_7d: None, - activity_window: None, - uptime_window: None, - metrics_window: None, - }) - }) - .collect() + let lnv2_urls = fetch_lnv2_registry(&api, lnv2_instance_id, &peer_ids) + .await + .unwrap_or_default(); + let gateway_conn = RealGatewayConnection { + api: GatewayApi::new(None, connectors), + }; + let now = Utc::now(); + let lnv2 = probe_lnv2_gateways( + &gateway_conn, + config.global.calculate_federation_id(), + lnv2_urls, + ) + .await + .iter() + .map(|gw| lnv2_gateway_info_from_probe(gw, now)) + .collect(); + + Ok(merge_protocols(lnv1, lnv2)) } impl FederationObserver { - /// Background task: poll the LN module on this federation for currently - /// registered gateways and persist them. Runs in a loop until cancelled. + /// Background task: poll the federation's LNv1 and LNv2 gateway registries + /// and persist what they contain. Runs in a loop until cancelled. pub async fn monitor_gateways( &self, federation_id: FederationId, @@ -164,6 +505,13 @@ impl FederationObserver { ) -> anyhow::Result<()> { const POLL_INTERVAL: Duration = Duration::from_secs(GATEWAY_POLL_INTERVAL_MINUTES * 60); + let instances = LnInstances::from_config(&config); + if instances.is_empty() { + info!("Federation {federation_id} has no Lightning module, not monitoring gateways"); + // Returning would make the supervisor restart this task every 30s + return futures::future::pending().await; + } + let peers = config .global .api_endpoints @@ -171,23 +519,18 @@ impl FederationObserver { .map(|(&peer_id, peer_url)| (peer_id, peer_url.url.clone())) .collect(); let api = DynGlobalApi::new(self.connectors().clone(), peers, None)?; + // Kept across polls so connections to gateways are reused + let gateway_conn = RealGatewayConnection { + api: GatewayApi::new(None, self.connectors().clone()), + }; - let ln_instance_id = config - .modules - .iter() - .find_map(|(&instance_id, module)| { - (module.kind.as_str() == "ln").then_some(instance_id) - }) - .context("No LN module found in federation config")?; - - let peer_ids: Vec = - config.global.api_endpoints.keys().copied().collect(); + let peer_ids: Vec = config.global.api_endpoints.keys().copied().collect(); let mut interval = tokio::time::interval(POLL_INTERVAL); loop { interval.tick().await; if let Err(e) = self - .fetch_and_store_gateways(federation_id, &api, ln_instance_id, &peer_ids) + .fetch_and_store_gateways(federation_id, &api, &gateway_conn, instances, &peer_ids) .await { warn!( @@ -202,46 +545,34 @@ impl FederationObserver { &self, federation_id: FederationId, api: &DynGlobalApi, - ln_instance_id: ModuleInstanceId, - peer_ids: &[fedimint_core::PeerId], + gateway_conn: &RealGatewayConnection, + instances: LnInstances, + peer_ids: &[PeerId], ) -> anyhow::Result<()> { - // Query all peers and merge by gateway_id — each guardian has their own - // registry so we take the union across all peers. - let mut merged: HashMap = HashMap::new(); - let mut successful_peer_queries: u32 = 0; - let peer_results = join_all(peer_ids.iter().copied().map(|peer_id| async move { - let result: anyhow::Result> = api - .with_module(ln_instance_id) - .request_single_peer( - LIST_GATEWAYS_ENDPOINT.to_owned(), - ApiRequestErased::default(), - peer_id, - ) - .await - .map_err(anyhow::Error::from) - .and_then(|v| serde_json::from_value(v).map_err(anyhow::Error::from)); - (peer_id, result) - })) - .await; - - for (peer_id, result) in peer_results { - match result { - Ok(gateways) => { - successful_peer_queries += 1; - for gw in gateways { - merged.entry(gw.info.gateway_id.to_string()).or_insert(gw); - } - } - Err(e) => { - warn!( - "Failed to fetch gateways from peer {} for {}: {:?}", - peer_id, federation_id, e - ); + // `None` for a protocol means its registry state is unknown this poll (module + // absent or no peer answered), so no snapshot is written for it. + let lnv1 = match instances.lnv1 { + Some(instance_id) => { + let registry = fetch_lnv1_registry(api, instance_id, peer_ids).await; + if registry.is_none() { + warn!("No peer answered the LNv1 gateway registry query for {federation_id}"); } + registry } - } + None => None, + }; + let lnv2 = match instances.lnv2 { + Some(instance_id) => match fetch_lnv2_registry(api, instance_id, peer_ids).await { + Some(urls) => Some(probe_lnv2_gateways(gateway_conn, federation_id, urls).await), + None => { + warn!("No peer answered the LNv2 gateway registry query for {federation_id}"); + None + } + }, + None => None, + }; - if successful_peer_queries == 0 { + if lnv1.is_none() && lnv2.is_none() { bail!( "No successful gateway registry responses from any federation peer for {}", federation_id @@ -253,96 +584,13 @@ impl FederationObserver { let now = chrono::Utc::now(); let federation_id_bytes = federation_id.consensus_encode_to_vec(); - let mut gateway_ids = Vec::with_capacity(merged.len()); - let mut node_pub_keys = Vec::with_capacity(merged.len()); - let mut api_endpoints = Vec::with_capacity(merged.len()); - let mut lightning_aliases = Vec::with_capacity(merged.len()); - let mut vetted_flags = Vec::with_capacity(merged.len()); - let mut raw_announcements = Vec::with_capacity(merged.len()); - - for (gateway_id, gw) in &merged { - gateway_ids.push(gateway_id.clone()); - node_pub_keys.push(gw.info.node_pub_key.to_string()); - api_endpoints.push(gw.info.api.to_string()); - lightning_aliases.push(gw.info.lightning_alias.clone()); - vetted_flags.push(gw.vetted); - raw_announcements.push(serde_json::to_string(gw)?); + if let Some(lnv1) = &lnv1 { + store_lnv1_gateways(&dbtx, &federation_id_bytes, now, lnv1).await?; } - - if !gateway_ids.is_empty() { - dbtx.execute( - "INSERT INTO gateways - (federation_id, gateway_id, node_pub_key, api_endpoint, - lightning_alias, vetted, raw, first_seen, last_seen) - SELECT - $1, - gw.gateway_id, - gw.node_pub_key, - gw.api_endpoint, - gw.lightning_alias, - gw.vetted, - gw.raw_json::jsonb, - $2, - $2 - FROM UNNEST( - $3::text[], - $4::text[], - $5::text[], - $6::text[], - $7::boolean[], - $8::text[] - ) AS gw(gateway_id, node_pub_key, api_endpoint, lightning_alias, vetted, raw_json) - ON CONFLICT (federation_id, gateway_id) DO UPDATE - SET node_pub_key = EXCLUDED.node_pub_key, - api_endpoint = EXCLUDED.api_endpoint, - lightning_alias = EXCLUDED.lightning_alias, - vetted = EXCLUDED.vetted, - raw = EXCLUDED.raw, - last_seen = EXCLUDED.last_seen", - &[ - &federation_id_bytes, - &now, - &gateway_ids, - &node_pub_keys, - &api_endpoints, - &lightning_aliases, - &vetted_flags, - &raw_announcements, - ], - ) - .await?; + if let Some(lnv2) = &lnv2 { + store_lnv2_gateways(&dbtx, &federation_id_bytes, now, lnv2).await?; } - dbtx.execute( - "WITH current_gateway_ids AS ( - SELECT UNNEST($3::text[]) AS gateway_id - ), - all_gateway_ids AS ( - SELECT gateway_id, TRUE AS is_seen - FROM current_gateway_ids - UNION - SELECT g.gateway_id, FALSE AS is_seen - FROM gateways g - WHERE g.federation_id = $1 - AND NOT EXISTS ( - SELECT 1 - FROM current_gateway_ids c - WHERE c.gateway_id = g.gateway_id - ) - ) - INSERT INTO gateway_poll_snapshots - (federation_id, gateway_id, poll_time, is_seen) - SELECT - $1, - a.gateway_id, - $2, - a.is_seen - FROM all_gateway_ids a - ON CONFLICT DO NOTHING", - &[&federation_id_bytes, &now, &gateway_ids], - ) - .await?; - let prune_interval_secs = GATEWAY_PRUNE_INTERVAL_HOURS * 60 * 60; let should_prune = now.timestamp().rem_euclid(prune_interval_secs) < (GATEWAY_POLL_INTERVAL_MINUTES as i64 * 60); @@ -361,8 +609,11 @@ impl FederationObserver { dbtx.commit().await?; info!( - "Stored {} seen gateway(s), persisted poll snapshots for federation {}, deleted {} old snapshots", - merged.len(), federation_id, deleted_snapshots + "Stored {} LNv1 and {} LNv2 gateway(s), persisted poll snapshots for federation {}, deleted {} old snapshots", + lnv1.as_ref().map_or(0, Vec::len), + lnv2.as_ref().map_or(0, Vec::len), + federation_id, + deleted_snapshots ); Ok(()) } @@ -374,14 +625,18 @@ impl FederationObserver { ) -> anyhow::Result> { #[derive(postgres_from_row::FromRow)] struct GatewayRow { + protocol: String, gateway_id: String, - node_pub_key: String, - lightning_alias: String, + node_pub_key: Option, + lightning_alias: Option, api_endpoint: String, vetted: bool, raw: serde_json::Value, first_seen: chrono::DateTime, last_seen: chrono::DateTime, + module_public_key: Option, + routing_status: Option, + routing_checked_at: Option>, } #[derive(postgres_from_row::FromRow)] @@ -395,6 +650,7 @@ impl FederationObserver { #[derive(postgres_from_row::FromRow)] struct GatewayUptimeRow { + protocol: String, gateway_id: String, seen_samples: i64, total_samples: i64, @@ -408,10 +664,10 @@ impl FederationObserver { let rows = query::( &conn, - "SELECT gateway_id, node_pub_key, lightning_alias, api_endpoint, vetted, raw, first_seen, last_seen + "SELECT protocol, gateway_id, node_pub_key, lightning_alias, api_endpoint, vetted, raw, + first_seen, last_seen, module_public_key, routing_status, routing_checked_at FROM gateways - WHERE federation_id = $1 - ORDER BY last_seen DESC", + WHERE federation_id = $1", &[&federation_id_bytes], ) .await?; @@ -539,18 +795,19 @@ impl FederationObserver { let uptime_rows = query::( &conn, "SELECT + protocol, gateway_id, COUNT(*) FILTER (WHERE is_seen)::bigint AS seen_samples, COUNT(*)::bigint AS total_samples FROM gateway_poll_snapshots WHERE federation_id = $1 AND poll_time >= $2 - GROUP BY gateway_id", + GROUP BY protocol, gateway_id", &[&federation_id_bytes, &window_start_utc], ) .await?; - let uptime_by_gateway_id: HashMap = uptime_rows + let uptime_by_gateway_id: HashMap<(String, String), GatewayUptimeMetrics> = uptime_rows .into_iter() .map(|row| { let seen_samples = row.seen_samples.max(0) as u64; @@ -565,7 +822,7 @@ impl FederationObserver { 0.0 }; ( - row.gateway_id, + (row.protocol, row.gateway_id), GatewayUptimeMetrics { sample_count: total_samples, seen_samples, @@ -577,36 +834,87 @@ impl FederationObserver { }) .collect(); - Ok(rows - .into_iter() - .map(|r| { - let activity_window = r - .raw - .pointer("/info/gateway_redeem_key") - .and_then(|v| v.as_str()) - .and_then(|gateway_key| activity_by_gateway_key.get(gateway_key).cloned()); - let uptime_window = uptime_by_gateway_id.get(&r.gateway_id).cloned(); - - GatewayInfo { - activity_7d: if window == GatewayMetricsWindow::D7 { - activity_window.clone() - } else { - None - }, - activity_window, - uptime_window, - metrics_window: Some(metrics_window.clone()), - gateway_id: r.gateway_id, - node_pub_key: r.node_pub_key, - lightning_alias: r.lightning_alias, - api_endpoint: r.api_endpoint, - vetted: r.vetted, - raw: Some(r.raw), - first_seen: Some(r.first_seen), - last_seen: Some(r.last_seen), + let mut lnv1 = Vec::new(); + let mut lnv2 = Vec::new(); + for r in rows { + let uptime_window = uptime_by_gateway_id + .get(&(r.protocol.clone(), r.gateway_id.clone())) + .cloned(); + + match parse_protocol(&r.protocol) { + Some(GatewayProtocol::Lnv1) => { + let activity_window = r + .raw + .pointer("/info/gateway_redeem_key") + .and_then(|v| v.as_str()) + .and_then(|gateway_key| activity_by_gateway_key.get(gateway_key).cloned()); + + lnv1.push(GatewayInfo { + activity_7d: if window == GatewayMetricsWindow::D7 { + activity_window.clone() + } else { + None + }, + activity_window, + uptime_window, + metrics_window: Some(metrics_window.clone()), + gateway_id: r.gateway_id, + node_pub_key: r.node_pub_key.unwrap_or_default(), + lightning_alias: r.lightning_alias.unwrap_or_default(), + api_endpoint: r.api_endpoint, + vetted: r.vetted, + raw: Some(r.raw), + first_seen: Some(r.first_seen), + last_seen: Some(r.last_seen), + protocols: Some(vec![GatewayProtocol::Lnv1]), + lnv2: None, + }); } - }) - .collect()) + Some(GatewayProtocol::Lnv2) => { + // Fees come from the last successful probe, which is what `raw` holds + let routing_info: Option = r + .raw + .get("routing_info") + .and_then(|info| serde_json::from_value(info.clone()).ok()); + let routing_status = r + .routing_status + .as_deref() + .and_then(parse_routing_status) + .unwrap_or(Lnv2RoutingStatus::Unreachable); + + let mut info = lnv2_gateway_info( + r.gateway_id, + r.api_endpoint, + routing_status, + r.routing_checked_at, + r.node_pub_key, + r.lightning_alias, + r.module_public_key, + routing_info.as_ref(), + r.raw, + ); + // LNv2 activity attribution is not implemented yet + info.metrics_window = Some(metrics_window.clone()); + info.first_seen = Some(r.first_seen); + info.last_seen = Some(r.last_seen); + info.uptime_window = uptime_window.clone(); + if let Some(details) = &mut info.lnv2 { + details.uptime_window = uptime_window; + } + lnv2.push(info); + } + None => { + warn!( + "Ignoring gateway {} with unknown protocol {}", + r.gateway_id, r.protocol + ); + } + } + } + + let mut gateways = merge_protocols(lnv1, lnv2); + gateways.sort_by_key(|gateway| std::cmp::Reverse(gateway.last_seen)); + Ok(gateways) } async fn federation_gateway_uptime_trend( @@ -685,3 +993,427 @@ pub(super) async fn get_federation_gateway_uptime_trend( .await? .into()) } + +async fn store_lnv1_gateways( + dbtx: &Transaction<'_>, + federation_id_bytes: &[u8], + now: DateTime, + gateways: &[LightningGatewayAnnouncement], +) -> anyhow::Result<()> { + let mut gateway_ids = Vec::with_capacity(gateways.len()); + let mut node_pub_keys = Vec::with_capacity(gateways.len()); + let mut api_endpoints = Vec::with_capacity(gateways.len()); + let mut lightning_aliases = Vec::with_capacity(gateways.len()); + let mut vetted_flags = Vec::with_capacity(gateways.len()); + let mut raw_announcements = Vec::with_capacity(gateways.len()); + + for gw in gateways { + gateway_ids.push(gw.info.gateway_id.to_string()); + node_pub_keys.push(gw.info.node_pub_key.to_string()); + api_endpoints.push(gw.info.api.to_string()); + lightning_aliases.push(gw.info.lightning_alias.clone()); + vetted_flags.push(gw.vetted); + raw_announcements.push(serde_json::to_string(gw)?); + } + + if !gateway_ids.is_empty() { + dbtx.execute( + "INSERT INTO gateways + (federation_id, protocol, gateway_id, node_pub_key, api_endpoint, + lightning_alias, vetted, raw, first_seen, last_seen) + SELECT + $1, + 'lnv1', + gw.gateway_id, + gw.node_pub_key, + gw.api_endpoint, + gw.lightning_alias, + gw.vetted, + gw.raw_json::jsonb, + $2, + $2 + FROM UNNEST( + $3::text[], + $4::text[], + $5::text[], + $6::text[], + $7::boolean[], + $8::text[] + ) AS gw(gateway_id, node_pub_key, api_endpoint, lightning_alias, vetted, raw_json) + ON CONFLICT (federation_id, protocol, gateway_id) DO UPDATE + SET node_pub_key = EXCLUDED.node_pub_key, + api_endpoint = EXCLUDED.api_endpoint, + lightning_alias = EXCLUDED.lightning_alias, + vetted = EXCLUDED.vetted, + raw = EXCLUDED.raw, + last_seen = EXCLUDED.last_seen", + &[ + &federation_id_bytes, + &now, + &gateway_ids, + &node_pub_keys, + &api_endpoints, + &lightning_aliases, + &vetted_flags, + &raw_announcements, + ], + ) + .await?; + } + + let routing_ok = vec![None; gateway_ids.len()]; + store_poll_snapshots( + dbtx, + federation_id_bytes, + now, + GatewayProtocol::Lnv1, + &gateway_ids, + &routing_ok, + ) + .await +} + +async fn store_lnv2_gateways( + dbtx: &Transaction<'_>, + federation_id_bytes: &[u8], + now: DateTime, + gateways: &[Lnv2Gateway], +) -> anyhow::Result<()> { + let mut gateway_ids = Vec::with_capacity(gateways.len()); + let mut lightning_public_keys = Vec::with_capacity(gateways.len()); + let mut api_endpoints = Vec::with_capacity(gateways.len()); + let mut lightning_aliases = Vec::with_capacity(gateways.len()); + let mut raw_entries = Vec::with_capacity(gateways.len()); + let mut module_public_keys = Vec::with_capacity(gateways.len()); + let mut routing_statuses = Vec::with_capacity(gateways.len()); + let mut routing_ok = Vec::with_capacity(gateways.len()); + + for gw in gateways { + let routing_info = gw.probe.routing_info(); + gateway_ids.push(gw.gateway_id.clone()); + lightning_public_keys.push(routing_info.map(|info| info.lightning_public_key.to_string())); + api_endpoints.push(gw.api_endpoint.to_string()); + lightning_aliases.push(routing_info.and_then(|info| info.lightning_alias.clone())); + raw_entries.push(serde_json::to_string(&lnv2_raw(gw))?); + module_public_keys.push(routing_info.map(|info| info.module_public_key.to_string())); + routing_statuses.push(routing_status_str(gw.probe.status())); + routing_ok.push(Some(routing_info.is_some())); + } + + if !gateway_ids.is_empty() { + // A failed or negative probe must not erase the last identity we verified, so + // keys, alias and raw routing info are only overwritten by a successful probe. + dbtx.execute( + "INSERT INTO gateways + (federation_id, protocol, gateway_id, node_pub_key, api_endpoint, + lightning_alias, vetted, raw, first_seen, last_seen, + module_public_key, routing_status, routing_checked_at) + SELECT + $1, + 'lnv2', + gw.gateway_id, + gw.node_pub_key, + gw.api_endpoint, + gw.lightning_alias, + FALSE, + gw.raw_json::jsonb, + $2, + $2, + gw.module_public_key, + gw.routing_status, + $2 + FROM UNNEST( + $3::text[], + $4::text[], + $5::text[], + $6::text[], + $7::text[], + $8::text[], + $9::text[] + ) AS gw(gateway_id, node_pub_key, api_endpoint, lightning_alias, raw_json, + module_public_key, routing_status) + ON CONFLICT (federation_id, protocol, gateway_id) DO UPDATE + SET node_pub_key = COALESCE(EXCLUDED.node_pub_key, gateways.node_pub_key), + api_endpoint = EXCLUDED.api_endpoint, + lightning_alias = COALESCE(EXCLUDED.lightning_alias, gateways.lightning_alias), + raw = CASE WHEN EXCLUDED.routing_status = 'ok' + THEN EXCLUDED.raw + ELSE gateways.raw + END, + module_public_key = COALESCE(EXCLUDED.module_public_key, gateways.module_public_key), + routing_status = EXCLUDED.routing_status, + routing_checked_at = EXCLUDED.routing_checked_at, + last_seen = EXCLUDED.last_seen", + &[ + &federation_id_bytes, + &now, + &gateway_ids, + &lightning_public_keys, + &api_endpoints, + &lightning_aliases, + &raw_entries, + &module_public_keys, + &routing_statuses, + ], + ) + .await?; + + dbtx.execute( + "INSERT INTO lnv2_gateway_keys + (federation_id, gateway_id, module_public_key, lightning_public_key, + first_seen, last_seen) + SELECT $1, k.gateway_id, k.module_public_key, k.lightning_public_key, $2, $2 + FROM UNNEST($3::text[], $4::text[], $5::text[]) + AS k(gateway_id, module_public_key, lightning_public_key) + WHERE k.module_public_key IS NOT NULL + ON CONFLICT (federation_id, gateway_id, module_public_key) DO UPDATE + SET lightning_public_key = EXCLUDED.lightning_public_key, + last_seen = EXCLUDED.last_seen", + &[ + &federation_id_bytes, + &now, + &gateway_ids, + &module_public_keys, + &lightning_public_keys, + ], + ) + .await?; + } + + store_poll_snapshots( + dbtx, + federation_id_bytes, + now, + GatewayProtocol::Lnv2, + &gateway_ids, + &routing_ok, + ) + .await +} + +/// Records one poll: every gateway currently in the protocol's registry as +/// seen, and every previously known gateway of that protocol as not seen. +async fn store_poll_snapshots( + dbtx: &Transaction<'_>, + federation_id_bytes: &[u8], + now: DateTime, + protocol: GatewayProtocol, + gateway_ids: &[String], + routing_ok: &[Option], +) -> anyhow::Result<()> { + dbtx.execute( + "WITH current_gateways AS ( + SELECT * + FROM UNNEST($4::text[], $5::boolean[]) AS c(gateway_id, routing_ok) + ), + all_gateways AS ( + SELECT gateway_id, TRUE AS is_seen, routing_ok + FROM current_gateways + UNION ALL + SELECT g.gateway_id, FALSE AS is_seen, NULL::boolean AS routing_ok + FROM gateways g + WHERE g.federation_id = $1 + AND g.protocol = $3 + AND NOT EXISTS ( + SELECT 1 + FROM current_gateways c + WHERE c.gateway_id = g.gateway_id + ) + ) + INSERT INTO gateway_poll_snapshots + (federation_id, protocol, gateway_id, poll_time, is_seen, routing_ok) + SELECT + $1, + $3, + a.gateway_id, + $2, + a.is_seen, + a.routing_ok + FROM all_gateways a + ON CONFLICT DO NOTHING", + &[ + &federation_id_bytes, + &now, + &protocol_str(protocol), + &gateway_ids, + &routing_ok, + ], + ) + .await?; + Ok(()) +} + +fn protocol_str(protocol: GatewayProtocol) -> &'static str { + match protocol { + GatewayProtocol::Lnv1 => "lnv1", + GatewayProtocol::Lnv2 => "lnv2", + } +} + +fn parse_protocol(protocol: &str) -> Option { + match protocol { + "lnv1" => Some(GatewayProtocol::Lnv1), + "lnv2" => Some(GatewayProtocol::Lnv2), + _ => None, + } +} + +#[cfg(test)] +mod tests { + use fedimint_core::core::{ModuleInstanceId, ModuleKind}; + use fmo_api_types::{GatewayInfo, GatewayProtocol, Lnv2RoutingStatus}; + + use super::{lnv2_gateway_info, merge_protocols, normalize_gateway_url, LnInstances}; + + const NODE_A: &str = "02aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"; + const NODE_B: &str = "03bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb"; + + fn lnv1(gateway_id: &str, node_pub_key: &str, api_endpoint: &str) -> GatewayInfo { + GatewayInfo { + gateway_id: gateway_id.to_owned(), + node_pub_key: node_pub_key.to_owned(), + lightning_alias: "alias".to_owned(), + api_endpoint: api_endpoint.to_owned(), + vetted: true, + raw: None, + first_seen: None, + last_seen: None, + activity_7d: None, + activity_window: None, + uptime_window: None, + metrics_window: None, + protocols: Some(vec![GatewayProtocol::Lnv1]), + lnv2: None, + } + } + + fn lnv2(api_endpoint: &str, lightning_public_key: Option<&str>) -> GatewayInfo { + lnv2_gateway_info( + normalize_gateway_url(api_endpoint), + api_endpoint.to_owned(), + if lightning_public_key.is_some() { + Lnv2RoutingStatus::Ok + } else { + Lnv2RoutingStatus::Unreachable + }, + None, + lightning_public_key.map(str::to_owned), + None, + None, + None, + serde_json::Value::Null, + ) + } + + fn protocols(gateway: &GatewayInfo) -> Vec { + gateway.protocols.clone().unwrap_or_default() + } + + #[test] + fn normalize_gateway_url_ignores_trailing_slash_and_v1() { + let expected = "https://gw.example.com"; + for url in [ + "https://gw.example.com", + "https://gw.example.com/", + "https://gw.example.com/v1", + "https://gw.example.com/v1/", + ] { + assert_eq!(normalize_gateway_url(url), expected, "{url}"); + } + assert_eq!( + normalize_gateway_url("https://gw.example.com/v12"), + "https://gw.example.com/v12" + ); + } + + #[test] + fn ln_instances_detects_either_both_or_neither() { + let ln = ModuleKind::from_static_str("ln"); + let lnv2 = ModuleKind::from_static_str("lnv2"); + let mint = ModuleKind::from_static_str("mint"); + fn kinds(modules: &[(ModuleInstanceId, &ModuleKind)]) -> LnInstances { + LnInstances::from_kinds(modules.iter().copied()) + } + + let none = kinds(&[(0, &mint)]); + assert!(none.is_empty()); + + let v1_only = kinds(&[(0, &mint), (1, &ln)]); + assert_eq!((v1_only.lnv1, v1_only.lnv2), (Some(1), None)); + + let v2_only = kinds(&[(0, &mint), (3, &lnv2)]); + assert_eq!((v2_only.lnv1, v2_only.lnv2), (None, Some(3))); + + let both = kinds(&[(1, &ln), (0, &mint), (2, &lnv2)]); + assert_eq!((both.lnv1, both.lnv2), (Some(1), Some(2))); + } + + #[test] + fn merge_single_protocol_inputs_unchanged() { + let merged = merge_protocols(vec![lnv1("g1", NODE_A, "https://a.example.com/v1")], vec![]); + assert_eq!(merged.len(), 1); + assert_eq!(protocols(&merged[0]), vec![GatewayProtocol::Lnv1]); + + let merged = merge_protocols(vec![], vec![lnv2("https://a.example.com", Some(NODE_A))]); + assert_eq!(merged.len(), 1); + assert_eq!(protocols(&merged[0]), vec![GatewayProtocol::Lnv2]); + } + + #[test] + fn merge_verified_dual_protocol_gateway() { + let merged = merge_protocols( + vec![lnv1("g1", NODE_A, "https://a.example.com/v1")], + vec![lnv2("https://a.example.com/", Some(NODE_A))], + ); + assert_eq!(merged.len(), 1); + assert_eq!(merged[0].gateway_id, "g1"); + assert_eq!( + protocols(&merged[0]), + vec![GatewayProtocol::Lnv1, GatewayProtocol::Lnv2] + ); + assert_eq!( + merged[0] + .lnv2 + .as_ref() + .map(|details| details.routing_status), + Some(Lnv2RoutingStatus::Ok) + ); + } + + #[test] + fn no_merge_without_both_key_and_url_match() { + // Same node, different API + let merged = merge_protocols( + vec![lnv1("g1", NODE_A, "https://a.example.com")], + vec![lnv2("https://other.example.com", Some(NODE_A))], + ); + assert_eq!(merged.len(), 2); + + // Same API, different node + let merged = merge_protocols( + vec![lnv1("g1", NODE_A, "https://a.example.com")], + vec![lnv2("https://a.example.com", Some(NODE_B))], + ); + assert_eq!(merged.len(), 2); + + // Identity unknown because the probe failed + let merged = merge_protocols( + vec![lnv1("g1", NODE_A, "https://a.example.com")], + vec![lnv2("https://a.example.com", None)], + ); + assert_eq!(merged.len(), 2); + } + + #[test] + fn no_merge_when_match_is_ambiguous() { + let merged = merge_protocols( + vec![ + lnv1("g1", NODE_A, "https://a.example.com"), + lnv1("g2", NODE_A, "https://a.example.com/v1"), + ], + vec![lnv2("https://a.example.com", Some(NODE_A))], + ); + assert_eq!(merged.len(), 3); + assert!(merged.iter().all(|gateway| protocols(gateway).len() == 1)); + } +} diff --git a/fmo_server/src/federation/observer.rs b/fmo_server/src/federation/observer.rs index a38978d..1ebf146 100644 --- a/fmo_server/src/federation/observer.rs +++ b/fmo_server/src/federation/observer.rs @@ -219,6 +219,7 @@ impl FederationObserver { ), migration!("/schema/v9.sql"), migration!("/schema/v10.sql"), + migration!("/schema/v11.sql"), ]; for (index, migration) in migrations.iter().enumerate() {