@@ -137,9 +137,14 @@ impl<T: AggregateExport> AggregateExportRuntime<T> {
137137 . inner
138138 . entry ( record. tag_map . clone ( ) )
139139 . and_modify ( |v| {
140- v. time = self . store_time ;
141- v. sum += record. value ;
142- v. diff = record. value ;
140+ if v. time < self . store_time {
141+ v. time = self . store_time ;
142+ v. sum += record. value ;
143+ v. diff = record. value ;
144+ } else {
145+ v. sum += record. value ;
146+ v. diff += record. value ;
147+ }
143148 } )
144149 . or_insert ( CounterStoreValue {
145150 time : self . store_time ,
@@ -160,3 +165,100 @@ impl<T: AggregateExport> AggregateExportRuntime<T> {
160165 }
161166 }
162167}
168+
169+ #[ cfg( test) ]
170+ mod tests {
171+ use super :: * ;
172+ use std:: sync:: Arc ;
173+ use tokio:: sync:: mpsc;
174+ use vey_types:: metrics:: MetricTagMap ;
175+
176+ struct TestExporter {
177+ counters : AHashMap < MetricName , AHashMap < Arc < MetricTagMap > , ( MetricValue , MetricValue ) > > ,
178+ }
179+
180+ impl AggregateExport for TestExporter {
181+ fn emit_interval ( & self ) -> Duration {
182+ Duration :: from_secs ( 10 )
183+ }
184+
185+ fn emit_gauge (
186+ & mut self ,
187+ _name : & MetricName ,
188+ _values : & AHashMap < Arc < MetricTagMap > , GaugeStoreValue > ,
189+ ) {
190+ }
191+
192+ fn emit_counter (
193+ & mut self ,
194+ name : & MetricName ,
195+ values : & AHashMap < Arc < MetricTagMap > , CounterStoreValue > ,
196+ ) {
197+ let map = self . counters . entry ( name. clone ( ) ) . or_default ( ) ;
198+ for ( tags, v) in values {
199+ map. insert ( tags. clone ( ) , ( v. sum , v. diff ) ) ;
200+ }
201+ }
202+ }
203+
204+ #[ test]
205+ fn test_counter_diff_accumulation_and_reset ( ) {
206+ let ( _tx, rx) = mpsc:: unbounded_channel ( ) ;
207+ let exporter = TestExporter {
208+ counters : AHashMap :: default ( ) ,
209+ } ;
210+ let mut runtime = AggregateExportRuntime :: new ( exporter, rx) ;
211+
212+ let name = Arc :: new ( MetricName :: parse ( "test.counter" ) . unwrap ( ) ) ;
213+ let tag_map = Arc :: new ( MetricTagMap :: default ( ) ) ;
214+
215+ // Interval 1 - First record
216+ runtime. add_record ( MetricRecord {
217+ name : name. clone ( ) ,
218+ tag_map : tag_map. clone ( ) ,
219+ r#type : MetricType :: Counter ,
220+ value : MetricValue :: Signed ( 10 ) ,
221+ } ) ;
222+
223+ // Interval 1 - Second record (same interval)
224+ runtime. add_record ( MetricRecord {
225+ name : name. clone ( ) ,
226+ tag_map : tag_map. clone ( ) ,
227+ r#type : MetricType :: Counter ,
228+ value : MetricValue :: Signed ( 5 ) ,
229+ } ) ;
230+
231+ // Verify state before retain
232+ let counter_entry = & runtime. counter . get ( & name) . unwrap ( ) . inner [ & tag_map] ;
233+ assert_eq ! ( counter_entry. sum, MetricValue :: Signed ( 15 ) ) ;
234+ assert_eq ! ( counter_entry. diff, MetricValue :: Signed ( 15 ) ) ;
235+
236+ // Simulate tick / interval transition
237+ runtime. retain ( ) ;
238+
239+ // Interval 2 - First record in new interval
240+ runtime. add_record ( MetricRecord {
241+ name : name. clone ( ) ,
242+ tag_map : tag_map. clone ( ) ,
243+ r#type : MetricType :: Counter ,
244+ value : MetricValue :: Signed ( 3 ) ,
245+ } ) ;
246+
247+ // Verify diff reset for new interval, sum accumulated
248+ let counter_entry = & runtime. counter . get ( & name) . unwrap ( ) . inner [ & tag_map] ;
249+ assert_eq ! ( counter_entry. sum, MetricValue :: Signed ( 18 ) ) ;
250+ assert_eq ! ( counter_entry. diff, MetricValue :: Signed ( 3 ) ) ;
251+
252+ // Interval 2 - Second record in new interval
253+ runtime. add_record ( MetricRecord {
254+ name : name. clone ( ) ,
255+ tag_map : tag_map. clone ( ) ,
256+ r#type : MetricType :: Counter ,
257+ value : MetricValue :: Signed ( 7 ) ,
258+ } ) ;
259+
260+ let counter_entry = & runtime. counter . get ( & name) . unwrap ( ) . inner [ & tag_map] ;
261+ assert_eq ! ( counter_entry. sum, MetricValue :: Signed ( 25 ) ) ;
262+ assert_eq ! ( counter_entry. diff, MetricValue :: Signed ( 10 ) ) ;
263+ }
264+ }
0 commit comments