diff --git a/docs/api-guide.md b/docs/api-guide.md index 6e430653..4f97c7a9 100644 --- a/docs/api-guide.md +++ b/docs/api-guide.md @@ -81,8 +81,9 @@ all the caller's restrictions. Revoking the caller's token does not revoke these Copies made with `derive-macaroon` share the original token's ID. Revoking that ID blocks all those copies. The server cannot list copies made locally. -Revocation and expiry block new requests. Existing event streams stay open until the client -disconnects or the server stops. Reconnecting requires a valid token. +Revocation and expiry block new requests. They also end existing event streams: the server +closes the stream with `UNAUTHENTICATED` ("Macaroon revoked" or "Macaroon expired") and does +not send later events. Reconnecting requires a valid token. See [Macaroon Management](#macaroon-management) for the RPCs and [Operations](operations.md#macaroons) for storage and recovery. @@ -111,7 +112,7 @@ RPCs with no permission mapping return `UNIMPLEMENTED`, even for admin tokens. | `messages:verify` | Verify message signatures | | `graph:read` | Read network graph data | | `utilities:read` | Decode invoices and offers | -| `events:read` | Subscribe to the event stream | +| `events:read` | Subscribe to the event streams | | `macaroons:manage` | Create, list, and revoke macaroons within your permissions | MCP provides token management through `create_macaroon`, `list_macaroons`, `revoke_macaroon`, @@ -309,11 +310,18 @@ See [Pagination](#pagination) below for how to page through results. ### Event Streaming -| RPC | Description | -|-------------------|-------------------------------------------------------------| -| `SubscribeEvents` | **Server-streaming.** Subscribe to real-time payment and channel events | +| RPC | Description | +|-----------------------------|-----------------------------------------------------------------------------| +| `SubscribeEvents` | **Server-streaming.** Subscribe to real-time payment and channel events | +| `SubscribeChannelEvents` | **Server-streaming.** Subscribe to real-time channel events only | +| `SubscribePaymentEvents` | **Server-streaming.** Subscribe to real-time payment events only | +| `SubscribeForwardingEvents` | **Server-streaming.** Subscribe to real-time payment forwarding events only | -`SubscribeEvents` returns a stream of `EventEnvelope` messages. Each envelope contains one of: +`SubscribeEvents` returns a stream of `EventEnvelope` messages. Each envelope contains one of the +events below. `SubscribePaymentEvents` delivers only `PaymentReceived`, `PaymentSuccessful`, +`PaymentFailed`, and `PaymentClaimable`. `SubscribeForwardingEvents` delivers only +`PaymentForwarded`. `SubscribeChannelEvents` delivers only `ChannelStateChanged`, +`SpliceNegotiated`, and `SpliceNegotiationFailed`. | Event | When | |---------------------|-----------------------------------------------------------------------| @@ -327,7 +335,7 @@ See [Pagination](#pagination) below for how to page through results. | `SpliceNegotiationFailed` | A channel splice negotiation round failed | > [!WARNING] -> `SubscribeEvents` is a best-effort stream of new events. Events are not persisted for +> All event streams are best-effort streams of new events. Events are not persisted for > subscribers, cannot be replayed after reconnecting, and have no client acknowledgement. > Acceptance by the server's broadcast channel does not guarantee that a client received or > processed an event. diff --git a/e2e-tests/tests/e2e.rs b/e2e-tests/tests/e2e.rs index 68ba0500..7859bcad 100644 --- a/e2e-tests/tests/e2e.rs +++ b/e2e-tests/tests/e2e.rs @@ -865,6 +865,86 @@ async fn test_subscribe_events_channel_state_lifecycle_pending_ready_closed() { assert_eq!(closed_b.closure_initiator, ChannelClosureInitiator::Remote as i32); } +#[tokio::test] +async fn test_subscribe_channel_and_payment_events() { + let bitcoind = TestBitcoind::new(); + let server_a = LdkServerHandle::start(&bitcoind).await; + let server_b = LdkServerHandle::start(&bitcoind).await; + + let mut channel_events = server_a.client().subscribe_channel_events().await.unwrap(); + let mut payment_events = server_a.client().subscribe_payment_events().await.unwrap(); + + let user_channel_id = setup_funded_channel(&bitcoind, &server_a, &server_b, 100_000).await; + let payment_id = send_bolt11_payment(&server_a, &server_b, 10_000_000).await; + close_channel(&server_a, &server_b, &user_channel_id).await; + mine_and_sync(&bitcoind, &[&server_a, &server_b], 6).await; + + // The channel opened before the payment was sent, so the first payment stream event shows + // that channel events were filtered out. + let event = wait_for_event(&mut payment_events, |_| true).await; + match event.event { + Some(Event::PaymentSuccessful(e)) => assert_eq!(e.payment_id, payment_id), + other => panic!("expected PaymentSuccessful event, got {other:?}"), + } + + // The payment was sent between the channel open and close, so it must have been filtered out. + let mut states = Vec::new(); + while states.last() != Some(&(ChannelState::Closed as i32)) { + let event = wait_for_event(&mut channel_events, |_| true).await; + match event.event { + Some(Event::ChannelStateChanged(e)) => { + assert_eq!(e.user_channel_id, user_channel_id); + states.push(e.state); + }, + other => panic!("expected ChannelStateChanged event, got {other:?}"), + } + } + assert_eq!( + states, + [ChannelState::Pending as i32, ChannelState::Ready as i32, ChannelState::Closed as i32] + ); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 1)] +async fn test_subscribe_forwarding_events() { + let bitcoind = TestBitcoind::new(); + let server_a = LdkServerHandle::start(&bitcoind).await; + let server_b = LdkServerHandle::start(&bitcoind).await; + let server_c = LdkServerHandle::start(&bitcoind).await; + + let mut forwarding_events = server_b.client().subscribe_forwarding_events().await.unwrap(); + let mut payment_events = server_b.client().subscribe_payment_events().await.unwrap(); + + // A -> B -> C + setup_funded_channel(&bitcoind, &server_a, &server_b, 1_000_000).await; + setup_funded_channel(&bitcoind, &server_b, &server_c, 1_000_000).await; + wait_for_usable_channel(server_c.client(), &bitcoind, Duration::from_secs(60)).await; + wait_for_channels(&server_b, 2, Duration::from_secs(60)).await; + wait_for_gossip(&server_a, 2, Duration::from_secs(60)).await; + + // B's own payments are sent before and after the forward so each stream must skip the other. + let first_payment_id = send_bolt11_payment(&server_b, &server_c, 10_000_000).await; + send_bolt11_payment(&server_a, &server_c, 10_000_000).await; + let second_payment_id = send_bolt11_payment(&server_b, &server_c, 10_000_000).await; + + let event = wait_for_event(&mut forwarding_events, |_| true).await; + match event.event { + Some(Event::PaymentForwarded(e)) => { + assert_eq!(e.prev_htlcs.len(), 1); + assert_eq!(e.next_htlcs.len(), 1); + }, + other => panic!("expected PaymentForwarded event, got {other:?}"), + } + + for payment_id in [first_payment_id, second_payment_id] { + let event = wait_for_event(&mut payment_events, |_| true).await; + match event.event { + Some(Event::PaymentSuccessful(e)) => assert_eq!(e.payment_id, payment_id), + other => panic!("expected PaymentSuccessful event, got {other:?}"), + } + } +} + #[tokio::test] async fn test_subscribe_events_channel_state_lifecycle_pending_ready_force_closed() { let bitcoind = TestBitcoind::new(); diff --git a/e2e-tests/tests/macaroons.rs b/e2e-tests/tests/macaroons.rs index 2cd0d9f4..d50a6dc5 100644 --- a/e2e-tests/tests/macaroons.rs +++ b/e2e-tests/tests/macaroons.rs @@ -12,7 +12,7 @@ use std::time::Duration; use e2e_tests::{ mine_and_sync, run_cli, setup_funded_channel, wait_for_event, LdkServerHandle, TestBitcoind, }; -use ldk_server_client::client::LdkServerClient; +use ldk_server_client::client::{EventStream, LdkServerClient}; use ldk_server_client::error::LdkServerErrorCode::{ AuthError, AuthorizationError, InvalidRequestError, }; @@ -213,7 +213,7 @@ async fn test_macaroon_expiry() { .unwrap() .caveats .contains(&expiry_caveat)); - let events = expiring.subscribe_events().await.unwrap(); + let mut events = expiring.subscribe_events().await.unwrap(); tokio::time::sleep(expiry.duration_since(std::time::SystemTime::now()).unwrap_or_default()) .await; assert_eq!( @@ -221,7 +221,8 @@ async fn test_macaroon_expiry() { AuthorizationError ); assert_eq!(expiring.subscribe_events().await.err().unwrap().error_code, AuthorizationError); - drop(events); + // The open stream ends at expiry instead of waiting for the client to disconnect. + assert_stream_ends_without_events(&mut events, "Macaroon expired").await; } #[tokio::test] @@ -307,7 +308,7 @@ fn test_offline_macaroon_derivation() { } #[tokio::test] -async fn test_revoking_a_root_keeps_existing_event_streams_open() { +async fn test_revoking_a_root_ends_existing_event_streams() { let bitcoind = TestBitcoind::new(); let server_a = LdkServerHandle::start(&bitcoind).await; let server_b = LdkServerHandle::start(&bitcoind).await; @@ -316,6 +317,7 @@ async fn test_revoking_a_root_keeps_existing_event_streams_open() { run_cli(&server_a, &["create-macaroon", "reader", "--permissions", "events:read"]); let client = client_with_macaroon(&server_a, created["token"].as_str().unwrap().to_string()); let mut events = client.subscribe_events().await.unwrap(); + let mut admin_events = server_a.client().subscribe_events().await.unwrap(); run_cli(&server_a, &["revoke-macaroon", created["macaroon"]["id"].as_str().unwrap()]); assert_eq!( @@ -323,10 +325,10 @@ async fn test_revoking_a_root_keeps_existing_event_streams_open() { AuthError ); - // An event created after revocation must still reach the existing subscription. + // An event created after revocation must not reach the existing subscription. run_cli(&server_a, &["close-channel", &channel_id, server_b.node_id()]); mine_and_sync(&bitcoind, &[&server_a, &server_b], 6).await; - wait_for_event(&mut events, |event| { + wait_for_event(&mut admin_events, |event| { matches!( event, Event::ChannelStateChanged(channel_event) @@ -335,4 +337,20 @@ async fn test_revoking_a_root_keeps_existing_event_streams_open() { ) }) .await; + assert_stream_ends_without_events(&mut events, "Macaroon revoked").await; +} + +async fn assert_stream_ends_without_events(events: &mut EventStream, message: &str) { + tokio::time::timeout(Duration::from_secs(10), async { + let error = events + .next_message() + .await + .expect("Stream must end with an error status") + .expect_err("Stream must not forward events after its macaroon is invalid"); + assert_eq!(error.error_code, AuthError); + assert_eq!(error.message, message); + assert!(events.next_message().await.is_none()); + }) + .await + .expect("Timed out waiting for the event stream to end"); } diff --git a/ldk-server-cli/src/pay_wait.rs b/ldk-server-cli/src/pay_wait.rs index f5032a03..cbd85cb9 100644 --- a/ldk-server-cli/src/pay_wait.rs +++ b/ldk-server-cli/src/pay_wait.rs @@ -46,7 +46,7 @@ fn exit_with_payment(code: i32, message: &str, payment: &GetPaymentDetailsRespon pub(crate) async fn pay_and_wait( client: &LdkServerClient, request: UnifiedSendRequest, timeout: Option, ) { - let mut events = client.subscribe_events().await.map_err(handle_error).unwrap(); + let mut events = client.subscribe_payment_events().await.map_err(handle_error).unwrap(); let response = client.unified_send(request).await.map_err(handle_error).unwrap(); diff --git a/ldk-server-client/src/client.rs b/ldk-server-client/src/client.rs index 54315a0c..aec3733f 100644 --- a/ldk-server-client/src/client.rs +++ b/ldk-server-client/src/client.rs @@ -46,9 +46,10 @@ use ldk_server_grpc::api::{ OnchainReceiveResponse, OnchainSendRequest, OnchainSendResponse, OpenChannelRequest, OpenChannelResponse, RevokeMacaroonRequest, RevokeMacaroonResponse, SignMessageRequest, SignMessageResponse, SpliceInRequest, SpliceInResponse, SpliceOutRequest, SpliceOutResponse, - SpontaneousSendRequest, SpontaneousSendResponse, SubscribeEventsRequest, UnifiedSendRequest, - UnifiedSendResponse, UpdateChannelConfigRequest, UpdateChannelConfigResponse, - VerifySignatureRequest, VerifySignatureResponse, + SpontaneousSendRequest, SpontaneousSendResponse, SubscribeChannelEventsRequest, + SubscribeEventsRequest, SubscribeForwardingEventsRequest, SubscribePaymentEventsRequest, + UnifiedSendRequest, UnifiedSendResponse, UpdateChannelConfigRequest, + UpdateChannelConfigResponse, VerifySignatureRequest, VerifySignatureResponse, }; use ldk_server_grpc::endpoints::{ BOLT11_CLAIM_FOR_ID_PATH, BOLT11_FAIL_FOR_ID_PATH, BOLT11_RECEIVE_FOR_HASH_PATH, @@ -67,7 +68,8 @@ use ldk_server_grpc::endpoints::{ LIST_CHANNEL_PAIR_FORWARDING_STATS_PATH, LIST_FORWARDED_PAYMENTS_PATH, LIST_MACAROONS_PATH, LIST_PAYMENTS_PATH, LIST_PEERS_PATH, ONCHAIN_BUMP_FEE_PATH, ONCHAIN_RECEIVE_PATH, ONCHAIN_SEND_PATH, OPEN_CHANNEL_PATH, REVOKE_MACAROON_PATH, SIGN_MESSAGE_PATH, SPLICE_IN_PATH, - SPLICE_OUT_PATH, SPONTANEOUS_SEND_PATH, SUBSCRIBE_EVENTS_PATH, UNIFIED_SEND_PATH, + SPLICE_OUT_PATH, SPONTANEOUS_SEND_PATH, SUBSCRIBE_CHANNEL_EVENTS_PATH, SUBSCRIBE_EVENTS_PATH, + SUBSCRIBE_FORWARDING_EVENTS_PATH, SUBSCRIBE_PAYMENT_EVENTS_PATH, UNIFIED_SEND_PATH, UPDATE_CHANNEL_CONFIG_PATH, VERIFY_SIGNATURE_PATH, }; use ldk_server_grpc::events::EventEnvelope; @@ -563,6 +565,33 @@ impl LdkServerClient { self.grpc_server_streaming(&SubscribeEventsRequest {}, SUBSCRIBE_EVENTS_PATH).await } + /// Subscribe to a stream of channel events via server-streaming gRPC. + /// + /// Returns an [`EventStream`] that only yields channel state change and splice events. + pub async fn subscribe_channel_events(&self) -> Result { + self.grpc_server_streaming(&SubscribeChannelEventsRequest {}, SUBSCRIBE_CHANNEL_EVENTS_PATH) + .await + } + + /// Subscribe to a stream of payment events via server-streaming gRPC. + /// + /// Returns an [`EventStream`] that only yields payment events. + pub async fn subscribe_payment_events(&self) -> Result { + self.grpc_server_streaming(&SubscribePaymentEventsRequest {}, SUBSCRIBE_PAYMENT_EVENTS_PATH) + .await + } + + /// Subscribe to a stream of payment forwarding events via server-streaming gRPC. + /// + /// Returns an [`EventStream`] that only yields payment forwarded events. + pub async fn subscribe_forwarding_events(&self) -> Result { + self.grpc_server_streaming( + &SubscribeForwardingEventsRequest {}, + SUBSCRIBE_FORWARDING_EVENTS_PATH, + ) + .await + } + fn request_macaroon(&self, method: &str, body: &[u8]) -> Result { crate::macaroon::bind_macaroon_to_request(&self.macaroon, method, body) .map_err(|message| LdkServerError::new(InternalError, message)) diff --git a/ldk-server-grpc/src/api.rs b/ldk-server-grpc/src/api.rs index 99f0cfba..84c4a52c 100644 --- a/ldk-server-grpc/src/api.rs +++ b/ldk-server-grpc/src/api.rs @@ -1688,6 +1688,39 @@ pub struct DecodeOfferResponse { #[allow(clippy::derive_partial_eq_without_eq)] #[derive(Clone, PartialEq, ::prost::Message)] pub struct SubscribeEventsRequest {} +/// Subscribe to a best-effort stream of new channel events. +/// +/// Only `ChannelStateChanged`, `SpliceNegotiated`, and `SpliceNegotiationFailed` events are +/// delivered. The same delivery guarantees as `SubscribeEventsRequest` apply. +#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] +#[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))] +#[cfg_attr(feature = "serde", serde(default))] +#[allow(clippy::derive_partial_eq_without_eq)] +#[derive(Clone, PartialEq, ::prost::Message)] +pub struct SubscribeChannelEventsRequest {} +/// Subscribe to a best-effort stream of new payment events. +/// +/// Only `PaymentReceived`, `PaymentSuccessful`, `PaymentFailed`, and `PaymentClaimable` events are +/// delivered. The same delivery guarantees as `SubscribeEventsRequest` apply. +/// +/// If a PaymentClaimable event is missed and the payment is not otherwise claimed or failed, LDK +/// Node automatically fails the HTLC backward at its claim_deadline. +#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] +#[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))] +#[cfg_attr(feature = "serde", serde(default))] +#[allow(clippy::derive_partial_eq_without_eq)] +#[derive(Clone, PartialEq, ::prost::Message)] +pub struct SubscribePaymentEventsRequest {} +/// Subscribe to a best-effort stream of new payment forwarding events. +/// +/// Only `PaymentForwarded` events are delivered. The same delivery guarantees as +/// `SubscribeEventsRequest` apply. +#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] +#[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))] +#[cfg_attr(feature = "serde", serde(default))] +#[allow(clippy::derive_partial_eq_without_eq)] +#[derive(Clone, PartialEq, ::prost::Message)] +pub struct SubscribeForwardingEventsRequest {} /// Macaroon details, without the token or root key. #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] #[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))] diff --git a/ldk-server-grpc/src/endpoints.rs b/ldk-server-grpc/src/endpoints.rs index 8d9223bf..2b27ded3 100644 --- a/ldk-server-grpc/src/endpoints.rs +++ b/ldk-server-grpc/src/endpoints.rs @@ -59,6 +59,9 @@ pub const DECODE_INVOICE_PATH: &str = "DecodeInvoice"; pub const DECODE_OFFER_PATH: &str = "DecodeOffer"; pub const GET_METRICS_PATH: &str = "metrics"; pub const SUBSCRIBE_EVENTS_PATH: &str = "SubscribeEvents"; +pub const SUBSCRIBE_CHANNEL_EVENTS_PATH: &str = "SubscribeChannelEvents"; +pub const SUBSCRIBE_PAYMENT_EVENTS_PATH: &str = "SubscribePaymentEvents"; +pub const SUBSCRIBE_FORWARDING_EVENTS_PATH: &str = "SubscribeForwardingEvents"; pub const GET_FORWARDED_PAYMENT_DETAILS_PATH: &str = "GetForwardedPaymentDetails"; pub const GET_FORWARDED_PAYMENT_TRACKING_MODE_PATH: &str = "GetForwardedPaymentTrackingMode"; pub const GET_CHANNEL_FORWARDING_STATS_PATH: &str = "GetChannelForwardingStats"; diff --git a/ldk-server-grpc/src/proto/api.proto b/ldk-server-grpc/src/proto/api.proto index 8580162d..735d0f40 100644 --- a/ldk-server-grpc/src/proto/api.proto +++ b/ldk-server-grpc/src/proto/api.proto @@ -1189,6 +1189,27 @@ message DecodeOfferResponse { // Node automatically fails the HTLC backward at its claim_deadline. message SubscribeEventsRequest {} +// Subscribe to a best-effort stream of new channel events. +// +// Only `ChannelStateChanged`, `SpliceNegotiated`, and `SpliceNegotiationFailed` events are +// delivered. The same delivery guarantees as `SubscribeEventsRequest` apply. +message SubscribeChannelEventsRequest {} + +// Subscribe to a best-effort stream of new payment events. +// +// Only `PaymentReceived`, `PaymentSuccessful`, `PaymentFailed`, and `PaymentClaimable` events are +// delivered. The same delivery guarantees as `SubscribeEventsRequest` apply. +// +// If a PaymentClaimable event is missed and the payment is not otherwise claimed or failed, LDK +// Node automatically fails the HTLC backward at its claim_deadline. +message SubscribePaymentEventsRequest {} + +// Subscribe to a best-effort stream of new payment forwarding events. +// +// Only `PaymentForwarded` events are delivered. The same delivery guarantees as +// `SubscribeEventsRequest` apply. +message SubscribeForwardingEventsRequest {} + // Macaroon details, without the token or root key. message MacaroonInfo { // The hex ID used to revoke this macaroon. @@ -1350,6 +1371,12 @@ service LightningNode { rpc GraphGetNode(GraphGetNodeRequest) returns (GraphGetNodeResponse); // Subscribe to a stream of server events. rpc SubscribeEvents(SubscribeEventsRequest) returns (stream events.EventEnvelope); + // Subscribe to a stream of channel events. + rpc SubscribeChannelEvents(SubscribeChannelEventsRequest) returns (stream events.EventEnvelope); + // Subscribe to a stream of payment events. + rpc SubscribePaymentEvents(SubscribePaymentEventsRequest) returns (stream events.EventEnvelope); + // Subscribe to a stream of payment forwarding events. + rpc SubscribeForwardingEvents(SubscribeForwardingEventsRequest) returns (stream events.EventEnvelope); // Create a macaroon. Requires macaroons:manage or admin permission. rpc CreateMacaroon(CreateMacaroonRequest) returns (CreateMacaroonResponse); // List macaroons. Requires macaroons:manage or admin permission. diff --git a/ldk-server/src/macaroons/authorization.rs b/ldk-server/src/macaroons/authorization.rs index 35be51c4..c12a1dff 100644 --- a/ldk-server/src/macaroons/authorization.rs +++ b/ldk-server/src/macaroons/authorization.rs @@ -26,8 +26,9 @@ use ldk_server_grpc::endpoints::{ LIST_FORWARDED_PAYMENTS_PATH, LIST_MACAROONS_PATH, LIST_PAYMENTS_PATH, LIST_PEERS_PATH, ONCHAIN_BUMP_FEE_PATH, ONCHAIN_RECEIVE_PATH, ONCHAIN_SEND_PATH, OPEN_CHANNEL_PATH, REVOKE_MACAROON_PATH, SIGN_MESSAGE_PATH, SPLICE_IN_PATH, SPLICE_OUT_PATH, - SPONTANEOUS_SEND_PATH, SUBSCRIBE_EVENTS_PATH, UNIFIED_SEND_PATH, UPDATE_CHANNEL_CONFIG_PATH, - VERIFY_SIGNATURE_PATH, + SPONTANEOUS_SEND_PATH, SUBSCRIBE_CHANNEL_EVENTS_PATH, SUBSCRIBE_EVENTS_PATH, + SUBSCRIBE_FORWARDING_EVENTS_PATH, SUBSCRIBE_PAYMENT_EVENTS_PATH, UNIFIED_SEND_PATH, + UPDATE_CHANNEL_CONFIG_PATH, VERIFY_SIGNATURE_PATH, }; use ldk_server_grpc::permissions::{ CHANNELS_FORCE_CLOSE_PERMISSION, CHANNELS_MANAGE_PERMISSION, CHANNELS_READ_PERMISSION, @@ -105,7 +106,10 @@ pub(crate) fn method_authorization(method: &str) -> MethodAuthorization { DECODE_INVOICE_PATH | DECODE_OFFER_PATH => { MethodAuthorization::Permission(UTILITIES_READ_PERMISSION) }, - SUBSCRIBE_EVENTS_PATH => MethodAuthorization::Permission(EVENTS_READ_PERMISSION), + SUBSCRIBE_EVENTS_PATH + | SUBSCRIBE_CHANNEL_EVENTS_PATH + | SUBSCRIBE_PAYMENT_EVENTS_PATH + | SUBSCRIBE_FORWARDING_EVENTS_PATH => MethodAuthorization::Permission(EVENTS_READ_PERMISSION), CREATE_MACAROON_PATH | LIST_MACAROONS_PATH | REVOKE_MACAROON_PATH => { MethodAuthorization::Permission(MACAROONS_MANAGE_PERMISSION) }, @@ -178,6 +182,9 @@ mod tests { ("GraphListNodes", Some("graph:read")), ("GraphGetNode", Some("graph:read")), ("SubscribeEvents", Some("events:read")), + ("SubscribeChannelEvents", Some("events:read")), + ("SubscribePaymentEvents", Some("events:read")), + ("SubscribeForwardingEvents", Some("events:read")), ("CreateMacaroon", Some("macaroons:manage")), ("ListMacaroons", Some("macaroons:manage")), ("RevokeMacaroon", Some("macaroons:manage")), diff --git a/ldk-server/src/macaroons/policy.rs b/ldk-server/src/macaroons/policy.rs index bf50f57d..e184968a 100644 --- a/ldk-server/src/macaroons/policy.rs +++ b/ldk-server/src/macaroons/policy.rs @@ -37,6 +37,14 @@ impl MacaroonInfo { pub(crate) fn allows(&self, permission: &str) -> bool { self.is_admin() || self.permissions.contains(permission) } + + /// The earliest `time-before` caveat, in seconds since the Unix epoch. + pub(crate) fn expiry(&self) -> Option { + self.caveats + .iter() + .filter_map(|caveat| caveat.strip_prefix("time-before = ")?.parse().ok()) + .min() + } } pub(super) fn mint_token(info: &MacaroonInfo, secret: &str) -> Result { diff --git a/ldk-server/src/macaroons/store.rs b/ldk-server/src/macaroons/store.rs index 8111059c..546a4167 100644 --- a/ldk-server/src/macaroons/store.rs +++ b/ldk-server/src/macaroons/store.rs @@ -20,6 +20,7 @@ use hex::FromHex; use ldk_server_grpc::endpoints::{CREATE_MACAROON_PATH, REVOKE_MACAROON_PATH}; use ldk_server_grpc::permissions::ADMIN_PERMISSION; use ldk_server_macaroons::{Macaroon, RequestBinding, MAX_MACAROON_BYTES}; +use tokio::sync::watch; use super::persistence::{ compute_root_id, generate_secret, is_hex, record_from_stored, write_private_file, @@ -55,6 +56,8 @@ pub(crate) struct MacaroonStore { roots: RwLock>>, management: Mutex<()>, directory: PathBuf, + // Notified after a root is removed so open event streams can recheck their credential. + revocations: watch::Sender<()>, } impl MacaroonStore { @@ -66,8 +69,12 @@ impl MacaroonStore { create_dir_all_private(&directory)?; fs::set_permissions(&directory, fs::Permissions::from_mode(0o700))?; - let mut store = - Self { roots: RwLock::new(HashMap::new()), management: Mutex::new(()), directory }; + let mut store = Self { + roots: RwLock::new(HashMap::new()), + management: Mutex::new(()), + directory, + revocations: watch::Sender::new(()), + }; store.load_root_files()?; if store.roots_mut()?.is_empty() { store.create_initial_admin()?; @@ -222,14 +229,36 @@ impl MacaroonStore { } // Reading the body may take time. Recheck revocation and caveat expiry before // admitting the request, including before opening an event subscription. - if !self.roots.read().map_err(|_| store_lock_error())?.contains_key(&request.info.id) { - return Err(auth_error("Invalid macaroon credentials")); + self.check_still_authorized_at(&request.info, method, now)?; + Ok(request.info) + } + + /// Recheck an admitted credential for revocation and caveat expiry. + /// + /// Long-lived event streams call this before forwarding each event. It only takes the + /// `roots` read lock, never the `management` mutex. + pub(crate) fn check_still_authorized( + &self, info: &MacaroonInfo, method: &str, + ) -> Result<(), LdkServerError> { + self.check_still_authorized_at(info, method, unix_time()?) + } + + fn check_still_authorized_at( + &self, info: &MacaroonInfo, method: &str, now: u64, + ) -> Result<(), LdkServerError> { + if !self.roots.read().map_err(|_| store_lock_error())?.contains_key(&info.id) { + return Err(auth_error("Macaroon revoked")); } - let mut permissions = request.info.permissions.clone(); - for caveat in &request.info.caveats { + let mut permissions = info.permissions.clone(); + for caveat in &info.caveats { check_caveat_at(caveat, method, &mut permissions, now)?; } - Ok(request.info) + Ok(()) + } + + /// Returns a receiver that is notified each time a root is revoked. + pub(crate) fn subscribe_revocations(&self) -> watch::Receiver<()> { + self.revocations.subscribe() } fn authenticate_caveats( @@ -366,6 +395,7 @@ impl MacaroonStore { Err(error) => return Err(internal_error(error)), } self.roots.write().map_err(|_| store_lock_error())?.remove(&id); + self.revocations.send_replace(()); File::open(&self.directory) .and_then(|directory| directory.sync_all()) .map_err(internal_error)?; diff --git a/ldk-server/src/macaroons/store/tests/requests.rs b/ldk-server/src/macaroons/store/tests/requests.rs index 36109fd2..a85c045e 100644 --- a/ldk-server/src/macaroons/store/tests/requests.rs +++ b/ldk-server/src/macaroons/store/tests/requests.rs @@ -279,6 +279,38 @@ fn revocation_and_freshness_are_rechecked_after_reading_the_body() { ); } +#[test] +fn admitted_credentials_are_rechecked_for_revocation_and_expiry() { + let (_directory, store) = test_store("still-authorized"); + let admin = store.authenticate(CREATE_MACAROON_PATH, Some(&admin_token(&store))).unwrap(); + let reader = store.create_root("reader", vec![EVENTS_READ_PERMISSION.into()], &admin).unwrap(); + let expiry = now() + 3600; + let credential = restrict( + &reader.token, + &[&format!("time-before = {}", expiry + 60), &format!("time-before = {expiry}")], + ); + let info = store.authenticate(SUBSCRIBE_EVENTS_PATH, Some(&credential)).unwrap(); + assert_eq!(info.expiry(), Some(expiry)); + store.check_still_authorized(&info, SUBSCRIBE_EVENTS_PATH).unwrap(); + store.check_still_authorized_at(&info, SUBSCRIBE_EVENTS_PATH, expiry - 1).unwrap(); + let expired = + store.check_still_authorized_at(&info, SUBSCRIBE_EVENTS_PATH, expiry).unwrap_err(); + assert_eq!(expired.error_code, LdkServerErrorCode::AuthorizationError); + assert_eq!(expired.message, "Macaroon expired"); + + let mut revocations = store.subscribe_revocations(); + assert!(!revocations.has_changed().unwrap()); + store.revoke_root(&reader.info.id, &admin).unwrap(); + assert!(revocations.has_changed().unwrap()); + revocations.mark_unchanged(); + let revoked = store.check_still_authorized(&info, SUBSCRIBE_EVENTS_PATH).unwrap_err(); + assert_eq!(revoked.error_code, LdkServerErrorCode::AuthError); + assert_eq!(revoked.message, "Macaroon revoked"); + // Revoking one root does not affect credentials from other roots. + store.check_still_authorized(&admin, SUBSCRIBE_EVENTS_PATH).unwrap(); + assert!(!revocations.has_changed().unwrap()); +} + // Deliberately bypass reusable-token checks to test rejection of multiple proofs. fn append_request_proof(token: &str, method: &str, body: &[u8], timestamp: u64) -> String { let proof = RequestBinding::new(method, body, timestamp); diff --git a/ldk-server/src/service.rs b/ldk-server/src/service.rs index a6a307e7..75640376 100644 --- a/ldk-server/src/service.rs +++ b/ldk-server/src/service.rs @@ -10,6 +10,7 @@ use std::future::Future; use std::pin::Pin; use std::sync::Arc; +use std::time::{Duration, SystemTime, UNIX_EPOCH}; use http_body_util::{BodyExt, Limited}; use hyper::body::Incoming; @@ -33,10 +34,11 @@ use ldk_server_grpc::endpoints::{ LIST_FORWARDED_PAYMENTS_PATH, LIST_MACAROONS_PATH, LIST_PAYMENTS_PATH, LIST_PEERS_PATH, ONCHAIN_BUMP_FEE_PATH, ONCHAIN_RECEIVE_PATH, ONCHAIN_SEND_PATH, OPEN_CHANNEL_PATH, REVOKE_MACAROON_PATH, SIGN_MESSAGE_PATH, SPLICE_IN_PATH, SPLICE_OUT_PATH, - SPONTANEOUS_SEND_PATH, SUBSCRIBE_EVENTS_PATH, UNIFIED_SEND_PATH, UPDATE_CHANNEL_CONFIG_PATH, - VERIFY_SIGNATURE_PATH, + SPONTANEOUS_SEND_PATH, SUBSCRIBE_CHANNEL_EVENTS_PATH, SUBSCRIBE_EVENTS_PATH, + SUBSCRIBE_FORWARDING_EVENTS_PATH, SUBSCRIBE_PAYMENT_EVENTS_PATH, UNIFIED_SEND_PATH, + UPDATE_CHANNEL_CONFIG_PATH, VERIFY_SIGNATURE_PATH, }; -use ldk_server_grpc::events::EventEnvelope; +use ldk_server_grpc::events::{event_envelope, EventEnvelope}; use ldk_server_grpc::grpc::{ decode_grpc_body, encode_grpc_frame, grpc_error_response, grpc_response, parse_grpc_timeout, validate_grpc_request, GrpcBody, GrpcStatus, GRPC_STATUS_DEADLINE_EXCEEDED, @@ -45,7 +47,7 @@ use ldk_server_grpc::grpc::{ GRPC_STATUS_UNIMPLEMENTED, }; use prost::Message; -use tokio::sync::{broadcast, mpsc}; +use tokio::sync::{broadcast, mpsc, watch}; use crate::api::bolt11_claim_for_id::handle_bolt11_claim_for_id_request; use crate::api::bolt11_fail_for_id::handle_bolt11_fail_for_id_request; @@ -221,7 +223,13 @@ impl Service> for NodeService { }, }; - let is_streaming = method == SUBSCRIBE_EVENTS_PATH; + let is_streaming = matches!( + method.as_str(), + SUBSCRIBE_EVENTS_PATH + | SUBSCRIBE_CHANNEL_EVENTS_PATH + | SUBSCRIBE_PAYMENT_EVENTS_PATH + | SUBSCRIBE_FORWARDING_EVENTS_PATH + ); let macaroon_store = Arc::clone(&self.macaroon_store); let event_sender = self.event_sender.clone(); let shutdown_rx = self.shutdown_rx.clone(); @@ -443,51 +451,38 @@ impl Service> for NodeService { DECODE_OFFER_PATH => { handle_grpc_unary(context, body_bytes, handle_decode_offer_request).await }, - SUBSCRIBE_EVENTS_PATH => { - // Authorization applies when the subscription starts; revocation does not close it. - let mut shutdown_rx = shutdown_rx; - let mut rx = event_sender.subscribe(); - let (tx, mpsc_rx) = mpsc::channel::>(64); - tokio::spawn(async move { - loop { - tokio::select! { - biased; - _ = shutdown_rx.changed() => { - let _ = tx - .send(Err(GrpcStatus::new( - GRPC_STATUS_UNAVAILABLE, - "server shutting down", - ))) - .await; - break; - }, - result = rx.recv() => { - match result { - Ok(event) => { - let frame = encode_grpc_frame(&event.encode_to_vec()); - if tx.send(Ok(frame)).await.is_err() { - break; // client disconnected - } - }, - Err(tokio::sync::broadcast::error::RecvError::Lagged(_)) => { - continue; // skip missed events, keep streaming - }, - Err(tokio::sync::broadcast::error::RecvError::Closed) => { - let _ = tx - .send(Err(GrpcStatus::new( - GRPC_STATUS_UNAVAILABLE, - "server shutting down", - ))) - .await; - break; - }, - } - } - } - } - }); - Ok(grpc_response(GrpcBody::Stream { rx: mpsc_rx, done: false })) - }, + SUBSCRIBE_EVENTS_PATH => Ok(handle_grpc_event_stream( + macaroon_store, + issuer, + SUBSCRIBE_EVENTS_PATH, + event_sender, + shutdown_rx, + None, + )), + SUBSCRIBE_CHANNEL_EVENTS_PATH => Ok(handle_grpc_event_stream( + macaroon_store, + issuer, + SUBSCRIBE_CHANNEL_EVENTS_PATH, + event_sender, + shutdown_rx, + Some(EventKind::Channel), + )), + SUBSCRIBE_PAYMENT_EVENTS_PATH => Ok(handle_grpc_event_stream( + macaroon_store, + issuer, + SUBSCRIBE_PAYMENT_EVENTS_PATH, + event_sender, + shutdown_rx, + Some(EventKind::Payment), + )), + SUBSCRIBE_FORWARDING_EVENTS_PATH => Ok(handle_grpc_event_stream( + macaroon_store, + issuer, + SUBSCRIBE_FORWARDING_EVENTS_PATH, + event_sender, + shutdown_rx, + Some(EventKind::Forwarding), + )), CREATE_MACAROON_PATH => { let store = Arc::clone(&macaroon_store); handle_grpc_unary(context, body_bytes, move |_context, request| { @@ -575,6 +570,113 @@ async fn handle_grpc_unary< } } +/// Streams events from the broadcast channel to the client. If `kind` is set, only events of that +/// kind are sent. +/// +/// The subscriber's credential is rechecked before each event is sent, after any root is revoked, +/// when its earliest `time-before` caveat passes, and after the subscriber lags. If the check +/// fails, the stream ends with UNAUTHENTICATED. +fn handle_grpc_event_stream( + store: Arc, issuer: Arc, method: &'static str, + event_sender: broadcast::Sender, mut shutdown_rx: watch::Receiver, + kind: Option, +) -> Response { + let mut rx = event_sender.subscribe(); + let (tx, mpsc_rx) = mpsc::channel::>(64); + // Subscribe before the first check so a revocation between the two is not missed. + let mut revocations = store.subscribe_revocations(); + // An expiry too far away to represent never fires; per-event checks still apply. + let until_expiry = issuer + .expiry() + .and_then(|expiry| UNIX_EPOCH.checked_add(Duration::from_secs(expiry))) + .map(|expiry| expiry.duration_since(SystemTime::now()).unwrap_or_default()); + tokio::spawn(async move { + let expiry_timer = tokio::time::sleep(until_expiry.unwrap_or_default()); + tokio::pin!(expiry_timer); + let check = || { + store.check_still_authorized(&issuer, method).map_err(|error| match error.error_code { + LdkServerErrorCode::AuthError | LdkServerErrorCode::AuthorizationError => { + GrpcStatus::new(GRPC_STATUS_UNAUTHENTICATED, error.message) + }, + _ => ldk_error_to_grpc_status(error), + }) + }; + if let Err(status) = check() { + let _ = tx.send(Err(status)).await; + return; + } + loop { + let event = tokio::select! { + biased; + _ = shutdown_rx.changed() => { + let _ = tx + .send(Err(GrpcStatus::new(GRPC_STATUS_UNAVAILABLE, "server shutting down"))) + .await; + break; + }, + // Filtered events are never sent, so a failed send cannot be relied on to detect + // a disconnected client. + _ = tx.closed() => break, + Ok(()) = revocations.changed() => None, + () = &mut expiry_timer, if until_expiry.is_some() => { + // Wall-clock time may lag the timer; retry until the caveat fails. + expiry_timer.as_mut().reset(tokio::time::Instant::now() + Duration::from_secs(1)); + None + }, + result = rx.recv() => match result { + Ok(event) => { + if kind.is_some() && event.event.as_ref().map(event_kind) != kind { + continue; + } + Some(event) + }, + // Skip missed events, but recheck the credential before continuing. + Err(broadcast::error::RecvError::Lagged(_)) => None, + Err(broadcast::error::RecvError::Closed) => { + let _ = tx + .send(Err(GrpcStatus::new(GRPC_STATUS_UNAVAILABLE, "server shutting down"))) + .await; + break; + }, + }, + }; + if let Err(status) = check() { + let _ = tx.send(Err(status)).await; + break; + } + if let Some(event) = event { + let frame = encode_grpc_frame(&event.encode_to_vec()); + if tx.send(Ok(frame)).await.is_err() { + break; // client disconnected + } + } + } + }); + grpc_response(GrpcBody::Stream { rx: mpsc_rx, done: false }) +} + +/// The kinds of events that can be subscribed to separately. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +enum EventKind { + Channel, + Payment, + Forwarding, +} + +/// Exhaustive so that every new event must be assigned a kind. +fn event_kind(event: &event_envelope::Event) -> EventKind { + match event { + event_envelope::Event::ChannelStateChanged(_) + | event_envelope::Event::SpliceNegotiated(_) + | event_envelope::Event::SpliceNegotiationFailed(_) => EventKind::Channel, + event_envelope::Event::PaymentReceived(_) + | event_envelope::Event::PaymentSuccessful(_) + | event_envelope::Event::PaymentFailed(_) + | event_envelope::Event::PaymentClaimable(_) => EventKind::Payment, + event_envelope::Event::PaymentForwarded(_) => EventKind::Forwarding, + } +} + fn request_content_length(headers: &HeaderMap) -> Result, GrpcStatus> { let Some(content_length) = headers.get("content-length") else { return Ok(None); @@ -871,6 +973,99 @@ mod tests { assert_eq!(error.message, "Macaroon expired"); } + type EventStreamReceiver = mpsc::Receiver>; + + const SUBSCRIPTIONS: [(&str, Option); 4] = [ + (SUBSCRIBE_EVENTS_PATH, None), + (SUBSCRIBE_CHANNEL_EVENTS_PATH, Some(EventKind::Channel)), + (SUBSCRIBE_PAYMENT_EVENTS_PATH, Some(EventKind::Payment)), + (SUBSCRIBE_FORWARDING_EVENTS_PATH, Some(EventKind::Forwarding)), + ]; + + fn open_event_stream( + store: &Arc, credential: &str, + event_sender: &broadcast::Sender, method: &'static str, + kind: Option, + ) -> (watch::Sender, EventStreamReceiver) { + let issuer = store.authenticate(method, Some(credential)).unwrap(); + let (shutdown_tx, shutdown_rx) = watch::channel(false); + let response = handle_grpc_event_stream( + Arc::clone(store), + issuer, + method, + event_sender.clone(), + shutdown_rx, + kind, + ); + let GrpcBody::Stream { rx, .. } = response.into_body() else { + panic!("Event subscriptions must return a streaming body"); + }; + (shutdown_tx, rx) + } + + async fn next_status(stream: &mut EventStreamReceiver) -> GrpcStatus { + tokio::time::timeout(Duration::from_secs(10), stream.recv()) + .await + .expect("Timed out waiting for the stream to end") + .expect("Stream closed without a status") + .expect_err("Stream forwarded an event after its credential became invalid") + } + + #[tokio::test] + async fn revoking_a_root_ends_its_event_streams() { + use ldk_server_grpc::permissions::EVENTS_READ_PERMISSION; + + let (_directory, store) = test_store("stream-revocation"); + let store = Arc::new(store); + let admin = store.authenticate(CREATE_MACAROON_PATH, Some(&admin_token(&store))).unwrap(); + let reader = + store.create_root("reader", vec![EVENTS_READ_PERMISSION.into()], &admin).unwrap(); + let (event_sender, _) = broadcast::channel(16); + let mut streams: Vec<_> = SUBSCRIPTIONS + .into_iter() + .map(|(method, kind)| { + open_event_stream(&store, &reader.token, &event_sender, method, kind) + }) + .collect(); + event_sender.send(EventEnvelope::default()).unwrap(); + let (_, unfiltered) = &mut streams[0]; + unfiltered.recv().await.unwrap().unwrap(); + + // Idle streams of every kind end on revocation without waiting for another event. + store.revoke_root(&reader.info.id, &admin).unwrap(); + for ((method, _), (_, stream)) in SUBSCRIPTIONS.into_iter().zip(&mut streams) { + let status = next_status(stream).await; + assert_eq!(status.code, GRPC_STATUS_UNAUTHENTICATED, "{method}"); + assert_eq!(status.message, "Macaroon revoked", "{method}"); + } + let _ = event_sender.send(EventEnvelope::default()); + for (_, stream) in &mut streams { + assert!(stream.recv().await.is_none()); + } + } + + #[tokio::test] + async fn expired_credentials_end_their_event_streams() { + use crate::macaroons::test_util::{now, restrict}; + + let (_directory, store) = test_store("stream-expiry"); + let store = Arc::new(store); + let expiry = now() + 2; + let credential = restrict(&admin_token(&store), &[&format!("time-before = {expiry}")]); + let (event_sender, _) = broadcast::channel(16); + let (_shutdown_tx, mut stream) = + open_event_stream(&store, &credential, &event_sender, SUBSCRIBE_EVENTS_PATH, None); + event_sender.send(EventEnvelope::default()).unwrap(); + stream.recv().await.unwrap().unwrap(); + + let status = next_status(&mut stream).await; + assert!(now() >= expiry, "Stream ended before its credential expired"); + assert_eq!(status.code, GRPC_STATUS_UNAUTHENTICATED); + assert_eq!(status.message, "Macaroon expired"); + let _ = event_sender.send(EventEnvelope::default()); + assert!(stream.recv().await.is_none()); + } + #[test] fn test_request_content_length_missing() { let headers = HeaderMap::new(); @@ -921,4 +1116,34 @@ mod tests { assert_eq!(err.code, GRPC_STATUS_INVALID_ARGUMENT); assert_eq!(err.message, "Request body length does not match content-length"); } + + #[tokio::test] + async fn filtered_event_stream_stops_when_client_disconnects() { + let (_directory, store) = test_store("stream-disconnect"); + let issuer = + store.authenticate(SUBSCRIBE_CHANNEL_EVENTS_PATH, Some(&admin_token(&store))).unwrap(); + let (event_sender, _) = broadcast::channel(16); + let (_shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(false); + let response = handle_grpc_event_stream( + Arc::new(store), + issuer, + SUBSCRIBE_CHANNEL_EVENTS_PATH, + event_sender.clone(), + shutdown_rx, + Some(EventKind::Channel), + ); + assert_eq!(event_sender.receiver_count(), 1); + + // Dropping the response disconnects the client. Filtered events are never sent, so the + // stream task must notice the disconnect without relying on a failed send. + drop(response); + event_sender.send(EventEnvelope::default()).unwrap(); + tokio::time::timeout(std::time::Duration::from_secs(5), async { + while event_sender.receiver_count() > 0 { + tokio::time::sleep(std::time::Duration::from_millis(10)).await; + } + }) + .await + .expect("event stream task did not stop after the client disconnected"); + } }