Skip to content

Commit b4def02

Browse files
committed
feat(stovepipe): record processing request state
Summary: Intent: - Retain the admitted request state before dependent Build work begins. - Let redelivery repair a missing processing occurrence from durable Request state. Changes: - Persist processing logs after a successful state transition and on already-processing retries. - Block hook and Build publication until materialization succeeds. - Wire the shared request-log materializer into Process and cover transition and retry failures. Revert Plan: - Revert this change to stop recording processing-state request-log entries. --- <sub>Generated by the 🪄 [pr-create](https://sg.uberinternal.com/code.uber.internal/uber-code/devexp-agent-marketplace/-/blob/claude-code/plugins/dev/uber-dev/skills/pr-create/SKILL.md) skill in devexp-agent-marketplace</sub>
1 parent 944e26f commit b4def02

6 files changed

Lines changed: 124 additions & 18 deletions

File tree

‎service/stovepipe/server/BUILD.bazel‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -81,6 +81,7 @@ go_test(
8181
"//api/base/hook:go_default_library",
8282
"//platform/consumer:go_default_library",
8383
"//stovepipe/controller/dlq:go_default_library",
84+
"//stovepipe/core/requestlog:go_default_library",
8485
"@com_github_stretchr_testify//assert:go_default_library",
8586
"@com_github_stretchr_testify//require:go_default_library",
8687
"@com_github_uber_go_tally//:go_default_library",

‎service/stovepipe/server/main.go‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -304,7 +304,7 @@ func run() error {
304304

305305
storageFty := storageFactory{backend: store}
306306
materializer := requestlog.NewMaterializer(scope)
307-
primaryCount, err := registerPrimaryControllers(primaryConsumer, logger.Sugar(), scope, storageFty, registry, sourceControl, brf, hookResolver{})
307+
primaryCount, err := registerPrimaryControllers(primaryConsumer, logger.Sugar(), scope, storageFty, materializer, registry, sourceControl, brf, hookResolver{})
308308
if err != nil {
309309
return err
310310
}
@@ -416,6 +416,7 @@ func registerPrimaryControllers(
416416
logger *zap.SugaredLogger,
417417
scope tally.Scope,
418418
store storage.Factory,
419+
materializer requestlog.Materializer,
419420
registry consumer.TopicRegistry,
420421
sourceControl sourcecontrol.Factory,
421422
brf buildrunner.Factory,
@@ -427,6 +428,7 @@ func registerPrimaryControllers(
427428
logger,
428429
scope,
429430
store,
431+
materializer,
430432
queueconfigdefault.NewStore(),
431433
sourceControl,
432434
registry,

‎service/stovepipe/server/main_test.go‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,7 @@ import (
2525
basehook "github.com/uber/submitqueue/api/base/hook"
2626
"github.com/uber/submitqueue/platform/consumer"
2727
"github.com/uber/submitqueue/stovepipe/controller/dlq"
28+
"github.com/uber/submitqueue/stovepipe/core/requestlog"
2829
"go.uber.org/zap/zaptest"
2930
)
3031

@@ -56,7 +57,7 @@ func registeredControllers(t *testing.T) (consumer.TopicRegistry, []consumer.Con
5657
primary := &recordingConsumer{}
5758
deadLetter := &recordingConsumer{}
5859

59-
_, err = registerPrimaryControllers(primary, logger, tally.NoopScope, store, registry,
60+
_, err = registerPrimaryControllers(primary, logger, tally.NoopScope, store, requestlog.NewMaterializer(tally.NoopScope), registry,
6061
fakeSourceControlFactory{}, fakeBuildRunnerFactory{}, hookResolver{})
6162
require.NoError(t, err)
6263

‎stovepipe/controller/process/BUILD.bazel‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,7 @@ go_library(
1515
"//stovepipe/core/hookevent:go_default_library",
1616
"//stovepipe/core/loader:go_default_library",
1717
"//stovepipe/core/messagequeue:go_default_library",
18+
"//stovepipe/core/requestlog:go_default_library",
1819
"//stovepipe/entity:go_default_library",
1920
"//stovepipe/extension/queueconfig:go_default_library",
2021
"//stovepipe/extension/sourcecontrol:go_default_library",
@@ -38,6 +39,8 @@ go_test(
3839
"//platform/metrics:go_default_library",
3940
"//stovepipe/core/hookevent:go_default_library",
4041
"//stovepipe/core/messagequeue:go_default_library",
42+
"//stovepipe/core/requestlog:go_default_library",
43+
"//stovepipe/core/requestlog/mock:go_default_library",
4144
"//stovepipe/entity:go_default_library",
4245
"//stovepipe/extension/queueconfig/default:go_default_library",
4346
"//stovepipe/extension/sourcecontrol:go_default_library",

‎stovepipe/controller/process/process.go‎

Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -33,6 +33,7 @@ import (
3333
"github.com/uber/submitqueue/stovepipe/core/hookevent"
3434
"github.com/uber/submitqueue/stovepipe/core/loader"
3535
stovepipemq "github.com/uber/submitqueue/stovepipe/core/messagequeue"
36+
"github.com/uber/submitqueue/stovepipe/core/requestlog"
3637
"github.com/uber/submitqueue/stovepipe/entity"
3738
"github.com/uber/submitqueue/stovepipe/extension/queueconfig"
3839
"github.com/uber/submitqueue/stovepipe/extension/sourcecontrol"
@@ -47,6 +48,7 @@ type Controller struct {
4748
logger *zap.SugaredLogger
4849
metricsScope tally.Scope
4950
stores storage.Factory
51+
materializer requestlog.Materializer
5052
queueConfigs queueconfig.Store
5153
sourceControl sourcecontrol.Factory
5254
registry consumer.TopicRegistry
@@ -65,6 +67,7 @@ func NewController(
6567
logger *zap.SugaredLogger,
6668
scope tally.Scope,
6769
stores storage.Factory,
70+
materializer requestlog.Materializer,
6871
queueConfigs queueconfig.Store,
6972
sourceControl sourcecontrol.Factory,
7073
registry consumer.TopicRegistry,
@@ -75,6 +78,7 @@ func NewController(
7578
logger: logger.Named("process_controller"),
7679
metricsScope: scope.SubScope("process_controller"),
7780
stores: stores,
81+
materializer: materializer,
7882
queueConfigs: queueConfigs,
7983
sourceControl: sourceControl,
8084
registry: registry,
@@ -116,6 +120,9 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) er
116120

117121
switch request.State {
118122
case entity.RequestStateProcessing:
123+
if err := c.persistProcessingLog(ctx, store, request); err != nil {
124+
return err
125+
}
119126
// Announce here as well as at admit: this is the only path a redelivery
120127
// takes once the transition is durable, so an admit that failed after
121128
// persisting would otherwise lose the start event for good. The event id
@@ -263,6 +270,10 @@ func (c *Controller) admitLatestHead(ctx context.Context, store storage.Storage,
263270
return nil
264271
}
265272

273+
if err := c.persistProcessingLog(ctx, store, request); err != nil {
274+
return err
275+
}
276+
266277
if err := c.publishHookEvent(ctx, request, hookevent.NewValidationRepositoryStarted(request)); err != nil {
267278
return err
268279
}
@@ -285,6 +296,14 @@ func (c *Controller) admitLatestHead(ctx context.Context, store storage.Storage,
285296
return nil
286297
}
287298

299+
func (c *Controller) persistProcessingLog(ctx context.Context, store storage.Storage, request entity.Request) error {
300+
log := requestlog.NewRequestStateLog(request, entity.RequestOutcomeReasonUnknown)
301+
if err := c.materializer.PersistLog(ctx, store, log); err != nil {
302+
return fmt.Errorf("failed to record processing state for request %s: %w", request.ID, err)
303+
}
304+
return nil
305+
}
306+
288307
// deriveBuildStrategy chooses the validation scope and baseline from the queue's last-known-good commit.
289308
// The caller resolves source control once and persists the returned values only after successfully claiming a build slot.
290309
func (c *Controller) deriveBuildStrategy(ctx context.Context, sc sourcecontrol.SourceControl, queueRow entity.Queue, request entity.Request) (strategy entity.BuildStrategy, baseURI string, err error) {

‎stovepipe/controller/process/process_test.go‎

Lines changed: 96 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -31,6 +31,8 @@ import (
3131
"github.com/uber/submitqueue/platform/metrics"
3232
"github.com/uber/submitqueue/stovepipe/core/hookevent"
3333
stovepipemq "github.com/uber/submitqueue/stovepipe/core/messagequeue"
34+
"github.com/uber/submitqueue/stovepipe/core/requestlog"
35+
requestlogmock "github.com/uber/submitqueue/stovepipe/core/requestlog/mock"
3436
"github.com/uber/submitqueue/stovepipe/entity"
3537
queueconfigdefault "github.com/uber/submitqueue/stovepipe/extension/queueconfig/default"
3638
"github.com/uber/submitqueue/stovepipe/extension/sourcecontrol"
@@ -58,6 +60,8 @@ type processMocks struct {
5860
queueStore *storagemock.MockQueueStore
5961
sourceFactory *sourcecontrolmock.MockFactory
6062
sourceControl *sourcecontrolmock.MockSourceControl
63+
store *storagemock.MockStorage
64+
materializer *requestlogmock.MockMaterializer
6165
publisher *mqmock.MockPublisher
6266
}
6367

@@ -79,12 +83,13 @@ func newControllerWithScope(t *testing.T, ctrl *gomock.Controller, scope tally.S
7983
queueStore: storagemock.NewMockQueueStore(ctrl),
8084
sourceFactory: sourcecontrolmock.NewMockFactory(ctrl),
8185
sourceControl: sourcecontrolmock.NewMockSourceControl(ctrl),
86+
store: storagemock.NewMockStorage(ctrl),
87+
materializer: requestlogmock.NewMockMaterializer(ctrl),
8288
publisher: mqmock.NewMockPublisher(ctrl),
8389
}
8490

85-
store := storagemock.NewMockStorage(ctrl)
86-
store.EXPECT().GetRequestStore().Return(m.reqStore).AnyTimes()
87-
store.EXPECT().GetQueueStore().Return(m.queueStore).AnyTimes()
91+
m.store.EXPECT().GetRequestStore().Return(m.reqStore).AnyTimes()
92+
m.store.EXPECT().GetQueueStore().Return(m.queueStore).AnyTimes()
8893

8994
queue := mqmock.NewMockQueue(ctrl)
9095
queue.EXPECT().Publisher().Return(m.publisher).AnyTimes()
@@ -98,7 +103,8 @@ func newControllerWithScope(t *testing.T, ctrl *gomock.Controller, scope tally.S
98103
c := NewController(
99104
zap.NewNop().Sugar(),
100105
scope,
101-
staticStorageFactory{store: store},
106+
staticStorageFactory{store: m.store},
107+
m.materializer,
102108
queueconfigdefault.NewStore(),
103109
m.sourceFactory,
104110
registry,
@@ -141,6 +147,7 @@ func TestProcessBuildPublishRequiresRegisteredTopic(t *testing.T) {
141147
m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(entity.Request{
142148
ID: testID, Queue: testQueue, State: entity.RequestStateProcessing, Version: 2,
143149
}, nil)
150+
expectProcessingLog(t, m, testID, 2)
144151

145152
err := c.Process(queueContext(testQueue), delivery(t, ctrl, processPayload(t, testID)))
146153

@@ -167,19 +174,41 @@ func expectAdmit(t *testing.T, m processMocks, id string) {
167174
expectStartAnnounceAndBuildPublish(t, m, id)
168175
}
169176

170-
// expectStartAnnounceAndBuildPublish expects the pair of publishes every admitted
171-
// request produces, in the order the controller performs them.
172177
func expectStartAnnounceAndBuildPublish(t *testing.T, m processMocks, id string) {
173178
t.Helper()
179+
expectProcessingLogAndHandoff(t, m, id, 2)
180+
}
181+
182+
func expectProcessingLogAndHandoff(t *testing.T, m processMocks, id string, version int32) {
183+
t.Helper()
174184

175-
expectStartValidationAnnounce(t, m, id)
176-
expectBuildPublish(t, m, id)
185+
logCall := expectProcessingLog(t, m, id, version)
186+
hookCall := expectStartValidationAnnounce(t, m, id)
187+
buildCall := expectBuildPublish(t, m, id)
188+
hookCall.After(logCall)
189+
buildCall.After(hookCall)
177190
}
178191

179-
func expectStartValidationAnnounce(t *testing.T, m processMocks, id string) {
192+
func expectProcessingLog(t *testing.T, m processMocks, id string, version int32) *gomock.Call {
180193
t.Helper()
181194

182-
m.publisher.EXPECT().
195+
request := entity.Request{
196+
ID: id,
197+
Queue: testQueue,
198+
State: entity.RequestStateProcessing,
199+
Version: version,
200+
}
201+
return m.materializer.EXPECT().PersistLog(
202+
gomock.Any(),
203+
m.store,
204+
requestlog.NewRequestStateLog(request, entity.RequestOutcomeReasonUnknown),
205+
).Return(nil)
206+
}
207+
208+
func expectStartValidationAnnounce(t *testing.T, m processMocks, id string) *gomock.Call {
209+
t.Helper()
210+
211+
return m.publisher.EXPECT().
183212
Publish(gomock.Any(), "stovepipe-hook", gomock.AssignableToTypeOf(entityqueue.Message{})).
184213
DoAndReturn(func(_ context.Context, _ string, msg entityqueue.Message) error {
185214
assert.Equal(t, id, msg.PartitionKey)
@@ -194,10 +223,10 @@ func expectStartValidationAnnounce(t *testing.T, m processMocks, id string) {
194223
})
195224
}
196225

197-
func expectBuildPublish(t *testing.T, m processMocks, id string) {
226+
func expectBuildPublish(t *testing.T, m processMocks, id string) *gomock.Call {
198227
t.Helper()
199228

200-
m.publisher.EXPECT().
229+
return m.publisher.EXPECT().
201230
Publish(gomock.Any(), "build", gomock.AssignableToTypeOf(entityqueue.Message{})).
202231
DoAndReturn(func(_ context.Context, _ string, msg entityqueue.Message) error {
203232
assert.Equal(t, id, msg.ID)
@@ -498,6 +527,21 @@ func TestProcess(t *testing.T) {
498527
expectStartAnnounceAndBuildPublish(t, m, testID)
499528
},
500529
},
530+
{
531+
name: "processing redelivery stops before handoff when log persistence fails",
532+
wantErr: true,
533+
setup: func(m processMocks) {
534+
request := entity.Request{
535+
ID: testID, Queue: testQueue, State: entity.RequestStateProcessing, Version: 2,
536+
}
537+
m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(request, nil)
538+
m.materializer.EXPECT().PersistLog(
539+
gomock.Any(),
540+
m.store,
541+
requestlog.NewRequestStateLog(request, entity.RequestOutcomeReasonUnknown),
542+
).Return(errors.New("db down"))
543+
},
544+
},
501545
{
502546
name: "unknown state is acked without retry",
503547
setup: func(m processMocks) {
@@ -539,9 +583,11 @@ func TestProcess(t *testing.T) {
539583
updatedRequest.State = entity.RequestStateProcessing
540584
updatedRequest.BuildStrategy = entity.BuildStrategyFull
541585
m.reqStore.EXPECT().Update(gomock.Any(), updatedRequest, int32(1), int32(2)).Return(nil)
586+
logCall := expectProcessingLog(t, m, testID, 2)
542587
m.publisher.EXPECT().
543588
Publish(gomock.Any(), "stovepipe-hook", gomock.AssignableToTypeOf(entityqueue.Message{})).
544-
Return(errors.New("queue unavailable"))
589+
Return(errors.New("queue unavailable")).
590+
After(logCall)
545591
},
546592
},
547593
{
@@ -565,10 +611,44 @@ func TestProcess(t *testing.T) {
565611
updatedRequest.State = entity.RequestStateProcessing
566612
updatedRequest.BuildStrategy = entity.BuildStrategyFull
567613
m.reqStore.EXPECT().Update(gomock.Any(), updatedRequest, int32(1), int32(2)).Return(nil)
568-
expectStartValidationAnnounce(t, m, testID)
614+
logCall := expectProcessingLog(t, m, testID, 2)
615+
hookCall := expectStartValidationAnnounce(t, m, testID)
569616
m.publisher.EXPECT().
570617
Publish(gomock.Any(), "build", gomock.AssignableToTypeOf(entityqueue.Message{})).
571-
Return(errors.New("queue unavailable"))
618+
Return(errors.New("queue unavailable")).
619+
After(hookCall)
620+
hookCall.After(logCall)
621+
},
622+
},
623+
{
624+
name: "processing log failure retains admitted request and claimed slot",
625+
wantErr: true,
626+
setup: func(m processMocks) {
627+
m.reqStore.EXPECT().Get(gomock.Any(), testID).Return(acceptedRequest(testID), nil)
628+
m.queueStore.EXPECT().Get(gomock.Any(), testQueue).Return(entity.Queue{
629+
Name: testQueue,
630+
LatestRequestID: testID,
631+
Version: 1,
632+
}, nil)
633+
updatedQueue := entity.Queue{
634+
Name: testQueue,
635+
LatestRequestID: testID,
636+
InFlightCount: 1,
637+
Version: 1,
638+
}
639+
m.queueStore.EXPECT().Update(gomock.Any(), updatedQueue, int32(1), int32(2)).Return(nil)
640+
updatedRequest := acceptedRequest(testID)
641+
updatedRequest.State = entity.RequestStateProcessing
642+
updatedRequest.BuildStrategy = entity.BuildStrategyFull
643+
m.reqStore.EXPECT().Update(gomock.Any(), updatedRequest, int32(1), int32(2)).Return(nil)
644+
request := entity.Request{
645+
ID: testID, Queue: testQueue, State: entity.RequestStateProcessing, Version: 2,
646+
}
647+
m.materializer.EXPECT().PersistLog(
648+
gomock.Any(),
649+
m.store,
650+
requestlog.NewRequestStateLog(request, entity.RequestOutcomeReasonUnknown),
651+
).Return(errors.New("db down"))
572652
},
573653
},
574654
{
@@ -710,7 +790,7 @@ func TestProcess(t *testing.T) {
710790
retry.BuildStrategy = entity.BuildStrategyIncrementalSinceGreen
711791
retry.BaseURI = lastGreenURI
712792
m.reqStore.EXPECT().Update(gomock.Any(), retry, int32(2), int32(3)).Return(nil)
713-
expectStartAnnounceAndBuildPublish(t, m, testID)
793+
expectProcessingLogAndHandoff(t, m, testID, 3)
714794
},
715795
},
716796
{

0 commit comments

Comments
 (0)