Skip to content

Commit b9ec32a

Browse files
committed
feat(stovepipe): tag controller metrics by queue
Summary: Attribute Stovepipe controller metrics to the logical queue carried with each message. Intent: - Make primary and DLQ controller metrics filterable by queue. - Keep message-derived metric context request-local and narrowly scoped. Changes: - Publish queue_name in Message.Metadata at each Stovepipe message handoff. - Copy only queue_name from delivered metadata into context. - Add metrics.TagsFromContext to append allowlisted context tags and migrate affected metrics. - Preserve payload queue fields as the authority for storage and validation, and omit the tag for older messages without metadata. Test Plan: - ./tool/bazel test //platform/base/messagequeue:go_default_test //platform/consumer:go_default_test //platform/metrics:go_default_test //stovepipe/controller:go_default_test //stovepipe/controller/build:go_default_test //stovepipe/controller/buildsignal:go_default_test //stovepipe/controller/dlq:go_default_test //stovepipe/controller/process:go_default_test //stovepipe/controller/record:go_default_test Revert Plan: - Revert this commit to restore the prior untagged controller metrics. --- <sub>Generated by the 🪄 [pr-create](https://sg.uberinternal.com/code.uber.internal/uber-code/devexp-agent-marketplace/-/blob/claude-code/plugins/dev/uber-dev/skills/pr-create/SKILL.md) skill in devexp-agent-marketplace</sub>
1 parent c93310d commit b9ec32a

28 files changed

Lines changed: 445 additions & 187 deletions

platform/base/messagequeue/BUILD.bazel

Lines changed: 8 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -2,14 +2,20 @@ load("@rules_go//go:def.bzl", "go_library", "go_test")
22

33
go_library(
44
name = "go_default_library",
5-
srcs = ["message.go"],
5+
srcs = [
6+
"context.go",
7+
"message.go",
8+
],
69
importpath = "github.com/uber/submitqueue/platform/base/messagequeue",
710
visibility = ["//visibility:public"],
811
)
912

1013
go_test(
1114
name = "go_default_test",
12-
srcs = ["message_test.go"],
15+
srcs = [
16+
"context_test.go",
17+
"message_test.go",
18+
],
1319
embed = [":go_default_library"],
1420
deps = ["@com_github_stretchr_testify//assert:go_default_library"],
1521
)
Lines changed: 38 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,38 @@
1+
// Copyright (c) 2025 Uber Technologies, Inc.
2+
//
3+
// Licensed under the Apache License, Version 2.0 (the "License");
4+
// you may not use this file except in compliance with the License.
5+
// You may obtain a copy of the License at
6+
//
7+
// http://www.apache.org/licenses/LICENSE-2.0
8+
//
9+
// Unless required by applicable law or agreed to in writing, software
10+
// distributed under the License is distributed on an "AS IS" BASIS,
11+
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12+
// See the License for the specific language governing permissions and
13+
// limitations under the License.
14+
15+
package messagequeue
16+
17+
import (
18+
"context"
19+
)
20+
21+
// MetadataKeyQueueName carries the logical queue name independently of the
22+
// transport partition key. Producers set it on Message.Metadata so consumers
23+
// can attribute work before decoding the payload.
24+
const MetadataKeyQueueName = "queue_name"
25+
26+
type queueNameContextKey struct{}
27+
28+
// WithQueueName returns a child context containing the logical queue name of
29+
// the delivered message.
30+
func WithQueueName(ctx context.Context, queueName string) context.Context {
31+
return context.WithValue(ctx, queueNameContextKey{}, queueName)
32+
}
33+
34+
// QueueName returns the delivered message's logical queue name from ctx.
35+
func QueueName(ctx context.Context) (string, bool) {
36+
queueName, ok := ctx.Value(queueNameContextKey{}).(string)
37+
return queueName, ok
38+
}
Lines changed: 33 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,33 @@
1+
// Copyright (c) 2025 Uber Technologies, Inc.
2+
//
3+
// Licensed under the Apache License, Version 2.0 (the "License");
4+
// you may not use this file except in compliance with the License.
5+
// You may obtain a copy of the License at
6+
//
7+
// http://www.apache.org/licenses/LICENSE-2.0
8+
//
9+
// Unless required by applicable law or agreed to in writing, software
10+
// distributed under the License is distributed on an "AS IS" BASIS,
11+
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12+
// See the License for the specific language governing permissions and
13+
// limitations under the License.
14+
15+
package messagequeue
16+
17+
import (
18+
"context"
19+
"testing"
20+
21+
"github.com/stretchr/testify/assert"
22+
)
23+
24+
func TestQueueNameContext(t *testing.T) {
25+
ctx := WithQueueName(context.Background(), "monorepo/main")
26+
27+
queueName, ok := QueueName(ctx)
28+
assert.True(t, ok)
29+
assert.Equal(t, "monorepo/main", queueName)
30+
31+
_, ok = QueueName(context.Background())
32+
assert.False(t, ok)
33+
}

platform/consumer/consumer.go

Lines changed: 15 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,7 @@ import (
2222
"time"
2323

2424
"github.com/uber-go/tally"
25+
entityqueue "github.com/uber/submitqueue/platform/base/messagequeue"
2526
"github.com/uber/submitqueue/platform/errs"
2627
"github.com/uber/submitqueue/platform/extension/consumergate"
2728
extqueue "github.com/uber/submitqueue/platform/extension/messagequeue"
@@ -371,6 +372,9 @@ func (m *consumer) processPartition(ctx context.Context, controller Controller,
371372
func (m *consumer) processDelivery(ctx context.Context, controller Controller, delivery extqueue.Delivery, controllerScope tally.Scope) {
372373
const opName = "process"
373374

375+
msg := delivery.Message()
376+
ctx = entityqueue.WithQueueName(ctx, msg.Metadata[entityqueue.MetadataKeyQueueName])
377+
374378
// Consumer gate: a delivery whose gate is closed is recorded as parked and
375379
// postponed (barrier + re-check on redelivery); a false return also covers
376380
// shutdown-while-checking, where the delivery is left in flight so its
@@ -380,7 +384,6 @@ func (m *consumer) processDelivery(ctx context.Context, controller Controller, d
380384
return
381385
}
382386

383-
msg := delivery.Message()
384387
topicKey := controller.TopicKey()
385388

386389
m.logger.Debugw("processing delivery",
@@ -396,7 +399,7 @@ func (m *consumer) processDelivery(ctx context.Context, controller Controller, d
396399

397400
// Call controller with wrapped delivery
398401
start := time.Now()
399-
op := metrics.Begin(controllerScope, opName, metrics.LongLatencyBuckets)
402+
op := metrics.Begin(controllerScope, opName, metrics.LongLatencyBuckets, metrics.TagsFromContext(ctx)...)
400403
err := controller.Process(ctx, wrapped)
401404

402405
elapsed := time.Since(start)
@@ -419,7 +422,7 @@ func (m *consumer) processDelivery(ctx context.Context, controller Controller, d
419422
// A failure outcome wins over a recorded hold — a hold is only honored
420423
// on success, so retry accounting and dead-lettering stay meaningful.
421424
if wrapped.held {
422-
metrics.NamedCounter(controllerScope, opName, "hold_ignored", 1)
425+
metrics.NamedCounter(controllerScope, opName, "hold_ignored", 1, metrics.TagsFromContext(ctx)...)
423426
m.logger.Warnw("hold recorded but controller returned error, failure outcome wins",
424427
"controller", controller.Name(),
425428
"topic_key", topicKey,
@@ -451,7 +454,7 @@ func (m *consumer) processDelivery(ctx context.Context, controller Controller, d
451454
)
452455

453456
// Reject moves to DLQ (or acks if DLQ disabled)
454-
rejectOp := metrics.Begin(controllerScope, "reject", metrics.StorageLatencyBuckets)
457+
rejectOp := metrics.Begin(controllerScope, "reject", metrics.StorageLatencyBuckets, metrics.TagsFromContext(ctx)...)
455458
rejectErr := delivery.Reject(ctx, controllerFailure)
456459
rejectOp.Complete(rejectErr)
457460
if rejectErr != nil {
@@ -485,7 +488,7 @@ func (m *consumer) processDelivery(ctx context.Context, controller Controller, d
485488
// Nack requeues immediately - the visibility timeout spaces retries.
486489
// The failure travels with it so that the attempt which finally spends
487490
// the retry budget can dead-letter saying why.
488-
nackOp := metrics.Begin(controllerScope, "nack", metrics.StorageLatencyBuckets)
491+
nackOp := metrics.Begin(controllerScope, "nack", metrics.StorageLatencyBuckets, metrics.TagsFromContext(ctx)...)
489492
nackErr := delivery.Nack(ctx, controllerFailure)
490493
nackOp.Complete(nackErr)
491494
if nackErr != nil {
@@ -505,7 +508,7 @@ func (m *consumer) processDelivery(ctx context.Context, controller Controller, d
505508
// the visibility timeout lapses into a normal redelivery, so the hold
506509
// loop's liveness never depends on this write succeeding.
507510
if wrapped.held {
508-
postponeOp := metrics.Begin(controllerScope, "postpone", metrics.StorageLatencyBuckets)
511+
postponeOp := metrics.Begin(controllerScope, "postpone", metrics.StorageLatencyBuckets, metrics.TagsFromContext(ctx)...)
509512
postponeErr := delivery.Postpone(ctx, wrapped.holdDelayMs)
510513
postponeOp.Complete(postponeErr)
511514
if postponeErr != nil {
@@ -530,7 +533,7 @@ func (m *consumer) processDelivery(ctx context.Context, controller Controller, d
530533
}
531534

532535
// Controller succeeded - ack message
533-
ackOp := metrics.Begin(controllerScope, "ack", metrics.StorageLatencyBuckets)
536+
ackOp := metrics.Begin(controllerScope, "ack", metrics.StorageLatencyBuckets, metrics.TagsFromContext(ctx)...)
534537
ackErr := delivery.Ack(ctx)
535538
ackOp.Complete(ackErr)
536539
if ackErr != nil {
@@ -578,7 +581,7 @@ func (m *consumer) checkGate(ctx context.Context, controller Controller, deliver
578581
// into a normal redelivery.
579582
return false
580583
}
581-
metrics.NamedCounter(scope, opName, "enter_errors", 1)
584+
metrics.NamedCounter(scope, opName, "enter_errors", 1, metrics.TagsFromContext(ctx)...)
582585
m.logger.Errorw("gate check failed, failing open",
583586
"consumer_group", consumerGroup,
584587
"topic", topic,
@@ -600,7 +603,7 @@ func (m *consumer) checkGate(ctx context.Context, controller Controller, deliver
600603
// earlier re-check, the gate has opened and the record must go so
601604
// observers see an empty parked set. A no-op when never parked.
602605
if unparkErr := entry.Unpark(ctx, descriptor); unparkErr != nil {
603-
metrics.NamedCounter(scope, opName, "unpark_errors", 1)
606+
metrics.NamedCounter(scope, opName, "unpark_errors", 1, metrics.TagsFromContext(ctx)...)
604607
m.logger.Warnw("failed to remove parked record on admit",
605608
"consumer_group", consumerGroup,
606609
"topic", topic,
@@ -615,7 +618,7 @@ func (m *consumer) checkGate(ctx context.Context, controller Controller, deliver
615618
// partition waits behind it (barrier) and the gate is re-checked on
616619
// redelivery without burning retry budget.
617620
if parkErr := entry.Park(ctx, descriptor); parkErr != nil {
618-
metrics.NamedCounter(scope, opName, "park_errors", 1)
621+
metrics.NamedCounter(scope, opName, "park_errors", 1, metrics.TagsFromContext(ctx)...)
619622
m.logger.Warnw("failed to write parked record, postponing anyway",
620623
"consumer_group", consumerGroup,
621624
"topic", topic,
@@ -624,8 +627,8 @@ func (m *consumer) checkGate(ctx context.Context, controller Controller, deliver
624627
)
625628
}
626629

627-
metrics.NamedCounter(scope, opName, "parked", 1)
628-
postponeOp := metrics.Begin(scope, "postpone", metrics.StorageLatencyBuckets)
630+
metrics.NamedCounter(scope, opName, "parked", 1, metrics.TagsFromContext(ctx)...)
631+
postponeOp := metrics.Begin(scope, "postpone", metrics.StorageLatencyBuckets, metrics.TagsFromContext(ctx)...)
629632
postponeErr := delivery.Postpone(ctx, m.gateRecheckDelayMs)
630633
postponeOp.Complete(postponeErr)
631634
if postponeErr != nil {

platform/consumer/consumer_test.go

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -286,10 +286,12 @@ func TestConsumer_ProcessDelivery_Success(t *testing.T) {
286286
c := New(logger, tally.NoopScope, reg, errs.NewClassifierProcessor(), consumergatenoop.New())
287287

288288
handledMsg := ""
289+
handledQueue := ""
289290
handler := &testController{}
290291
setupController(handler, "test-handler", testTopicKeyStart, "test-group",
291292
func(ctx context.Context, delivery Delivery) error {
292293
handledMsg = delivery.Message().ID
294+
handledQueue, _ = entityqueue.QueueName(ctx)
293295
return nil
294296
},
295297
)
@@ -303,14 +305,17 @@ func TestConsumer_ProcessDelivery_Success(t *testing.T) {
303305
err = c.Start(ctx)
304306
require.NoError(t, err)
305307

306-
msg := entityqueue.NewMessage("test-msg-1", []byte("payload"), "partition1", nil)
308+
msg := entityqueue.NewMessage("test-msg-1", []byte("payload"), "partition1", map[string]string{
309+
entityqueue.MetadataKeyQueueName: "monorepo/main",
310+
})
307311
mockDel := queuemock.NewMockDelivery(ctrl)
308312
done := setupDelivery(mockDel, msg, nil, nil)
309313

310314
deliveryChan <- mockDel
311315
<-done
312316

313317
assert.Equal(t, "test-msg-1", handledMsg)
318+
assert.Equal(t, "monorepo/main", handledQueue)
314319

315320
err = c.Stop(30000)
316321
require.NoError(t, err)

platform/metrics/BUILD.bazel

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -5,14 +5,18 @@ go_library(
55
srcs = ["metrics.go"],
66
importpath = "github.com/uber/submitqueue/platform/metrics",
77
visibility = ["//visibility:public"],
8-
deps = ["@com_github_uber_go_tally//:go_default_library"],
8+
deps = [
9+
"//platform/base/messagequeue:go_default_library",
10+
"@com_github_uber_go_tally//:go_default_library",
11+
],
912
)
1013

1114
go_test(
1215
name = "go_default_test",
1316
srcs = ["metrics_test.go"],
1417
embed = [":go_default_library"],
1518
deps = [
19+
"//platform/base/messagequeue:go_default_library",
1620
"@com_github_stretchr_testify//assert:go_default_library",
1721
"@com_github_uber_go_tally//:go_default_library",
1822
],

platform/metrics/metrics.go

Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,7 @@ import (
2020
"time"
2121

2222
"github.com/uber-go/tally"
23+
entityqueue "github.com/uber/submitqueue/platform/base/messagequeue"
2324
)
2425

2526
// Tag is a key-value pair attached to a metric for dimensional filtering.
@@ -35,6 +36,21 @@ func NewTag(key, value string) Tag {
3536
return Tag{Key: key, Value: value}
3637
}
3738

39+
// TagsFromContext appends the allowlisted metric tags carried by ctx. Additional
40+
// tags are preserved, and context-derived tags win when the same key is present.
41+
// Messages published before a context value was introduced return only the
42+
// additional tags.
43+
func TagsFromContext(ctx context.Context, tags ...Tag) []Tag {
44+
queueName, ok := entityqueue.QueueName(ctx)
45+
if !ok || queueName == "" {
46+
return tags
47+
}
48+
49+
contextTags := make([]Tag, 0, len(tags)+1)
50+
contextTags = append(contextTags, tags...)
51+
return append(contextTags, NewTag("queue", queueName))
52+
}
53+
3854
// Common duration bucket sets for latency histograms. Operations differ widely
3955
// in expected latency, so there is no single default — pick the set whose range
4056
// matches the operation and pass it to Begin or NamedHistogram. Buckets

platform/metrics/metrics_test.go

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,7 @@ import (
2222

2323
"github.com/stretchr/testify/assert"
2424
"github.com/uber-go/tally"
25+
entityqueue "github.com/uber/submitqueue/platform/base/messagequeue"
2526
)
2627

2728
func TestBegin_EmitsStart(t *testing.T) {
@@ -158,6 +159,19 @@ func TestNamedGauge(t *testing.T) {
158159
assert.Equal(t, float64(42), g.Value())
159160
}
160161

162+
func TestTagsFromContext(t *testing.T) {
163+
ctx := entityqueue.WithQueueName(context.Background(), "monorepo/main")
164+
165+
tags := TagsFromContext(ctx, NewTag("result", "success"), NewTag("queue", "wrong"))
166+
assert.Equal(t, []Tag{
167+
NewTag("result", "success"),
168+
NewTag("queue", "wrong"),
169+
NewTag("queue", "monorepo/main"),
170+
}, tags)
171+
assert.Empty(t, TagsFromContext(context.Background()))
172+
assert.Empty(t, TagsFromContext(entityqueue.WithQueueName(context.Background(), "")))
173+
}
174+
161175
func TestLatencyBuckets_Sorted(t *testing.T) {
162176
sets := map[string]tally.DurationBuckets{
163177
"FastLatencyBuckets": FastLatencyBuckets,

stovepipe/controller/BUILD.bazel

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -10,6 +10,7 @@ go_library(
1010
visibility = ["//visibility:public"],
1111
deps = [
1212
"//api/stovepipe/protopb:go_default_library",
13+
"//platform/base/messagequeue:go_default_library",
1314
"//platform/consumer:go_default_library",
1415
"//platform/errs:go_default_library",
1516
"//platform/extension/counter:go_default_library",
@@ -33,6 +34,7 @@ go_test(
3334
embed = [":go_default_library"],
3435
deps = [
3536
"//api/stovepipe/protopb:go_default_library",
37+
"//platform/base/messagequeue:go_default_library",
3638
"//platform/consumer:go_default_library",
3739
"//platform/extension/counter:go_default_library",
3840
"//platform/extension/counter/mock:go_default_library",

stovepipe/controller/build/BUILD.bazel

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,7 @@ go_library(
66
importpath = "github.com/uber/submitqueue/stovepipe/controller/build",
77
visibility = ["//visibility:public"],
88
deps = [
9+
"//platform/base/messagequeue:go_default_library",
910
"//platform/consumer:go_default_library",
1011
"//platform/errs:go_default_library",
1112
"//platform/metrics:go_default_library",

0 commit comments

Comments
 (0)