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-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/service.rs b/ldk-server/src/service.rs index a6a307e7..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(); @@ -444,50 +451,23 @@ impl Service> for NodeService { 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 })) - }, + 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| { @@ -575,6 +555,85 @@ 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. +fn handle_grpc_event_stream( + event_sender: broadcast::Sender, + mut shutdown_rx: tokio::sync::watch::Receiver, kind: Option, +) -> 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; + }, + // 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; + }, + } + } + } + } + }); + 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); @@ -921,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"); + } }