Skip to content

Commit af820bd

Browse files
committed
refactor(storage): replace complete request summaries
Persist every non-key RequestSummary field while retaining optimistic version guards and full-row integration coverage. Jira Issues CODEM-204
1 parent efcd134 commit af820bd

4 files changed

Lines changed: 127 additions & 37 deletions

File tree

submitqueue/extension/storage/mysql/request_summary_store.go

Lines changed: 9 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -45,7 +45,7 @@ func (s *requestSummaryStore) Create(ctx context.Context, summary entity.Request
4545

4646
changeURIsJSON, metadataJSON, err := marshalSummaryJSON(summary.ChangeURIs, summary.Metadata)
4747
if err != nil {
48-
return fmt.Errorf("failed to marshal request summary request_id=%s: %w", summary.RequestID, err)
48+
return fmt.Errorf("failed to marshal request summary metadata request_id=%s: %w", summary.RequestID, err)
4949
}
5050

5151
_, err = s.db.ExecContext(ctx, `
@@ -103,19 +103,20 @@ func (s *requestSummaryStore) Update(ctx context.Context, summary entity.Request
103103
op := metrics.Begin(s.scope, "update", metrics.StorageLatencyBuckets)
104104
defer func() { op.Complete(retErr) }()
105105

106-
metadata := normalizeMetadata(summary.Metadata)
107-
metadataJSON, err := json.Marshal(metadata)
106+
changeURIsJSON, metadataJSON, err := marshalSummaryJSON(summary.ChangeURIs, summary.Metadata)
108107
if err != nil {
109-
return fmt.Errorf("failed to marshal request summary metadata request_id=%s: %w", summary.RequestID, err)
108+
return fmt.Errorf("failed to marshal request summary request_id=%s: %w", summary.RequestID, err)
110109
}
111110

112111
result, err := s.db.ExecContext(ctx, `
113112
UPDATE request_summary
114-
SET status = ?, request_version = ?, status_timestamp_ms = ?,
115-
version = ?, last_error = ?, metadata = ?
113+
SET queue = ?, change_uris = ?, received_at_ms = ?, status = ?,
114+
request_version = ?, status_timestamp_ms = ?, version = ?,
115+
last_error = ?, metadata = ?
116116
WHERE request_id = ? AND version = ?`,
117-
summary.Status, summary.RequestVersion, summary.StatusTimestampMs,
118-
newVersion, summary.LastError, metadataJSON,
117+
summary.Queue, changeURIsJSON, summary.ReceivedAtMs, summary.Status,
118+
summary.RequestVersion, summary.StatusTimestampMs, newVersion,
119+
summary.LastError, metadataJSON,
119120
summary.RequestID, oldVersion,
120121
)
121122
if err != nil {

submitqueue/extension/storage/mysql/request_summary_store_test.go

Lines changed: 55 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -199,51 +199,93 @@ func TestRequestSummaryStore_Get(t *testing.T) {
199199
func TestRequestSummaryStore_Update(t *testing.T) {
200200
summary := entity.RequestSummary{
201201
RequestID: "monorepo/1",
202-
Queue: "monorepo",
203-
ReceivedAtMs: 1000,
202+
Queue: "monorepo-updated",
203+
ChangeURIs: []string{"github://github.example.com/uber/submitqueue/pull/456/cafebabe"},
204+
ReceivedAtMs: 1500,
204205
Status: entity.RequestStatusValidated,
205206
RequestVersion: 2,
206207
StatusTimestampMs: 2000,
207-
LastError: "",
208+
Version: 1,
209+
LastError: "validation detail",
210+
Metadata: map[string]string{"result": "validated"},
208211
}
209212
const oldVersion, newVersion = int32(1), int32(2)
210213

211214
tests := []struct {
212215
name string
216+
summary entity.RequestSummary
213217
setup func(mock sqlmock.Sqlmock)
214218
wantErr bool
215219
wantErrIs error
216220
}{
217221
{
218-
name: "success",
222+
name: "success",
223+
summary: summary,
219224
setup: func(mock sqlmock.Sqlmock) {
220225
mock.ExpectExec("UPDATE request_summary").
221-
WithArgs(summary.Status, summary.RequestVersion, summary.StatusTimestampMs,
222-
newVersion, summary.LastError, sqlmock.AnyArg(), summary.RequestID, oldVersion).
226+
WithArgs(summary.Queue, []byte(`["github://github.example.com/uber/submitqueue/pull/456/cafebabe"]`),
227+
summary.ReceivedAtMs, summary.Status, summary.RequestVersion, summary.StatusTimestampMs,
228+
newVersion, summary.LastError, []byte(`{"result":"validated"}`), summary.RequestID, oldVersion).
223229
WillReturnResult(sqlmock.NewResult(0, 1))
224230
},
225231
},
226232
{
227-
name: "version mismatch",
233+
name: "version mismatch",
234+
summary: summary,
228235
setup: func(mock sqlmock.Sqlmock) {
229236
mock.ExpectExec("UPDATE request_summary").
230-
WithArgs(summary.Status, summary.RequestVersion, summary.StatusTimestampMs,
231-
newVersion, summary.LastError, sqlmock.AnyArg(), summary.RequestID, oldVersion).
237+
WithArgs(summary.Queue, []byte(`["github://github.example.com/uber/submitqueue/pull/456/cafebabe"]`),
238+
summary.ReceivedAtMs, summary.Status, summary.RequestVersion, summary.StatusTimestampMs,
239+
newVersion, summary.LastError, []byte(`{"result":"validated"}`), summary.RequestID, oldVersion).
232240
WillReturnResult(sqlmock.NewResult(0, 0))
233241
},
234242
wantErr: true,
235243
wantErrIs: storage.ErrVersionMismatch,
236244
},
237245
{
238-
name: "exec error",
246+
name: "exec error",
247+
summary: summary,
239248
setup: func(mock sqlmock.Sqlmock) {
240249
mock.ExpectExec("UPDATE request_summary").
241-
WithArgs(summary.Status, summary.RequestVersion, summary.StatusTimestampMs,
242-
newVersion, summary.LastError, sqlmock.AnyArg(), summary.RequestID, oldVersion).
250+
WithArgs(summary.Queue, []byte(`["github://github.example.com/uber/submitqueue/pull/456/cafebabe"]`),
251+
summary.ReceivedAtMs, summary.Status, summary.RequestVersion, summary.StatusTimestampMs,
252+
newVersion, summary.LastError, []byte(`{"result":"validated"}`), summary.RequestID, oldVersion).
243253
WillReturnError(fmt.Errorf("connection reset"))
244254
},
245255
wantErr: true,
246256
},
257+
{
258+
name: "rows affected error",
259+
summary: summary,
260+
setup: func(mock sqlmock.Sqlmock) {
261+
mock.ExpectExec("UPDATE request_summary").
262+
WithArgs(summary.Queue, []byte(`["github://github.example.com/uber/submitqueue/pull/456/cafebabe"]`),
263+
summary.ReceivedAtMs, summary.Status, summary.RequestVersion, summary.StatusTimestampMs,
264+
newVersion, summary.LastError, []byte(`{"result":"validated"}`), summary.RequestID, oldVersion).
265+
WillReturnResult(sqlmock.NewErrorResult(fmt.Errorf("rows unavailable")))
266+
},
267+
wantErr: true,
268+
},
269+
{
270+
name: "nil collections normalize to empty JSON",
271+
summary: entity.RequestSummary{
272+
RequestID: summary.RequestID,
273+
Queue: summary.Queue,
274+
ReceivedAtMs: summary.ReceivedAtMs,
275+
Status: summary.Status,
276+
RequestVersion: summary.RequestVersion,
277+
StatusTimestampMs: summary.StatusTimestampMs,
278+
Version: summary.Version,
279+
LastError: summary.LastError,
280+
},
281+
setup: func(mock sqlmock.Sqlmock) {
282+
mock.ExpectExec("UPDATE request_summary").
283+
WithArgs(summary.Queue, []byte(`[]`), summary.ReceivedAtMs, summary.Status,
284+
summary.RequestVersion, summary.StatusTimestampMs, newVersion, summary.LastError,
285+
[]byte(`{}`), summary.RequestID, oldVersion).
286+
WillReturnResult(sqlmock.NewResult(0, 1))
287+
},
288+
},
247289
}
248290

249291
for _, tt := range tests {
@@ -253,7 +295,7 @@ func TestRequestSummaryStore_Update(t *testing.T) {
253295

254296
tt.setup(mock)
255297

256-
err := store.Update(context.Background(), summary, oldVersion, newVersion)
298+
err := store.Update(context.Background(), tt.summary, oldVersion, newVersion)
257299
if tt.wantErr {
258300
require.Error(t, err)
259301
if tt.wantErrIs != nil {

submitqueue/extension/storage/request_summary_store.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -31,7 +31,7 @@ type RequestSummaryStore interface {
3131
// Get returns the summary for requestID, or ErrNotFound when absent.
3232
Get(ctx context.Context, requestID string) (entity.RequestSummary, error)
3333

34-
// Update conditionally replaces the mutable status fields when the persisted projection version equals oldVersion.
34+
// Update conditionally replaces every non-key field when the persisted projection version equals oldVersion.
3535
// The store writes newVersion exactly as supplied and returns ErrVersionMismatch when the guard does not match.
3636
Update(ctx context.Context, summary entity.RequestSummary, oldVersion, newVersion int32) error
3737
}

test/integration/submitqueue/extension/storage/suite.go

Lines changed: 62 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -442,33 +442,80 @@ func (s *StorageContractSuite) TestStorage_RequestSummaryCreateGetAndCAS() {
442442
ctx := s.ctx
443443
summary := entity.RequestSummary{
444444
RequestID: "summary/1", Queue: "summary-q", ChangeURIs: nil, ReceivedAtMs: 100,
445-
Status: entity.RequestStatusAccepted, StatusTimestampMs: 100, Version: 1, Metadata: nil,
445+
Status: entity.RequestStatusAccepted, RequestVersion: 1, StatusTimestampMs: 100, Version: 1,
446+
LastError: "", Metadata: nil,
446447
}
448+
store := s.storage.GetRequestSummaryStore()
447449

448-
require.NoError(t, s.storage.GetRequestSummaryStore().Create(ctx, summary))
449-
require.ErrorIs(t, s.storage.GetRequestSummaryStore().Create(ctx, summary), storage.ErrAlreadyExists)
450+
require.NoError(t, store.Create(ctx, summary))
451+
require.ErrorIs(t, store.Create(ctx, summary), storage.ErrAlreadyExists)
450452

451-
got, err := s.storage.GetRequestSummaryStore().Get(ctx, summary.RequestID)
453+
got, err := store.Get(ctx, summary.RequestID)
452454
require.NoError(t, err)
453-
assert.NotNil(t, got.ChangeURIs)
454-
assert.NotNil(t, got.Metadata)
455-
_, err = s.storage.GetRequestSummaryStore().Get(ctx, "summary/missing")
455+
assert.Equal(t, []string{}, got.ChangeURIs)
456+
assert.Equal(t, map[string]string{}, got.Metadata)
457+
_, err = store.Get(ctx, "summary/missing")
456458
require.ErrorIs(t, err, storage.ErrNotFound)
457459

460+
got.Queue = "summary-q-updated"
461+
got.ChangeURIs = []string{"change/updated"}
462+
got.ReceivedAtMs = 200
458463
got.Status = entity.RequestStatusLanded
459464
got.RequestVersion = 2
460-
got.StatusTimestampMs = 200
465+
got.StatusTimestampMs = 300
461466
got.LastError = "terminal detail"
462467
got.Metadata = map[string]string{"source": "test"}
463-
require.NoError(t, s.storage.GetRequestSummaryStore().Update(ctx, got, 1, 2))
464-
require.ErrorIs(t, s.storage.GetRequestSummaryStore().Update(ctx, got, 1, 3), storage.ErrVersionMismatch)
468+
require.NoError(t, store.Update(ctx, got, 1, 2))
465469

466-
updated, err := s.storage.GetRequestSummaryStore().Get(ctx, summary.RequestID)
470+
updated, err := store.Get(ctx, summary.RequestID)
467471
require.NoError(t, err)
468-
assert.Equal(t, int32(2), updated.Version)
469-
assert.Equal(t, entity.RequestStatusLanded, updated.Status)
470-
assert.Equal(t, "terminal detail", updated.LastError)
471-
assert.Equal(t, map[string]string{"source": "test"}, updated.Metadata)
472+
assert.Equal(t, entity.RequestSummary{
473+
RequestID: summary.RequestID,
474+
Queue: "summary-q-updated",
475+
ChangeURIs: []string{"change/updated"},
476+
ReceivedAtMs: 200,
477+
Status: entity.RequestStatusLanded,
478+
RequestVersion: 2,
479+
StatusTimestampMs: 300,
480+
Version: 2,
481+
LastError: "terminal detail",
482+
Metadata: map[string]string{"source": "test"},
483+
}, updated)
484+
485+
stale := updated
486+
stale.Queue = "stale-q"
487+
stale.ChangeURIs = []string{"change/stale"}
488+
stale.ReceivedAtMs = 400
489+
stale.Status = entity.RequestStatusError
490+
stale.RequestVersion = 3
491+
stale.StatusTimestampMs = 500
492+
stale.LastError = "stale detail"
493+
stale.Metadata = map[string]string{"source": "stale"}
494+
require.ErrorIs(t, store.Update(ctx, stale, 1, 3), storage.ErrVersionMismatch)
495+
496+
afterStale, err := store.Get(ctx, summary.RequestID)
497+
require.NoError(t, err)
498+
assert.Equal(t, updated, afterStale)
499+
500+
updated.ChangeURIs = nil
501+
updated.Metadata = nil
502+
require.NoError(t, store.Update(ctx, updated, 2, 3))
503+
504+
normalizedNil, err := store.Get(ctx, summary.RequestID)
505+
require.NoError(t, err)
506+
assert.Equal(t, []string{}, normalizedNil.ChangeURIs)
507+
assert.Equal(t, map[string]string{}, normalizedNil.Metadata)
508+
assert.Equal(t, int32(3), normalizedNil.Version)
509+
510+
normalizedNil.ChangeURIs = []string{}
511+
normalizedNil.Metadata = map[string]string{}
512+
require.NoError(t, store.Update(ctx, normalizedNil, 3, 4))
513+
514+
normalizedEmpty, err := store.Get(ctx, summary.RequestID)
515+
require.NoError(t, err)
516+
assert.Equal(t, []string{}, normalizedEmpty.ChangeURIs)
517+
assert.Equal(t, map[string]string{}, normalizedEmpty.Metadata)
518+
assert.Equal(t, int32(4), normalizedEmpty.Version)
472519
}
473520

474521
func (s *StorageContractSuite) TestStorage_RequestQueueSummaryListAndCursor() {

0 commit comments

Comments
 (0)