From 3f8c82a0b60093c9e4ada590704f672310a95a03 Mon Sep 17 00:00:00 2001 From: mnoah1 Date: Tue, 15 Sep 2026 17:08:36 +0000 Subject: [PATCH 1/2] feat(stovepipe): persist build admission deadline Summary: Intent: - Establish durable queue state shared by time-based logical admission policies. Changes: - Add the build admission deadline to the Queue entity and MySQL backend. - Cover optimistic storage updates and deadline round-tripping. --- stovepipe/entity/queue.go | 4 ++ .../extension/storage/mysql/queue_store.go | 11 ++-- .../storage/mysql/queue_store_test.go | 61 ++++++++++--------- .../extension/storage/mysql/schema/queue.sql | 1 + .../stovepipe/extension/storage/suite.go | 2 + 5 files changed, 46 insertions(+), 33 deletions(-) diff --git a/stovepipe/entity/queue.go b/stovepipe/entity/queue.go index 23d31efba..a73669568 100644 --- a/stovepipe/entity/queue.go +++ b/stovepipe/entity/queue.go @@ -39,6 +39,10 @@ type Queue struct { // InFlightCount is the number of trunk validations admitted by process but not yet terminal. InFlightCount int32 `json:"in_flight_count"` + // BuildAdmissionNotBeforeMs is the earliest Unix-millisecond timestamp at which another + // build may be admitted. Zero means no time-based admission restriction. + BuildAdmissionNotBeforeMs int64 `json:"build_admission_not_before_ms"` + // LatestRequestID is the request id of the newest head ingest accepted for this queue. // Empty until the first request is created. Coalescing compares IDs via CompareRequestID. LatestRequestID string `json:"latest_request_id"` diff --git a/stovepipe/extension/storage/mysql/queue_store.go b/stovepipe/extension/storage/mysql/queue_store.go index aaf3cb169..e547d2304 100644 --- a/stovepipe/extension/storage/mysql/queue_store.go +++ b/stovepipe/extension/storage/mysql/queue_store.go @@ -49,11 +49,12 @@ func (q *queueStore) Create(ctx context.Context, queue entity.Queue) (retErr err } _, err := q.db.ExecContext(ctx, - `INSERT INTO queue (name, last_green_uri, in_flight_count, latest_request_id, version, last_green_request_id) - VALUES (?, ?, ?, ?, ?, ?)`, + `INSERT INTO queue (name, last_green_uri, in_flight_count, build_admission_not_before_ms, latest_request_id, version, last_green_request_id) + VALUES (?, ?, ?, ?, ?, ?, ?)`, queue.Name, queue.LastGreenURI, queue.InFlightCount, + queue.BuildAdmissionNotBeforeMs, queue.LatestRequestID, queue.Version, queue.LastGreenRequestID, @@ -78,12 +79,13 @@ func (q *queueStore) Get(ctx context.Context, name string) (ret entity.Queue, re var queue entity.Queue err := q.db.QueryRowContext(ctx, - "SELECT name, last_green_uri, in_flight_count, latest_request_id, version, last_green_request_id FROM queue WHERE name = ?", + "SELECT name, last_green_uri, in_flight_count, build_admission_not_before_ms, latest_request_id, version, last_green_request_id FROM queue WHERE name = ?", name, ).Scan( &queue.Name, &queue.LastGreenURI, &queue.InFlightCount, + &queue.BuildAdmissionNotBeforeMs, &queue.LatestRequestID, &queue.Version, &queue.LastGreenRequestID, @@ -111,10 +113,11 @@ func (q *queueStore) Update(ctx context.Context, queue entity.Queue, oldVersion, result, err := q.db.ExecContext(ctx, `UPDATE queue - SET last_green_uri = ?, in_flight_count = ?, latest_request_id = ?, version = ?, last_green_request_id = ? + SET last_green_uri = ?, in_flight_count = ?, build_admission_not_before_ms = ?, latest_request_id = ?, version = ?, last_green_request_id = ? WHERE name = ? AND version = ?`, queue.LastGreenURI, queue.InFlightCount, + queue.BuildAdmissionNotBeforeMs, queue.LatestRequestID, newVersion, queue.LastGreenRequestID, diff --git a/stovepipe/extension/storage/mysql/queue_store_test.go b/stovepipe/extension/storage/mysql/queue_store_test.go index 4bf6be9aa..a7c1f0b37 100644 --- a/stovepipe/extension/storage/mysql/queue_store_test.go +++ b/stovepipe/extension/storage/mysql/queue_store_test.go @@ -42,12 +42,13 @@ func setupQueueStoreTest(t *testing.T) (*sql.DB, sqlmock.Sqlmock, storage.QueueS func TestQueueStore_Create(t *testing.T) { queue := entity.Queue{ - Name: "monorepo/main", - LastGreenURI: "git://remote/monorepo/main/green", - LastGreenRequestID: "request/monorepo/main/1", - InFlightCount: 0, - LatestRequestID: "request/monorepo/main/1", - Version: 1, + Name: "monorepo/main", + LastGreenURI: "git://remote/monorepo/main/green", + LastGreenRequestID: "request/monorepo/main/1", + InFlightCount: 0, + BuildAdmissionNotBeforeMs: 1234, + LatestRequestID: "request/monorepo/main/1", + Version: 1, } tests := []struct { @@ -60,7 +61,7 @@ func TestQueueStore_Create(t *testing.T) { name: "success", setup: func(mock sqlmock.Sqlmock) { mock.ExpectExec("INSERT INTO queue"). - WithArgs(queue.Name, queue.LastGreenURI, queue.InFlightCount, queue.LatestRequestID, queue.Version, queue.LastGreenRequestID). + WithArgs(queue.Name, queue.LastGreenURI, queue.InFlightCount, queue.BuildAdmissionNotBeforeMs, queue.LatestRequestID, queue.Version, queue.LastGreenRequestID). WillReturnResult(sqlmock.NewResult(0, 1)) }, }, @@ -68,7 +69,7 @@ func TestQueueStore_Create(t *testing.T) { name: "duplicate name returns ErrAlreadyExists", setup: func(mock sqlmock.Sqlmock) { mock.ExpectExec("INSERT INTO queue"). - WithArgs(queue.Name, queue.LastGreenURI, queue.InFlightCount, queue.LatestRequestID, queue.Version, queue.LastGreenRequestID). + WithArgs(queue.Name, queue.LastGreenURI, queue.InFlightCount, queue.BuildAdmissionNotBeforeMs, queue.LatestRequestID, queue.Version, queue.LastGreenRequestID). WillReturnError(&mysql.MySQLError{Number: mysqlErrDuplicateEntry}) }, wantErr: true, @@ -78,7 +79,7 @@ func TestQueueStore_Create(t *testing.T) { name: "other exec error", setup: func(mock sqlmock.Sqlmock) { mock.ExpectExec("INSERT INTO queue"). - WithArgs(queue.Name, queue.LastGreenURI, queue.InFlightCount, queue.LatestRequestID, queue.Version, queue.LastGreenRequestID). + WithArgs(queue.Name, queue.LastGreenURI, queue.InFlightCount, queue.BuildAdmissionNotBeforeMs, queue.LatestRequestID, queue.Version, queue.LastGreenRequestID). WillReturnError(fmt.Errorf("connection reset")) }, wantErr: true, @@ -108,12 +109,13 @@ func TestQueueStore_Create(t *testing.T) { func TestQueueStore_Get(t *testing.T) { want := entity.Queue{ - Name: "monorepo/main", - LastGreenURI: "git://remote/monorepo/main/green", - LastGreenRequestID: "request/monorepo/main/2", - InFlightCount: 2, - LatestRequestID: "request/monorepo/main/3", - Version: 3, + Name: "monorepo/main", + LastGreenURI: "git://remote/monorepo/main/green", + LastGreenRequestID: "request/monorepo/main/2", + InFlightCount: 2, + BuildAdmissionNotBeforeMs: 5678, + LatestRequestID: "request/monorepo/main/3", + Version: 3, } tests := []struct { @@ -128,9 +130,9 @@ func TestQueueStore_Get(t *testing.T) { name: "found", queueName: want.Name, setup: func(mock sqlmock.Sqlmock) { - rows := sqlmock.NewRows([]string{"name", "last_green_uri", "in_flight_count", "latest_request_id", "version", "last_green_request_id"}). - AddRow(want.Name, want.LastGreenURI, want.InFlightCount, want.LatestRequestID, want.Version, want.LastGreenRequestID) - mock.ExpectQuery("SELECT name, last_green_uri, in_flight_count, latest_request_id, version, last_green_request_id FROM queue"). + rows := sqlmock.NewRows([]string{"name", "last_green_uri", "in_flight_count", "build_admission_not_before_ms", "latest_request_id", "version", "last_green_request_id"}). + AddRow(want.Name, want.LastGreenURI, want.InFlightCount, want.BuildAdmissionNotBeforeMs, want.LatestRequestID, want.Version, want.LastGreenRequestID) + mock.ExpectQuery("SELECT name, last_green_uri, in_flight_count, build_admission_not_before_ms, latest_request_id, version, last_green_request_id FROM queue"). WithArgs(want.Name). WillReturnRows(rows) }, @@ -140,7 +142,7 @@ func TestQueueStore_Get(t *testing.T) { name: "not found", queueName: want.Name, setup: func(mock sqlmock.Sqlmock) { - mock.ExpectQuery("SELECT name, last_green_uri, in_flight_count, latest_request_id, version, last_green_request_id FROM queue"). + mock.ExpectQuery("SELECT name, last_green_uri, in_flight_count, build_admission_not_before_ms, latest_request_id, version, last_green_request_id FROM queue"). WithArgs(want.Name). WillReturnError(sql.ErrNoRows) }, @@ -151,7 +153,7 @@ func TestQueueStore_Get(t *testing.T) { name: "query error", queueName: want.Name, setup: func(mock sqlmock.Sqlmock) { - mock.ExpectQuery("SELECT name, last_green_uri, in_flight_count, latest_request_id, version, last_green_request_id FROM queue"). + mock.ExpectQuery("SELECT name, last_green_uri, in_flight_count, build_admission_not_before_ms, latest_request_id, version, last_green_request_id FROM queue"). WithArgs(want.Name). WillReturnError(fmt.Errorf("connection reset")) }, @@ -183,11 +185,12 @@ func TestQueueStore_Get(t *testing.T) { func TestQueueStore_Update(t *testing.T) { queue := entity.Queue{ - Name: "monorepo/main", - LastGreenURI: "git://remote/monorepo/main/green", - LastGreenRequestID: "request/monorepo/main/1", - InFlightCount: 1, - LatestRequestID: "request/monorepo/main/2", + Name: "monorepo/main", + LastGreenURI: "git://remote/monorepo/main/green", + LastGreenRequestID: "request/monorepo/main/1", + InFlightCount: 1, + BuildAdmissionNotBeforeMs: 9012, + LatestRequestID: "request/monorepo/main/2", } const oldVersion, newVersion = int32(1), int32(2) @@ -201,7 +204,7 @@ func TestQueueStore_Update(t *testing.T) { name: "success", setup: func(mock sqlmock.Sqlmock) { mock.ExpectExec("UPDATE queue"). - WithArgs(queue.LastGreenURI, queue.InFlightCount, queue.LatestRequestID, newVersion, queue.LastGreenRequestID, queue.Name, oldVersion). + WithArgs(queue.LastGreenURI, queue.InFlightCount, queue.BuildAdmissionNotBeforeMs, queue.LatestRequestID, newVersion, queue.LastGreenRequestID, queue.Name, oldVersion). WillReturnResult(sqlmock.NewResult(0, 1)) }, }, @@ -209,7 +212,7 @@ func TestQueueStore_Update(t *testing.T) { name: "version mismatch", setup: func(mock sqlmock.Sqlmock) { mock.ExpectExec("UPDATE queue"). - WithArgs(queue.LastGreenURI, queue.InFlightCount, queue.LatestRequestID, newVersion, queue.LastGreenRequestID, queue.Name, oldVersion). + WithArgs(queue.LastGreenURI, queue.InFlightCount, queue.BuildAdmissionNotBeforeMs, queue.LatestRequestID, newVersion, queue.LastGreenRequestID, queue.Name, oldVersion). WillReturnResult(sqlmock.NewResult(0, 0)) }, wantErr: true, @@ -219,7 +222,7 @@ func TestQueueStore_Update(t *testing.T) { name: "exec error", setup: func(mock sqlmock.Sqlmock) { mock.ExpectExec("UPDATE queue"). - WithArgs(queue.LastGreenURI, queue.InFlightCount, queue.LatestRequestID, newVersion, queue.LastGreenRequestID, queue.Name, oldVersion). + WithArgs(queue.LastGreenURI, queue.InFlightCount, queue.BuildAdmissionNotBeforeMs, queue.LatestRequestID, newVersion, queue.LastGreenRequestID, queue.Name, oldVersion). WillReturnError(fmt.Errorf("connection reset")) }, wantErr: true, @@ -228,7 +231,7 @@ func TestQueueStore_Update(t *testing.T) { name: "rows affected error", setup: func(mock sqlmock.Sqlmock) { mock.ExpectExec("UPDATE queue"). - WithArgs(queue.LastGreenURI, queue.InFlightCount, queue.LatestRequestID, newVersion, queue.LastGreenRequestID, queue.Name, oldVersion). + WithArgs(queue.LastGreenURI, queue.InFlightCount, queue.BuildAdmissionNotBeforeMs, queue.LatestRequestID, newVersion, queue.LastGreenRequestID, queue.Name, oldVersion). WillReturnResult(sqlmock.NewErrorResult(fmt.Errorf("driver error"))) }, wantErr: true, diff --git a/stovepipe/extension/storage/mysql/schema/queue.sql b/stovepipe/extension/storage/mysql/schema/queue.sql index c10463f0f..ec1a377e3 100644 --- a/stovepipe/extension/storage/mysql/schema/queue.sql +++ b/stovepipe/extension/storage/mysql/schema/queue.sql @@ -7,5 +7,6 @@ CREATE TABLE IF NOT EXISTS queue ( latest_request_id VARCHAR(255) NOT NULL DEFAULT '', version INT NOT NULL, last_green_request_id VARCHAR(255) NOT NULL DEFAULT '', + build_admission_not_before_ms BIGINT NOT NULL DEFAULT 0, PRIMARY KEY (name) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4; diff --git a/test/integration/stovepipe/extension/storage/suite.go b/test/integration/stovepipe/extension/storage/suite.go index 8e5665093..b5d5a3ff1 100644 --- a/test/integration/stovepipe/extension/storage/suite.go +++ b/test/integration/stovepipe/extension/storage/suite.go @@ -142,6 +142,7 @@ func (s *QueueStoreContractSuite) TestQueueStore_UpdateCAS() { updated.LastGreenRequestID = "request/contract/update-cas/41" updated.LatestRequestID = "request/contract/update-cas/42" updated.InFlightCount = 1 + updated.BuildAdmissionNotBeforeMs = 123456789 require.NoError(t, s.storeFor(name).Update(s.ctx, updated, 1, 2)) got, err := s.storeFor(name).Get(s.ctx, name) @@ -150,6 +151,7 @@ func (s *QueueStoreContractSuite) TestQueueStore_UpdateCAS() { assert.Equal(t, "request/contract/update-cas/41", got.LastGreenRequestID) assert.Equal(t, "request/contract/update-cas/42", got.LatestRequestID) assert.Equal(t, int32(1), got.InFlightCount) + assert.Equal(t, int64(123456789), got.BuildAdmissionNotBeforeMs) assert.Equal(t, int32(2), got.Version) err = s.storeFor(name).Update(s.ctx, updated, 1, 2) From f7d2bab5e4dff5aef333d102498c20232d8d9ff4 Mon Sep 17 00:00:00 2001 From: mnoah1 Date: Tue, 15 Sep 2026 21:09:15 +0000 Subject: [PATCH 2/2] style(stovepipe): align queue schema columns --- stovepipe/extension/storage/mysql/schema/queue.sql | 12 ++++++------ 1 file changed, 6 insertions(+), 6 deletions(-) diff --git a/stovepipe/extension/storage/mysql/schema/queue.sql b/stovepipe/extension/storage/mysql/schema/queue.sql index ec1a377e3..3811558f1 100644 --- a/stovepipe/extension/storage/mysql/schema/queue.sql +++ b/stovepipe/extension/storage/mysql/schema/queue.sql @@ -1,12 +1,12 @@ -- queue holds per-queue coordination state for the validation pipeline: the last-green -- bookmark, in-flight gate count, and latest-request id pointer. CREATE TABLE IF NOT EXISTS queue ( - name VARCHAR(255) NOT NULL, - last_green_uri VARCHAR(255) NOT NULL DEFAULT '', - in_flight_count INT NOT NULL DEFAULT 0, - latest_request_id VARCHAR(255) NOT NULL DEFAULT '', - version INT NOT NULL, - last_green_request_id VARCHAR(255) NOT NULL DEFAULT '', + name VARCHAR(255) NOT NULL, + last_green_uri VARCHAR(255) NOT NULL DEFAULT '', + in_flight_count INT NOT NULL DEFAULT 0, + latest_request_id VARCHAR(255) NOT NULL DEFAULT '', + version INT NOT NULL, + last_green_request_id VARCHAR(255) NOT NULL DEFAULT '', build_admission_not_before_ms BIGINT NOT NULL DEFAULT 0, PRIMARY KEY (name) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;