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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 3 additions & 3 deletions doc/rfc/stovepipe/list-api.md
Original file line number Diff line number Diff line change
@@ -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

Expand Down Expand Up @@ -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.

Expand Down
5 changes: 5 additions & 0 deletions stovepipe/controller/BUILD.bazel
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down Expand Up @@ -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",
],
Expand Down
133 changes: 133 additions & 0 deletions stovepipe/controller/list.go
Original file line number Diff line number Diff line change
@@ -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)
Comment thread
mnoah1 marked this conversation as resolved.
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
}
117 changes: 117 additions & 0 deletions stovepipe/controller/list_consistency_test.go
Original file line number Diff line number Diff line change
@@ -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)
})
}
}
98 changes: 98 additions & 0 deletions stovepipe/controller/list_pagination.go
Original file line number Diff line number Diff line change
@@ -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
}
Loading
Loading