Skip to content

Commit 23fc1db

Browse files
committed
refactor(storage): replace complete queue summaries
Persist change URIs with the existing RequestQueueSummary replacement update and expand guarded round-trip coverage. Jira Issues CODEM-204
1 parent 4332a2d commit 23fc1db

4 files changed

Lines changed: 107 additions & 19 deletions

File tree

submitqueue/extension/storage/mysql/request_queue_summary_store.go

Lines changed: 5 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -17,7 +17,6 @@ package mysql
1717
import (
1818
"context"
1919
"database/sql"
20-
"encoding/json"
2120
"errors"
2221
"fmt"
2322

@@ -45,7 +44,7 @@ func (s *requestQueueSummaryStore) Create(ctx context.Context, summary entity.Re
4544

4645
changeURIsJSON, metadataJSON, err := marshalSummaryJSON(summary.ChangeURIs, summary.Metadata)
4746
if err != nil {
48-
return fmt.Errorf("failed to marshal queue summary request_id=%s: %w", summary.RequestID, err)
47+
return fmt.Errorf("failed to marshal queue summary metadata request_id=%s: %w", summary.RequestID, err)
4948
}
5049
_, err = s.db.ExecContext(ctx, `
5150
INSERT INTO request_summary_by_queue (
@@ -93,15 +92,15 @@ func (s *requestQueueSummaryStore) Update(ctx context.Context, summary entity.Re
9392
op := metrics.Begin(s.scope, "update", metrics.StorageLatencyBuckets)
9493
defer func() { op.Complete(retErr) }()
9594

96-
metadataJSON, err := json.Marshal(normalizeMetadata(summary.Metadata))
95+
changeURIsJSON, metadataJSON, err := marshalSummaryJSON(summary.ChangeURIs, summary.Metadata)
9796
if err != nil {
98-
return fmt.Errorf("failed to marshal queue summary metadata request_id=%s: %w", summary.RequestID, err)
97+
return fmt.Errorf("failed to marshal queue summary request_id=%s: %w", summary.RequestID, err)
9998
}
10099
result, err := s.db.ExecContext(ctx, `
101100
UPDATE request_summary_by_queue
102-
SET status = ?, version = ?, last_error = ?, metadata = ?
101+
SET change_uris = ?, status = ?, version = ?, last_error = ?, metadata = ?
103102
WHERE queue = ? AND received_at_ms = ? AND request_id = ? AND version = ?`,
104-
summary.Status, newVersion, summary.LastError, metadataJSON,
103+
changeURIsJSON, summary.Status, newVersion, summary.LastError, metadataJSON,
105104
summary.Queue, summary.ReceivedAtMs, summary.RequestID, oldVersion,
106105
)
107106
if err != nil {

submitqueue/extension/storage/mysql/request_queue_summary_store_test.go

Lines changed: 66 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -185,48 +185,103 @@ func TestRequestQueueSummaryStore_Update(t *testing.T) {
185185
summary := entity.RequestQueueSummary{
186186
RequestID: "monorepo/1",
187187
Queue: "monorepo",
188+
ChangeURIs: []string{"github://github.example.com/uber/submitqueue/pull/456/cafebabe"},
188189
ReceivedAtMs: 1000,
189190
Status: entity.RequestStatusValidated,
190-
LastError: "",
191+
LastError: "validation detail",
192+
Metadata: map[string]string{"result": "validated"},
191193
}
192194
const oldVersion, newVersion = int32(1), int32(2)
193195

194196
tests := []struct {
195197
name string
198+
summary entity.RequestQueueSummary
196199
setup func(mock sqlmock.Sqlmock)
197200
wantErr bool
198201
wantErrIs error
199202
}{
200203
{
201-
name: "success",
204+
name: "success replaces all non-key fields",
205+
summary: summary,
206+
setup: func(mock sqlmock.Sqlmock) {
207+
mock.ExpectExec("SET change_uris = \\?, status = \\?, version = \\?, last_error = \\?, metadata = \\?").
208+
WithArgs([]byte(`["github://github.example.com/uber/submitqueue/pull/456/cafebabe"]`),
209+
summary.Status, newVersion, summary.LastError, []byte(`{"result":"validated"}`),
210+
summary.Queue, summary.ReceivedAtMs, summary.RequestID, oldVersion).
211+
WillReturnResult(sqlmock.NewResult(0, 1))
212+
},
213+
},
214+
{
215+
name: "success normalizes nil collections",
216+
summary: entity.RequestQueueSummary{
217+
RequestID: summary.RequestID,
218+
Queue: summary.Queue,
219+
ReceivedAtMs: summary.ReceivedAtMs,
220+
Status: summary.Status,
221+
LastError: summary.LastError,
222+
},
223+
setup: func(mock sqlmock.Sqlmock) {
224+
mock.ExpectExec("SET change_uris = \\?, status = \\?, version = \\?, last_error = \\?, metadata = \\?").
225+
WithArgs([]byte(`[]`), summary.Status, newVersion, summary.LastError, []byte(`{}`),
226+
summary.Queue, summary.ReceivedAtMs, summary.RequestID, oldVersion).
227+
WillReturnResult(sqlmock.NewResult(0, 1))
228+
},
229+
},
230+
{
231+
name: "success persists empty collections",
232+
summary: entity.RequestQueueSummary{
233+
RequestID: summary.RequestID,
234+
Queue: summary.Queue,
235+
ChangeURIs: []string{},
236+
ReceivedAtMs: summary.ReceivedAtMs,
237+
Status: summary.Status,
238+
LastError: summary.LastError,
239+
Metadata: map[string]string{},
240+
},
202241
setup: func(mock sqlmock.Sqlmock) {
203-
mock.ExpectExec("UPDATE request_summary_by_queue").
204-
WithArgs(summary.Status, newVersion, summary.LastError, sqlmock.AnyArg(),
242+
mock.ExpectExec("SET change_uris = \\?, status = \\?, version = \\?, last_error = \\?, metadata = \\?").
243+
WithArgs([]byte(`[]`), summary.Status, newVersion, summary.LastError, []byte(`{}`),
205244
summary.Queue, summary.ReceivedAtMs, summary.RequestID, oldVersion).
206245
WillReturnResult(sqlmock.NewResult(0, 1))
207246
},
208247
},
209248
{
210-
name: "version mismatch",
249+
name: "version mismatch",
250+
summary: summary,
211251
setup: func(mock sqlmock.Sqlmock) {
212-
mock.ExpectExec("UPDATE request_summary_by_queue").
213-
WithArgs(summary.Status, newVersion, summary.LastError, sqlmock.AnyArg(),
252+
mock.ExpectExec("SET change_uris = \\?, status = \\?, version = \\?, last_error = \\?, metadata = \\?").
253+
WithArgs([]byte(`["github://github.example.com/uber/submitqueue/pull/456/cafebabe"]`),
254+
summary.Status, newVersion, summary.LastError, []byte(`{"result":"validated"}`),
214255
summary.Queue, summary.ReceivedAtMs, summary.RequestID, oldVersion).
215256
WillReturnResult(sqlmock.NewResult(0, 0))
216257
},
217258
wantErr: true,
218259
wantErrIs: storage.ErrVersionMismatch,
219260
},
220261
{
221-
name: "exec error",
262+
name: "exec error",
263+
summary: summary,
222264
setup: func(mock sqlmock.Sqlmock) {
223-
mock.ExpectExec("UPDATE request_summary_by_queue").
224-
WithArgs(summary.Status, newVersion, summary.LastError, sqlmock.AnyArg(),
265+
mock.ExpectExec("SET change_uris = \\?, status = \\?, version = \\?, last_error = \\?, metadata = \\?").
266+
WithArgs([]byte(`["github://github.example.com/uber/submitqueue/pull/456/cafebabe"]`),
267+
summary.Status, newVersion, summary.LastError, []byte(`{"result":"validated"}`),
225268
summary.Queue, summary.ReceivedAtMs, summary.RequestID, oldVersion).
226269
WillReturnError(fmt.Errorf("connection reset"))
227270
},
228271
wantErr: true,
229272
},
273+
{
274+
name: "rows affected error",
275+
summary: summary,
276+
setup: func(mock sqlmock.Sqlmock) {
277+
mock.ExpectExec("SET change_uris = \\?, status = \\?, version = \\?, last_error = \\?, metadata = \\?").
278+
WithArgs([]byte(`["github://github.example.com/uber/submitqueue/pull/456/cafebabe"]`),
279+
summary.Status, newVersion, summary.LastError, []byte(`{"result":"validated"}`),
280+
summary.Queue, summary.ReceivedAtMs, summary.RequestID, oldVersion).
281+
WillReturnResult(sqlmock.NewErrorResult(fmt.Errorf("driver error")))
282+
},
283+
wantErr: true,
284+
},
230285
}
231286

232287
for _, tt := range tests {
@@ -236,7 +291,7 @@ func TestRequestQueueSummaryStore_Update(t *testing.T) {
236291

237292
tt.setup(mock)
238293

239-
err := store.Update(context.Background(), summary, oldVersion, newVersion)
294+
err := store.Update(context.Background(), tt.summary, oldVersion, newVersion)
240295
if tt.wantErr {
241296
require.Error(t, err)
242297
if tt.wantErrIs != nil {

submitqueue/extension/storage/request_queue_summary_store.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -54,7 +54,7 @@ type RequestQueueSummaryStore interface {
5454
// Get returns the row identified by its full primary key, or ErrNotFound when absent.
5555
Get(ctx context.Context, queue string, receivedAtMs int64, requestID string) (entity.RequestQueueSummary, error)
5656

57-
// Update conditionally replaces mutable fields when the persisted projection version equals oldVersion.
57+
// Update conditionally replaces all non-key fields when the persisted projection version equals oldVersion.
5858
// The store writes newVersion exactly as supplied and returns ErrVersionMismatch when the guard does not match.
5959
Update(ctx context.Context, summary entity.RequestQueueSummary, oldVersion, newVersion int32) error
6060

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

Lines changed: 35 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -540,15 +540,49 @@ func (s *StorageContractSuite) TestStorage_RequestQueueSummaryListAndCursor() {
540540
require.ErrorIs(t, err, storage.ErrNotFound)
541541

542542
got.Status = entity.RequestStatusLanded
543+
got.ChangeURIs = []string{"uri/replacement/1", "uri/replacement/2"}
543544
got.LastError = "done"
544545
got.Metadata = map[string]string{"result": "landed"}
545546
require.NoError(t, store.Update(ctx, got, 1, 2))
546-
require.ErrorIs(t, store.Update(ctx, got, 1, 3), storage.ErrVersionMismatch)
547547
updated, err := store.Get(ctx, got.Queue, got.ReceivedAtMs, got.RequestID)
548548
require.NoError(t, err)
549549
assert.Equal(t, int32(2), updated.Version)
550+
assert.Equal(t, []string{"uri/replacement/1", "uri/replacement/2"}, updated.ChangeURIs)
550551
assert.Equal(t, entity.RequestStatusLanded, updated.Status)
551552
assert.Equal(t, "done", updated.LastError)
553+
assert.Equal(t, map[string]string{"result": "landed"}, updated.Metadata)
554+
555+
stale := updated
556+
stale.ChangeURIs = []string{}
557+
stale.Status = entity.RequestStatusError
558+
stale.LastError = "stale"
559+
stale.Metadata = map[string]string{}
560+
require.ErrorIs(t, store.Update(ctx, stale, 1, 3), storage.ErrVersionMismatch)
561+
unchanged, err := store.Get(ctx, got.Queue, got.ReceivedAtMs, got.RequestID)
562+
require.NoError(t, err)
563+
assert.Equal(t, updated, unchanged)
564+
565+
updated.ChangeURIs = nil
566+
updated.Metadata = nil
567+
require.NoError(t, store.Update(ctx, updated, 2, 3))
568+
normalized, err := store.Get(ctx, got.Queue, got.ReceivedAtMs, got.RequestID)
569+
require.NoError(t, err)
570+
assert.Equal(t, int32(3), normalized.Version)
571+
assert.NotNil(t, normalized.ChangeURIs)
572+
assert.Empty(t, normalized.ChangeURIs)
573+
assert.NotNil(t, normalized.Metadata)
574+
assert.Empty(t, normalized.Metadata)
575+
576+
normalized.ChangeURIs = []string{}
577+
normalized.Metadata = map[string]string{}
578+
require.NoError(t, store.Update(ctx, normalized, 3, 4))
579+
emptyCollections, err := store.Get(ctx, got.Queue, got.ReceivedAtMs, got.RequestID)
580+
require.NoError(t, err)
581+
assert.Equal(t, int32(4), emptyCollections.Version)
582+
assert.NotNil(t, emptyCollections.ChangeURIs)
583+
assert.Empty(t, emptyCollections.ChangeURIs)
584+
assert.NotNil(t, emptyCollections.Metadata)
585+
assert.Empty(t, emptyCollections.Metadata)
552586

553587
firstPage, err := store.List(ctx, storage.RequestQueueSummaryQuery{
554588
Queue: "queue-summary", ReceivedAtOrAfterMs: 50, ReceivedBeforeMs: 250, Limit: 2,

0 commit comments

Comments
 (0)