From 3d58be11a190a8fb3554e14b5fea74f07d18309e Mon Sep 17 00:00:00 2001 From: benthecarman Date: Mon, 28 Sep 2026 17:16:42 -0500 Subject: [PATCH 1/4] 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/4] 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/4] 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(); From 80060fed0bf6c333babd7306966f567d64552f98 Mon Sep 17 00:00:00 2001 From: benthecarman Date: Mon, 5 Oct 2026 18:03:56 -0500 Subject: [PATCH 4/4] End event streams when their macaroon is invalid Event subscriptions were authorized only when they started. After a root was revoked or a time-before caveat passed, new requests were rejected, but an open stream kept forwarding events, including payment preimages, until the client disconnected. The shared event stream helper now rechecks the credential against the store before forwarding each event and after the subscriber lags, so every subscription kind is covered. Revocation notifies open streams and the earliest expiry arms a timer, so idle streams also close promptly. A failed check ends the stream with UNAUTHENTICATED ("Macaroon revoked" or "Macaroon expired"). The recheck shares its root and caveat checks with request admission so the two cannot drift apart, and takes only the roots read lock. Co-Authored-By: Claude Opus 5.5 --- docs/api-guide.md | 5 +- e2e-tests/tests/macaroons.rs | 30 ++- ldk-server/src/macaroons/policy.rs | 8 + ldk-server/src/macaroons/store.rs | 44 +++- .../src/macaroons/store/tests/requests.rs | 32 +++ ldk-server/src/service.rs | 221 +++++++++++++++--- 6 files changed, 287 insertions(+), 53 deletions(-) diff --git a/docs/api-guide.md b/docs/api-guide.md index 449b5ede..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. 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/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 6ce311a8..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; @@ -46,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; @@ -450,20 +451,34 @@ 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, None)) - }, + 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), @@ -558,53 +573,81 @@ 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. /// -/// Authorization applies when the subscription starts; revocation does not close it. +/// 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( - event_sender: broadcast::Sender, - mut shutdown_rx: tokio::sync::watch::Receiver, kind: Option, + 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 { - tokio::select! { + let event = tokio::select! { biased; _ = shutdown_rx.changed() => { let _ = tx - .send(Err(GrpcStatus::new( - GRPC_STATUS_UNAVAILABLE, - "server shutting down", - ))) + .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, - 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 - } - }, - 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(()) = 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 } } } @@ -930,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(); @@ -983,10 +1119,19 @@ mod tests { #[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(event_sender.clone(), shutdown_rx, Some(EventKind::Channel)); + 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