Skip to content

Commit 9adfbb2

Browse files
authored
feat(stovepipe): add request logs for outcomes and lifecycle events (#666)
## Summary This PR builds on #665, which records Process-owned request states. Intent: - Retain terminal request outcomes and the durable build/fact milestones that explain them. - Let queue redelivery repair any missing occurrence before downstream work continues. Changes: - Persist succeeded, failed, and cancelled request states with their build outcome reasons. - Record build_triggered, build_finished, and validation_fact_recorded events with stable identities and bounded metadata. - Keep each source write ahead of event materialization and materialization ahead of downstream publication or derived work. - Wire the shared materializer into Build, BuildSignal, and Record. ## Test Plan - Run the focused Build, BuildSignal, Record, request-log materializer, and server wiring Bazel tests. ## Revert Plan - Revert this change to stop recording terminal outcomes and lifecycle events. --- <sub>Generated by the 🪄 [pr-create](https://sg.uberinternal.com/code.uber.internal/uber-code/devexp-agent-marketplace/-/blob/claude-code/plugins/dev/uber-dev/skills/pr-create/SKILL.md) skill in devexp-agent-marketplace</sub> ## Issues
1 parent c19d95e commit 9adfbb2

13 files changed

Lines changed: 410 additions & 54 deletions

File tree

‎service/stovepipe/server/main.go‎

Lines changed: 6 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -308,7 +308,7 @@ func run() error {
308308
if err != nil {
309309
return err
310310
}
311-
dlqCount, err := registerDLQControllers(dlqConsumer, logger.Sugar(), scope, storageFty, registry, sourceControl)
311+
dlqCount, err := registerDLQControllers(dlqConsumer, logger.Sugar(), scope, storageFty, materializer, registry, sourceControl)
312312
if err != nil {
313313
return err
314314
}
@@ -440,19 +440,19 @@ func registerPrimaryControllers(
440440
}
441441
count++
442442

443-
buildController := build.NewController(logger, scope, store, brf, registry, stovepipemq.TopicKeyBuild, "stovepipe-build")
443+
buildController := build.NewController(logger, scope, store, materializer, brf, registry, stovepipemq.TopicKeyBuild, "stovepipe-build")
444444
if err := c.Register(buildController); err != nil {
445445
return count, fmt.Errorf("failed to register build controller: %w", err)
446446
}
447447
count++
448448

449-
buildSignalController := buildsignal.NewController(logger, scope, store, brf, registry, stovepipemq.TopicKeyBuildSignal, "stovepipe-buildsignal")
449+
buildSignalController := buildsignal.NewController(logger, scope, store, materializer, brf, registry, stovepipemq.TopicKeyBuildSignal, "stovepipe-buildsignal")
450450
if err := c.Register(buildSignalController); err != nil {
451451
return count, fmt.Errorf("failed to register buildsignal controller: %w", err)
452452
}
453453
count++
454454

455-
recordController := record.NewController(logger, scope, store, sourceControl, registry, stovepipemq.TopicKeyRecord, "stovepipe-record")
455+
recordController := record.NewController(logger, scope, store, materializer, sourceControl, registry, stovepipemq.TopicKeyRecord, "stovepipe-record")
456456
if err := c.Register(recordController); err != nil {
457457
return count, fmt.Errorf("failed to register record controller: %w", err)
458458
}
@@ -474,6 +474,7 @@ func registerDLQControllers(
474474
logger *zap.SugaredLogger,
475475
scope tally.Scope,
476476
store storage.Factory,
477+
materializer requestlog.Materializer,
477478
registry consumer.TopicRegistry,
478479
sourceControl sourcecontrol.Factory,
479480
) (int, error) {
@@ -497,7 +498,7 @@ func registerDLQControllers(
497498
}
498499
count++
499500

500-
recordDLQController := record.NewController(logger, scope, store, sourceControl, registry, dlq.TopicKey(stovepipemq.TopicKeyRecord), "stovepipe-record-dlq")
501+
recordDLQController := record.NewController(logger, scope, store, materializer, sourceControl, registry, dlq.TopicKey(stovepipemq.TopicKeyRecord), "stovepipe-record-dlq")
501502
if err := c.Register(recordDLQController); err != nil {
502503
return count, fmt.Errorf("failed to register record dlq controller: %w", err)
503504
}

‎service/stovepipe/server/main_test.go‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -61,7 +61,7 @@ func registeredControllers(t *testing.T) (consumer.TopicRegistry, []consumer.Con
6161
fakeSourceControlFactory{}, fakeBuildRunnerFactory{}, hookResolver{})
6262
require.NoError(t, err)
6363

64-
_, err = registerDLQControllers(deadLetter, logger, tally.NoopScope, store, registry,
64+
_, err = registerDLQControllers(deadLetter, logger, tally.NoopScope, store, requestlog.NewMaterializer(tally.NoopScope), registry,
6565
fakeSourceControlFactory{})
6666
require.NoError(t, err)
6767

‎stovepipe/controller/build/BUILD.bazel‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -12,6 +12,7 @@ go_library(
1212
"//platform/publish:go_default_library",
1313
"//stovepipe/core/loader:go_default_library",
1414
"//stovepipe/core/messagequeue:go_default_library",
15+
"//stovepipe/core/requestlog:go_default_library",
1516
"//stovepipe/entity:go_default_library",
1617
"//stovepipe/extension/buildrunner:go_default_library",
1718
"//stovepipe/extension/storage:go_default_library",
@@ -32,6 +33,8 @@ go_test(
3233
"//platform/extension/messagequeue/mock:go_default_library",
3334
"//platform/metrics:go_default_library",
3435
"//stovepipe/core/messagequeue:go_default_library",
36+
"//stovepipe/core/requestlog:go_default_library",
37+
"//stovepipe/core/requestlog/mock:go_default_library",
3538
"//stovepipe/entity:go_default_library",
3639
"//stovepipe/extension/buildrunner:go_default_library",
3740
"//stovepipe/extension/buildrunner/mock:go_default_library",

‎stovepipe/controller/build/build.go‎

Lines changed: 21 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -20,7 +20,6 @@ package build
2020

2121
import (
2222
"context"
23-
"errors"
2423
"fmt"
2524

2625
"github.com/uber-go/tally"
@@ -30,6 +29,7 @@ import (
3029
"github.com/uber/submitqueue/platform/publish"
3130
"github.com/uber/submitqueue/stovepipe/core/loader"
3231
stovepipemq "github.com/uber/submitqueue/stovepipe/core/messagequeue"
32+
"github.com/uber/submitqueue/stovepipe/core/requestlog"
3333
"github.com/uber/submitqueue/stovepipe/entity"
3434
"github.com/uber/submitqueue/stovepipe/extension/buildrunner"
3535
"github.com/uber/submitqueue/stovepipe/extension/storage"
@@ -43,6 +43,7 @@ type Controller struct {
4343
logger *zap.SugaredLogger
4444
metricsScope tally.Scope
4545
stores storage.Factory
46+
materializer requestlog.Materializer
4647
buildRunners buildrunner.Factory
4748
registry consumer.TopicRegistry
4849
topicKey consumer.TopicKey
@@ -60,6 +61,7 @@ func NewController(
6061
logger *zap.SugaredLogger,
6162
scope tally.Scope,
6263
stores storage.Factory,
64+
materializer requestlog.Materializer,
6365
buildRunners buildrunner.Factory,
6466
registry consumer.TopicRegistry,
6567
topicKey consumer.TopicKey,
@@ -69,6 +71,7 @@ func NewController(
6971
logger: logger.Named("build_controller"),
7072
metricsScope: scope.SubScope("build_controller"),
7173
stores: stores,
74+
materializer: materializer,
7275
buildRunners: buildRunners,
7376
registry: registry,
7477
topicKey: topicKey,
@@ -141,9 +144,12 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) er
141144
Status: entity.BuildStatusAccepted,
142145
Version: 1,
143146
}
144-
if err := store.GetBuildStore().Create(ctx, build); err != nil && !errors.Is(err, storage.ErrAlreadyExists) {
147+
if err := store.GetBuildStore().Create(ctx, build); err != nil {
145148
return fmt.Errorf("failed to persist build %s: %w", build.ID, err)
146149
}
150+
if err := c.persistBuildTriggeredLog(ctx, store, request, build.ID); err != nil {
151+
return err
152+
}
147153

148154
if err := c.publishBuildSignal(ctx, build.ID, request.Queue); err != nil {
149155
return fmt.Errorf("failed to publish build signal for %s: %w", build.ID, err)
@@ -158,6 +164,19 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) er
158164
return nil
159165
}
160166

167+
func (c *Controller) persistBuildTriggeredLog(ctx context.Context, store storage.Storage, request entity.Request, buildID string) error {
168+
log := requestlog.NewRequestEventLog(
169+
request,
170+
entity.RequestEventBuildTriggered,
171+
buildID,
172+
map[string]string{requestlog.MetadataKeyBuildID: buildID},
173+
)
174+
if err := c.materializer.PersistLog(ctx, store, log); err != nil {
175+
return fmt.Errorf("failed to record build %s trigger for request %s: %w", buildID, request.ID, err)
176+
}
177+
return nil
178+
}
179+
161180
// loadRequest returns the request for id.
162181
func (c *Controller) loadRequest(ctx context.Context, store storage.Storage, id string) (entity.Request, error) {
163182
return loader.ByID(ctx, id, store.GetRequestStore().Get, "request")

‎stovepipe/controller/build/build_test.go‎

Lines changed: 56 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -29,6 +29,8 @@ import (
2929
mqmock "github.com/uber/submitqueue/platform/extension/messagequeue/mock"
3030
"github.com/uber/submitqueue/platform/metrics"
3131
stovepipemq "github.com/uber/submitqueue/stovepipe/core/messagequeue"
32+
"github.com/uber/submitqueue/stovepipe/core/requestlog"
33+
requestlogmock "github.com/uber/submitqueue/stovepipe/core/requestlog/mock"
3234
"github.com/uber/submitqueue/stovepipe/entity"
3335
"github.com/uber/submitqueue/stovepipe/extension/buildrunner"
3436
buildrunnermock "github.com/uber/submitqueue/stovepipe/extension/buildrunner/mock"
@@ -55,6 +57,8 @@ func queueContext() context.Context {
5557
type buildMocks struct {
5658
reqStore *storagemock.MockRequestStore
5759
buildStore *storagemock.MockBuildStore
60+
store *storagemock.MockStorage
61+
materializer *requestlogmock.MockMaterializer
5862
runnerFactory *buildrunnermock.MockFactory
5963
runner *buildrunnermock.MockBuildRunner
6064
publisher *mqmock.MockPublisher
@@ -74,15 +78,16 @@ func newController(t *testing.T, ctrl *gomock.Controller) (*Controller, buildMoc
7478
m := buildMocks{
7579
reqStore: storagemock.NewMockRequestStore(ctrl),
7680
buildStore: storagemock.NewMockBuildStore(ctrl),
81+
store: storagemock.NewMockStorage(ctrl),
82+
materializer: requestlogmock.NewMockMaterializer(ctrl),
7783
runnerFactory: buildrunnermock.NewMockFactory(ctrl),
7884
runner: buildrunnermock.NewMockBuildRunner(ctrl),
7985
publisher: mqmock.NewMockPublisher(ctrl),
8086
metricsScope: scope,
8187
}
8288

83-
store := storagemock.NewMockStorage(ctrl)
84-
store.EXPECT().GetRequestStore().Return(m.reqStore).AnyTimes()
85-
store.EXPECT().GetBuildStore().Return(m.buildStore).AnyTimes()
89+
m.store.EXPECT().GetRequestStore().Return(m.reqStore).AnyTimes()
90+
m.store.EXPECT().GetBuildStore().Return(m.buildStore).AnyTimes()
8691

8792
queue := mqmock.NewMockQueue(ctrl)
8893
queue.EXPECT().Publisher().Return(m.publisher).AnyTimes()
@@ -92,10 +97,24 @@ func newController(t *testing.T, ctrl *gomock.Controller) (*Controller, buildMoc
9297
})
9398
require.NoError(t, err)
9499

95-
c := NewController(zap.NewNop().Sugar(), scope, staticStorageFactory{store: store}, m.runnerFactory, registry, stovepipemq.TopicKeyBuild, "stovepipe-build")
100+
c := NewController(zap.NewNop().Sugar(), scope, staticStorageFactory{store: m.store}, m.materializer, m.runnerFactory, registry, stovepipemq.TopicKeyBuild, "stovepipe-build")
96101
return c, m
97102
}
98103

104+
func expectBuildTriggered(m buildMocks) *gomock.Call {
105+
request := entity.Request{ID: testID, Queue: testQueue}
106+
return m.materializer.EXPECT().PersistLog(
107+
gomock.Any(),
108+
m.store,
109+
requestlog.NewRequestEventLog(
110+
request,
111+
entity.RequestEventBuildTriggered,
112+
testBuildID,
113+
map[string]string{requestlog.MetadataKeyBuildID: testBuildID},
114+
),
115+
).Return(nil)
116+
}
117+
99118
func TestProcessTagsMetricsWithQueue(t *testing.T) {
100119
ctrl := gomock.NewController(t)
101120
c, m := newController(t, ctrl)
@@ -185,8 +204,9 @@ func TestProcess(t *testing.T) {
185204
Status: entity.BuildStatusAccepted,
186205
Version: 1,
187206
}
188-
m.buildStore.EXPECT().Create(gomock.Any(), build).Return(nil)
189-
m.publisher.EXPECT().Publish(gomock.Any(), "buildsignal", gomock.Any()).Return(nil)
207+
createCall := m.buildStore.EXPECT().Create(gomock.Any(), build).Return(nil)
208+
logCall := expectBuildTriggered(m).After(createCall)
209+
m.publisher.EXPECT().Publish(gomock.Any(), "buildsignal", gomock.Any()).Return(nil).After(logCall)
190210
},
191211
},
192212
{
@@ -202,8 +222,9 @@ func TestProcess(t *testing.T) {
202222
Status: entity.BuildStatusAccepted,
203223
Version: 1,
204224
}
205-
m.buildStore.EXPECT().Create(gomock.Any(), build).Return(nil)
206-
m.publisher.EXPECT().Publish(gomock.Any(), "buildsignal", gomock.Any()).Return(nil)
225+
createCall := m.buildStore.EXPECT().Create(gomock.Any(), build).Return(nil)
226+
logCall := expectBuildTriggered(m).After(createCall)
227+
m.publisher.EXPECT().Publish(gomock.Any(), "buildsignal", gomock.Any()).Return(nil).After(logCall)
207228
},
208229
},
209230
{
@@ -286,14 +307,14 @@ func TestProcess(t *testing.T) {
286307
},
287308
},
288309
{
289-
name: "already exists on create is swallowed and publish still happens",
310+
name: "duplicate runner build id is rejected",
311+
wantErr: true,
290312
setup: func(m buildMocks) {
291313
req := processingRequest(entity.BuildStrategyFull, "")
292314
m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(req, nil)
293315
m.runnerFactory.EXPECT().For(buildrunner.Config{QueueName: testQueue}).Return(m.runner, nil)
294316
m.runner.EXPECT().Trigger(gomock.Any(), "", testHeadURI, entity.BuildMetadata(nil)).Return(entity.BuildID{ID: testBuildID}, nil)
295317
m.buildStore.EXPECT().Create(gomock.Any(), gomock.Any()).Return(storage.ErrAlreadyExists)
296-
m.publisher.EXPECT().Publish(gomock.Any(), "buildsignal", gomock.Any()).Return(nil)
297318
},
298319
},
299320
{
@@ -308,6 +329,28 @@ func TestProcess(t *testing.T) {
308329
m.buildStore.EXPECT().Create(gomock.Any(), gomock.Any()).Return(errors.New("db down"))
309330
},
310331
},
332+
{
333+
name: "event persistence failure stops before buildsignal",
334+
wantErr: true,
335+
setup: func(m buildMocks) {
336+
req := processingRequest(entity.BuildStrategyFull, "")
337+
m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(req, nil)
338+
m.runnerFactory.EXPECT().For(buildrunner.Config{QueueName: testQueue}).Return(m.runner, nil)
339+
m.runner.EXPECT().Trigger(gomock.Any(), "", testHeadURI, entity.BuildMetadata(nil)).Return(entity.BuildID{ID: testBuildID}, nil)
340+
createCall := m.buildStore.EXPECT().Create(gomock.Any(), gomock.Any()).Return(nil)
341+
request := entity.Request{ID: testID, Queue: testQueue}
342+
m.materializer.EXPECT().PersistLog(
343+
gomock.Any(),
344+
m.store,
345+
requestlog.NewRequestEventLog(
346+
request,
347+
entity.RequestEventBuildTriggered,
348+
testBuildID,
349+
map[string]string{requestlog.MetadataKeyBuildID: testBuildID},
350+
),
351+
).Return(errors.New("db down")).After(createCall)
352+
},
353+
},
311354
{
312355
name: "publish failure is not retryable",
313356
wantErr: true,
@@ -317,8 +360,9 @@ func TestProcess(t *testing.T) {
317360
m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(req, nil)
318361
m.runnerFactory.EXPECT().For(buildrunner.Config{QueueName: testQueue}).Return(m.runner, nil)
319362
m.runner.EXPECT().Trigger(gomock.Any(), "", testHeadURI, entity.BuildMetadata(nil)).Return(entity.BuildID{ID: testBuildID}, nil)
320-
m.buildStore.EXPECT().Create(gomock.Any(), gomock.Any()).Return(nil)
321-
m.publisher.EXPECT().Publish(gomock.Any(), "buildsignal", gomock.Any()).Return(errors.New("queue down"))
363+
createCall := m.buildStore.EXPECT().Create(gomock.Any(), gomock.Any()).Return(nil)
364+
logCall := expectBuildTriggered(m).After(createCall)
365+
m.publisher.EXPECT().Publish(gomock.Any(), "buildsignal", gomock.Any()).Return(errors.New("queue down")).After(logCall)
322366
},
323367
},
324368
{

‎stovepipe/controller/buildsignal/BUILD.bazel‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -11,6 +11,7 @@ go_library(
1111
"//platform/publish:go_default_library",
1212
"//stovepipe/core/loader:go_default_library",
1313
"//stovepipe/core/messagequeue:go_default_library",
14+
"//stovepipe/core/requestlog:go_default_library",
1415
"//stovepipe/entity:go_default_library",
1516
"//stovepipe/extension/buildrunner:go_default_library",
1617
"//stovepipe/extension/storage:go_default_library",
@@ -31,6 +32,8 @@ go_test(
3132
"//platform/extension/messagequeue/mock:go_default_library",
3233
"//platform/metrics:go_default_library",
3334
"//stovepipe/core/messagequeue:go_default_library",
35+
"//stovepipe/core/requestlog:go_default_library",
36+
"//stovepipe/core/requestlog/mock:go_default_library",
3437
"//stovepipe/entity:go_default_library",
3538
"//stovepipe/extension/buildrunner:go_default_library",
3639
"//stovepipe/extension/buildrunner/mock:go_default_library",

0 commit comments

Comments
 (0)