Skip to content

Commit b023b73

Browse files
authored
Merge branch 'main' into remove-mysql-support
2 parents 8298207 + 5da9b37 commit b023b73

4 files changed

Lines changed: 53 additions & 10 deletions

File tree

‎service/datasetworker/datasetworker.go‎

Lines changed: 14 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -299,13 +299,20 @@ func (w *Thread) run(ctx context.Context) (retErr error) {
299299
}
300300

301301
w.stateMonitor.AddJob(job.ID, workCancel)
302-
switch job.Type {
303-
case model.Scan:
304-
err = w.scan(workCtx, *job.Attachment)
305-
case model.Pack:
306-
err = w.pack(workCtx, *job)
307-
case model.DagGen:
308-
err = w.ExportDag(workCtx, *job)
302+
// belt-and-suspenders: findJob filters attachment_id IS NOT NULL, but if a job slips through
303+
// without its Attachment preloaded (e.g. row was orphaned between claim and preload) we'd
304+
// otherwise nil-deref *job.Attachment below. Mark it errored and move on.
305+
if job.Attachment == nil {
306+
err = errors.Errorf("job %d has no attachment (orphaned by prep deletion)", job.ID)
307+
} else {
308+
switch job.Type {
309+
case model.Scan:
310+
err = w.scan(workCtx, *job.Attachment)
311+
case model.Pack:
312+
err = w.pack(workCtx, *job)
313+
case model.DagGen:
314+
err = w.ExportDag(workCtx, *job)
315+
}
309316
}
310317
w.stateMonitor.RemoveJob(job.ID)
311318
if workCtx.Err() != nil && ctx.Err() == nil {

‎service/datasetworker/find.go‎

Lines changed: 6 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -33,12 +33,16 @@ func (w *Thread) findJob(ctx context.Context, typesOrdered []model.JobType) (*mo
3333
for _, jobType := range typesOrdered {
3434
err := database.DoRetry(ctx, func() error {
3535
return db.Transaction(func(db *gorm.DB) error {
36-
// First, lock and claim the job without preloading (preload can interfere with locking)
36+
// First, lock and claim the job without preloading (preload can interfere with locking).
37+
// attachment_id IS NOT NULL skips orphans: jobs.attachment_id is ON DELETE SET NULL, so
38+
// deleting a prep leaves its jobs behind with state=ready and a null FK. Claiming one
39+
// would nil-deref *job.Attachment in the run loop. The healthcheck sweep deletes these
40+
// eventually but it races worker pickup -- filter here too.
3741
err := db.Clauses(clause.Locking{
3842
Strength: "UPDATE",
3943
Options: "SKIP LOCKED",
4044
}).
41-
Where("type = ? AND (state = ? OR (state = ? AND worker_id IS NULL))", jobType, model.Ready, model.Processing).
45+
Where("type = ? AND attachment_id IS NOT NULL AND (state = ? OR (state = ? AND worker_id IS NULL))", jobType, model.Ready, model.Processing).
4246
First(&job).Error
4347
if err != nil {
4448
if errors.Is(err, gorm.ErrRecordNotFound) {

‎service/datasetworker/find_test.go‎

Lines changed: 32 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -13,6 +13,38 @@ import (
1313
"gorm.io/gorm"
1414
)
1515

16+
// A job whose preparation was deleted has attachment_id NULL (SET NULL cascade)
17+
// and stays Ready until the healthcheck reaper sweeps it. findJob must skip these
18+
// rather than claim them and nil-deref *job.Attachment downstream.
19+
func TestFindWorkSkipsOrphanedJob(t *testing.T) {
20+
testutil.All(t, func(ctx context.Context, t *testing.T, db *gorm.DB) {
21+
thread := &Thread{
22+
dbNoContext: db,
23+
config: Config{EnableScan: true, EnablePack: true, EnableDag: true},
24+
logger: logger.With("test", true),
25+
id: uuid.New(),
26+
}
27+
28+
_, err := healthcheck.Register(ctx, thread.dbNoContext, thread.id, model.DatasetWorker, true)
29+
require.NoError(t, err)
30+
31+
for _, jt := range []model.JobType{model.Scan, model.Pack, model.DagGen} {
32+
err = db.Create(&model.Job{
33+
AttachmentID: nil,
34+
State: model.Ready,
35+
Type: jt,
36+
}).Error
37+
require.NoError(t, err)
38+
}
39+
40+
for _, jt := range []model.JobType{model.Scan, model.Pack, model.DagGen} {
41+
found, err := thread.findJob(ctx, []model.JobType{jt})
42+
require.NoError(t, err)
43+
require.Nil(t, found, "orphaned %s job must not be claimed", jt)
44+
}
45+
})
46+
}
47+
1648
func TestFindPackWork(t *testing.T) {
1749
testutil.All(t, func(ctx context.Context, t *testing.T, db *gorm.DB) {
1850
thread := &Thread{

‎version.json‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,3 +1,3 @@
11
{
2-
"version": "v0.7.0-RC1"
2+
"version": "v1.0.0-RC1"
33
}

0 commit comments

Comments
 (0)