Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
19 changes: 1 addition & 18 deletions service/submitqueue/gateway/server/mapper/list.go
Original file line number Diff line number Diff line change
Expand Up @@ -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),
}
}
17 changes: 12 additions & 5 deletions service/submitqueue/gateway/server/mapper/list_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
}
2 changes: 1 addition & 1 deletion submitqueue/entity/list.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
}
2 changes: 1 addition & 1 deletion submitqueue/entity/request_summary.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
1 change: 1 addition & 0 deletions submitqueue/gateway/controller/BUILD.bazel
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
48 changes: 36 additions & 12 deletions submitqueue/gateway/controller/list.go
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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,
Expand All @@ -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},
Expand Down Expand Up @@ -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
Expand Down
90 changes: 90 additions & 0 deletions submitqueue/gateway/controller/list_consistency_test.go
Original file line number Diff line number Diff line change
@@ -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)
}
Loading
Loading