From 74e3cc30b4d5af365efeb5e8fed9643049726621 Mon Sep 17 00:00:00 2001 From: mnoah1 Date: Wed, 7 Oct 2026 15:45:53 +0000 Subject: [PATCH] feat(submitqueue): add immutable receipt lookup Summary: First step toward keeping one full request projection in request_summary, with a small time-based lookup instead of duplicating lifecycle data for List. All required response data and receipt time are already captured. ### Changes - Add request_receipt with only (queue, received_at_ms, request_id), its immutable storage contract, and gateway wiring. - Ensure receipt keys after public projection writes, including unchanged/stale-log retries. Keep accepting receipts hidden and surface partial-write failures for existing retry handling. - Keep List reads and the existing queue-summary projection writes in place for a safe cutover. - Add storage, materializer, and integration coverage; document the transition in the RFC. ### Migration plan 1. Deploy the additive table before rolling out these writes. This PR populates receipt keys alongside the existing projection; it does not switch List reads. 2. Upgrade all writers, then backfill existing List membership from request_summary_by_queue keys. Validate keys against public authoritative summaries, exclude accepting receipts, and verify coverage before cutover. 3. In a follow-up, switch List to scan receipt keys and point-read request_summary. Preserve the API, time bounds, ordering, and page tokens; check page latency and keep old projection writes during the rollback window. 4. Once cutover is stable, stop duplicate projection writes, remove RequestQueueSummary and its store/mapper, and retire request_summary_by_queue after the rollback window. The final model stores one full projection plus immutable lookup keys. Backfill tooling, read cutover, and duplicate-data removal are follow-up work. This preserves currently listed history without reconstructing requests predating the original read-model rollout. Test Plan: Before the internal rollout, provision the new table and confirm newly public requests produce matching receipt keys while List continues serving the old projection. Existing automated tests cover retries, concurrency, hidden receipts, and pagination. Revert Plan: Revert the writer rollout; List still uses the existing projection. Leave the additive table in place until no writers depend on it. API Changes: No RPC/proto changes. Gateway storage implementations must provide the new RequestReceiptStore accessor. --- submitqueue/entity/BUILD.bazel | 1 + submitqueue/entity/request_receipt.go | 25 +++ submitqueue/extension/storage/BUILD.bazel | 1 + .../extension/storage/mock/BUILD.bazel | 1 + .../mock/request_receipt_store_mock.go | 72 ++++++++ .../storage/request_receipt_store.go | 54 ++++++ submitqueue/gateway/controller/land_test.go | 10 ++ .../gateway/controller/log/log_test.go | 5 + .../controller/storage_fixture_test.go | 4 + submitqueue/gateway/core/request/BUILD.bazel | 2 + .../gateway/core/request/materializer.go | 6 + .../gateway/core/request/materializer_test.go | 16 +- .../gateway/core/request/request_receipt.go | 35 ++++ .../core/request/request_receipt_test.go | 159 ++++++++++++++++ .../extension/storage/mock/storage_mock.go | 14 ++ .../extension/storage/mysql/BUILD.bazel | 2 + .../storage/mysql/request_receipt_store.go | 134 ++++++++++++++ .../mysql/request_receipt_store_test.go | 170 ++++++++++++++++++ .../storage/mysql/schema/request_receipt.sql | 8 + .../extension/storage/mysql/storage.go | 7 + .../gateway/extension/storage/storage.go | 5 +- .../submitqueue/extension/storage/BUILD.bazel | 6 +- .../extension/storage/request_receipt.go | 112 ++++++++++++ .../submitqueue/gateway/suite_test.go | 3 + 24 files changed, 837 insertions(+), 15 deletions(-) create mode 100644 submitqueue/entity/request_receipt.go create mode 100644 submitqueue/extension/storage/mock/request_receipt_store_mock.go create mode 100644 submitqueue/extension/storage/request_receipt_store.go create mode 100644 submitqueue/gateway/core/request/request_receipt.go create mode 100644 submitqueue/gateway/core/request/request_receipt_test.go create mode 100644 submitqueue/gateway/extension/storage/mysql/request_receipt_store.go create mode 100644 submitqueue/gateway/extension/storage/mysql/request_receipt_store_test.go create mode 100644 submitqueue/gateway/extension/storage/mysql/schema/request_receipt.sql create mode 100644 test/integration/submitqueue/extension/storage/request_receipt.go diff --git a/submitqueue/entity/BUILD.bazel b/submitqueue/entity/BUILD.bazel index b0890f6f2..d3fcf8eb0 100644 --- a/submitqueue/entity/BUILD.bazel +++ b/submitqueue/entity/BUILD.bazel @@ -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", diff --git a/submitqueue/entity/request_receipt.go b/submitqueue/entity/request_receipt.go new file mode 100644 index 000000000..8e22c2ad4 --- /dev/null +++ b/submitqueue/entity/request_receipt.go @@ -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 +} diff --git a/submitqueue/extension/storage/BUILD.bazel b/submitqueue/extension/storage/BUILD.bazel index 9dcf2e624..d26ba2d1c 100644 --- a/submitqueue/extension/storage/BUILD.bazel +++ b/submitqueue/extension/storage/BUILD.bazel @@ -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", diff --git a/submitqueue/extension/storage/mock/BUILD.bazel b/submitqueue/extension/storage/mock/BUILD.bazel index 94f252551..aed6e27f8 100644 --- a/submitqueue/extension/storage/mock/BUILD.bazel +++ b/submitqueue/extension/storage/mock/BUILD.bazel @@ -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", diff --git a/submitqueue/extension/storage/mock/request_receipt_store_mock.go b/submitqueue/extension/storage/mock/request_receipt_store_mock.go new file mode 100644 index 000000000..ee76046f8 --- /dev/null +++ b/submitqueue/extension/storage/mock/request_receipt_store_mock.go @@ -0,0 +1,72 @@ +// Code generated by MockGen. DO NOT EDIT. +// Source: request_receipt_store.go +// +// Generated by this command: +// +// mockgen -source=request_receipt_store.go -destination=mock/request_receipt_store_mock.go -package=mock +// + +// Package mock is a generated GoMock package. +package mock + +import ( + context "context" + reflect "reflect" + + entity "github.com/uber/submitqueue/submitqueue/entity" + storage "github.com/uber/submitqueue/submitqueue/extension/storage" + gomock "go.uber.org/mock/gomock" +) + +// MockRequestReceiptStore is a mock of RequestReceiptStore interface. +type MockRequestReceiptStore struct { + ctrl *gomock.Controller + recorder *MockRequestReceiptStoreMockRecorder + isgomock struct{} +} + +// MockRequestReceiptStoreMockRecorder is the mock recorder for MockRequestReceiptStore. +type MockRequestReceiptStoreMockRecorder struct { + mock *MockRequestReceiptStore +} + +// NewMockRequestReceiptStore creates a new mock instance. +func NewMockRequestReceiptStore(ctrl *gomock.Controller) *MockRequestReceiptStore { + mock := &MockRequestReceiptStore{ctrl: ctrl} + mock.recorder = &MockRequestReceiptStoreMockRecorder{mock} + return mock +} + +// EXPECT returns an object that allows the caller to indicate expected use. +func (m *MockRequestReceiptStore) EXPECT() *MockRequestReceiptStoreMockRecorder { + return m.recorder +} + +// Create mocks base method. +func (m *MockRequestReceiptStore) Create(ctx context.Context, receipt entity.RequestReceipt) error { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "Create", ctx, receipt) + ret0, _ := ret[0].(error) + return ret0 +} + +// Create indicates an expected call of Create. +func (mr *MockRequestReceiptStoreMockRecorder) Create(ctx, receipt any) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Create", reflect.TypeOf((*MockRequestReceiptStore)(nil).Create), ctx, receipt) +} + +// List mocks base method. +func (m *MockRequestReceiptStore) List(ctx context.Context, bounds storage.RequestReceiptRange) ([]entity.RequestReceipt, error) { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "List", ctx, bounds) + ret0, _ := ret[0].([]entity.RequestReceipt) + ret1, _ := ret[1].(error) + return ret0, ret1 +} + +// List indicates an expected call of List. +func (mr *MockRequestReceiptStoreMockRecorder) List(ctx, bounds any) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "List", reflect.TypeOf((*MockRequestReceiptStore)(nil).List), ctx, bounds) +} diff --git a/submitqueue/extension/storage/request_receipt_store.go b/submitqueue/extension/storage/request_receipt_store.go new file mode 100644 index 000000000..740688bac --- /dev/null +++ b/submitqueue/extension/storage/request_receipt_store.go @@ -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) +} diff --git a/submitqueue/gateway/controller/land_test.go b/submitqueue/gateway/controller/land_test.go index 481ae8b4a..aab5b3602 100644 --- a/submitqueue/gateway/controller/land_test.go +++ b/submitqueue/gateway/controller/land_test.go @@ -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) @@ -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() @@ -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) @@ -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", diff --git a/submitqueue/gateway/controller/log/log_test.go b/submitqueue/gateway/controller/log/log_test.go index 41207b809..91f6844fc 100644 --- a/submitqueue/gateway/controller/log/log_test.go +++ b/submitqueue/gateway/controller/log/log_test.go @@ -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() @@ -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 } diff --git a/submitqueue/gateway/controller/storage_fixture_test.go b/submitqueue/gateway/controller/storage_fixture_test.go index 39d04a2b2..11d492530 100644 --- a/submitqueue/gateway/controller/storage_fixture_test.go +++ b/submitqueue/gateway/controller/storage_fixture_test.go @@ -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 @@ -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() @@ -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() diff --git a/submitqueue/gateway/core/request/BUILD.bazel b/submitqueue/gateway/core/request/BUILD.bazel index f5b153036..8f355ffa7 100644 --- a/submitqueue/gateway/core/request/BUILD.bazel +++ b/submitqueue/gateway/core/request/BUILD.bazel @@ -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"], @@ -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"], diff --git a/submitqueue/gateway/core/request/materializer.go b/submitqueue/gateway/core/request/materializer.go index 2b149ca61..6d44e6861 100644 --- a/submitqueue/gateway/core/request/materializer.go +++ b/submitqueue/gateway/core/request/materializer.go @@ -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 } } diff --git a/submitqueue/gateway/core/request/materializer_test.go b/submitqueue/gateway/core/request/materializer_test.go index 9704edd46..3d83033b4 100644 --- a/submitqueue/gateway/core/request/materializer_test.go +++ b/submitqueue/gateway/core/request/materializer_test.go @@ -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" ) @@ -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 { diff --git a/submitqueue/gateway/core/request/request_receipt.go b/submitqueue/gateway/core/request/request_receipt.go new file mode 100644 index 000000000..bb2f10bc4 --- /dev/null +++ b/submitqueue/gateway/core/request/request_receipt.go @@ -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 +} diff --git a/submitqueue/gateway/core/request/request_receipt_test.go b/submitqueue/gateway/core/request/request_receipt_test.go new file mode 100644 index 000000000..0428a6ff1 --- /dev/null +++ b/submitqueue/gateway/core/request/request_receipt_test.go @@ -0,0 +1,159 @@ +// 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" + "testing" + + "github.com/stretchr/testify/require" + "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" +) + +type materializerReceiptFixture struct { + materializer *Materializer + summaries *storagemock.MockRequestSummaryStore + queueSummaries *storagemock.MockRequestQueueSummaryStore + uris *storagemock.MockRequestURIStore + logs *storagemock.MockRequestLogStore + receipts *storagemock.MockRequestReceiptStore +} + +func newMaterializerReceiptFixture(ctrl *gomock.Controller) materializerReceiptFixture { + f := materializerReceiptFixture{ + summaries: storagemock.NewMockRequestSummaryStore(ctrl), queueSummaries: storagemock.NewMockRequestQueueSummaryStore(ctrl), + uris: storagemock.NewMockRequestURIStore(ctrl), logs: storagemock.NewMockRequestLogStore(ctrl), + receipts: storagemock.NewMockRequestReceiptStore(ctrl), + } + stores := gwstoragemock.NewMockStorage(ctrl) + stores.EXPECT().GetRequestSummaryStore().Return(f.summaries).AnyTimes() + stores.EXPECT().GetRequestQueueSummaryStore().Return(f.queueSummaries).AnyTimes() + stores.EXPECT().GetRequestURIStore().Return(f.uris).AnyTimes() + stores.EXPECT().GetRequestLogStore().Return(f.logs).AnyTimes() + stores.EXPECT().GetRequestReceiptStore().Return(f.receipts).AnyTimes() + factory := gwstoragemock.NewMockFactory(ctrl) + factory.EXPECT().For(gomock.Any()).Return(stores, nil).AnyTimes() + f.materializer = NewMaterializer(factory) + return f +} + +func TestMaterializer_ActivatesReceiptMapping(t *testing.T) { + for _, status := range []entity.RequestStatus{entity.RequestStatusAccepted, entity.RequestStatusStarted, entity.RequestStatusLanded} { + t.Run(string(status), func(t *testing.T) { + f := newMaterializerReceiptFixture(gomock.NewController(t)) + current := testRequestSummary() + log := entity.RequestLog{Queue: current.Queue, RequestID: current.RequestID, Type: entity.RequestLogTypeStatus, Status: status, TimestampMs: 20} + updated := current + updated.Status, updated.StatusTimestampMs, updated.Version = status, log.TimestampMs, 2 + gomock.InOrder( + f.logs.EXPECT().Insert(gomock.Any(), log).Return(nil), + f.summaries.EXPECT().Get(gomock.Any(), current.RequestID).Return(current, nil), + f.summaries.EXPECT().Update(gomock.Any(), gomock.Any(), int32(1), int32(2)).Return(nil), + f.queueSummaries.EXPECT().Get(gomock.Any(), current.ReceivedAtMs, current.RequestID).Return(entity.RequestQueueSummary{}, storage.ErrNotFound), + f.uris.EXPECT().Create(gomock.Any(), entity.RequestURI{Queue: current.Queue, ChangeURI: current.ChangeURIs[0], ReceivedAtMs: current.ReceivedAtMs, RequestID: current.RequestID}).Return(nil), + f.uris.EXPECT().Create(gomock.Any(), entity.RequestURI{Queue: current.Queue, ChangeURI: current.ChangeURIs[1], ReceivedAtMs: current.ReceivedAtMs, RequestID: current.RequestID}).Return(nil), + f.queueSummaries.EXPECT().Create(gomock.Any(), queueSummaryFromSummary(updated)).Return(nil), + f.receipts.EXPECT().Create(gomock.Any(), receiptFromTestSummary(current)).Return(nil), + ) + require.NoError(t, f.materializer.PersistLog(context.Background(), log)) + }) + } +} + +func TestMaterializer_EnsuresReceiptForUnchangedSummary(t *testing.T) { + for _, tt := range []struct { + name string + log entity.RequestLog + createErr error + }{ + {"late accepted", entity.RequestLog{Type: entity.RequestLogTypeStatus, Status: entity.RequestStatusAccepted, TimestampMs: 30}, nil}, + {"duplicate mapping", entity.RequestLog{Type: entity.RequestLogTypeStatus, Status: entity.RequestStatusLanded, TimestampMs: 20, RequestVersion: 2}, storage.ErrAlreadyExists}, + {"audit event", entity.RequestLog{Type: entity.RequestLogTypeEvent, Event: entity.RequestEventBuilt, TimestampMs: 30}, nil}, + } { + t.Run(tt.name, func(t *testing.T) { + f := newMaterializerReceiptFixture(gomock.NewController(t)) + current := testRequestSummary() + current.Status, current.RequestVersion, current.StatusTimestampMs, current.Version = entity.RequestStatusLanded, 2, 20, 2 + log := tt.log + log.Queue, log.RequestID = current.Queue, current.RequestID + f.logs.EXPECT().Insert(gomock.Any(), log).Return(nil) + f.summaries.EXPECT().Get(gomock.Any(), current.RequestID).Return(current, nil) + f.queueSummaries.EXPECT().Get(gomock.Any(), current.ReceivedAtMs, current.RequestID).Return(queueSummaryFromSummary(current), nil) + f.receipts.EXPECT().Create(gomock.Any(), receiptFromTestSummary(current)).Return(tt.createErr) + require.NoError(t, f.materializer.PersistLog(context.Background(), log)) + }) + } +} + +func TestMaterializer_RetriesReceiptMappingFailure(t *testing.T) { + f := newMaterializerReceiptFixture(gomock.NewController(t)) + current := testRequestSummary() + current.Status = entity.RequestStatusAccepted + log := entity.RequestLog{Queue: current.Queue, RequestID: current.RequestID, Type: entity.RequestLogTypeStatus, Status: entity.RequestStatusStarted, TimestampMs: 20} + updated := current + updated.Status, updated.StatusTimestampMs, updated.Version = log.Status, log.TimestampMs, 2 + f.logs.EXPECT().Insert(gomock.Any(), log).Return(nil).Times(2) + writeErr := errors.New("receipt write failed") + gomock.InOrder( + f.summaries.EXPECT().Get(gomock.Any(), current.RequestID).Return(current, nil), + f.summaries.EXPECT().Update(gomock.Any(), gomock.Any(), int32(1), int32(2)).Return(nil), + f.queueSummaries.EXPECT().Get(gomock.Any(), current.ReceivedAtMs, current.RequestID).Return(queueSummaryFromSummary(current), nil), + f.queueSummaries.EXPECT().Update(gomock.Any(), queueSummaryFromSummary(updated), int32(1), int32(2)).Return(nil), + f.receipts.EXPECT().Create(gomock.Any(), receiptFromTestSummary(current)).Return(writeErr), + f.summaries.EXPECT().Get(gomock.Any(), current.RequestID).Return(updated, nil), + f.queueSummaries.EXPECT().Get(gomock.Any(), current.ReceivedAtMs, current.RequestID).Return(queueSummaryFromSummary(updated), nil), + f.receipts.EXPECT().Create(gomock.Any(), receiptFromTestSummary(current)).Return(nil), + ) + require.ErrorIs(t, f.materializer.PersistLog(context.Background(), log), writeErr) + require.NoError(t, f.materializer.PersistLog(context.Background(), log)) +} + +func TestMaterializer_DoesNotActivateAcceptingReceipts(t *testing.T) { + for _, log := range []entity.RequestLog{ + {Type: entity.RequestLogTypeStatus, Status: entity.RequestStatusAccepting, TimestampMs: 20}, + {Type: entity.RequestLogTypeEvent, Event: entity.RequestEventBuilding, TimestampMs: 20}, + } { + t.Run(string(log.Type), func(t *testing.T) { + f := newMaterializerReceiptFixture(gomock.NewController(t)) + current := testRequestSummary() + log.Queue, log.RequestID = current.Queue, current.RequestID + f.logs.EXPECT().Insert(gomock.Any(), log).Return(nil) + f.summaries.EXPECT().Get(gomock.Any(), current.RequestID).Return(current, nil) + require.NoError(t, f.materializer.PersistLog(context.Background(), log)) + }) + } +} + +func TestMaterializer_DoesNotCreateReceiptAfterPublicProjectionFailure(t *testing.T) { + f := newMaterializerReceiptFixture(gomock.NewController(t)) + current := testRequestSummary() + current.Status = entity.RequestStatusAccepted + log := entity.RequestLog{Queue: current.Queue, RequestID: current.RequestID, Type: entity.RequestLogTypeStatus, Status: current.Status, TimestampMs: current.StatusTimestampMs} + f.logs.EXPECT().Insert(gomock.Any(), log).Return(nil) + f.summaries.EXPECT().Get(gomock.Any(), current.RequestID).Return(current, nil) + f.queueSummaries.EXPECT().Get(gomock.Any(), current.ReceivedAtMs, current.RequestID).Return(entity.RequestQueueSummary{}, storage.ErrNotFound) + writeErr := errors.New("URI write failed") + f.uris.EXPECT().Create(gomock.Any(), gomock.Any()).Return(writeErr) + require.ErrorIs(t, f.materializer.PersistLog(context.Background(), log), writeErr) +} + +func receiptFromTestSummary(summary entity.RequestSummary) entity.RequestReceipt { + return entity.RequestReceipt{Queue: summary.Queue, ReceivedAtMs: summary.ReceivedAtMs, RequestID: summary.RequestID} +} diff --git a/submitqueue/gateway/extension/storage/mock/storage_mock.go b/submitqueue/gateway/extension/storage/mock/storage_mock.go index 0f15384ab..32e927082 100644 --- a/submitqueue/gateway/extension/storage/mock/storage_mock.go +++ b/submitqueue/gateway/extension/storage/mock/storage_mock.go @@ -108,6 +108,20 @@ func (mr *MockStorageMockRecorder) GetRequestQueueSummaryStore() *gomock.Call { return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "GetRequestQueueSummaryStore", reflect.TypeOf((*MockStorage)(nil).GetRequestQueueSummaryStore)) } +// GetRequestReceiptStore mocks base method. +func (m *MockStorage) GetRequestReceiptStore() storage.RequestReceiptStore { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "GetRequestReceiptStore") + ret0, _ := ret[0].(storage.RequestReceiptStore) + return ret0 +} + +// GetRequestReceiptStore indicates an expected call of GetRequestReceiptStore. +func (mr *MockStorageMockRecorder) GetRequestReceiptStore() *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "GetRequestReceiptStore", reflect.TypeOf((*MockStorage)(nil).GetRequestReceiptStore)) +} + // GetRequestSummaryStore mocks base method. func (m *MockStorage) GetRequestSummaryStore() storage.RequestSummaryStore { m.ctrl.T.Helper() diff --git a/submitqueue/gateway/extension/storage/mysql/BUILD.bazel b/submitqueue/gateway/extension/storage/mysql/BUILD.bazel index b4fba78e6..893100c93 100644 --- a/submitqueue/gateway/extension/storage/mysql/BUILD.bazel +++ b/submitqueue/gateway/extension/storage/mysql/BUILD.bazel @@ -5,6 +5,7 @@ go_library( srcs = [ "request_log_store.go", "request_queue_summary_store.go", + "request_receipt_store.go", "request_summary_store.go", "request_uri_store.go", "storage.go", @@ -26,6 +27,7 @@ go_test( srcs = [ "request_log_store_test.go", "request_queue_summary_store_test.go", + "request_receipt_store_test.go", "request_summary_store_test.go", "request_uri_store_test.go", "storage_test.go", diff --git a/submitqueue/gateway/extension/storage/mysql/request_receipt_store.go b/submitqueue/gateway/extension/storage/mysql/request_receipt_store.go new file mode 100644 index 000000000..34e666262 --- /dev/null +++ b/submitqueue/gateway/extension/storage/mysql/request_receipt_store.go @@ -0,0 +1,134 @@ +// 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 mysql + +import ( + "context" + "database/sql" + "errors" + "fmt" + + "github.com/go-sql-driver/mysql" + "github.com/uber-go/tally" + "github.com/uber/submitqueue/platform/metrics" + "github.com/uber/submitqueue/submitqueue/entity" + "github.com/uber/submitqueue/submitqueue/extension/storage" +) + +const listRequestReceiptsQuery = ` + SELECT queue, received_at_ms, request_id + FROM request_receipt + WHERE queue = ? AND received_at_ms >= ? AND received_at_ms < ? + ORDER BY received_at_ms DESC, request_id DESC LIMIT ?` + +const listRequestReceiptsBeforeCursorQuery = ` + SELECT queue, received_at_ms, request_id + FROM request_receipt + WHERE queue = ? AND received_at_ms >= ? AND received_at_ms < ? + AND (received_at_ms < ? OR (received_at_ms = ? AND request_id < ?)) + ORDER BY received_at_ms DESC, request_id DESC LIMIT ?` + +type requestReceiptStore struct { + db *sql.DB + scope tally.Scope + queue string +} + +// NewRequestReceiptStore creates a MySQL-backed, queue-scoped RequestReceiptStore. +func NewRequestReceiptStore(db *sql.DB, scope tally.Scope, queue string) storage.RequestReceiptStore { + return &requestReceiptStore{db: db, scope: scope, queue: queue} +} + +func (r *requestReceiptStore) Create(ctx context.Context, receipt entity.RequestReceipt) (retErr error) { + op := metrics.Begin(r.scope, "create", metrics.StorageLatencyBuckets) + defer func() { op.Complete(retErr) }() + + if receipt.Queue != r.queue { + return fmt.Errorf("request receipt queue %q does not match the store's bound queue %q", receipt.Queue, r.queue) + } + if receipt.ReceivedAtMs <= 0 || receipt.RequestID == "" { + return fmt.Errorf("request receipt requires a positive timestamp and nonempty request ID") + } + _, err := r.db.ExecContext(ctx, ` + INSERT INTO request_receipt (queue, received_at_ms, request_id) + VALUES (?, ?, ?)`, receipt.Queue, receipt.ReceivedAtMs, receipt.RequestID) + if err != nil { + var mysqlErr *mysql.MySQLError + if errors.As(err, &mysqlErr) && mysqlErr.Number == mysqlErrDuplicateEntry { + return fmt.Errorf("request receipt request_id=%q: %w", receipt.RequestID, storage.ErrAlreadyExists) + } + return fmt.Errorf("failed to insert request receipt request_id=%q: %w", receipt.RequestID, err) + } + return nil +} + +func (r *requestReceiptStore) List(ctx context.Context, bounds storage.RequestReceiptRange) (ret []entity.RequestReceipt, retErr error) { + op := metrics.Begin(r.scope, "list", metrics.StorageLatencyBuckets) + defer func() { op.Complete(retErr) }() + + if err := validateRequestReceiptRange(bounds); err != nil { + return nil, err + } + return r.queryRequestReceipts(ctx, bounds) +} + +func validateRequestReceiptRange(bounds storage.RequestReceiptRange) error { + if bounds.ReceivedBeforeMs <= bounds.ReceivedAtOrAfterMs || bounds.Limit <= 0 { + return fmt.Errorf("request receipt range requires ordered bounds and a positive limit") + } + cursor := bounds.Before + if cursor.ReceivedAtMs < 0 || (cursor.ReceivedAtMs == 0) != (cursor.RequestID == "") { + return fmt.Errorf("request receipt cursor requires a positive timestamp and nonempty request ID, or the zero value") + } + return nil +} + +func (r *requestReceiptStore) queryRequestReceipts(ctx context.Context, bounds storage.RequestReceiptRange) ([]entity.RequestReceipt, error) { + var ( + rows *sql.Rows + err error + ) + cursor := bounds.Before + if cursor.ReceivedAtMs == 0 { + rows, err = r.db.QueryContext(ctx, listRequestReceiptsQuery, + r.queue, bounds.ReceivedAtOrAfterMs, bounds.ReceivedBeforeMs, bounds.Limit, + ) + } else { + rows, err = r.db.QueryContext(ctx, listRequestReceiptsBeforeCursorQuery, + r.queue, bounds.ReceivedAtOrAfterMs, bounds.ReceivedBeforeMs, + cursor.ReceivedAtMs, cursor.ReceivedAtMs, cursor.RequestID, bounds.Limit, + ) + } + if err != nil { + return nil, fmt.Errorf("failed to list request receipts queue=%q: %w", r.queue, err) + } + defer rows.Close() + return scanRequestReceiptRows(rows, r.queue) +} + +func scanRequestReceiptRows(rows *sql.Rows, queue string) ([]entity.RequestReceipt, error) { + receipts := make([]entity.RequestReceipt, 0) + for rows.Next() { + var receipt entity.RequestReceipt + if err := rows.Scan(&receipt.Queue, &receipt.ReceivedAtMs, &receipt.RequestID); err != nil { + return nil, fmt.Errorf("failed to scan request receipt queue=%q: %w", queue, err) + } + receipts = append(receipts, receipt) + } + if err := rows.Err(); err != nil { + return nil, fmt.Errorf("failed to iterate request receipts queue=%q: %w", queue, err) + } + return receipts, nil +} diff --git a/submitqueue/gateway/extension/storage/mysql/request_receipt_store_test.go b/submitqueue/gateway/extension/storage/mysql/request_receipt_store_test.go new file mode 100644 index 000000000..008c2c26b --- /dev/null +++ b/submitqueue/gateway/extension/storage/mysql/request_receipt_store_test.go @@ -0,0 +1,170 @@ +// 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 mysql + +import ( + "context" + "database/sql/driver" + "errors" + "testing" + + "github.com/DATA-DOG/go-sqlmock" + "github.com/go-sql-driver/mysql" + "github.com/stretchr/testify/require" + "github.com/uber/submitqueue/submitqueue/entity" + "github.com/uber/submitqueue/submitqueue/extension/storage" +) + +func newRequestReceiptStoreTest(t *testing.T) (sqlmock.Sqlmock, storage.RequestReceiptStore) { + t.Helper() + db, mock, err := sqlmock.New(sqlmock.QueryMatcherOption(sqlmock.QueryMatcherEqual)) + require.NoError(t, err) + t.Cleanup(func() { db.Close() }) + return mock, NewRequestReceiptStore(db, testMetrics(), "monorepo/main") +} + +func TestRequestReceiptStoreCreate(t *testing.T) { + receipt := entity.RequestReceipt{Queue: "monorepo/main", ReceivedAtMs: 1000, RequestID: "request/monorepo/main/1"} + writeErr := errors.New("write failed") + for _, tt := range []struct { + name string + err error + want error + }{ + {name: "created"}, + {name: "duplicate", err: &mysql.MySQLError{Number: mysqlErrDuplicateEntry}, want: storage.ErrAlreadyExists}, + {name: "write failure", err: writeErr, want: writeErr}, + } { + t.Run(tt.name, func(t *testing.T) { + mock, store := newRequestReceiptStoreTest(t) + write := mock.ExpectExec("INSERT INTO request_receipt (queue, received_at_ms, request_id) VALUES (?, ?, ?)"). + WithArgs(receipt.Queue, receipt.ReceivedAtMs, receipt.RequestID) + if tt.err != nil { + write.WillReturnError(tt.err) + } else { + write.WillReturnResult(sqlmock.NewResult(0, 1)) + } + err := store.Create(context.Background(), receipt) + if tt.want != nil { + require.ErrorIs(t, err, tt.want) + } else { + require.NoError(t, err) + } + require.NoError(t, mock.ExpectationsWereMet()) + }) + } +} + +func TestRequestReceiptStoreRejectsInvalidMappings(t *testing.T) { + for _, tt := range []struct { + name string + receipt entity.RequestReceipt + }{ + {"wrong queue", entity.RequestReceipt{Queue: "other", ReceivedAtMs: 1000, RequestID: "request/1"}}, + {"unknown time", entity.RequestReceipt{Queue: "monorepo/main", RequestID: "request/1"}}, + {"negative time", entity.RequestReceipt{Queue: "monorepo/main", ReceivedAtMs: -1, RequestID: "request/1"}}, + {"missing ID", entity.RequestReceipt{Queue: "monorepo/main", ReceivedAtMs: 1000}}, + } { + t.Run(tt.name, func(t *testing.T) { + mock, store := newRequestReceiptStoreTest(t) + require.Error(t, store.Create(context.Background(), tt.receipt)) + require.NoError(t, mock.ExpectationsWereMet()) + }) + } +} + +func TestRequestReceiptStoreList(t *testing.T) { + const rangeQuery = `SELECT queue, received_at_ms, request_id FROM request_receipt + WHERE queue = ? AND received_at_ms >= ? AND received_at_ms < ?` + const cursorCondition = ` AND (received_at_ms < ? OR (received_at_ms = ? AND request_id < ?))` + const orderAndLimit = ` ORDER BY received_at_ms DESC, request_id DESC LIMIT ?` + queryErr := errors.New("query failed") + rowErr := errors.New("iteration failed") + first := entity.RequestReceipt{Queue: "monorepo/main", ReceivedAtMs: 2000, RequestID: "request/9"} + second := entity.RequestReceipt{Queue: "monorepo/main", ReceivedAtMs: 2000, RequestID: "request/10"} + for _, tt := range []struct { + name string + cursor storage.RequestReceiptCursor + rows *sqlmock.Rows + err error + want []entity.RequestReceipt + fails bool + }{ + {"first page", storage.RequestReceiptCursor{}, receiptRows(first, second), nil, []entity.RequestReceipt{first, second}, false}, + {"continuation", storage.RequestReceiptCursor{ReceivedAtMs: first.ReceivedAtMs, RequestID: first.RequestID}, receiptRows(second), nil, []entity.RequestReceipt{second}, false}, + {"empty", storage.RequestReceiptCursor{}, receiptRows(), nil, []entity.RequestReceipt{}, false}, + {"query failure", storage.RequestReceiptCursor{}, nil, queryErr, nil, true}, + {"scan failure", storage.RequestReceiptCursor{}, receiptRows().AddRow("monorepo/main", "not a timestamp", "request/1"), nil, nil, true}, + {"iteration failure", storage.RequestReceiptCursor{}, receiptRows(first, second).RowError(1, rowErr), nil, nil, true}, + } { + t.Run(tt.name, func(t *testing.T) { + mock, store := newRequestReceiptStoreTest(t) + bounds := storage.RequestReceiptRange{ReceivedAtOrAfterMs: 1000, ReceivedBeforeMs: 3000, Before: tt.cursor, Limit: 2} + query := rangeQuery + args := []driver.Value{"monorepo/main", int64(1000), int64(3000)} + if tt.cursor.ReceivedAtMs != 0 { + query += cursorCondition + args = append(args, tt.cursor.ReceivedAtMs, tt.cursor.ReceivedAtMs, tt.cursor.RequestID) + } + query += orderAndLimit + args = append(args, 2) + read := mock.ExpectQuery(query).WithArgs(args...) + if tt.err != nil { + read.WillReturnError(tt.err) + } else { + read.WillReturnRows(tt.rows).RowsWillBeClosed() + } + got, err := store.List(context.Background(), bounds) + if tt.fails { + require.Error(t, err) + require.Nil(t, got) + } else { + require.NoError(t, err) + require.Equal(t, tt.want, got) + } + require.NoError(t, mock.ExpectationsWereMet()) + }) + } +} + +func receiptRows(receipts ...entity.RequestReceipt) *sqlmock.Rows { + rows := sqlmock.NewRows([]string{"queue", "received_at_ms", "request_id"}) + for _, receipt := range receipts { + rows.AddRow(receipt.Queue, receipt.ReceivedAtMs, receipt.RequestID) + } + return rows +} + +func TestRequestReceiptStoreRejectsInvalidRanges(t *testing.T) { + for _, tt := range []struct { + name string + bounds storage.RequestReceiptRange + }{ + {"empty range", storage.RequestReceiptRange{ReceivedAtOrAfterMs: 1000, ReceivedBeforeMs: 1000, Limit: 1}}, + {"inverted range", storage.RequestReceiptRange{ReceivedAtOrAfterMs: 2000, ReceivedBeforeMs: 1000, Limit: 1}}, + {"zero limit", storage.RequestReceiptRange{ReceivedBeforeMs: 2000}}, + {"negative limit", storage.RequestReceiptRange{ReceivedBeforeMs: 2000, Limit: -1}}, + {"cursor missing ID", storage.RequestReceiptRange{ReceivedBeforeMs: 2000, Limit: 1, Before: storage.RequestReceiptCursor{ReceivedAtMs: 1000}}}, + {"cursor missing time", storage.RequestReceiptRange{ReceivedBeforeMs: 2000, Limit: 1, Before: storage.RequestReceiptCursor{RequestID: "request/1"}}}, + {"cursor negative time", storage.RequestReceiptRange{ReceivedBeforeMs: 2000, Limit: 1, Before: storage.RequestReceiptCursor{ReceivedAtMs: -1, RequestID: "request/1"}}}, + } { + t.Run(tt.name, func(t *testing.T) { + mock, store := newRequestReceiptStoreTest(t) + _, err := store.List(context.Background(), tt.bounds) + require.Error(t, err) + require.NoError(t, mock.ExpectationsWereMet()) + }) + } +} diff --git a/submitqueue/gateway/extension/storage/mysql/schema/request_receipt.sql b/submitqueue/gateway/extension/storage/mysql/schema/request_receipt.sql new file mode 100644 index 000000000..44eec724b --- /dev/null +++ b/submitqueue/gateway/extension/storage/mysql/schema/request_receipt.sql @@ -0,0 +1,8 @@ +-- Immutable lookup keys; the full projection remains in request_summary. +-- Match the existing queue projection's key types to preserve List ordering. +CREATE TABLE IF NOT EXISTS request_receipt ( + queue VARCHAR(255) NOT NULL, + received_at_ms BIGINT NOT NULL, + request_id VARCHAR(255) NOT NULL, + PRIMARY KEY (queue, received_at_ms, request_id) +) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4; diff --git a/submitqueue/gateway/extension/storage/mysql/storage.go b/submitqueue/gateway/extension/storage/mysql/storage.go index 897713299..ea69d368c 100644 --- a/submitqueue/gateway/extension/storage/mysql/storage.go +++ b/submitqueue/gateway/extension/storage/mysql/storage.go @@ -51,6 +51,7 @@ func (s *Storage) For(queueName string) (storage.Storage, error) { return nil, fmt.Errorf("queue name must not be empty") } return &boundStorage{ + requestReceiptStore: NewRequestReceiptStore(s.db, s.scope.SubScope("request_receipt_store"), queueName), requestQueueStore: NewRequestQueueSummaryStore(s.db, s.scope.SubScope("request_queue_summary_store"), queueName), requestSummaryStore: NewRequestSummaryStore(s.db, s.scope.SubScope("request_summary_store"), queueName), requestLogStore: NewRequestLogStore(s.db, s.scope.SubScope("request_log_store"), queueName), @@ -65,6 +66,7 @@ func (s *Storage) Close() error { // boundStorage is the queue-scoped store aggregate returned by For. type boundStorage struct { + requestReceiptStore basestorage.RequestReceiptStore requestQueueStore basestorage.RequestQueueSummaryStore requestSummaryStore basestorage.RequestSummaryStore requestLogStore basestorage.RequestLogStore @@ -74,6 +76,11 @@ type boundStorage struct { // Verify boundStorage implements the queue-scoped aggregate at compile time. var _ storage.Storage = (*boundStorage)(nil) +// GetRequestReceiptStore returns the bound MySQL-backed RequestReceiptStore. +func (f *boundStorage) GetRequestReceiptStore() basestorage.RequestReceiptStore { + return f.requestReceiptStore +} + // GetRequestQueueSummaryStore returns the bound MySQL-backed RequestQueueSummaryStore. func (f *boundStorage) GetRequestQueueSummaryStore() basestorage.RequestQueueSummaryStore { return f.requestQueueStore diff --git a/submitqueue/gateway/extension/storage/storage.go b/submitqueue/gateway/extension/storage/storage.go index 760492157..df230180f 100644 --- a/submitqueue/gateway/extension/storage/storage.go +++ b/submitqueue/gateway/extension/storage/storage.go @@ -13,7 +13,7 @@ // limitations under the License. // Package storage resolves the gateway's queue-scoped stores: the append-only -// request log and the three read models behind request-summary retrieval and +// request log and the read models behind request-summary retrieval and // List. The store contracts themselves stay in // submitqueue/extension/storage — only the aggregate is service-scoped, so // what the gateway can reach is narrower than what the domain defines. @@ -50,6 +50,9 @@ type Storage interface { // GetRequestSummaryStore returns the RequestSummaryStore instance. GetRequestSummaryStore() basestorage.RequestSummaryStore + // GetRequestReceiptStore returns the RequestReceiptStore instance. + GetRequestReceiptStore() basestorage.RequestReceiptStore + // GetRequestQueueSummaryStore returns the RequestQueueSummaryStore instance. GetRequestQueueSummaryStore() basestorage.RequestQueueSummaryStore diff --git a/test/integration/submitqueue/extension/storage/BUILD.bazel b/test/integration/submitqueue/extension/storage/BUILD.bazel index e844c7337..a7a2f1efe 100644 --- a/test/integration/submitqueue/extension/storage/BUILD.bazel +++ b/test/integration/submitqueue/extension/storage/BUILD.bazel @@ -2,7 +2,10 @@ load("@rules_go//go:def.bzl", "go_library") go_library( name = "go_default_library", - srcs = ["suite.go"], + srcs = [ + "request_receipt.go", + "suite.go", + ], importpath = "github.com/uber/submitqueue/test/integration/submitqueue/extension/storage", visibility = ["//visibility:public"], deps = [ @@ -10,6 +13,7 @@ go_library( "//platform/base/mergestrategy:go_default_library", "//submitqueue/entity:go_default_library", "//submitqueue/extension/storage:go_default_library", + "//submitqueue/gateway/core/request:go_default_library", "//submitqueue/gateway/extension/storage:go_default_library", "//submitqueue/orchestrator/extension/storage:go_default_library", "//test/testutil:go_default_library", diff --git a/test/integration/submitqueue/extension/storage/request_receipt.go b/test/integration/submitqueue/extension/storage/request_receipt.go new file mode 100644 index 000000000..f88c6619e --- /dev/null +++ b/test/integration/submitqueue/extension/storage/request_receipt.go @@ -0,0 +1,112 @@ +// 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 + +import ( + "sync" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "github.com/uber/submitqueue/submitqueue/entity" + "github.com/uber/submitqueue/submitqueue/extension/storage" + requestcore "github.com/uber/submitqueue/submitqueue/gateway/core/request" +) + +func (s *StorageContractSuite) TestStorage_RequestReceiptListAndCursor() { + const queue = "receipt-list" + store := s.forGatewayQueue(queue).GetRequestReceiptStore() + lower := entity.RequestReceipt{Queue: queue, ReceivedAtMs: 100, RequestID: "1"} + ten := entity.RequestReceipt{Queue: queue, ReceivedAtMs: 200, RequestID: "10"} + nine := entity.RequestReceipt{Queue: queue, ReceivedAtMs: 200, RequestID: "9"} + upper := entity.RequestReceipt{Queue: queue, ReceivedAtMs: 300, RequestID: "2"} + for _, receipt := range []entity.RequestReceipt{ten, upper, lower, nine} { + require.NoError(s.T(), store.Create(s.ctx, receipt)) + } + require.ErrorIs(s.T(), store.Create(s.ctx, nine), storage.ErrAlreadyExists) + other := s.forGatewayQueue("receipt-list-other").GetRequestReceiptStore() + require.NoError(s.T(), other.Create(s.ctx, entity.RequestReceipt{Queue: "receipt-list-other", ReceivedAtMs: nine.ReceivedAtMs, RequestID: nine.RequestID})) + require.Error(s.T(), other.Create(s.ctx, nine)) + + for _, tt := range []struct { + name string + bounds storage.RequestReceiptRange + want []entity.RequestReceipt + }{ + {"bounded string order", storage.RequestReceiptRange{ReceivedAtOrAfterMs: 100, ReceivedBeforeMs: 300, Limit: 10}, []entity.RequestReceipt{nine, ten, lower}}, + {"first page", storage.RequestReceiptRange{ReceivedAtOrAfterMs: 100, ReceivedBeforeMs: 300, Limit: 1}, []entity.RequestReceipt{nine}}, + {"continuation within timestamp tie", storage.RequestReceiptRange{ReceivedAtOrAfterMs: 100, ReceivedBeforeMs: 300, Before: storage.RequestReceiptCursor{ReceivedAtMs: 200, RequestID: "9"}, Limit: 1}, []entity.RequestReceipt{ten}}, + {"continuation across timestamps", storage.RequestReceiptRange{ReceivedAtOrAfterMs: 100, ReceivedBeforeMs: 300, Before: storage.RequestReceiptCursor{ReceivedAtMs: 200, RequestID: "10"}, Limit: 1}, []entity.RequestReceipt{lower}}, + {"exhausted", storage.RequestReceiptRange{ReceivedAtOrAfterMs: 100, ReceivedBeforeMs: 300, Before: storage.RequestReceiptCursor{ReceivedAtMs: 100, RequestID: "1"}, Limit: 1}, []entity.RequestReceipt{}}, + {"empty", storage.RequestReceiptRange{ReceivedBeforeMs: 100, Limit: 10}, []entity.RequestReceipt{}}, + } { + s.Run(tt.name, func() { + got, err := store.List(s.ctx, tt.bounds) + require.NoError(s.T(), err) + assert.Equal(s.T(), tt.want, got) + }) + } + got, err := other.List(s.ctx, storage.RequestReceiptRange{ReceivedBeforeMs: 300, Limit: 10}) + require.NoError(s.T(), err) + assert.Equal(s.T(), []entity.RequestReceipt{{Queue: "receipt-list-other", ReceivedAtMs: nine.ReceivedAtMs, RequestID: nine.RequestID}}, got) +} + +func (s *StorageContractSuite) TestStorage_RequestReceiptMaterialization() { + const queue = "receipt-materialization" + stores := s.forGatewayQueue(queue) + summary := entity.RequestSummary{ + Queue: queue, RequestID: "1", ReceivedAtMs: 100, ChangeURIs: []string{"uri/receipt"}, + Status: entity.RequestStatusAccepting, StatusTimestampMs: 100, Version: 1, + } + require.NoError(s.T(), stores.GetRequestSummaryStore().Create(s.ctx, summary)) + materializer := requestcore.NewMaterializer(s.gatewayFactory) + bounds := storage.RequestReceiptRange{ReceivedBeforeMs: 1000, Limit: 10} + require.NoError(s.T(), materializer.PersistLog(s.ctx, entity.RequestLog{ + Queue: queue, RequestID: summary.RequestID, TimestampMs: 150, + Type: entity.RequestLogTypeEvent, Event: entity.RequestEventBuilding, + })) + got, err := stores.GetRequestReceiptStore().List(s.ctx, bounds) + require.NoError(s.T(), err) + assert.Empty(s.T(), got) + + logs := []entity.RequestLog{ + {Queue: queue, RequestID: summary.RequestID, Type: entity.RequestLogTypeStatus, Status: entity.RequestStatusStarted, TimestampMs: 200}, + {Queue: queue, RequestID: summary.RequestID, Type: entity.RequestLogTypeStatus, Status: entity.RequestStatusLanded, TimestampMs: 300, RequestVersion: 2}, + } + var writes sync.WaitGroup + results := make(chan error, len(logs)) + for _, log := range logs { + writes.Add(1) + go func() { + defer writes.Done() + results <- materializer.PersistLog(s.ctx, log) + }() + } + writes.Wait() + close(results) + for err := range results { + require.NoError(s.T(), err) + } + require.NoError(s.T(), materializer.PersistLog(s.ctx, entity.RequestLog{ + Queue: queue, RequestID: summary.RequestID, Type: entity.RequestLogTypeStatus, + Status: entity.RequestStatusAccepted, TimestampMs: 400, + })) + got, err = stores.GetRequestReceiptStore().List(s.ctx, bounds) + require.NoError(s.T(), err) + assert.Equal(s.T(), []entity.RequestReceipt{{Queue: queue, RequestID: summary.RequestID, ReceivedAtMs: 100}}, got) + current, err := stores.GetRequestSummaryStore().Get(s.ctx, summary.RequestID) + require.NoError(s.T(), err) + assert.Equal(s.T(), entity.RequestStatusLanded, current.Status) + assert.Equal(s.T(), summary.ReceivedAtMs, current.ReceivedAtMs) +} diff --git a/test/integration/submitqueue/gateway/suite_test.go b/test/integration/submitqueue/gateway/suite_test.go index 096190785..b80f5eb36 100644 --- a/test/integration/submitqueue/gateway/suite_test.go +++ b/test/integration/submitqueue/gateway/suite_test.go @@ -204,6 +204,7 @@ func (s *GatewayIntegrationSuite) TestListAPI() { RequestID: summary.RequestID, Queue: summary.Queue, TimestampMs: summary.StatusTimestampMs, + Type: entity.RequestLogTypeStatus, Status: publicStatus, Metadata: map[string]string{}, })) @@ -213,12 +214,14 @@ func (s *GatewayIntegrationSuite) TestListAPI() { require.NoError(t, err) require.Len(t, resp.Requests, 1) assert.Equal(t, "902", resp.Requests[0].Sqid) + assert.Equal(t, string(entity.RequestStatusLanded), resp.Requests[0].Status) require.NotEmpty(t, resp.NextPageToken) resp, err = s.client.List(s.ctx, &pb.ListRequest{Queue: "test-queue", ReceivedAtOrAfterMs: 50, ReceivedBeforeMs: 250, PageSize: 1, PageToken: resp.NextPageToken}) require.NoError(t, err) require.Len(t, resp.Requests, 1) assert.Equal(t, "901", resp.Requests[0].Sqid) + assert.Equal(t, string(entity.RequestStatusAccepted), resp.Requests[0].Status) assert.Empty(t, resp.NextPageToken) }