11use crate :: logging:: Logger ;
2- use crate :: store:: { self , PaymentId } ;
2+ use crate :: store:: { self , MppOutcome , PaymentId , TxMetadataStore , TxType } ;
33
44use crate :: dyn_store:: DynStore ;
55use ldk_node:: bitcoin:: hashes:: Hash ;
@@ -15,7 +15,7 @@ use ldk_node::lightning_types::payment::{PaymentHash, PaymentPreimage};
1515use ldk_node:: payment:: { ConfirmationStatus , PaymentKind } ;
1616use ldk_node:: { CustomTlvRecord , UserChannelId } ;
1717
18- use std:: collections:: VecDeque ;
18+ use std:: collections:: { HashMap , VecDeque } ;
1919use std:: sync:: Arc ;
2020use std:: task:: { Poll , Waker } ;
2121use std:: time:: SystemTime ;
@@ -209,21 +209,114 @@ impl_writeable_tlv_based_enum!(Event,
209209/// [`Wallet`]: [`crate::Wallet`]
210210pub struct EventQueue {
211211 queue : Arc < Mutex < VecDeque < Event > > > ,
212+ pending_mpp_events : Arc < Mutex < HashMap < PaymentHash , Vec < Event > > > > ,
212213 waker : Arc < Mutex < Option < Waker > > > ,
213214 kv_store : Arc < dyn DynStore > ,
215+ tx_metadata : TxMetadataStore ,
214216 logger : Arc < Logger > ,
215217}
216218
217219impl EventQueue {
218- pub ( crate ) fn new ( kv_store : Arc < dyn DynStore > , logger : Arc < Logger > ) -> Self {
220+ pub ( crate ) fn new (
221+ kv_store : Arc < dyn DynStore > , tx_metadata : TxMetadataStore , logger : Arc < Logger > ,
222+ ) -> Self {
219223 let queue = Arc :: new ( Mutex :: new ( VecDeque :: new ( ) ) ) ;
224+ let pending_mpp_events = Arc :: new ( Mutex :: new ( HashMap :: new ( ) ) ) ;
220225 let waker = Arc :: new ( Mutex :: new ( None ) ) ;
221- Self { queue, waker, kv_store, logger }
226+ Self { queue, pending_mpp_events, waker, kv_store, tx_metadata, logger }
227+ }
228+
229+ /// Starts buffering terminal events for a multi-path payment while its leg metadata is being
230+ /// registered.
231+ pub ( crate ) async fn begin_mpp_setup ( & self , payment_hash : PaymentHash ) {
232+ self . pending_mpp_events . lock ( ) . await . entry ( payment_hash) . or_default ( ) ;
233+ }
234+
235+ /// Stops buffering terminal events for a multi-path payment and re-processes any events that
236+ /// arrived before the leg metadata was registered.
237+ pub ( crate ) async fn finish_mpp_setup (
238+ & self , payment_hash : PaymentHash ,
239+ ) -> Result < ( ) , ldk_node:: lightning:: io:: Error > {
240+ let pending_events = self . pending_mpp_events . lock ( ) . await . remove ( & payment_hash) ;
241+ if let Some ( events) = pending_events {
242+ for event in events {
243+ self . add_event ( event) . await ?;
244+ }
245+ }
246+ Ok ( ( ) )
222247 }
223248
224249 pub ( crate ) async fn add_event (
225250 & self , event : Event ,
226251 ) -> Result < ( ) , ldk_node:: lightning:: io:: Error > {
252+ // Outgoing payments split across the trusted and lightning wallets emit a terminal event per
253+ // leg. Record each leg's result onto the shared, persisted metadata; the leg that completes
254+ // the payment yields the single combined event we surface instead of the per-leg ones.
255+ match & event {
256+ Event :: PaymentSuccessful { payment_id, payment_preimage, fee_paid_msat, .. }
257+ if self . is_mpp_leg ( payment_id) =>
258+ {
259+ let combined = self
260+ . tx_metadata
261+ . record_mpp_leg (
262+ * payment_id,
263+ Some ( ( fee_paid_msat. unwrap_or ( 0 ) , payment_preimage. 0 ) ) ,
264+ )
265+ . await ;
266+ return self . push_combined_mpp ( combined) . await ;
267+ } ,
268+ Event :: PaymentFailed { payment_id, .. } if self . is_mpp_leg ( payment_id) => {
269+ let combined = self . tx_metadata . record_mpp_leg ( * payment_id, None ) . await ;
270+ return self . push_combined_mpp ( combined) . await ;
271+ } ,
272+ _ => { } ,
273+ }
274+
275+ if let Some ( payment_hash) = terminal_payment_hash ( & event) {
276+ let mut pending_mpp_events = self . pending_mpp_events . lock ( ) . await ;
277+ if let Some ( events) = pending_mpp_events. get_mut ( & payment_hash) {
278+ events. push ( event) ;
279+ return Ok ( ( ) ) ;
280+ }
281+ }
282+
283+ self . push_event ( event) . await
284+ }
285+
286+ /// Whether `id` identifies a leg of a multi-path payment.
287+ fn is_mpp_leg ( & self , id : & PaymentId ) -> bool {
288+ matches ! ( self . tx_metadata. read( ) . get( id) . map( |m| m. ty) , Some ( TxType :: MppPayment { .. } ) )
289+ }
290+
291+ /// Surfaces the single combined event for a multi-path payment, or nothing if the payment is
292+ /// still waiting on its other leg (or already produced its combined event).
293+ async fn push_combined_mpp (
294+ & self , combined : Option < ( PaymentId , MppOutcome ) > ,
295+ ) -> Result < ( ) , ldk_node:: lightning:: io:: Error > {
296+ match combined {
297+ Some ( ( surface_id, MppOutcome :: Succeeded { payment_hash, preimage, fee_msat } ) ) => {
298+ self . push_event ( Event :: PaymentSuccessful {
299+ payment_id : surface_id,
300+ payment_hash : PaymentHash ( payment_hash) ,
301+ payment_preimage : PaymentPreimage ( preimage) ,
302+ fee_paid_msat : Some ( fee_msat) ,
303+ } )
304+ . await
305+ } ,
306+ Some ( ( surface_id, MppOutcome :: Failed { payment_hash } ) ) => {
307+ self . push_event ( Event :: PaymentFailed {
308+ payment_id : surface_id,
309+ payment_hash : Some ( PaymentHash ( payment_hash) ) ,
310+ reason : None ,
311+ } )
312+ . await
313+ } ,
314+ None => Ok ( ( ) ) ,
315+ }
316+ }
317+
318+ /// Appends an event to the queue and persists it, waking any pending consumer.
319+ async fn push_event ( & self , event : Event ) -> Result < ( ) , ldk_node:: lightning:: io:: Error > {
227320 {
228321 let mut locked_queue = self . queue . lock ( ) . await ;
229322 locked_queue. push_back ( event) ;
@@ -285,6 +378,104 @@ impl EventQueue {
285378 }
286379}
287380
381+ fn terminal_payment_hash ( event : & Event ) -> Option < PaymentHash > {
382+ match event {
383+ Event :: PaymentSuccessful { payment_hash, .. } => Some ( * payment_hash) ,
384+ Event :: PaymentFailed { payment_hash, .. } => * payment_hash,
385+ _ => None ,
386+ }
387+ }
388+
389+ #[ cfg( test) ]
390+ mod tests {
391+ use super :: * ;
392+ use crate :: logging:: LoggerType ;
393+ use crate :: store:: { PaymentId , PaymentType , TxMetadata , TxMetadataStore , TxType } ;
394+ use ldk_node:: io:: sqlite_store:: SqliteStore ;
395+ use std:: path:: PathBuf ;
396+ use std:: time:: { Duration , UNIX_EPOCH } ;
397+
398+ fn temp_sqlite_store ( ) -> ( PathBuf , Arc < dyn DynStore > ) {
399+ let path = std:: env:: temp_dir ( ) . join ( format ! (
400+ "orange-sdk-event-mpp-buffer-test-{}" ,
401+ SystemTime :: now( ) . duration_since( UNIX_EPOCH ) . unwrap( ) . as_nanos( )
402+ ) ) ;
403+ let store = SqliteStore :: new ( path. clone ( ) , Some ( "orange.sqlite" . to_string ( ) ) , None )
404+ . expect ( "sqlite store" ) ;
405+ ( path, Arc :: new ( store) )
406+ }
407+
408+ fn mpp_metadata ( surface_id : PaymentId , lightning_leg : [ u8 ; 32 ] ) -> TxMetadata {
409+ TxMetadata {
410+ ty : TxType :: MppPayment {
411+ surface_id,
412+ lightning_leg,
413+ total_amount_msat : 200_000 ,
414+ ty : PaymentType :: OutgoingLightningBolt11 { payment_preimage : None } ,
415+ trusted_fee_msat : None ,
416+ lightning_fee_msat : None ,
417+ preimage : None ,
418+ failed : false ,
419+ finalized : false ,
420+ } ,
421+ time : Duration :: from_secs ( 1 ) ,
422+ }
423+ }
424+
425+ #[ tokio:: test]
426+ async fn pending_mpp_setup_buffers_terminal_events_until_metadata_exists ( ) {
427+ let ( _path, store) = temp_sqlite_store ( ) ;
428+ let tx_metadata = TxMetadataStore :: new ( Arc :: clone ( & store) ) . await ;
429+ let queue = EventQueue :: new (
430+ store,
431+ tx_metadata. clone ( ) ,
432+ Arc :: new ( Logger :: new ( & LoggerType :: LogFacade ) . expect ( "logger" ) ) ,
433+ ) ;
434+
435+ let payment_hash = PaymentHash ( [ 3u8 ; 32 ] ) ;
436+ let surface_id = PaymentId :: Trusted ( [ 7u8 ; 32 ] ) ;
437+ let lightning_id = PaymentId :: SelfCustodial ( payment_hash. 0 ) ;
438+ let preimage = PaymentPreimage ( [ 1u8 ; 32 ] ) ;
439+
440+ queue. begin_mpp_setup ( payment_hash) . await ;
441+ queue
442+ . add_event ( Event :: PaymentSuccessful {
443+ payment_id : surface_id,
444+ payment_hash,
445+ payment_preimage : preimage,
446+ fee_paid_msat : Some ( 1_000 ) ,
447+ } )
448+ . await
449+ . expect ( "buffer event" ) ;
450+ assert_eq ! ( queue. next_event( ) . await , None ) ;
451+
452+ tx_metadata. insert ( surface_id, mpp_metadata ( surface_id, payment_hash. 0 ) ) . await ;
453+ tx_metadata. upsert ( lightning_id, mpp_metadata ( surface_id, payment_hash. 0 ) ) . await ;
454+ queue. finish_mpp_setup ( payment_hash) . await . expect ( "replay buffered events" ) ;
455+ assert_eq ! ( queue. next_event( ) . await , None ) ;
456+
457+ queue
458+ . add_event ( Event :: PaymentSuccessful {
459+ payment_id : lightning_id,
460+ payment_hash,
461+ payment_preimage : preimage,
462+ fee_paid_msat : Some ( 2_000 ) ,
463+ } )
464+ . await
465+ . expect ( "complete mpp" ) ;
466+
467+ assert_eq ! (
468+ queue. next_event( ) . await ,
469+ Some ( Event :: PaymentSuccessful {
470+ payment_id: surface_id,
471+ payment_hash,
472+ payment_preimage: preimage,
473+ fee_paid_msat: Some ( 3_000 ) ,
474+ } )
475+ ) ;
476+ }
477+ }
478+
288479struct EventQueueSerWrapper < ' a > ( & ' a VecDeque < Event > ) ;
289480
290481impl Writeable for EventQueueSerWrapper < ' _ > {
0 commit comments