1use core::fmt;
2use std::{
3 collections::BTreeMap,
4 sync::atomic::{AtomicBool, AtomicU64, Ordering},
5};
6
7use platform_utils::time::Instant;
8use serde::Serialize;
9use tokio::sync::{Mutex, RwLock};
10use tracing::info;
11use uuid::Uuid;
12
13use crate::{DepositInfo, LightningAddressInfo, Payment, sdk::RuntimeEvent};
14
15#[allow(clippy::large_enum_variant)]
17#[derive(Debug, Clone, Serialize)]
18#[cfg_attr(feature = "uniffi", derive(uniffi::Enum))]
19pub enum SdkEvent {
20 Synced,
22 UnclaimedDeposits {
25 unclaimed_deposits: Vec<DepositInfo>,
26 },
27 ClaimedDeposits { claimed_deposits: Vec<DepositInfo> },
30 PaymentSucceeded { payment: Payment },
33 PaymentPending { payment: Payment },
36 PaymentFailed { payment: Payment },
38 AutoOptimization {
45 optimization_event: AutoOptimizationEvent,
47 },
48 LightningAddressChanged {
51 lightning_address: Option<LightningAddressInfo>,
52 },
53 NewDeposits { new_deposits: Vec<DepositInfo> },
56 UnilateralExitStateChanged,
59}
60
61impl SdkEvent {
62 pub(crate) fn from_payment(payment: Payment) -> Self {
63 match payment.status {
64 crate::PaymentStatus::Completed => SdkEvent::PaymentSucceeded { payment },
65 crate::PaymentStatus::Pending => SdkEvent::PaymentPending { payment },
66 crate::PaymentStatus::Failed => SdkEvent::PaymentFailed { payment },
67 }
68 }
69}
70
71impl fmt::Display for SdkEvent {
72 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
73 match self {
74 SdkEvent::Synced => write!(f, "Synced"),
75 SdkEvent::UnclaimedDeposits { unclaimed_deposits } => {
76 write!(f, "UnclaimedDeposits: {unclaimed_deposits:?}")
77 }
78 SdkEvent::ClaimedDeposits { claimed_deposits } => {
79 write!(f, "ClaimedDeposits: {claimed_deposits:?}")
80 }
81 SdkEvent::PaymentSucceeded { payment } => {
82 write!(f, "PaymentSucceeded: {payment:?}")
83 }
84 SdkEvent::PaymentPending { payment } => {
85 write!(f, "PaymentPending: {payment:?}")
86 }
87 SdkEvent::PaymentFailed { payment } => {
88 write!(f, "PaymentFailed: {payment:?}")
89 }
90 SdkEvent::AutoOptimization {
91 optimization_event: event,
92 } => {
93 write!(f, "AutoOptimization: {event:?}")
94 }
95 SdkEvent::LightningAddressChanged { lightning_address } => {
96 write!(f, "LightningAddressChanged: {lightning_address:?}")
97 }
98 SdkEvent::NewDeposits { new_deposits } => {
99 write!(f, "NewDeposits: {new_deposits:?}")
100 }
101 SdkEvent::UnilateralExitStateChanged => write!(f, "UnilateralExitStateChanged"),
102 }
103 }
104}
105
106#[derive(Debug, Clone, Serialize)]
107#[cfg_attr(feature = "uniffi", derive(uniffi::Enum))]
108pub enum AutoOptimizationEvent {
109 Started { total_rounds: u32 },
111 RoundCompleted {
113 current_round: u32,
114 total_rounds: u32,
115 },
116 Completed,
118 Cancelled,
120 Failed { error: String },
122 Skipped,
124}
125
126#[allow(clippy::struct_excessive_bools)]
127#[derive(Debug, Default)]
128pub struct InternalSyncedEvent {
129 pub wallet: bool,
130 pub wallet_state: bool,
131 pub deposits: bool,
132 pub lnurl_metadata: bool,
133 pub storage_incoming: Option<u32>,
134}
135
136impl InternalSyncedEvent {
137 pub fn any(&self) -> bool {
138 self.wallet
139 || self.wallet_state
140 || self.deposits
141 || self.lnurl_metadata
142 || self.storage_incoming.is_some()
143 }
144
145 pub fn any_non_zero(&self) -> bool {
146 self.wallet
147 || self.wallet_state
148 || self.deposits
149 || self.lnurl_metadata
150 || self.storage_incoming.is_some_and(|v| v > 0)
151 }
152
153 pub fn merge(&self, other: &InternalSyncedEvent) -> Self {
154 Self {
155 wallet: self.wallet || other.wallet,
156 wallet_state: self.wallet_state || other.wallet_state,
157 deposits: self.deposits || other.deposits,
158 lnurl_metadata: self.lnurl_metadata || other.lnurl_metadata,
159 storage_incoming: self
160 .storage_incoming
161 .zip(other.storage_incoming)
162 .map(|(a, b)| a.saturating_add(b))
163 .or(self.storage_incoming)
164 .or(other.storage_incoming),
165 }
166 }
167}
168
169#[cfg_attr(feature = "uniffi", uniffi::export(callback_interface))]
171#[macros::async_trait]
172pub trait EventListener: Send + Sync {
173 async fn on_event(&self, event: SdkEvent);
175}
176
177#[macros::async_trait]
185pub trait EventMiddleware: Send + Sync {
186 async fn process(&self, event: SdkEvent) -> Option<SdkEvent>;
188}
189
190pub struct EventEmitter {
197 has_real_time_sync: bool,
198 rtsync_failed: AtomicBool,
199 listener_index: AtomicU64,
200 runtime_event_handlers: RwLock<Vec<Box<dyn RuntimeEventHandler>>>,
201 internal_listeners: RwLock<BTreeMap<String, Box<dyn EventListener>>>,
203 middleware: RwLock<Vec<Box<dyn EventMiddleware>>>,
205 external_listeners: RwLock<BTreeMap<String, Box<dyn EventListener>>>,
207 synced_event_buffer: Mutex<Option<InternalSyncedEvent>>,
208}
209
210#[macros::async_trait]
216pub(crate) trait RuntimeEventHandler: Send + Sync {
217 async fn handle(&self, emitter: &EventEmitter, event: RuntimeEvent);
218}
219
220impl EventEmitter {
221 pub fn new(has_real_time_sync: bool) -> Self {
223 Self {
224 has_real_time_sync,
225 rtsync_failed: AtomicBool::new(false),
226 listener_index: AtomicU64::new(0),
227 runtime_event_handlers: RwLock::new(Vec::new()),
228 internal_listeners: RwLock::new(BTreeMap::new()),
229 middleware: RwLock::new(Vec::new()),
230 external_listeners: RwLock::new(BTreeMap::new()),
231 synced_event_buffer: Mutex::new(Some(InternalSyncedEvent::default())),
232 }
233 }
234
235 pub async fn add_external_listener(&self, listener: Box<dyn EventListener>) -> String {
245 let index = self.listener_index.fetch_add(1, Ordering::Relaxed);
246 let id = format!("listener_{}-{}", index, Uuid::new_v4());
247 let mut listeners = self.external_listeners.write().await;
248 listeners.insert(id.clone(), listener);
249 id
250 }
251
252 pub async fn remove_external_listener(&self, id: &str) -> bool {
262 let mut listeners = self.external_listeners.write().await;
263 listeners.remove(id).is_some()
264 }
265
266 pub async fn clear_external_listeners(&self) {
272 let mut listeners = self.external_listeners.write().await;
273 listeners.clear();
274 }
275
276 pub async fn add_internal_listener(&self, listener: Box<dyn EventListener>) -> String {
281 let index = self.listener_index.fetch_add(1, Ordering::Relaxed);
282 let id = format!("internal_{}-{}", index, Uuid::new_v4());
283 let mut listeners = self.internal_listeners.write().await;
284 listeners.insert(id.clone(), listener);
285 id
286 }
287
288 pub async fn remove_internal_listener(&self, id: &str) -> bool {
290 let mut listeners = self.internal_listeners.write().await;
291 listeners.remove(id).is_some()
292 }
293
294 pub async fn add_middleware(&self, middleware: Box<dyn EventMiddleware>) {
304 let mut mw = self.middleware.write().await;
305 mw.push(middleware);
306 }
307
308 pub(crate) async fn add_runtime_event_handler(&self, handler: Box<dyn RuntimeEventHandler>) {
309 let mut handlers = self.runtime_event_handlers.write().await;
310 handlers.push(handler);
311 }
312
313 pub(crate) async fn emit_runtime_event(&self, event: RuntimeEvent) {
314 let handlers = self.runtime_event_handlers.read().await;
315 for handler in handlers.iter() {
316 handler.handle(self, event.clone()).await;
317 }
318 }
319
320 pub async fn emit(&self, event: &SdkEvent) {
325 let start = Instant::now();
326 let event_label = format!("{event}");
327 let mut internal_total = std::time::Duration::ZERO;
328 let mut middleware_total = std::time::Duration::ZERO;
329 let mut external_total = std::time::Duration::ZERO;
330
331 let internal = self.internal_listeners.read().await;
333 let internal_count = internal.len();
334 for (id, listener) in internal.iter() {
335 let t = Instant::now();
336 listener.on_event(event.clone()).await;
337 let dt = t.elapsed();
338 internal_total = internal_total.saturating_add(dt);
339 info!("emit({event_label}) internal listener {id}: {dt:?}");
340 }
341 drop(internal);
342
343 let mut event = Some(event.clone());
345 let middleware = self.middleware.read().await;
346 let middleware_count = middleware.len();
347 for (i, mw) in middleware.iter().enumerate() {
348 if let Some(e) = event {
349 let t = Instant::now();
350 event = mw.process(e).await;
351 let dt = t.elapsed();
352 middleware_total = middleware_total.saturating_add(dt);
353 info!("emit({event_label}) middleware #{i}: {dt:?}");
354 } else {
355 break;
356 }
357 }
358 drop(middleware);
359
360 let mut external_count = 0;
362 if let Some(ref event) = event {
363 let listeners = self.external_listeners.read().await;
364 external_count = listeners.len();
365 for (id, listener) in listeners.iter() {
366 let t = Instant::now();
367 listener.on_event(event.clone()).await;
368 let dt = t.elapsed();
369 external_total = external_total.saturating_add(dt);
370 info!("emit({event_label}) external listener {id}: {dt:?}");
371 }
372 }
373
374 info!(
375 "emit({event_label}) completed in {:?} (internal[{}]={:?}, middleware[{}]={:?}, external[{}]={:?})",
376 start.elapsed(),
377 internal_count,
378 internal_total,
379 middleware_count,
380 middleware_total,
381 external_count,
382 external_total
383 );
384 }
385
386 pub async fn emit_synced(&self, synced: &InternalSyncedEvent) {
387 if !synced.any() {
388 return;
390 }
391
392 let mut mtx = self.synced_event_buffer.lock().await;
393
394 let is_first_event = if let Some(buffered) = &*mtx {
395 let merged = buffered.merge(synced);
396
397 if merged.wallet
400 && (!self.has_real_time_sync
401 || merged.storage_incoming.is_some()
402 || self.rtsync_failed.load(Ordering::Relaxed))
403 {
404 *mtx = None;
405 } else {
406 *mtx = Some(merged);
407 return;
408 }
409
410 true
411 } else {
412 false
413 };
414
415 drop(mtx);
416
417 if !is_first_event && !synced.any_non_zero() {
419 return;
420 }
421
422 self.emit(&SdkEvent::Synced).await;
424 }
425
426 pub async fn notify_rtsync_failed(&self) {
430 self.rtsync_failed.store(true, Ordering::Relaxed);
431
432 let mut mtx = self.synced_event_buffer.lock().await;
433 if let Some(buffered) = &*mtx
434 && buffered.wallet
435 {
436 *mtx = None;
437 drop(mtx);
438 self.emit(&SdkEvent::Synced).await;
439 }
440 }
441}
442
443impl Default for EventEmitter {
444 fn default() -> Self {
445 Self::new(false)
446 }
447}
448
449#[cfg(test)]
450mod tests {
451 use super::*;
452 use std::sync::Arc;
453 use std::sync::atomic::{AtomicBool, Ordering};
454
455 use macros::async_test_all;
456
457 #[cfg(feature = "browser-tests")]
458 wasm_bindgen_test::wasm_bindgen_test_configure!(run_in_browser);
459
460 struct TestListener {
461 received: Arc<AtomicBool>,
462 }
463
464 #[macros::async_trait]
465 impl EventListener for TestListener {
466 async fn on_event(&self, _event: SdkEvent) {
467 self.received.store(true, Ordering::Relaxed);
468 }
469 }
470
471 #[async_test_all]
472 async fn test_event_emission() {
473 let emitter = EventEmitter::new(false);
474 let received = Arc::new(AtomicBool::new(false));
475
476 let listener = Box::new(TestListener {
478 received: received.clone(),
479 });
480
481 let _ = emitter.add_external_listener(listener).await;
482
483 let event = SdkEvent::Synced {};
484
485 emitter.emit(&event).await;
486
487 assert!(received.load(Ordering::Relaxed));
489 }
490
491 #[async_test_all]
492 async fn test_remove_listener() {
493 let emitter = EventEmitter::new(false);
494
495 let received1 = Arc::new(AtomicBool::new(false));
497 let received2 = Arc::new(AtomicBool::new(false));
498
499 let listener1 = Box::new(TestListener {
501 received: received1.clone(),
502 });
503
504 let listener2 = Box::new(TestListener {
505 received: received2.clone(),
506 });
507
508 let id1 = emitter.add_external_listener(listener1).await;
509 let id2 = emitter.add_external_listener(listener2).await;
510
511 assert!(emitter.remove_external_listener(&id1).await);
513
514 let event = SdkEvent::Synced {};
516 emitter.emit(&event).await;
517
518 assert!(!received1.load(Ordering::Relaxed));
520
521 assert!(received2.load(Ordering::Relaxed));
523
524 assert!(emitter.remove_external_listener(&id2).await);
526
527 assert!(!emitter.remove_external_listener("non-existent-id").await);
529 }
530
531 #[async_test_all]
532 async fn test_clear_external_listeners() {
533 let emitter = EventEmitter::new(false);
534
535 let received1 = Arc::new(AtomicBool::new(false));
536 let received2 = Arc::new(AtomicBool::new(false));
537
538 let id1 = emitter
539 .add_external_listener(Box::new(TestListener {
540 received: received1.clone(),
541 }))
542 .await;
543 let id2 = emitter
544 .add_external_listener(Box::new(TestListener {
545 received: received2.clone(),
546 }))
547 .await;
548
549 emitter.clear_external_listeners().await;
550
551 emitter.emit(&SdkEvent::Synced).await;
552
553 assert!(!received1.load(Ordering::Relaxed));
554 assert!(!received2.load(Ordering::Relaxed));
555
556 assert!(!emitter.remove_external_listener(&id1).await);
558 assert!(!emitter.remove_external_listener(&id2).await);
559 }
560
561 #[async_test_all]
562 async fn test_synced_event_only_emitted_with_wallet_sync() {
563 let emitter = EventEmitter::new(false);
564 let received = Arc::new(AtomicBool::new(false));
565
566 let listener = Box::new(TestListener {
567 received: received.clone(),
568 });
569
570 emitter.add_external_listener(listener).await;
571
572 emitter
574 .emit_synced(&InternalSyncedEvent {
575 wallet: false,
576 wallet_state: true,
577 deposits: true,
578 lnurl_metadata: true,
579 storage_incoming: None,
580 })
581 .await;
582
583 assert!(!received.load(Ordering::Relaxed));
584
585 emitter
587 .emit_synced(&InternalSyncedEvent {
588 wallet: true,
589 wallet_state: false,
590 deposits: false,
591 lnurl_metadata: false,
592 storage_incoming: Some(1),
593 })
594 .await;
595
596 assert!(received.load(Ordering::Relaxed));
597 }
598
599 #[async_test_all]
600 async fn test_has_real_time_sync_synced_event_only_emitted_with_wallet_and_storage_sync() {
601 let emitter = EventEmitter::new(true);
602 let received = Arc::new(AtomicBool::new(false));
603
604 let listener = Box::new(TestListener {
605 received: received.clone(),
606 });
607
608 emitter.add_external_listener(listener).await;
609
610 emitter
612 .emit_synced(&InternalSyncedEvent {
613 wallet: false,
614 wallet_state: false,
615 deposits: false,
616 lnurl_metadata: false,
617 storage_incoming: Some(0),
618 })
619 .await;
620
621 assert!(!received.load(Ordering::Relaxed));
622
623 emitter
625 .emit_synced(&InternalSyncedEvent {
626 wallet: true,
627 wallet_state: false,
628 deposits: false,
629 lnurl_metadata: false,
630 storage_incoming: None,
631 })
632 .await;
633
634 assert!(received.load(Ordering::Relaxed));
635 }
636
637 #[async_test_all]
638 async fn test_has_real_time_sync_synced_event_only_emitted_with_wallet_and_storage_sync_reverse()
639 {
640 let emitter = EventEmitter::new(true);
641 let received = Arc::new(AtomicBool::new(false));
642
643 let listener = Box::new(TestListener {
644 received: received.clone(),
645 });
646
647 emitter.add_external_listener(listener).await;
648
649 emitter
651 .emit_synced(&InternalSyncedEvent {
652 wallet: true,
653 wallet_state: false,
654 deposits: false,
655 lnurl_metadata: false,
656 storage_incoming: None,
657 })
658 .await;
659
660 assert!(!received.load(Ordering::Relaxed));
661
662 emitter
664 .emit_synced(&InternalSyncedEvent {
665 wallet: false,
666 wallet_state: false,
667 deposits: false,
668 lnurl_metadata: false,
669 storage_incoming: Some(0),
670 })
671 .await;
672
673 assert!(received.load(Ordering::Relaxed));
674 }
675
676 #[async_test_all]
677 async fn test_rtsync_failed_emits_synced_on_wallet_alone() {
678 let emitter = EventEmitter::new(true);
679 let received = Arc::new(AtomicBool::new(false));
680
681 let listener = Box::new(TestListener {
682 received: received.clone(),
683 });
684
685 emitter.add_external_listener(listener).await;
686
687 emitter
689 .emit_synced(&InternalSyncedEvent {
690 wallet: true,
691 wallet_state: false,
692 deposits: false,
693 lnurl_metadata: false,
694 storage_incoming: None,
695 })
696 .await;
697
698 assert!(!received.load(Ordering::Relaxed));
699
700 emitter.notify_rtsync_failed().await;
702
703 assert!(received.load(Ordering::Relaxed));
704 }
705
706 #[async_test_all]
707 async fn test_rtsync_failed_before_wallet_sync_emits_on_wallet() {
708 let emitter = EventEmitter::new(true);
709 let received = Arc::new(AtomicBool::new(false));
710
711 let listener = Box::new(TestListener {
712 received: received.clone(),
713 });
714
715 emitter.add_external_listener(listener).await;
716
717 emitter.notify_rtsync_failed().await;
719
720 assert!(!received.load(Ordering::Relaxed));
721
722 emitter
724 .emit_synced(&InternalSyncedEvent {
725 wallet: true,
726 wallet_state: false,
727 deposits: false,
728 lnurl_metadata: false,
729 storage_incoming: None,
730 })
731 .await;
732
733 assert!(received.load(Ordering::Relaxed));
734 }
735
736 #[async_test_all]
737 async fn test_synced_event_buffers_until_wallet_sync() {
738 let emitter = EventEmitter::new(false);
739 let received = Arc::new(AtomicBool::new(false));
740
741 let listener = Box::new(TestListener {
742 received: received.clone(),
743 });
744
745 emitter.add_external_listener(listener).await;
746
747 emitter
749 .emit_synced(&InternalSyncedEvent {
750 wallet: false,
751 wallet_state: true,
752 deposits: false,
753 lnurl_metadata: false,
754 storage_incoming: None,
755 })
756 .await;
757
758 assert!(!received.load(Ordering::Relaxed));
759
760 emitter
761 .emit_synced(&InternalSyncedEvent {
762 wallet: false,
763 wallet_state: false,
764 deposits: true,
765 lnurl_metadata: false,
766 storage_incoming: None,
767 })
768 .await;
769
770 assert!(!received.load(Ordering::Relaxed));
771
772 emitter
773 .emit_synced(&InternalSyncedEvent {
774 wallet: false,
775 wallet_state: false,
776 deposits: false,
777 lnurl_metadata: true,
778 storage_incoming: None,
779 })
780 .await;
781
782 assert!(!received.load(Ordering::Relaxed));
783
784 emitter
785 .emit_synced(&InternalSyncedEvent {
786 wallet: false,
787 wallet_state: false,
788 deposits: false,
789 lnurl_metadata: false,
790 storage_incoming: None,
791 })
792 .await;
793
794 assert!(!received.load(Ordering::Relaxed));
795
796 emitter
798 .emit_synced(&InternalSyncedEvent {
799 wallet: true,
800 wallet_state: false,
801 deposits: false,
802 lnurl_metadata: false,
803 storage_incoming: None,
804 })
805 .await;
806
807 assert!(received.load(Ordering::Relaxed));
808 }
809
810 #[async_test_all]
811 async fn test_synced_event_all_true() {
812 let emitter = EventEmitter::new(false);
813 let received = Arc::new(AtomicBool::new(false));
814
815 let listener = Box::new(TestListener {
816 received: received.clone(),
817 });
818
819 emitter.add_external_listener(listener).await;
820
821 emitter
823 .emit_synced(&InternalSyncedEvent {
824 wallet: true,
825 wallet_state: true,
826 deposits: true,
827 lnurl_metadata: true,
828 storage_incoming: Some(1),
829 })
830 .await;
831
832 assert!(received.load(Ordering::Relaxed));
833 }
834
835 #[async_test_all]
836 async fn test_synced_event_empty_does_not_emit() {
837 let emitter = EventEmitter::new(false);
838 let received = Arc::new(AtomicBool::new(false));
839
840 let listener = Box::new(TestListener {
841 received: received.clone(),
842 });
843
844 emitter.add_external_listener(listener).await;
845
846 emitter
848 .emit_synced(&InternalSyncedEvent {
849 wallet: false,
850 wallet_state: false,
851 deposits: false,
852 lnurl_metadata: false,
853 storage_incoming: None,
854 })
855 .await;
856
857 assert!(!received.load(Ordering::Relaxed));
858 }
859
860 #[async_test_all]
861 async fn test_subsequent_syncs_after_wallet_emit_immediately() {
862 use std::sync::atomic::AtomicUsize;
863
864 struct CountingListener {
865 count: Arc<AtomicUsize>,
866 }
867
868 #[macros::async_trait]
869 impl EventListener for CountingListener {
870 async fn on_event(&self, event: SdkEvent) {
871 if matches!(event, SdkEvent::Synced) {
872 self.count.fetch_add(1, Ordering::Relaxed);
873 }
874 }
875 }
876
877 let emitter = EventEmitter::new(true);
878 let count = Arc::new(AtomicUsize::new(0));
879
880 let listener = Box::new(CountingListener {
881 count: count.clone(),
882 });
883
884 emitter.add_external_listener(listener).await;
885
886 emitter
888 .emit_synced(&InternalSyncedEvent {
889 wallet: true,
890 wallet_state: false,
891 deposits: false,
892 lnurl_metadata: false,
893 storage_incoming: Some(0),
894 })
895 .await;
896
897 assert_eq!(count.load(Ordering::Relaxed), 1);
898
899 emitter
901 .emit_synced(&InternalSyncedEvent {
902 wallet: false,
903 wallet_state: true,
904 deposits: false,
905 lnurl_metadata: false,
906 storage_incoming: None,
907 })
908 .await;
909
910 assert_eq!(count.load(Ordering::Relaxed), 2);
911
912 emitter
914 .emit_synced(&InternalSyncedEvent {
915 wallet: false,
916 wallet_state: false,
917 deposits: true,
918 lnurl_metadata: false,
919 storage_incoming: None,
920 })
921 .await;
922
923 assert_eq!(count.load(Ordering::Relaxed), 3);
924
925 emitter
926 .emit_synced(&InternalSyncedEvent {
927 wallet: false,
928 wallet_state: false,
929 deposits: false,
930 lnurl_metadata: true,
931 storage_incoming: None,
932 })
933 .await;
934
935 assert_eq!(count.load(Ordering::Relaxed), 4);
936
937 emitter
938 .emit_synced(&InternalSyncedEvent {
939 wallet: false,
940 wallet_state: false,
941 deposits: false,
942 lnurl_metadata: false,
943 storage_incoming: Some(1),
944 })
945 .await;
946
947 assert_eq!(count.load(Ordering::Relaxed), 5);
948
949 emitter
951 .emit_synced(&InternalSyncedEvent {
952 wallet: false,
953 wallet_state: false,
954 deposits: false,
955 lnurl_metadata: false,
956 storage_incoming: Some(0),
957 })
958 .await;
959
960 assert_eq!(count.load(Ordering::Relaxed), 5);
961
962 emitter
963 .emit_synced(&InternalSyncedEvent {
964 wallet: true,
965 wallet_state: false,
966 deposits: false,
967 lnurl_metadata: false,
968 storage_incoming: None,
969 })
970 .await;
971
972 assert_eq!(count.load(Ordering::Relaxed), 6);
973 }
974
975 struct RecordingListener {
979 events: Arc<Mutex<Vec<String>>>,
980 }
981
982 impl RecordingListener {
983 fn new() -> (Self, Arc<Mutex<Vec<String>>>) {
984 let events = Arc::new(Mutex::new(Vec::new()));
985 (
986 Self {
987 events: events.clone(),
988 },
989 events,
990 )
991 }
992 }
993
994 #[macros::async_trait]
995 impl EventListener for RecordingListener {
996 async fn on_event(&self, event: SdkEvent) {
997 self.events.lock().await.push(format!("{event}"));
998 }
999 }
1000
1001 struct SuppressSyncedMiddleware;
1003
1004 #[macros::async_trait]
1005 impl EventMiddleware for SuppressSyncedMiddleware {
1006 async fn process(&self, event: SdkEvent) -> Option<SdkEvent> {
1007 match event {
1008 SdkEvent::Synced => None,
1009 other => Some(other),
1010 }
1011 }
1012 }
1013
1014 struct DowngradePaymentMiddleware;
1016
1017 #[macros::async_trait]
1018 impl EventMiddleware for DowngradePaymentMiddleware {
1019 async fn process(&self, event: SdkEvent) -> Option<SdkEvent> {
1020 match event {
1021 SdkEvent::PaymentSucceeded { payment } => {
1022 Some(SdkEvent::PaymentPending { payment })
1023 }
1024 other => Some(other),
1025 }
1026 }
1027 }
1028
1029 struct SuppressAllMiddleware;
1031
1032 #[macros::async_trait]
1033 impl EventMiddleware for SuppressAllMiddleware {
1034 async fn process(&self, _event: SdkEvent) -> Option<SdkEvent> {
1035 None
1036 }
1037 }
1038
1039 fn test_payment() -> Payment {
1040 Payment {
1041 id: "test-id".to_string(),
1042 payment_type: crate::PaymentType::Receive,
1043 status: crate::PaymentStatus::Completed,
1044 amount: 1000,
1045 fees: 10,
1046 timestamp: 123_456,
1047 method: crate::PaymentMethod::Spark,
1048 details: None,
1049 conversion_details: None,
1050 }
1051 }
1052
1053 #[async_test_all]
1056 async fn test_internal_listener_receives_events() {
1057 let emitter = EventEmitter::new(false);
1058 let (listener, events) = RecordingListener::new();
1059
1060 emitter.add_internal_listener(Box::new(listener)).await;
1061
1062 emitter.emit(&SdkEvent::Synced).await;
1063
1064 let log = events.lock().await;
1065 assert_eq!(log.len(), 1);
1066 assert!(log[0].contains("Synced"));
1067 }
1068
1069 #[async_test_all]
1070 async fn test_remove_internal_listener() {
1071 let emitter = EventEmitter::new(false);
1072 let (listener, events) = RecordingListener::new();
1073
1074 let id = emitter.add_internal_listener(Box::new(listener)).await;
1075
1076 assert!(emitter.remove_internal_listener(&id).await);
1077 assert!(!emitter.remove_internal_listener(&id).await);
1078
1079 emitter.emit(&SdkEvent::Synced).await;
1080
1081 assert!(events.lock().await.is_empty());
1082 }
1083
1084 #[async_test_all]
1087 async fn test_middleware_suppresses_event_for_external() {
1088 let emitter = EventEmitter::new(false);
1089 let (ext, ext_events) = RecordingListener::new();
1090
1091 emitter.add_external_listener(Box::new(ext)).await;
1092 emitter
1093 .add_middleware(Box::new(SuppressSyncedMiddleware))
1094 .await;
1095
1096 emitter.emit(&SdkEvent::Synced).await;
1097
1098 assert!(ext_events.lock().await.is_empty());
1100 }
1101
1102 #[async_test_all]
1103 async fn test_middleware_transforms_event() {
1104 let emitter = EventEmitter::new(false);
1105 let (ext, ext_events) = RecordingListener::new();
1106
1107 emitter.add_external_listener(Box::new(ext)).await;
1108 emitter
1109 .add_middleware(Box::new(DowngradePaymentMiddleware))
1110 .await;
1111
1112 let event = SdkEvent::PaymentSucceeded {
1113 payment: test_payment(),
1114 };
1115 emitter.emit(&event).await;
1116
1117 let log = ext_events.lock().await;
1118 assert_eq!(log.len(), 1);
1119 assert!(log[0].contains("PaymentPending"));
1120 }
1121
1122 #[async_test_all]
1123 async fn test_middleware_passthrough_unmatched_events() {
1124 let emitter = EventEmitter::new(false);
1125 let (ext, ext_events) = RecordingListener::new();
1126
1127 emitter.add_external_listener(Box::new(ext)).await;
1128 emitter
1129 .add_middleware(Box::new(SuppressSyncedMiddleware))
1130 .await;
1131
1132 emitter.emit(&SdkEvent::Synced).await;
1134 let event = SdkEvent::PaymentSucceeded {
1135 payment: test_payment(),
1136 };
1137 emitter.emit(&event).await;
1138
1139 let log = ext_events.lock().await;
1140 assert_eq!(log.len(), 1);
1141 assert!(log[0].contains("PaymentSucceeded"));
1142 }
1143
1144 #[async_test_all]
1145 async fn test_middleware_chain_ordering() {
1146 let emitter = EventEmitter::new(false);
1147 let (ext, ext_events) = RecordingListener::new();
1148
1149 emitter.add_external_listener(Box::new(ext)).await;
1150
1151 emitter
1153 .add_middleware(Box::new(DowngradePaymentMiddleware))
1154 .await;
1155 emitter
1157 .add_middleware(Box::new(SuppressSyncedMiddleware))
1158 .await;
1159
1160 let event = SdkEvent::PaymentSucceeded {
1161 payment: test_payment(),
1162 };
1163 emitter.emit(&event).await;
1164
1165 let log = ext_events.lock().await;
1166 assert_eq!(log.len(), 1);
1167 assert!(log[0].contains("PaymentPending"));
1168 }
1169
1170 #[async_test_all]
1171 async fn test_suppress_all_middleware_stops_chain() {
1172 let emitter = EventEmitter::new(false);
1173 let (ext, ext_events) = RecordingListener::new();
1174
1175 emitter.add_external_listener(Box::new(ext)).await;
1176
1177 emitter
1179 .add_middleware(Box::new(SuppressAllMiddleware))
1180 .await;
1181 emitter
1182 .add_middleware(Box::new(DowngradePaymentMiddleware))
1183 .await;
1184
1185 emitter.emit(&SdkEvent::Synced).await;
1186 let event = SdkEvent::PaymentSucceeded {
1187 payment: test_payment(),
1188 };
1189 emitter.emit(&event).await;
1190
1191 assert!(ext_events.lock().await.is_empty());
1192 }
1193
1194 #[async_test_all]
1197 async fn test_three_phase_emit_flow() {
1198 let emitter = EventEmitter::new(false);
1199 let (int, int_events) = RecordingListener::new();
1200 let (ext, ext_events) = RecordingListener::new();
1201
1202 emitter.add_internal_listener(Box::new(int)).await;
1203 emitter.add_external_listener(Box::new(ext)).await;
1204 emitter
1205 .add_middleware(Box::new(SuppressSyncedMiddleware))
1206 .await;
1207
1208 emitter.emit(&SdkEvent::Synced).await;
1210
1211 assert_eq!(int_events.lock().await.len(), 1);
1212 assert!(ext_events.lock().await.is_empty());
1213
1214 let event = SdkEvent::PaymentSucceeded {
1216 payment: test_payment(),
1217 };
1218 emitter.emit(&event).await;
1219
1220 assert_eq!(int_events.lock().await.len(), 2);
1221 assert_eq!(ext_events.lock().await.len(), 1);
1222 }
1223
1224 #[async_test_all]
1225 async fn test_internal_sees_raw_event_external_sees_transformed() {
1226 let emitter = EventEmitter::new(false);
1227 let (int, int_events) = RecordingListener::new();
1228 let (ext, ext_events) = RecordingListener::new();
1229
1230 emitter.add_internal_listener(Box::new(int)).await;
1231 emitter.add_external_listener(Box::new(ext)).await;
1232 emitter
1233 .add_middleware(Box::new(DowngradePaymentMiddleware))
1234 .await;
1235
1236 let event = SdkEvent::PaymentSucceeded {
1237 payment: test_payment(),
1238 };
1239 emitter.emit(&event).await;
1240
1241 let int_log = int_events.lock().await;
1242 let ext_log = ext_events.lock().await;
1243
1244 assert_eq!(int_log.len(), 1);
1246 assert!(int_log[0].contains("PaymentSucceeded"));
1247
1248 assert_eq!(ext_log.len(), 1);
1250 assert!(ext_log[0].contains("PaymentPending"));
1251 }
1252
1253 #[async_test_all]
1254 async fn test_no_listeners_no_middleware_does_not_panic() {
1255 let emitter = EventEmitter::new(false);
1256 emitter.emit(&SdkEvent::Synced).await;
1257 }
1259
1260 #[async_test_all]
1261 async fn test_empty_event_does_not_emit_after_wallet_sync() {
1262 use std::sync::atomic::AtomicUsize;
1263
1264 struct CountingListener {
1265 count: Arc<AtomicUsize>,
1266 }
1267
1268 #[macros::async_trait]
1269 impl EventListener for CountingListener {
1270 async fn on_event(&self, event: SdkEvent) {
1271 if matches!(event, SdkEvent::Synced) {
1272 self.count.fetch_add(1, Ordering::Relaxed);
1273 }
1274 }
1275 }
1276
1277 let emitter = EventEmitter::new(false);
1278 let count = Arc::new(AtomicUsize::new(0));
1279
1280 let listener = Box::new(CountingListener {
1281 count: count.clone(),
1282 });
1283
1284 emitter.add_external_listener(listener).await;
1285
1286 emitter
1288 .emit_synced(&InternalSyncedEvent {
1289 wallet: true,
1290 wallet_state: false,
1291 deposits: false,
1292 lnurl_metadata: false,
1293 storage_incoming: None,
1294 })
1295 .await;
1296
1297 assert_eq!(count.load(Ordering::Relaxed), 1);
1298
1299 emitter
1301 .emit_synced(&InternalSyncedEvent {
1302 wallet: false,
1303 wallet_state: false,
1304 deposits: false,
1305 lnurl_metadata: false,
1306 storage_incoming: None,
1307 })
1308 .await;
1309
1310 assert_eq!(count.load(Ordering::Relaxed), 1); emitter
1314 .emit_synced(&InternalSyncedEvent {
1315 wallet: false,
1316 wallet_state: true,
1317 deposits: false,
1318 lnurl_metadata: false,
1319 storage_incoming: None,
1320 })
1321 .await;
1322
1323 assert_eq!(count.load(Ordering::Relaxed), 2); }
1325}