From 3d58be11a190a8fb3554e14b5fea74f07d18309e Mon Sep 17 00:00:00 2001 From: benthecarman Date: Mon, 28 Sep 2026 17:16:42 -0500 Subject: [PATCH 1/3] Extract event stream handling into a helper Move the SubscribeEvents streaming loop out of the request dispatch match into its own function so it can be reused by additional event streaming RPCs. No behavior change. Co-Authored-By: Claude Opus 5.5 --- ldk-server/src/service.rs | 96 +++++++++++++++++++++------------------ 1 file changed, 51 insertions(+), 45 deletions(-) diff --git a/ldk-server/src/service.rs b/ldk-server/src/service.rs index a6a307e7..bddc09e7 100644 --- a/ldk-server/src/service.rs +++ b/ldk-server/src/service.rs @@ -443,51 +443,7 @@ 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(event_sender, shutdown_rx)), CREATE_MACAROON_PATH => { let store = Arc::clone(&macaroon_store); handle_grpc_unary(context, body_bytes, move |_context, request| { @@ -575,6 +531,56 @@ async fn handle_grpc_unary< } } +/// Streams events from the broadcast channel to the client. +/// +/// Authorization applies when the subscription starts; revocation does not close it. +fn handle_grpc_event_stream( + event_sender: broadcast::Sender, + mut shutdown_rx: tokio::sync::watch::Receiver, +) -> Response { + 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; + }, + } + } + } + } + }); + grpc_response(GrpcBody::Stream { rx: mpsc_rx, done: false }) +} + fn request_content_length(headers: &HeaderMap) -> Result, GrpcStatus> { let Some(content_length) = headers.get("content-length") else { return Ok(None); From b49bdfaac4b89e206b43ed517c9c858740705711 Mon Sep 17 00:00:00 2001 From: benthecarman Date: Mon, 28 Sep 2026 17:16:42 -0500 Subject: [PATCH 2/3] Add channel, payment, and forwarding subscriptions SubscribeEvents delivers every event, so clients that only care about one kind of event have to receive and discard everything else. Add SubscribeChannelEvents, SubscribePaymentEvents, and SubscribeForwardingEvents RPCs, which stream only the matching subset of events. SubscribeEvents is unchanged. Channel events are ChannelStateChanged, SpliceNegotiated, and SpliceNegotiationFailed. Payment events are PaymentReceived, PaymentSuccessful, PaymentFailed, and PaymentClaimable. Forwarding events are PaymentForwarded; they get their own stream because they are the highest-volume event on a routing node and are not this node's own payments. Based on an earlier contribution that added an event kind filter to SubscribeEventsRequest; reworked into separate RPCs per review. Co-authored-by: Ekong Jemimah Co-Authored-By: Claude Opus 5.5 --- docs/api-guide.md | 19 +++-- e2e-tests/tests/e2e.rs | 80 +++++++++++++++++++++ ldk-server-client/src/client.rs | 37 ++++++++-- ldk-server-grpc/src/api.rs | 33 +++++++++ ldk-server-grpc/src/endpoints.rs | 3 + ldk-server-grpc/src/proto/api.proto | 27 +++++++ ldk-server/src/macaroons/authorization.rs | 13 +++- ldk-server/src/service.rs | 88 +++++++++++++++++++++-- 8 files changed, 280 insertions(+), 20 deletions(-) diff --git a/docs/api-guide.md b/docs/api-guide.md index 6e430653..449b5ede 100644 --- a/docs/api-guide.md +++ b/docs/api-guide.md @@ -111,7 +111,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 +309,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 +334,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/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/service.rs b/ldk-server/src/service.rs index bddc09e7..6ce311a8 100644 --- a/ldk-server/src/service.rs +++ b/ldk-server/src/service.rs @@ -33,10 +33,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, @@ -221,7 +222,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,7 +450,24 @@ impl Service> for NodeService { DECODE_OFFER_PATH => { handle_grpc_unary(context, body_bytes, handle_decode_offer_request).await }, - SUBSCRIBE_EVENTS_PATH => Ok(handle_grpc_event_stream(event_sender, shutdown_rx)), + SUBSCRIBE_EVENTS_PATH => { + Ok(handle_grpc_event_stream(event_sender, shutdown_rx, None)) + }, + SUBSCRIBE_CHANNEL_EVENTS_PATH => Ok(handle_grpc_event_stream( + event_sender, + shutdown_rx, + Some(EventKind::Channel), + )), + SUBSCRIBE_PAYMENT_EVENTS_PATH => Ok(handle_grpc_event_stream( + event_sender, + shutdown_rx, + Some(EventKind::Payment), + )), + SUBSCRIBE_FORWARDING_EVENTS_PATH => Ok(handle_grpc_event_stream( + 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| { @@ -531,12 +555,13 @@ async fn handle_grpc_unary< } } -/// Streams events from the broadcast channel to the client. +/// Streams events from the broadcast channel to the client. If `kind` is set, only events of that +/// kind are sent. /// /// Authorization applies when the subscription starts; revocation does not close it. fn handle_grpc_event_stream( event_sender: broadcast::Sender, - mut shutdown_rx: tokio::sync::watch::Receiver, + mut shutdown_rx: tokio::sync::watch::Receiver, kind: Option, ) -> Response { let mut rx = event_sender.subscribe(); let (tx, mpsc_rx) = mpsc::channel::>(64); @@ -553,9 +578,15 @@ fn handle_grpc_event_stream( .await; break; }, + // Filtered events are never sent, so a failed send cannot be relied on to detect + // a disconnected client. + _ = tx.closed() => break, result = rx.recv() => { match result { Ok(event) => { + if kind.is_some() && event.event.as_ref().map(event_kind) != kind { + continue; + } let frame = encode_grpc_frame(&event.encode_to_vec()); if tx.send(Ok(frame)).await.is_err() { break; // client disconnected @@ -581,6 +612,28 @@ fn handle_grpc_event_stream( 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); @@ -927,4 +980,25 @@ 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 (event_sender, _) = broadcast::channel(16); + let (_shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(false); + let response = + handle_grpc_event_stream(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"); + } } From a145008520c39e113706283a5f8abedf4a691305 Mon Sep 17 00:00:00 2001 From: benthecarman Date: Mon, 28 Sep 2026 14:31:19 -0500 Subject: [PATCH 3/3] Use payment event stream for pay --wait pay --wait only looks at PaymentSuccessful and PaymentFailed events, so subscribe to payment events rather than every server event. Co-Authored-By: Claude Opus 5.5 --- ldk-server-cli/src/pay_wait.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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();