H-141: harden worker shutdown persistence (#24)
Co-authored-by: Rogee <rogee@ipao.vip>
This commit is contained in:
@@ -2,6 +2,7 @@ package dispatch
|
||||
|
||||
import (
|
||||
"context"
|
||||
"path/filepath"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
"time"
|
||||
@@ -28,7 +29,7 @@ func (l *dispatchWorkerListener) OnEvent(ctx context.Context, event *channel.Cha
|
||||
|
||||
func newDispatchWorkerDB(t *testing.T) *gorm.DB {
|
||||
t.Helper()
|
||||
db, err := gorm.Open(sqlite.Open("file:dispatch-worker?mode=memory&cache=shared"), &gorm.Config{Logger: logger.Default.LogMode(logger.Silent)})
|
||||
db, err := gorm.Open(sqlite.Open(filepath.Join(t.TempDir(), "dispatch-worker.db")), &gorm.Config{Logger: logger.Default.LogMode(logger.Silent)})
|
||||
if err != nil {
|
||||
t.Fatalf("open sqlite: %v", err)
|
||||
}
|
||||
@@ -41,7 +42,6 @@ func newDispatchWorkerDB(t *testing.T) *gorm.DB {
|
||||
t.Fatalf("migrate background jobs: %v", err)
|
||||
}
|
||||
t.Cleanup(func() {
|
||||
db.Exec("DELETE FROM background_jobs")
|
||||
sqlDB.Close()
|
||||
})
|
||||
return db
|
||||
|
||||
@@ -387,6 +387,9 @@ func (wp *WorkerPool) ProcessOne(ctx context.Context) (bool, error) {
|
||||
if wp.db == nil {
|
||||
return false, ErrWorkerDatabaseRequired
|
||||
}
|
||||
if err := ctx.Err(); err != nil {
|
||||
return false, err
|
||||
}
|
||||
job, err := wp.claimNext(ctx)
|
||||
if err != nil {
|
||||
if errors.Is(err, gorm.ErrRecordNotFound) {
|
||||
@@ -695,6 +698,7 @@ func (wp *WorkerPool) perform(ctx context.Context, job *model.BackgroundJob) err
|
||||
if err := handler(ctx, job); err != nil {
|
||||
return wp.fail(ctx, job, err)
|
||||
}
|
||||
ctx = context.WithoutCancel(ctx)
|
||||
finishedAt := wp.now()
|
||||
updates := map[string]any{
|
||||
"status": model.BackgroundJobStatusCompleted,
|
||||
@@ -707,6 +711,7 @@ func (wp *WorkerPool) perform(ctx context.Context, job *model.BackgroundJob) err
|
||||
}
|
||||
|
||||
func (wp *WorkerPool) fail(ctx context.Context, job *model.BackgroundJob, err error) error {
|
||||
ctx = context.WithoutCancel(ctx)
|
||||
now := wp.now()
|
||||
updates := map[string]any{
|
||||
"locked_at": nil,
|
||||
|
||||
@@ -4,6 +4,7 @@ import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"path/filepath"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
"time"
|
||||
@@ -18,7 +19,7 @@ import (
|
||||
|
||||
func newWorkerTestDB(t *testing.T) *gorm.DB {
|
||||
t.Helper()
|
||||
db, err := gorm.Open(sqlite.Open("file:worker-test?mode=memory&cache=shared"), &gorm.Config{Logger: logger.Default.LogMode(logger.Silent)})
|
||||
db, err := gorm.Open(sqlite.Open(filepath.Join(t.TempDir(), "worker.db")), &gorm.Config{Logger: logger.Default.LogMode(logger.Silent)})
|
||||
if err != nil {
|
||||
t.Fatalf("open sqlite: %v", err)
|
||||
}
|
||||
@@ -31,7 +32,6 @@ func newWorkerTestDB(t *testing.T) *gorm.DB {
|
||||
t.Fatalf("migrate background jobs: %v", err)
|
||||
}
|
||||
t.Cleanup(func() {
|
||||
db.Exec("DELETE FROM background_jobs")
|
||||
sqlDB.Close()
|
||||
})
|
||||
return db
|
||||
@@ -242,6 +242,36 @@ func TestWorkerPoolStartAndStopProcessJobs(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestWorkerPoolStopPersistsCancelledJobRetry(t *testing.T) {
|
||||
db := newWorkerTestDB(t)
|
||||
wp := NewWorkerPoolWithOptions(db, WithPollInterval(5*time.Millisecond), WithBackoff(func(int) time.Duration { return 0 }))
|
||||
started := make(chan struct{})
|
||||
wp.Register("cancelled", func(ctx context.Context, job *model.BackgroundJob) error {
|
||||
close(started)
|
||||
<-ctx.Done()
|
||||
return ctx.Err()
|
||||
})
|
||||
job, err := wp.Enqueue(context.Background(), "cancelled", nil)
|
||||
if err != nil {
|
||||
t.Fatalf("enqueue: %v", err)
|
||||
}
|
||||
if err := wp.Start(); err != nil {
|
||||
t.Fatalf("start worker: %v", err)
|
||||
}
|
||||
select {
|
||||
case <-started:
|
||||
case <-time.After(time.Second):
|
||||
t.Fatal("worker did not start job")
|
||||
}
|
||||
if err := wp.Stop(); err != nil {
|
||||
t.Fatalf("stop worker: %v", err)
|
||||
}
|
||||
reloaded := loadJob(t, db, job.ID)
|
||||
if reloaded.Status != model.BackgroundJobStatusRetrying || reloaded.LockedAt != nil || reloaded.LockedBy != "" || reloaded.LastError != context.Canceled.Error() {
|
||||
t.Fatalf("expected cancelled job to be retryable and unlocked: %+v", reloaded)
|
||||
}
|
||||
}
|
||||
|
||||
// --- Redis Stream path tests ---
|
||||
|
||||
func newMiniRedis(t *testing.T) (*miniredis.Miniredis, *redis.Client) {
|
||||
|
||||
Reference in New Issue
Block a user