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) }