@@ -49,9 +49,18 @@ func defaultEndpoint() fwkdl.Endpoint {
4949
5050var (
5151 endpoint = defaultEndpoint ()
52- sources = []fwkdl.PollingDispatcher {& datasourcemocks.MetricsDataSource {}}
52+ sources = scheduleEveryTick ( []fwkdl.PollingDispatcher {& datasourcemocks.MetricsDataSource {}})
5353)
5454
55+ // scheduleEveryTick wraps each dispatcher with PeriodTicks=1 for collector tests.
56+ func scheduleEveryTick (dispatchers []fwkdl.PollingDispatcher ) []ScheduledDispatcher {
57+ out := make ([]ScheduledDispatcher , len (dispatchers ))
58+ for i , d := range dispatchers {
59+ out [i ] = ScheduledDispatcher {Dispatcher : d , PeriodTicks : 1 }
60+ }
61+ return out
62+ }
63+
5564// Mock PollingDispatchers for collector tests. Each tracks its own
5665// invocation count and (for dataSource) runs bound extractors with metric
5766// instrumentation. Same behavior the framework collector did before the
@@ -121,13 +130,14 @@ func TestCollectorStartInputs(t *testing.T) {
121130 tests := []struct {
122131 name string
123132 ctxCanceled bool
124- sources []fwkdl. PollingDispatcher
133+ sources []ScheduledDispatcher
125134 wantErr bool
126135 wantErrIs error
127136 }{
128137 {name : "valid sources, live ctx" , sources : sources },
129- {name : "empty sources" , sources : []fwkdl.PollingDispatcher {}, wantErr : true },
130- {name : "nil source" , sources : []fwkdl.PollingDispatcher {nil }, wantErr : true },
138+ {name : "empty sources" , sources : []ScheduledDispatcher {}, wantErr : true },
139+ {name : "nil source" , sources : []ScheduledDispatcher {{Dispatcher : nil , PeriodTicks : 1 }}, wantErr : true },
140+ {name : "periodTicks zero" , sources : []ScheduledDispatcher {{Dispatcher : & datasourcemocks.MetricsDataSource {}, PeriodTicks : 0 }}, wantErr : true },
131141 {name : "cancelled parent ctx" , ctxCanceled : true , sources : sources , wantErr : true , wantErrIs : context .Canceled },
132142 }
133143
@@ -182,7 +192,7 @@ func TestCollectorStop(t *testing.T) {
182192 setup : func (t * testing.T ) * Collector {
183193 c := NewCollector ()
184194 ticker := mocks .NewTicker ()
185- _ = c .Start (context .Background (), ticker , endpoint , []fwkdl. PollingDispatcher {})
195+ _ = c .Start (context .Background (), ticker , endpoint , []ScheduledDispatcher {})
186196 return c
187197 },
188198 },
@@ -213,7 +223,7 @@ func TestCollectorCollectsOnTicks(t *testing.T) {
213223 c := NewCollector ()
214224 ticker := mocks .NewTicker ()
215225
216- require .NoError (t , c .Start (context .Background (), ticker , endpoint , []fwkdl.PollingDispatcher {source }))
226+ require .NoError (t , c .Start (context .Background (), ticker , endpoint , scheduleEveryTick ( []fwkdl.PollingDispatcher {source }) ))
217227 defer c .Stop ()
218228
219229 ticker .Tick ()
@@ -224,6 +234,30 @@ func TestCollectorCollectsOnTicks(t *testing.T) {
224234 }, 1 * time .Second , 2 * time .Millisecond , "expected 2 collections" )
225235}
226236
237+ // TestCollectorRespectsPeriodTicks confirms slower sources skip ticks.
238+ func TestCollectorRespectsPeriodTicks (t * testing.T ) {
239+ fast := & errSource {kind : "fast" }
240+ slow := & errSource {kind : "slow" }
241+
242+ c := NewCollector ()
243+ ticker := mocks .NewTicker ()
244+ require .NoError (t , c .Start (context .Background (), ticker , endpoint , []ScheduledDispatcher {
245+ {Dispatcher : fast , PeriodTicks : 1 },
246+ {Dispatcher : slow , PeriodTicks : 4 },
247+ }))
248+ defer c .Stop ()
249+
250+ for i := 0 ; i < 8 ; i ++ {
251+ ticker .Tick ()
252+ }
253+
254+ require .Eventually (t , func () bool {
255+ return atomic .LoadInt64 (& fast .CallCount ) == 8 && atomic .LoadInt64 (& slow .CallCount ) == 2
256+ }, 1 * time .Second , 2 * time .Millisecond ,
257+ "fast=%d slow=%d want fast=8 slow=2" ,
258+ atomic .LoadInt64 (& fast .CallCount ), atomic .LoadInt64 (& slow .CallCount ))
259+ }
260+
227261// TestCollectorErrorMetrics confirms Poll/Extract errors increment per-event
228262// counters (no transition dedup) and successes do not.
229263func TestCollectorErrorMetrics (t * testing.T ) {
@@ -285,7 +319,7 @@ func TestCollectorErrorMetrics(t *testing.T) {
285319
286320 c := NewCollector ()
287321 ticker := mocks .NewTicker ()
288- require .NoError (t , c .Start (context .Background (), ticker , endpoint , []fwkdl.PollingDispatcher {src }))
322+ require .NoError (t , c .Start (context .Background (), ticker , endpoint , scheduleEveryTick ( []fwkdl.PollingDispatcher {src }) ))
289323 defer c .Stop ()
290324
291325 for i := 0 ; i < tt .ticks ; i ++ {
@@ -324,7 +358,7 @@ func TestCollectorRapidStartStopRaceFree(t *testing.T) {
324358 for i := 0 ; i < 100 ; i ++ {
325359 c := NewCollector ()
326360 ticker := mocks .NewTicker ()
327- require .NoError (t , c .Start (context .Background (), ticker , endpoint , []fwkdl.PollingDispatcher {src }))
361+ require .NoError (t , c .Start (context .Background (), ticker , endpoint , scheduleEveryTick ( []fwkdl.PollingDispatcher {src }) ))
328362 ticker .Tick ()
329363 c .Stop ()
330364 }
0 commit comments