From b50531ac728e3a88021844d65a638eae060f6b9c Mon Sep 17 00:00:00 2001 From: mnoah1 Date: Mon, 5 Oct 2026 20:15:50 +0000 Subject: [PATCH] feat(stovepipe): persist List summary fields MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Summary: Intent: - Prepare existing request summaries for the approved List contract. - This is a standalone summary change; the List endpoint remains pending. Changes: - Persist acceptance time (0 when unknown) and the existing typed outcome reason. - Repair missing fields on replay without regressing newer state; retain CAS retries and idempotency. - Append the SQL columns and extend behavioral/storage tests. Deployment: - Additive migration is covered internally through UQL and must precede the new binary; no migration files are included here. --- 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 --- .../core/requestlog/materializer_test.go | 14 +- stovepipe/core/requestlog/request_summary.go | 81 +++-- .../core/requestlog/request_summary_test.go | 300 +++++++++++++++--- stovepipe/entity/request_summary.go | 4 + .../storage/mysql/request_summary_store.go | 14 +- .../mysql/request_summary_store_test.go | 16 +- .../storage/mysql/schema/request_summary.sql | 2 + .../extension/storage/mysql/storage_test.go | 4 +- 8 files changed, 353 insertions(+), 82 deletions(-) diff --git a/stovepipe/core/requestlog/materializer_test.go b/stovepipe/core/requestlog/materializer_test.go index 5b02ad25b..326f6fd60 100644 --- a/stovepipe/core/requestlog/materializer_test.go +++ b/stovepipe/core/requestlog/materializer_test.go @@ -105,7 +105,18 @@ func TestMaterializerPersistRequestStateLog(t *testing.T) { request.BaseURI = testBaseURI } if !tt.wantErr { - expectRequestSummaryUpdate(summaries) + summaries.EXPECT().Get(gomock.Any(), testRequestID).Return(seededRequestSummary(), nil) + summaries.EXPECT().Update(gomock.Any(), gomock.Any(), int32(1), int32(2)).DoAndReturn( + func(_ context.Context, summary entity.RequestSummary, _, _ int32) error { + assert.Equal(t, tt.outcomeReason, summary.OutcomeReason) + if tt.state == entity.RequestStateAccepted { + assert.Equal(t, testNowMs, summary.AcceptedAtMs) + } else { + assert.Zero(t, summary.AcceptedAtMs) + } + return nil + }, + ) store.EXPECT().Create(gomock.Any(), gomock.Any()).DoAndReturn(func(_ context.Context, entry entity.RequestLog) error { assert.NotEmpty(t, entry.ID) assert.Equal(t, testNowMs, entry.TimestampMs) @@ -157,6 +168,7 @@ func TestMaterializerExistingIdenticalOccurrenceIsSuccess(t *testing.T) { current.State = entity.RequestStateAccepted current.RequestVersion = 1 current.StateTimestampMs = testNowMs - 1000 + current.AcceptedAtMs = testNowMs - 1000 current.Version = 2 summaries.EXPECT().Get(gomock.Any(), testRequestID).Return(current, nil) diff --git a/stovepipe/core/requestlog/request_summary.go b/stovepipe/core/requestlog/request_summary.go index 3c4a2a336..8683f1057 100644 --- a/stovepipe/core/requestlog/request_summary.go +++ b/stovepipe/core/requestlog/request_summary.go @@ -81,40 +81,71 @@ func updateExistingRequestSummary( current entity.RequestSummary, log entity.RequestLog, ) error { - if current.RequestID != log.RequestID || current.Queue != log.Queue { - return fmt.Errorf( - "request summary identity conflicts with state log log_request_id=%q summary_request_id=%q log_queue=%q summary_queue=%q", - log.RequestID, current.RequestID, log.Queue, current.Queue, - ) + updated, err := projectRequestSummary(current, log) + if err != nil { + return err } - if current.RequestVersion > log.RequestVersion { + if updated == current { return nil } - if current.RequestVersion == log.RequestVersion { - if requestSummaryMatchesLog(current, log) { - return nil - } - return fmt.Errorf("request summary conflicts with state log request_id=%q version=%d", log.RequestID, log.RequestVersion) - } - - updated := current - updated.State = log.State - updated.RequestVersion = log.RequestVersion - updated.StateTimestampMs = log.TimestampMs - if baseURI, ok := log.Metadata[MetadataKeyBaseURI]; ok { - updated.BaseURI = baseURI - } oldVersion := current.Version newVersion := oldVersion + 1 - updated.Version = oldVersion if err := summaryStore.Update(ctx, updated, oldVersion, newVersion); err != nil { return fmt.Errorf("failed to update request summary request_id=%q: %w", log.RequestID, err) } return nil } -func requestSummaryMatchesLog(summary entity.RequestSummary, log entity.RequestLog) bool { - if summary.State != log.State || summary.StateTimestampMs != log.TimestampMs { +func projectRequestSummary(current entity.RequestSummary, log entity.RequestLog) (entity.RequestSummary, error) { + updated, err := projectRequestSummaryFacts(current, log) + if err != nil { + return entity.RequestSummary{}, err + } + return projectRequestSummaryState(updated, log) +} + +// Facts may be filled from any retained version; known values must agree. +func projectRequestSummaryFacts(summary entity.RequestSummary, log entity.RequestLog) (entity.RequestSummary, error) { + if summary.RequestID != log.RequestID || summary.Queue != log.Queue { + return entity.RequestSummary{}, fmt.Errorf( + "request summary identity conflicts with state log log_request_id=%q summary_request_id=%q log_queue=%q summary_queue=%q", + log.RequestID, summary.RequestID, log.Queue, summary.Queue, + ) + } + if log.State != entity.RequestStateAccepted { + return summary, nil + } + if summary.AcceptedAtMs != 0 && summary.AcceptedAtMs != log.TimestampMs { + return entity.RequestSummary{}, fmt.Errorf("request summary acceptance time conflicts with retained log request_id=%q", log.RequestID) + } + summary.AcceptedAtMs = log.TimestampMs + return summary, nil +} + +func projectRequestSummaryState(summary entity.RequestSummary, log entity.RequestLog) (entity.RequestSummary, error) { + if log.RequestVersion < summary.RequestVersion { + return summary, nil + } + if log.RequestVersion == summary.RequestVersion { + if !requestSummaryStateMatchesLog(summary, log) { + return entity.RequestSummary{}, fmt.Errorf("request summary conflicts with state log request_id=%q version=%d", log.RequestID, log.RequestVersion) + } + summary.OutcomeReason = log.OutcomeReason + return summary, nil + } + summary.State = log.State + summary.RequestVersion = log.RequestVersion + summary.StateTimestampMs = log.TimestampMs + summary.OutcomeReason = log.OutcomeReason + if baseURI, ok := log.Metadata[MetadataKeyBaseURI]; ok { + summary.BaseURI = baseURI + } + return summary, nil +} + +func requestSummaryStateMatchesLog(summary entity.RequestSummary, log entity.RequestLog) bool { + if summary.State != log.State || summary.StateTimestampMs != log.TimestampMs || + (summary.OutcomeReason != entity.RequestOutcomeReasonUnknown && summary.OutcomeReason != log.OutcomeReason) { return false } baseURI, updatesBaseURI := log.Metadata[MetadataKeyBaseURI] @@ -129,6 +160,10 @@ func requestSummaryFromRequestAndLog(request entity.Request, log entity.RequestL State: log.State, RequestVersion: log.RequestVersion, StateTimestampMs: log.TimestampMs, + OutcomeReason: log.OutcomeReason, + } + if log.State == entity.RequestStateAccepted { + summary.AcceptedAtMs = log.TimestampMs } if baseURI, ok := log.Metadata[MetadataKeyBaseURI]; ok { summary.BaseURI = baseURI diff --git a/stovepipe/core/requestlog/request_summary_test.go b/stovepipe/core/requestlog/request_summary_test.go index b6a033228..23aea201d 100644 --- a/stovepipe/core/requestlog/request_summary_test.go +++ b/stovepipe/core/requestlog/request_summary_test.go @@ -18,8 +18,8 @@ import ( "context" "errors" "testing" + "time" - "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" "github.com/uber/submitqueue/stovepipe/entity" "github.com/uber/submitqueue/stovepipe/extension/storage" @@ -48,24 +48,19 @@ func TestMaterializerRepairsSummaryAfterPartialWrite(t *testing.T) { } log := NewRequestStateLog(request, entity.RequestOutcomeReasonUnknown) - var retained entity.RequestLog - logs.EXPECT().Create(gomock.Any(), gomock.Any()).DoAndReturn(func(_ context.Context, candidate entity.RequestLog) error { - retained = candidate - return nil - }) + retained := log + retained.TimestampMs = testNowMs + logs.EXPECT().Create(gomock.Any(), retained).Return(nil) summaries.EXPECT().Get(gomock.Any(), testRequestID).Return(entity.RequestSummary{}, errors.New("summary unavailable")) require.Error(t, materializer.PersistLog(context.Background(), stores, log)) logs.EXPECT().Create(gomock.Any(), gomock.Any()).Return(storage.ErrAlreadyExists) logs.EXPECT().Get(gomock.Any(), testRequestID, retained.ID).Return(retained, nil) summaries.EXPECT().Get(gomock.Any(), testRequestID).Return(seededRequestSummary(), nil) - summaries.EXPECT().Update(gomock.Any(), gomock.Any(), int32(1), int32(2)).DoAndReturn(func(_ context.Context, summary entity.RequestSummary, _, _ int32) error { - assert.Equal(t, entity.RequestSummary{ - RequestID: testRequestID, Queue: testQueue, URI: testRequestURI, BaseURI: testBaseURI, - State: entity.RequestStateProcessing, RequestVersion: 2, StateTimestampMs: testNowMs, Version: 1, - }, summary) - return nil - }) + summaries.EXPECT().Update(gomock.Any(), entity.RequestSummary{ + RequestID: testRequestID, Queue: testQueue, URI: testRequestURI, BaseURI: testBaseURI, + State: entity.RequestStateProcessing, RequestVersion: 2, StateTimestampMs: testNowMs, Version: 1, + }, int32(1), int32(2)).Return(nil) require.NoError(t, materializer.PersistLog(context.Background(), stores, log)) } @@ -81,14 +76,30 @@ func TestMaterializerCreatesMissingSummaryFromRequest(t *testing.T) { BuildStrategy: entity.BuildStrategyIncrementalSinceGreen, BaseURI: testBaseURI, State: entity.RequestStateProcessing, Version: 2, }, nil) - summaries.EXPECT().Create(gomock.Any(), gomock.Any()).DoAndReturn(func(_ context.Context, summary entity.RequestSummary) error { - assert.Equal(t, entity.RequestSummary{ - RequestID: testRequestID, Queue: testQueue, URI: testRequestURI, - State: entity.RequestStateAccepted, RequestVersion: 1, StateTimestampMs: testNowMs, Version: 1, - }, summary) - return nil - }) + summaries.EXPECT().Create(gomock.Any(), entity.RequestSummary{ + RequestID: testRequestID, Queue: testQueue, URI: testRequestURI, + State: entity.RequestStateAccepted, RequestVersion: 1, StateTimestampMs: testNowMs, Version: 1, + AcceptedAtMs: testNowMs, + }).Return(nil) + + require.NoError(t, materializer.PersistLog(context.Background(), stores, log)) +} +func TestMaterializerCreatesMissingSummaryFromTerminalLog(t *testing.T) { + materializer, stores, logs, summaries, requests := newTestMaterializer(t) + request := entity.Request{ + ID: testRequestID, Queue: testQueue, URI: testRequestURI, + State: entity.RequestStateFailed, Version: 3, + } + log := NewRequestStateLog(request, entity.RequestOutcomeReasonBuildFailed) + logs.EXPECT().Create(gomock.Any(), gomock.Any()).Return(nil) + summaries.EXPECT().Get(gomock.Any(), testRequestID).Return(entity.RequestSummary{}, storage.ErrNotFound) + requests.EXPECT().Get(gomock.Any(), testRequestID).Return(request, nil) + summaries.EXPECT().Create(gomock.Any(), entity.RequestSummary{ + RequestID: testRequestID, Queue: testQueue, URI: testRequestURI, + State: entity.RequestStateFailed, RequestVersion: 3, StateTimestampMs: testNowMs, + OutcomeReason: entity.RequestOutcomeReasonBuildFailed, Version: 1, + }).Return(nil) require.NoError(t, materializer.PersistLog(context.Background(), stores, log)) } @@ -107,6 +118,7 @@ func TestMaterializerMissingSummaryCreateRace(t *testing.T) { current.State = entity.RequestStateAccepted current.RequestVersion = 1 current.StateTimestampMs = testNowMs + current.AcceptedAtMs = testNowMs summaries.EXPECT().Get(gomock.Any(), testRequestID).Return(current, nil) require.NoError(t, materializer.PersistLog(context.Background(), stores, log)) @@ -177,15 +189,10 @@ func TestMaterializerProjectsOnlyNewerRequestState(t *testing.T) { RequestID: testRequestID, Queue: testQueue, URI: testRequestURI, State: entity.RequestStateAccepted, RequestVersion: 1, StateTimestampMs: testNowMs - 1, Version: 4, }, nil) - summaries.EXPECT().Update(gomock.Any(), gomock.Any(), int32(4), int32(5)).DoAndReturn( - func(_ context.Context, updated entity.RequestSummary, _, _ int32) error { - assert.Equal(t, entity.RequestStateProcessing, updated.State) - assert.Equal(t, testBaseURI, updated.BaseURI) - assert.Equal(t, int32(2), updated.RequestVersion) - assert.Equal(t, int32(4), updated.Version) - return nil - }, - ) + summaries.EXPECT().Update(gomock.Any(), entity.RequestSummary{ + RequestID: testRequestID, Queue: testQueue, URI: testRequestURI, BaseURI: testBaseURI, + State: entity.RequestStateProcessing, RequestVersion: 2, StateTimestampMs: testNowMs, Version: 4, + }, int32(4), int32(5)).Return(nil) require.NoError(t, materializer.PersistLog(context.Background(), stores, log)) } @@ -196,25 +203,25 @@ func TestMaterializerRetriesRequestSummaryVersionMismatch(t *testing.T) { ID: testRequestID, Queue: testQueue, URI: testRequestURI, BuildStrategy: entity.BuildStrategyFull, State: entity.RequestStateProcessing, Version: 2, }, entity.RequestOutcomeReasonUnknown) - requestLogs.EXPECT().Create(gomock.Any(), gomock.Any()).Return(nil) - current := seededRequestSummary() current.State = entity.RequestStateAccepted current.RequestVersion = 1 current.StateTimestampMs = testNowMs - 1 current.Version = 4 - summaries.EXPECT().Get(gomock.Any(), testRequestID).Return(current, nil) - summaries.EXPECT().Update(gomock.Any(), gomock.Any(), int32(4), int32(5)).Return(storage.ErrVersionMismatch) - winner := current winner.Version = 5 - summaries.EXPECT().Get(gomock.Any(), testRequestID).Return(winner, nil) - summaries.EXPECT().Update(gomock.Any(), gomock.Any(), int32(5), int32(6)).DoAndReturn( - func(_ context.Context, updated entity.RequestSummary, _, _ int32) error { - assert.Equal(t, int32(5), updated.Version) - assert.Equal(t, entity.RequestStateProcessing, updated.State) - return nil - }, + updated := winner + updated.State = log.State + updated.RequestVersion = log.RequestVersion + updated.StateTimestampMs = testNowMs + firstAttempt := updated + firstAttempt.Version = current.Version + gomock.InOrder( + requestLogs.EXPECT().Create(gomock.Any(), gomock.Any()).Return(nil), + summaries.EXPECT().Get(gomock.Any(), testRequestID).Return(current, nil), + summaries.EXPECT().Update(gomock.Any(), firstAttempt, int32(4), int32(5)).Return(storage.ErrVersionMismatch), + summaries.EXPECT().Get(gomock.Any(), testRequestID).Return(winner, nil), + summaries.EXPECT().Update(gomock.Any(), updated, int32(5), int32(6)).Return(nil), ) require.NoError(t, materializer.PersistLog(context.Background(), stores, log)) @@ -265,6 +272,7 @@ func TestMaterializerIgnoresOlderStateForSummary(t *testing.T) { summaries.EXPECT().Get(gomock.Any(), testRequestID).Return(entity.RequestSummary{ RequestID: testRequestID, Queue: testQueue, URI: testRequestURI, BaseURI: testBaseURI, State: entity.RequestStateProcessing, RequestVersion: 2, StateTimestampMs: testNowMs + 1, Version: 2, + AcceptedAtMs: testNowMs, }, nil) require.NoError(t, materializer.PersistLog(context.Background(), stores, NewRequestStateLog(request, entity.RequestOutcomeReasonUnknown))) @@ -284,13 +292,213 @@ func TestMaterializerPreservesBaseURIWhenLogOmitsIt(t *testing.T) { current.RequestVersion = 2 current.StateTimestampMs = testNowMs - 1 current.Version = 4 + updated := current + updated.State = log.State + updated.RequestVersion = log.RequestVersion + updated.StateTimestampMs = testNowMs + updated.OutcomeReason = log.OutcomeReason summaries.EXPECT().Get(gomock.Any(), testRequestID).Return(current, nil) - summaries.EXPECT().Update(gomock.Any(), gomock.Any(), int32(4), int32(5)).DoAndReturn( - func(_ context.Context, updated entity.RequestSummary, _, _ int32) error { - assert.Equal(t, testBaseURI, updated.BaseURI) - return nil - }, + summaries.EXPECT().Update(gomock.Any(), updated, int32(4), int32(5)).Return(nil) + + require.NoError(t, materializer.PersistLog(context.Background(), stores, log)) +} + +func TestUpdateRequestSummaryAcceptanceTime(t *testing.T) { + log := entity.RequestLog{ + RequestID: testRequestID, Queue: testQueue, State: entity.RequestStateAccepted, + RequestVersion: 1, TimestampMs: testNowMs - 1000, + } + terminal := seededRequestSummary() + terminal.BaseURI = testBaseURI + terminal.State = entity.RequestStateFailed + terminal.RequestVersion = 3 + terminal.StateTimestampMs = testNowMs + terminal.OutcomeReason = entity.RequestOutcomeReasonBuildFailed + terminal.Version = 4 + accepted := seededRequestSummary() + accepted.State = log.State + accepted.RequestVersion = log.RequestVersion + accepted.StateTimestampMs = log.TimestampMs + + tests := []struct { + name string + current entity.RequestSummary + acceptedAtMs int64 + wantErr bool + }{ + {name: "older acceptance repairs timestamp only", current: terminal}, + {name: "known acceptance is idempotent", current: terminal, acceptedAtMs: log.TimestampMs}, + {name: "same version acceptance repairs timestamp", current: accepted}, + {name: "same version known acceptance is idempotent", current: accepted, acceptedAtMs: log.TimestampMs}, + {name: "conflicting acceptance is rejected", current: terminal, acceptedAtMs: log.TimestampMs - 1, wantErr: true}, + {name: "same version conflicting acceptance is rejected", current: accepted, acceptedAtMs: log.TimestampMs - 1, wantErr: true}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + current := tt.current + current.AcceptedAtMs = tt.acceptedAtMs + want := current + want.AcceptedAtMs = log.TimestampMs + summaries := storagemock.NewMockRequestSummaryStore(gomock.NewController(t)) + if !tt.wantErr && want != current { + summaries.EXPECT().Update(gomock.Any(), want, current.Version, current.Version+1).Return(nil) + } + + err := updateExistingRequestSummary(context.Background(), summaries, current, log) + if tt.wantErr { + require.Error(t, err) + } else { + require.NoError(t, err) + } + }) + } +} + +func TestUpdateRequestSummaryOutcomeReason(t *testing.T) { + baseline := seededRequestSummary() + baseline.BaseURI = testBaseURI + baseline.State = entity.RequestStateFailed + baseline.RequestVersion = 3 + baseline.StateTimestampMs = testNowMs + baseline.AcceptedAtMs = testNowMs - 1000 + baseline.Version = 4 + buildFailed := entity.RequestOutcomeReasonBuildFailed + processingFailed := entity.RequestOutcomeReasonProcessingFailed + + tests := []struct { + name string + currentReason entity.RequestOutcomeReason + logVersion int32 + wantReason entity.RequestOutcomeReason + wantErr bool + }{ + {name: "same version repairs reason", logVersion: 3, wantReason: buildFailed}, + {name: "known reason is idempotent", currentReason: buildFailed, logVersion: 3, wantReason: buildFailed}, + {name: "older reason cannot overwrite winner", currentReason: processingFailed, logVersion: 2, wantReason: processingFailed}, + {name: "conflicting reason is rejected", currentReason: processingFailed, logVersion: 3, wantErr: true}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + current := baseline + current.OutcomeReason = tt.currentReason + log := entity.RequestLog{ + RequestID: testRequestID, Queue: testQueue, State: entity.RequestStateFailed, + RequestVersion: tt.logVersion, TimestampMs: testNowMs, OutcomeReason: buildFailed, + } + want := current + want.OutcomeReason = tt.wantReason + summaries := storagemock.NewMockRequestSummaryStore(gomock.NewController(t)) + if !tt.wantErr && want != current { + summaries.EXPECT().Update(gomock.Any(), want, current.Version, current.Version+1).Return(nil) + } + + err := updateExistingRequestSummary(context.Background(), summaries, current, log) + if tt.wantErr { + require.Error(t, err) + } else { + require.NoError(t, err) + } + }) + } +} + +func TestUpdateRequestSummaryNewerState(t *testing.T) { + current := seededRequestSummary() + current.BaseURI = testBaseURI + current.State = entity.RequestStateAccepted + current.RequestVersion = 1 + current.StateTimestampMs = testNowMs - 1000 + current.AcceptedAtMs = testNowMs - 1000 + current.Version = 4 + + tests := []struct { + name string + metadata map[string]string + wantBaseURI string + }{ + {name: "omitted baseline is preserved", wantBaseURI: testBaseURI}, + {name: "full build clears baseline", metadata: map[string]string{MetadataKeyBaseURI: ""}}, + {name: "new baseline replaces previous", metadata: map[string]string{MetadataKeyBaseURI: "git://repo/new-base"}, wantBaseURI: "git://repo/new-base"}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + log := entity.RequestLog{ + RequestID: testRequestID, Queue: testQueue, State: entity.RequestStateFailed, + RequestVersion: 3, TimestampMs: testNowMs, OutcomeReason: entity.RequestOutcomeReasonBuildFailed, + Metadata: tt.metadata, + } + want := current + want.State = entity.RequestStateFailed + want.RequestVersion = 3 + want.StateTimestampMs = testNowMs + want.OutcomeReason = entity.RequestOutcomeReasonBuildFailed + want.BaseURI = tt.wantBaseURI + summaries := storagemock.NewMockRequestSummaryStore(gomock.NewController(t)) + summaries.EXPECT().Update(gomock.Any(), want, current.Version, current.Version+1).Return(nil) + + require.NoError(t, updateExistingRequestSummary(context.Background(), summaries, current, log)) + }) + } +} + +func TestMaterializerRetriesAcceptanceRepairAfterNewerStateWins(t *testing.T) { + materializer, stores, logs, summaries, _ := newTestMaterializer(t) + log := NewRequestStateLog(entity.Request{ + ID: testRequestID, Queue: testQueue, State: entity.RequestStateAccepted, Version: 1, + }, entity.RequestOutcomeReasonUnknown) + log.TimestampMs = testNowMs - 1000 + current := seededRequestSummary() + current.State = entity.RequestStateProcessing + current.RequestVersion = 2 + current.StateTimestampMs = testNowMs + updated := current + updated.AcceptedAtMs = log.TimestampMs + winner := current + winner.State = entity.RequestStateFailed + winner.RequestVersion = 3 + winner.StateTimestampMs++ + winner.OutcomeReason = entity.RequestOutcomeReasonBuildFailed + winner.Version = 2 + updatedWinner := winner + updatedWinner.AcceptedAtMs = log.TimestampMs + gomock.InOrder( + logs.EXPECT().Create(gomock.Any(), log).Return(storage.ErrAlreadyExists), + logs.EXPECT().Get(gomock.Any(), testRequestID, log.ID).Return(log, nil), + summaries.EXPECT().Get(gomock.Any(), testRequestID).Return(current, nil), + summaries.EXPECT().Update(gomock.Any(), updated, int32(1), int32(2)).Return(storage.ErrVersionMismatch), + summaries.EXPECT().Get(gomock.Any(), testRequestID).Return(winner, nil), + summaries.EXPECT().Update(gomock.Any(), updatedWinner, int32(2), int32(3)).Return(nil), ) + require.NoError(t, materializer.PersistLog(context.Background(), stores, log)) +} +func TestMaterializerRetriesAcceptanceRepairAfterPartialWrite(t *testing.T) { + materializer, stores, logs, summaries, _ := newTestMaterializer(t) + log := NewRequestStateLog(entity.Request{ + ID: testRequestID, Queue: testQueue, State: entity.RequestStateAccepted, Version: 1, + }, entity.RequestOutcomeReasonUnknown) + retained := log + retained.TimestampMs = testNowMs + gomock.InOrder( + logs.EXPECT().Create(gomock.Any(), retained).Return(nil), + summaries.EXPECT().Get(gomock.Any(), testRequestID).Return(entity.RequestSummary{}, errors.New("unavailable")), + ) + require.Error(t, materializer.PersistLog(context.Background(), stores, log)) + + materializer.now = func() time.Time { return time.UnixMilli(testNowMs + 1000) } + candidate := log + candidate.TimestampMs = testNowMs + 1000 + current := seededRequestSummary() + current.State = entity.RequestStateProcessing + current.RequestVersion = 2 + current.StateTimestampMs = testNowMs + 500 + updated := current + updated.AcceptedAtMs = retained.TimestampMs + gomock.InOrder( + logs.EXPECT().Create(gomock.Any(), candidate).Return(storage.ErrAlreadyExists), + logs.EXPECT().Get(gomock.Any(), testRequestID, log.ID).Return(retained, nil), + summaries.EXPECT().Get(gomock.Any(), testRequestID).Return(current, nil), + summaries.EXPECT().Update(gomock.Any(), updated, int32(1), int32(2)).Return(nil), + ) require.NoError(t, materializer.PersistLog(context.Background(), stores, log)) } diff --git a/stovepipe/entity/request_summary.go b/stovepipe/entity/request_summary.go index b42c198ef..18d2e5d44 100644 --- a/stovepipe/entity/request_summary.go +++ b/stovepipe/entity/request_summary.go @@ -30,6 +30,10 @@ type RequestSummary struct { RequestVersion int32 // StateTimestampMs is when the represented state was first retained, in Unix milliseconds. StateTimestampMs int64 + // AcceptedAtMs is the original acceptance timestamp in Unix milliseconds; zero means unknown and a positive value is immutable. + AcceptedAtMs int64 + // OutcomeReason is the reason for the represented state; empty means unavailable or inapplicable. + OutcomeReason RequestOutcomeReason // Version is the optimistic-lock version of this materialized view. Version int32 } diff --git a/stovepipe/extension/storage/mysql/request_summary_store.go b/stovepipe/extension/storage/mysql/request_summary_store.go index 697adc77f..965dfc297 100644 --- a/stovepipe/extension/storage/mysql/request_summary_store.go +++ b/stovepipe/extension/storage/mysql/request_summary_store.go @@ -48,8 +48,8 @@ func (s *requestSummaryStore) Create(ctx context.Context, summary entity.Request _, err := s.db.ExecContext(ctx, ` INSERT INTO request_summary ( - queue, request_id, uri, base_uri, state, request_version, state_timestamp_ms, version - ) VALUES (?, ?, ?, ?, ?, ?, ?, ?)`, + queue, request_id, uri, base_uri, state, request_version, state_timestamp_ms, accepted_at_ms, outcome_reason, version + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`, summary.Queue, summary.RequestID, summary.URI, @@ -57,6 +57,8 @@ func (s *requestSummaryStore) Create(ctx context.Context, summary entity.Request summary.State, summary.RequestVersion, summary.StateTimestampMs, + summary.AcceptedAtMs, + summary.OutcomeReason, summary.Version, ) if err != nil { @@ -73,7 +75,7 @@ func (s *requestSummaryStore) Get(ctx context.Context, requestID string) (ret en defer func() { op.Complete(retErr) }() err := s.db.QueryRowContext(ctx, ` - SELECT queue, request_id, uri, base_uri, state, request_version, state_timestamp_ms, version + SELECT queue, request_id, uri, base_uri, state, request_version, state_timestamp_ms, accepted_at_ms, outcome_reason, version FROM request_summary WHERE queue = ? AND request_id = ?`, s.queue, requestID, @@ -85,6 +87,8 @@ func (s *requestSummaryStore) Get(ctx context.Context, requestID string) (ret en &ret.State, &ret.RequestVersion, &ret.StateTimestampMs, + &ret.AcceptedAtMs, + &ret.OutcomeReason, &ret.Version, ) if errors.Is(err, sql.ErrNoRows) { @@ -106,12 +110,14 @@ func (s *requestSummaryStore) Update(ctx context.Context, summary entity.Request result, err := s.db.ExecContext(ctx, ` UPDATE request_summary - SET base_uri = ?, state = ?, request_version = ?, state_timestamp_ms = ?, version = ? + SET base_uri = ?, state = ?, request_version = ?, state_timestamp_ms = ?, accepted_at_ms = ?, outcome_reason = ?, version = ? WHERE queue = ? AND request_id = ? AND version = ?`, summary.BaseURI, summary.State, summary.RequestVersion, summary.StateTimestampMs, + summary.AcceptedAtMs, + summary.OutcomeReason, newVersion, summary.Queue, summary.RequestID, diff --git a/stovepipe/extension/storage/mysql/request_summary_store_test.go b/stovepipe/extension/storage/mysql/request_summary_store_test.go index beaaf01d5..c79669048 100644 --- a/stovepipe/extension/storage/mysql/request_summary_store_test.go +++ b/stovepipe/extension/storage/mysql/request_summary_store_test.go @@ -37,9 +37,11 @@ func testRequestSummary() entity.RequestSummary { Queue: testRequestSummaryQueue, URI: "git://repo/head", BaseURI: "git://repo/base", - State: entity.RequestStateProcessing, - RequestVersion: 2, + State: entity.RequestStateFailed, + RequestVersion: 3, StateTimestampMs: 1735689600000, + AcceptedAtMs: 1735689599000, + OutcomeReason: entity.RequestOutcomeReasonProcessingFailed, Version: 1, } } @@ -65,7 +67,7 @@ func TestRequestSummaryStoreCreate(t *testing.T) { summary: summary, setup: func(mock sqlmock.Sqlmock) { mock.ExpectExec("INSERT INTO request_summary"). - WithArgs(summary.Queue, summary.RequestID, summary.URI, summary.BaseURI, summary.State, summary.RequestVersion, summary.StateTimestampMs, summary.Version). + WithArgs(summary.Queue, summary.RequestID, summary.URI, summary.BaseURI, summary.State, summary.RequestVersion, summary.StateTimestampMs, summary.AcceptedAtMs, summary.OutcomeReason, summary.Version). WillReturnResult(sqlmock.NewResult(0, 1)) }, }, @@ -74,7 +76,7 @@ func TestRequestSummaryStoreCreate(t *testing.T) { summary: summary, setup: func(mock sqlmock.Sqlmock) { mock.ExpectExec("INSERT INTO request_summary"). - WithArgs(summary.Queue, summary.RequestID, summary.URI, summary.BaseURI, summary.State, summary.RequestVersion, summary.StateTimestampMs, summary.Version). + WithArgs(summary.Queue, summary.RequestID, summary.URI, summary.BaseURI, summary.State, summary.RequestVersion, summary.StateTimestampMs, summary.AcceptedAtMs, summary.OutcomeReason, summary.Version). WillReturnError(&mysql.MySQLError{Number: mysqlErrDuplicateEntry}) }, wantErr: true, @@ -124,8 +126,8 @@ func TestRequestSummaryStoreGet(t *testing.T) { { name: "found", setup: func(mock sqlmock.Sqlmock) { - rows := sqlmock.NewRows([]string{"queue", "request_id", "uri", "base_uri", "state", "request_version", "state_timestamp_ms", "version"}). - AddRow(want.Queue, want.RequestID, want.URI, want.BaseURI, want.State, want.RequestVersion, want.StateTimestampMs, want.Version) + rows := sqlmock.NewRows([]string{"queue", "request_id", "uri", "base_uri", "state", "request_version", "state_timestamp_ms", "accepted_at_ms", "outcome_reason", "version"}). + AddRow(want.Queue, want.RequestID, want.URI, want.BaseURI, want.State, want.RequestVersion, want.StateTimestampMs, want.AcceptedAtMs, want.OutcomeReason, want.Version) mock.ExpectQuery("SELECT queue, request_id, uri, base_uri, state").WithArgs(testRequestSummaryQueue, want.RequestID).WillReturnRows(rows) }, }, @@ -178,7 +180,7 @@ func TestRequestSummaryStoreUpdate(t *testing.T) { db, mock, store := setupRequestSummaryStoreTest(t) defer db.Close() expectation := mock.ExpectExec("UPDATE request_summary"). - WithArgs(summary.BaseURI, summary.State, summary.RequestVersion, summary.StateTimestampMs, newVersion, summary.Queue, summary.RequestID, oldVersion) + WithArgs(summary.BaseURI, summary.State, summary.RequestVersion, summary.StateTimestampMs, summary.AcceptedAtMs, summary.OutcomeReason, newVersion, summary.Queue, summary.RequestID, oldVersion) if tt.execErr != nil { expectation.WillReturnError(tt.execErr) } else { diff --git a/stovepipe/extension/storage/mysql/schema/request_summary.sql b/stovepipe/extension/storage/mysql/schema/request_summary.sql index ce096a511..0b10438c5 100644 --- a/stovepipe/extension/storage/mysql/schema/request_summary.sql +++ b/stovepipe/extension/storage/mysql/schema/request_summary.sql @@ -10,5 +10,7 @@ CREATE TABLE IF NOT EXISTS request_summary ( request_version INT NOT NULL, state_timestamp_ms BIGINT NOT NULL, version INT NOT NULL, + accepted_at_ms BIGINT NOT NULL DEFAULT 0, + outcome_reason VARCHAR(64) NOT NULL DEFAULT '', PRIMARY KEY (queue, request_id) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4; diff --git a/test/integration/stovepipe/extension/storage/mysql/storage_test.go b/test/integration/stovepipe/extension/storage/mysql/storage_test.go index aa243abe2..aa53019df 100644 --- a/test/integration/stovepipe/extension/storage/mysql/storage_test.go +++ b/test/integration/stovepipe/extension/storage/mysql/storage_test.go @@ -142,6 +142,7 @@ func (s *MySQLRequestSummaryStoreSuite) TestCreateGetAndUpdate() { summary := entity.RequestSummary{ RequestID: "1", Queue: "monorepo/main", URI: "git://repo/head", State: entity.RequestStateAccepted, RequestVersion: 1, StateTimestampMs: 1000, Version: 1, + AcceptedAtMs: 1000, } require.NoError(s.T(), s.store.Create(s.ctx, summary)) @@ -151,7 +152,8 @@ func (s *MySQLRequestSummaryStoreSuite) TestCreateGetAndUpdate() { updated := summary updated.BaseURI = "git://repo/base" - updated.State = entity.RequestStateProcessing + updated.State = entity.RequestStateFailed + updated.OutcomeReason = entity.RequestOutcomeReasonProcessingFailed updated.RequestVersion = 2 updated.StateTimestampMs = 2000 require.NoError(s.T(), s.store.Update(s.ctx, updated, 1, 2))