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
1 change: 1 addition & 0 deletions submitqueue/entity/BUILD.bazel
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@ go_library(
"request_batch.go",
"request_history.go",
"request_log.go",
"request_receipt.go",
"request_summary.go",
"speculation.go",
"subject.go",
Expand Down
25 changes: 25 additions & 0 deletions submitqueue/entity/request_receipt.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,25 @@
// 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

// RequestReceipt is an immutable lookup key for a publicly visible request.
type RequestReceipt struct {
// Queue is the queue supplied at receipt.
Queue string
// ReceivedAtMs is the immutable, positive receipt timestamp in Unix milliseconds.
ReceivedAtMs int64
// RequestID identifies the request within Queue.
RequestID string
}
1 change: 1 addition & 0 deletions submitqueue/extension/storage/BUILD.bazel
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@ go_library(
"request_batch_store.go",
"request_log_store.go",
"request_queue_summary_store.go",
"request_receipt_store.go",
"request_store.go",
"request_summary_store.go",
"request_uri_store.go",
Expand Down
1 change: 1 addition & 0 deletions submitqueue/extension/storage/mock/BUILD.bazel
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@ go_library(
"request_batch_store_mock.go",
"request_log_store_mock.go",
"request_queue_summary_store_mock.go",
"request_receipt_store_mock.go",
"request_store_mock.go",
"request_summary_store_mock.go",
"request_uri_store_mock.go",
Expand Down

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

54 changes: 54 additions & 0 deletions submitqueue/extension/storage/request_receipt_store.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,54 @@
// 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 storage

//go:generate mockgen -source=request_receipt_store.go -destination=mock/request_receipt_store_mock.go -package=mock

import (
"context"

"github.com/uber/submitqueue/submitqueue/entity"
)

// RequestReceiptCursor is an exclusive position in descending receipt-time order.
type RequestReceiptCursor struct {
// ReceivedAtMs is the positive receipt timestamp in Unix milliseconds; zero with an empty RequestID starts the range.
ReceivedAtMs int64
// RequestID breaks timestamp ties in descending string order, not numeric order.
RequestID string
}

// RequestReceiptRange bounds a primary-key scan within the store's queue.
type RequestReceiptRange struct {
// ReceivedAtOrAfterMs is the inclusive lower receipt-time bound in Unix milliseconds.
ReceivedAtOrAfterMs int64
// ReceivedBeforeMs is the exclusive upper receipt-time bound and must exceed the lower bound.
ReceivedBeforeMs int64
// Before is the exclusive continuation boundary; its zero value starts the range.
Before RequestReceiptCursor
// Limit is the positive maximum number of mappings to return.
Limit int
}

// RequestReceiptStore retains immutable receipt mappings in its bound queue.
type RequestReceiptStore interface {
// Create rejects queue mismatches, nonpositive receipt times, and empty request IDs.
// An existing composite key returns ErrAlreadyExists.
Create(ctx context.Context, receipt entity.RequestReceipt) error

// List scans the primary-key range in descending (received_at_ms, request_id) order.
// It returns at most Limit mappings, an empty slice when none match, and an error for invalid ranges.
List(ctx context.Context, bounds RequestReceiptRange) ([]entity.RequestReceipt, error)
}
10 changes: 10 additions & 0 deletions submitqueue/gateway/controller/land_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -348,6 +348,7 @@ func TestLand_PublishesToQueue(t *testing.T) {
var materializedSummary entity.RequestSummary
var persistedMapping entity.RequestURI
var persistedQueueSummary entity.RequestQueueSummary
var persistedReceipt entity.RequestReceipt
var persistedLog entity.RequestLog

ctrl := gomock.NewController(t)
Expand All @@ -359,8 +360,10 @@ func TestLand_PublishesToQueue(t *testing.T) {
summaryStore := storagemock.NewMockRequestSummaryStore(ctrl)
uriStore := storagemock.NewMockRequestURIStore(ctrl)
queueStore := storagemock.NewMockRequestQueueSummaryStore(ctrl)
receiptStore := storagemock.NewMockRequestReceiptStore(ctrl)
logStore := storagemock.NewMockRequestLogStore(ctrl)
store.EXPECT().GetRequestQueueSummaryStore().Return(queueStore).AnyTimes()
store.EXPECT().GetRequestReceiptStore().Return(receiptStore).AnyTimes()
store.EXPECT().GetRequestSummaryStore().Return(summaryStore).AnyTimes()
store.EXPECT().GetRequestLogStore().Return(logStore).AnyTimes()
store.EXPECT().GetRequestURIStore().Return(uriStore).AnyTimes()
Expand Down Expand Up @@ -418,6 +421,12 @@ func TestLand_PublishesToQueue(t *testing.T) {
return nil
},
),
receiptStore.EXPECT().Create(gomock.Any(), gomock.Any()).DoAndReturn(
func(_ context.Context, receipt entity.RequestReceipt) error {
persistedReceipt = receipt
return nil
},
),
)

controller := NewLandController(zap.NewNop().Sugar(), tally.NoopScope, staticCounterFactory{counter: cnt}, factoryForStorage(ctrl, store), materializer, noopQueueConfigStore(ctrl), registry)
Expand All @@ -444,6 +453,7 @@ func TestLand_PublishesToQueue(t *testing.T) {
Metadata: map[string]string{},
}, receiptSummary)
assert.Positive(t, receiptSummary.ReceivedAtMs)
assert.Equal(t, entity.RequestReceipt{Queue: req.Queue, ReceivedAtMs: receiptSummary.ReceivedAtMs, RequestID: result.ID}, persistedReceipt)
assert.Equal(t, entity.RequestLog{
RequestID: "123",
Queue: "test-queue",
Expand Down
5 changes: 5 additions & 0 deletions submitqueue/gateway/controller/log/log_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -141,8 +141,10 @@ func newLogControllerStore(ctrl *gomock.Controller, insertErr, getErr, updateErr
logStore := storagemock.NewMockRequestLogStore(ctrl)
summaryStore := storagemock.NewMockRequestSummaryStore(ctrl)
queueStore := storagemock.NewMockRequestQueueSummaryStore(ctrl)
receiptStore := storagemock.NewMockRequestReceiptStore(ctrl)
uriStore := storagemock.NewMockRequestURIStore(ctrl)
store.EXPECT().GetRequestQueueSummaryStore().Return(queueStore).AnyTimes()
store.EXPECT().GetRequestReceiptStore().Return(receiptStore).AnyTimes()
store.EXPECT().GetRequestSummaryStore().Return(summaryStore).AnyTimes()
store.EXPECT().GetRequestLogStore().Return(logStore).AnyTimes()
store.EXPECT().GetRequestURIStore().Return(uriStore).AnyTimes()
Expand All @@ -169,6 +171,9 @@ func newLogControllerStore(ctrl *gomock.Controller, insertErr, getErr, updateErr
Status: entity.RequestStatusAccepted, Version: 1, Metadata: map[string]string{},
}, nil)
queueStore.EXPECT().Update(gomock.Any(), gomock.Any(), int32(1), int32(2)).Return(queueErr)
if queueErr == nil {
receiptStore.EXPECT().Create(gomock.Any(), entity.RequestReceipt{Queue: "test-queue", ReceivedAtMs: 1, RequestID: "2"}).Return(nil)
}
return materializer
}

Expand Down
4 changes: 4 additions & 0 deletions submitqueue/gateway/controller/storage_fixture_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,7 @@ type controllerStorageFixture struct {
storage *gwstoragemock.MockStorage
summaryStore *storagemock.MockRequestSummaryStore
queueStore *storagemock.MockRequestQueueSummaryStore
receiptStore *storagemock.MockRequestReceiptStore
uriStore *storagemock.MockRequestURIStore
logStore *storagemock.MockRequestLogStore
mu sync.Mutex
Expand All @@ -47,12 +48,14 @@ func newControllerStorageFixture(ctrl *gomock.Controller) *controllerStorageFixt
storage: gwstoragemock.NewMockStorage(ctrl),
summaryStore: storagemock.NewMockRequestSummaryStore(ctrl),
queueStore: storagemock.NewMockRequestQueueSummaryStore(ctrl),
receiptStore: storagemock.NewMockRequestReceiptStore(ctrl),
uriStore: storagemock.NewMockRequestURIStore(ctrl),
logStore: storagemock.NewMockRequestLogStore(ctrl),
summaries: make(map[string]entity.RequestSummary),
queueSummaries: make(map[string]entity.RequestQueueSummary),
}
fixture.storage.EXPECT().GetRequestQueueSummaryStore().Return(fixture.queueStore).AnyTimes()
fixture.storage.EXPECT().GetRequestReceiptStore().Return(fixture.receiptStore).AnyTimes()
fixture.storage.EXPECT().GetRequestSummaryStore().Return(fixture.summaryStore).AnyTimes()
fixture.storage.EXPECT().GetRequestLogStore().Return(fixture.logStore).AnyTimes()
fixture.storage.EXPECT().GetRequestURIStore().Return(fixture.uriStore).AnyTimes()
Expand Down Expand Up @@ -127,6 +130,7 @@ func newControllerStorageFixture(ctrl *gomock.Controller) *controllerStorageFixt
}).AnyTimes()

fixture.uriStore.EXPECT().Create(gomock.Any(), gomock.Any()).Return(nil).AnyTimes()
fixture.receiptStore.EXPECT().Create(gomock.Any(), gomock.Any()).Return(nil).AnyTimes()
fixture.logStore.EXPECT().Insert(gomock.Any(), gomock.Any()).DoAndReturn(func(_ context.Context, log entity.RequestLog) error {
fixture.mu.Lock()
defer fixture.mu.Unlock()
Expand Down
2 changes: 2 additions & 0 deletions submitqueue/gateway/core/request/BUILD.bazel
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ go_library(
srcs = [
"materializer.go",
"request.go",
"request_receipt.go",
],
importpath = "github.com/uber/submitqueue/submitqueue/gateway/core/request",
visibility = ["//visibility:public"],
Expand All @@ -19,6 +20,7 @@ go_test(
name = "go_default_test",
srcs = [
"materializer_test.go",
"request_receipt_test.go",
"request_test.go",
],
embed = [":go_default_library"],
Expand Down
6 changes: 6 additions & 0 deletions submitqueue/gateway/core/request/materializer.go
Original file line number Diff line number Diff line change
Expand Up @@ -80,9 +80,15 @@ func (m *Materializer) PersistLog(ctx context.Context, log entity.RequestLog) er
summary = updated
}

if summary.Status == entity.RequestStatusAccepting {
return nil
}
if err := m.repairPublicProjections(ctx, stores, summary); err != nil {
return err
}
if err := ensureRequestReceiptMapping(ctx, stores.GetRequestReceiptStore(), summary); err != nil {
return err
}
return nil
}
}
Expand Down
16 changes: 3 additions & 13 deletions submitqueue/gateway/core/request/materializer_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,6 @@ import (
"github.com/uber/submitqueue/submitqueue/entity"
"github.com/uber/submitqueue/submitqueue/extension/storage"
storagemock "github.com/uber/submitqueue/submitqueue/extension/storage/mock"
gwstoragemock "github.com/uber/submitqueue/submitqueue/gateway/extension/storage/mock"
"go.uber.org/mock/gomock"
)

Expand Down Expand Up @@ -313,18 +312,9 @@ func TestLogWins(t *testing.T) {
}

func materializerStores(ctrl *gomock.Controller) (*Materializer, *storagemock.MockRequestSummaryStore, *storagemock.MockRequestQueueSummaryStore, *storagemock.MockRequestURIStore, *storagemock.MockRequestLogStore) {
summaryStore := storagemock.NewMockRequestSummaryStore(ctrl)
queueStore := storagemock.NewMockRequestQueueSummaryStore(ctrl)
uriStore := storagemock.NewMockRequestURIStore(ctrl)
logStore := storagemock.NewMockRequestLogStore(ctrl)
queueScoped := gwstoragemock.NewMockStorage(ctrl)
queueScoped.EXPECT().GetRequestQueueSummaryStore().Return(queueStore).AnyTimes()
queueScoped.EXPECT().GetRequestSummaryStore().Return(summaryStore).AnyTimes()
queueScoped.EXPECT().GetRequestURIStore().Return(uriStore).AnyTimes()
queueScoped.EXPECT().GetRequestLogStore().Return(logStore).AnyTimes()
factory := gwstoragemock.NewMockFactory(ctrl)
factory.EXPECT().For(gomock.Any()).Return(queueScoped, nil).AnyTimes()
return NewMaterializer(factory), summaryStore, queueStore, uriStore, logStore
fixture := newMaterializerReceiptFixture(ctrl)
fixture.receipts.EXPECT().Create(gomock.Any(), gomock.Any()).Return(nil).AnyTimes()
return fixture.materializer, fixture.summaries, fixture.queueSummaries, fixture.uris, fixture.logs
}

func testRequestSummary() entity.RequestSummary {
Expand Down
35 changes: 35 additions & 0 deletions submitqueue/gateway/core/request/request_receipt.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,35 @@
// 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 request

import (
"context"
"errors"
"fmt"

"github.com/uber/submitqueue/submitqueue/entity"
basestorage "github.com/uber/submitqueue/submitqueue/extension/storage"
)

func ensureRequestReceiptMapping(ctx context.Context, receipts basestorage.RequestReceiptStore, summary entity.RequestSummary) error {
receipt := entity.RequestReceipt{
Queue: summary.Queue, ReceivedAtMs: summary.ReceivedAtMs, RequestID: summary.RequestID,
}
// Unchanged summaries still need this write after a partially completed attempt.
if err := receipts.Create(ctx, receipt); err != nil && !errors.Is(err, basestorage.ErrAlreadyExists) {
return fmt.Errorf("failed to ensure request receipt mapping queue=%q request_id=%q: %w", summary.Queue, summary.RequestID, err)
}
return nil
}
Loading
Loading