Skip to content

Commit df0321c

Browse files
committed
feat(stovepipe): implement List controller and pagination
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.
1 parent de36e11 commit df0321c

10 files changed

Lines changed: 839 additions & 3 deletions

File tree

‎doc/rfc/stovepipe/list-api.md‎

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
11
# Stovepipe List API
22

3-
The [protobuf contract](../../../api/stovepipe/proto/stovepipe.proto) is included for review; controller and storage implementation are deferred.
3+
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.
44

55
## Proposal
66

@@ -45,9 +45,9 @@ This is the wire representation of the existing domain `RequestSummary`, not ano
4545

4646
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.
4747

48-
## Data and Storage Work Required
48+
## Data and Storage
4949

50-
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.
50+
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.
5151

5252
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.
5353

‎stovepipe/controller/BUILD.bazel‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,8 @@ go_library(
55
srcs = [
66
"get_project_status_by_uri.go",
77
"ingest.go",
8+
"list.go",
9+
"list_pagination.go",
810
"ping.go",
911
"read_errors.go",
1012
"request_history.go",
@@ -34,6 +36,9 @@ go_test(
3436
srcs = [
3537
"get_project_status_by_uri_test.go",
3638
"ingest_test.go",
39+
"list_consistency_test.go",
40+
"list_pagination_test.go",
41+
"list_test.go",
3742
"ping_test.go",
3843
"request_history_test.go",
3944
],

‎stovepipe/controller/list.go‎

Lines changed: 133 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,133 @@
1+
// Copyright (c) 2026 Uber Technologies, Inc.
2+
//
3+
// Licensed under the Apache License, Version 2.0 (the "License");
4+
// you may not use this file except in compliance with the License.
5+
// You may obtain a copy of the License at
6+
//
7+
// http://www.apache.org/licenses/LICENSE-2.0
8+
//
9+
// Unless required by applicable law or agreed to in writing, software
10+
// distributed under the License is distributed on an "AS IS" BASIS,
11+
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12+
// See the License for the specific language governing permissions and
13+
// limitations under the License.
14+
15+
package controller
16+
17+
import (
18+
"context"
19+
"fmt"
20+
"time"
21+
22+
"github.com/uber-go/tally"
23+
"github.com/uber/submitqueue/platform/metrics"
24+
"github.com/uber/submitqueue/stovepipe/entity"
25+
"github.com/uber/submitqueue/stovepipe/extension/storage"
26+
"go.uber.org/zap"
27+
)
28+
29+
const (
30+
defaultListPageSize = 50
31+
maxListPageSize = 200
32+
)
33+
34+
// ListController handles queue-scoped acceptance-time listing.
35+
type ListController interface {
36+
List(ctx context.Context, req entity.ListRequest) (entity.ListResult, error)
37+
}
38+
39+
type listController struct {
40+
logger *zap.SugaredLogger
41+
metricsScope tally.Scope
42+
stores storage.Factory
43+
configuredQueues map[string]struct{}
44+
now func() time.Time
45+
}
46+
47+
// NewListController creates a controller restricted to the supplied configured queues.
48+
func NewListController(logger *zap.SugaredLogger, scope tally.Scope, stores storage.Factory, configuredQueues []string) ListController {
49+
queues := make(map[string]struct{}, len(configuredQueues))
50+
for _, queue := range configuredQueues {
51+
queues[queue] = struct{}{}
52+
}
53+
return &listController{
54+
logger: logger, metricsScope: scope.SubScope("list_controller"), stores: stores,
55+
configuredQueues: queues, now: time.Now,
56+
}
57+
}
58+
59+
// List returns current summaries in descending acceptance-time order within a half-open time window.
60+
// Omitted bounds resolve to [0, server now) on the first page and to the token's window on continuations.
61+
func (c *listController) List(ctx context.Context, req entity.ListRequest) (result entity.ListResult, retErr error) {
62+
op := metrics.Begin(c.metricsScope, "list", metrics.StorageLatencyBuckets, metrics.TagsFromContext(ctx)...)
63+
defer func() { op.Complete(retErr) }()
64+
65+
if req.Queue == "" {
66+
return entity.ListResult{}, fmt.Errorf("List requires a queue: %w", ErrInvalidRequest)
67+
}
68+
if _, ok := c.configuredQueues[req.Queue]; !ok {
69+
return entity.ListResult{}, fmt.Errorf("List queue %q is not configured: %w", req.Queue, ErrInvalidRequest)
70+
}
71+
query, err := c.resolveListRange(req)
72+
if err != nil {
73+
return entity.ListResult{}, err
74+
}
75+
stores, err := c.stores.For(storage.Config{QueueName: req.Queue})
76+
if err != nil {
77+
return entity.ListResult{}, fmt.Errorf("List failed to resolve storage for queue %q: %w", req.Queue, err)
78+
}
79+
mappings, err := stores.GetRequestAcceptanceStore().List(ctx, query)
80+
if err != nil {
81+
return entity.ListResult{}, fmt.Errorf("List failed to read acceptance mappings for queue %q: %w", req.Queue, err)
82+
}
83+
pageSize := query.Limit - 1
84+
visible := mappings[:min(len(mappings), pageSize)]
85+
result.Requests = make([]entity.RequestSummary, 0, len(visible))
86+
if len(visible) > 0 {
87+
summaries := stores.GetRequestSummaryStore()
88+
for _, mapping := range visible {
89+
summary, err := readListSummary(ctx, summaries, req.Queue, mapping)
90+
if err != nil {
91+
return entity.ListResult{}, err
92+
}
93+
result.Requests = append(result.Requests, summary)
94+
}
95+
}
96+
if len(mappings) > pageSize {
97+
last := visible[len(visible)-1]
98+
nextToken, err := encodeListPageToken(listPageToken{
99+
Version: listPageTokenVersion, Queue: req.Queue,
100+
AcceptedAtOrAfterMs: query.AcceptedAtOrAfterMs, AcceptedBeforeMs: query.AcceptedBeforeMs,
101+
LastAcceptedAtMs: last.AcceptedAtMs, LastRequestID: last.RequestID,
102+
})
103+
if err != nil {
104+
return entity.ListResult{}, fmt.Errorf("List failed to encode continuation: %w", err)
105+
}
106+
result.NextPageToken = nextToken
107+
}
108+
c.logger.Debugw("queue requests listed", "queue", req.Queue, "request_count", len(result.Requests), "has_next_page", result.NextPageToken != "")
109+
return result, nil
110+
}
111+
112+
func readListSummary(ctx context.Context, summaries storage.RequestSummaryStore, queue string, mapping entity.RequestAcceptance) (entity.RequestSummary, error) {
113+
if mapping.Queue != queue || mapping.AcceptedAtMs <= 0 || mapping.RequestID == "" {
114+
return entity.RequestSummary{}, &ListConsistencyError{
115+
Queue: queue, RequestID: mapping.RequestID, Reason: "invalid acceptance mapping",
116+
}
117+
}
118+
summary, err := summaries.Get(ctx, mapping.RequestID)
119+
if err != nil {
120+
if storage.IsNotFound(err) {
121+
return entity.RequestSummary{}, &ListConsistencyError{
122+
Queue: queue, RequestID: mapping.RequestID, Reason: "acceptance mapping has no summary", Err: err,
123+
}
124+
}
125+
return entity.RequestSummary{}, fmt.Errorf("List failed to read summary queue=%q request_id=%q: %w", queue, mapping.RequestID, err)
126+
}
127+
if summary.Queue != queue || summary.RequestID != mapping.RequestID || summary.AcceptedAtMs != mapping.AcceptedAtMs {
128+
return entity.RequestSummary{}, &ListConsistencyError{
129+
Queue: queue, RequestID: mapping.RequestID, Reason: "acceptance mapping disagrees with summary",
130+
}
131+
}
132+
return summary, nil
133+
}
Lines changed: 117 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,117 @@
1+
// Copyright (c) 2026 Uber Technologies, Inc.
2+
//
3+
// Licensed under the Apache License, Version 2.0 (the "License");
4+
// you may not use this file except in compliance with the License.
5+
// You may obtain a copy of the License at
6+
//
7+
// http://www.apache.org/licenses/LICENSE-2.0
8+
//
9+
// Unless required by applicable law or agreed to in writing, software
10+
// distributed under the License is distributed on an "AS IS" BASIS,
11+
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12+
// See the License for the specific language governing permissions and
13+
// limitations under the License.
14+
15+
package controller
16+
17+
import (
18+
"context"
19+
"fmt"
20+
"testing"
21+
22+
"github.com/stretchr/testify/require"
23+
"github.com/uber/submitqueue/platform/errs"
24+
"github.com/uber/submitqueue/stovepipe/entity"
25+
"github.com/uber/submitqueue/stovepipe/extension/storage"
26+
"go.uber.org/mock/gomock"
27+
)
28+
29+
func TestListRejectsInconsistentRecordsWithoutPartialResults(t *testing.T) {
30+
first := listTestSummary(900, "9")
31+
second := listTestSummary(800, "1")
32+
mapping := entity.RequestAcceptance{Queue: second.Queue, AcceptedAtMs: second.AcceptedAtMs, RequestID: second.RequestID}
33+
for _, tt := range []struct {
34+
name string
35+
change func(*entity.RequestSummary)
36+
readErr error
37+
wantReason string
38+
}{
39+
{
40+
name: "wrong queue", change: func(summary *entity.RequestSummary) { summary.Queue = "other/main" },
41+
wantReason: "acceptance mapping disagrees with summary",
42+
},
43+
{
44+
name: "wrong request", change: func(summary *entity.RequestSummary) { summary.RequestID = first.RequestID },
45+
wantReason: "acceptance mapping disagrees with summary",
46+
},
47+
{
48+
name: "unknown acceptance", change: func(summary *entity.RequestSummary) { summary.AcceptedAtMs = 0 },
49+
wantReason: "acceptance mapping disagrees with summary",
50+
},
51+
{
52+
name: "different acceptance", change: func(summary *entity.RequestSummary) { summary.AcceptedAtMs++ },
53+
wantReason: "acceptance mapping disagrees with summary",
54+
},
55+
{
56+
name: "missing summary", readErr: fmt.Errorf("summary read failed: %w", storage.ErrNotFound),
57+
wantReason: "acceptance mapping has no summary",
58+
},
59+
} {
60+
t.Run(tt.name, func(t *testing.T) {
61+
f := newListTestFixture(t)
62+
corrupt := second
63+
if tt.change != nil {
64+
tt.change(&corrupt)
65+
}
66+
f.factory.EXPECT().For(gomock.Any()).Return(f.stores, nil)
67+
f.acceptances.EXPECT().List(gomock.Any(), gomock.Any()).Return([]entity.RequestAcceptance{
68+
{Queue: first.Queue, AcceptedAtMs: first.AcceptedAtMs, RequestID: first.RequestID}, mapping,
69+
}, nil)
70+
gomock.InOrder(
71+
f.summaries.EXPECT().Get(gomock.Any(), first.RequestID).Return(first, nil),
72+
f.summaries.EXPECT().Get(gomock.Any(), second.RequestID).Return(corrupt, tt.readErr),
73+
)
74+
result, err := f.controller.List(context.Background(), entity.ListRequest{Queue: first.Queue})
75+
require.True(t, IsListConsistency(err))
76+
require.True(t, IsListConsistency(fmt.Errorf("List failed: %w", err)))
77+
var consistency *ListConsistencyError
78+
require.ErrorAs(t, err, &consistency)
79+
require.Equal(t, first.Queue, consistency.Queue)
80+
require.Equal(t, second.RequestID, consistency.RequestID)
81+
require.Equal(t, tt.wantReason, consistency.Reason)
82+
require.Equal(t, tt.readErr, consistency.Err)
83+
if tt.readErr != nil {
84+
require.ErrorIs(t, err, tt.readErr)
85+
require.ErrorIs(t, err, storage.ErrNotFound)
86+
}
87+
require.False(t, errs.IsRetryable(err))
88+
require.False(t, errs.IsUserError(err))
89+
require.Equal(t, entity.ListResult{}, result)
90+
})
91+
}
92+
}
93+
94+
func TestListRejectsInvalidMappingBeforeSummaryLookup(t *testing.T) {
95+
for _, mapping := range []entity.RequestAcceptance{
96+
{Queue: "other/main", AcceptedAtMs: 900, RequestID: "request/other/main/1"},
97+
{Queue: "monorepo/main", AcceptedAtMs: 0, RequestID: "request/monorepo/main/1"},
98+
{Queue: "monorepo/main", AcceptedAtMs: 900},
99+
} {
100+
t.Run(mapping.RequestID, func(t *testing.T) {
101+
f := newListTestFixture(t)
102+
f.factory.EXPECT().For(gomock.Any()).Return(f.stores, nil)
103+
f.acceptances.EXPECT().List(gomock.Any(), gomock.Any()).Return([]entity.RequestAcceptance{mapping}, nil)
104+
result, err := f.controller.List(context.Background(), entity.ListRequest{Queue: "monorepo/main"})
105+
require.True(t, IsListConsistency(err))
106+
var consistency *ListConsistencyError
107+
require.ErrorAs(t, err, &consistency)
108+
require.Equal(t, "monorepo/main", consistency.Queue)
109+
require.Equal(t, mapping.RequestID, consistency.RequestID)
110+
require.Equal(t, "invalid acceptance mapping", consistency.Reason)
111+
require.Nil(t, consistency.Err)
112+
require.False(t, errs.IsRetryable(err))
113+
require.False(t, errs.IsUserError(err))
114+
require.Equal(t, entity.ListResult{}, result)
115+
})
116+
}
117+
}
Lines changed: 98 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,98 @@
1+
// Copyright (c) 2026 Uber Technologies, Inc.
2+
//
3+
// Licensed under the Apache License, Version 2.0 (the "License");
4+
// you may not use this file except in compliance with the License.
5+
// You may obtain a copy of the License at
6+
//
7+
// http://www.apache.org/licenses/LICENSE-2.0
8+
//
9+
// Unless required by applicable law or agreed to in writing, software
10+
// distributed under the License is distributed on an "AS IS" BASIS,
11+
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12+
// See the License for the specific language governing permissions and
13+
// limitations under the License.
14+
15+
package controller
16+
17+
import (
18+
"encoding/base64"
19+
"encoding/json"
20+
"fmt"
21+
22+
"github.com/uber/submitqueue/stovepipe/entity"
23+
"github.com/uber/submitqueue/stovepipe/extension/storage"
24+
)
25+
26+
const listPageTokenVersion = 1
27+
28+
type listPageToken struct {
29+
Version int `json:"version"`
30+
Queue string `json:"queue"`
31+
AcceptedAtOrAfterMs int64 `json:"accepted_at_or_after_ms"`
32+
AcceptedBeforeMs int64 `json:"accepted_before_ms"`
33+
LastAcceptedAtMs int64 `json:"last_accepted_at_ms"`
34+
LastRequestID string `json:"last_request_id"`
35+
}
36+
37+
func (c *listController) resolveListRange(req entity.ListRequest) (storage.RequestAcceptanceRange, error) {
38+
if req.PageSize < 0 || req.PageSize > maxListPageSize {
39+
return storage.RequestAcceptanceRange{}, fmt.Errorf("List page_size must be between 0 and %d: %w", maxListPageSize, ErrInvalidRequest)
40+
}
41+
pageSize := int(req.PageSize)
42+
if pageSize == 0 {
43+
pageSize = defaultListPageSize
44+
}
45+
query := storage.RequestAcceptanceRange{Limit: pageSize + 1}
46+
if req.PageToken != "" {
47+
token, err := decodeListPageToken(req.PageToken)
48+
if err != nil {
49+
return storage.RequestAcceptanceRange{}, err
50+
}
51+
if token.Queue != req.Queue ||
52+
(req.HasAcceptedAtOrAfterMs && req.AcceptedAtOrAfterMs != token.AcceptedAtOrAfterMs) ||
53+
(req.HasAcceptedBeforeMs && req.AcceptedBeforeMs != token.AcceptedBeforeMs) {
54+
return storage.RequestAcceptanceRange{}, fmt.Errorf("List page_token does not match the queue and time bounds: %w", ErrInvalidRequest)
55+
}
56+
query.AcceptedAtOrAfterMs = token.AcceptedAtOrAfterMs
57+
query.AcceptedBeforeMs = token.AcceptedBeforeMs
58+
query.Before = storage.RequestAcceptanceCursor{AcceptedAtMs: token.LastAcceptedAtMs, RequestID: token.LastRequestID}
59+
return query, nil
60+
}
61+
if req.HasAcceptedAtOrAfterMs {
62+
query.AcceptedAtOrAfterMs = req.AcceptedAtOrAfterMs
63+
}
64+
if req.HasAcceptedBeforeMs {
65+
query.AcceptedBeforeMs = req.AcceptedBeforeMs
66+
} else {
67+
query.AcceptedBeforeMs = c.now().UnixMilli()
68+
}
69+
if query.AcceptedAtOrAfterMs < 0 || query.AcceptedBeforeMs <= query.AcceptedAtOrAfterMs {
70+
return storage.RequestAcceptanceRange{}, fmt.Errorf("List requires 0 <= accepted_at_or_after_ms < accepted_before_ms: %w", ErrInvalidRequest)
71+
}
72+
return query, nil
73+
}
74+
75+
func encodeListPageToken(token listPageToken) (string, error) {
76+
contents, err := json.Marshal(token)
77+
if err != nil {
78+
return "", err
79+
}
80+
return base64.RawURLEncoding.EncodeToString(contents), nil
81+
}
82+
83+
func decodeListPageToken(encoded string) (listPageToken, error) {
84+
contents, err := base64.RawURLEncoding.Strict().DecodeString(encoded)
85+
if err != nil {
86+
return listPageToken{}, fmt.Errorf("List invalid page_token encoding: %w", ErrInvalidRequest)
87+
}
88+
var token listPageToken
89+
if err := json.Unmarshal(contents, &token); err != nil {
90+
return listPageToken{}, fmt.Errorf("List invalid page_token contents: %w", ErrInvalidRequest)
91+
}
92+
if token.Version != listPageTokenVersion || token.Queue == "" || token.LastRequestID == "" ||
93+
token.AcceptedAtOrAfterMs < 0 || token.AcceptedBeforeMs <= token.AcceptedAtOrAfterMs ||
94+
token.LastAcceptedAtMs <= 0 || token.LastAcceptedAtMs < token.AcceptedAtOrAfterMs || token.LastAcceptedAtMs >= token.AcceptedBeforeMs {
95+
return listPageToken{}, fmt.Errorf("List invalid page_token fields: %w", ErrInvalidRequest)
96+
}
97+
return token, nil
98+
}

0 commit comments

Comments
 (0)