Skip to content

Commit 5b20a3e

Browse files
authored
Update job and worker delete and stale update functionality (#9)
* add init functions for job and worker * fix commit sql * update cli with formatted examples, add manager commands * add cluster example * add manager server * update go db handlers to new sql functions * fix gosec issues * fix formatting * synchronize with sql * update sql to unified functions * synchronize sql submodule * try adding fetch depth for tagging issue * sync sql * update sql, update dependencies * update dependencies * add replace statement for sql, update imports * update delete job to delete from archive, update stale job update to status queued * add delete stale workers functionality, add defaults to master settings * add tests for new functions
1 parent 0d17fe6 commit 5b20a3e

15 files changed

Lines changed: 535 additions & 61 deletions

database/dbJob.go

Lines changed: 19 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -381,15 +381,14 @@ func (r JobDBHandler) UpdateJobFinal(job *model.Job) (*model.Job, error) {
381381
return archivedJob, nil
382382
}
383383

384-
// UpdateStaleJobs updates all jobs to CANCELLED status where the assigned worker is STOPPED
385-
// based on the provided threshold. It returns the number of jobs that were updated.
386-
// Jobs are considered stale if their assigned worker has STOPPED status and the worker's
387-
// updated_at timestamp is older than the threshold.
384+
// UpdateStaleJobs updates all jobs to QUEUED status where the assigned worker is STOPPED
385+
// so they can be picked up by available workers again. It returns the number of jobs that were updated.
386+
// Jobs are considered stale if their assigned worker has STOPPED status.
388387
func (r JobDBHandler) UpdateStaleJobs() (int, error) {
389388
var affectedRows int
390389
err := r.db.Instance.QueryRow(
391390
`SELECT update_stale_jobs($1, $2, $3, $4, $5)`,
392-
model.JobStatusCancelled,
391+
model.JobStatusQueued,
393392
model.JobStatusSucceeded,
394393
model.JobStatusCancelled,
395394
model.JobStatusFailed,
@@ -402,19 +401,6 @@ func (r JobDBHandler) UpdateStaleJobs() (int, error) {
402401
return affectedRows, nil
403402
}
404403

405-
// DeleteJob deletes a job record from the database based on its RID.
406-
func (r JobDBHandler) DeleteJob(rid uuid.UUID) error {
407-
_, err := r.db.Instance.Exec(
408-
`SELECT delete_job($1::UUID)`,
409-
rid,
410-
)
411-
if err != nil {
412-
return helper.NewError("exec", err)
413-
}
414-
415-
return nil
416-
}
417-
418404
// SelectJob retrieves a single job record from the database based on its RID.
419405
func (r JobDBHandler) SelectJob(rid uuid.UUID) (*model.Job, error) {
420406
row := r.db.Instance.QueryRow(
@@ -635,6 +621,21 @@ func (r JobDBHandler) RemoveRetentionArchive() error {
635621
return nil
636622
}
637623

624+
// DeleteJob deletes a job record from the job archive based on its RID.
625+
// We only delete jobs from the archive as queued and running jobs should be cancelled first.
626+
// Cancelling a job will move it to the archive with CANCELLED status.
627+
func (r JobDBHandler) DeleteJob(rid uuid.UUID) error {
628+
_, err := r.db.Instance.Exec(
629+
`SELECT delete_job($1::UUID)`,
630+
rid,
631+
)
632+
if err != nil {
633+
return helper.NewError("exec", err)
634+
}
635+
636+
return nil
637+
}
638+
638639
// SelectJobFromArchive retrieves a single archived job record from the database based on its RID.
639640
func (r JobDBHandler) SelectJobFromArchive(rid uuid.UUID) (*model.Job, error) {
640641
row := r.db.Instance.QueryRow(

database/dbJob_test.go

Lines changed: 20 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -346,7 +346,7 @@ func TestUpdateStaleJobs(t *testing.T) {
346346
workerDbHandler, err := NewWorkerDBHandler(database, true)
347347
assert.NoError(t, err, "Expected NewWorkerDBHandler to not return an error")
348348

349-
t.Run("Cancel jobs with stopped workers", func(t *testing.T) {
349+
t.Run("Requeue jobs with stopped workers", func(t *testing.T) {
350350
// Create workers and set one to STOPPED
351351
worker1, err := model.NewWorker("worker-1", 3)
352352
require.NoError(t, err)
@@ -364,13 +364,13 @@ func TestUpdateStaleJobs(t *testing.T) {
364364

365365
// Create jobs: QUEUED (stopped worker), SUCCEEDED (stopped worker), QUEUED (ready worker)
366366
testCases := []struct {
367-
status string
368-
workerRID string
369-
shouldCancel bool
367+
status string
368+
workerRID string
369+
shouldRequeue bool
370370
}{
371-
{model.JobStatusQueued, "stopped", true}, // Should be cancelled
372-
{model.JobStatusSucceeded, "stopped", false}, // Should not be cancelled (final status)
373-
{model.JobStatusQueued, "ready", false}, // Should not be cancelled (ready worker)
371+
{model.JobStatusQueued, "stopped", true}, // Should be requeued
372+
{model.JobStatusSucceeded, "stopped", false}, // Should not be requeued (final status)
373+
{model.JobStatusQueued, "ready", false}, // Should not be requeued (ready worker)
374374
}
375375

376376
jobs := make([]*model.Job, len(testCases))
@@ -399,14 +399,14 @@ func TestUpdateStaleJobs(t *testing.T) {
399399
// Test UpdateStaleJobs
400400
updatedCount, err := jobDbHandler.UpdateStaleJobs()
401401
assert.NoError(t, err)
402-
assert.Equal(t, 1, updatedCount, "Expected 1 job to be cancelled")
402+
assert.Equal(t, 1, updatedCount, "Expected 1 job to be requeued")
403403

404404
// Verify results
405405
for i, tc := range testCases {
406406
updatedJob, err := jobDbHandler.SelectJob(jobs[i].RID)
407407
require.NoError(t, err)
408-
if tc.shouldCancel {
409-
assert.Equal(t, model.JobStatusCancelled, updatedJob.Status)
408+
if tc.shouldRequeue {
409+
assert.Equal(t, model.JobStatusQueued, updatedJob.Status)
410410
} else {
411411
assert.Equal(t, tc.status, updatedJob.Status)
412412
}
@@ -459,12 +459,19 @@ func TestJobDeleteJob(t *testing.T) {
459459
insertedJob, err := jobDbHandler.InsertJob(job)
460460
require.NoError(t, err, "Expected InsertJob to not return an error")
461461

462+
// First, archive the job by marking it as completed (this moves it from job to job_archive)
463+
insertedJob.Status = model.JobStatusSucceeded
464+
archivedJob, err := jobDbHandler.UpdateJobFinal(insertedJob)
465+
require.NoError(t, err, "Expected UpdateJobFinal to not return an error")
466+
require.NotNil(t, archivedJob, "Expected archived job to not be nil")
467+
468+
// Now delete it from the archive
462469
err = jobDbHandler.DeleteJob(insertedJob.RID)
463470
assert.NoError(t, err, "Expected DeleteJob to not return an error")
464471

465-
// Verify that the job no longer exists
466-
deletedJob, err := jobDbHandler.SelectJob(insertedJob.RID)
467-
require.Error(t, err, "Expected SelectJob to return an error")
472+
// Verify that the job no longer exists in the archive
473+
deletedJob, err := jobDbHandler.SelectJobFromArchive(insertedJob.RID)
474+
require.Error(t, err, "Expected SelectJobFromArchive to return an error")
468475
assert.Contains(t, err.Error(), sql.ErrNoRows.Error(), "Expected error to contain sql.ErrNoRows for deleted job")
469476
assert.Nil(t, deletedJob, "Expected deleted job to be nil")
470477
}

database/dbMaster_test.go

Lines changed: 6 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -101,8 +101,8 @@ func TestMasterUpdateMaster(t *testing.T) {
101101
RID: uuid.New(),
102102
}
103103
settings := &model.MasterSettings{
104-
RetentionArchive: 30,
105-
MasterLockTimeout: 10 * time.Minute,
104+
JobDeleteThreshold: 30,
105+
MasterLockTimeout: 10 * time.Minute,
106106
}
107107
masterOld, err := workerDbHandler.UpdateMaster(worker1, settings)
108108
assert.NoError(t, err, "Expected no error updating master")
@@ -117,7 +117,7 @@ func TestMasterUpdateMaster(t *testing.T) {
117117
assert.NotNil(t, master, "Expected master to not be nil")
118118
assert.Equal(t, worker1.ID, master.WorkerID, "Expected master worker ID to match worker ID")
119119
assert.Equal(t, worker1.RID, master.WorkerRID, "Expected master worker RID to match worker RID")
120-
assert.Equal(t, settings.RetentionArchive, master.Settings.RetentionArchive, "Expected master retention archive to match")
120+
assert.Equal(t, settings.JobDeleteThreshold, master.Settings.JobDeleteThreshold, "Expected master retention archive to match")
121121

122122
worker2 := &model.Worker{
123123
ID: 2,
@@ -145,8 +145,8 @@ func TestMasterSelectMaster(t *testing.T) {
145145
RID: uuid.New(),
146146
}
147147
settings := &model.MasterSettings{
148-
RetentionArchive: 30,
149-
MasterLockTimeout: 10 * time.Minute,
148+
JobDeleteThreshold: 30,
149+
MasterLockTimeout: 10 * time.Minute,
150150
}
151151
master, err := workerDbHandler.UpdateMaster(worker, settings)
152152
require.NoError(t, err, "Expected no error updating master")
@@ -158,5 +158,5 @@ func TestMasterSelectMaster(t *testing.T) {
158158
assert.Equal(t, 1, master.ID, "Expected master ID to be 1")
159159
assert.Equal(t, worker.ID, master.WorkerID, "Expected master worker ID to be 1")
160160
assert.Equal(t, worker.RID, master.WorkerRID, "Expected master worker RID to not equal to worker RID")
161-
assert.Equal(t, settings.RetentionArchive, master.Settings.RetentionArchive, "Expected master retention archive to be 30")
161+
assert.Equal(t, settings.JobDeleteThreshold, master.Settings.JobDeleteThreshold, "Expected master retention archive to be 30")
162162
}

database/dbWorker.go

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,7 @@ type WorkerDBHandlerFunctions interface {
2222
UpdateWorker(worker *model.Worker) (*model.Worker, error)
2323
UpdateStaleWorkers(staleThreshold time.Duration) (int, error)
2424
DeleteWorker(rid uuid.UUID) error
25+
DeleteStaleWorkers(deleteThreshold time.Duration) (int, error)
2526
SelectWorker(rid uuid.UUID) (*model.Worker, error)
2627
SelectAllWorkers(lastID int, entries int) ([]*model.Worker, error)
2728
SelectAllWorkersBySearch(search string, lastID int, entries int) ([]*model.Worker, error)
@@ -235,6 +236,23 @@ func (r WorkerDBHandler) DeleteWorker(rid uuid.UUID) error {
235236
return nil
236237
}
237238

239+
// DeleteStaleWorkers deletes workers that have been in STOPPED status for longer than the deleteThreshold.
240+
// It returns the number of workers that were deleted.
241+
func (r WorkerDBHandler) DeleteStaleWorkers(deleteThreshold time.Duration) (int, error) {
242+
cutoffTime := time.Now().UTC().Add(-deleteThreshold)
243+
244+
var rowsAffected int
245+
err := r.db.Instance.QueryRow(
246+
`SELECT delete_stale_workers($1)`,
247+
cutoffTime,
248+
).Scan(&rowsAffected)
249+
if err != nil {
250+
return 0, helper.NewError("delete stale workers", err)
251+
}
252+
253+
return rowsAffected, nil
254+
}
255+
238256
// SelectWorker retrieves a single worker record from the database based on its RID.
239257
// It returns the worker record.
240258
// If the worker is not found or an error occurs during the query, it returns an error.

database/dbWorker_test.go

Lines changed: 185 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -251,6 +251,191 @@ func TestWorkerDeleteWorker(t *testing.T) {
251251
assert.Nil(t, deletedWorker, "Expected deleted worker to be nil")
252252
}
253253

254+
func TestDeleteStaleWorkers(t *testing.T) {
255+
helper.SetTestDatabaseConfigEnvs(t, dbPort)
256+
dbConfig, err := helper.NewDatabaseConfiguration()
257+
if err != nil {
258+
t.Fatalf("failed to create database configuration: %v", err)
259+
}
260+
database := helper.NewTestDatabase(dbConfig)
261+
262+
workerDbHandler, err := NewWorkerDBHandler(database, true)
263+
assert.NoError(t, err, "Expected NewWorkerDBHandler to not return an error")
264+
265+
t.Run("Delete stale workers with different statuses", func(t *testing.T) {
266+
// Create workers: STOPPED (stale), STOPPED (fresh), READY (stale), RUNNING (stale)
267+
testCases := []struct {
268+
status string
269+
makeStale bool
270+
shouldDelete bool
271+
}{
272+
{model.WorkerStatusStopped, true, true}, // Should be deleted
273+
{model.WorkerStatusStopped, false, false}, // Should not be deleted (fresh)
274+
{model.WorkerStatusReady, true, false}, // Should not be deleted (not STOPPED)
275+
{model.WorkerStatusRunning, true, false}, // Should not be deleted (not STOPPED)
276+
}
277+
278+
workers := make([]*model.Worker, len(testCases))
279+
280+
for i, tc := range testCases {
281+
testWorker, err := model.NewWorker(fmt.Sprintf("worker-%d", i), 3)
282+
require.NoError(t, err, "Expected to create test worker %d", i)
283+
284+
insertedWorker, err := workerDbHandler.InsertWorker(testWorker)
285+
require.NoError(t, err, "Expected to insert worker %d", i)
286+
287+
insertedWorker.Status = tc.status
288+
updatedWorker, err := workerDbHandler.UpdateWorker(insertedWorker)
289+
require.NoError(t, err, "Expected to update worker %d status", i)
290+
workers[i] = updatedWorker
291+
}
292+
293+
// Make some workers stale (older than 1 hour)
294+
staleTime := time.Now().UTC().Add(-1 * time.Hour)
295+
for i, tc := range testCases {
296+
if tc.makeStale {
297+
_, err = database.Instance.Exec(
298+
"UPDATE worker SET updated_at = $1 WHERE rid = $2",
299+
staleTime, workers[i].RID,
300+
)
301+
require.NoError(t, err, "Expected to make worker %d stale", i)
302+
}
303+
}
304+
305+
// Test DeleteStaleWorkers - should delete only stale STOPPED workers
306+
deleteThreshold := 10 * time.Minute
307+
deletedCount, err := workerDbHandler.DeleteStaleWorkers(deleteThreshold)
308+
assert.NoError(t, err, "Expected DeleteStaleWorkers to complete successfully")
309+
assert.Equal(t, 1, deletedCount, "Expected 1 worker to be deleted (stale STOPPED)")
310+
311+
// Verify only the stale STOPPED worker was deleted
312+
for i, tc := range testCases {
313+
worker, err := workerDbHandler.SelectWorker(workers[i].RID)
314+
315+
if tc.shouldDelete {
316+
assert.Error(t, err, "Expected worker %d to be deleted", i)
317+
assert.Nil(t, worker, "Expected deleted worker %d to be nil", i)
318+
} else {
319+
assert.NoError(t, err, "Expected worker %d to still exist", i)
320+
assert.NotNil(t, worker, "Expected worker %d to not be nil", i)
321+
assert.Equal(t, tc.status, worker.Status, "Expected worker %d status to remain unchanged", i)
322+
}
323+
}
324+
325+
// Clean up remaining workers
326+
for i, tc := range testCases {
327+
if !tc.shouldDelete {
328+
err = workerDbHandler.DeleteWorker(workers[i].RID)
329+
assert.NoError(t, err, "Expected to delete remaining worker %d", i)
330+
}
331+
}
332+
})
333+
334+
t.Run("No stale workers to delete", func(t *testing.T) {
335+
// Create fresh workers with different statuses
336+
statuses := []string{model.WorkerStatusReady, model.WorkerStatusRunning, model.WorkerStatusStopped}
337+
workers := make([]*model.Worker, len(statuses))
338+
339+
for i, status := range statuses {
340+
testWorker, err := model.NewWorker(fmt.Sprintf("fresh-worker-%d", i), 3)
341+
require.NoError(t, err, "Expected to create fresh worker %d", i)
342+
343+
insertedWorker, err := workerDbHandler.InsertWorker(testWorker)
344+
require.NoError(t, err, "Expected to insert fresh worker %d", i)
345+
346+
insertedWorker.Status = status
347+
updatedWorker, err := workerDbHandler.UpdateWorker(insertedWorker)
348+
require.NoError(t, err, "Expected to update fresh worker %d status", i)
349+
workers[i] = updatedWorker
350+
}
351+
352+
// Test with short threshold - no workers should be deleted
353+
deletedCount, err := workerDbHandler.DeleteStaleWorkers(10 * time.Second)
354+
assert.NoError(t, err, "Expected DeleteStaleWorkers to complete successfully")
355+
assert.Equal(t, 0, deletedCount, "Expected no workers to be deleted")
356+
357+
// Verify all workers still exist
358+
for i, worker := range workers {
359+
existingWorker, err := workerDbHandler.SelectWorker(worker.RID)
360+
assert.NoError(t, err, "Expected fresh worker %d to still exist", i)
361+
assert.NotNil(t, existingWorker, "Expected fresh worker %d to not be nil", i)
362+
}
363+
364+
// Clean up
365+
for i, worker := range workers {
366+
err = workerDbHandler.DeleteWorker(worker.RID)
367+
assert.NoError(t, err, "Expected to delete fresh worker %d", i)
368+
}
369+
})
370+
371+
t.Run("Delete only STOPPED workers older than threshold", func(t *testing.T) {
372+
// Create multiple STOPPED workers with different timestamps
373+
testCases := []struct {
374+
name string
375+
ageMinutes int
376+
shouldDelete bool
377+
}{
378+
{"very-old-stopped", 120, true}, // 2 hours old, should be deleted
379+
{"old-stopped", 30, true}, // 30 minutes old, should be deleted
380+
{"recent-stopped", 5, false}, // 5 minutes old, should not be deleted
381+
{"fresh-stopped", 1, false}, // 1 minute old, should not be deleted
382+
}
383+
384+
workers := make([]*model.Worker, len(testCases))
385+
386+
for i, tc := range testCases {
387+
testWorker, err := model.NewWorker(tc.name, 3)
388+
require.NoError(t, err, "Expected to create worker %s", tc.name)
389+
390+
insertedWorker, err := workerDbHandler.InsertWorker(testWorker)
391+
require.NoError(t, err, "Expected to insert worker %s", tc.name)
392+
393+
// Set status to STOPPED
394+
insertedWorker.Status = model.WorkerStatusStopped
395+
updatedWorker, err := workerDbHandler.UpdateWorker(insertedWorker)
396+
require.NoError(t, err, "Expected to update worker %s status", tc.name)
397+
398+
// Set custom timestamp
399+
timestamp := time.Now().UTC().Add(-time.Duration(tc.ageMinutes) * time.Minute)
400+
_, err = database.Instance.Exec(
401+
"UPDATE worker SET updated_at = $1 WHERE rid = $2",
402+
timestamp, updatedWorker.RID,
403+
)
404+
require.NoError(t, err, "Expected to set custom timestamp for worker %s", tc.name)
405+
406+
workers[i] = updatedWorker
407+
}
408+
409+
// Test with 15 minute threshold
410+
deleteThreshold := 15 * time.Minute
411+
deletedCount, err := workerDbHandler.DeleteStaleWorkers(deleteThreshold)
412+
assert.NoError(t, err, "Expected DeleteStaleWorkers to complete successfully")
413+
assert.Equal(t, 2, deletedCount, "Expected 2 workers to be deleted (older than 15 minutes)")
414+
415+
// Verify results
416+
for i, tc := range testCases {
417+
worker, err := workerDbHandler.SelectWorker(workers[i].RID)
418+
419+
if tc.shouldDelete {
420+
assert.Error(t, err, "Expected worker %s to be deleted", tc.name)
421+
assert.Nil(t, worker, "Expected deleted worker %s to be nil", tc.name)
422+
} else {
423+
assert.NoError(t, err, "Expected worker %s to still exist", tc.name)
424+
assert.NotNil(t, worker, "Expected worker %s to not be nil", tc.name)
425+
assert.Equal(t, model.WorkerStatusStopped, worker.Status, "Expected worker %s to remain STOPPED", tc.name)
426+
}
427+
}
428+
429+
// Clean up remaining workers
430+
for i, tc := range testCases {
431+
if !tc.shouldDelete {
432+
err = workerDbHandler.DeleteWorker(workers[i].RID)
433+
assert.NoError(t, err, "Expected to delete remaining worker %s", tc.name)
434+
}
435+
}
436+
})
437+
}
438+
254439
func TestWorkerSelectWorker(t *testing.T) {
255440
helper.SetTestDatabaseConfigEnvs(t, dbPort)
256441
dbConfig, err := helper.NewDatabaseConfiguration()

0 commit comments

Comments
 (0)