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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions stovepipe/entity/queue.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"`
Expand Down
11 changes: 7 additions & 4 deletions stovepipe/extension/storage/mysql/queue_store.go
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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,
Expand Down Expand Up @@ -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,
Expand Down
61 changes: 32 additions & 29 deletions stovepipe/extension/storage/mysql/queue_store_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand All @@ -60,15 +61,15 @@ 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))
},
},
{
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,
Expand All @@ -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,
Expand Down Expand Up @@ -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 {
Expand All @@ -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)
},
Expand All @@ -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)
},
Expand All @@ -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"))
},
Expand Down Expand Up @@ -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)

Expand All @@ -201,15 +204,15 @@ 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))
},
},
{
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,
Expand All @@ -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,
Expand All @@ -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,
Expand Down
13 changes: 7 additions & 6 deletions stovepipe/extension/storage/mysql/schema/queue.sql
Original file line number Diff line number Diff line change
@@ -1,11 +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,

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Need a default during the initial migration (cannot have NOT NULL without a default value), but can remove later after existing tables are updated. New field must be last in the table.

PRIMARY KEY (name)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
2 changes: 2 additions & 0 deletions test/integration/stovepipe/extension/storage/suite.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand All @@ -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)
Expand Down
Loading