From 8705a55e2b7572093c887442fa1b9c3ca0f2d258 Mon Sep 17 00:00:00 2001 From: mnoah1 Date: Tue, 6 Oct 2026 15:40:58 +0000 Subject: [PATCH] fix(stovepipe): tag storage metrics by queue Summary: Intent: - Make queue-filtered storage activity and latency visible; storage metrics currently omit the queue tag. Changes: - Tag the metric scope in Storage.For before constructing all seven queue-bound MySQL stores. - Add regression coverage for interleaved queues, preserved parent tags, and success, error, and cancellation outcomes. Rollout: Consumers must update their SubmitQueue dependency and deploy to emit the new queue-tagged series. Metric names and database behavior are unchanged; historical untagged metrics are not backfilled. --- stovepipe/extension/storage/mysql/storage.go | 17 +- .../extension/storage/mysql/storage_test.go | 147 ++++++++++++++++++ 2 files changed, 156 insertions(+), 8 deletions(-) diff --git a/stovepipe/extension/storage/mysql/storage.go b/stovepipe/extension/storage/mysql/storage.go index b1d3476a2..c8c714e6c 100644 --- a/stovepipe/extension/storage/mysql/storage.go +++ b/stovepipe/extension/storage/mysql/storage.go @@ -39,19 +39,20 @@ func NewStorage(db *sql.DB, scope tally.Scope) (*Storage, error) { // For returns the queue-scoped store aggregate bound to queueName over the // shared pool. Every store the aggregate hands back reads and writes only that -// queue's records. +// queue's records and tags its metrics with queue=queueName. func (s *Storage) For(queueName string) (storage.Storage, error) { if queueName == "" { return nil, fmt.Errorf("queue name must not be empty") } + queueScope := s.scope.Tagged(map[string]string{"queue": queueName}) return &mysqlStorage{ - requestStore: NewRequestStore(s.db, s.scope.SubScope("request_store"), queueName), - requestURIStore: NewRequestURIStore(s.db, s.scope.SubScope("request_uri_store"), queueName), - requestLogStore: NewRequestLogStore(s.db, s.scope.SubScope("request_log_store"), queueName), - requestSummaryStore: NewRequestSummaryStore(s.db, s.scope.SubScope("request_summary_store"), queueName), - queueStore: NewQueueStore(s.db, s.scope.SubScope("queue_store"), queueName), - buildStore: NewBuildStore(s.db, s.scope.SubScope("build_store"), queueName), - validationFactStore: NewValidationFactStore(s.db, s.scope.SubScope("validation_fact_store"), queueName), + requestStore: NewRequestStore(s.db, queueScope.SubScope("request_store"), queueName), + requestURIStore: NewRequestURIStore(s.db, queueScope.SubScope("request_uri_store"), queueName), + requestLogStore: NewRequestLogStore(s.db, queueScope.SubScope("request_log_store"), queueName), + requestSummaryStore: NewRequestSummaryStore(s.db, queueScope.SubScope("request_summary_store"), queueName), + queueStore: NewQueueStore(s.db, queueScope.SubScope("queue_store"), queueName), + buildStore: NewBuildStore(s.db, queueScope.SubScope("build_store"), queueName), + validationFactStore: NewValidationFactStore(s.db, queueScope.SubScope("validation_fact_store"), queueName), }, nil } diff --git a/stovepipe/extension/storage/mysql/storage_test.go b/stovepipe/extension/storage/mysql/storage_test.go index bba835941..37c65123a 100644 --- a/stovepipe/extension/storage/mysql/storage_test.go +++ b/stovepipe/extension/storage/mysql/storage_test.go @@ -15,12 +15,15 @@ package mysql import ( + "context" + "errors" "testing" "github.com/DATA-DOG/go-sqlmock" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" "github.com/uber-go/tally" + "github.com/uber/submitqueue/stovepipe/extension/storage" ) // testMetrics returns a test metrics scope for use in tests. @@ -61,3 +64,147 @@ func TestMysqlStorage_Close(t *testing.T) { require.NoError(t, s.Close()) require.NoError(t, mock.ExpectationsWereMet()) } + +func TestStorageForQueueMetrics(t *testing.T) { + tests := []struct { + name string + op string + read func(storage.Storage, string) error + }{ + { + name: "request_store", + op: "get", + read: func(bound storage.Storage, _ string) error { + _, err := bound.GetRequestStore().Get(context.Background(), "request-id") + return err + }, + }, + { + name: "request_uri_store", + op: "get_id_by_uri", + read: func(bound storage.Storage, _ string) error { + _, err := bound.GetRequestURIStore().GetIDByURI(context.Background(), "change-uri") + return err + }, + }, + { + name: "request_log_store", + op: "get", + read: func(bound storage.Storage, _ string) error { + _, err := bound.GetRequestLogStore().Get(context.Background(), "request-id", "log-id") + return err + }, + }, + { + name: "request_summary_store", + op: "get", + read: func(bound storage.Storage, _ string) error { + _, err := bound.GetRequestSummaryStore().Get(context.Background(), "request-id") + return err + }, + }, + { + name: "queue_store", + op: "get", + read: func(bound storage.Storage, queue string) error { + _, err := bound.GetQueueStore().Get(context.Background(), queue) + return err + }, + }, + { + name: "build_store", + op: "get", + read: func(bound storage.Storage, _ string) error { + _, err := bound.GetBuildStore().Get(context.Background(), "build-id") + return err + }, + }, + { + name: "validation_fact_store", + op: "get", + read: func(bound storage.Storage, _ string) error { + _, err := bound.GetValidationFactStore().Get(context.Background(), "change-uri", "project") + return err + }, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + db, mock, err := sqlmock.New() + require.NoError(t, err) + defer db.Close() + + scope := tally.NewTestScope("storage", map[string]string{"component": "test"}) + s, err := NewStorage(db, scope) + require.NoError(t, err) + dbErr := errors.New("database unavailable") + for _, queue := range []string{"queue-a", "queue-b", "queue-a"} { + bound, err := s.For(queue) + require.NoError(t, err) + mock.ExpectQuery("SELECT").WillReturnError(dbErr) + require.ErrorIs(t, tt.read(bound, queue), dbErr) + } + + scope.Counter("unbound").Inc(1) + snapshot := scope.Snapshot() + require.Len(t, snapshot.Counters(), 3) + require.Contains(t, snapshot.Counters(), "storage.unbound+component=test") + require.Len(t, snapshot.Histograms(), 2) + for queue, count := range map[string]int64{"queue-a": 2, "queue-b": 1} { + metric := "storage." + tt.name + "." + tt.op + start := metric + ".start+component=test,queue=" + queue + counter, ok := snapshot.Counters()[start] + require.True(t, ok, "missing metric %s", start) + assert.Equal(t, count, counter.Value()) + finish := metric + ".finish+component=test,queue=" + queue + ",result=error" + histogram, ok := snapshot.Histograms()[finish] + require.True(t, ok, "missing metric %s", finish) + var samples int64 + for _, value := range histogram.Durations() { + samples += value + } + assert.Equal(t, count, samples) + } + require.NoError(t, mock.ExpectationsWereMet()) + }) + } +} + +func TestStorageForQueueMetricsOutcomes(t *testing.T) { + for _, tt := range []struct { + result string + err error + }{ + {result: "success"}, + {result: "error", err: errors.New("database unavailable")}, + {result: "cancel", err: context.Canceled}, + } { + t.Run(tt.result, func(t *testing.T) { + db, mock, err := sqlmock.New() + require.NoError(t, err) + defer db.Close() + scope := tally.NewTestScope("storage", nil) + s, err := NewStorage(db, scope) + require.NoError(t, err) + bound, err := s.For("queue-a") + require.NoError(t, err) + expectation := mock.ExpectExec("INSERT INTO request_uri").WithArgs("queue-a", "change-uri", "request-id", requestURIInitialVersion) + if tt.err != nil { + expectation.WillReturnError(tt.err) + } else { + expectation.WillReturnResult(sqlmock.NewResult(1, 1)) + } + err = bound.GetRequestURIStore().Create(context.Background(), "change-uri", "request-id") + if tt.err != nil { + require.ErrorIs(t, err, tt.err) + } else { + require.NoError(t, err) + } + snapshot := scope.Snapshot() + require.Contains(t, snapshot.Counters(), "storage.request_uri_store.create.start+queue=queue-a") + require.Contains(t, snapshot.Histograms(), "storage.request_uri_store.create.finish+queue=queue-a,result="+tt.result) + require.NoError(t, mock.ExpectationsWereMet()) + }) + } +}