Skip to content

Commit 50fc3cb

Browse files
committed
feat(stovepipe): tag controller metrics by queue
Scope metrics from decoded pipeline messages to their logical queue so operators can filter controller health without splitting controller instances.
1 parent e47fee7 commit 50fc3cb

8 files changed

Lines changed: 68 additions & 4 deletions

File tree

stovepipe/controller/build/build.go

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -88,6 +88,7 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) er
8888
// Non-retryable: a malformed message will never succeed regardless of retries.
8989
return fmt.Errorf("failed to deserialize build request: %w", err)
9090
}
91+
c = c.forQueue(br.GetQueueName())
9192

9293
store, err := c.stores.For(storage.Config{QueueName: br.GetQueueName()})
9394
if err != nil {
@@ -192,3 +193,12 @@ func (c *Controller) TopicKey() consumer.TopicKey {
192193
func (c *Controller) ConsumerGroup() string {
193194
return c.consumerGroup
194195
}
196+
197+
func (c *Controller) forQueue(queue string) *Controller {
198+
if queue == "" {
199+
return c
200+
}
201+
scoped := *c
202+
scoped.metricsScope = c.metricsScope.Tagged(map[string]string{"queue": queue})
203+
return &scoped
204+
}

stovepipe/controller/buildsignal/buildsignal.go

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -108,6 +108,7 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) er
108108
// Non-retryable: a malformed message will never succeed regardless of retries.
109109
return fmt.Errorf("failed to deserialize build signal: %w", err)
110110
}
111+
c = c.forQueue(sig.GetQueueName())
111112

112113
store, err := c.stores.For(storage.Config{QueueName: sig.GetQueueName()})
113114
if err != nil {
@@ -376,3 +377,12 @@ func (c *Controller) TopicKey() consumer.TopicKey {
376377
func (c *Controller) ConsumerGroup() string {
377378
return c.consumerGroup
378379
}
380+
381+
func (c *Controller) forQueue(queue string) *Controller {
382+
if queue == "" {
383+
return c
384+
}
385+
scoped := *c
386+
scoped.metricsScope = c.metricsScope.Tagged(map[string]string{"queue": queue})
387+
return &scoped
388+
}

stovepipe/controller/dlq/buildsignal.go

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -95,6 +95,7 @@ func (c *BuildSignalController) Process(ctx context.Context, delivery consumer.D
9595
// without saying so.
9696
return fmt.Errorf("failed to decode dlq payload: %w", err)
9797
}
98+
c = c.forQueue(sig.GetQueueName())
9899
if sig.Id == "" {
99100
metrics.NamedCounter(c.metricsScope, _buildSignalOpName, "empty_id_errors", 1)
100101
return fmt.Errorf("dlq payload decoded to empty build id")
@@ -165,3 +166,12 @@ func (c *BuildSignalController) TopicKey() consumer.TopicKey {
165166
func (c *BuildSignalController) ConsumerGroup() string {
166167
return c.consumerGroup
167168
}
169+
170+
func (c *BuildSignalController) forQueue(queue string) *BuildSignalController {
171+
if queue == "" {
172+
return c
173+
}
174+
scoped := *c
175+
scoped.metricsScope = c.metricsScope.Tagged(map[string]string{"queue": queue})
176+
return &scoped
177+
}

stovepipe/controller/dlq/request.go

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -83,6 +83,7 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) er
8383
// non-terminal.
8484
return fmt.Errorf("failed to decode dlq payload: %w", err)
8585
}
86+
c = c.forQueue(pr.GetQueueName())
8687
if pr.Id == "" {
8788
metrics.NamedCounter(c.metricsScope, _opName, "empty_id_errors", 1)
8889
return fmt.Errorf("dlq payload decoded to empty request id")
@@ -126,3 +127,12 @@ func (c *Controller) TopicKey() consumer.TopicKey {
126127
func (c *Controller) ConsumerGroup() string {
127128
return c.consumerGroup
128129
}
130+
131+
func (c *Controller) forQueue(queue string) *Controller {
132+
if queue == "" {
133+
return c
134+
}
135+
scoped := *c
136+
scoped.metricsScope = c.metricsScope.Tagged(map[string]string{"queue": queue})
137+
return &scoped
138+
}

stovepipe/controller/ingest.go

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -94,7 +94,11 @@ func NewIngestController(
9494
func (c *IngestController) Ingest(ctx context.Context, req entity.IngestRequest) (result entity.IngestResult, retErr error) {
9595
const opName = "ingest"
9696

97-
op := metrics.Begin(c.metricsScope, opName, metrics.LongLatencyBuckets)
97+
metricsScope := c.metricsScope
98+
if req.Queue != "" {
99+
metricsScope = metricsScope.Tagged(map[string]string{"queue": req.Queue})
100+
}
101+
op := metrics.Begin(metricsScope, opName, metrics.LongLatencyBuckets)
98102
defer func() { op.Complete(retErr) }()
99103

100104
if req.Queue == "" {

stovepipe/controller/process/process.go

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -91,6 +91,7 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) er
9191
// Non-retryable: a malformed message will never succeed regardless of retries.
9292
return fmt.Errorf("failed to deserialize process request: %w", err)
9393
}
94+
c = c.forQueue(pr.GetQueueName())
9495

9596
store, err := c.stores.For(storage.Config{QueueName: pr.GetQueueName()})
9697
if err != nil {
@@ -489,3 +490,12 @@ func (c *Controller) TopicKey() consumer.TopicKey {
489490
func (c *Controller) ConsumerGroup() string {
490491
return c.consumerGroup
491492
}
493+
494+
func (c *Controller) forQueue(queue string) *Controller {
495+
if queue == "" {
496+
return c
497+
}
498+
scoped := *c
499+
scoped.metricsScope = c.metricsScope.Tagged(map[string]string{"queue": queue})
500+
return &scoped
501+
}

stovepipe/controller/process/process_test.go

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -109,7 +109,7 @@ func delivery(t *testing.T, ctrl *gomock.Controller, payload []byte) *consumermo
109109

110110
func processPayload(t *testing.T, id string) []byte {
111111
t.Helper()
112-
b, err := stovepipemq.Marshal(&stovepipemq.ProcessRequest{Id: id})
112+
b, err := stovepipemq.Marshal(&stovepipemq.ProcessRequest{Id: id, QueueName: testQueue})
113113
require.NoError(t, err)
114114
return b
115115
}
@@ -313,7 +313,7 @@ func TestProcessEmitsAdmittedStrategyMetric(t *testing.T) {
313313

314314
require.NoError(t, c.Process(context.Background(), delivery(t, ctrl, processPayload(t, testID))))
315315

316-
counter, ok := scope.Snapshot().Counters()["test.process_controller.process.admitted+strategy=full"]
316+
counter, ok := scope.Snapshot().Counters()["test.process_controller.process.admitted+queue=monorepo/main,strategy=full"]
317317
require.True(t, ok)
318318
assert.Equal(t, int64(1), counter.Value())
319319
}
@@ -338,7 +338,7 @@ func TestProcessEmitsSourceControlResolutionMetric(t *testing.T) {
338338

339339
require.Error(t, c.Process(context.Background(), delivery(t, ctrl, processPayload(t, testID))))
340340

341-
counter, ok := scope.Snapshot().Counters()["test.process_controller.process.source_control_errors+stage=resolve"]
341+
counter, ok := scope.Snapshot().Counters()["test.process_controller.process.source_control_errors+queue=monorepo/main,stage=resolve"]
342342
require.True(t, ok)
343343
assert.Equal(t, int64(1), counter.Value())
344344
}

stovepipe/controller/record/record.go

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -99,6 +99,7 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) er
9999
// Non-retryable: a malformed message will never succeed regardless of retries.
100100
return fmt.Errorf("failed to deserialize record: %w", err)
101101
}
102+
c = c.forQueue(rec.GetQueueName())
102103

103104
store, err := c.stores.For(storage.Config{QueueName: rec.GetQueueName()})
104105
if err != nil {
@@ -489,3 +490,12 @@ func (c *Controller) TopicKey() consumer.TopicKey {
489490
func (c *Controller) ConsumerGroup() string {
490491
return c.consumerGroup
491492
}
493+
494+
func (c *Controller) forQueue(queue string) *Controller {
495+
if queue == "" {
496+
return c
497+
}
498+
scoped := *c
499+
scoped.metricsScope = c.metricsScope.Tagged(map[string]string{"queue": queue})
500+
return &scoped
501+
}

0 commit comments

Comments
 (0)