diff --git a/doc/rfc/stovepipe/list-api.md b/doc/rfc/stovepipe/list-api.md index ef83d1a0..cbaf27f1 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 38ff0fb0..6f25456a 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 00000000..f799c31d --- /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 00000000..5fd40174 --- /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 00000000..2bbcea6d --- /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 00000000..bda6630f --- /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 00000000..11b9e35b --- /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 5271c4bb..2555f6a0 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 c216fdfa..b38416e0 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 00000000..44f09066 --- /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 +}