Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
19 changes: 13 additions & 6 deletions docs/api-guide.md
Original file line number Diff line number Diff line change
Expand Up @@ -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`,
Expand Down Expand Up @@ -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 |
|---------------------|-----------------------------------------------------------------------|
Expand All @@ -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.
Expand Down
80 changes: 80 additions & 0 deletions e2e-tests/tests/e2e.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand Down
2 changes: 1 addition & 1 deletion ldk-server-cli/src/pay_wait.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<Duration>,
) {
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();

Expand Down
37 changes: 33 additions & 4 deletions ldk-server-client/src/client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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;
Expand Down Expand Up @@ -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<EventStream, LdkServerError> {
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<EventStream, LdkServerError> {
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<EventStream, LdkServerError> {
self.grpc_server_streaming(
&SubscribeForwardingEventsRequest {},
SUBSCRIBE_FORWARDING_EVENTS_PATH,
)
.await
}

fn request_macaroon(&self, method: &str, body: &[u8]) -> Result<String, LdkServerError> {
crate::macaroon::bind_macaroon_to_request(&self.macaroon, method, body)
.map_err(|message| LdkServerError::new(InternalError, message))
Expand Down
33 changes: 33 additions & 0 deletions ldk-server-grpc/src/api.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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"))]
Expand Down
3 changes: 3 additions & 0 deletions ldk-server-grpc/src/endpoints.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down
27 changes: 27 additions & 0 deletions ldk-server-grpc/src/proto/api.proto
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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.
Expand Down
13 changes: 10 additions & 3 deletions ldk-server/src/macaroons/authorization.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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)
},
Expand Down Expand Up @@ -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")),
Expand Down
Loading
Loading