Skip to content

Commit bbe75df

Browse files
committed
feat(stovepipe): add Queue.LastGreenRequestID so the last-green bookmark can only move forward
1 parent 74a734a commit bbe75df

5 files changed

Lines changed: 54 additions & 39 deletions

File tree

stovepipe/entity/queue.go

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -31,6 +31,11 @@ type Queue struct {
3131
// whole-repo validation recorded green (health degree 0). Empty until the first such outcome.
3232
LastGreenURI string `json:"last_green_uri"`
3333

34+
// LastGreenRequestID is the request that established LastGreenURI. The bookmark only
35+
// moves forward: a green outcome adopts the pair only when this id is empty or older
36+
// than the candidate's, compared by ingest order via CompareRequestID.
37+
LastGreenRequestID string `json:"last_green_request_id"`
38+
3439
// InFlightCount is the number of trunk validations admitted by process but not yet terminal.
3540
InFlightCount int32 `json:"in_flight_count"`
3641

stovepipe/extension/storage/mysql/queue_store.go

Lines changed: 7 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -49,13 +49,14 @@ func (q *queueStore) Create(ctx context.Context, queue entity.Queue) (retErr err
4949
}
5050

5151
_, err := q.db.ExecContext(ctx,
52-
`INSERT INTO queue (name, last_green_uri, in_flight_count, latest_request_id, version)
53-
VALUES (?, ?, ?, ?, ?)`,
52+
`INSERT INTO queue (name, last_green_uri, in_flight_count, latest_request_id, version, last_green_request_id)
53+
VALUES (?, ?, ?, ?, ?, ?)`,
5454
queue.Name,
5555
queue.LastGreenURI,
5656
queue.InFlightCount,
5757
queue.LatestRequestID,
5858
queue.Version,
59+
queue.LastGreenRequestID,
5960
)
6061
if err != nil {
6162
if isDuplicateEntry(err) {
@@ -77,14 +78,15 @@ func (q *queueStore) Get(ctx context.Context, name string) (ret entity.Queue, re
7778

7879
var queue entity.Queue
7980
err := q.db.QueryRowContext(ctx,
80-
"SELECT name, last_green_uri, in_flight_count, latest_request_id, version FROM queue WHERE name = ?",
81+
"SELECT name, last_green_uri, in_flight_count, latest_request_id, version, last_green_request_id FROM queue WHERE name = ?",
8182
name,
8283
).Scan(
8384
&queue.Name,
8485
&queue.LastGreenURI,
8586
&queue.InFlightCount,
8687
&queue.LatestRequestID,
8788
&queue.Version,
89+
&queue.LastGreenRequestID,
8890
)
8991

9092
if errors.Is(err, sql.ErrNoRows) {
@@ -109,12 +111,13 @@ func (q *queueStore) Update(ctx context.Context, queue entity.Queue, oldVersion,
109111

110112
result, err := q.db.ExecContext(ctx,
111113
`UPDATE queue
112-
SET last_green_uri = ?, in_flight_count = ?, latest_request_id = ?, version = ?
114+
SET last_green_uri = ?, in_flight_count = ?, latest_request_id = ?, version = ?, last_green_request_id = ?
113115
WHERE name = ? AND version = ?`,
114116
queue.LastGreenURI,
115117
queue.InFlightCount,
116118
queue.LatestRequestID,
117119
newVersion,
120+
queue.LastGreenRequestID,
118121
queue.Name,
119122
oldVersion,
120123
)

stovepipe/extension/storage/mysql/queue_store_test.go

Lines changed: 29 additions & 26 deletions
Original file line numberDiff line numberDiff line change
@@ -42,11 +42,12 @@ func setupQueueStoreTest(t *testing.T) (*sql.DB, sqlmock.Sqlmock, storage.QueueS
4242

4343
func TestQueueStore_Create(t *testing.T) {
4444
queue := entity.Queue{
45-
Name: "monorepo/main",
46-
LastGreenURI: "git://remote/monorepo/main/green",
47-
InFlightCount: 0,
48-
LatestRequestID: "request/monorepo/main/1",
49-
Version: 1,
45+
Name: "monorepo/main",
46+
LastGreenURI: "git://remote/monorepo/main/green",
47+
LastGreenRequestID: "request/monorepo/main/1",
48+
InFlightCount: 0,
49+
LatestRequestID: "request/monorepo/main/1",
50+
Version: 1,
5051
}
5152

5253
tests := []struct {
@@ -59,15 +60,15 @@ func TestQueueStore_Create(t *testing.T) {
5960
name: "success",
6061
setup: func(mock sqlmock.Sqlmock) {
6162
mock.ExpectExec("INSERT INTO queue").
62-
WithArgs(queue.Name, queue.LastGreenURI, queue.InFlightCount, queue.LatestRequestID, queue.Version).
63+
WithArgs(queue.Name, queue.LastGreenURI, queue.InFlightCount, queue.LatestRequestID, queue.Version, queue.LastGreenRequestID).
6364
WillReturnResult(sqlmock.NewResult(0, 1))
6465
},
6566
},
6667
{
6768
name: "duplicate name returns ErrAlreadyExists",
6869
setup: func(mock sqlmock.Sqlmock) {
6970
mock.ExpectExec("INSERT INTO queue").
70-
WithArgs(queue.Name, queue.LastGreenURI, queue.InFlightCount, queue.LatestRequestID, queue.Version).
71+
WithArgs(queue.Name, queue.LastGreenURI, queue.InFlightCount, queue.LatestRequestID, queue.Version, queue.LastGreenRequestID).
7172
WillReturnError(&mysql.MySQLError{Number: mysqlErrDuplicateEntry})
7273
},
7374
wantErr: true,
@@ -77,7 +78,7 @@ func TestQueueStore_Create(t *testing.T) {
7778
name: "other exec error",
7879
setup: func(mock sqlmock.Sqlmock) {
7980
mock.ExpectExec("INSERT INTO queue").
80-
WithArgs(queue.Name, queue.LastGreenURI, queue.InFlightCount, queue.LatestRequestID, queue.Version).
81+
WithArgs(queue.Name, queue.LastGreenURI, queue.InFlightCount, queue.LatestRequestID, queue.Version, queue.LastGreenRequestID).
8182
WillReturnError(fmt.Errorf("connection reset"))
8283
},
8384
wantErr: true,
@@ -107,11 +108,12 @@ func TestQueueStore_Create(t *testing.T) {
107108

108109
func TestQueueStore_Get(t *testing.T) {
109110
want := entity.Queue{
110-
Name: "monorepo/main",
111-
LastGreenURI: "git://remote/monorepo/main/green",
112-
InFlightCount: 2,
113-
LatestRequestID: "request/monorepo/main/3",
114-
Version: 3,
111+
Name: "monorepo/main",
112+
LastGreenURI: "git://remote/monorepo/main/green",
113+
LastGreenRequestID: "request/monorepo/main/2",
114+
InFlightCount: 2,
115+
LatestRequestID: "request/monorepo/main/3",
116+
Version: 3,
115117
}
116118

117119
tests := []struct {
@@ -126,9 +128,9 @@ func TestQueueStore_Get(t *testing.T) {
126128
name: "found",
127129
queueName: want.Name,
128130
setup: func(mock sqlmock.Sqlmock) {
129-
rows := sqlmock.NewRows([]string{"name", "last_green_uri", "in_flight_count", "latest_request_id", "version"}).
130-
AddRow(want.Name, want.LastGreenURI, want.InFlightCount, want.LatestRequestID, want.Version)
131-
mock.ExpectQuery("SELECT name, last_green_uri, in_flight_count, latest_request_id, version FROM queue").
131+
rows := sqlmock.NewRows([]string{"name", "last_green_uri", "in_flight_count", "latest_request_id", "version", "last_green_request_id"}).
132+
AddRow(want.Name, want.LastGreenURI, want.InFlightCount, want.LatestRequestID, want.Version, want.LastGreenRequestID)
133+
mock.ExpectQuery("SELECT name, last_green_uri, in_flight_count, latest_request_id, version, last_green_request_id FROM queue").
132134
WithArgs(want.Name).
133135
WillReturnRows(rows)
134136
},
@@ -138,7 +140,7 @@ func TestQueueStore_Get(t *testing.T) {
138140
name: "not found",
139141
queueName: want.Name,
140142
setup: func(mock sqlmock.Sqlmock) {
141-
mock.ExpectQuery("SELECT name, last_green_uri, in_flight_count, latest_request_id, version FROM queue").
143+
mock.ExpectQuery("SELECT name, last_green_uri, in_flight_count, latest_request_id, version, last_green_request_id FROM queue").
142144
WithArgs(want.Name).
143145
WillReturnError(sql.ErrNoRows)
144146
},
@@ -149,7 +151,7 @@ func TestQueueStore_Get(t *testing.T) {
149151
name: "query error",
150152
queueName: want.Name,
151153
setup: func(mock sqlmock.Sqlmock) {
152-
mock.ExpectQuery("SELECT name, last_green_uri, in_flight_count, latest_request_id, version FROM queue").
154+
mock.ExpectQuery("SELECT name, last_green_uri, in_flight_count, latest_request_id, version, last_green_request_id FROM queue").
153155
WithArgs(want.Name).
154156
WillReturnError(fmt.Errorf("connection reset"))
155157
},
@@ -181,10 +183,11 @@ func TestQueueStore_Get(t *testing.T) {
181183

182184
func TestQueueStore_Update(t *testing.T) {
183185
queue := entity.Queue{
184-
Name: "monorepo/main",
185-
LastGreenURI: "git://remote/monorepo/main/green",
186-
InFlightCount: 1,
187-
LatestRequestID: "request/monorepo/main/2",
186+
Name: "monorepo/main",
187+
LastGreenURI: "git://remote/monorepo/main/green",
188+
LastGreenRequestID: "request/monorepo/main/1",
189+
InFlightCount: 1,
190+
LatestRequestID: "request/monorepo/main/2",
188191
}
189192
const oldVersion, newVersion = int32(1), int32(2)
190193

@@ -198,15 +201,15 @@ func TestQueueStore_Update(t *testing.T) {
198201
name: "success",
199202
setup: func(mock sqlmock.Sqlmock) {
200203
mock.ExpectExec("UPDATE queue").
201-
WithArgs(queue.LastGreenURI, queue.InFlightCount, queue.LatestRequestID, newVersion, queue.Name, oldVersion).
204+
WithArgs(queue.LastGreenURI, queue.InFlightCount, queue.LatestRequestID, newVersion, queue.LastGreenRequestID, queue.Name, oldVersion).
202205
WillReturnResult(sqlmock.NewResult(0, 1))
203206
},
204207
},
205208
{
206209
name: "version mismatch",
207210
setup: func(mock sqlmock.Sqlmock) {
208211
mock.ExpectExec("UPDATE queue").
209-
WithArgs(queue.LastGreenURI, queue.InFlightCount, queue.LatestRequestID, newVersion, queue.Name, oldVersion).
212+
WithArgs(queue.LastGreenURI, queue.InFlightCount, queue.LatestRequestID, newVersion, queue.LastGreenRequestID, queue.Name, oldVersion).
210213
WillReturnResult(sqlmock.NewResult(0, 0))
211214
},
212215
wantErr: true,
@@ -216,7 +219,7 @@ func TestQueueStore_Update(t *testing.T) {
216219
name: "exec error",
217220
setup: func(mock sqlmock.Sqlmock) {
218221
mock.ExpectExec("UPDATE queue").
219-
WithArgs(queue.LastGreenURI, queue.InFlightCount, queue.LatestRequestID, newVersion, queue.Name, oldVersion).
222+
WithArgs(queue.LastGreenURI, queue.InFlightCount, queue.LatestRequestID, newVersion, queue.LastGreenRequestID, queue.Name, oldVersion).
220223
WillReturnError(fmt.Errorf("connection reset"))
221224
},
222225
wantErr: true,
@@ -225,7 +228,7 @@ func TestQueueStore_Update(t *testing.T) {
225228
name: "rows affected error",
226229
setup: func(mock sqlmock.Sqlmock) {
227230
mock.ExpectExec("UPDATE queue").
228-
WithArgs(queue.LastGreenURI, queue.InFlightCount, queue.LatestRequestID, newVersion, queue.Name, oldVersion).
231+
WithArgs(queue.LastGreenURI, queue.InFlightCount, queue.LatestRequestID, newVersion, queue.LastGreenRequestID, queue.Name, oldVersion).
229232
WillReturnResult(sqlmock.NewErrorResult(fmt.Errorf("driver error")))
230233
},
231234
wantErr: true,
Lines changed: 6 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -1,10 +1,11 @@
11
-- queue holds per-queue coordination state for the validation pipeline: the last-green
22
-- bookmark, in-flight gate count, and latest-request id pointer.
33
CREATE TABLE IF NOT EXISTS queue (
4-
name VARCHAR(255) NOT NULL,
5-
last_green_uri VARCHAR(255) NOT NULL DEFAULT '',
6-
in_flight_count INT NOT NULL DEFAULT 0,
7-
latest_request_id VARCHAR(255) NOT NULL DEFAULT '',
8-
version INT NOT NULL,
4+
name VARCHAR(255) NOT NULL,
5+
last_green_uri VARCHAR(255) NOT NULL DEFAULT '',
6+
in_flight_count INT NOT NULL DEFAULT 0,
7+
latest_request_id VARCHAR(255) NOT NULL DEFAULT '',
8+
version INT NOT NULL,
9+
last_green_request_id VARCHAR(255) NOT NULL DEFAULT '',
910
PRIMARY KEY (name)
1011
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;

test/integration/stovepipe/extension/storage/suite.go

Lines changed: 7 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -87,10 +87,11 @@ func (s *QueueStoreContractSuite) TestQueueStore_CreateWithFields() {
8787
const name = "contract/defaults"
8888

8989
toCreate := entity.Queue{
90-
Name: name,
91-
LastGreenURI: "git://remote/monorepo/main/green-bbbb",
92-
LatestRequestID: "request/contract/defaults/99",
93-
Version: 1,
90+
Name: name,
91+
LastGreenURI: "git://remote/monorepo/main/green-bbbb",
92+
LastGreenRequestID: "request/contract/defaults/98",
93+
LatestRequestID: "request/contract/defaults/99",
94+
Version: 1,
9495
}
9596
require.NoError(t, s.storeFor(name).Create(s.ctx, toCreate))
9697

@@ -138,13 +139,15 @@ func (s *QueueStoreContractSuite) TestQueueStore_UpdateCAS() {
138139

139140
updated := created
140141
updated.LastGreenURI = "git://remote/monorepo/main/green-cccc"
142+
updated.LastGreenRequestID = "request/contract/update-cas/41"
141143
updated.LatestRequestID = "request/contract/update-cas/42"
142144
updated.InFlightCount = 1
143145
require.NoError(t, s.storeFor(name).Update(s.ctx, updated, 1, 2))
144146

145147
got, err := s.storeFor(name).Get(s.ctx, name)
146148
require.NoError(t, err)
147149
assert.Equal(t, updated.LastGreenURI, got.LastGreenURI)
150+
assert.Equal(t, "request/contract/update-cas/41", got.LastGreenRequestID)
148151
assert.Equal(t, "request/contract/update-cas/42", got.LatestRequestID)
149152
assert.Equal(t, int32(1), got.InFlightCount)
150153
assert.Equal(t, int32(2), got.Version)

0 commit comments

Comments
 (0)