From 7ec48a524e5d1c2ac5d7002b5577dbdffe8635b2 Mon Sep 17 00:00:00 2001 From: mnoah1 Date: Tue, 6 Oct 2026 20:50:02 +0000 Subject: [PATCH] feat(stovepipe): implement List controller and pagination MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Summary: Build on #782 with queue-scoped List behavior: half-open acceptance windows, default page size 50 (maximum 200), versioned continuation tokens, and point reads of current request summaries. Tokens preserve the resolved window and bytewise ordering cursor across pages. Inconsistent mappings fail the page without partial results. RPC/server wiring remains a follow-up. ## Issue Links None — follow-up to the approved List API RFC (#775); no separate ticket. Test Plan: All 22 Stovepipe and service unit targets pass. make fmt gazelle lint check-tidy check-gazelle passes. Revert Plan: Revert this commit; no schema changes or exposed RPC behavior are introduced here. API Changes: Add the internal List request/result types and transport-independent controller. --- doc/rfc/stovepipe/list-api.md | 6 +- stovepipe/controller/BUILD.bazel | 5 + stovepipe/controller/list.go | 133 +++++++++++ stovepipe/controller/list_consistency_test.go | 117 ++++++++++ stovepipe/controller/list_pagination.go | 98 ++++++++ stovepipe/controller/list_pagination_test.go | 218 ++++++++++++++++++ stovepipe/controller/list_test.go | 197 ++++++++++++++++ stovepipe/controller/read_errors.go | 26 +++ stovepipe/entity/BUILD.bazel | 1 + stovepipe/entity/list.go | 41 ++++ 10 files changed, 839 insertions(+), 3 deletions(-) create mode 100644 stovepipe/controller/list.go create mode 100644 stovepipe/controller/list_consistency_test.go create mode 100644 stovepipe/controller/list_pagination.go create mode 100644 stovepipe/controller/list_pagination_test.go create mode 100644 stovepipe/controller/list_test.go create mode 100644 stovepipe/entity/list.go diff --git a/doc/rfc/stovepipe/list-api.md b/doc/rfc/stovepipe/list-api.md index ef83d1a0d..cbaf27f1a 100644 --- a/doc/rfc/stovepipe/list-api.md +++ b/doc/rfc/stovepipe/list-api.md @@ -1,6 +1,6 @@ # Stovepipe List API -The [protobuf contract](../../../api/stovepipe/proto/stovepipe.proto) is included for review; controller and storage implementation are deferred. +The [protobuf contract](../../../api/stovepipe/proto/stovepipe.proto) defines the approved API. Summary fields, acceptance-mapping storage, and the List controller are implemented; RPC/server wiring and end-to-end coverage remain. ## Proposal @@ -45,9 +45,9 @@ This is the wire representation of the existing domain `RequestSummary`, not ano The domain already defines a typed [RequestState](../../../stovepipe/entity/request.go). The wire field remains a string to match existing status/history APIs. States and reasons use the existing [public vocabulary](request-log.md#outcome-reasons); clients tolerate future values. Duplicate Ingest calls resolving to the same request produce one row. -## Data and Storage Work Required +## Data and Storage -The first six fields already exist in [RequestSummary](../../../stovepipe/entity/request_summary.go). Acceptance time and outcome reason already exist in retained logs but must be added to the summary. Acceptance time is acceptance-log time, not first RPC receipt time. No new producer signal or per-queue counter is needed. +All listed fields are materialized in [RequestSummary](../../../stovepipe/entity/request_summary.go), including acceptance time and outcome reason from retained logs. Acceptance time is acceptance-log time, not first RPC receipt time. No new producer signal or per-queue counter is needed. Listing uses an immutable mapping keyed by `(queue, accepted_at_ms, request_id)` for known positive acceptance times, then point-reads the corresponding summaries. This adds one small record per request and bounded extra reads, not another mutable status projection. Existing summary keys and public request IDs remain unchanged; no numeric-key migration is needed. diff --git a/stovepipe/controller/BUILD.bazel b/stovepipe/controller/BUILD.bazel index 38ff0fb08..6f25456aa 100644 --- a/stovepipe/controller/BUILD.bazel +++ b/stovepipe/controller/BUILD.bazel @@ -5,6 +5,8 @@ go_library( srcs = [ "get_project_status_by_uri.go", "ingest.go", + "list.go", + "list_pagination.go", "ping.go", "read_errors.go", "request_history.go", @@ -34,6 +36,9 @@ go_test( srcs = [ "get_project_status_by_uri_test.go", "ingest_test.go", + "list_consistency_test.go", + "list_pagination_test.go", + "list_test.go", "ping_test.go", "request_history_test.go", ], diff --git a/stovepipe/controller/list.go b/stovepipe/controller/list.go new file mode 100644 index 000000000..f799c31d6 --- /dev/null +++ b/stovepipe/controller/list.go @@ -0,0 +1,133 @@ +// 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" + "time" + + "github.com/uber-go/tally" + "github.com/uber/submitqueue/platform/metrics" + "github.com/uber/submitqueue/stovepipe/entity" + "github.com/uber/submitqueue/stovepipe/extension/storage" + "go.uber.org/zap" +) + +const ( + defaultListPageSize = 50 + maxListPageSize = 200 +) + +// ListController handles queue-scoped acceptance-time listing. +type ListController interface { + List(ctx context.Context, req entity.ListRequest) (entity.ListResult, error) +} + +type listController struct { + logger *zap.SugaredLogger + metricsScope tally.Scope + stores storage.Factory + configuredQueues map[string]struct{} + now func() time.Time +} + +// NewListController creates a controller restricted to the supplied configured queues. +func NewListController(logger *zap.SugaredLogger, scope tally.Scope, stores storage.Factory, configuredQueues []string) ListController { + queues := make(map[string]struct{}, len(configuredQueues)) + for _, queue := range configuredQueues { + queues[queue] = struct{}{} + } + return &listController{ + logger: logger, metricsScope: scope.SubScope("list_controller"), stores: stores, + configuredQueues: queues, now: time.Now, + } +} + +// List returns current summaries in descending acceptance-time order within a half-open time window. +// Omitted bounds resolve to [0, server now) on the first page and to the token's window on continuations. +func (c *listController) List(ctx context.Context, req entity.ListRequest) (result entity.ListResult, retErr error) { + op := metrics.Begin(c.metricsScope, "list", metrics.StorageLatencyBuckets, metrics.TagsFromContext(ctx)...) + defer func() { op.Complete(retErr) }() + + if req.Queue == "" { + return entity.ListResult{}, fmt.Errorf("List requires a queue: %w", ErrInvalidRequest) + } + if _, ok := c.configuredQueues[req.Queue]; !ok { + return entity.ListResult{}, fmt.Errorf("List queue %q is not configured: %w", req.Queue, ErrInvalidRequest) + } + query, err := c.resolveListRange(req) + if err != nil { + return entity.ListResult{}, err + } + stores, err := c.stores.For(storage.Config{QueueName: req.Queue}) + if err != nil { + return entity.ListResult{}, fmt.Errorf("List failed to resolve storage for queue %q: %w", req.Queue, err) + } + mappings, err := stores.GetRequestAcceptanceStore().List(ctx, query) + if err != nil { + return entity.ListResult{}, fmt.Errorf("List failed to read acceptance mappings for queue %q: %w", req.Queue, err) + } + pageSize := query.Limit - 1 + visible := mappings[:min(len(mappings), pageSize)] + result.Requests = make([]entity.RequestSummary, 0, len(visible)) + if len(visible) > 0 { + summaries := stores.GetRequestSummaryStore() + for _, mapping := range visible { + summary, err := readListSummary(ctx, summaries, req.Queue, mapping) + if err != nil { + return entity.ListResult{}, err + } + result.Requests = append(result.Requests, summary) + } + } + if len(mappings) > pageSize { + last := visible[len(visible)-1] + nextToken, err := encodeListPageToken(listPageToken{ + Version: listPageTokenVersion, Queue: req.Queue, + AcceptedAtOrAfterMs: query.AcceptedAtOrAfterMs, AcceptedBeforeMs: query.AcceptedBeforeMs, + LastAcceptedAtMs: last.AcceptedAtMs, LastRequestID: last.RequestID, + }) + if err != nil { + return entity.ListResult{}, fmt.Errorf("List failed to encode continuation: %w", err) + } + result.NextPageToken = nextToken + } + 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 storage.RequestSummaryStore, queue string, mapping entity.RequestAcceptance) (entity.RequestSummary, error) { + if mapping.Queue != queue || mapping.AcceptedAtMs <= 0 || mapping.RequestID == "" { + return entity.RequestSummary{}, &ListConsistencyError{ + Queue: queue, RequestID: mapping.RequestID, Reason: "invalid acceptance mapping", + } + } + summary, err := summaries.Get(ctx, mapping.RequestID) + if err != nil { + if storage.IsNotFound(err) { + return entity.RequestSummary{}, &ListConsistencyError{ + Queue: queue, RequestID: mapping.RequestID, Reason: "acceptance mapping has no summary", Err: err, + } + } + return entity.RequestSummary{}, fmt.Errorf("List failed to read summary queue=%q request_id=%q: %w", queue, mapping.RequestID, err) + } + if summary.Queue != queue || summary.RequestID != mapping.RequestID || summary.AcceptedAtMs != mapping.AcceptedAtMs { + return entity.RequestSummary{}, &ListConsistencyError{ + Queue: queue, RequestID: mapping.RequestID, Reason: "acceptance mapping disagrees with summary", + } + } + return summary, nil +} diff --git a/stovepipe/controller/list_consistency_test.go b/stovepipe/controller/list_consistency_test.go new file mode 100644 index 000000000..5fd40174c --- /dev/null +++ b/stovepipe/controller/list_consistency_test.go @@ -0,0 +1,117 @@ +// 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/require" + "github.com/uber/submitqueue/platform/errs" + "github.com/uber/submitqueue/stovepipe/entity" + "github.com/uber/submitqueue/stovepipe/extension/storage" + "go.uber.org/mock/gomock" +) + +func TestListRejectsInconsistentRecordsWithoutPartialResults(t *testing.T) { + first := listTestSummary(900, "9") + second := listTestSummary(800, "1") + mapping := entity.RequestAcceptance{Queue: second.Queue, AcceptedAtMs: second.AcceptedAtMs, RequestID: second.RequestID} + for _, tt := range []struct { + name string + change func(*entity.RequestSummary) + readErr error + wantReason string + }{ + { + name: "wrong queue", change: func(summary *entity.RequestSummary) { summary.Queue = "other/main" }, + wantReason: "acceptance mapping disagrees with summary", + }, + { + name: "wrong request", change: func(summary *entity.RequestSummary) { summary.RequestID = first.RequestID }, + wantReason: "acceptance mapping disagrees with summary", + }, + { + name: "unknown acceptance", change: func(summary *entity.RequestSummary) { summary.AcceptedAtMs = 0 }, + wantReason: "acceptance mapping disagrees with summary", + }, + { + name: "different acceptance", change: func(summary *entity.RequestSummary) { summary.AcceptedAtMs++ }, + wantReason: "acceptance mapping disagrees with summary", + }, + { + name: "missing summary", readErr: fmt.Errorf("summary read failed: %w", storage.ErrNotFound), + wantReason: "acceptance mapping has no summary", + }, + } { + t.Run(tt.name, func(t *testing.T) { + f := newListTestFixture(t) + corrupt := second + if tt.change != nil { + tt.change(&corrupt) + } + f.factory.EXPECT().For(gomock.Any()).Return(f.stores, nil) + f.acceptances.EXPECT().List(gomock.Any(), gomock.Any()).Return([]entity.RequestAcceptance{ + {Queue: first.Queue, AcceptedAtMs: first.AcceptedAtMs, RequestID: first.RequestID}, mapping, + }, nil) + gomock.InOrder( + f.summaries.EXPECT().Get(gomock.Any(), first.RequestID).Return(first, nil), + f.summaries.EXPECT().Get(gomock.Any(), second.RequestID).Return(corrupt, tt.readErr), + ) + result, err := f.controller.List(context.Background(), entity.ListRequest{Queue: first.Queue}) + require.True(t, IsListConsistency(err)) + require.True(t, IsListConsistency(fmt.Errorf("List failed: %w", err))) + var consistency *ListConsistencyError + require.ErrorAs(t, err, &consistency) + require.Equal(t, first.Queue, consistency.Queue) + require.Equal(t, second.RequestID, consistency.RequestID) + require.Equal(t, tt.wantReason, consistency.Reason) + require.Equal(t, tt.readErr, consistency.Err) + if tt.readErr != nil { + require.ErrorIs(t, err, tt.readErr) + require.ErrorIs(t, err, storage.ErrNotFound) + } + require.False(t, errs.IsRetryable(err)) + require.False(t, errs.IsUserError(err)) + require.Equal(t, entity.ListResult{}, result) + }) + } +} + +func TestListRejectsInvalidMappingBeforeSummaryLookup(t *testing.T) { + for _, mapping := range []entity.RequestAcceptance{ + {Queue: "other/main", AcceptedAtMs: 900, RequestID: "request/other/main/1"}, + {Queue: "monorepo/main", AcceptedAtMs: 0, RequestID: "request/monorepo/main/1"}, + {Queue: "monorepo/main", AcceptedAtMs: 900}, + } { + t.Run(mapping.RequestID, func(t *testing.T) { + f := newListTestFixture(t) + f.factory.EXPECT().For(gomock.Any()).Return(f.stores, nil) + f.acceptances.EXPECT().List(gomock.Any(), gomock.Any()).Return([]entity.RequestAcceptance{mapping}, nil) + result, err := f.controller.List(context.Background(), entity.ListRequest{Queue: "monorepo/main"}) + require.True(t, IsListConsistency(err)) + var consistency *ListConsistencyError + require.ErrorAs(t, err, &consistency) + require.Equal(t, "monorepo/main", consistency.Queue) + require.Equal(t, mapping.RequestID, consistency.RequestID) + require.Equal(t, "invalid acceptance mapping", consistency.Reason) + require.Nil(t, consistency.Err) + require.False(t, errs.IsRetryable(err)) + require.False(t, errs.IsUserError(err)) + require.Equal(t, entity.ListResult{}, result) + }) + } +} diff --git a/stovepipe/controller/list_pagination.go b/stovepipe/controller/list_pagination.go new file mode 100644 index 000000000..2bbcea6d2 --- /dev/null +++ b/stovepipe/controller/list_pagination.go @@ -0,0 +1,98 @@ +// 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 ( + "encoding/base64" + "encoding/json" + "fmt" + + "github.com/uber/submitqueue/stovepipe/entity" + "github.com/uber/submitqueue/stovepipe/extension/storage" +) + +const listPageTokenVersion = 1 + +type listPageToken struct { + Version int `json:"version"` + Queue string `json:"queue"` + AcceptedAtOrAfterMs int64 `json:"accepted_at_or_after_ms"` + AcceptedBeforeMs int64 `json:"accepted_before_ms"` + LastAcceptedAtMs int64 `json:"last_accepted_at_ms"` + LastRequestID string `json:"last_request_id"` +} + +func (c *listController) resolveListRange(req entity.ListRequest) (storage.RequestAcceptanceRange, error) { + if req.PageSize < 0 || req.PageSize > maxListPageSize { + return storage.RequestAcceptanceRange{}, fmt.Errorf("List page_size must be between 0 and %d: %w", maxListPageSize, ErrInvalidRequest) + } + pageSize := int(req.PageSize) + if pageSize == 0 { + pageSize = defaultListPageSize + } + query := storage.RequestAcceptanceRange{Limit: pageSize + 1} + if req.PageToken != "" { + token, err := decodeListPageToken(req.PageToken) + if err != nil { + return storage.RequestAcceptanceRange{}, err + } + if token.Queue != req.Queue || + (req.HasAcceptedAtOrAfterMs && req.AcceptedAtOrAfterMs != token.AcceptedAtOrAfterMs) || + (req.HasAcceptedBeforeMs && req.AcceptedBeforeMs != token.AcceptedBeforeMs) { + return storage.RequestAcceptanceRange{}, fmt.Errorf("List page_token does not match the queue and time bounds: %w", ErrInvalidRequest) + } + query.AcceptedAtOrAfterMs = token.AcceptedAtOrAfterMs + query.AcceptedBeforeMs = token.AcceptedBeforeMs + query.Before = storage.RequestAcceptanceCursor{AcceptedAtMs: token.LastAcceptedAtMs, RequestID: token.LastRequestID} + return query, nil + } + if req.HasAcceptedAtOrAfterMs { + query.AcceptedAtOrAfterMs = req.AcceptedAtOrAfterMs + } + if req.HasAcceptedBeforeMs { + query.AcceptedBeforeMs = req.AcceptedBeforeMs + } else { + query.AcceptedBeforeMs = c.now().UnixMilli() + } + if query.AcceptedAtOrAfterMs < 0 || query.AcceptedBeforeMs <= query.AcceptedAtOrAfterMs { + return storage.RequestAcceptanceRange{}, fmt.Errorf("List requires 0 <= accepted_at_or_after_ms < accepted_before_ms: %w", ErrInvalidRequest) + } + return query, nil +} + +func encodeListPageToken(token listPageToken) (string, error) { + contents, err := json.Marshal(token) + if err != nil { + return "", err + } + return base64.RawURLEncoding.EncodeToString(contents), nil +} + +func decodeListPageToken(encoded string) (listPageToken, error) { + contents, err := base64.RawURLEncoding.Strict().DecodeString(encoded) + if err != nil { + return listPageToken{}, fmt.Errorf("List invalid page_token encoding: %w", ErrInvalidRequest) + } + var token listPageToken + if err := json.Unmarshal(contents, &token); err != nil { + return listPageToken{}, fmt.Errorf("List invalid page_token contents: %w", ErrInvalidRequest) + } + if token.Version != listPageTokenVersion || token.Queue == "" || token.LastRequestID == "" || + token.AcceptedAtOrAfterMs < 0 || token.AcceptedBeforeMs <= token.AcceptedAtOrAfterMs || + token.LastAcceptedAtMs <= 0 || token.LastAcceptedAtMs < token.AcceptedAtOrAfterMs || token.LastAcceptedAtMs >= token.AcceptedBeforeMs { + return listPageToken{}, fmt.Errorf("List invalid page_token fields: %w", ErrInvalidRequest) + } + return token, nil +} diff --git a/stovepipe/controller/list_pagination_test.go b/stovepipe/controller/list_pagination_test.go new file mode 100644 index 000000000..bda6630ff --- /dev/null +++ b/stovepipe/controller/list_pagination_test.go @@ -0,0 +1,218 @@ +// 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 ( + "encoding/base64" + "testing" + "time" + + "github.com/stretchr/testify/require" + "github.com/uber/submitqueue/stovepipe/entity" + "github.com/uber/submitqueue/stovepipe/extension/storage" +) + +func explicitListWindow(lower, upper int64) entity.ListRequest { + return entity.ListRequest{ + Queue: "monorepo/main", AcceptedAtOrAfterMs: lower, HasAcceptedAtOrAfterMs: true, + AcceptedBeforeMs: upper, HasAcceptedBeforeMs: true, + } +} + +func TestResolveListRangeFirstPage(t *testing.T) { + for _, tt := range []struct { + name string + request entity.ListRequest + lower, upper int64 + limit int + clockCalls int + }{ + {"omitted bounds", entity.ListRequest{Queue: "monorepo/main"}, 0, 1000, 51, 1}, + {"lower only", entity.ListRequest{AcceptedAtOrAfterMs: 100, HasAcceptedAtOrAfterMs: true}, 100, 1000, 51, 1}, + {"upper only", entity.ListRequest{AcceptedBeforeMs: 2000, HasAcceptedBeforeMs: true}, 0, 2000, 51, 0}, + {"explicit window", explicitListWindow(100, 900), 100, 900, 51, 0}, + {"explicit zero lower", explicitListWindow(0, 900), 0, 900, 51, 0}, + {"future window", explicitListWindow(2000, 3000), 2000, 3000, 51, 0}, + {"one row", entity.ListRequest{PageSize: 1}, 0, 1000, 2, 1}, + {"maximum page", entity.ListRequest{PageSize: 200}, 0, 1000, 201, 1}, + } { + t.Run(tt.name, func(t *testing.T) { + clockCalls := 0 + controller := listController{now: func() time.Time { + clockCalls++ + return time.UnixMilli(1000) + }} + query, err := controller.resolveListRange(tt.request) + require.NoError(t, err) + require.Equal(t, storage.RequestAcceptanceRange{ + AcceptedAtOrAfterMs: tt.lower, AcceptedBeforeMs: tt.upper, Limit: tt.limit, + }, query) + require.Equal(t, tt.clockCalls, clockCalls) + }) + } +} + +func TestResolveListRangeRejectsInvalidRequest(t *testing.T) { + for _, tt := range []struct { + name string + request entity.ListRequest + }{ + {"negative lower", explicitListWindow(-1, 1000)}, + {"explicit zero upper", explicitListWindow(0, 0)}, + {"negative upper", explicitListWindow(0, -1)}, + {"equal bounds", explicitListWindow(1000, 1000)}, + {"inverted bounds", explicitListWindow(2000, 1000)}, + {"lower at default upper", entity.ListRequest{AcceptedAtOrAfterMs: 1000, HasAcceptedAtOrAfterMs: true}}, + {"negative page size", entity.ListRequest{PageSize: -1}}, + {"oversized page", entity.ListRequest{PageSize: 201}}, + } { + t.Run(tt.name, func(t *testing.T) { + controller := listController{now: func() time.Time { return time.UnixMilli(1000) }} + _, err := controller.resolveListRange(tt.request) + require.ErrorIs(t, err, ErrInvalidRequest) + }) + } +} + +func validListToken() listPageToken { + return listPageToken{ + Version: listPageTokenVersion, Queue: "monorepo/main", AcceptedAtOrAfterMs: 0, AcceptedBeforeMs: 1000, + LastAcceptedAtMs: 800, LastRequestID: "request/monorepo/main/9", + } +} + +func mustEncodeListToken(t *testing.T, token listPageToken) string { + t.Helper() + encoded, err := encodeListPageToken(token) + require.NoError(t, err) + return encoded +} + +func TestResolveListRangeContinuation(t *testing.T) { + token := validListToken() + encoded := mustEncodeListToken(t, token) + for _, tt := range []struct { + name string + request entity.ListRequest + limit int + }{ + {"bounds omitted", entity.ListRequest{Queue: token.Queue}, 51}, + {"explicit matching bounds", explicitListWindow(0, 1000), 51}, + {"explicit matching lower", entity.ListRequest{Queue: token.Queue, HasAcceptedAtOrAfterMs: true}, 51}, + {"explicit matching upper", entity.ListRequest{Queue: token.Queue, HasAcceptedBeforeMs: true, AcceptedBeforeMs: 1000}, 51}, + {"changed page size", entity.ListRequest{Queue: token.Queue, PageSize: 200}, 201}, + } { + t.Run(tt.name, func(t *testing.T) { + req := tt.request + req.PageToken = encoded + clockCalls := 0 + controller := listController{now: func() time.Time { + clockCalls++ + return time.UnixMilli(5000) + }} + query, err := controller.resolveListRange(req) + require.NoError(t, err) + require.Zero(t, clockCalls) + require.Equal(t, storage.RequestAcceptanceRange{ + AcceptedAtOrAfterMs: 0, AcceptedBeforeMs: 1000, Limit: tt.limit, + Before: storage.RequestAcceptanceCursor{AcceptedAtMs: 800, RequestID: token.LastRequestID}, + }, query) + }) + } +} + +func TestResolveListRangeRejectsMismatchedContinuation(t *testing.T) { + for _, tt := range []struct { + name string + request entity.ListRequest + tokenLower int64 + }{ + {"queue", entity.ListRequest{Queue: "other/main"}, 0}, + {"lower", explicitListWindow(1, 1000), 0}, + {"upper", explicitListWindow(0, 1001), 0}, + {"explicit zero lower", explicitListWindow(0, 1000), 100}, + } { + t.Run(tt.name, func(t *testing.T) { + req := tt.request + token := validListToken() + token.AcceptedAtOrAfterMs = tt.tokenLower + req.PageToken = mustEncodeListToken(t, token) + controller := listController{} + _, err := controller.resolveListRange(req) + require.ErrorIs(t, err, ErrInvalidRequest) + }) + } +} + +func TestResolveListRangePreservesCursorBoundaries(t *testing.T) { + for _, tt := range []struct { + name string + lower, upper, timestamp int64 + requestID string + }{ + {"inclusive lower", 800, 1000, 800, "request/monorepo/main/9"}, + {"opaque bytewise ID", 0, 1000, 800, "request/monorepo/main/é<&> "}, + {"large integer timestamps", 1<<63 - 100, 1<<63 - 1, 1<<63 - 2, "request/monorepo/main/9"}, + } { + t.Run(tt.name, func(t *testing.T) { + token := listPageToken{ + Version: listPageTokenVersion, Queue: "monorepo/main", + AcceptedAtOrAfterMs: tt.lower, AcceptedBeforeMs: tt.upper, + LastAcceptedAtMs: tt.timestamp, LastRequestID: tt.requestID, + } + controller := listController{} + query, err := controller.resolveListRange(entity.ListRequest{Queue: token.Queue, PageToken: mustEncodeListToken(t, token)}) + require.NoError(t, err) + require.Equal(t, storage.RequestAcceptanceRange{ + AcceptedAtOrAfterMs: tt.lower, AcceptedBeforeMs: tt.upper, Limit: 51, + Before: storage.RequestAcceptanceCursor{AcceptedAtMs: tt.timestamp, RequestID: tt.requestID}, + }, query) + }) + } +} + +func TestDecodeListPageTokenRejectsMalformedInput(t *testing.T) { + for _, contents := range []string{"not-json", "null", "{}", "[]", `{"version":"1"}`} { + t.Run(contents, func(t *testing.T) { + _, err := decodeListPageToken(base64.RawURLEncoding.EncodeToString([]byte(contents))) + require.ErrorIs(t, err, ErrInvalidRequest) + }) + } + _, err := decodeListPageToken("%%%") + require.ErrorIs(t, err, ErrInvalidRequest) +} + +func TestDecodeListPageTokenRejectsInvalidFields(t *testing.T) { + for _, tt := range []struct { + name string + change func(*listPageToken) + }{ + {"unknown version", func(token *listPageToken) { token.Version++ }}, + {"missing queue", func(token *listPageToken) { token.Queue = "" }}, + {"missing ID", func(token *listPageToken) { token.LastRequestID = "" }}, + {"negative lower", func(token *listPageToken) { token.AcceptedAtOrAfterMs = -1 }}, + {"inverted window", func(token *listPageToken) { token.AcceptedBeforeMs = -1 }}, + {"unknown acceptance time", func(token *listPageToken) { token.LastAcceptedAtMs = 0 }}, + {"cursor below window", func(token *listPageToken) { token.AcceptedAtOrAfterMs = 900 }}, + {"cursor at upper", func(token *listPageToken) { token.LastAcceptedAtMs = 1000 }}, + } { + t.Run(tt.name, func(t *testing.T) { + token := validListToken() + tt.change(&token) + _, err := decodeListPageToken(mustEncodeListToken(t, token)) + require.ErrorIs(t, err, ErrInvalidRequest) + }) + } +} diff --git a/stovepipe/controller/list_test.go b/stovepipe/controller/list_test.go new file mode 100644 index 000000000..11b9e35bc --- /dev/null +++ b/stovepipe/controller/list_test.go @@ -0,0 +1,197 @@ +// 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" + "errors" + "testing" + "time" + + "github.com/stretchr/testify/require" + "github.com/uber-go/tally" + "github.com/uber/submitqueue/platform/errs" + "github.com/uber/submitqueue/stovepipe/entity" + "github.com/uber/submitqueue/stovepipe/extension/storage" + storagemock "github.com/uber/submitqueue/stovepipe/extension/storage/mock" + "go.uber.org/mock/gomock" + "go.uber.org/zap" +) + +type listTestFixture struct { + controller *listController + factory *storagemock.MockFactory + stores *storagemock.MockStorage + acceptances *storagemock.MockRequestAcceptanceStore + summaries *storagemock.MockRequestSummaryStore +} + +func newListTestFixture(t *testing.T) listTestFixture { + t.Helper() + ctrl := gomock.NewController(t) + f := listTestFixture{ + factory: storagemock.NewMockFactory(ctrl), + stores: storagemock.NewMockStorage(ctrl), + acceptances: storagemock.NewMockRequestAcceptanceStore(ctrl), + summaries: storagemock.NewMockRequestSummaryStore(ctrl), + } + f.controller = NewListController(zap.NewNop().Sugar(), tally.NoopScope, f.factory, []string{"monorepo/main"}).(*listController) + f.controller.now = func() time.Time { return time.UnixMilli(1000) } + f.stores.EXPECT().GetRequestAcceptanceStore().Return(f.acceptances).AnyTimes() + f.stores.EXPECT().GetRequestSummaryStore().Return(f.summaries).AnyTimes() + return f +} + +func listTestSummary(timestamp int64, suffix string) entity.RequestSummary { + return entity.RequestSummary{ + Queue: "monorepo/main", RequestID: "request/monorepo/main/" + suffix, + URI: "git://repo/" + suffix, BaseURI: "git://repo/base", AcceptedAtMs: timestamp, + State: entity.RequestStateProcessing, RequestVersion: 2, StateTimestampMs: 2000, Version: 2, + } +} + +func (f listTestFixture) expectPage(query storage.RequestAcceptanceRange, summaries []entity.RequestSummary) { + mappings := make([]entity.RequestAcceptance, 0, len(summaries)) + for _, summary := range summaries { + mappings = append(mappings, entity.RequestAcceptance{Queue: summary.Queue, AcceptedAtMs: summary.AcceptedAtMs, RequestID: summary.RequestID}) + } + f.factory.EXPECT().For(storage.Config{QueueName: "monorepo/main"}).Return(f.stores, nil) + f.acceptances.EXPECT().List(gomock.Any(), query).Return(mappings, nil) + for _, summary := range summaries[:min(len(summaries), query.Limit-1)] { + f.summaries.EXPECT().Get(gomock.Any(), summary.RequestID).Return(summary, nil) + } +} + +func TestListReadsOnePage(t *testing.T) { + first, second, lookahead := listTestSummary(900, "9"), listTestSummary(900, "10"), listTestSummary(800, "1") + for _, tt := range []struct { + name string + rows []entity.RequestSummary + want []entity.RequestSummary + more bool + }{ + {"no matches", nil, []entity.RequestSummary{}, false}, + {"short page", []entity.RequestSummary{first}, []entity.RequestSummary{first}, false}, + {"exact page", []entity.RequestSummary{first, second}, []entity.RequestSummary{first, second}, false}, + {"lookahead", []entity.RequestSummary{first, second, lookahead}, []entity.RequestSummary{first, second}, true}, + } { + t.Run(tt.name, func(t *testing.T) { + f := newListTestFixture(t) + f.expectPage(storage.RequestAcceptanceRange{AcceptedBeforeMs: 1000, Limit: 3}, tt.rows) + got, err := f.controller.List(context.Background(), entity.ListRequest{Queue: "monorepo/main", PageSize: 2}) + require.NoError(t, err) + require.Equal(t, tt.want, got.Requests) + if !tt.more { + require.Empty(t, got.NextPageToken) + return + } + token, err := decodeListPageToken(got.NextPageToken) + require.NoError(t, err) + require.Equal(t, listPageToken{ + Version: listPageTokenVersion, Queue: "monorepo/main", AcceptedBeforeMs: 1000, + LastAcceptedAtMs: second.AcceptedAtMs, LastRequestID: second.RequestID, + }, token) + }) + } +} + +func TestListContinuationPreservesWindowAndReadsCurrentState(t *testing.T) { + f := newListTestFixture(t) + ctx := context.Background() + first, second, third := listTestSummary(900, "9"), listTestSummary(900, "10"), listTestSummary(800, "1") + f.expectPage(storage.RequestAcceptanceRange{AcceptedBeforeMs: 1000, Limit: 2}, []entity.RequestSummary{first, second}) + page, err := f.controller.List(ctx, entity.ListRequest{Queue: "monorepo/main", PageSize: 1}) + require.NoError(t, err) + require.Equal(t, []entity.RequestSummary{first}, page.Requests) + require.NotEmpty(t, page.NextPageToken) + + second.State = entity.RequestStateFailed + second.OutcomeReason = entity.RequestOutcomeReasonBuildFailed + second.StateTimestampMs = 5000 + second.RequestVersion = 3 + second.Version = 3 + f.controller.now = func() time.Time { return time.UnixMilli(5000) } + f.expectPage(storage.RequestAcceptanceRange{ + AcceptedBeforeMs: 1000, Limit: 3, + Before: storage.RequestAcceptanceCursor{AcceptedAtMs: first.AcceptedAtMs, RequestID: first.RequestID}, + }, []entity.RequestSummary{second, third}) + page, err = f.controller.List(ctx, entity.ListRequest{Queue: "monorepo/main", PageToken: page.NextPageToken, PageSize: 2}) + require.NoError(t, err) + require.Equal(t, []entity.RequestSummary{second, third}, page.Requests) + require.Empty(t, page.NextPageToken) +} + +func TestListRejectsInvalidInputBeforeStorageAccess(t *testing.T) { + for _, tt := range []struct { + name string + request entity.ListRequest + }{ + {"empty queue", entity.ListRequest{}}, + {"unconfigured queue", entity.ListRequest{Queue: "other/main"}}, + {"negative bounds", explicitListWindow(-1, 1000)}, + {"invalid page size", entity.ListRequest{Queue: "monorepo/main", PageSize: 201}}, + {"malformed token", entity.ListRequest{Queue: "monorepo/main", PageToken: "%%%"}}, + {"mismatched bounds", entity.ListRequest{Queue: "monorepo/main", PageToken: mustEncodeListToken(t, validListToken()), HasAcceptedBeforeMs: true, AcceptedBeforeMs: 1001}}, + } { + t.Run(tt.name, func(t *testing.T) { + f := newListTestFixture(t) + result, err := f.controller.List(context.Background(), tt.request) + require.ErrorIs(t, err, ErrInvalidRequest) + require.True(t, IsInvalidRequest(err)) + require.True(t, errs.IsUserError(err)) + require.Equal(t, entity.ListResult{}, result) + }) + } +} + +func TestListRechecksConfiguredQueueOnContinuation(t *testing.T) { + f := newListTestFixture(t) + controller := NewListController(zap.NewNop().Sugar(), tally.NoopScope, f.factory, []string{"other/main"}) + _, err := controller.List(context.Background(), entity.ListRequest{ + Queue: "monorepo/main", PageToken: mustEncodeListToken(t, validListToken()), + }) + require.ErrorIs(t, err, ErrInvalidRequest) +} + +func TestListPreservesStorageFailures(t *testing.T) { + failure := errs.NewRetryableError(errors.New("unavailable")) + for _, tt := range []struct { + name string + resolveErr, mappingErr, summaryErr, want error + }{ + {name: "resolve", resolveErr: failure, want: failure}, + {name: "mapping", mappingErr: failure, want: failure}, + {name: "summary", summaryErr: failure, want: failure}, + } { + t.Run(tt.name, func(t *testing.T) { + f := newListTestFixture(t) + f.factory.EXPECT().For(storage.Config{QueueName: "monorepo/main"}).Return(f.stores, tt.resolveErr) + if tt.resolveErr == nil { + f.acceptances.EXPECT().List(gomock.Any(), gomock.Any()).Return([]entity.RequestAcceptance{ + {Queue: "monorepo/main", AcceptedAtMs: 900, RequestID: "request/monorepo/main/1"}, + }, tt.mappingErr) + if tt.mappingErr == nil { + f.summaries.EXPECT().Get(gomock.Any(), "request/monorepo/main/1").Return(entity.RequestSummary{}, tt.summaryErr) + } + } + result, err := f.controller.List(context.Background(), entity.ListRequest{Queue: "monorepo/main"}) + require.Equal(t, entity.ListResult{}, result) + require.ErrorIs(t, err, tt.want) + require.False(t, IsListConsistency(err)) + require.True(t, errs.IsRetryable(err)) + require.False(t, errs.IsUserError(err)) + }) + } +} diff --git a/stovepipe/controller/read_errors.go b/stovepipe/controller/read_errors.go index 5271c4bb3..2555f6a02 100644 --- a/stovepipe/controller/read_errors.go +++ b/stovepipe/controller/read_errors.go @@ -80,3 +80,29 @@ func IsRequestHistoryNotFound(err error) bool { var byURI *RequestHistoryByURINotFoundError return errors.As(err, &byURI) } + +// ListConsistencyError indicates an invalid acceptance mapping, a missing summary, +// or disagreement between the two. It represents an infrastructure failure, not a user lookup miss. +type ListConsistencyError struct { + // Queue is the selected queue. + Queue string + // RequestID identifies the mapped request; empty when the mapping lacks an ID. + RequestID string + // Reason describes the inconsistency in the mapping or summary. + Reason string + // Err is the underlying summary-read failure, if any. + Err error +} + +func (e *ListConsistencyError) Error() string { + return fmt.Sprintf("List %s queue=%q request_id=%q", e.Reason, e.Queue, e.RequestID) +} + +// Unwrap preserves the underlying summary-read failure, when present. +func (e *ListConsistencyError) Unwrap() error { return e.Err } + +// IsListConsistency reports whether err represents inconsistent listing records. +func IsListConsistency(err error) bool { + var target *ListConsistencyError + return errors.As(err, &target) +} diff --git a/stovepipe/entity/BUILD.bazel b/stovepipe/entity/BUILD.bazel index c216fdfa1..b38416e00 100644 --- a/stovepipe/entity/BUILD.bazel +++ b/stovepipe/entity/BUILD.bazel @@ -5,6 +5,7 @@ go_library( srcs = [ "build.go", "ingest.go", + "list.go", "project_status.go", "queue.go", "queue_config.go", diff --git a/stovepipe/entity/list.go b/stovepipe/entity/list.go new file mode 100644 index 000000000..44f090665 --- /dev/null +++ b/stovepipe/entity/list.go @@ -0,0 +1,41 @@ +// 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 entity + +// ListRequest selects current request summaries from one queue by acceptance time. +type ListRequest struct { + // Queue is the required configured queue containing the requests. + Queue string + // AcceptedAtOrAfterMs is the inclusive Unix millisecond bound when HasAcceptedAtOrAfterMs is true. + AcceptedAtOrAfterMs int64 + // HasAcceptedAtOrAfterMs distinguishes an explicit lower bound from an omitted one. + HasAcceptedAtOrAfterMs bool + // AcceptedBeforeMs is the exclusive Unix millisecond bound when HasAcceptedBeforeMs is true. + AcceptedBeforeMs int64 + // HasAcceptedBeforeMs distinguishes an explicit upper bound from an omitted one. + HasAcceptedBeforeMs bool + // PageSize is the maximum number of summaries, from 1 to 200; zero selects 50. + PageSize int32 + // PageToken is an opaque continuation bound to the queue and resolved time window. + PageToken string +} + +// ListResult is one page of current summaries, not a snapshot across pages. +type ListResult struct { + // Requests contains summaries ordered by descending acceptance time, then descending bytewise request ID. + Requests []RequestSummary + // NextPageToken continues after the last returned key; empty means no further mapping was observed. + NextPageToken string +}