Skip to content

Commit a1dec26

Browse files
committed
Code review
1 parent c2d519c commit a1dec26

5 files changed

Lines changed: 76 additions & 118 deletions

File tree

doc/rfc/stovepipe/steps/process.md

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -205,7 +205,7 @@ Transitions use the repo's optimistic-locking pattern: compute `newVersion = old
205205

206206
New key/value-shaped operations (single-key reads/writes, no server-side filtering or aggregation):
207207

208-
- **`QueueStore`** (new): `GetOrCreate(ctx, name, defaults)`, `Get(ctx, name)`, and `Update(ctx, queue, oldVersion, newVersion)` (CAS). Ingest `GetOrCreate`s and CASes `latest_request_seq`; `process` CASes `in_flight_count`; `record` CASes `last_green_uri` + `in_flight_count`.
208+
- **`QueueStore`** (new): `Create(ctx, queue)`, `Get(ctx, name)`, and `Update(ctx, queue, oldVersion, newVersion)` (CAS). Callers orchestrate get-or-create; ingest CASes `latest_request_seq`; `process` CASes `in_flight_count`; `record` CASes `last_green_uri` + `in_flight_count`.
209209
- **`RequestStore`**: no new methods — the added `Request` fields ride the existing `Create`/`Update` CAS.
210210

211211
No "list requests by queue/state" query is introduced; coalescing uses the single-row `latest_request_seq` pointer instead, keeping the contract satisfiable by a plain KV backend.

stovepipe/extension/storage/mock/queue_store_mock.go

Lines changed: 17 additions & 24 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

stovepipe/extension/storage/mysql/queue_store.go

Lines changed: 17 additions & 41 deletions
Original file line numberDiff line numberDiff line change
@@ -36,32 +36,27 @@ func NewQueueStore(db *sql.DB, scope tally.Scope) storage.QueueStore {
3636
return &queueStore{db: db, scope: scope}
3737
}
3838

39-
// GetOrCreate returns the queue row for name, creating it with defaults if absent.
40-
func (q *queueStore) GetOrCreate(ctx context.Context, name string, defaults entity.Queue) (ret entity.Queue, retErr error) {
41-
op := metrics.Begin(q.scope, "get_or_create")
39+
// Create persists a new queue row. Returns ErrAlreadyExists if the name already exists.
40+
func (q *queueStore) Create(ctx context.Context, queue entity.Queue) (retErr error) {
41+
op := metrics.Begin(q.scope, "create")
4242
defer func() { op.Complete(retErr) }()
4343

44-
queue, err := q.Get(ctx, name)
45-
if err == nil {
46-
return queue, nil
47-
}
48-
if !storage.IsNotFound(err) {
49-
return entity.Queue{}, err
50-
}
51-
52-
toCreate := defaults
53-
toCreate.Name = name
54-
if toCreate.Version == 0 {
55-
toCreate.Version = 1
56-
}
57-
58-
if err := q.create(ctx, toCreate); err != nil {
59-
if errors.Is(err, storage.ErrAlreadyExists) {
60-
return q.Get(ctx, name)
44+
_, err := q.db.ExecContext(ctx,
45+
`INSERT INTO queue (name, last_green_uri, in_flight_count, latest_request_seq, version)
46+
VALUES (?, ?, ?, ?, ?)`,
47+
queue.Name,
48+
queue.LastGreenURI,
49+
queue.InFlightCount,
50+
queue.LatestRequestSeq,
51+
queue.Version,
52+
)
53+
if err != nil {
54+
if isDuplicateEntry(err) {
55+
return fmt.Errorf("queue name=%s: %w", queue.Name, storage.ErrAlreadyExists)
6156
}
62-
return entity.Queue{}, err
57+
return fmt.Errorf("failed to insert queue name=%s: %w", queue.Name, err)
6358
}
64-
return toCreate, nil
59+
return nil
6560
}
6661

6762
// Get retrieves a queue by name. Returns ErrNotFound if the queue is not found.
@@ -132,22 +127,3 @@ func (q *queueStore) Update(ctx context.Context, queue entity.Queue, oldVersion,
132127

133128
return nil
134129
}
135-
136-
func (q *queueStore) create(ctx context.Context, queue entity.Queue) error {
137-
_, err := q.db.ExecContext(ctx,
138-
`INSERT INTO queue (name, last_green_uri, in_flight_count, latest_request_seq, version)
139-
VALUES (?, ?, ?, ?, ?)`,
140-
queue.Name,
141-
queue.LastGreenURI,
142-
queue.InFlightCount,
143-
queue.LatestRequestSeq,
144-
queue.Version,
145-
)
146-
if err != nil {
147-
if isDuplicateEntry(err) {
148-
return fmt.Errorf("queue name=%s: %w", queue.Name, storage.ErrAlreadyExists)
149-
}
150-
return fmt.Errorf("failed to insert queue name=%s: %w", queue.Name, err)
151-
}
152-
return nil
153-
}

stovepipe/extension/storage/queue_store.go

Lines changed: 3 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -24,10 +24,9 @@ import (
2424

2525
// QueueStore persists per-queue coordination rows, keyed by queue name.
2626
type QueueStore interface {
27-
// GetOrCreate returns the queue row for name, creating it with defaults if absent.
28-
// defaults supplies initial field values for a new row (Name is set from name).
29-
// On a create race, re-reads and returns the canonical row.
30-
GetOrCreate(ctx context.Context, name string, defaults entity.Queue) (entity.Queue, error)
27+
// Create persists a new queue row. queue.Name must be set. Returns ErrAlreadyExists
28+
// if a row with the same name already exists.
29+
Create(ctx context.Context, queue entity.Queue) error
3130

3231
// Get retrieves a queue by name. Returns ErrNotFound if the queue is not found.
3332
Get(ctx context.Context, name string) (entity.Queue, error)

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

Lines changed: 38 additions & 48 deletions
Original file line numberDiff line numberDiff line change
@@ -53,74 +53,68 @@ func (s *QueueStoreContractSuite) queueDefaults() entity.Queue {
5353
return entity.Queue{Version: 1}
5454
}
5555

56-
// TestQueueStore_GetOrCreateCreates verifies GetOrCreate inserts a new row with zero-value runtime fields.
57-
func (s *QueueStoreContractSuite) TestQueueStore_GetOrCreateCreates() {
56+
// TestQueueStore_Create verifies Create inserts a new row with caller-supplied fields.
57+
func (s *QueueStoreContractSuite) TestQueueStore_Create() {
5858
t := s.T()
5959
const name = "contract/create"
6060

61-
got, err := s.queueStore.GetOrCreate(s.ctx, name, s.queueDefaults())
61+
require.NoError(t, s.queueStore.Create(s.ctx, entity.Queue{
62+
Name: name,
63+
Version: 1,
64+
}))
65+
66+
got, err := s.queueStore.Get(s.ctx, name)
6267
require.NoError(t, err)
6368
assert.Equal(t, entity.Queue{
6469
Name: name,
6570
LatestRequestSeq: 0,
6671
Version: 1,
6772
}, got)
6873

69-
s.log.Logf("GetOrCreateCreates passed: created queue %s", name)
74+
s.log.Logf("Create passed: created queue %s", name)
7075
}
7176

72-
// TestQueueStore_GetOrCreateWithDefaults verifies caller-supplied defaults are persisted on create.
73-
func (s *QueueStoreContractSuite) TestQueueStore_GetOrCreateWithDefaults() {
77+
// TestQueueStore_CreateWithFields verifies caller-supplied initial field values are persisted.
78+
func (s *QueueStoreContractSuite) TestQueueStore_CreateWithFields() {
7479
t := s.T()
7580
const name = "contract/defaults"
7681

77-
defaults := entity.Queue{
82+
toCreate := entity.Queue{
83+
Name: name,
7884
LastGreenURI: "git://remote/monorepo/main/green-bbbb",
7985
LatestRequestSeq: 99,
8086
Version: 1,
8187
}
88+
require.NoError(t, s.queueStore.Create(s.ctx, toCreate))
8289

83-
got, err := s.queueStore.GetOrCreate(s.ctx, name, defaults)
90+
got, err := s.queueStore.Get(s.ctx, name)
8491
require.NoError(t, err)
85-
assert.Equal(t, entity.Queue{
86-
Name: name,
87-
LastGreenURI: defaults.LastGreenURI,
88-
LatestRequestSeq: defaults.LatestRequestSeq,
89-
Version: 1,
90-
}, got)
92+
assert.Equal(t, toCreate, got)
9193

92-
s.log.Logf("GetOrCreateWithDefaults passed: persisted defaults for queue %s", name)
94+
s.log.Logf("CreateWithFields passed: persisted fields for queue %s", name)
9395
}
9496

95-
// TestQueueStore_GetOrCreateIdempotent verifies a second GetOrCreate returns the existing row unchanged.
96-
func (s *QueueStoreContractSuite) TestQueueStore_GetOrCreateIdempotent() {
97+
// TestQueueStore_CreateAlreadyExists verifies a duplicate Create returns ErrAlreadyExists.
98+
func (s *QueueStoreContractSuite) TestQueueStore_CreateAlreadyExists() {
9799
t := s.T()
98-
const name = "contract/idempotent"
100+
const name = "contract/already-exists"
99101

100-
first, err := s.queueStore.GetOrCreate(s.ctx, name, s.queueDefaults())
101-
require.NoError(t, err)
102+
first := entity.Queue{Name: name, LatestRequestSeq: 3, Version: 1}
103+
require.NoError(t, s.queueStore.Create(s.ctx, first))
102104

103-
second, err := s.queueStore.GetOrCreate(s.ctx, name, entity.Queue{
104-
LastGreenURI: "git://remote/monorepo/main/ignored-on-hit",
105+
err := s.queueStore.Create(s.ctx, entity.Queue{
106+
Name: name,
107+
LastGreenURI: "git://remote/monorepo/main/ignored-on-race",
105108
LatestRequestSeq: 500,
106109
Version: 1,
107110
})
108-
require.NoError(t, err)
109-
assert.Equal(t, first, second)
110-
111-
s.log.Logf("GetOrCreateIdempotent passed: queue %s", name)
112-
}
113-
114-
// TestQueueStore_GetOrCreateDefaultVersion verifies GetOrCreate writes version=1 when defaults.Version is zero.
115-
func (s *QueueStoreContractSuite) TestQueueStore_GetOrCreateDefaultVersion() {
116-
t := s.T()
117-
const name = "contract/default-version"
111+
assert.ErrorIs(t, err, storage.ErrAlreadyExists)
118112

119-
got, err := s.queueStore.GetOrCreate(s.ctx, name, entity.Queue{})
113+
got, err := s.queueStore.Get(s.ctx, name)
120114
require.NoError(t, err)
121-
assert.Equal(t, int32(1), got.Version)
115+
assert.Equal(t, first, got)
122116

123-
s.log.Logf("GetOrCreateDefaultVersion passed: queue %s", name)
117+
s.log.Logf("CreateAlreadyExists passed: queue %s", name)
124118
}
125119

126120
// TestQueueStore_GetNotFound verifies Get returns ErrNotFound for a missing queue.
@@ -138,8 +132,8 @@ func (s *QueueStoreContractSuite) TestQueueStore_UpdateCAS() {
138132
t := s.T()
139133
const name = "contract/update-cas"
140134

141-
created, err := s.queueStore.GetOrCreate(s.ctx, name, s.queueDefaults())
142-
require.NoError(t, err)
135+
created := entity.Queue{Name: name, Version: 1}
136+
require.NoError(t, s.queueStore.Create(s.ctx, created))
143137

144138
updated := created
145139
updated.LastGreenURI = "git://remote/monorepo/main/green-cccc"
@@ -175,17 +169,12 @@ func (s *QueueStoreContractSuite) TestQueueStore_UpdateSequentialCAS() {
175169
t := s.T()
176170
const name = "contract/sequential-cas"
177171

178-
created, err := s.queueStore.GetOrCreate(s.ctx, name, s.queueDefaults())
179-
require.NoError(t, err)
180-
require.Equal(t, int32(1), created.Version)
172+
require.NoError(t, s.queueStore.Create(s.ctx, entity.Queue{Name: name, Version: 1}))
181173

182-
v2 := created
183-
v2.LatestRequestSeq = 10
174+
v2 := entity.Queue{Name: name, LatestRequestSeq: 10, Version: 1}
184175
require.NoError(t, s.queueStore.Update(s.ctx, v2, 1, 2))
185176

186-
v3 := v2
187-
v3.Version = 2
188-
v3.InFlightCount = 1
177+
v3 := entity.Queue{Name: name, LatestRequestSeq: 10, InFlightCount: 1, Version: 2}
189178
require.NoError(t, s.queueStore.Update(s.ctx, v3, 2, 3))
190179

191180
got, err := s.queueStore.Get(s.ctx, name)
@@ -205,9 +194,10 @@ func (s *QueueStoreContractSuite) TestQueueStore_QueueIsolation() {
205194
nameB = "contract/isolation-b"
206195
)
207196

208-
_, err := s.queueStore.GetOrCreate(s.ctx, nameA, s.queueDefaults())
209-
require.NoError(t, err)
210-
baseline, err := s.queueStore.GetOrCreate(s.ctx, nameB, s.queueDefaults())
197+
require.NoError(t, s.queueStore.Create(s.ctx, entity.Queue{Name: nameA, Version: 1}))
198+
require.NoError(t, s.queueStore.Create(s.ctx, entity.Queue{Name: nameB, Version: 1}))
199+
200+
baseline, err := s.queueStore.Get(s.ctx, nameB)
211201
require.NoError(t, err)
212202

213203
updatedA := entity.Queue{

0 commit comments

Comments
 (0)