@@ -25,6 +25,9 @@ use core::task::{Context, Poll, Waker};
2525use aimdb_core:: buffer:: { Buffer , BufferCfg , BufferReader , DynBuffer } ;
2626use aimdb_core:: DbError ;
2727
28+ #[ cfg( feature = "observability" ) ]
29+ use aimdb_core:: buffer:: { BufferCounters , BufferMetrics , BufferMetricsSnapshot } ;
30+
2831// ============================================================================
2932// Buffer
3033// ============================================================================
@@ -37,6 +40,8 @@ use aimdb_core::DbError;
3740/// construction time.
3841pub struct WasmBuffer < T > {
3942 inner : Rc < RefCell < WasmBufferInner < T > > > ,
43+ #[ cfg( feature = "observability" ) ]
44+ metrics : Rc < BufferCounters > ,
4045}
4146
4247// SAFETY: wasm32 is single-threaded — Rc<RefCell<…>> cannot be accessed concurrently
@@ -81,6 +86,12 @@ impl<T: Clone + Send + 'static> Buffer<T> for WasmBuffer<T> {
8186 type Reader = WasmBufferReader < T > ;
8287
8388 fn new ( cfg : & BufferCfg ) -> Self {
89+ #[ cfg( feature = "observability" ) ]
90+ let capacity = match cfg {
91+ BufferCfg :: SpmcRing { capacity } => * capacity,
92+ BufferCfg :: SingleLatest | BufferCfg :: Mailbox => 1 ,
93+ } ;
94+
8495 let inner = match cfg {
8596 BufferCfg :: SpmcRing { capacity } => WasmBufferInner :: SpmcRing {
8697 ring : VecDeque :: with_capacity ( * capacity) ,
@@ -101,10 +112,15 @@ impl<T: Clone + Send + 'static> Buffer<T> for WasmBuffer<T> {
101112
102113 WasmBuffer {
103114 inner : Rc :: new ( RefCell :: new ( inner) ) ,
115+ #[ cfg( feature = "observability" ) ]
116+ metrics : Rc :: new ( BufferCounters :: new ( capacity) ) ,
104117 }
105118 }
106119
107120 fn push ( & self , value : T ) {
121+ #[ cfg( feature = "observability" ) ]
122+ self . metrics . increment_produced ( ) ;
123+
108124 let mut inner = self . inner . borrow_mut ( ) ;
109125 match & mut * inner {
110126 WasmBufferInner :: SpmcRing {
@@ -163,6 +179,8 @@ impl<T: Clone + Send + 'static> Buffer<T> for WasmBuffer<T> {
163179 WasmBufferReader {
164180 buffer : Rc :: clone ( & self . inner ) ,
165181 state,
182+ #[ cfg( feature = "observability" ) ]
183+ metrics : Rc :: clone ( & self . metrics ) ,
166184 }
167185 }
168186}
@@ -192,6 +210,34 @@ impl<T: Clone + Send + 'static> DynBuffer<T> for WasmBuffer<T> {
192210 WasmBufferInner :: SpmcRing { .. } => None ,
193211 }
194212 }
213+
214+ #[ cfg( feature = "observability" ) ]
215+ fn metrics_snapshot ( & self ) -> Option < BufferMetricsSnapshot > {
216+ Some ( <Self as BufferMetrics >:: metrics ( self ) )
217+ }
218+
219+ #[ cfg( feature = "observability" ) ]
220+ fn reset_metrics ( & self ) {
221+ <Self as BufferMetrics >:: reset_metrics ( self ) ;
222+ }
223+ }
224+
225+ #[ cfg( feature = "observability" ) ]
226+ impl < T : Clone + Send + ' static > BufferMetrics for WasmBuffer < T > {
227+ fn metrics ( & self ) -> BufferMetricsSnapshot {
228+ let current_occupancy = match & * self . inner . borrow ( ) {
229+ WasmBufferInner :: SpmcRing { ring, .. } => ring. len ( ) ,
230+ WasmBufferInner :: SingleLatest { value, .. } => usize:: from ( value. is_some ( ) ) ,
231+ WasmBufferInner :: Mailbox { slot, .. } => usize:: from ( slot. is_some ( ) ) ,
232+ } ;
233+
234+ self . metrics
235+ . snapshot ( ( current_occupancy, self . metrics . capacity ( ) ) )
236+ }
237+
238+ fn reset_metrics ( & self ) {
239+ self . metrics . reset ( ) ;
240+ }
195241}
196242
197243// ============================================================================
@@ -205,6 +251,8 @@ impl<T: Clone + Send + 'static> DynBuffer<T> for WasmBuffer<T> {
205251pub struct WasmBufferReader < T > {
206252 buffer : Rc < RefCell < WasmBufferInner < T > > > ,
207253 state : ReaderState ,
254+ #[ cfg( feature = "observability" ) ]
255+ metrics : Rc < BufferCounters > ,
208256}
209257
210258// SAFETY: wasm32 is single-threaded — no concurrent access possible
@@ -276,6 +324,8 @@ impl<T: Clone + Send + 'static> BufferReader<T> for WasmBufferReader<T> {
276324 // Reader fell behind — skip to oldest available
277325 let lag_count = oldest_seq - * read_seq;
278326 * read_seq = oldest_seq;
327+ #[ cfg( feature = "observability" ) ]
328+ self . metrics . add_dropped ( lag_count) ;
279329 return Err ( DbError :: BufferLagged {
280330 lag_count,
281331 buffer_name : alloc:: string:: String :: from ( "wasm ring" ) ,
@@ -285,6 +335,8 @@ impl<T: Clone + Send + 'static> BufferReader<T> for WasmBufferReader<T> {
285335 let offset = ( * read_seq - oldest_seq) as usize ;
286336 let value = ring[ offset] . clone ( ) ;
287337 * read_seq += 1 ;
338+ #[ cfg( feature = "observability" ) ]
339+ self . metrics . increment_consumed ( ) ;
288340 Ok ( value)
289341 }
290342 (
@@ -297,6 +349,8 @@ impl<T: Clone + Send + 'static> BufferReader<T> for WasmBufferReader<T> {
297349 match value {
298350 Some ( v) => {
299351 * last_seen_version = * version;
352+ #[ cfg( feature = "observability" ) ]
353+ self . metrics . increment_consumed ( ) ;
300354 Ok ( v. clone ( ) )
301355 }
302356 // Unreachable given the invariants (value is `Some` iff
@@ -306,9 +360,14 @@ impl<T: Clone + Send + 'static> BufferReader<T> for WasmBufferReader<T> {
306360 None => Err ( DbError :: BufferEmpty ) ,
307361 }
308362 }
309- ( WasmBufferInner :: Mailbox { slot, .. } , ReaderState :: Mailbox ) => {
310- slot. take ( ) . ok_or ( DbError :: BufferEmpty )
311- }
363+ ( WasmBufferInner :: Mailbox { slot, .. } , ReaderState :: Mailbox ) => match slot. take ( ) {
364+ Some ( value) => {
365+ #[ cfg( feature = "observability" ) ]
366+ self . metrics . increment_consumed ( ) ;
367+ Ok ( value)
368+ }
369+ None => Err ( DbError :: BufferEmpty ) ,
370+ } ,
312371 _ => unreachable ! ( "reader state mismatch" ) ,
313372 }
314373 }
@@ -587,4 +646,101 @@ mod tests {
587646 assert_eq ! ( buffer. peek( ) , None ) ;
588647 }
589648 }
649+
650+ #[ cfg( feature = "observability" ) ]
651+ mod metrics_tests {
652+ use super :: super :: * ;
653+ use aimdb_core:: buffer:: { BufferMetrics , DynBuffer } ;
654+
655+ #[ test]
656+ fn metrics_snapshot_is_available_for_all_buffer_types ( ) {
657+ let ring = WasmBuffer :: < i32 > :: new ( & BufferCfg :: SpmcRing { capacity : 4 } ) ;
658+ let latest = WasmBuffer :: < i32 > :: new ( & BufferCfg :: SingleLatest ) ;
659+ let mailbox = WasmBuffer :: < i32 > :: new ( & BufferCfg :: Mailbox ) ;
660+
661+ assert ! ( DynBuffer :: metrics_snapshot( & ring) . is_some( ) ) ;
662+ assert ! ( DynBuffer :: metrics_snapshot( & latest) . is_some( ) ) ;
663+ assert ! ( DynBuffer :: metrics_snapshot( & mailbox) . is_some( ) ) ;
664+ }
665+
666+ #[ test]
667+ fn spmc_ring_counts_produced_consumed_and_reader_lag ( ) {
668+ let buffer = WasmBuffer :: < i32 > :: new ( & BufferCfg :: SpmcRing { capacity : 3 } ) ;
669+ let mut reader = buffer. subscribe ( ) ;
670+
671+ for value in 0 ..10 {
672+ Buffer :: push ( & buffer, value) ;
673+ }
674+
675+ assert ! ( matches!(
676+ reader. try_recv( ) ,
677+ Err ( DbError :: BufferLagged { lag_count: 7 , .. } )
678+ ) ) ;
679+
680+ let metrics = buffer. metrics ( ) ;
681+ assert_eq ! ( metrics. produced_count, 10 ) ;
682+ assert_eq ! ( metrics. consumed_count, 0 ) ;
683+ assert_eq ! ( metrics. dropped_count, 7 ) ;
684+ assert_eq ! ( metrics. occupancy, ( 3 , 3 ) ) ;
685+
686+ assert_eq ! ( reader. try_recv( ) . unwrap( ) , 7 ) ;
687+ let metrics = buffer. metrics ( ) ;
688+ assert_eq ! ( metrics. consumed_count, 1 ) ;
689+ assert_eq ! ( metrics. dropped_count, 7 ) ;
690+ }
691+
692+ #[ test]
693+ fn single_latest_overwrites_are_not_drops_and_occupancy_stays_present ( ) {
694+ let buffer = WasmBuffer :: < i32 > :: new ( & BufferCfg :: SingleLatest ) ;
695+ let mut reader = buffer. subscribe ( ) ;
696+
697+ Buffer :: push ( & buffer, 1 ) ;
698+ Buffer :: push ( & buffer, 2 ) ;
699+ Buffer :: push ( & buffer, 3 ) ;
700+
701+ assert_eq ! ( reader. try_recv( ) . unwrap( ) , 3 ) ;
702+ let metrics = buffer. metrics ( ) ;
703+ assert_eq ! ( metrics. produced_count, 3 ) ;
704+ assert_eq ! ( metrics. consumed_count, 1 ) ;
705+ assert_eq ! ( metrics. dropped_count, 0 ) ;
706+ assert_eq ! ( metrics. occupancy, ( 1 , 1 ) ) ;
707+ }
708+
709+ #[ test]
710+ fn mailbox_overwrites_are_not_drops_and_receive_clears_occupancy ( ) {
711+ let buffer = WasmBuffer :: < i32 > :: new ( & BufferCfg :: Mailbox ) ;
712+ let mut reader = buffer. subscribe ( ) ;
713+
714+ Buffer :: push ( & buffer, 1 ) ;
715+ Buffer :: push ( & buffer, 2 ) ;
716+
717+ let metrics = buffer. metrics ( ) ;
718+ assert_eq ! ( metrics. produced_count, 2 ) ;
719+ assert_eq ! ( metrics. consumed_count, 0 ) ;
720+ assert_eq ! ( metrics. dropped_count, 0 ) ;
721+ assert_eq ! ( metrics. occupancy, ( 1 , 1 ) ) ;
722+
723+ assert_eq ! ( reader. try_recv( ) . unwrap( ) , 2 ) ;
724+ let metrics = buffer. metrics ( ) ;
725+ assert_eq ! ( metrics. consumed_count, 1 ) ;
726+ assert_eq ! ( metrics. dropped_count, 0 ) ;
727+ assert_eq ! ( metrics. occupancy, ( 0 , 1 ) ) ;
728+ }
729+
730+ #[ test]
731+ fn reset_metrics_preserves_live_occupancy_and_capacity ( ) {
732+ let buffer = WasmBuffer :: < i32 > :: new ( & BufferCfg :: SingleLatest ) ;
733+ let mut reader = buffer. subscribe ( ) ;
734+
735+ Buffer :: push ( & buffer, 42 ) ;
736+ assert_eq ! ( reader. try_recv( ) . unwrap( ) , 42 ) ;
737+ DynBuffer :: reset_metrics ( & buffer) ;
738+
739+ let metrics = buffer. metrics ( ) ;
740+ assert_eq ! ( metrics. produced_count, 0 ) ;
741+ assert_eq ! ( metrics. consumed_count, 0 ) ;
742+ assert_eq ! ( metrics. dropped_count, 0 ) ;
743+ assert_eq ! ( metrics. occupancy, ( 1 , 1 ) ) ;
744+ }
745+ }
590746}
0 commit comments