From 981c311b66977952e66651d1ce04fdb47a8fa5a6 Mon Sep 17 00:00:00 2001 From: mnoah1 Date: Wed, 7 Oct 2026 16:07:15 +0000 Subject: [PATCH] feat(submitqueue): read List from receipt lookup MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Summary: Builds on #799 to serve List from immutable request_receipt keys and authoritative request_summary rows, instead of the duplicated queue-summary projection. - Scan receipt keys, then point-read summaries for visible results; keep bounds, string-ID ordering, and page-token format unchanged. - Use the existing RequestSummary response type and wire mapper. - Fail the page on missing, mismatched, or hidden summaries using the existing InternalConsistencyError. - Add controller and gateway integration coverage; update the RFC. Keep legacy projection writes for rollback. Deploy only after #799 reaches all writers. No backfill or legacy-read fallback: requests without receipt keys are absent from List. Duplicate-write and old-table removal remain follow-up work. None — no tracking issue was supplied. Test Plan: After CD rollout, confirm new Land requests appear in List with the same summary fields as ID lookup, check continuation pages, and monitor page latency. Revert Plan: Revert the read cutover to restore the old List path. Legacy projection writes remain active; this PR makes no schema changes. API Changes: No RPC/proto changes. The Go ListResult now contains RequestSummary values. --- .../submitqueue/gateway/server/mapper/list.go | 19 +- .../gateway/server/mapper/list_test.go | 17 +- submitqueue/entity/list.go | 2 +- submitqueue/entity/request_summary.go | 2 +- submitqueue/gateway/controller/BUILD.bazel | 1 + submitqueue/gateway/controller/list.go | 48 +++-- .../controller/list_consistency_test.go | 90 +++++++++ submitqueue/gateway/controller/list_test.go | 191 ++++++++++++------ .../submitqueue/gateway/suite_test.go | 61 +++++- 9 files changed, 329 insertions(+), 102 deletions(-) create mode 100644 submitqueue/gateway/controller/list_consistency_test.go diff --git a/service/submitqueue/gateway/server/mapper/list.go b/service/submitqueue/gateway/server/mapper/list.go index aca2613a7..a99a741f4 100644 --- a/service/submitqueue/gateway/server/mapper/list.go +++ b/service/submitqueue/gateway/server/mapper/list.go @@ -32,25 +32,8 @@ func ProtoToListRequest(req *pb.ListRequest) entity.ListRequest { // ListResultToProto maps the domain result to the wire response. func ListResultToProto(result entity.ListResult) *pb.ListResponse { - requests := make([]*pb.RequestSummary, 0, len(result.Requests)) - for _, summary := range result.Requests { - requests = append(requests, RequestQueueSummaryToProto(summary)) - } return &pb.ListResponse{ - Requests: requests, + Requests: RequestSummariesToProto(result.Requests), NextPageToken: result.NextPageToken, } } - -// RequestQueueSummaryToProto maps a queue projection to the wire request summary. -func RequestQueueSummaryToProto(summary entity.RequestQueueSummary) *pb.RequestSummary { - return &pb.RequestSummary{ - Sqid: summary.RequestID, - Queue: summary.Queue, - ChangeUris: append([]string{}, summary.ChangeURIs...), - ReceivedAtMs: summary.ReceivedAtMs, - Status: string(summary.Status), - LastError: summary.LastError, - Metadata: cloneStringMap(summary.Metadata), - } -} diff --git a/service/submitqueue/gateway/server/mapper/list_test.go b/service/submitqueue/gateway/server/mapper/list_test.go index dc9e77a38..1768e6001 100644 --- a/service/submitqueue/gateway/server/mapper/list_test.go +++ b/service/submitqueue/gateway/server/mapper/list_test.go @@ -42,19 +42,26 @@ func TestProtoToListRequest(t *testing.T) { func TestListResultToProto(t *testing.T) { result := entity.ListResult{ - Requests: []entity.RequestQueueSummary{{ + Requests: []entity.RequestSummary{{ RequestID: "1", Queue: "q", ChangeURIs: []string{"github://uber/repo/pull/1/abc"}, ReceivedAtMs: 100, - Status: entity.RequestStatusAccepted, - Metadata: map[string]string{}, + Status: entity.RequestStatusError, + LastError: "build failed", + Metadata: map[string]string{"build": "url"}, }}, NextPageToken: "next", } response := ListResultToProto(result) - assert.Equal(t, "1", response.Requests[0].Sqid) - assert.Equal(t, "next", response.NextPageToken) + assert.Equal(t, &pb.ListResponse{ + Requests: []*pb.RequestSummary{{ + Sqid: "1", Queue: "q", ChangeUris: []string{"github://uber/repo/pull/1/abc"}, + ReceivedAtMs: 100, Status: string(entity.RequestStatusError), LastError: "build failed", + Metadata: map[string]string{"build": "url"}, + }}, + NextPageToken: "next", + }, response) } diff --git a/submitqueue/entity/list.go b/submitqueue/entity/list.go index ffef78a7e..fd10ebc3a 100644 --- a/submitqueue/entity/list.go +++ b/submitqueue/entity/list.go @@ -31,7 +31,7 @@ type ListRequest struct { // ListResult contains one page of queue receipt history. type ListResult struct { // Requests are ordered by receipt time descending, then request ID descending. - Requests []RequestQueueSummary + Requests []RequestSummary // NextPageToken is an opaque continuation token. Empty means this is the last page. NextPageToken string } diff --git a/submitqueue/entity/request_summary.go b/submitqueue/entity/request_summary.go index f45051f17..f0d4730f2 100644 --- a/submitqueue/entity/request_summary.go +++ b/submitqueue/entity/request_summary.go @@ -57,7 +57,7 @@ type RequestSummary struct { Metadata map[string]string } -// RequestQueueSummary is the queue-ordered projection returned by List. +// RequestQueueSummary is a queue-ordered copy of a public request summary. type RequestQueueSummary struct { // RequestID is the canonical decimal request identifier, unique within Queue. RequestID string diff --git a/submitqueue/gateway/controller/BUILD.bazel b/submitqueue/gateway/controller/BUILD.bazel index d26b44d02..a56cdec6b 100644 --- a/submitqueue/gateway/controller/BUILD.bazel +++ b/submitqueue/gateway/controller/BUILD.bazel @@ -39,6 +39,7 @@ go_test( srcs = [ "cancel_test.go", "land_test.go", + "list_consistency_test.go", "list_queues_test.go", "list_test.go", "ping_test.go", diff --git a/submitqueue/gateway/controller/list.go b/submitqueue/gateway/controller/list.go index 8985d0f24..bbc86fc41 100644 --- a/submitqueue/gateway/controller/list.go +++ b/submitqueue/gateway/controller/list.go @@ -99,7 +99,7 @@ func (c *listController) List(ctx context.Context, req entity.ListRequest) (resu return entity.ListResult{}, fmt.Errorf("failed to resolve storage for queue %q: %w", req.Queue, err) } - query := basestorage.RequestQueueSummaryQuery{ + query := basestorage.RequestReceiptRange{ ReceivedAtOrAfterMs: req.ReceivedAtOrAfterMs, ReceivedBeforeMs: req.ReceivedBeforeMs, Limit: pageSize + 1, @@ -112,19 +112,27 @@ func (c *listController) List(ctx context.Context, req entity.ListRequest) (resu if token.Queue != req.Queue || token.ReceivedAtOrAfterMs != req.ReceivedAtOrAfterMs || token.ReceivedBeforeMs != req.ReceivedBeforeMs { return entity.ListResult{}, fmt.Errorf("page token does not match query: %w", ErrInvalidRequest) } - query.HasCursor = true - query.Cursor = basestorage.RequestQueueSummaryCursor{ReceivedAtMs: token.LastReceivedAtMs, RequestID: token.LastRequestID} + query.Before = basestorage.RequestReceiptCursor{ReceivedAtMs: token.LastReceivedAtMs, RequestID: token.LastRequestID} } - summaries, err := store.GetRequestQueueSummaryStore().List(ctx, query) + receipts, err := store.GetRequestReceiptStore().List(ctx, query) if err != nil { - return entity.ListResult{}, fmt.Errorf("failed to list queue=%s: %w", req.Queue, err) + return entity.ListResult{}, fmt.Errorf("failed to list request receipts queue=%q: %w", req.Queue, err) + } + + visible := receipts[:min(len(receipts), pageSize)] + result.Requests = make([]entity.RequestSummary, 0, len(visible)) + if len(visible) > 0 { + summaries := store.GetRequestSummaryStore() + for _, receipt := range visible { + summary, err := readListSummary(ctx, summaries, req.Queue, receipt) + if err != nil { + return entity.ListResult{}, err + } + result.Requests = append(result.Requests, summary) + } } - - visible := summaries - result = entity.ListResult{Requests: make([]entity.RequestQueueSummary, 0, min(len(visible), pageSize))} - if len(visible) > pageSize { - visible = visible[:pageSize] + if len(receipts) > pageSize { last := visible[len(visible)-1] result.NextPageToken = encodeListPageToken(listPageToken{ Queue: req.Queue, @@ -134,11 +142,27 @@ func (c *listController) List(ctx context.Context, req entity.ListRequest) (resu LastRequestID: last.RequestID, }) } - result.Requests = append(result.Requests, visible...) c.logger.Debugw("queue requests listed", "queue", req.Queue, "request_count", len(result.Requests), "has_next_page", result.NextPageToken != "") return result, nil } +func readListSummary(ctx context.Context, summaries basestorage.RequestSummaryStore, queue string, receipt entity.RequestReceipt) (entity.RequestSummary, error) { + if receipt.Queue != queue || receipt.ReceivedAtMs <= 0 || receipt.RequestID == "" { + return entity.RequestSummary{}, &InternalConsistencyError{Message: fmt.Sprintf("invalid request receipt queue=%q request_id=%q", queue, receipt.RequestID)} + } + summary, err := summaries.Get(ctx, receipt.RequestID) + if err != nil { + if basestorage.IsNotFound(err) { + return entity.RequestSummary{}, &InternalConsistencyError{Message: fmt.Sprintf("request summary missing for receipt queue=%q request_id=%q", queue, receipt.RequestID)} + } + return entity.RequestSummary{}, fmt.Errorf("failed to read request summary queue=%q request_id=%q: %w", queue, receipt.RequestID, err) + } + if summary.Queue != queue || summary.RequestID != receipt.RequestID || summary.ReceivedAtMs != receipt.ReceivedAtMs || summary.Status == entity.RequestStatusAccepting { + return entity.RequestSummary{}, &InternalConsistencyError{Message: fmt.Sprintf("request receipt disagrees with public summary queue=%q request_id=%q", queue, receipt.RequestID)} + } + return summary, nil +} + func encodeListPageToken(token listPageToken) string { values := url.Values{ "queue": {token.Queue}, @@ -178,7 +202,7 @@ func decodeListPageToken(encoded string) (listPageToken, error) { LastReceivedAtMs: lastReceivedAtMs, LastRequestID: values.Get("last_request_id"), } - if token.Queue == "" || token.LastRequestID == "" || token.ReceivedAtOrAfterMs >= token.ReceivedBeforeMs { + if token.Queue == "" || token.LastRequestID == "" || token.LastReceivedAtMs <= 0 || token.ReceivedAtOrAfterMs >= token.ReceivedBeforeMs { return listPageToken{}, fmt.Errorf("invalid token fields") } return token, nil diff --git a/submitqueue/gateway/controller/list_consistency_test.go b/submitqueue/gateway/controller/list_consistency_test.go new file mode 100644 index 000000000..3d3a56e26 --- /dev/null +++ b/submitqueue/gateway/controller/list_consistency_test.go @@ -0,0 +1,90 @@ +// Copyright (c) 2026 Uber Technologies, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package controller + +import ( + "context" + "fmt" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "github.com/uber/submitqueue/platform/errs" + "github.com/uber/submitqueue/submitqueue/entity" + basestorage "github.com/uber/submitqueue/submitqueue/extension/storage" + "go.uber.org/mock/gomock" +) + +func TestList_ConsistencyErrors(t *testing.T) { + receipt := entity.RequestReceipt{Queue: "q", RequestID: "1", ReceivedAtMs: 100} + summary := entity.RequestSummary{Queue: "q", RequestID: "1", ReceivedAtMs: 100, Status: entity.RequestStatusAccepted} + tests := []struct { + name string + changeReceipt func(*entity.RequestReceipt) + changeSummary func(*entity.RequestSummary) + readErr error + }{ + {name: "receipt from another queue", changeReceipt: func(r *entity.RequestReceipt) { r.Queue = "other" }}, + {name: "receipt has no request ID", changeReceipt: func(r *entity.RequestReceipt) { r.RequestID = "" }}, + {name: "receipt has no time", changeReceipt: func(r *entity.RequestReceipt) { r.ReceivedAtMs = 0 }}, + {name: "receipt has negative time", changeReceipt: func(r *entity.RequestReceipt) { r.ReceivedAtMs = -1 }}, + {name: "missing summary", readErr: fmt.Errorf("get summary: %w", basestorage.ErrNotFound)}, + {name: "summary from another queue", changeSummary: func(s *entity.RequestSummary) { s.Queue = "other" }}, + {name: "summary for another request", changeSummary: func(s *entity.RequestSummary) { s.RequestID = "2" }}, + {name: "summary has another receipt time", changeSummary: func(s *entity.RequestSummary) { s.ReceivedAtMs = 101 }}, + {name: "hidden accepting summary", changeSummary: func(s *entity.RequestSummary) { s.Status = entity.RequestStatusAccepting }}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + f := newListTestFixture(t) + badReceipt, badSummary := receipt, summary + if tt.changeReceipt != nil { + tt.changeReceipt(&badReceipt) + } else { + if tt.changeSummary != nil { + tt.changeSummary(&badSummary) + } + f.summaries.EXPECT().Get(gomock.Any(), "1").Return(badSummary, tt.readErr) + } + f.queueConfigs.EXPECT().Get(gomock.Any(), "q").Return(entity.QueueConfig{}, nil) + f.receipts.EXPECT().List(gomock.Any(), gomock.Any()).Return([]entity.RequestReceipt{badReceipt}, nil) + + result, err := f.controller.List(context.Background(), entity.ListRequest{Queue: "q", ReceivedAtOrAfterMs: 1, ReceivedBeforeMs: 200}) + + require.Error(t, err) + assert.True(t, IsInternalConsistency(err)) + assert.False(t, errs.IsUserError(err)) + assert.Equal(t, entity.ListResult{}, result) + }) + } +} + +func TestList_DoesNotReturnPartialPage(t *testing.T) { + f := newListTestFixture(t) + f.queueConfigs.EXPECT().Get(gomock.Any(), "q").Return(entity.QueueConfig{}, nil) + f.receipts.EXPECT().List(gomock.Any(), gomock.Any()).Return([]entity.RequestReceipt{ + {Queue: "q", RequestID: "2", ReceivedAtMs: 100}, + {Queue: "q", RequestID: "1", ReceivedAtMs: 90}, + }, nil) + f.summaries.EXPECT().Get(gomock.Any(), "2").Return(entity.RequestSummary{ + Queue: "q", RequestID: "2", ReceivedAtMs: 100, Status: entity.RequestStatusAccepted, + }, nil) + f.summaries.EXPECT().Get(gomock.Any(), "1").Return(entity.RequestSummary{}, basestorage.ErrNotFound) + + result, err := f.controller.List(context.Background(), entity.ListRequest{Queue: "q", ReceivedAtOrAfterMs: 1, ReceivedBeforeMs: 200}) + + assert.True(t, IsInternalConsistency(err)) + assert.Equal(t, entity.ListResult{}, result) +} diff --git a/submitqueue/gateway/controller/list_test.go b/submitqueue/gateway/controller/list_test.go index b7c79ae27..79cb652f5 100644 --- a/submitqueue/gateway/controller/list_test.go +++ b/submitqueue/gateway/controller/list_test.go @@ -17,12 +17,13 @@ package controller import ( "context" "encoding/base64" - "fmt" + "errors" "testing" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" "github.com/uber-go/tally" + "github.com/uber/submitqueue/platform/errs" "github.com/uber/submitqueue/submitqueue/entity" "github.com/uber/submitqueue/submitqueue/extension/queueconfig" qcmock "github.com/uber/submitqueue/submitqueue/extension/queueconfig/mock" @@ -34,55 +35,103 @@ import ( "go.uber.org/zap" ) -// listFactoryFor wraps a queue-summary store in a storage.Factory that -// resolves every queue to an aggregate exposing it. -func listFactoryFor(ctrl *gomock.Controller, store basestorage.RequestQueueSummaryStore) storage.Factory { +type listTestFixture struct { + controller ListController + receipts *storagemock.MockRequestReceiptStore + summaries *storagemock.MockRequestSummaryStore + queueConfigs *qcmock.MockStore +} + +func newListTestFixture(t *testing.T) listTestFixture { + ctrl := gomock.NewController(t) + f := listTestFixture{ + receipts: storagemock.NewMockRequestReceiptStore(ctrl), summaries: storagemock.NewMockRequestSummaryStore(ctrl), + queueConfigs: qcmock.NewMockStore(ctrl), + } agg := gwstoragemock.NewMockStorage(ctrl) - agg.EXPECT().GetRequestQueueSummaryStore().Return(store).AnyTimes() - f := gwstoragemock.NewMockFactory(ctrl) - f.EXPECT().For(gomock.Any()).Return(agg, nil).AnyTimes() + agg.EXPECT().GetRequestReceiptStore().Return(f.receipts).AnyTimes() + agg.EXPECT().GetRequestSummaryStore().Return(f.summaries).AnyTimes() + factory := gwstoragemock.NewMockFactory(ctrl) + factory.EXPECT().For(storage.Config{QueueName: "q"}).Return(agg, nil).AnyTimes() + f.controller = NewListController(zap.NewNop().Sugar(), tally.NoopScope, factory, f.queueConfigs) return f } -func TestList_ReturnsPageAndCursor(t *testing.T) { - ctrl := gomock.NewController(t) - store := storagemock.NewMockRequestQueueSummaryStore(ctrl) - store.EXPECT().List(gomock.Any(), basestorage.RequestQueueSummaryQuery{ - ReceivedAtOrAfterMs: 100, ReceivedBeforeMs: 200, Limit: 3, - }).Return([]entity.RequestQueueSummary{ - {RequestID: "3", Queue: "q", ChangeURIs: []string{}, ReceivedAtMs: 190, Status: entity.RequestStatusAccepted, Metadata: map[string]string{}}, - {RequestID: "2", Queue: "q", ChangeURIs: []string{}, ReceivedAtMs: 180, Status: entity.RequestStatusLanded, Metadata: map[string]string{}}, - {RequestID: "1", Queue: "q", ChangeURIs: []string{}, ReceivedAtMs: 170, Status: entity.RequestStatusError, Metadata: map[string]string{}}, - }, nil) - controller := newConfiguredListController(ctrl, store) - - result, err := controller.List(context.Background(), entity.ListRequest{Queue: "q", ReceivedAtOrAfterMs: 100, ReceivedBeforeMs: 200, PageSize: 2}) +func TestList_ReceiptPages(t *testing.T) { + receipts := []entity.RequestReceipt{ + {RequestID: "9", Queue: "q", ReceivedAtMs: 190}, + {RequestID: "10", Queue: "q", ReceivedAtMs: 190}, + {RequestID: "1", Queue: "q", ReceivedAtMs: 170}, + } + tests := []struct { + name string + receipts []entity.RequestReceipt + wantIDs []string + wantNext bool + }{ + {name: "empty", receipts: nil, wantIDs: []string{}}, + {name: "partial final page", receipts: receipts[:1], wantIDs: []string{"9"}}, + {name: "full final page", receipts: receipts[:2], wantIDs: []string{"9", "10"}}, + {name: "timestamp ties use string order", receipts: receipts, wantIDs: []string{"9", "10"}, wantNext: true}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + f := newListTestFixture(t) + f.queueConfigs.EXPECT().Get(gomock.Any(), "q").Return(entity.QueueConfig{}, nil) + f.receipts.EXPECT().List(gomock.Any(), basestorage.RequestReceiptRange{ + ReceivedAtOrAfterMs: 100, ReceivedBeforeMs: 200, Limit: 3, + }).Return(tt.receipts, nil) + want := make([]entity.RequestSummary, 0, len(tt.wantIDs)) + for _, receipt := range tt.receipts[:len(tt.wantIDs)] { + summary := entity.RequestSummary{ + RequestID: receipt.RequestID, Queue: receipt.Queue, ReceivedAtMs: receipt.ReceivedAtMs, + ChangeURIs: []string{"uri/" + receipt.RequestID}, Status: entity.RequestStatusError, + LastError: "build failed", Metadata: map[string]string{"build": "url"}, Version: 5, + } + f.summaries.EXPECT().Get(gomock.Any(), receipt.RequestID).Return(summary, nil) + want = append(want, summary) + } - require.NoError(t, err) - require.Len(t, result.Requests, 2) - assert.Equal(t, []string{"3", "2"}, []string{result.Requests[0].RequestID, result.Requests[1].RequestID}) - assert.Equal(t, int64(190), result.Requests[0].ReceivedAtMs) - require.NotEmpty(t, result.NextPageToken) - token, err := decodeListPageToken(result.NextPageToken) - require.NoError(t, err) - assert.Equal(t, int64(180), token.LastReceivedAtMs) - assert.Equal(t, "2", token.LastRequestID) + result, err := f.controller.List(context.Background(), entity.ListRequest{ + Queue: "q", ReceivedAtOrAfterMs: 100, ReceivedBeforeMs: 200, PageSize: 2, + }) + + require.NoError(t, err) + assert.Equal(t, want, result.Requests) + ids := make([]string, 0, len(result.Requests)) + for _, summary := range result.Requests { + ids = append(ids, summary.RequestID) + } + assert.Equal(t, tt.wantIDs, ids) + if !tt.wantNext { + assert.Empty(t, result.NextPageToken) + return + } + token, err := decodeListPageToken(result.NextPageToken) + require.NoError(t, err) + assert.Equal(t, listPageToken{ + Queue: "q", ReceivedAtOrAfterMs: 100, ReceivedBeforeMs: 200, + LastReceivedAtMs: 190, LastRequestID: "10", + }, token) + }) + } } -func TestList_UsesCursor(t *testing.T) { - ctrl := gomock.NewController(t) - store := storagemock.NewMockRequestQueueSummaryStore(ctrl) - token := encodeListPageToken(listPageToken{Queue: "q", ReceivedAtOrAfterMs: 100, ReceivedBeforeMs: 200, LastReceivedAtMs: 180, LastRequestID: "2"}) - store.EXPECT().List(gomock.Any(), basestorage.RequestQueueSummaryQuery{ +func TestList_UsesExistingCursor(t *testing.T) { + f := newListTestFixture(t) + f.queueConfigs.EXPECT().Get(gomock.Any(), "q").Return(entity.QueueConfig{}, nil) + token := encodeListPageToken(listPageToken{Queue: "q", ReceivedAtOrAfterMs: 100, ReceivedBeforeMs: 200, LastReceivedAtMs: 190, LastRequestID: "10"}) + f.receipts.EXPECT().List(gomock.Any(), basestorage.RequestReceiptRange{ ReceivedAtOrAfterMs: 100, ReceivedBeforeMs: 200, Limit: 51, - HasCursor: true, Cursor: basestorage.RequestQueueSummaryCursor{ReceivedAtMs: 180, RequestID: "2"}, - }).Return([]entity.RequestQueueSummary{}, nil) - controller := newConfiguredListController(ctrl, store) + Before: basestorage.RequestReceiptCursor{ReceivedAtMs: 190, RequestID: "10"}, + }).Return([]entity.RequestReceipt{{Queue: "q", ReceivedAtMs: 170, RequestID: "1"}}, nil) + summary := entity.RequestSummary{Queue: "q", ReceivedAtMs: 170, RequestID: "1", Status: entity.RequestStatusLanded} + f.summaries.EXPECT().Get(gomock.Any(), "1").Return(summary, nil) - result, err := controller.List(context.Background(), entity.ListRequest{Queue: "q", ReceivedAtOrAfterMs: 100, ReceivedBeforeMs: 200, PageToken: token}) + result, err := f.controller.List(context.Background(), entity.ListRequest{Queue: "q", ReceivedAtOrAfterMs: 100, ReceivedBeforeMs: 200, PageToken: token}) require.NoError(t, err) - assert.Empty(t, result.Requests) + assert.Equal(t, []entity.RequestSummary{summary}, result.Requests) assert.Empty(t, result.NextPageToken) } @@ -90,60 +139,78 @@ func TestList_Errors(t *testing.T) { validToken := encodeListPageToken(listPageToken{Queue: "other", ReceivedAtOrAfterMs: 100, ReceivedBeforeMs: 200, LastReceivedAtMs: 150, LastRequestID: "1"}) invalidFieldsToken := base64.RawURLEncoding.EncodeToString([]byte("queue=q&received_at_or_after_ms=100&received_before_ms=200&last_received_at_ms=150")) invalidNumberToken := base64.RawURLEncoding.EncodeToString([]byte("queue=q&received_at_or_after_ms=x&received_before_ms=200&last_received_at_ms=150&last_request_id=q%2F1")) - backendErr := fmt.Errorf("store down") + zeroTimeToken := encodeListPageToken(listPageToken{Queue: "q", ReceivedAtOrAfterMs: 100, ReceivedBeforeMs: 200, LastRequestID: "1"}) + negativeTimeToken := encodeListPageToken(listPageToken{Queue: "q", ReceivedAtOrAfterMs: 100, ReceivedBeforeMs: 200, LastReceivedAtMs: -1, LastRequestID: "1"}) + backendErr := errors.New("store down") tests := []struct { name string request entity.ListRequest - setup func(*storagemock.MockRequestQueueSummaryStore) + setup func(listTestFixture) + queueErr error + wantErr error wantInvalid bool wantUnknown bool }{ {name: "empty queue", request: entity.ListRequest{ReceivedAtOrAfterMs: 1, ReceivedBeforeMs: 2}, wantInvalid: true}, - {name: "unknown queue", request: entity.ListRequest{Queue: "missing", ReceivedAtOrAfterMs: 1, ReceivedBeforeMs: 2}, wantUnknown: true}, + {name: "unknown queue", request: entity.ListRequest{Queue: "missing", ReceivedAtOrAfterMs: 1, ReceivedBeforeMs: 2}, queueErr: queueconfig.ErrNotFound, wantUnknown: true, wantInvalid: true}, + {name: "queue config failure", request: entity.ListRequest{Queue: "q", ReceivedAtOrAfterMs: 1, ReceivedBeforeMs: 2}, queueErr: backendErr, wantErr: backendErr}, {name: "invalid range", request: entity.ListRequest{Queue: "q", ReceivedAtOrAfterMs: 2, ReceivedBeforeMs: 2}, wantInvalid: true}, {name: "negative page size", request: entity.ListRequest{Queue: "q", ReceivedAtOrAfterMs: 1, ReceivedBeforeMs: 2, PageSize: -1}, wantInvalid: true}, {name: "page size above maximum", request: entity.ListRequest{Queue: "q", ReceivedAtOrAfterMs: 1, ReceivedBeforeMs: 2, PageSize: 201}, wantInvalid: true}, {name: "malformed token", request: entity.ListRequest{Queue: "q", ReceivedAtOrAfterMs: 1, ReceivedBeforeMs: 2, PageToken: "%%%"}, wantInvalid: true}, {name: "invalid token number", request: entity.ListRequest{Queue: "q", ReceivedAtOrAfterMs: 100, ReceivedBeforeMs: 200, PageToken: invalidNumberToken}, wantInvalid: true}, {name: "invalid token fields", request: entity.ListRequest{Queue: "q", ReceivedAtOrAfterMs: 100, ReceivedBeforeMs: 200, PageToken: invalidFieldsToken}, wantInvalid: true}, + {name: "zero cursor time", request: entity.ListRequest{Queue: "q", ReceivedAtOrAfterMs: 100, ReceivedBeforeMs: 200, PageToken: zeroTimeToken}, wantInvalid: true}, + {name: "negative cursor time", request: entity.ListRequest{Queue: "q", ReceivedAtOrAfterMs: 100, ReceivedBeforeMs: 200, PageToken: negativeTimeToken}, wantInvalid: true}, {name: "token query mismatch", request: entity.ListRequest{Queue: "q", ReceivedAtOrAfterMs: 100, ReceivedBeforeMs: 200, PageToken: validToken}, wantInvalid: true}, { - name: "store failure", - request: entity.ListRequest{Queue: "q", ReceivedAtOrAfterMs: 1, ReceivedBeforeMs: 2}, - setup: func(store *storagemock.MockRequestQueueSummaryStore) { - store.EXPECT().List(gomock.Any(), basestorage.RequestQueueSummaryQuery{ReceivedAtOrAfterMs: 1, ReceivedBeforeMs: 2, Limit: 51}).Return(nil, backendErr) + name: "receipt scan failure", request: entity.ListRequest{Queue: "q", ReceivedAtOrAfterMs: 1, ReceivedBeforeMs: 2}, wantErr: backendErr, + setup: func(f listTestFixture) { + f.receipts.EXPECT().List(gomock.Any(), basestorage.RequestReceiptRange{ReceivedAtOrAfterMs: 1, ReceivedBeforeMs: 2, Limit: 51}).Return(nil, backendErr) + }, + }, + { + name: "summary read failure", request: entity.ListRequest{Queue: "q", ReceivedAtOrAfterMs: 1, ReceivedBeforeMs: 2}, wantErr: backendErr, + setup: func(f listTestFixture) { + f.receipts.EXPECT().List(gomock.Any(), gomock.Any()).Return([]entity.RequestReceipt{{Queue: "q", ReceivedAtMs: 1, RequestID: "1"}}, nil) + f.summaries.EXPECT().Get(gomock.Any(), "1").Return(entity.RequestSummary{}, backendErr) }, }, } for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { - ctrl := gomock.NewController(t) - store := storagemock.NewMockRequestQueueSummaryStore(ctrl) - queueConfigs := qcmock.NewMockStore(ctrl) + f := newListTestFixture(t) if tt.request.Queue != "" { - if tt.wantUnknown { - queueConfigs.EXPECT().Get(gomock.Any(), tt.request.Queue).Return(entity.QueueConfig{}, queueconfig.ErrNotFound) - } else { - queueConfigs.EXPECT().Get(gomock.Any(), tt.request.Queue).Return(entity.QueueConfig{}, nil) - } + f.queueConfigs.EXPECT().Get(gomock.Any(), tt.request.Queue).Return(entity.QueueConfig{}, tt.queueErr) } if tt.setup != nil { - tt.setup(store) + tt.setup(f) } - controller := NewListController(zap.NewNop().Sugar(), tally.NoopScope, listFactoryFor(ctrl, store), queueConfigs) - _, err := controller.List(context.Background(), tt.request) + result, err := f.controller.List(context.Background(), tt.request) require.Error(t, err) - if tt.wantInvalid { - assert.True(t, IsInvalidRequest(err)) - } + assert.Equal(t, entity.ListResult{}, result) + assert.Equal(t, tt.wantInvalid, IsInvalidRequest(err)) assert.Equal(t, tt.wantUnknown, IsUnrecognizedQueue(err)) + if tt.wantErr != nil { + assert.ErrorIs(t, err, tt.wantErr) + assert.False(t, errs.IsUserError(err)) + } }) } } -func newConfiguredListController(ctrl *gomock.Controller, store basestorage.RequestQueueSummaryStore) ListController { +func TestList_StorageResolutionFailure(t *testing.T) { + ctrl := gomock.NewController(t) + factory := gwstoragemock.NewMockFactory(ctrl) + backendErr := errors.New("storage unavailable") + factory.EXPECT().For(storage.Config{QueueName: "q"}).Return(nil, backendErr) queueConfigs := qcmock.NewMockStore(ctrl) queueConfigs.EXPECT().Get(gomock.Any(), "q").Return(entity.QueueConfig{}, nil) - return NewListController(zap.NewNop().Sugar(), tally.NoopScope, listFactoryFor(ctrl, store), queueConfigs) + c := NewListController(zap.NewNop().Sugar(), tally.NoopScope, factory, queueConfigs) + + result, err := c.List(context.Background(), entity.ListRequest{Queue: "q", ReceivedAtOrAfterMs: 1, ReceivedBeforeMs: 2}) + + assert.ErrorIs(t, err, backendErr) + assert.Equal(t, entity.ListResult{}, result) } diff --git a/test/integration/submitqueue/gateway/suite_test.go b/test/integration/submitqueue/gateway/suite_test.go index b80f5eb36..9baef2281 100644 --- a/test/integration/submitqueue/gateway/suite_test.go +++ b/test/integration/submitqueue/gateway/suite_test.go @@ -183,9 +183,16 @@ func (s *GatewayIntegrationSuite) TestLandAPI() { err = s.queueDB.QueryRow("SELECT COUNT(*) FROM queue_messages WHERE tenant = ? AND id = ?", req.Queue, resp.Sqid).Scan(&msgCount) require.NoError(t, err, "failed to query queue messages") assert.Equal(t, 1, msgCount, "should have 1 message in queue") + + summary, err := s.client.GetRequestSummaryByID(s.ctx, &pb.GetRequestSummaryByIDRequest{Sqid: resp.Sqid, Queue: req.Queue}) + require.NoError(t, err) + listed, err := s.client.List(s.ctx, &pb.ListRequest{ + Queue: req.Queue, ReceivedAtOrAfterMs: summary.Request.ReceivedAtMs, ReceivedBeforeMs: summary.Request.ReceivedAtMs + 1, + }) + require.NoError(t, err) + assert.Contains(t, listed.Requests, summary.Request) } -// TestListAPI verifies the queue projection is exposed in deterministic receipt order. func (s *GatewayIntegrationSuite) TestListAPI() { t := s.T() store, err := mysqlstorage.NewStorage(s.db, tally.NoopScope) @@ -195,7 +202,8 @@ func (s *GatewayIntegrationSuite) TestListAPI() { require.NoError(t, err) for _, summary := range []entity.RequestSummary{ {RequestID: "901", Queue: "test-queue", ChangeURIs: []string{"uri/1"}, ReceivedAtMs: 100, Status: entity.RequestStatusAccepted, StatusTimestampMs: 100, Version: 1, Metadata: map[string]string{}}, - {RequestID: "902", Queue: "test-queue", ChangeURIs: []string{"uri/2"}, ReceivedAtMs: 200, Status: entity.RequestStatusLanded, StatusTimestampMs: 200, Version: 1, Metadata: map[string]string{}}, + {RequestID: "902", Queue: "test-queue", ChangeURIs: []string{"uri/2"}, ReceivedAtMs: 200, Status: entity.RequestStatusAccepted, StatusTimestampMs: 200, Version: 1, Metadata: map[string]string{}}, + {RequestID: "903", Queue: "test-queue", ChangeURIs: []string{"uri/3"}, ReceivedAtMs: 200, Status: entity.RequestStatusLanded, StatusTimestampMs: 200, Version: 1, Metadata: map[string]string{}}, } { publicStatus := summary.Status summary.Status = entity.RequestStatusAccepting @@ -209,20 +217,58 @@ func (s *GatewayIntegrationSuite) TestListAPI() { Metadata: map[string]string{}, })) } + oldOnly := entity.RequestSummary{RequestID: "905", Queue: "test-queue", ReceivedAtMs: 150, Status: entity.RequestStatusAccepted, Version: 1} + require.NoError(t, queueStore.GetRequestSummaryStore().Create(s.ctx, oldOnly)) + require.NoError(t, queueStore.GetRequestQueueSummaryStore().Create(s.ctx, entity.RequestQueueSummary{ + RequestID: oldOnly.RequestID, Queue: oldOnly.Queue, ReceivedAtMs: oldOnly.ReceivedAtMs, Status: oldOnly.Status, Version: oldOnly.Version, + })) + require.NoError(t, queueStore.GetRequestSummaryStore().Create(s.ctx, entity.RequestSummary{ + RequestID: "904", Queue: "test-queue", ReceivedAtMs: 180, Status: entity.RequestStatusAccepting, Version: 1, + })) resp, err := s.client.List(s.ctx, &pb.ListRequest{Queue: "test-queue", ReceivedAtOrAfterMs: 50, ReceivedBeforeMs: 250, PageSize: 1}) require.NoError(t, err) require.Len(t, resp.Requests, 1) - assert.Equal(t, "902", resp.Requests[0].Sqid) + assert.Equal(t, "903", resp.Requests[0].Sqid) assert.Equal(t, string(entity.RequestStatusLanded), resp.Requests[0].Status) require.NotEmpty(t, resp.NextPageToken) + // Simulate a newer authoritative write before the legacy projection catches up. + current, err := queueStore.GetRequestSummaryStore().Get(s.ctx, "902") + require.NoError(t, err) + updated := current + updated.Status = entity.RequestStatusError + updated.StatusTimestampMs = 300 + updated.LastError = "build failed" + updated.Metadata = map[string]string{"build": "url"} + require.NoError(t, queueStore.GetRequestSummaryStore().Update(s.ctx, updated, current.Version, current.Version+1)) + legacy, err := queueStore.GetRequestQueueSummaryStore().Get(s.ctx, 200, "902") + require.NoError(t, err) + assert.Equal(t, entity.RequestStatusAccepted, legacy.Status) + + resp, err = s.client.List(s.ctx, &pb.ListRequest{Queue: "test-queue", ReceivedAtOrAfterMs: 50, ReceivedBeforeMs: 250, PageSize: 1, PageToken: resp.NextPageToken}) + require.NoError(t, err) + require.Len(t, resp.Requests, 1) + assert.Equal(t, "902", resp.Requests[0].Sqid) + summary, err := s.client.GetRequestSummaryByID(s.ctx, &pb.GetRequestSummaryByIDRequest{Sqid: "902", Queue: "test-queue"}) + require.NoError(t, err) + assert.Equal(t, summary.Request, resp.Requests[0]) + assert.Equal(t, string(entity.RequestStatusError), resp.Requests[0].Status) + assert.Equal(t, "build failed", resp.Requests[0].LastError) + assert.Equal(t, map[string]string{"build": "url"}, resp.Requests[0].Metadata) + require.NotEmpty(t, resp.NextPageToken) + resp, err = s.client.List(s.ctx, &pb.ListRequest{Queue: "test-queue", ReceivedAtOrAfterMs: 50, ReceivedBeforeMs: 250, PageSize: 1, PageToken: resp.NextPageToken}) require.NoError(t, err) require.Len(t, resp.Requests, 1) assert.Equal(t, "901", resp.Requests[0].Sqid) assert.Equal(t, string(entity.RequestStatusAccepted), resp.Requests[0].Status) assert.Empty(t, resp.NextPageToken) + + resp, err = s.client.List(s.ctx, &pb.ListRequest{Queue: "test-queue", ReceivedAtOrAfterMs: 100, ReceivedBeforeMs: 200}) + require.NoError(t, err) + require.Len(t, resp.Requests, 1) + assert.Equal(t, "901", resp.Requests[0].Sqid) } // TestReadAPIErrorCodes verifies controller error classes reach stable gRPC codes. @@ -275,6 +321,15 @@ func (s *GatewayIntegrationSuite) TestReadAPIErrorCodes() { _, err = s.client.GetRequestSummaryByChangeURI(s.ctx, &pb.GetRequestSummaryByChangeURIRequest{ChangeUri: inconsistentChangeURI, Queue: "missing-summary"}) require.Error(t, err) assert.Equal(t, codes.Internal, status.Code(err)) + + receiptStore, err := store.For("test-queue") + require.NoError(t, err) + require.NoError(t, receiptStore.GetRequestReceiptStore().Create(s.ctx, entity.RequestReceipt{ + Queue: "test-queue", RequestID: "999", ReceivedAtMs: 500, + })) + _, err = s.client.List(s.ctx, &pb.ListRequest{Queue: "test-queue", ReceivedAtOrAfterMs: 500, ReceivedBeforeMs: 501}) + require.Error(t, err) + assert.Equal(t, codes.Internal, status.Code(err)) } // TestRequestLogConsumer verifies the gateway's log-topic consumer in isolation: