From 0a71da99714b655f6eed0cbaa164256813f2c3cc Mon Sep 17 00:00:00 2001 From: mnoah1 Date: Mon, 24 Aug 2026 19:21:39 +0000 Subject: [PATCH] feat(stovepipe): tag controller metrics by queue MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Summary: Intent: - Make controller metrics attributable to their logical queue. - Enforce request-local metric scoping without mutating shared controllers. Changes: - Add a metrics factory that resolves queue-bound scopes from explicit config. - Pass typed queue-scoped metrics through every primary and DLQ controller path. - Keep pre-decode failures explicitly unscoped and cover queue tags in tests. Test Plan: - Exercise primary and DLQ controller metric paths and assert emitted series carry the expected queue tag. - Run the affected packages with the Go race detector and targeted Bazel tests. Revert Plan: - Revert this commit to restore the prior unscoped controller metrics. --- 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 --- platform/metrics/metrics.go | 69 ++++++++++- platform/metrics/metrics_test.go | 29 +++++ stovepipe/controller/build/build.go | 39 +++--- stovepipe/controller/build/build_test.go | 20 ++- .../controller/buildsignal/buildsignal.go | 59 ++++----- .../buildsignal/buildsignal_test.go | 18 ++- stovepipe/controller/dlq/build.go | 31 ++--- stovepipe/controller/dlq/build_test.go | 30 +++-- stovepipe/controller/dlq/buildsignal.go | 37 +++--- stovepipe/controller/dlq/buildsignal_test.go | 40 +++--- stovepipe/controller/dlq/dlq_test.go | 38 +++--- stovepipe/controller/dlq/request.go | 29 ++--- stovepipe/controller/process/BUILD.bazel | 1 + stovepipe/controller/process/process.go | 91 +++++++------- stovepipe/controller/process/process_test.go | 15 ++- stovepipe/controller/record/record.go | 115 +++++++++--------- stovepipe/controller/record/record_test.go | 3 + 17 files changed, 414 insertions(+), 250 deletions(-) diff --git a/platform/metrics/metrics.go b/platform/metrics/metrics.go index 20f98c755..1c5761142 100644 --- a/platform/metrics/metrics.go +++ b/platform/metrics/metrics.go @@ -30,6 +30,69 @@ type Tag struct { Value string } +// Config identifies the logical queue for metrics returned by a Factory. +type Config struct { + QueueName string +} + +// Factory creates configured metric scopes from a shared Tally scope. +type Factory struct { + scope tally.Scope +} + +// NewFactory returns a Factory backed by scope. +func NewFactory(scope tally.Scope) Factory { + return Factory{scope: scope} +} + +// For returns a metric scope bound to the configured queue. +func (f Factory) For(config Config) Scope { + return Scope{ + scope: f.scope, + configuredTags: []Tag{NewTag("queue", config.QueueName)}, + } +} + +// Base returns an emitter on the factory's underlying scope without adding a +// queue tag. Existing namespaces and inherited tags remain intact. It is +// intended for failures that happen before a message's queue can be decoded. +func (f Factory) Base() Scope { + return Scope{scope: f.scope} +} + +// Scope emits named metrics on a factory's underlying Tally scope. Factory.For +// may bind configured tags that callers cannot override; Factory.Base returns +// the same type without adding tags. +type Scope struct { + scope tally.Scope + configuredTags []Tag +} + +// NamedCounter increments the {name}.{counter} counter by value. +func (s Scope) NamedCounter(name string, counter string, value int64, tags ...Tag) { + tagged(s.scope, s.withConfiguredTags(tags)).SubScope(name).Counter(counter).Inc(value) +} + +// NamedHistogram returns a tally.Histogram at {name}.{histogram} with the given +// bucket configuration. +func (s Scope) NamedHistogram(name string, histogram string, buckets tally.Buckets, tags ...Tag) tally.Histogram { + return tagged(s.scope, s.withConfiguredTags(tags)).SubScope(name).Histogram(histogram, buckets) +} + +// NamedGauge sets the {name}.{gauge} gauge to value. +func (s Scope) NamedGauge(name string, gauge string, value float64, tags ...Tag) { + tagged(s.scope, s.withConfiguredTags(tags)).SubScope(name).Gauge(gauge).Update(value) +} + +func (s Scope) withConfiguredTags(tags []Tag) []Tag { + if len(s.configuredTags) == 0 { + return tags + } + configuredTags := make([]Tag, 0, len(tags)+len(s.configuredTags)) + configuredTags = append(configuredTags, tags...) + return append(configuredTags, s.configuredTags...) +} + // NewTag creates a Tag with the given key and value. func NewTag(key, value string) Tag { return Tag{Key: key, Value: value} @@ -183,19 +246,19 @@ func (o Op) Complete(err error, tags ...Tag) { // NamedCounter increments the {name}.{counter} counter by value. func NamedCounter(scope tally.Scope, name string, counter string, value int64, tags ...Tag) { - tagged(scope, tags).SubScope(name).Counter(counter).Inc(value) + Scope{scope: scope}.NamedCounter(name, counter, value, tags...) } // NamedHistogram returns a tally.Histogram at {name}.{histogram} with the given // bucket configuration. Store the returned histogram and call RecordDuration or // RecordValue on each invocation. func NamedHistogram(scope tally.Scope, name string, histogram string, buckets tally.Buckets, tags ...Tag) tally.Histogram { - return tagged(scope, tags).SubScope(name).Histogram(histogram, buckets) + return Scope{scope: scope}.NamedHistogram(name, histogram, buckets, tags...) } // NamedGauge sets the {name}.{gauge} gauge to value. func NamedGauge(scope tally.Scope, name string, gauge string, value float64, tags ...Tag) { - tagged(scope, tags).SubScope(name).Gauge(gauge).Update(value) + Scope{scope: scope}.NamedGauge(name, gauge, value, tags...) } // tagsToMap converts a slice of Tag to a map for tally. diff --git a/platform/metrics/metrics_test.go b/platform/metrics/metrics_test.go index 3dd70c23e..7b5b72202 100644 --- a/platform/metrics/metrics_test.go +++ b/platform/metrics/metrics_test.go @@ -158,6 +158,35 @@ func TestNamedGauge(t *testing.T) { assert.Equal(t, float64(42), g.Value()) } +func TestFactory(t *testing.T) { + scope := tally.NewTestScope("", nil) + factory := NewFactory(scope) + queueScope := factory.For(Config{QueueName: "monorepo/main"}) + base := factory.Base() + + queueScope.NamedCounter("process", "attempts", 2, NewTag("result", "success"), NewTag("queue", "wrong")) + queueScope.NamedGauge("process", "in_flight", 3) + queueScope.NamedHistogram("process", "duration", StorageLatencyBuckets).RecordDuration(time.Second) + base.NamedCounter("process", "deserialize_errors", 1) + base.NamedGauge("process", "decode_in_flight", 1) + base.NamedHistogram("process", "decode_duration", StorageLatencyBuckets).RecordDuration(time.Second) + + snapshot := scope.Snapshot() + counter, ok := snapshot.Counters()["process.attempts+queue=monorepo/main,result=success"] + assert.True(t, ok) + assert.EqualValues(t, 2, counter.Value()) + _, ok = snapshot.Gauges()["process.in_flight+queue=monorepo/main"] + assert.True(t, ok) + _, ok = snapshot.Histograms()["process.duration+queue=monorepo/main"] + assert.True(t, ok) + _, ok = snapshot.Counters()["process.deserialize_errors+"] + assert.True(t, ok) + _, ok = snapshot.Gauges()["process.decode_in_flight+"] + assert.True(t, ok) + _, ok = snapshot.Histograms()["process.decode_duration+"] + assert.True(t, ok) +} + func TestLatencyBuckets_Sorted(t *testing.T) { sets := map[string]tally.DurationBuckets{ "FastLatencyBuckets": FastLatencyBuckets, diff --git a/stovepipe/controller/build/build.go b/stovepipe/controller/build/build.go index 8bd08c6cd..ede84faaa 100644 --- a/stovepipe/controller/build/build.go +++ b/stovepipe/controller/build/build.go @@ -40,13 +40,13 @@ import ( // triggers a build for its already-decided scope, and publishes the resulting // build id to buildsignal. Implements consumer.Controller. type Controller struct { - logger *zap.SugaredLogger - metricsScope tally.Scope - stores storage.Factory - buildRunners buildrunner.Factory - registry consumer.TopicRegistry - topicKey consumer.TopicKey - consumerGroup string + logger *zap.SugaredLogger + metricsFactory metrics.Factory + stores storage.Factory + buildRunners buildrunner.Factory + registry consumer.TopicRegistry + topicKey consumer.TopicKey + consumerGroup string } // Verify Controller implements consumer.Controller interface at compile time. @@ -66,13 +66,13 @@ func NewController( consumerGroup string, ) *Controller { return &Controller{ - logger: logger.Named("build_controller"), - metricsScope: scope.SubScope("build_controller"), - stores: stores, - buildRunners: buildRunners, - registry: registry, - topicKey: topicKey, - consumerGroup: consumerGroup, + logger: logger.Named("build_controller"), + metricsFactory: metrics.NewFactory(scope.SubScope("build_controller")), + stores: stores, + buildRunners: buildRunners, + registry: registry, + topicKey: topicKey, + consumerGroup: consumerGroup, } } @@ -84,28 +84,29 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) er br := &stovepipemq.BuildRequest{} if err := stovepipemq.Unmarshal(msg.Payload, br); err != nil { - metrics.NamedCounter(c.metricsScope, _opName, "deserialize_errors", 1) + c.metricsFactory.Base().NamedCounter(_opName, "deserialize_errors", 1) // Non-retryable: a malformed message will never succeed regardless of retries. return fmt.Errorf("failed to deserialize build request: %w", err) } + messageMetrics := c.metricsFactory.For(metrics.Config{QueueName: br.GetQueueName()}) store, err := c.stores.For(storage.Config{QueueName: br.GetQueueName()}) if err != nil { - metrics.NamedCounter(c.metricsScope, _opName, "storage_resolve_errors", 1) + messageMetrics.NamedCounter(_opName, "storage_resolve_errors", 1) // Non-retryable: a missing or unresolvable queue is a malformed message. return fmt.Errorf("failed to resolve storage for queue %q: %w", br.GetQueueName(), err) } request, err := c.loadRequest(ctx, store, br.Id) if err != nil { - metrics.NamedCounter(c.metricsScope, _opName, "storage_errors", 1) + messageMetrics.NamedCounter(_opName, "storage_errors", 1) return err } // The payload's queue must match the request's authoritative queue; a // mismatch is a malformed message. Non-retryable — reject to the DLQ. if br.GetQueueName() != "" && br.GetQueueName() != request.Queue { - metrics.NamedCounter(c.metricsScope, _opName, "queue_mismatch", 1) + messageMetrics.NamedCounter(_opName, "queue_mismatch", 1) return fmt.Errorf("payload queue %q does not match queue %q of request %s", br.GetQueueName(), request.Queue, request.ID) } @@ -123,7 +124,7 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) er // process decided the scope; build never re-derives incremental-vs-full. if request.BuildStrategy == entity.BuildStrategyUnknown { - metrics.NamedCounter(c.metricsScope, _opName, "strategy_not_visible", 1) + messageMetrics.NamedCounter(_opName, "strategy_not_visible", 1) return errs.NewRetryableError(fmt.Errorf("request %s has no build strategy yet", request.ID)) } baseURI := "" diff --git a/stovepipe/controller/build/build_test.go b/stovepipe/controller/build/build_test.go index 1e7d85207..bdb8db0a6 100644 --- a/stovepipe/controller/build/build_test.go +++ b/stovepipe/controller/build/build_test.go @@ -52,6 +52,7 @@ type buildMocks struct { runnerFactory *buildrunnermock.MockFactory runner *buildrunnermock.MockBuildRunner publisher *mqmock.MockPublisher + metricsScope tally.TestScope } // staticStorageFactory resolves every queue to one fixed store aggregate. @@ -63,12 +64,14 @@ func (f staticStorageFactory) For(storage.Config) (storage.Storage, error) { ret func newController(t *testing.T, ctrl *gomock.Controller) (*Controller, buildMocks) { t.Helper() + scope := tally.NewTestScope("test", nil) m := buildMocks{ reqStore: storagemock.NewMockRequestStore(ctrl), buildStore: storagemock.NewMockBuildStore(ctrl), runnerFactory: buildrunnermock.NewMockFactory(ctrl), runner: buildrunnermock.NewMockBuildRunner(ctrl), publisher: mqmock.NewMockPublisher(ctrl), + metricsScope: scope, } store := storagemock.NewMockStorage(ctrl) @@ -83,10 +86,23 @@ func newController(t *testing.T, ctrl *gomock.Controller) (*Controller, buildMoc }) require.NoError(t, err) - c := NewController(zap.NewNop().Sugar(), tally.NewTestScope("test", nil), staticStorageFactory{store: store}, m.runnerFactory, registry, stovepipemq.TopicKeyBuild, "stovepipe-build") + c := NewController(zap.NewNop().Sugar(), scope, staticStorageFactory{store: store}, m.runnerFactory, registry, stovepipemq.TopicKeyBuild, "stovepipe-build") return c, m } +func TestProcessTagsMetricsWithQueue(t *testing.T) { + ctrl := gomock.NewController(t) + c, m := newController(t, ctrl) + m.reqStore.EXPECT().Get(gomock.Any(), testID). + Return(processingRequest(entity.BuildStrategyUnknown, ""), nil) + m.runnerFactory.EXPECT().For(buildrunner.Config{QueueName: testQueue}).Return(m.runner, nil) + + require.Error(t, c.Process(context.Background(), delivery(t, ctrl, buildPayload(t, testID)))) + counter, ok := m.metricsScope.Snapshot().Counters()["test.build_controller.build.strategy_not_visible+queue=monorepo/main"] + require.True(t, ok) + assert.EqualValues(t, 1, counter.Value()) +} + func delivery(t *testing.T, ctrl *gomock.Controller, payload []byte) consumer.Delivery { t.Helper() d := consumermock.NewMockDelivery(ctrl) @@ -97,7 +113,7 @@ func delivery(t *testing.T, ctrl *gomock.Controller, payload []byte) consumer.De func buildPayload(t *testing.T, id string) []byte { t.Helper() - b, err := stovepipemq.Marshal(&stovepipemq.BuildRequest{Id: id}) + b, err := stovepipemq.Marshal(&stovepipemq.BuildRequest{Id: id, QueueName: testQueue}) require.NoError(t, err) return b } diff --git a/stovepipe/controller/buildsignal/buildsignal.go b/stovepipe/controller/buildsignal/buildsignal.go index c040d239c..fbc7706aa 100644 --- a/stovepipe/controller/buildsignal/buildsignal.go +++ b/stovepipe/controller/buildsignal/buildsignal.go @@ -57,13 +57,13 @@ var ( // the request, and publishes the request id to record. Implements // consumer.Controller. type Controller struct { - logger *zap.SugaredLogger - metricsScope tally.Scope - stores storage.Factory - buildRunners buildrunner.Factory - registry consumer.TopicRegistry - topicKey consumer.TopicKey - consumerGroup string + logger *zap.SugaredLogger + metricsFactory metrics.Factory + stores storage.Factory + buildRunners buildrunner.Factory + registry consumer.TopicRegistry + topicKey consumer.TopicKey + consumerGroup string } // Verify Controller implements consumer.Controller interface at compile time. @@ -83,13 +83,13 @@ func NewController( consumerGroup string, ) *Controller { return &Controller{ - logger: logger.Named("buildsignal_controller"), - metricsScope: scope.SubScope("buildsignal_controller"), - stores: stores, - buildRunners: buildRunners, - registry: registry, - topicKey: topicKey, - consumerGroup: consumerGroup, + logger: logger.Named("buildsignal_controller"), + metricsFactory: metrics.NewFactory(scope.SubScope("buildsignal_controller")), + stores: stores, + buildRunners: buildRunners, + registry: registry, + topicKey: topicKey, + consumerGroup: consumerGroup, } } @@ -104,34 +104,35 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) er sig := &stovepipemq.BuildSignal{} if err := stovepipemq.Unmarshal(msg.Payload, sig); err != nil { - metrics.NamedCounter(c.metricsScope, _opName, "deserialize_errors", 1) + c.metricsFactory.Base().NamedCounter(_opName, "deserialize_errors", 1) // Non-retryable: a malformed message will never succeed regardless of retries. return fmt.Errorf("failed to deserialize build signal: %w", err) } + messageMetrics := c.metricsFactory.For(metrics.Config{QueueName: sig.GetQueueName()}) store, err := c.stores.For(storage.Config{QueueName: sig.GetQueueName()}) if err != nil { - metrics.NamedCounter(c.metricsScope, _opName, "storage_resolve_errors", 1) + messageMetrics.NamedCounter(_opName, "storage_resolve_errors", 1) // Non-retryable: a missing or unresolvable queue is a malformed message. return fmt.Errorf("failed to resolve storage for queue %q: %w", sig.GetQueueName(), err) } build, err := c.loadBuild(ctx, store, sig.Id) if err != nil { - metrics.NamedCounter(c.metricsScope, _opName, "storage_errors", 1) + messageMetrics.NamedCounter(_opName, "storage_errors", 1) return err } request, err := c.loadRequest(ctx, store, build.RequestID) if err != nil { - metrics.NamedCounter(c.metricsScope, _opName, "storage_errors", 1) + messageMetrics.NamedCounter(_opName, "storage_errors", 1) return err } // The payload's queue must match the request's authoritative queue; a // mismatch is a malformed message. Non-retryable — reject to the DLQ. if sig.GetQueueName() != "" && sig.GetQueueName() != request.Queue { - metrics.NamedCounter(c.metricsScope, _opName, "queue_mismatch", 1) + messageMetrics.NamedCounter(_opName, "queue_mismatch", 1) return fmt.Errorf("payload queue %q does not match queue %q of request %s", sig.GetQueueName(), request.Queue, request.ID) } @@ -167,7 +168,7 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) er } if effective.IsTerminal() { - if err := c.finishRequest(ctx, store, &request, effective); err != nil { + if err := c.finishRequest(ctx, messageMetrics, store, &request, effective); err != nil { return err } 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 // the request non-terminal, so redelivery re-runs both steps and decrements again // — transiently over-admitting by one until releaseBuildSlot's zero clamp // reconverges, which is the failure mode this pipeline prefers. -func (c *Controller) finishRequest(ctx context.Context, store storage.Storage, request *entity.Request, status entity.BuildStatus) error { +func (c *Controller) finishRequest(ctx context.Context, messageMetrics metrics.Scope, store storage.Storage, request *entity.Request, status entity.BuildStatus) error { if request.State.HasBuildOutcome() { return nil } - if err := c.releaseBuildSlot(ctx, store, request.Queue); err != nil { - metrics.NamedCounter(c.metricsScope, _opName, "storage_errors", 1) + if err := c.releaseBuildSlot(ctx, messageMetrics, store, request.Queue); err != nil { + messageMetrics.NamedCounter(_opName, "storage_errors", 1) return err } - if err := c.markOutcome(ctx, store, request, outcomeState(status)); err != nil { - metrics.NamedCounter(c.metricsScope, _opName, "storage_errors", 1) + if err := c.markOutcome(ctx, messageMetrics, store, request, outcomeState(status)); err != nil { + messageMetrics.NamedCounter(_opName, "storage_errors", 1) return err } return nil @@ -243,7 +244,7 @@ func outcomeState(status entity.BuildStatus) entity.RequestState { // conflicts. First writer wins: once any outcome is recorded a later caller leaves it // alone, so duplicate builds for one request (which build.md accepts) cannot flip the // verdict back and forth. -func (c *Controller) markOutcome(ctx context.Context, store storage.Storage, request *entity.Request, state entity.RequestState) error { +func (c *Controller) markOutcome(ctx context.Context, messageMetrics metrics.Scope, store storage.Storage, request *entity.Request, state entity.RequestState) error { reqStore := store.GetRequestStore() for { @@ -267,7 +268,7 @@ func (c *Controller) markOutcome(ctx context.Context, store storage.Storage, req } updated.Version = newVersion *request = updated - metrics.NamedCounter(c.metricsScope, _opName, "outcomes", 1, + messageMetrics.NamedCounter(_opName, "outcomes", 1, metrics.NewTag("state", string(state)), ) return nil @@ -279,7 +280,7 @@ func (c *Controller) markOutcome(ctx context.Context, store storage.Storage, req // (preserving concurrent updates), clamps at zero, and retries on version conflicts. // Unlike process's unwind-path release this is not best-effort: the caller must not // mark the request terminal if the slot was not freed, so a hard failure is returned. -func (c *Controller) releaseBuildSlot(ctx context.Context, store storage.Storage, queueName string) error { +func (c *Controller) releaseBuildSlot(ctx context.Context, messageMetrics metrics.Scope, store storage.Storage, queueName string) error { queueStore := store.GetQueueStore() for { @@ -300,7 +301,7 @@ func (c *Controller) releaseBuildSlot(ctx context.Context, store storage.Storage } return fmt.Errorf("failed to release build slot for queue %s: %w", queueName, err) } - metrics.NamedCounter(c.metricsScope, _opName, "slot_released", 1) + messageMetrics.NamedCounter(_opName, "slot_released", 1) return nil } } diff --git a/stovepipe/controller/buildsignal/buildsignal_test.go b/stovepipe/controller/buildsignal/buildsignal_test.go index 022ff927a..4368088c4 100644 --- a/stovepipe/controller/buildsignal/buildsignal_test.go +++ b/stovepipe/controller/buildsignal/buildsignal_test.go @@ -52,6 +52,7 @@ type buildsignalMocks struct { runnerFactory *buildrunnermock.MockFactory runner *buildrunnermock.MockBuildRunner publisher *mqmock.MockPublisher + metricsScope tally.TestScope } // staticStorageFactory resolves every queue to one fixed store aggregate. @@ -63,6 +64,7 @@ func (f staticStorageFactory) For(storage.Config) (storage.Storage, error) { ret func newController(t *testing.T, ctrl *gomock.Controller) (*Controller, buildsignalMocks) { t.Helper() + scope := tally.NewTestScope("test", nil) m := buildsignalMocks{ reqStore: storagemock.NewMockRequestStore(ctrl), buildStore: storagemock.NewMockBuildStore(ctrl), @@ -70,6 +72,7 @@ func newController(t *testing.T, ctrl *gomock.Controller) (*Controller, buildsig runnerFactory: buildrunnermock.NewMockFactory(ctrl), runner: buildrunnermock.NewMockBuildRunner(ctrl), publisher: mqmock.NewMockPublisher(ctrl), + metricsScope: scope, } store := storagemock.NewMockStorage(ctrl) @@ -86,10 +89,21 @@ func newController(t *testing.T, ctrl *gomock.Controller) (*Controller, buildsig }) require.NoError(t, err) - c := NewController(zap.NewNop().Sugar(), tally.NewTestScope("test", nil), staticStorageFactory{store: store}, m.runnerFactory, registry, stovepipemq.TopicKeyBuildSignal, "stovepipe-buildsignal") + c := NewController(zap.NewNop().Sugar(), scope, staticStorageFactory{store: store}, m.runnerFactory, registry, stovepipemq.TopicKeyBuildSignal, "stovepipe-buildsignal") return c, m } +func TestProcessTagsMetricsWithQueue(t *testing.T) { + ctrl := gomock.NewController(t) + c, m := newController(t, ctrl) + m.buildStore.EXPECT().Get(gomock.Any(), testBuildID).Return(entity.Build{}, assert.AnError) + + require.Error(t, c.Process(context.Background(), delivery(t, ctrl, buildSignalPayload(t, testBuildID)))) + counter, ok := m.metricsScope.Snapshot().Counters()["test.buildsignal_controller.buildsignal.storage_errors+queue=monorepo/main"] + require.True(t, ok) + assert.EqualValues(t, 1, counter.Value()) +} + func delivery(t *testing.T, ctrl *gomock.Controller, payload []byte) *consumermock.MockDelivery { t.Helper() d := consumermock.NewMockDelivery(ctrl) @@ -100,7 +114,7 @@ func delivery(t *testing.T, ctrl *gomock.Controller, payload []byte) *consumermo func buildSignalPayload(t *testing.T, id string) []byte { t.Helper() - b, err := stovepipemq.Marshal(&stovepipemq.BuildSignal{Id: id}) + b, err := stovepipemq.Marshal(&stovepipemq.BuildSignal{Id: id, QueueName: testQueue}) require.NoError(t, err) return b } diff --git a/stovepipe/controller/dlq/build.go b/stovepipe/controller/dlq/build.go index 3afed6815..504a7f1ae 100644 --- a/stovepipe/controller/dlq/build.go +++ b/stovepipe/controller/dlq/build.go @@ -29,11 +29,11 @@ import ( // buildController reconciles build-stage dead letters by failing the request // whose build could not be triggered, persisted, or handed to the poll loop. type buildController struct { - logger *zap.SugaredLogger - metricsScope tally.Scope - stores storage.Factory - topicKey consumer.TopicKey - consumerGroup string + logger *zap.SugaredLogger + metricsFactory metrics.Factory + stores storage.Factory + topicKey consumer.TopicKey + consumerGroup string } var _ consumer.Controller = (*buildController)(nil) @@ -50,11 +50,11 @@ func NewDLQBuildController( ) consumer.Controller { name := string(topicKey) + "_controller" return &buildController{ - logger: logger.Named(name), - metricsScope: scope.SubScope(name), - stores: stores, - topicKey: topicKey, - consumerGroup: consumerGroup, + logger: logger.Named(name), + metricsFactory: metrics.NewFactory(scope.SubScope(name)), + stores: stores, + topicKey: topicKey, + consumerGroup: consumerGroup, } } @@ -62,17 +62,18 @@ func NewDLQBuildController( func (c *buildController) Process(ctx context.Context, delivery consumer.Delivery) error { buildRequest := &stovepipemq.BuildRequest{} if err := stovepipemq.Unmarshal(delivery.Message().Payload, buildRequest); err != nil { - metrics.NamedCounter(c.metricsScope, _buildOpName, "deserialize_errors", 1) + c.metricsFactory.Base().NamedCounter(_buildOpName, "deserialize_errors", 1) return fmt.Errorf("failed to decode dlq payload: %w", err) } + messageMetrics := c.metricsFactory.For(metrics.Config{QueueName: buildRequest.GetQueueName()}) if buildRequest.Id == "" { - metrics.NamedCounter(c.metricsScope, _buildOpName, "empty_id_errors", 1) + messageMetrics.NamedCounter(_buildOpName, "empty_id_errors", 1) return fmt.Errorf("build dlq payload decoded to empty request id") } store, err := c.stores.For(storage.Config{QueueName: buildRequest.GetQueueName()}) if err != nil { - metrics.NamedCounter(c.metricsScope, _buildOpName, "storage_resolve_errors", 1) + messageMetrics.NamedCounter(_buildOpName, "storage_resolve_errors", 1) return fmt.Errorf("failed to resolve storage for queue %q: %w", buildRequest.GetQueueName(), err) } @@ -86,10 +87,10 @@ func (c *buildController) Process(ctx context.Context, delivery consumer.Deliver ) if err := failRequest(ctx, store, c.logger, buildRequest.Id); err != nil { - metrics.NamedCounter(c.metricsScope, _buildOpName, "reconcile_errors", 1) + messageMetrics.NamedCounter(_buildOpName, "reconcile_errors", 1) return err } - metrics.NamedCounter(c.metricsScope, _buildOpName, "reconciled", 1) + messageMetrics.NamedCounter(_buildOpName, "reconciled", 1) return nil } diff --git a/stovepipe/controller/dlq/build_test.go b/stovepipe/controller/dlq/build_test.go index a725b521d..21a9670a2 100644 --- a/stovepipe/controller/dlq/build_test.go +++ b/stovepipe/controller/dlq/build_test.go @@ -32,15 +32,17 @@ import ( func newBuildController(t *testing.T, ctrl *gomock.Controller) (consumer.Controller, dlqMocks) { t.Helper() + scope := tally.NewTestScope("test", nil) m := dlqMocks{ - reqStore: storagemock.NewMockRequestStore(ctrl), - queueStore: storagemock.NewMockQueueStore(ctrl), + reqStore: storagemock.NewMockRequestStore(ctrl), + queueStore: storagemock.NewMockQueueStore(ctrl), + metricsScope: scope, } store := storagemock.NewMockStorage(ctrl) store.EXPECT().GetRequestStore().Return(m.reqStore).AnyTimes() store.EXPECT().GetQueueStore().Return(m.queueStore).AnyTimes() - c := NewDLQBuildController(zap.NewNop().Sugar(), tally.NewTestScope("test", nil), staticStorageFactory{store: store}, TopicKey(stovepipemq.TopicKeyBuild), "stovepipe-build-dlq") + c := NewDLQBuildController(zap.NewNop().Sugar(), scope, staticStorageFactory{store: store}, TopicKey(stovepipemq.TopicKeyBuild), "stovepipe-build-dlq") return c, m } @@ -53,10 +55,11 @@ func buildPayload(t *testing.T, id string) []byte { func TestBuildProcess(t *testing.T) { tests := []struct { - name string - payload []byte - setup func(m dlqMocks) - wantErr bool + name string + payload []byte + setup func(m dlqMocks) + wantErr bool + wantMetric string }{ { name: "processing request is failed", @@ -68,9 +71,15 @@ func TestBuildProcess(t *testing.T) { updated.State = entity.RequestStateFailed m.reqStore.EXPECT().Update(gomock.Any(), updated, int32(2), int32(3)).Return(nil) }, + wantMetric: "test.build_dlq_controller.build_dlq.reconciled+queue=monorepo/main", }, {name: "malformed payload is returned", payload: []byte("not-a-proto"), wantErr: true}, - {name: "empty request id is returned", payload: buildPayload(t, ""), wantErr: true}, + { + name: "empty request id is returned", + payload: buildPayload(t, ""), + wantErr: true, + wantMetric: "test.build_dlq_controller.build_dlq.empty_id_errors+queue=monorepo/main", + }, } for _, tt := range tests { @@ -86,6 +95,11 @@ func TestBuildProcess(t *testing.T) { payload = buildPayload(t, testID) } err := controller.Process(context.Background(), delivery(t, ctrl, payload)) + if tt.wantMetric != "" { + counter, ok := mocks.metricsScope.Snapshot().Counters()[tt.wantMetric] + require.True(t, ok) + assert.EqualValues(t, 1, counter.Value()) + } if tt.wantErr { require.Error(t, err) return diff --git a/stovepipe/controller/dlq/buildsignal.go b/stovepipe/controller/dlq/buildsignal.go index c68d2f2d0..4cb43ee57 100644 --- a/stovepipe/controller/dlq/buildsignal.go +++ b/stovepipe/controller/dlq/buildsignal.go @@ -50,11 +50,11 @@ const _buildSignalOpName = "buildsignal_dlq" // not Build.Status, and writing a terminal status here would claim we saw an // outcome we never saw. type buildSignalController struct { - logger *zap.SugaredLogger - metricsScope tally.Scope - stores storage.Factory - topicKey consumer.TopicKey - consumerGroup string + logger *zap.SugaredLogger + metricsFactory metrics.Factory + stores storage.Factory + topicKey consumer.TopicKey + consumerGroup string } // Verify buildSignalController implements consumer.Controller at compile time. @@ -72,11 +72,11 @@ func NewDLQBuildSignalController( ) consumer.Controller { name := string(topicKey) + "_controller" return &buildSignalController{ - logger: logger.Named(name), - metricsScope: scope.SubScope(name), - stores: stores, - topicKey: topicKey, - consumerGroup: consumerGroup, + logger: logger.Named(name), + metricsFactory: metrics.NewFactory(scope.SubScope(name)), + stores: stores, + topicKey: topicKey, + consumerGroup: consumerGroup, } } @@ -89,21 +89,22 @@ func (c *buildSignalController) Process(ctx context.Context, delivery consumer.D sig := &stovepipemq.BuildSignal{} if err := stovepipemq.Unmarshal(msg.Payload, sig); err != nil { - metrics.NamedCounter(c.metricsScope, _buildSignalOpName, "deserialize_errors", 1) + c.metricsFactory.Base().NamedCounter(_buildSignalOpName, "deserialize_errors", 1) // Retried rather than acked, for the same deployment-skew reason the // process reconciler gives: a newer producer's payload decodes fine once // the rollout finishes, and acking here would skip the slot release // without saying so. return fmt.Errorf("failed to decode dlq payload: %w", err) } + messageMetrics := c.metricsFactory.For(metrics.Config{QueueName: sig.GetQueueName()}) if sig.Id == "" { - metrics.NamedCounter(c.metricsScope, _buildSignalOpName, "empty_id_errors", 1) + messageMetrics.NamedCounter(_buildSignalOpName, "empty_id_errors", 1) return fmt.Errorf("dlq payload decoded to empty build id") } store, err := c.stores.For(storage.Config{QueueName: sig.GetQueueName()}) if err != nil { - metrics.NamedCounter(c.metricsScope, _buildSignalOpName, "storage_resolve_errors", 1) + messageMetrics.NamedCounter(_buildSignalOpName, "storage_resolve_errors", 1) // Non-retryable: a missing or unresolvable queue is a malformed message. return fmt.Errorf("failed to resolve storage for queue %q: %w", sig.GetQueueName(), err) } @@ -124,10 +125,10 @@ func (c *buildSignalController) Process(ctx context.Context, delivery consumer.D // Create. There is no request to recover from this payload; the build // stage's own DLQ handles the request that triggered it. c.logger.Warnw("dlq reconcile: build not found, skipping", "build_id", sig.Id) - metrics.NamedCounter(c.metricsScope, _buildSignalOpName, "build_not_found", 1) + messageMetrics.NamedCounter(_buildSignalOpName, "build_not_found", 1) return nil } - metrics.NamedCounter(c.metricsScope, _buildSignalOpName, "build_store_errors", 1) + messageMetrics.NamedCounter(_buildSignalOpName, "build_store_errors", 1) return fmt.Errorf("failed to get build %s: %w", sig.Id, err) } @@ -135,7 +136,7 @@ func (c *buildSignalController) Process(ctx context.Context, delivery consumer.D // Defensive: a build with no request has nothing to reconcile and no slot // to release. Ack it so the DLQ does not grow forever. c.logger.Errorw("dlq reconcile: build has empty request id, skipping", "build_id", sig.Id) - metrics.NamedCounter(c.metricsScope, _buildSignalOpName, "build_missing_request", 1) + messageMetrics.NamedCounter(_buildSignalOpName, "build_missing_request", 1) return nil } @@ -144,11 +145,11 @@ func (c *buildSignalController) Process(ctx context.Context, delivery consumer.D // triggers only once process has written the strategy, which lands in the same CAS // as accepted→processing, and processing exits only to a terminal outcome. if err := failRequest(ctx, store, c.logger, build.RequestID); err != nil { - metrics.NamedCounter(c.metricsScope, _buildSignalOpName, "reconcile_errors", 1) + messageMetrics.NamedCounter(_buildSignalOpName, "reconcile_errors", 1) return err } - metrics.NamedCounter(c.metricsScope, _buildSignalOpName, "reconciled", 1) + messageMetrics.NamedCounter(_buildSignalOpName, "reconciled", 1) return nil } diff --git a/stovepipe/controller/dlq/buildsignal_test.go b/stovepipe/controller/dlq/buildsignal_test.go index dbe88812f..8789522ad 100644 --- a/stovepipe/controller/dlq/buildsignal_test.go +++ b/stovepipe/controller/dlq/buildsignal_test.go @@ -33,18 +33,21 @@ import ( const testBuildID = "go-code-on-odin-submitqueue/builds/2867068" type buildSignalDLQMocks struct { - reqStore *storagemock.MockRequestStore - queueStore *storagemock.MockQueueStore - buildStore *storagemock.MockBuildStore + reqStore *storagemock.MockRequestStore + queueStore *storagemock.MockQueueStore + buildStore *storagemock.MockBuildStore + metricsScope tally.TestScope } func newBuildSignalController(t *testing.T, ctrl *gomock.Controller) (consumer.Controller, buildSignalDLQMocks) { t.Helper() + scope := tally.NewTestScope("test", nil) m := buildSignalDLQMocks{ - reqStore: storagemock.NewMockRequestStore(ctrl), - queueStore: storagemock.NewMockQueueStore(ctrl), - buildStore: storagemock.NewMockBuildStore(ctrl), + reqStore: storagemock.NewMockRequestStore(ctrl), + queueStore: storagemock.NewMockQueueStore(ctrl), + buildStore: storagemock.NewMockBuildStore(ctrl), + metricsScope: scope, } store := storagemock.NewMockStorage(ctrl) @@ -54,7 +57,7 @@ func newBuildSignalController(t *testing.T, ctrl *gomock.Controller) (consumer.C c := NewDLQBuildSignalController( zap.NewNop().Sugar(), - tally.NewTestScope("test", nil), + scope, staticStorageFactory{store: store}, TopicKey(stovepipemq.TopicKeyBuildSignal), "stovepipe-buildsignal-dlq", @@ -80,10 +83,11 @@ func build() entity.Build { func TestBuildSignalProcess(t *testing.T) { tests := []struct { - name string - payload []byte - setup func(m buildSignalDLQMocks) - wantErr bool + name string + payload []byte + setup func(m buildSignalDLQMocks) + wantErr bool + wantMetric string }{ { // The case this reconciler exists for: a poll message that @@ -147,10 +151,11 @@ func TestBuildSignalProcess(t *testing.T) { wantErr: true, }, { - name: "empty build id is returned as an error", - payload: buildSignalPayload(t, ""), - setup: func(buildSignalDLQMocks) {}, - wantErr: true, + name: "empty build id is returned as an error", + payload: buildSignalPayload(t, ""), + setup: func(buildSignalDLQMocks) {}, + wantErr: true, + wantMetric: "test.buildsignal_dlq_controller.buildsignal_dlq.empty_id_errors+queue=monorepo/main", }, } @@ -166,6 +171,11 @@ func TestBuildSignalProcess(t *testing.T) { } err := c.Process(context.Background(), delivery(t, ctrl, payload)) + if tt.wantMetric != "" { + counter, ok := m.metricsScope.Snapshot().Counters()[tt.wantMetric] + require.True(t, ok) + assert.EqualValues(t, 1, counter.Value()) + } if tt.wantErr { require.Error(t, err) diff --git a/stovepipe/controller/dlq/dlq_test.go b/stovepipe/controller/dlq/dlq_test.go index ed3962ea6..84e23a363 100644 --- a/stovepipe/controller/dlq/dlq_test.go +++ b/stovepipe/controller/dlq/dlq_test.go @@ -38,8 +38,9 @@ const ( ) type dlqMocks struct { - reqStore *storagemock.MockRequestStore - queueStore *storagemock.MockQueueStore + reqStore *storagemock.MockRequestStore + queueStore *storagemock.MockQueueStore + metricsScope tally.TestScope } // staticStorageFactory resolves every queue to one fixed store aggregate. @@ -51,16 +52,18 @@ func (f staticStorageFactory) For(storage.Config) (storage.Storage, error) { ret func newController(t *testing.T, ctrl *gomock.Controller) (consumer.Controller, dlqMocks) { t.Helper() + scope := tally.NewTestScope("test", nil) m := dlqMocks{ - reqStore: storagemock.NewMockRequestStore(ctrl), - queueStore: storagemock.NewMockQueueStore(ctrl), + reqStore: storagemock.NewMockRequestStore(ctrl), + queueStore: storagemock.NewMockQueueStore(ctrl), + metricsScope: scope, } store := storagemock.NewMockStorage(ctrl) store.EXPECT().GetRequestStore().Return(m.reqStore).AnyTimes() store.EXPECT().GetQueueStore().Return(m.queueStore).AnyTimes() - c := NewDLQRequestController(zap.NewNop().Sugar(), tally.NewTestScope("test", nil), staticStorageFactory{store: store}, TopicKey(stovepipemq.TopicKeyProcess), "stovepipe-process-dlq") + c := NewDLQRequestController(zap.NewNop().Sugar(), scope, staticStorageFactory{store: store}, TopicKey(stovepipemq.TopicKeyProcess), "stovepipe-process-dlq") return c, m } @@ -79,7 +82,7 @@ func delivery(t *testing.T, ctrl *gomock.Controller, payload []byte) consumer.De func processPayload(t *testing.T, id string) []byte { t.Helper() - b, err := stovepipemq.Marshal(&stovepipemq.ProcessRequest{Id: id}) + b, err := stovepipemq.Marshal(&stovepipemq.ProcessRequest{Id: id, QueueName: testQueue}) require.NoError(t, err) return b } @@ -95,10 +98,11 @@ func requestWithState(state entity.RequestState) entity.Request { func TestProcess(t *testing.T) { tests := []struct { - name string - payload []byte - setup func(m dlqMocks) - wantErr bool + name string + payload []byte + setup func(m dlqMocks) + wantErr bool + wantMetric string }{ { // No queue expectations: an accepted request never claimed a slot, @@ -194,10 +198,11 @@ func TestProcess(t *testing.T) { wantErr: true, }, { - name: "empty request id is not retryable", - payload: processPayload(t, ""), - setup: func(m dlqMocks) {}, - wantErr: true, + name: "empty request id is not retryable", + payload: processPayload(t, ""), + setup: func(m dlqMocks) {}, + wantErr: true, + wantMetric: "test.process_dlq_controller.process_dlq.empty_id_errors+queue=monorepo/main", }, } @@ -213,6 +218,11 @@ func TestProcess(t *testing.T) { } err := c.Process(context.Background(), delivery(t, ctrl, payload)) + if tt.wantMetric != "" { + counter, ok := m.metricsScope.Snapshot().Counters()[tt.wantMetric] + require.True(t, ok) + assert.EqualValues(t, 1, counter.Value()) + } if tt.wantErr { require.Error(t, err) diff --git a/stovepipe/controller/dlq/request.go b/stovepipe/controller/dlq/request.go index bdf8f0169..d3003dc12 100644 --- a/stovepipe/controller/dlq/request.go +++ b/stovepipe/controller/dlq/request.go @@ -31,11 +31,11 @@ import ( // the same ProcessRequest payload the primary process controller consumes, then drives // the referenced request to a terminal failed state via failRequest. type requestController struct { - logger *zap.SugaredLogger - metricsScope tally.Scope - stores storage.Factory - topicKey consumer.TopicKey - consumerGroup string + logger *zap.SugaredLogger + metricsFactory metrics.Factory + stores storage.Factory + topicKey consumer.TopicKey + consumerGroup string } // Verify requestController implements consumer.Controller at compile time. @@ -55,11 +55,11 @@ func NewDLQRequestController( ) consumer.Controller { name := string(topicKey) + "_controller" return &requestController{ - logger: logger.Named(name), - metricsScope: scope.SubScope(name), - stores: stores, - topicKey: topicKey, - consumerGroup: consumerGroup, + logger: logger.Named(name), + metricsFactory: metrics.NewFactory(scope.SubScope(name)), + stores: stores, + topicKey: topicKey, + consumerGroup: consumerGroup, } } @@ -72,7 +72,7 @@ func (c *requestController) Process(ctx context.Context, delivery consumer.Deliv pr := &stovepipemq.ProcessRequest{} if err := stovepipemq.Unmarshal(msg.Payload, pr); err != nil { - metrics.NamedCounter(c.metricsScope, _opName, "deserialize_errors", 1) + c.metricsFactory.Base().NamedCounter(_opName, "deserialize_errors", 1) // Decoding the same bytes normally fails deterministically, but this error is // still retried: the DLQ consumer's AlwaysRetryableProcessor (see Process doc) // classifies every error as retryable. That is deliberate — the recoverable @@ -84,14 +84,15 @@ func (c *requestController) Process(ctx context.Context, delivery consumer.Deliv // non-terminal. return fmt.Errorf("failed to decode dlq payload: %w", err) } + messageMetrics := c.metricsFactory.For(metrics.Config{QueueName: pr.GetQueueName()}) if pr.Id == "" { - metrics.NamedCounter(c.metricsScope, _opName, "empty_id_errors", 1) + messageMetrics.NamedCounter(_opName, "empty_id_errors", 1) return fmt.Errorf("dlq payload decoded to empty request id") } store, err := c.stores.For(storage.Config{QueueName: pr.GetQueueName()}) if err != nil { - metrics.NamedCounter(c.metricsScope, _opName, "storage_resolve_errors", 1) + messageMetrics.NamedCounter(_opName, "storage_resolve_errors", 1) // Non-retryable: a missing or unresolvable queue is a malformed message. return fmt.Errorf("failed to resolve storage for queue %q: %w", pr.GetQueueName(), err) } @@ -106,7 +107,7 @@ func (c *requestController) Process(ctx context.Context, delivery consumer.Deliv ) if err := failRequest(ctx, store, c.logger, pr.Id); err != nil { - metrics.NamedCounter(c.metricsScope, _opName, "reconcile_errors", 1) + messageMetrics.NamedCounter(_opName, "reconcile_errors", 1) return err } diff --git a/stovepipe/controller/process/BUILD.bazel b/stovepipe/controller/process/BUILD.bazel index 209f87f07..6db59c697 100644 --- a/stovepipe/controller/process/BUILD.bazel +++ b/stovepipe/controller/process/BUILD.bazel @@ -31,6 +31,7 @@ go_test( "//platform/consumer/mock:go_default_library", "//platform/errs:go_default_library", "//platform/extension/messagequeue/mock:go_default_library", + "//platform/metrics:go_default_library", "//stovepipe/core/messagequeue:go_default_library", "//stovepipe/entity:go_default_library", "//stovepipe/extension/queueconfig/default:go_default_library", diff --git a/stovepipe/controller/process/process.go b/stovepipe/controller/process/process.go index 0e71fd1c9..6f854bba0 100644 --- a/stovepipe/controller/process/process.go +++ b/stovepipe/controller/process/process.go @@ -41,14 +41,14 @@ import ( // referenced Request from storage, coalesces older heads, and admits the latest when // a slot is open. Implements consumer.Controller. type Controller struct { - logger *zap.SugaredLogger - metricsScope tally.Scope - stores storage.Factory - queueConfigs queueconfig.Store - sourceControl sourcecontrol.Factory - registry consumer.TopicRegistry - topicKey consumer.TopicKey - consumerGroup string + logger *zap.SugaredLogger + metricsFactory metrics.Factory + stores storage.Factory + queueConfigs queueconfig.Store + sourceControl sourcecontrol.Factory + registry consumer.TopicRegistry + topicKey consumer.TopicKey + consumerGroup string } // Verify Controller implements consumer.Controller interface at compile time. @@ -69,14 +69,14 @@ func NewController( consumerGroup string, ) *Controller { return &Controller{ - logger: logger.Named("process_controller"), - metricsScope: scope.SubScope("process_controller"), - stores: stores, - queueConfigs: queueConfigs, - sourceControl: sourceControl, - registry: registry, - topicKey: topicKey, - consumerGroup: consumerGroup, + logger: logger.Named("process_controller"), + metricsFactory: metrics.NewFactory(scope.SubScope("process_controller")), + stores: stores, + queueConfigs: queueConfigs, + sourceControl: sourceControl, + registry: registry, + topicKey: topicKey, + consumerGroup: consumerGroup, } } @@ -87,35 +87,36 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) er pr := &stovepipemq.ProcessRequest{} if err := stovepipemq.Unmarshal(msg.Payload, pr); err != nil { - metrics.NamedCounter(c.metricsScope, _opName, "deserialize_errors", 1) + c.metricsFactory.Base().NamedCounter(_opName, "deserialize_errors", 1) // Non-retryable: a malformed message will never succeed regardless of retries. return fmt.Errorf("failed to deserialize process request: %w", err) } + messageMetrics := c.metricsFactory.For(metrics.Config{QueueName: pr.GetQueueName()}) store, err := c.stores.For(storage.Config{QueueName: pr.GetQueueName()}) if err != nil { - metrics.NamedCounter(c.metricsScope, _opName, "storage_resolve_errors", 1) + messageMetrics.NamedCounter(_opName, "storage_resolve_errors", 1) // Non-retryable: a missing or unresolvable queue is a malformed message. return fmt.Errorf("failed to resolve storage for queue %q: %w", pr.GetQueueName(), err) } request, err := c.loadRequest(ctx, store, pr.Id) if err != nil { - metrics.NamedCounter(c.metricsScope, _opName, "storage_errors", 1) + messageMetrics.NamedCounter(_opName, "storage_errors", 1) return err } // The payload's queue must match the request's authoritative queue; a // mismatch is a malformed message. Non-retryable — reject to the DLQ. if pr.GetQueueName() != "" && pr.GetQueueName() != request.Queue { - metrics.NamedCounter(c.metricsScope, _opName, "queue_mismatch", 1) + messageMetrics.NamedCounter(_opName, "queue_mismatch", 1) return fmt.Errorf("payload queue %q does not match queue %q of request %s", pr.GetQueueName(), request.Queue, request.ID) } switch request.State { case entity.RequestStateProcessing: if err := c.publishBuild(ctx, request.ID, request.Queue); err != nil { - metrics.NamedCounter(c.metricsScope, _opName, "publish_errors", 1) + messageMetrics.NamedCounter(_opName, "publish_errors", 1) return fmt.Errorf("failed to publish request %s to build: %w", request.ID, err) } return nil @@ -124,7 +125,7 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) er // A stale redelivery has nothing left to do. return nil case entity.RequestStateAccepted: - return c.processAccepted(ctx, store, delivery, request) + return c.processAccepted(ctx, messageMetrics, store, delivery, request) default: c.logger.Warnw("ignored request in unexpected state", "request_id", request.ID, @@ -138,11 +139,11 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) er // processAccepted coalesces older heads against queue.latest_request_id, then admits // the latest head when a build slot is available. The delivery is threaded down so a // closed gate can hold it. -func (c *Controller) processAccepted(ctx context.Context, store storage.Storage, delivery consumer.Delivery, request entity.Request) error { +func (c *Controller) processAccepted(ctx context.Context, messageMetrics metrics.Scope, store storage.Storage, delivery consumer.Delivery, request entity.Request) error { queueRow, err := c.loadQueue(ctx, store, request.Queue) if err != nil { if !errs.IsRetryable(err) { - metrics.NamedCounter(c.metricsScope, _opName, "storage_errors", 1) + messageMetrics.NamedCounter(_opName, "storage_errors", 1) } return err } @@ -156,7 +157,7 @@ func (c *Controller) processAccepted(ctx context.Context, store storage.Storage, return nil } - superseded, err := c.coalesce(ctx, store, request, queueRow.LatestRequestID) + superseded, err := c.coalesce(ctx, messageMetrics, store, request, queueRow.LatestRequestID) if err != nil || superseded { return err } @@ -168,13 +169,13 @@ func (c *Controller) processAccepted(ctx context.Context, store storage.Storage, return fmt.Errorf("failed to load queue config for %s: %w", request.Queue, err) } - return c.admitLatestHead(ctx, store, delivery, request, queueRow, cfg) + return c.admitLatestHead(ctx, messageMetrics, store, delivery, request, queueRow, cfg) } // coalesce supersedes request when a newer head exists (RFC process step 5), returning // true so the caller acks. It returns false when request is still the latest head and // should proceed to the gate. Superseding consumes no build slot. -func (c *Controller) coalesce(ctx context.Context, store storage.Storage, request entity.Request, latestRequestID string) (bool, error) { +func (c *Controller) coalesce(ctx context.Context, messageMetrics metrics.Scope, store storage.Storage, request entity.Request, latestRequestID string) (bool, error) { cmp, err := entity.CompareRequestID(request.Queue, request.ID, latestRequestID) if err != nil { return false, fmt.Errorf("failed to compare request ids for queue %s: %w", request.Queue, err) @@ -183,10 +184,10 @@ func (c *Controller) coalesce(ctx context.Context, store storage.Storage, reques return false, nil } if err := c.supersedeRequest(ctx, store, request); err != nil { - metrics.NamedCounter(c.metricsScope, _opName, "storage_errors", 1) + messageMetrics.NamedCounter(_opName, "storage_errors", 1) return false, err } - metrics.NamedCounter(c.metricsScope, _opName, "superseded", 1) + messageMetrics.NamedCounter(_opName, "superseded", 1) c.logger.Infow("superseded request for newer head", "request_id", request.ID, "queue", request.Queue, @@ -199,7 +200,7 @@ func (c *Controller) coalesce(ctx context.Context, store storage.Storage, reques // slot, mark the request processing, and publish it to build. Every queue-row reload // re-runs coalesce-then-gate, so a slot is never spent on a now-stale head; a closed gate // defers by holding the delivery (redeliver after the gate wait delay) rather than failing. -func (c *Controller) admitLatestHead(ctx context.Context, store storage.Storage, delivery consumer.Delivery, request entity.Request, queueRow entity.Queue, cfg entity.QueueConfig) error { +func (c *Controller) admitLatestHead(ctx context.Context, messageMetrics metrics.Scope, store storage.Storage, delivery consumer.Delivery, request entity.Request, queueRow entity.Queue, cfg entity.QueueConfig) error { var sc sourcecontrol.SourceControl var strategy entity.BuildStrategy var baseURI string @@ -207,20 +208,20 @@ func (c *Controller) admitLatestHead(ctx context.Context, store storage.Storage, for { if queueRow.InFlightCount >= cfg.MaxConcurrent { - return c.holdForBuildSlot(delivery, request, queueRow.InFlightCount, cfg.GateWaitDelayMs) + return c.holdForBuildSlot(messageMetrics, delivery, request, queueRow.InFlightCount, cfg.GateWaitDelayMs) } if queueRow.LastGreenURI != "" && sc == nil { sc, err = c.sourceControl.For(sourcecontrol.Config{QueueName: request.Queue}) if err != nil { - metrics.NamedCounter(c.metricsScope, _opName, "source_control_errors", 1, + messageMetrics.NamedCounter(_opName, "source_control_errors", 1, metrics.NewTag("stage", "resolve"), ) return fmt.Errorf("failed to resolve source control for queue %s: %w", request.Queue, err) } } - strategy, baseURI, err = c.deriveBuildStrategy(ctx, sc, queueRow, request) + strategy, baseURI, err = c.deriveBuildStrategy(ctx, messageMetrics, sc, queueRow, request) if err != nil { return err } @@ -234,7 +235,7 @@ func (c *Controller) admitLatestHead(ctx context.Context, store storage.Storage, } // claimBuildSlot reloaded queueRow. Re-coalesce: supersede if a newer head arrived, // otherwise loop to re-check the gate. - superseded, err := c.coalesce(ctx, store, request, queueRow.LatestRequestID) + superseded, err := c.coalesce(ctx, messageMetrics, store, request, queueRow.LatestRequestID) if err != nil || superseded { return err } @@ -244,21 +245,21 @@ func (c *Controller) admitLatestHead(ctx context.Context, store storage.Storage, if err != nil { // Slot claimed but never admitted: release best-effort so the slot isn't leaked // (a redelivery would find the gate closed by its own claim and nothing decrements it). - c.releaseBuildSlot(ctx, store, request.Queue) + c.releaseBuildSlot(ctx, messageMetrics, store, request.Queue) return err } if !transitioned { // Lost the admit race: another delivery advanced this request. Release and skip. - c.releaseBuildSlot(ctx, store, request.Queue) + c.releaseBuildSlot(ctx, messageMetrics, store, request.Queue) return nil } if err := c.publishBuild(ctx, request.ID, request.Queue); err != nil { - metrics.NamedCounter(c.metricsScope, _opName, "publish_errors", 1) + messageMetrics.NamedCounter(_opName, "publish_errors", 1) return fmt.Errorf("failed to publish request %s to build: %w", request.ID, err) } - metrics.NamedCounter(c.metricsScope, _opName, "admitted", 1, + messageMetrics.NamedCounter(_opName, "admitted", 1, metrics.NewTag("strategy", string(request.BuildStrategy)), ) c.logger.Infow("admitted request to build", @@ -273,7 +274,7 @@ func (c *Controller) admitLatestHead(ctx context.Context, store storage.Storage, // deriveBuildStrategy chooses the validation scope and baseline from the queue's last-known-good commit. // The caller resolves source control once and persists the returned values only after successfully claiming a build slot. -func (c *Controller) deriveBuildStrategy(ctx context.Context, sc sourcecontrol.SourceControl, queueRow entity.Queue, request entity.Request) (strategy entity.BuildStrategy, baseURI string, err error) { +func (c *Controller) deriveBuildStrategy(ctx context.Context, messageMetrics metrics.Scope, sc sourcecontrol.SourceControl, queueRow entity.Queue, request entity.Request) (strategy entity.BuildStrategy, baseURI string, err error) { if queueRow.LastGreenURI == "" { return entity.BuildStrategyFull, "", nil } @@ -281,7 +282,7 @@ func (c *Controller) deriveBuildStrategy(ctx context.Context, sc sourcecontrol.S isAncestor, err := sc.IsAncestor(ctx, queueRow.LastGreenURI, request.URI) if err != nil { if sourcecontrol.IsNotFound(err) { - metrics.NamedCounter(c.metricsScope, _opName, "strategy_fallbacks", 1, + messageMetrics.NamedCounter(_opName, "strategy_fallbacks", 1, metrics.NewTag("reason", "unknown_ancestry"), ) c.logger.Warnw("last-green URI is not in request history; using full build", @@ -291,7 +292,7 @@ func (c *Controller) deriveBuildStrategy(ctx context.Context, sc sourcecontrol.S ) return entity.BuildStrategyFull, "", nil } - metrics.NamedCounter(c.metricsScope, _opName, "source_control_errors", 1, + messageMetrics.NamedCounter(_opName, "source_control_errors", 1, metrics.NewTag("stage", "ancestry"), ) return entity.BuildStrategyUnknown, "", fmt.Errorf("failed to check ancestry for queue %s: %w", request.Queue, err) @@ -365,7 +366,7 @@ func (c *Controller) markProcessing(ctx context.Context, store storage.Storage, // releaseBuildSlot CAS-decrements queue.in_flight_count to compensate a slot claimed but never // admitted. It decrements relatively (preserving a concurrent record decrement) and retries on // version conflicts. Best-effort: it only logs on a hard failure, since the caller is unwinding. -func (c *Controller) releaseBuildSlot(ctx context.Context, store storage.Storage, queueName string) { +func (c *Controller) releaseBuildSlot(ctx context.Context, messageMetrics metrics.Scope, store storage.Storage, queueName string) { queueStore := store.GetQueueStore() for { @@ -394,7 +395,7 @@ func (c *Controller) releaseBuildSlot(ctx context.Context, store storage.Storage ) return } - metrics.NamedCounter(c.metricsScope, _opName, "slot_released", 1) + messageMetrics.NamedCounter(_opName, "slot_released", 1) return } } @@ -430,9 +431,9 @@ func (c *Controller) supersedeRequest(ctx context.Context, store storage.Storage // delayMs and the gate is re-checked, without burning MaxAttempts — the partition // (keyed by queue name) waits with it. delayMs must be positive: a non-positive // hold would redeliver immediately and hot-loop the gate check. -func (c *Controller) holdForBuildSlot(delivery consumer.Delivery, request entity.Request, inFlightCount int32, delayMs int64) error { +func (c *Controller) holdForBuildSlot(messageMetrics metrics.Scope, delivery consumer.Delivery, request entity.Request, inFlightCount int32, delayMs int64) error { if delayMs <= 0 { - metrics.NamedCounter(c.metricsScope, _opName, "config_errors", 1) + messageMetrics.NamedCounter(_opName, "config_errors", 1) return fmt.Errorf("requires a positive gate wait delay for queue %s, got %dms", request.Queue, delayMs) } diff --git a/stovepipe/controller/process/process_test.go b/stovepipe/controller/process/process_test.go index 71aae9a0c..8e05d70fa 100644 --- a/stovepipe/controller/process/process_test.go +++ b/stovepipe/controller/process/process_test.go @@ -27,6 +27,7 @@ import ( consumermock "github.com/uber/submitqueue/platform/consumer/mock" "github.com/uber/submitqueue/platform/errs" mqmock "github.com/uber/submitqueue/platform/extension/messagequeue/mock" + "github.com/uber/submitqueue/platform/metrics" stovepipemq "github.com/uber/submitqueue/stovepipe/core/messagequeue" "github.com/uber/submitqueue/stovepipe/entity" queueconfigdefault "github.com/uber/submitqueue/stovepipe/extension/queueconfig/default" @@ -109,7 +110,7 @@ func delivery(t *testing.T, ctrl *gomock.Controller, payload []byte) *consumermo func processPayload(t *testing.T, id string) []byte { t.Helper() - b, err := stovepipemq.Marshal(&stovepipemq.ProcessRequest{Id: id}) + b, err := stovepipemq.Marshal(&stovepipemq.ProcessRequest{Id: id, QueueName: testQueue}) require.NoError(t, err) return b } @@ -236,7 +237,8 @@ func TestDeriveBuildStrategy(t *testing.T) { if tt.queue.LastGreenURI != "" { sc = m.sourceControl } - strategy, baseURI, err := c.deriveBuildStrategy(context.Background(), sc, tt.queue, acceptedRequest(testID)) + messageMetrics := c.metricsFactory.For(metrics.Config{QueueName: tt.queue.Name}) + strategy, baseURI, err := c.deriveBuildStrategy(context.Background(), messageMetrics, sc, tt.queue, acceptedRequest(testID)) if tt.wantErr { require.Error(t, err) @@ -281,6 +283,7 @@ func TestDeriveBuildStrategyEmitsSourceControlMetrics(t *testing.T) { _, _, err := c.deriveBuildStrategy( context.Background(), + c.metricsFactory.For(metrics.Config{QueueName: testQueue}), m.sourceControl, entity.Queue{Name: testQueue, LastGreenURI: lastGreenURI}, acceptedRequest(testID), @@ -291,7 +294,7 @@ func TestDeriveBuildStrategyEmitsSourceControlMetrics(t *testing.T) { } else { require.Error(t, err) } - counter, ok := scope.Snapshot().Counters()["test.process_controller.process."+tt.metricName+"+"+tt.metricTags] + counter, ok := scope.Snapshot().Counters()["test.process_controller.process."+tt.metricName+"+queue="+testQueue+","+tt.metricTags] require.True(t, ok) assert.Equal(t, int64(1), counter.Value()) }) @@ -313,7 +316,7 @@ func TestProcessEmitsAdmittedStrategyMetric(t *testing.T) { require.NoError(t, c.Process(context.Background(), delivery(t, ctrl, processPayload(t, testID)))) - counter, ok := scope.Snapshot().Counters()["test.process_controller.process.admitted+strategy=full"] + counter, ok := scope.Snapshot().Counters()["test.process_controller.process.admitted+queue=monorepo/main,strategy=full"] require.True(t, ok) assert.Equal(t, int64(1), counter.Value()) } @@ -338,7 +341,7 @@ func TestProcessEmitsSourceControlResolutionMetric(t *testing.T) { require.Error(t, c.Process(context.Background(), delivery(t, ctrl, processPayload(t, testID)))) - counter, ok := scope.Snapshot().Counters()["test.process_controller.process.source_control_errors+stage=resolve"] + counter, ok := scope.Snapshot().Counters()["test.process_controller.process.source_control_errors+queue=monorepo/main,stage=resolve"] require.True(t, ok) assert.Equal(t, int64(1), counter.Value()) } @@ -825,7 +828,7 @@ func TestHoldForBuildSlotRequiresPositiveDelay(t *testing.T) { d := consumermock.NewMockDelivery(ctrl) // No Hold expectation: the guard must reject before recording a hold. - err := c.holdForBuildSlot(d, acceptedRequest(testID), 1, 0) + err := c.holdForBuildSlot(c.metricsFactory.For(metrics.Config{QueueName: testQueue}), d, acceptedRequest(testID), 1, 0) require.Error(t, err) assert.False(t, errs.IsRetryable(err)) diff --git a/stovepipe/controller/record/record.go b/stovepipe/controller/record/record.go index 5f7727605..a6a9805af 100644 --- a/stovepipe/controller/record/record.go +++ b/stovepipe/controller/record/record.go @@ -44,12 +44,12 @@ import ( // when that fact is green advances the queue's last-green bookmark and promotes // the commit. Implements consumer.Controller. type Controller struct { - logger *zap.SugaredLogger - metricsScope tally.Scope - stores storage.Factory - sourceControl sourcecontrol.Factory - topicKey consumer.TopicKey - consumerGroup string + logger *zap.SugaredLogger + metricsFactory metrics.Factory + stores storage.Factory + sourceControl sourcecontrol.Factory + topicKey consumer.TopicKey + consumerGroup string } // Verify Controller implements consumer.Controller interface at compile time. @@ -74,12 +74,12 @@ func NewController( ) *Controller { name := string(topicKey) + "_controller" return &Controller{ - logger: logger.Named(name), - metricsScope: scope.SubScope(name), - stores: stores, - sourceControl: sourceControl, - topicKey: topicKey, - consumerGroup: consumerGroup, + logger: logger.Named(name), + metricsFactory: metrics.NewFactory(scope.SubScope(name)), + stores: stores, + sourceControl: sourceControl, + topicKey: topicKey, + consumerGroup: consumerGroup, } } @@ -96,50 +96,51 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) er rec := &stovepipemq.Record{} if err := stovepipemq.Unmarshal(msg.Payload, rec); err != nil { - metrics.NamedCounter(c.metricsScope, _opName, "deserialize_errors", 1) + c.metricsFactory.Base().NamedCounter(_opName, "deserialize_errors", 1) // Non-retryable: a malformed message will never succeed regardless of retries. return fmt.Errorf("failed to deserialize record: %w", err) } + messageMetrics := c.metricsFactory.For(metrics.Config{QueueName: rec.GetQueueName()}) store, err := c.stores.For(storage.Config{QueueName: rec.GetQueueName()}) if err != nil { - metrics.NamedCounter(c.metricsScope, _opName, "storage_resolve_errors", 1) + messageMetrics.NamedCounter(_opName, "storage_resolve_errors", 1) // Non-retryable: a missing or unresolvable queue is a malformed message. return fmt.Errorf("failed to resolve storage for queue %q: %w", rec.GetQueueName(), err) } request, err := c.loadRequest(ctx, store, rec.Id) if err != nil { - metrics.NamedCounter(c.metricsScope, _opName, "storage_errors", 1) + messageMetrics.NamedCounter(_opName, "storage_errors", 1) return err } // The payload's queue must match the request's authoritative queue; a // mismatch is a malformed message. Non-retryable — reject to the DLQ. if rec.GetQueueName() != "" && rec.GetQueueName() != request.Queue { - metrics.NamedCounter(c.metricsScope, _opName, "queue_mismatch", 1) + messageMetrics.NamedCounter(_opName, "queue_mismatch", 1) return fmt.Errorf("payload queue %q does not match queue %q of request %s", rec.GetQueueName(), request.Queue, request.ID) } switch request.State { case entity.RequestStateSucceeded, entity.RequestStateFailed: - fact, created, err := c.recordFact(ctx, store, request) + fact, created, err := c.recordFact(ctx, messageMetrics, store, request) if err != nil { return err } if !fact.IsGreen() { - metrics.NamedCounter(c.metricsScope, _opName, "not_green", 1) + messageMetrics.NamedCounter(_opName, "not_green", 1) // Only the writer of the fact reports the latency: a redelivery adopts // the stored fact instead, and a second sample would count one break // twice in the distribution. if created { - c.reportFailureDetectionLatency(ctx, request) + c.reportFailureDetectionLatency(ctx, messageMetrics, request) } return nil } - holdsBookmark, err := c.advanceLastGreen(ctx, store, request) + holdsBookmark, err := c.advanceLastGreen(ctx, messageMetrics, store, request) if err != nil { - metrics.NamedCounter(c.metricsScope, _opName, "storage_errors", 1) + messageMetrics.NamedCounter(_opName, "storage_errors", 1) return err } if !holdsBookmark { @@ -148,25 +149,25 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) er // backwards. return nil } - return c.promote(ctx, request) + return c.promote(ctx, messageMetrics, request) case entity.RequestStateCancelled: // A cancelled build decided nothing about the commit, so it establishes // no fact. The identity stays unclaimed; the next commit re-validates. - metrics.NamedCounter(c.metricsScope, _opName, "cancelled", 1) + messageMetrics.NamedCounter(_opName, "cancelled", 1) return nil case entity.RequestStateSuperseded: // Terminal without a build outcome. buildsignal never publishes for a // superseded request, so this is unreachable in practice. - metrics.NamedCounter(c.metricsScope, _opName, "superseded", 1) + messageMetrics.NamedCounter(_opName, "superseded", 1) return nil default: // Non-retryable: buildsignal publishes only after committing the // outcome, so a non-terminal request here is a broken invariant that // retrying cannot fix. - metrics.NamedCounter(c.metricsScope, _opName, "invariant_errors", 1) + messageMetrics.NamedCounter(_opName, "invariant_errors", 1) return fmt.Errorf("request %s reached record in non-terminal state %q", request.ID, request.State) } } @@ -179,7 +180,7 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) er // request, so a redelivery cannot reach a different verdict than the original. The // second return reports whether this call is the one that wrote the fact, which is // how a caller tells the original delivery from a redelivery. -func (c *Controller) recordFact(ctx context.Context, store storage.Storage, request entity.Request) (entity.ValidationFact, bool, error) { +func (c *Controller) recordFact(ctx context.Context, messageMetrics metrics.Scope, store storage.Storage, request entity.Request) (entity.ValidationFact, bool, error) { factStore := store.GetValidationFactStore() fact := entity.ValidationFact{ @@ -193,7 +194,7 @@ func (c *Controller) recordFact(ctx context.Context, store storage.Storage, requ err := factStore.Create(ctx, fact) switch { case err == nil: - metrics.NamedCounter(c.metricsScope, _opName, "fact_created", 1) + messageMetrics.NamedCounter(_opName, "fact_created", 1) c.logger.Infow("recorded validation fact", "queue", request.Queue, "request_id", request.ID, @@ -205,22 +206,22 @@ func (c *Controller) recordFact(ctx context.Context, store storage.Storage, requ case errors.Is(err, storage.ErrAlreadyExists): stored, getErr := factStore.Get(ctx, request.URI, wholeRepositoryProject) if getErr != nil { - metrics.NamedCounter(c.metricsScope, _opName, "storage_errors", 1) + messageMetrics.NamedCounter(_opName, "storage_errors", 1) return entity.ValidationFact{}, false, fmt.Errorf("failed to load the existing fact for uri %s: %w", request.URI, getErr) } if stored.RequestID != request.ID { // Two requests validating one URI would break the dedup ingest // enforces, so this is a broken invariant rather than a race to // resolve. Non-retryable: the stored fact is immutable. - metrics.NamedCounter(c.metricsScope, _opName, "invariant_errors", 1) + messageMetrics.NamedCounter(_opName, "invariant_errors", 1) return entity.ValidationFact{}, false, fmt.Errorf( "fact for uri %s is owned by request %s, not %s", request.URI, stored.RequestID, request.ID) } - metrics.NamedCounter(c.metricsScope, _opName, "fact_exists", 1) + messageMetrics.NamedCounter(_opName, "fact_exists", 1) return stored, false, nil default: - metrics.NamedCounter(c.metricsScope, _opName, "storage_errors", 1) + messageMetrics.NamedCounter(_opName, "storage_errors", 1) return entity.ValidationFact{}, false, fmt.Errorf("failed to create the fact for uri %s: %w", request.URI, err) } } @@ -235,27 +236,26 @@ func (c *Controller) recordFact(ctx context.Context, store storage.Storage, requ // source-control lookup cannot be moved off the delivery path onto a clock. It is // confined to failures and made once the fact is durable, and every way it can fail is // counted and swallowed so a reporting fault cannot retry an outcome already recorded. -func (c *Controller) reportFailureDetectionLatency(ctx context.Context, request entity.Request) { - queueTag := metrics.NewTag("queue", request.Queue) +func (c *Controller) reportFailureDetectionLatency(ctx context.Context, messageMetrics metrics.Scope, request entity.Request) { strategyTag := metrics.NewTag("strategy", string(request.BuildStrategy)) // Only a strategy that validates a delta pins a base commit, so a full build has // no baseline to measure from. Its failures are counted rather than timed: absent // here is the ordinary case, not a fault. if request.BaseURI == "" { - metrics.NamedCounter(c.metricsScope, _opName, "failure_detection_missing", 1, queueTag, strategyTag) + messageMetrics.NamedCounter(_opName, "failure_detection_missing", 1, strategyTag) return } sourceControl, err := c.sourceControl.For(sourcecontrol.Config{QueueName: request.Queue}) if err != nil { - c.failureDetectionUnobserved(request, "resolve_source_control", err) + c.failureDetectionUnobserved(messageMetrics, request, "resolve_source_control", err) return } info, err := sourceControl.ChangeInfo(ctx, request.BaseURI) if err != nil { - c.failureDetectionUnobserved(request, "get_change_info", err) + c.failureDetectionUnobserved(messageMetrics, request, "get_change_info", err) return } @@ -263,7 +263,7 @@ func (c *Controller) reportFailureDetectionLatency(ctx context.Context, request // broken extension contract rather than a lookup failure. Measuring from 1970 // would drop a decades-long sample into the distribution. if info.CreatedAt <= 0 { - c.failureDetectionUnobserved(request, "undated_change", nil) + c.failureDetectionUnobserved(messageMetrics, request, "undated_change", nil) return } @@ -271,21 +271,20 @@ func (c *Controller) reportFailureDetectionLatency(ctx context.Context, request // negative latency would corrupt the distribution rather than describe it. latency := time.Since(time.UnixMilli(info.CreatedAt)) if latency < 0 { - c.failureDetectionUnobserved(request, "future_change", nil) + c.failureDetectionUnobserved(messageMetrics, request, "future_change", nil) return } - metrics.NamedHistogram(c.metricsScope, _opName, "failure_detection_latency", metrics.ChangeAgeBuckets, - queueTag, strategyTag, + messageMetrics.NamedHistogram(_opName, "failure_detection_latency", metrics.ChangeAgeBuckets, + strategyTag, ).RecordDuration(latency) } // failureDetectionUnobserved counts a latency that could not be observed, tagged with // the step that failed so an unmeasurable failure can be told apart from a broken // dependency. -func (c *Controller) failureDetectionUnobserved(request entity.Request, step string, err error) { - metrics.NamedCounter(c.metricsScope, _opName, "failure_detection_errors", 1, - metrics.NewTag("queue", request.Queue), +func (c *Controller) failureDetectionUnobserved(messageMetrics metrics.Scope, request entity.Request, step string, err error) { + messageMetrics.NamedCounter(_opName, "failure_detection_errors", 1, metrics.NewTag("step", step), ) c.logger.Warnw("failed to observe how long the build failure went undetected", @@ -316,7 +315,7 @@ func degreeFor(state entity.RequestState) float64 { // advanced only after the green fact is durable. Losing the advance to a crash is // recoverable — the redelivery reloads the same fact and retries — whereas a // bookmark with no fact behind it would point at greenness nothing recorded. -func (c *Controller) advanceLastGreen(ctx context.Context, store storage.Storage, request entity.Request) (bool, error) { +func (c *Controller) advanceLastGreen(ctx context.Context, messageMetrics metrics.Scope, store storage.Storage, request entity.Request) (bool, error) { queueStore := store.GetQueueStore() for { @@ -348,13 +347,13 @@ func (c *Controller) advanceLastGreen(ctx context.Context, store storage.Storage return false, fmt.Errorf("failed to advance last green for queue %s: %w", request.Queue, err) } - metrics.NamedCounter(c.metricsScope, _opName, "last_green_advanced", 1) + messageMetrics.NamedCounter(_opName, "last_green_advanced", 1) c.logger.Infow("advanced last green bookmark", "queue", request.Queue, "request_id", request.ID, "last_green_uri", request.URI, ) - c.emitLastGreenTimestamp(ctx, request) + c.emitLastGreenTimestamp(ctx, messageMetrics, request) return true, nil } } @@ -363,12 +362,10 @@ func (c *Controller) advanceLastGreen(ctx context.Context, store storage.Storage // points at, once that bookmark is durable. Reporting is best-effort so an // observability failure cannot turn a successful record operation into a retry, // which is why each cause is counted and logged separately instead of returned. -func (c *Controller) emitLastGreenTimestamp(ctx context.Context, request entity.Request) { - queueTag := metrics.NewTag("queue", request.Queue) - +func (c *Controller) emitLastGreenTimestamp(ctx context.Context, messageMetrics metrics.Scope, request entity.Request) { sourceControl, err := c.sourceControl.For(sourcecontrol.Config{QueueName: request.Queue}) if err != nil { - metrics.NamedCounter(c.metricsScope, _opName, "last_green_timestamp_resolve_errors", 1, queueTag) + messageMetrics.NamedCounter(_opName, "last_green_timestamp_resolve_errors", 1) c.logger.Warnw("failed to resolve source control to report the last green timestamp", "queue", request.Queue, "error", err, @@ -378,7 +375,7 @@ func (c *Controller) emitLastGreenTimestamp(ctx context.Context, request entity. info, err := sourceControl.ChangeInfo(ctx, request.URI) if err != nil { - metrics.NamedCounter(c.metricsScope, _opName, "last_green_timestamp_errors", 1, queueTag) + messageMetrics.NamedCounter(_opName, "last_green_timestamp_errors", 1) c.logger.Warnw("failed to look up the last green change timestamp", "queue", request.Queue, "uri", request.URI, @@ -391,7 +388,7 @@ func (c *Controller) emitLastGreenTimestamp(ctx context.Context, request entity. // is a broken extension contract rather than a lookup failure. Emitting it // anyway would publish a 1970 timestamp and read as an infinitely stale queue. if info.CreatedAt <= 0 { - metrics.NamedCounter(c.metricsScope, _opName, "last_green_timestamp_invalid", 1, queueTag) + messageMetrics.NamedCounter(_opName, "last_green_timestamp_invalid", 1) c.logger.Warnw("source control reported no creation timestamp for the last green change", "queue", request.Queue, "uri", request.URI, @@ -402,12 +399,10 @@ func (c *Controller) emitLastGreenTimestamp(ctx context.Context, request entity. // The gauge carries the creation time as Unix seconds, so subtracting it // from the current time yields the age of the last-green change in seconds. - metrics.NamedGauge( - c.metricsScope, + messageMetrics.NamedGauge( _opName, "last_green_timestamp_seconds", float64(time.UnixMilli(info.CreatedAt).Unix()), - queueTag, ) } @@ -420,10 +415,10 @@ func (c *Controller) emitLastGreenTimestamp(ctx context.Context, request entity. // green fact is durable. Promotion is idempotent, so a redelivery repeats it // harmlessly. A commit that a rewritten history dropped from the ref cannot be // promoted by any retry, so that case is counted and skipped rather than failed. -func (c *Controller) promote(ctx context.Context, request entity.Request) error { +func (c *Controller) promote(ctx context.Context, messageMetrics metrics.Scope, request entity.Request) error { sc, err := c.sourceControl.For(sourcecontrol.Config{QueueName: request.Queue}) if err != nil { - metrics.NamedCounter(c.metricsScope, _opName, "source_control_errors", 1, + messageMetrics.NamedCounter(_opName, "source_control_errors", 1, metrics.NewTag("stage", "resolve"), ) return fmt.Errorf("failed to resolve source control for queue %s: %w", request.Queue, err) @@ -431,7 +426,7 @@ func (c *Controller) promote(ctx context.Context, request entity.Request) error if err := sc.Promote(ctx, request.URI); err != nil { if sourcecontrol.IsNotFound(err) { - metrics.NamedCounter(c.metricsScope, _opName, "promotions_skipped", 1, + messageMetrics.NamedCounter(_opName, "promotions_skipped", 1, metrics.NewTag("reason", "unknown_uri"), ) c.logger.Warnw("green commit is no longer on the queue's ref; skipping promotion", @@ -442,13 +437,13 @@ func (c *Controller) promote(ctx context.Context, request entity.Request) error return nil } - metrics.NamedCounter(c.metricsScope, _opName, "source_control_errors", 1, + messageMetrics.NamedCounter(_opName, "source_control_errors", 1, metrics.NewTag("stage", "promote"), ) return fmt.Errorf("failed to promote uri %s of queue %s: %w", request.URI, request.Queue, err) } - metrics.NamedCounter(c.metricsScope, _opName, "promotions", 1) + messageMetrics.NamedCounter(_opName, "promotions", 1) c.logger.Infow("promoted green commit", "queue", request.Queue, "request_id", request.ID, diff --git a/stovepipe/controller/record/record_test.go b/stovepipe/controller/record/record_test.go index 02aeda08a..b6b20a387 100644 --- a/stovepipe/controller/record/record_test.go +++ b/stovepipe/controller/record/record_test.go @@ -326,6 +326,9 @@ func TestProcess_RecordsBrokenFactWithoutAdvancing(t *testing.T) { require.NoError(t, c.Process(context.Background(), delivery(t, ctrl, recordPayload(t, testID)))) assert.Equal(t, entity.DegreeBroken, fact.Degree) assert.False(t, fact.IsGreen()) + counter, ok := m.metricsScope.Snapshot().Counters()["record_controller.record.not_green+queue=monorepo/main"] + require.True(t, ok) + assert.EqualValues(t, 1, counter.Value()) } func TestProcess_ReportsFailureDetectionLatency(t *testing.T) {