Skip to main content

breez_sdk_spark/
events.rs

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/// Events emitted by the SDK
16#[allow(clippy::large_enum_variant)]
17#[derive(Debug, Clone, Serialize)]
18#[cfg_attr(feature = "uniffi", derive(uniffi::Enum))]
19pub enum SdkEvent {
20    /// Emitted when the wallet has been synchronized with the network
21    Synced,
22    /// Emitted when the SDK was unable to claim deposits. Each deposit carries a
23    /// `claim_error` with the reason.
24    UnclaimedDeposits {
25        unclaimed_deposits: Vec<DepositInfo>,
26    },
27    /// Emitted when deposits were claimed into the wallet. The resulting payment
28    /// is emitted separately as `PaymentSucceeded`.
29    ClaimedDeposits { claimed_deposits: Vec<DepositInfo> },
30    /// Emitted when a payment completed. The cached balance is refreshed before
31    /// this event, so `get_info` returns the new value.
32    PaymentSucceeded { payment: Payment },
33    /// Emitted when a payment is in flight. The same payment is emitted again as
34    /// succeeded or failed once it settles.
35    PaymentPending { payment: Payment },
36    /// Emitted when a payment failed.
37    PaymentFailed { payment: Payment },
38    /// Emitted while the background auto-optimizer is running.
39    ///
40    /// Only fired from the auto path (enabled via
41    /// `LeafOptimizationConfig::auto_enabled`). Manually-triggered runs
42    /// via `BreezSdk::optimize_leaves` do not emit events: they return an
43    /// `OptimizationOutcome` instead.
44    AutoOptimization {
45        // Named with `optimization` prefix to avoid collision with `event` keyword in C#
46        optimization_event: AutoOptimizationEvent,
47    },
48    /// Emitted when the Lightning address changed on another device. The address
49    /// is unset when it was deleted.
50    LightningAddressChanged {
51        lightning_address: Option<LightningAddressInfo>,
52    },
53    /// Emitted when on-chain deposits are detected. Only deposits whose
54    /// `is_mature` is set have enough confirmations to be claimed.
55    NewDeposits { new_deposits: Vec<DepositInfo> },
56    /// Emitted when the data a unilateral exit is built from has changed, so an
57    /// exit state exported earlier is out of date and should be exported again.
58    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    /// Optimization has started with the given number of rounds.
110    Started { total_rounds: u32 },
111    /// A round has completed.
112    RoundCompleted {
113        current_round: u32,
114        total_rounds: u32,
115    },
116    /// Optimization completed successfully.
117    Completed,
118    /// Optimization was cancelled.
119    Cancelled,
120    /// Optimization failed with an error.
121    Failed { error: String },
122    /// Optimization was skipped because leaves are already optimal.
123    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/// Trait for event listeners
170#[cfg_attr(feature = "uniffi", uniffi::export(callback_interface))]
171#[macros::async_trait]
172pub trait EventListener: Send + Sync {
173    /// Called when an event occurs
174    async fn on_event(&self, event: SdkEvent);
175}
176
177/// Middleware that can intercept and transform events before they reach external listeners.
178///
179/// Middleware processes events in a chain. Each middleware receives the event from the
180/// previous one and can:
181/// - Pass it through unchanged: `Some(event)`
182/// - Transform it: `Some(modified_event)`
183/// - Suppress it: `None`
184#[macros::async_trait]
185pub trait EventMiddleware: Send + Sync {
186    /// Process an event. Return `Some` to forward (possibly modified), `None` to suppress.
187    async fn process(&self, event: SdkEvent) -> Option<SdkEvent>;
188}
189
190/// Event publisher that manages event listeners and middleware.
191///
192/// Events flow through three phases:
193/// 1. Internal listeners see raw events (SDK components like `ClientSyncListener`)
194/// 2. Middleware chain can transform or suppress events
195/// 3. External listeners see processed events (client event handlers)
196pub 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 see ALL events before middleware processing
202    internal_listeners: RwLock<BTreeMap<String, Box<dyn EventListener>>>,
203    /// Middleware chain that can transform/suppress events
204    middleware: RwLock<Vec<Box<dyn EventMiddleware>>>,
205    /// External listeners see events after middleware processing
206    external_listeners: RwLock<BTreeMap<String, Box<dyn EventListener>>>,
207    synced_event_buffer: Mutex<Option<InternalSyncedEvent>>,
208}
209
210/// Handlers registered here are owned by the `EventEmitter` for its whole
211/// lifetime, and the emitter is itself owned by the SDK object graph. To keep
212/// that graph droppable, a handler must not hold a strong reference back to
213/// the `EventEmitter` (directly or via a `BreezSdk` clone); the emitter is
214/// passed into `handle` instead.
215#[macros::async_trait]
216pub(crate) trait RuntimeEventHandler: Send + Sync {
217    async fn handle(&self, emitter: &EventEmitter, event: RuntimeEvent);
218}
219
220impl EventEmitter {
221    /// Create a new event emitter
222    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    /// Add an external listener to receive events
236    ///
237    /// # Arguments
238    ///
239    /// * `listener` - The listener to add
240    ///
241    /// # Returns
242    ///
243    /// A unique identifier for the listener, which can be used to remove it later
244    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    /// Remove an external listener by its ID
253    ///
254    /// # Arguments
255    ///
256    /// * `id` - The ID returned from `add_listener`
257    ///
258    /// # Returns
259    ///
260    /// `true` if the listener was found and removed, `false` otherwise
261    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    /// Remove all external listeners.
267    ///
268    /// Listeners are owned by the emitter, so a listener that references the
269    /// SDK pins the whole instance; dropping them here on disconnect makes
270    /// the instance releasable regardless of what listeners capture.
271    pub async fn clear_external_listeners(&self) {
272        let mut listeners = self.external_listeners.write().await;
273        listeners.clear();
274    }
275
276    /// Add an internal listener that sees all raw events before middleware processing.
277    ///
278    /// Used by SDK components (e.g., `ClientSyncListener`) that need to observe events
279    /// that middleware may suppress.
280    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    /// Remove an internal listener by its ID
289    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    /// Add middleware to the event processing chain.
295    ///
296    /// Middleware can transform or suppress events before they reach external listeners.
297    ///
298    /// Middleware is owned by the emitter for its whole lifetime, so it must
299    /// not hold a reference back to the emitter (directly or transitively),
300    /// or the SDK object graph can never be dropped. Code that needs to emit
301    /// must receive the emitter from its caller instead (see
302    /// `RuntimeEventHandler` and `TokenConverter::convert`).
303    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    /// Emit an event through the three-phase pipeline:
321    /// 1. Internal listeners see the raw event
322    /// 2. Middleware chain can transform or suppress
323    /// 3. External listeners see the processed event
324    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        // Phase 1: Internal listeners see raw event
332        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        // Phase 2: Middleware chain
344        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        // Phase 3: External listeners see processed event
361        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            // Nothing to emit
389            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            // The first synced event emitted should at least have the wallet synced.
398            // Subsequent events might have only partial syncs.
399            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        // Only emit zero real-time syncs on the first event.
418        if !is_first_event && !synced.any_non_zero() {
419            return;
420        }
421
422        // Emit the merged event
423        self.emit(&SdkEvent::Synced).await;
424    }
425
426    /// Notify that real-time sync has failed. If the first synced event is still
427    /// buffered and the wallet has already synced, release it immediately instead
428    /// of waiting for a remote pull that may never arrive.
429    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        // Create the listener with a shared reference to the atomic boolean
477        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        // Check if event was received using the shared reference
488        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        // Create shared atomic booleans to track event reception
496        let received1 = Arc::new(AtomicBool::new(false));
497        let received2 = Arc::new(AtomicBool::new(false));
498
499        // Create listeners with their own shared references
500        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        // Remove the first listener
512        assert!(emitter.remove_external_listener(&id1).await);
513
514        // Emit an event
515        let event = SdkEvent::Synced {};
516        emitter.emit(&event).await;
517
518        // The first listener should not receive the event
519        assert!(!received1.load(Ordering::Relaxed));
520
521        // The second listener should receive the event
522        assert!(received2.load(Ordering::Relaxed));
523
524        // Remove the second listener
525        assert!(emitter.remove_external_listener(&id2).await);
526
527        // Try to remove a non-existent listener
528        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        // Already cleared, so individual removal finds nothing
557        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        // Emit synced event without wallet sync - should NOT emit Synced
573        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        // Emit synced event with wallet sync - should emit Synced
586        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        // Emit synced event with storage
611        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        // Emit synced event with wallet sync - should emit Synced
624        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        // Emit synced event with wallet sync
650        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        // Emit synced event with storage - should emit Synced
663        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        // Wallet synced but rtsync hasn't failed yet — should NOT emit
688        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        // rtsync fails — should immediately release the buffered event
701        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        // rtsync fails before wallet syncs — nothing to release yet
718        emitter.notify_rtsync_failed().await;
719
720        assert!(!received.load(Ordering::Relaxed));
721
722        // Wallet syncs — should emit immediately (rtsync already marked failed)
723        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        // Emit multiple partial syncs without wallet sync
748        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        // Finally emit wallet sync - should emit Synced
797        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        // Emit synced event with wallet and other components - should emit Synced
822        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        // Emit empty synced event - should NOT emit Synced
847        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        // First sync with wallet - should emit
887        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        // Subsequent partial sync without wallet - should emit (buffer cleared after first wallet sync)
900        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        // Another partial sync - should emit
913        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        // storage_incoming with Some(0) - should NOT emit after first sync
950        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    // ── Helpers for middleware / internal listener tests ──
976
977    /// Listener that records all received events
978    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    /// Middleware that suppresses all Synced events
1002    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    /// Middleware that replaces `PaymentSucceeded` with `PaymentPending`
1015    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    /// Middleware that suppresses all events
1030    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    // ── Internal listener tests ──
1054
1055    #[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    // ── Middleware tests ──
1085
1086    #[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        // External should NOT see the suppressed event
1099        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        // Synced is suppressed, PaymentSucceeded passes through
1133        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        // First middleware transforms PaymentSucceeded → PaymentPending
1152        emitter
1153            .add_middleware(Box::new(DowngradePaymentMiddleware))
1154            .await;
1155        // Second middleware suppresses Synced (doesn't affect PaymentPending)
1156        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        // SuppressAll first — nothing should reach the next middleware or external
1178        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    // ── Three-phase flow tests ──
1195
1196    #[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        // Synced: internal sees it, middleware suppresses it, external doesn't
1209        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        // PaymentSucceeded: both see it (middleware passes it through)
1215        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        // Internal sees the original PaymentSucceeded
1245        assert_eq!(int_log.len(), 1);
1246        assert!(int_log[0].contains("PaymentSucceeded"));
1247
1248        // External sees the transformed PaymentPending
1249        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        // Should not panic
1258    }
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        // First sync with wallet - should emit
1287        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        // Empty sync after wallet sync - should NOT emit (all fields false)
1300        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); // Count should remain 1
1311
1312        // Another non-empty sync - should emit
1313        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); // Now count should be 2
1324    }
1325}