Skip to content

Commit 71eb860

Browse files
committed
feat(stovepipe): tag controller metrics by queue
Summary: Intent: - Make controller metrics attributable to their logical queue. - Keep queue tagging request-local without mutating shared controllers. Changes: - Pass queue as an ordinary metrics tag through primary and DLQ controller paths. - Reuse the existing package-level metric helpers and controller Tally scopes. - Keep pre-decode failures untagged and cover queue tags in tests. Test Plan: - Run targeted Bazel tests for platform metrics and the affected primary and DLQ controllers. Revert Plan: - Revert this commit to restore the prior 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 fee08f1 commit 71eb860

15 files changed

Lines changed: 219 additions & 136 deletions

File tree

stovepipe/controller/build/build.go

Lines changed: 5 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -88,24 +88,25 @@ 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+
queueTag := metrics.NewTag("queue", br.GetQueueName())
9192

9293
store, err := c.stores.For(storage.Config{QueueName: br.GetQueueName()})
9394
if err != nil {
94-
metrics.NamedCounter(c.metricsScope, _opName, "storage_resolve_errors", 1)
95+
metrics.NamedCounter(c.metricsScope, _opName, "storage_resolve_errors", 1, queueTag)
9596
// Non-retryable: a missing or unresolvable queue is a malformed message.
9697
return fmt.Errorf("failed to resolve storage for queue %q: %w", br.GetQueueName(), err)
9798
}
9899

99100
request, err := c.loadRequest(ctx, store, br.Id)
100101
if err != nil {
101-
metrics.NamedCounter(c.metricsScope, _opName, "storage_errors", 1)
102+
metrics.NamedCounter(c.metricsScope, _opName, "storage_errors", 1, queueTag)
102103
return err
103104
}
104105

105106
// The payload's queue must match the request's authoritative queue; a
106107
// mismatch is a malformed message. Non-retryable — reject to the DLQ.
107108
if br.GetQueueName() != "" && br.GetQueueName() != request.Queue {
108-
metrics.NamedCounter(c.metricsScope, _opName, "queue_mismatch", 1)
109+
metrics.NamedCounter(c.metricsScope, _opName, "queue_mismatch", 1, queueTag)
109110
return fmt.Errorf("payload queue %q does not match queue %q of request %s", br.GetQueueName(), request.Queue, request.ID)
110111
}
111112

@@ -123,7 +124,7 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) er
123124

124125
// process decided the scope; build never re-derives incremental-vs-full.
125126
if request.BuildStrategy == entity.BuildStrategyUnknown {
126-
metrics.NamedCounter(c.metricsScope, _opName, "strategy_not_visible", 1)
127+
metrics.NamedCounter(c.metricsScope, _opName, "strategy_not_visible", 1, queueTag)
127128
return errs.NewRetryableError(fmt.Errorf("request %s has no build strategy yet", request.ID))
128129
}
129130
baseURI := ""

stovepipe/controller/build/build_test.go

Lines changed: 18 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -52,6 +52,7 @@ type buildMocks struct {
5252
runnerFactory *buildrunnermock.MockFactory
5353
runner *buildrunnermock.MockBuildRunner
5454
publisher *mqmock.MockPublisher
55+
metricsScope tally.TestScope
5556
}
5657

5758
// staticStorageFactory resolves every queue to one fixed store aggregate.
@@ -63,12 +64,14 @@ func (f staticStorageFactory) For(storage.Config) (storage.Storage, error) { ret
6364
func newController(t *testing.T, ctrl *gomock.Controller) (*Controller, buildMocks) {
6465
t.Helper()
6566

67+
scope := tally.NewTestScope("test", nil)
6668
m := buildMocks{
6769
reqStore: storagemock.NewMockRequestStore(ctrl),
6870
buildStore: storagemock.NewMockBuildStore(ctrl),
6971
runnerFactory: buildrunnermock.NewMockFactory(ctrl),
7072
runner: buildrunnermock.NewMockBuildRunner(ctrl),
7173
publisher: mqmock.NewMockPublisher(ctrl),
74+
metricsScope: scope,
7275
}
7376

7477
store := storagemock.NewMockStorage(ctrl)
@@ -83,10 +86,23 @@ func newController(t *testing.T, ctrl *gomock.Controller) (*Controller, buildMoc
8386
})
8487
require.NoError(t, err)
8588

86-
c := NewController(zap.NewNop().Sugar(), tally.NewTestScope("test", nil), staticStorageFactory{store: store}, m.runnerFactory, registry, stovepipemq.TopicKeyBuild, "stovepipe-build")
89+
c := NewController(zap.NewNop().Sugar(), scope, staticStorageFactory{store: store}, m.runnerFactory, registry, stovepipemq.TopicKeyBuild, "stovepipe-build")
8790
return c, m
8891
}
8992

93+
func TestProcessTagsMetricsWithQueue(t *testing.T) {
94+
ctrl := gomock.NewController(t)
95+
c, m := newController(t, ctrl)
96+
m.reqStore.EXPECT().Get(gomock.Any(), testID).
97+
Return(processingRequest(entity.BuildStrategyUnknown, ""), nil)
98+
m.runnerFactory.EXPECT().For(buildrunner.Config{QueueName: testQueue}).Return(m.runner, nil)
99+
100+
require.Error(t, c.Process(context.Background(), delivery(t, ctrl, buildPayload(t, testID))))
101+
counter, ok := m.metricsScope.Snapshot().Counters()["test.build_controller.build.strategy_not_visible+queue=monorepo/main"]
102+
require.True(t, ok)
103+
assert.EqualValues(t, 1, counter.Value())
104+
}
105+
90106
func delivery(t *testing.T, ctrl *gomock.Controller, payload []byte) consumer.Delivery {
91107
t.Helper()
92108
d := consumermock.NewMockDelivery(ctrl)
@@ -97,7 +113,7 @@ func delivery(t *testing.T, ctrl *gomock.Controller, payload []byte) consumer.De
97113

98114
func buildPayload(t *testing.T, id string) []byte {
99115
t.Helper()
100-
b, err := stovepipemq.Marshal(&stovepipemq.BuildRequest{Id: id})
116+
b, err := stovepipemq.Marshal(&stovepipemq.BuildRequest{Id: id, QueueName: testQueue})
101117
require.NoError(t, err)
102118
return b
103119
}

stovepipe/controller/buildsignal/buildsignal.go

Lines changed: 15 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -108,30 +108,31 @@ 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+
queueTag := metrics.NewTag("queue", sig.GetQueueName())
111112

112113
store, err := c.stores.For(storage.Config{QueueName: sig.GetQueueName()})
113114
if err != nil {
114-
metrics.NamedCounter(c.metricsScope, _opName, "storage_resolve_errors", 1)
115+
metrics.NamedCounter(c.metricsScope, _opName, "storage_resolve_errors", 1, queueTag)
115116
// Non-retryable: a missing or unresolvable queue is a malformed message.
116117
return fmt.Errorf("failed to resolve storage for queue %q: %w", sig.GetQueueName(), err)
117118
}
118119

119120
build, err := c.loadBuild(ctx, store, sig.Id)
120121
if err != nil {
121-
metrics.NamedCounter(c.metricsScope, _opName, "storage_errors", 1)
122+
metrics.NamedCounter(c.metricsScope, _opName, "storage_errors", 1, queueTag)
122123
return err
123124
}
124125

125126
request, err := c.loadRequest(ctx, store, build.RequestID)
126127
if err != nil {
127-
metrics.NamedCounter(c.metricsScope, _opName, "storage_errors", 1)
128+
metrics.NamedCounter(c.metricsScope, _opName, "storage_errors", 1, queueTag)
128129
return err
129130
}
130131

131132
// The payload's queue must match the request's authoritative queue; a
132133
// mismatch is a malformed message. Non-retryable — reject to the DLQ.
133134
if sig.GetQueueName() != "" && sig.GetQueueName() != request.Queue {
134-
metrics.NamedCounter(c.metricsScope, _opName, "queue_mismatch", 1)
135+
metrics.NamedCounter(c.metricsScope, _opName, "queue_mismatch", 1, queueTag)
135136
return fmt.Errorf("payload queue %q does not match queue %q of request %s", sig.GetQueueName(), request.Queue, request.ID)
136137
}
137138

@@ -167,7 +168,7 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) er
167168
}
168169

169170
if effective.IsTerminal() {
170-
if err := c.finishRequest(ctx, store, &request, effective); err != nil {
171+
if err := c.finishRequest(ctx, queueTag, store, &request, effective); err != nil {
171172
return err
172173
}
173174
if err := c.publishRecord(ctx, request.ID, request.Queue); err != nil {
@@ -208,18 +209,18 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) er
208209
// the request non-terminal, so redelivery re-runs both steps and decrements again
209210
// — transiently over-admitting by one until releaseBuildSlot's zero clamp
210211
// reconverges, which is the failure mode this pipeline prefers.
211-
func (c *Controller) finishRequest(ctx context.Context, store storage.Storage, request *entity.Request, status entity.BuildStatus) error {
212+
func (c *Controller) finishRequest(ctx context.Context, queueTag metrics.Tag, store storage.Storage, request *entity.Request, status entity.BuildStatus) error {
212213
if request.State.HasBuildOutcome() {
213214
return nil
214215
}
215216

216-
if err := c.releaseBuildSlot(ctx, store, request.Queue); err != nil {
217-
metrics.NamedCounter(c.metricsScope, _opName, "storage_errors", 1)
217+
if err := c.releaseBuildSlot(ctx, queueTag, store, request.Queue); err != nil {
218+
metrics.NamedCounter(c.metricsScope, _opName, "storage_errors", 1, queueTag)
218219
return err
219220
}
220221

221-
if err := c.markOutcome(ctx, store, request, outcomeState(status)); err != nil {
222-
metrics.NamedCounter(c.metricsScope, _opName, "storage_errors", 1)
222+
if err := c.markOutcome(ctx, queueTag, store, request, outcomeState(status)); err != nil {
223+
metrics.NamedCounter(c.metricsScope, _opName, "storage_errors", 1, queueTag)
223224
return err
224225
}
225226
return nil
@@ -243,7 +244,7 @@ func outcomeState(status entity.BuildStatus) entity.RequestState {
243244
// conflicts. First writer wins: once any outcome is recorded a later caller leaves it
244245
// alone, so duplicate builds for one request (which build.md accepts) cannot flip the
245246
// verdict back and forth.
246-
func (c *Controller) markOutcome(ctx context.Context, store storage.Storage, request *entity.Request, state entity.RequestState) error {
247+
func (c *Controller) markOutcome(ctx context.Context, queueTag metrics.Tag, store storage.Storage, request *entity.Request, state entity.RequestState) error {
247248
reqStore := store.GetRequestStore()
248249

249250
for {
@@ -268,6 +269,7 @@ func (c *Controller) markOutcome(ctx context.Context, store storage.Storage, req
268269
updated.Version = newVersion
269270
*request = updated
270271
metrics.NamedCounter(c.metricsScope, _opName, "outcomes", 1,
272+
queueTag,
271273
metrics.NewTag("state", string(state)),
272274
)
273275
return nil
@@ -279,7 +281,7 @@ func (c *Controller) markOutcome(ctx context.Context, store storage.Storage, req
279281
// (preserving concurrent updates), clamps at zero, and retries on version conflicts.
280282
// Unlike process's unwind-path release this is not best-effort: the caller must not
281283
// mark the request terminal if the slot was not freed, so a hard failure is returned.
282-
func (c *Controller) releaseBuildSlot(ctx context.Context, store storage.Storage, queueName string) error {
284+
func (c *Controller) releaseBuildSlot(ctx context.Context, queueTag metrics.Tag, store storage.Storage, queueName string) error {
283285
queueStore := store.GetQueueStore()
284286

285287
for {
@@ -300,7 +302,7 @@ func (c *Controller) releaseBuildSlot(ctx context.Context, store storage.Storage
300302
}
301303
return fmt.Errorf("failed to release build slot for queue %s: %w", queueName, err)
302304
}
303-
metrics.NamedCounter(c.metricsScope, _opName, "slot_released", 1)
305+
metrics.NamedCounter(c.metricsScope, _opName, "slot_released", 1, queueTag)
304306
return nil
305307
}
306308
}

stovepipe/controller/buildsignal/buildsignal_test.go

Lines changed: 16 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -52,6 +52,7 @@ type buildsignalMocks struct {
5252
runnerFactory *buildrunnermock.MockFactory
5353
runner *buildrunnermock.MockBuildRunner
5454
publisher *mqmock.MockPublisher
55+
metricsScope tally.TestScope
5556
}
5657

5758
// staticStorageFactory resolves every queue to one fixed store aggregate.
@@ -63,13 +64,15 @@ func (f staticStorageFactory) For(storage.Config) (storage.Storage, error) { ret
6364
func newController(t *testing.T, ctrl *gomock.Controller) (*Controller, buildsignalMocks) {
6465
t.Helper()
6566

67+
scope := tally.NewTestScope("test", nil)
6668
m := buildsignalMocks{
6769
reqStore: storagemock.NewMockRequestStore(ctrl),
6870
buildStore: storagemock.NewMockBuildStore(ctrl),
6971
queueStore: storagemock.NewMockQueueStore(ctrl),
7072
runnerFactory: buildrunnermock.NewMockFactory(ctrl),
7173
runner: buildrunnermock.NewMockBuildRunner(ctrl),
7274
publisher: mqmock.NewMockPublisher(ctrl),
75+
metricsScope: scope,
7376
}
7477

7578
store := storagemock.NewMockStorage(ctrl)
@@ -86,10 +89,21 @@ func newController(t *testing.T, ctrl *gomock.Controller) (*Controller, buildsig
8689
})
8790
require.NoError(t, err)
8891

89-
c := NewController(zap.NewNop().Sugar(), tally.NewTestScope("test", nil), staticStorageFactory{store: store}, m.runnerFactory, registry, stovepipemq.TopicKeyBuildSignal, "stovepipe-buildsignal")
92+
c := NewController(zap.NewNop().Sugar(), scope, staticStorageFactory{store: store}, m.runnerFactory, registry, stovepipemq.TopicKeyBuildSignal, "stovepipe-buildsignal")
9093
return c, m
9194
}
9295

96+
func TestProcessTagsMetricsWithQueue(t *testing.T) {
97+
ctrl := gomock.NewController(t)
98+
c, m := newController(t, ctrl)
99+
m.buildStore.EXPECT().Get(gomock.Any(), testBuildID).Return(entity.Build{}, assert.AnError)
100+
101+
require.Error(t, c.Process(context.Background(), delivery(t, ctrl, buildSignalPayload(t, testBuildID))))
102+
counter, ok := m.metricsScope.Snapshot().Counters()["test.buildsignal_controller.buildsignal.storage_errors+queue=monorepo/main"]
103+
require.True(t, ok)
104+
assert.EqualValues(t, 1, counter.Value())
105+
}
106+
93107
func delivery(t *testing.T, ctrl *gomock.Controller, payload []byte) *consumermock.MockDelivery {
94108
t.Helper()
95109
d := consumermock.NewMockDelivery(ctrl)
@@ -100,7 +114,7 @@ func delivery(t *testing.T, ctrl *gomock.Controller, payload []byte) *consumermo
100114

101115
func buildSignalPayload(t *testing.T, id string) []byte {
102116
t.Helper()
103-
b, err := stovepipemq.Marshal(&stovepipemq.BuildSignal{Id: id})
117+
b, err := stovepipemq.Marshal(&stovepipemq.BuildSignal{Id: id, QueueName: testQueue})
104118
require.NoError(t, err)
105119
return b
106120
}

stovepipe/controller/dlq/build.go

Lines changed: 5 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -65,14 +65,15 @@ func (c *buildController) Process(ctx context.Context, delivery consumer.Deliver
6565
metrics.NamedCounter(c.metricsScope, _buildOpName, "deserialize_errors", 1)
6666
return fmt.Errorf("failed to decode dlq payload: %w", err)
6767
}
68+
queueTag := metrics.NewTag("queue", buildRequest.GetQueueName())
6869
if buildRequest.Id == "" {
69-
metrics.NamedCounter(c.metricsScope, _buildOpName, "empty_id_errors", 1)
70+
metrics.NamedCounter(c.metricsScope, _buildOpName, "empty_id_errors", 1, queueTag)
7071
return fmt.Errorf("build dlq payload decoded to empty request id")
7172
}
7273

7374
store, err := c.stores.For(storage.Config{QueueName: buildRequest.GetQueueName()})
7475
if err != nil {
75-
metrics.NamedCounter(c.metricsScope, _buildOpName, "storage_resolve_errors", 1)
76+
metrics.NamedCounter(c.metricsScope, _buildOpName, "storage_resolve_errors", 1, queueTag)
7677
return fmt.Errorf("failed to resolve storage for queue %q: %w", buildRequest.GetQueueName(), err)
7778
}
7879

@@ -86,10 +87,10 @@ func (c *buildController) Process(ctx context.Context, delivery consumer.Deliver
8687
)
8788

8889
if err := failRequest(ctx, store, c.logger, buildRequest.Id); err != nil {
89-
metrics.NamedCounter(c.metricsScope, _buildOpName, "reconcile_errors", 1)
90+
metrics.NamedCounter(c.metricsScope, _buildOpName, "reconcile_errors", 1, queueTag)
9091
return err
9192
}
92-
metrics.NamedCounter(c.metricsScope, _buildOpName, "reconciled", 1)
93+
metrics.NamedCounter(c.metricsScope, _buildOpName, "reconciled", 1, queueTag)
9394
return nil
9495
}
9596

stovepipe/controller/dlq/build_test.go

Lines changed: 22 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -32,15 +32,17 @@ import (
3232
func newBuildController(t *testing.T, ctrl *gomock.Controller) (consumer.Controller, dlqMocks) {
3333
t.Helper()
3434

35+
scope := tally.NewTestScope("test", nil)
3536
m := dlqMocks{
36-
reqStore: storagemock.NewMockRequestStore(ctrl),
37-
queueStore: storagemock.NewMockQueueStore(ctrl),
37+
reqStore: storagemock.NewMockRequestStore(ctrl),
38+
queueStore: storagemock.NewMockQueueStore(ctrl),
39+
metricsScope: scope,
3840
}
3941
store := storagemock.NewMockStorage(ctrl)
4042
store.EXPECT().GetRequestStore().Return(m.reqStore).AnyTimes()
4143
store.EXPECT().GetQueueStore().Return(m.queueStore).AnyTimes()
4244

43-
c := NewDLQBuildController(zap.NewNop().Sugar(), tally.NewTestScope("test", nil), staticStorageFactory{store: store}, TopicKey(stovepipemq.TopicKeyBuild), "stovepipe-build-dlq")
45+
c := NewDLQBuildController(zap.NewNop().Sugar(), scope, staticStorageFactory{store: store}, TopicKey(stovepipemq.TopicKeyBuild), "stovepipe-build-dlq")
4446
return c, m
4547
}
4648

@@ -53,10 +55,11 @@ func buildPayload(t *testing.T, id string) []byte {
5355

5456
func TestBuildProcess(t *testing.T) {
5557
tests := []struct {
56-
name string
57-
payload []byte
58-
setup func(m dlqMocks)
59-
wantErr bool
58+
name string
59+
payload []byte
60+
setup func(m dlqMocks)
61+
wantErr bool
62+
wantMetric string
6063
}{
6164
{
6265
name: "processing request is failed",
@@ -68,9 +71,15 @@ func TestBuildProcess(t *testing.T) {
6871
updated.State = entity.RequestStateFailed
6972
m.reqStore.EXPECT().Update(gomock.Any(), updated, int32(2), int32(3)).Return(nil)
7073
},
74+
wantMetric: "test.build_dlq_controller.build_dlq.reconciled+queue=monorepo/main",
7175
},
7276
{name: "malformed payload is returned", payload: []byte("not-a-proto"), wantErr: true},
73-
{name: "empty request id is returned", payload: buildPayload(t, ""), wantErr: true},
77+
{
78+
name: "empty request id is returned",
79+
payload: buildPayload(t, ""),
80+
wantErr: true,
81+
wantMetric: "test.build_dlq_controller.build_dlq.empty_id_errors+queue=monorepo/main",
82+
},
7483
}
7584

7685
for _, tt := range tests {
@@ -86,6 +95,11 @@ func TestBuildProcess(t *testing.T) {
8695
payload = buildPayload(t, testID)
8796
}
8897
err := controller.Process(context.Background(), delivery(t, ctrl, payload))
98+
if tt.wantMetric != "" {
99+
counter, ok := mocks.metricsScope.Snapshot().Counters()[tt.wantMetric]
100+
require.True(t, ok)
101+
assert.EqualValues(t, 1, counter.Value())
102+
}
89103
if tt.wantErr {
90104
require.Error(t, err)
91105
return

0 commit comments

Comments
 (0)