Skip to content

Commit b282822

Browse files
committed
refactor(stovepipe): bind freshness reporters to queues
Resolve storage and source control through a queue-scoped reporter factory so individual reporters operate on concrete queue dependencies.
1 parent dac4f91 commit b282822

5 files changed

Lines changed: 78 additions & 30 deletions

File tree

stovepipe/controller/record/record.go

Lines changed: 11 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -45,7 +45,7 @@ type Controller struct {
4545
logger *zap.SugaredLogger
4646
metricsScope tally.Scope
4747
stores storage.Factory
48-
reporter observability.Reporter
48+
reporters observability.Factory
4949
topicKey consumer.TopicKey
5050
consumerGroup string
5151
}
@@ -64,10 +64,10 @@ const wholeRepositoryProject = ""
6464
// Option configures a Controller.
6565
type Option func(*Controller)
6666

67-
// WithReporter configures best-effort queue observability reporting.
68-
func WithReporter(reporter observability.Reporter) Option {
67+
// WithReporterFactory configures best-effort queue observability reporting.
68+
func WithReporterFactory(reporters observability.Factory) Option {
6969
return func(c *Controller) {
70-
c.reporter = reporter
70+
c.reporters = reporters
7171
}
7272
}
7373

@@ -129,8 +129,13 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) er
129129
metrics.NamedCounter(c.metricsScope, _opName, "queue_mismatch", 1)
130130
return fmt.Errorf("payload queue %q does not match queue %q of request %s", rec.GetQueueName(), request.Queue, request.ID)
131131
}
132-
if c.reporter != nil {
133-
defer c.reporter.Report(ctx, request.Queue)
132+
if c.reporters != nil {
133+
reporter, err := c.reporters.For(observability.Config{QueueName: request.Queue})
134+
if err != nil {
135+
metrics.NamedCounter(c.metricsScope, _opName, "reporter_resolve_errors", 1)
136+
} else {
137+
defer reporter.Report(ctx)
138+
}
134139
}
135140

136141
switch request.State {

stovepipe/extension/observability/lastgreen/BUILD.bazel

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,7 @@ go_test(
1919
embed = [":go_default_library"],
2020
deps = [
2121
"//stovepipe/entity:go_default_library",
22+
"//stovepipe/extension/observability:go_default_library",
2223
"//stovepipe/extension/sourcecontrol:go_default_library",
2324
"//stovepipe/extension/sourcecontrol/mock:go_default_library",
2425
"//stovepipe/extension/storage:go_default_library",

stovepipe/extension/observability/lastgreen/lastgreen.go

Lines changed: 49 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,7 @@ package lastgreen
33

44
import (
55
"context"
6+
"fmt"
67
"time"
78

89
"github.com/uber-go/tally"
@@ -11,32 +12,67 @@ import (
1112
"github.com/uber/submitqueue/stovepipe/extension/storage"
1213
)
1314

15+
type factory struct {
16+
scope tally.Scope
17+
stores storage.Factory
18+
sourceControls sourcecontrol.Factory
19+
}
20+
1421
type reporter struct {
1522
scope tally.Scope
16-
stores storage.Factory
17-
sourceControl sourcecontrol.Factory
23+
queue string
24+
store storage.Storage
25+
sourceControl sourcecontrol.SourceControl
1826
}
1927

28+
var _ observability.Factory = (*factory)(nil)
2029
var _ observability.Reporter = (*reporter)(nil)
2130

31+
// NewFactory creates queue-scoped last-known-green reporters.
32+
func NewFactory(
33+
scope tally.Scope,
34+
stores storage.Factory,
35+
sourceControls sourcecontrol.Factory,
36+
) observability.Factory {
37+
return &factory{
38+
scope: scope,
39+
stores: stores,
40+
sourceControls: sourceControls,
41+
}
42+
}
43+
44+
// For resolves the concrete dependencies needed to report for one queue.
45+
func (f *factory) For(cfg observability.Config) (observability.Reporter, error) {
46+
store, err := f.stores.For(storage.Config{QueueName: cfg.QueueName})
47+
if err != nil {
48+
return nil, fmt.Errorf("resolve storage for queue %q: %w", cfg.QueueName, err)
49+
}
50+
sourceControl, err := f.sourceControls.For(sourcecontrol.Config{QueueName: cfg.QueueName})
51+
if err != nil {
52+
return nil, fmt.Errorf("resolve source control for queue %q: %w", cfg.QueueName, err)
53+
}
54+
return New(f.scope, cfg.QueueName, store, sourceControl), nil
55+
}
56+
2257
// New creates a Reporter for a queue's last-known-green age.
23-
func New(scope tally.Scope, stores storage.Factory, sourceControl sourcecontrol.Factory) observability.Reporter {
58+
func New(
59+
scope tally.Scope,
60+
queue string,
61+
store storage.Storage,
62+
sourceControl sourcecontrol.SourceControl,
63+
) observability.Reporter {
2464
return &reporter{
2565
scope: scope.SubScope("last_green"),
26-
stores: stores,
66+
queue: queue,
67+
store: store,
2768
sourceControl: sourceControl,
2869
}
2970
}
3071

3172
// Report updates the queue's current last-known-green age gauge.
32-
func (r *reporter) Report(ctx context.Context, queue string) {
33-
tags := map[string]string{"queue": queue}
34-
store, err := r.stores.For(storage.Config{QueueName: queue})
35-
if err != nil {
36-
r.error(tags, "resolve_storage")
37-
return
38-
}
39-
queueRow, err := store.GetQueueStore().Get(ctx, queue)
73+
func (r *reporter) Report(ctx context.Context) {
74+
tags := map[string]string{"queue": r.queue}
75+
queueRow, err := r.store.GetQueueStore().Get(ctx, r.queue)
4076
if err != nil {
4177
r.error(tags, "get_queue")
4278
return
@@ -45,12 +81,7 @@ func (r *reporter) Report(ctx context.Context, queue string) {
4581
r.scope.Tagged(tags).Counter("age_missing").Inc(1)
4682
return
4783
}
48-
control, err := r.sourceControl.For(sourcecontrol.Config{QueueName: queue})
49-
if err != nil {
50-
r.error(tags, "resolve_source_control")
51-
return
52-
}
53-
info, err := control.ChangeInfo(ctx, queueRow.LastGreenURI)
84+
info, err := r.sourceControl.ChangeInfo(ctx, queueRow.LastGreenURI)
5485
if err != nil || info.CreatedAt.IsZero() {
5586
r.error(tags, "get_change_info")
5687
return

stovepipe/extension/observability/lastgreen/lastgreen_test.go

Lines changed: 5 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,7 @@ import (
99
"github.com/stretchr/testify/require"
1010
"github.com/uber-go/tally"
1111
"github.com/uber/submitqueue/stovepipe/entity"
12+
"github.com/uber/submitqueue/stovepipe/extension/observability"
1213
sourcecontrol "github.com/uber/submitqueue/stovepipe/extension/sourcecontrol"
1314
sourcecontrolmock "github.com/uber/submitqueue/stovepipe/extension/sourcecontrol/mock"
1415
"github.com/uber/submitqueue/stovepipe/extension/storage"
@@ -39,7 +40,9 @@ func TestReport_EmitsLastGreenAge(t *testing.T) {
3940
CreatedAt: createdAt,
4041
}, nil)
4142

42-
New(scope, stores, sourceControls).Report(context.Background(), testQueue)
43+
reporter, err := NewFactory(scope, stores, sourceControls).For(observability.Config{QueueName: testQueue})
44+
require.NoError(t, err)
45+
reporter.Report(context.Background())
4346

4447
gauge, ok := scope.Snapshot().Gauges()["stovepipe.last_green.age_seconds+queue=monorepo/main"]
4548
require.True(t, ok)
@@ -49,15 +52,13 @@ func TestReport_EmitsLastGreenAge(t *testing.T) {
4952
func TestReport_RecordsMissingLastGreen(t *testing.T) {
5053
ctrl := gomock.NewController(t)
5154
scope := tally.NewTestScope("stovepipe", nil)
52-
stores := storagemock.NewMockFactory(ctrl)
5355
store := storagemock.NewMockStorage(ctrl)
5456
queueStore := storagemock.NewMockQueueStore(ctrl)
5557

56-
stores.EXPECT().For(storage.Config{QueueName: testQueue}).Return(store, nil)
5758
store.EXPECT().GetQueueStore().Return(queueStore)
5859
queueStore.EXPECT().Get(gomock.Any(), testQueue).Return(entity.Queue{Name: testQueue}, nil)
5960

60-
New(scope, stores, sourcecontrolmock.NewMockFactory(ctrl)).Report(context.Background(), testQueue)
61+
New(scope, testQueue, store, sourcecontrolmock.NewMockSourceControl(ctrl)).Report(context.Background())
6162

6263
counter, ok := scope.Snapshot().Counters()["stovepipe.last_green.age_missing+queue=monorepo/main"]
6364
require.True(t, ok)

stovepipe/extension/observability/observability.go

Lines changed: 12 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -3,7 +3,17 @@ package observability
33

44
import "context"
55

6-
// Reporter emits best-effort observability data for a queue.
6+
// Reporter emits best-effort observability data for one queue.
77
type Reporter interface {
8-
Report(context.Context, string)
8+
Report(context.Context)
9+
}
10+
11+
// Config identifies the queue whose observability reporter is requested.
12+
type Config struct {
13+
QueueName string
14+
}
15+
16+
// Factory resolves a Reporter bound to one queue.
17+
type Factory interface {
18+
For(Config) (Reporter, error)
919
}

0 commit comments

Comments
 (0)