* fix(HH-565): complete assignment parity regressions * fix(HH-565): scope capacity exclusions per agent --------- Co-authored-by: Rogee <rogee@ipao.vip>
795 lines
27 KiB
Go
795 lines
27 KiB
Go
package worker
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"errors"
|
|
"path/filepath"
|
|
"sync"
|
|
"sync/atomic"
|
|
"testing"
|
|
"time"
|
|
|
|
"github.com/alicebob/miniredis/v2"
|
|
"github.com/gochat/gochat/internal/model"
|
|
"github.com/redis/go-redis/v9"
|
|
"github.com/stretchr/testify/require"
|
|
"gorm.io/driver/sqlite"
|
|
"gorm.io/gorm"
|
|
"gorm.io/gorm/logger"
|
|
)
|
|
|
|
func newWorkerTestDB(t *testing.T) *gorm.DB {
|
|
t.Helper()
|
|
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)
|
|
}
|
|
sqlDB, err := db.DB()
|
|
if err != nil {
|
|
t.Fatalf("sqlite db handle: %v", err)
|
|
}
|
|
sqlDB.SetMaxOpenConns(1)
|
|
if err := db.AutoMigrate(&model.BackgroundJob{}); err != nil {
|
|
t.Fatalf("migrate background jobs: %v", err)
|
|
}
|
|
t.Cleanup(func() {
|
|
sqlDB.Close()
|
|
})
|
|
return db
|
|
}
|
|
|
|
func loadJob(t *testing.T, db *gorm.DB, id uint) model.BackgroundJob {
|
|
t.Helper()
|
|
var job model.BackgroundJob
|
|
if err := db.First(&job, id).Error; err != nil {
|
|
t.Fatalf("load job: %v", err)
|
|
}
|
|
return job
|
|
}
|
|
|
|
func TestWorkerPoolEnqueueStoresPayloadAndIdempotency(t *testing.T) {
|
|
db := newWorkerTestDB(t)
|
|
now := time.Date(2026, 6, 5, 10, 0, 0, 0, time.UTC)
|
|
wp := NewWorkerPoolWithOptions(db, WithNow(func() time.Time { return now }))
|
|
|
|
job, err := wp.Enqueue(context.Background(), "contact_export", map[string]any{"account_id": 1}, WithIdempotencyKey("contact-export:1"), WithQueue("exports"), WithMaxAttempts(5), WithPriority(10))
|
|
if err != nil {
|
|
t.Fatalf("enqueue: %v", err)
|
|
}
|
|
duplicate, err := wp.Enqueue(context.Background(), "contact_export", map[string]any{"account_id": 2}, WithIdempotencyKey("contact-export:1"))
|
|
if err != nil {
|
|
t.Fatalf("enqueue duplicate: %v", err)
|
|
}
|
|
if duplicate.ID != job.ID {
|
|
t.Fatalf("expected duplicate enqueue to return existing job %d, got %d", job.ID, duplicate.ID)
|
|
}
|
|
|
|
reloaded := loadJob(t, db, job.ID)
|
|
if reloaded.Queue != "exports" || reloaded.JobType != "contact_export" || reloaded.Status != model.BackgroundJobStatusQueued || reloaded.MaxAttempts != 5 || reloaded.Priority != 10 {
|
|
t.Fatalf("unexpected job fields: %+v", reloaded)
|
|
}
|
|
var payload map[string]any
|
|
if err := json.Unmarshal(reloaded.Payload, &payload); err != nil {
|
|
t.Fatalf("unmarshal payload: %v", err)
|
|
}
|
|
if payload["account_id"].(float64) != 1 {
|
|
t.Fatalf("unexpected payload: %s", string(reloaded.Payload))
|
|
}
|
|
}
|
|
|
|
func TestWorkerPoolProcessOneCompletesDueJob(t *testing.T) {
|
|
db := newWorkerTestDB(t)
|
|
now := time.Date(2026, 6, 5, 10, 0, 0, 0, time.UTC)
|
|
wp := NewWorkerPoolWithOptions(db, WithNow(func() time.Time { return now }), WithWorkerID("test-worker"))
|
|
var handled atomic.Int32
|
|
wp.Register("send_reply", func(ctx context.Context, job *model.BackgroundJob) error {
|
|
handled.Add(1)
|
|
if job.Attempts != 1 || job.LockedBy != "test-worker" || job.LockedAt == nil {
|
|
t.Fatalf("job was not claimed before handler: %+v", job)
|
|
}
|
|
return nil
|
|
})
|
|
job, err := wp.Enqueue(context.Background(), "send_reply", map[string]any{"message_id": 7})
|
|
if err != nil {
|
|
t.Fatalf("enqueue: %v", err)
|
|
}
|
|
|
|
processed, err := wp.ProcessOne(context.Background())
|
|
if err != nil {
|
|
t.Fatalf("process one: %v", err)
|
|
}
|
|
if !processed || handled.Load() != 1 {
|
|
t.Fatalf("expected one handled job, processed=%v handled=%d", processed, handled.Load())
|
|
}
|
|
reloaded := loadJob(t, db, job.ID)
|
|
if reloaded.Status != model.BackgroundJobStatusCompleted || reloaded.FinishedAt == nil || reloaded.LockedAt != nil || reloaded.LockedBy != "" {
|
|
t.Fatalf("expected completed unlocked job: %+v", reloaded)
|
|
}
|
|
}
|
|
|
|
func TestWorkerPoolRetriesThenDeadLettersFailures(t *testing.T) {
|
|
db := newWorkerTestDB(t)
|
|
now := time.Date(2026, 6, 5, 10, 0, 0, 0, time.UTC)
|
|
wp := NewWorkerPoolWithOptions(db, WithNow(func() time.Time { return now }), WithBackoff(func(attempt int) time.Duration { return 0 }))
|
|
boom := errors.New("provider timeout")
|
|
wp.Register("webhook_delivery", func(ctx context.Context, job *model.BackgroundJob) error { return boom })
|
|
job, err := wp.Enqueue(context.Background(), "webhook_delivery", nil, WithMaxAttempts(2))
|
|
if err != nil {
|
|
t.Fatalf("enqueue: %v", err)
|
|
}
|
|
|
|
processed, err := wp.ProcessOne(context.Background())
|
|
if !processed || !errors.Is(err, boom) {
|
|
t.Fatalf("expected first failure, processed=%v err=%v", processed, err)
|
|
}
|
|
reloaded := loadJob(t, db, job.ID)
|
|
if reloaded.Status != model.BackgroundJobStatusRetrying || reloaded.Attempts != 1 || reloaded.LastError != boom.Error() || reloaded.FailedAt != nil {
|
|
t.Fatalf("expected retrying job after first failure: %+v", reloaded)
|
|
}
|
|
|
|
processed, err = wp.ProcessOne(context.Background())
|
|
if !processed || !errors.Is(err, boom) {
|
|
t.Fatalf("expected second failure, processed=%v err=%v", processed, err)
|
|
}
|
|
reloaded = loadJob(t, db, job.ID)
|
|
if reloaded.Status != model.BackgroundJobStatusDead || reloaded.Attempts != 2 || reloaded.FailedAt == nil || reloaded.LockedAt != nil {
|
|
t.Fatalf("expected dead-lettered job: %+v", reloaded)
|
|
}
|
|
}
|
|
|
|
func TestWorkerPoolPermanentErrorDeadLettersImmediately(t *testing.T) {
|
|
db := newWorkerTestDB(t)
|
|
wp := NewWorkerPoolWithOptions(db, WithBackoff(func(attempt int) time.Duration { return 0 }))
|
|
wp.Register("permanent", func(context.Context, *model.BackgroundJob) error {
|
|
return Permanent(errors.New("invalid webhook payload"))
|
|
})
|
|
job, err := wp.Enqueue(context.Background(), "permanent", nil, WithMaxAttempts(10))
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
processed, err := wp.ProcessOne(context.Background())
|
|
if !processed || err == nil {
|
|
t.Fatalf("processed=%v err=%v", processed, err)
|
|
}
|
|
reloaded := loadJob(t, db, job.ID)
|
|
if reloaded.Status != model.BackgroundJobStatusDead || reloaded.Attempts != 1 {
|
|
t.Fatalf("expected immediate dead letter, got %+v", reloaded)
|
|
}
|
|
}
|
|
|
|
func TestWorkerPoolRespectsScheduleAndQueues(t *testing.T) {
|
|
db := newWorkerTestDB(t)
|
|
now := time.Date(2026, 6, 5, 10, 0, 0, 0, time.UTC)
|
|
wp := NewWorkerPoolWithOptions(db, WithNow(func() time.Time { return now }), WithQueues("critical"))
|
|
wp.Register("sla_scan", func(ctx context.Context, job *model.BackgroundJob) error { return nil })
|
|
if _, err := wp.Enqueue(context.Background(), "sla_scan", nil, WithQueue("default")); err != nil {
|
|
t.Fatalf("enqueue default: %v", err)
|
|
}
|
|
if _, err := wp.Enqueue(context.Background(), "sla_scan", nil, WithQueue("critical"), WithScheduledAt(now.Add(time.Hour))); err != nil {
|
|
t.Fatalf("enqueue future: %v", err)
|
|
}
|
|
processed, err := wp.ProcessOne(context.Background())
|
|
if err != nil || processed {
|
|
t.Fatalf("expected no eligible job, processed=%v err=%v", processed, err)
|
|
}
|
|
}
|
|
|
|
func TestWorkerPoolRequeuesStaleRunningJobs(t *testing.T) {
|
|
db := newWorkerTestDB(t)
|
|
now := time.Date(2026, 6, 5, 10, 0, 0, 0, time.UTC)
|
|
wp := NewWorkerPoolWithOptions(db, WithNow(func() time.Time { return now }), WithStaleLockTimeout(time.Minute))
|
|
lockedAt := now.Add(-2 * time.Minute)
|
|
job := model.BackgroundJob{
|
|
Queue: model.DefaultBackgroundJobQueue,
|
|
JobType: "captain_document_sync",
|
|
Payload: json.RawMessage(`{}`),
|
|
Status: model.BackgroundJobStatusRunning,
|
|
MaxAttempts: 3,
|
|
ScheduledAt: now.Add(-time.Hour),
|
|
LockedAt: &lockedAt,
|
|
LockedBy: "dead-worker",
|
|
}
|
|
if err := db.Create(&job).Error; err != nil {
|
|
t.Fatalf("create stale job: %v", err)
|
|
}
|
|
|
|
count, err := wp.RequeueStaleJobs(context.Background())
|
|
if err != nil {
|
|
t.Fatalf("requeue stale: %v", err)
|
|
}
|
|
if count != 1 {
|
|
t.Fatalf("expected 1 stale job requeued, got %d", count)
|
|
}
|
|
reloaded := loadJob(t, db, job.ID)
|
|
if reloaded.Status != model.BackgroundJobStatusRetrying || reloaded.LockedAt != nil || reloaded.LockedBy != "" {
|
|
t.Fatalf("expected retrying unlocked stale job: %+v", reloaded)
|
|
}
|
|
}
|
|
|
|
func TestWorkerPoolStartAndStopProcessJobs(t *testing.T) {
|
|
db := newWorkerTestDB(t)
|
|
now := time.Date(2026, 6, 5, 10, 0, 0, 0, time.UTC)
|
|
wp := NewWorkerPoolWithOptions(db, WithNow(func() time.Time { return now }), WithPollInterval(5*time.Millisecond))
|
|
var handled atomic.Int32
|
|
wp.Register("event_dispatch", func(ctx context.Context, job *model.BackgroundJob) error {
|
|
handled.Add(1)
|
|
return nil
|
|
})
|
|
job, err := wp.Enqueue(context.Background(), "event_dispatch", map[string]any{"event": "conversation_created"})
|
|
if err != nil {
|
|
t.Fatalf("enqueue: %v", err)
|
|
}
|
|
if err := wp.Start(); err != nil {
|
|
t.Fatalf("start worker: %v", err)
|
|
}
|
|
t.Cleanup(func() {
|
|
if err := wp.Stop(); err != nil {
|
|
t.Errorf("stop worker: %v", err)
|
|
}
|
|
})
|
|
|
|
deadline := time.Now().Add(time.Second)
|
|
for time.Now().Before(deadline) {
|
|
if handled.Load() == 1 {
|
|
break
|
|
}
|
|
time.Sleep(10 * time.Millisecond)
|
|
}
|
|
if handled.Load() != 1 {
|
|
t.Fatalf("worker loop did not process job")
|
|
}
|
|
if err := wp.Stop(); err != nil {
|
|
t.Fatalf("stop worker: %v", err)
|
|
}
|
|
reloaded := loadJob(t, db, job.ID)
|
|
if reloaded.Status != model.BackgroundJobStatusCompleted {
|
|
t.Fatalf("expected completed job after worker loop: %+v", reloaded)
|
|
}
|
|
}
|
|
|
|
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)
|
|
}
|
|
}
|
|
|
|
func TestWorkerPoolShutdownDrainsActiveJob(t *testing.T) {
|
|
db := newWorkerTestDB(t)
|
|
wp := NewWorkerPoolWithOptions(db, WithPollInterval(5*time.Millisecond))
|
|
started := make(chan struct{})
|
|
release := make(chan struct{})
|
|
wp.Register("drain", func(ctx context.Context, job *model.BackgroundJob) error {
|
|
close(started)
|
|
select {
|
|
case <-release:
|
|
return nil
|
|
case <-ctx.Done():
|
|
return ctx.Err()
|
|
}
|
|
})
|
|
job, err := wp.Enqueue(context.Background(), "drain", nil)
|
|
require.NoError(t, err)
|
|
require.NoError(t, wp.Start())
|
|
select {
|
|
case <-started:
|
|
case <-time.After(time.Second):
|
|
t.Fatal("worker did not start job")
|
|
}
|
|
|
|
done := make(chan error, 1)
|
|
ctx, cancel := context.WithTimeout(context.Background(), time.Second)
|
|
defer cancel()
|
|
go func() { done <- wp.Shutdown(ctx) }()
|
|
select {
|
|
case err := <-done:
|
|
t.Fatalf("shutdown returned before active job drained: %v", err)
|
|
case <-time.After(30 * time.Millisecond):
|
|
}
|
|
close(release)
|
|
require.NoError(t, <-done)
|
|
require.Equal(t, model.BackgroundJobStatusCompleted, loadJob(t, db, job.ID).Status)
|
|
}
|
|
|
|
func TestWorkerPoolShutdownDeadlineCancelsActiveJob(t *testing.T) {
|
|
db := newWorkerTestDB(t)
|
|
wp := NewWorkerPoolWithOptions(db, WithPollInterval(5*time.Millisecond), WithBackoff(func(int) time.Duration { return 0 }))
|
|
started := make(chan struct{})
|
|
cancelled := make(chan struct{})
|
|
release := make(chan struct{})
|
|
wp.Register("deadline", func(ctx context.Context, job *model.BackgroundJob) error {
|
|
close(started)
|
|
<-ctx.Done()
|
|
close(cancelled)
|
|
<-release
|
|
return ctx.Err()
|
|
})
|
|
job, err := wp.Enqueue(context.Background(), "deadline", nil)
|
|
require.NoError(t, err)
|
|
require.NoError(t, wp.Start())
|
|
select {
|
|
case <-started:
|
|
case <-time.After(time.Second):
|
|
t.Fatal("worker did not start job")
|
|
}
|
|
|
|
ctx, cancel := context.WithTimeout(context.Background(), 20*time.Millisecond)
|
|
defer cancel()
|
|
done := make(chan error, 1)
|
|
go func() { done <- wp.Shutdown(ctx) }()
|
|
<-cancelled
|
|
select {
|
|
case err := <-done:
|
|
t.Fatalf("shutdown returned before cancelled handler released: %v", err)
|
|
default:
|
|
}
|
|
close(release)
|
|
require.ErrorIs(t, <-done, context.DeadlineExceeded)
|
|
require.Equal(t, model.BackgroundJobStatusRetrying, loadJob(t, db, job.ID).Status)
|
|
}
|
|
|
|
// --- Redis Stream path tests ---
|
|
|
|
func newMiniRedis(t *testing.T) (*miniredis.Miniredis, *redis.Client) {
|
|
t.Helper()
|
|
mr := miniredis.RunT(t)
|
|
rdb := redis.NewClient(&redis.Options{Addr: mr.Addr()})
|
|
t.Cleanup(func() { rdb.Close() })
|
|
return mr, rdb
|
|
}
|
|
|
|
func newRedisWorkerPool(t *testing.T, db *gorm.DB, rdb redis.UniversalClient, opts ...Option) *WorkerPool {
|
|
t.Helper()
|
|
defaults := []Option{
|
|
WithRedisClient(rdb),
|
|
WithNow(func() time.Time { return time.Date(2026, 7, 10, 12, 0, 0, 0, time.UTC) }),
|
|
WithBlockTimeout(200 * time.Millisecond),
|
|
WithSweepInterval(100 * time.Millisecond),
|
|
WithWorkerID("test-redis-worker"),
|
|
}
|
|
wp := NewWorkerPoolWithOptions(db, append(defaults, opts...)...)
|
|
return wp
|
|
}
|
|
|
|
func TestRedisMessageRejectedWhenShutdownFollowsBatchCheck(t *testing.T) {
|
|
db := newWorkerTestDB(t)
|
|
_, rdb := newMiniRedis(t)
|
|
wp := newRedisWorkerPool(t, db, rdb)
|
|
var handled atomic.Int32
|
|
wp.Register("shutdown_race", func(context.Context, *model.BackgroundJob) error {
|
|
handled.Add(1)
|
|
return nil
|
|
})
|
|
job, err := wp.Enqueue(context.Background(), "shutdown_race", nil)
|
|
require.NoError(t, err)
|
|
|
|
lifecycleCtx, cancel := context.WithCancel(context.Background())
|
|
require.NoError(t, lifecycleCtx.Err())
|
|
cancel() // SIGTERM after the batch-level lifecycle check.
|
|
wp.processRedisMessage(lifecycleCtx, context.Background(), wp.streamKeyFor(job.Queue), redis.XMessage{
|
|
ID: "1-0", Values: map[string]any{"job_id": job.ID},
|
|
})
|
|
|
|
require.Zero(t, handled.Load())
|
|
require.Equal(t, model.BackgroundJobStatusQueued, loadJob(t, db, job.ID).Status)
|
|
}
|
|
|
|
func TestRedisClaimCanceledWhenShutdownRacesClaim(t *testing.T) {
|
|
db := newWorkerTestDB(t)
|
|
_, rdb := newMiniRedis(t)
|
|
wp := newRedisWorkerPool(t, db, rdb, WithBlockTimeout(10*time.Millisecond))
|
|
require.NoError(t, wp.Start())
|
|
|
|
claimStarted := make(chan struct{})
|
|
var once sync.Once
|
|
require.NoError(t, db.Callback().Update().Before("gorm:update").Register("test:block_claim", func(tx *gorm.DB) {
|
|
once.Do(func() { close(claimStarted) })
|
|
<-tx.Statement.Context.Done()
|
|
}))
|
|
t.Cleanup(func() { db.Callback().Update().Remove("test:block_claim") })
|
|
job, err := wp.Enqueue(context.Background(), "shutdown_claim_race", nil)
|
|
require.NoError(t, err)
|
|
<-claimStarted
|
|
|
|
ctx, cancel := context.WithTimeout(context.Background(), 20*time.Millisecond)
|
|
defer cancel()
|
|
require.ErrorIs(t, wp.Shutdown(ctx), context.DeadlineExceeded)
|
|
|
|
require.Equal(t, model.BackgroundJobStatusQueued, loadJob(t, db, job.ID).Status)
|
|
}
|
|
|
|
func TestRedisFailureLifecycleTransitionsOncePerClaim(t *testing.T) {
|
|
db := newWorkerTestDB(t)
|
|
_, rdb := newMiniRedis(t)
|
|
now := time.Date(2026, 7, 10, 12, 0, 0, 0, time.UTC)
|
|
wp := newRedisWorkerPool(t, db, rdb,
|
|
WithNow(func() time.Time { return now }),
|
|
WithBackoff(func(int) time.Duration { return time.Hour }),
|
|
)
|
|
boom := errors.New("provider timeout")
|
|
var handled, failed atomic.Int32
|
|
wp.Register("failing_stream_job", func(context.Context, *model.BackgroundJob) error {
|
|
handled.Add(1)
|
|
return boom
|
|
})
|
|
wp.RegisterFailureHandler("failing_stream_job", func(context.Context, *model.BackgroundJob, time.Duration) error {
|
|
failed.Add(1)
|
|
return nil
|
|
})
|
|
wp.ensureConsumerGroups(context.Background())
|
|
job, err := wp.Enqueue(context.Background(), "failing_stream_job", nil, WithMaxAttempts(2))
|
|
require.NoError(t, err)
|
|
stream := wp.streamKeyFor(job.Queue)
|
|
|
|
wp.processRedisMessage(context.Background(), context.Background(), stream, redis.XMessage{ID: "1-0", Values: map[string]any{"job_id": job.ID}})
|
|
retrying := loadJob(t, db, job.ID)
|
|
require.Equal(t, model.BackgroundJobStatusRetrying, retrying.Status)
|
|
require.Equal(t, 1, retrying.Attempts)
|
|
require.EqualValues(t, 1, handled.Load())
|
|
require.EqualValues(t, 1, failed.Load())
|
|
|
|
wp.processRedisMessage(context.Background(), context.Background(), stream, redis.XMessage{ID: "1-1", Values: map[string]any{"job_id": job.ID}})
|
|
require.Equal(t, model.BackgroundJobStatusRetrying, loadJob(t, db, job.ID).Status)
|
|
require.EqualValues(t, 1, handled.Load(), "duplicate delivery during backoff must not rerun the handler")
|
|
require.EqualValues(t, 1, failed.Load(), "duplicate delivery during backoff must not repeat failure side effects")
|
|
|
|
now = now.Add(time.Hour)
|
|
wp.processRedisMessage(context.Background(), context.Background(), stream, redis.XMessage{ID: "2-0", Values: map[string]any{"job_id": job.ID}})
|
|
dead := loadJob(t, db, job.ID)
|
|
require.Equal(t, model.BackgroundJobStatusDead, dead.Status)
|
|
require.Equal(t, 2, dead.Attempts)
|
|
require.NotNil(t, dead.FailedAt)
|
|
require.EqualValues(t, 2, handled.Load())
|
|
require.EqualValues(t, 2, failed.Load())
|
|
|
|
wp.processRedisMessage(context.Background(), context.Background(), stream, redis.XMessage{ID: "2-1", Values: map[string]any{"job_id": job.ID}})
|
|
require.Equal(t, model.BackgroundJobStatusDead, loadJob(t, db, job.ID).Status)
|
|
require.EqualValues(t, 2, handled.Load(), "dead jobs must not be executed again")
|
|
require.EqualValues(t, 2, failed.Load(), "dead jobs must not repeat failure side effects")
|
|
}
|
|
|
|
// TestRedisEnqueueAndProcessEndToEnd verifies the full Redis path: Enqueue
|
|
// XADDs to the stream, XREADGROUP picks it up, DB claim succeeds, handler
|
|
// runs, and the job reaches "completed" status.
|
|
func TestRedisEnqueueAndProcessEndToEnd(t *testing.T) {
|
|
db := newWorkerTestDB(t)
|
|
_, rdb := newMiniRedis(t)
|
|
now := time.Date(2026, 7, 10, 12, 0, 0, 0, time.UTC)
|
|
wp := newRedisWorkerPool(t, db, rdb, WithNow(func() time.Time { return now }))
|
|
|
|
var handled atomic.Int32
|
|
wp.Register("test_job", func(ctx context.Context, job *model.BackgroundJob) error {
|
|
handled.Add(1)
|
|
if job.LockedBy != "test-redis-worker" {
|
|
t.Fatalf("expected locked_by=test-redis-worker, got %s", job.LockedBy)
|
|
}
|
|
return nil
|
|
})
|
|
|
|
if err := wp.Start(); err != nil {
|
|
t.Fatalf("start: %v", err)
|
|
}
|
|
t.Cleanup(func() {
|
|
if err := wp.Stop(); err != nil {
|
|
t.Errorf("stop worker: %v", err)
|
|
}
|
|
})
|
|
|
|
job, err := wp.Enqueue(context.Background(), "test_job", map[string]any{"k": "v"})
|
|
if err != nil {
|
|
t.Fatalf("enqueue: %v", err)
|
|
}
|
|
|
|
deadline := time.Now().Add(3 * time.Second)
|
|
for time.Now().Before(deadline) {
|
|
if handled.Load() == 1 {
|
|
break
|
|
}
|
|
time.Sleep(10 * time.Millisecond)
|
|
}
|
|
if handled.Load() != 1 {
|
|
t.Fatalf("handler was not invoked via Redis path")
|
|
}
|
|
|
|
reloaded := loadJob(t, db, job.ID)
|
|
if reloaded.Status != model.BackgroundJobStatusCompleted {
|
|
t.Fatalf("expected completed, got %s", reloaded.Status)
|
|
}
|
|
if reloaded.LockedAt != nil || reloaded.LockedBy != "" {
|
|
t.Fatalf("expected cleared lock, got locked_by=%s locked_at=%v", reloaded.LockedBy, reloaded.LockedAt)
|
|
}
|
|
}
|
|
|
|
// TestRedisFallbackToDBPolling verifies that when rdb is nil the pool falls
|
|
// back to DB polling and existing ProcessOne path works.
|
|
func TestRedisFallbackToDBPolling(t *testing.T) {
|
|
db := newWorkerTestDB(t)
|
|
wp := NewWorkerPoolWithOptions(db,
|
|
WithNow(func() time.Time { return time.Date(2026, 7, 10, 12, 0, 0, 0, time.UTC) }),
|
|
WithPollInterval(5*time.Millisecond),
|
|
)
|
|
// rdb is nil → DB polling fallback
|
|
|
|
var handled atomic.Int32
|
|
wp.Register("fallback_job", func(ctx context.Context, job *model.BackgroundJob) error {
|
|
handled.Add(1)
|
|
return nil
|
|
})
|
|
|
|
if err := wp.Start(); err != nil {
|
|
t.Fatalf("start: %v", err)
|
|
}
|
|
t.Cleanup(func() {
|
|
if err := wp.Stop(); err != nil {
|
|
t.Errorf("stop worker: %v", err)
|
|
}
|
|
})
|
|
|
|
if _, err := wp.Enqueue(context.Background(), "fallback_job", nil); err != nil {
|
|
t.Fatalf("enqueue: %v", err)
|
|
}
|
|
|
|
deadline := time.Now().Add(2 * time.Second)
|
|
for time.Now().Before(deadline) {
|
|
if handled.Load() == 1 {
|
|
break
|
|
}
|
|
time.Sleep(10 * time.Millisecond)
|
|
}
|
|
if handled.Load() != 1 {
|
|
t.Fatalf("fallback DB polling did not process job")
|
|
}
|
|
}
|
|
|
|
// TestRedisSweepPicksUpDueDelayedJob verifies that a job enqueued with a
|
|
// future scheduled_at is not XADD'd immediately, but the sweep goroutine
|
|
// pushes it to Redis once it matures.
|
|
func TestRedisSweepPicksUpDueDelayedJob(t *testing.T) {
|
|
db := newWorkerTestDB(t)
|
|
mr, rdb := newMiniRedis(t)
|
|
baseTime := time.Date(2026, 7, 10, 12, 0, 0, 0, time.UTC)
|
|
var elapsed atomic.Int64
|
|
wp := newRedisWorkerPool(t, db, rdb,
|
|
WithNow(func() time.Time { return baseTime.Add(time.Duration(elapsed.Load())) }),
|
|
WithSweepInterval(50*time.Millisecond),
|
|
WithBlockTimeout(50*time.Millisecond),
|
|
)
|
|
|
|
var handled atomic.Int32
|
|
wp.Register("delayed_job", func(ctx context.Context, job *model.BackgroundJob) error {
|
|
handled.Add(1)
|
|
return nil
|
|
})
|
|
|
|
if err := wp.Start(); err != nil {
|
|
t.Fatalf("start: %v", err)
|
|
}
|
|
defer func() { require.NoError(t, wp.Stop()) }()
|
|
|
|
// Enqueue a job scheduled 200ms in the future.
|
|
_, err := wp.Enqueue(context.Background(), "delayed_job", nil, WithScheduledAt(baseTime.Add(200*time.Millisecond)))
|
|
if err != nil {
|
|
t.Fatalf("enqueue: %v", err)
|
|
}
|
|
|
|
// At this point the stream should be empty (no immediate XADD for future jobs).
|
|
// Verify no messages in the default stream yet.
|
|
streamKey := "gochat:jobs:default"
|
|
streamLen := func() int { entries, _ := mr.Stream(streamKey); return len(entries) }()
|
|
if streamLen != 0 {
|
|
t.Fatalf("expected 0 messages in stream before due time, got %d", streamLen)
|
|
}
|
|
|
|
// Advance the mock clock past the scheduled_at so the job becomes due.
|
|
elapsed.Store(int64(time.Second))
|
|
|
|
// The sweep goroutine should now see the job as due, push it to Redis,
|
|
// and the worker should process it.
|
|
deadline := time.Now().Add(5 * time.Second)
|
|
for time.Now().Before(deadline) {
|
|
if handled.Load() == 1 {
|
|
break
|
|
}
|
|
time.Sleep(20 * time.Millisecond)
|
|
}
|
|
if handled.Load() != 1 {
|
|
t.Fatalf("sweep did not pick up delayed job")
|
|
}
|
|
}
|
|
|
|
// TestRedisMultiConsumerCompetition verifies that when two WorkerPool
|
|
// instances share the same consumer group, a job is only executed once.
|
|
func TestRedisMultiConsumerCompetition(t *testing.T) {
|
|
db := newWorkerTestDB(t)
|
|
_, rdb := newMiniRedis(t)
|
|
now := time.Date(2026, 7, 10, 12, 0, 0, 0, time.UTC)
|
|
|
|
var handled atomic.Int32
|
|
handler := func(ctx context.Context, job *model.BackgroundJob) error {
|
|
handled.Add(1)
|
|
time.Sleep(50 * time.Millisecond) // simulate work to widen the race window
|
|
return nil
|
|
}
|
|
|
|
wp1 := newRedisWorkerPool(t, db, rdb, WithNow(func() time.Time { return now }), WithWorkerID("worker-1"))
|
|
wp1.Register("compete_job", handler)
|
|
wp2 := newRedisWorkerPool(t, db, rdb, WithNow(func() time.Time { return now }), WithWorkerID("worker-2"))
|
|
wp2.Register("compete_job", handler)
|
|
|
|
if err := wp1.Start(); err != nil {
|
|
t.Fatalf("start wp1: %v", err)
|
|
}
|
|
t.Cleanup(func() {
|
|
if err := wp1.Stop(); err != nil {
|
|
t.Errorf("stop worker 1: %v", err)
|
|
}
|
|
})
|
|
if err := wp2.Start(); err != nil {
|
|
t.Fatalf("start wp2: %v", err)
|
|
}
|
|
t.Cleanup(func() {
|
|
if err := wp2.Stop(); err != nil {
|
|
t.Errorf("stop worker 2: %v", err)
|
|
}
|
|
})
|
|
|
|
// Enqueue 5 jobs, each gets XADD'd; both workers compete for them.
|
|
for i := 0; i < 5; i++ {
|
|
if _, err := wp1.Enqueue(context.Background(), "compete_job", map[string]any{"i": i}); err != nil {
|
|
t.Fatalf("enqueue %d: %v", i, err)
|
|
}
|
|
}
|
|
|
|
// Wait for all 5 jobs to be handled.
|
|
deadline := time.Now().Add(10 * time.Second)
|
|
for time.Now().Before(deadline) {
|
|
if handled.Load() == 5 {
|
|
break
|
|
}
|
|
time.Sleep(20 * time.Millisecond)
|
|
}
|
|
if handled.Load() != 5 {
|
|
t.Fatalf("expected 5 jobs handled exactly once, got %d", handled.Load())
|
|
}
|
|
|
|
// Verify no job was double-processed: all should be completed, none running/retrying.
|
|
// Poll until all DB records reach completed state (handlers may still be finishing).
|
|
var stuck int64
|
|
checkDeadline := time.Now().Add(5 * time.Second)
|
|
for time.Now().Before(checkDeadline) {
|
|
db.Model(&model.BackgroundJob{}).Where("status != ?", model.BackgroundJobStatusCompleted).Count(&stuck)
|
|
if stuck == 0 {
|
|
break
|
|
}
|
|
time.Sleep(50 * time.Millisecond)
|
|
}
|
|
if stuck != 0 {
|
|
t.Fatalf("expected all jobs completed, found %d in non-completed state", stuck)
|
|
}
|
|
}
|
|
|
|
// TestRedisConsumerGroupCreationIdempotent verifies that calling Start()
|
|
// multiple times (or multiple pools) does not error on XGROUP CREATE.
|
|
func TestRedisConsumerGroupCreationIdempotent(t *testing.T) {
|
|
db := newWorkerTestDB(t)
|
|
_, rdb := newMiniRedis(t)
|
|
|
|
wp := newRedisWorkerPool(t, db, rdb)
|
|
wp.Register("noop", func(ctx context.Context, job *model.BackgroundJob) error { return nil })
|
|
|
|
// ensureConsumerGroups should succeed and create groups.
|
|
ctx := context.Background()
|
|
wp.ensureConsumerGroups(ctx)
|
|
|
|
// Calling again should be idempotent (BUSYGROUP is silently ignored).
|
|
wp.ensureConsumerGroups(ctx)
|
|
|
|
// Verify groups exist on each stream.
|
|
for _, queue := range wp.allQueues() {
|
|
streamKey := wp.streamKeyFor(queue)
|
|
groups, err := rdb.XInfoGroups(ctx, streamKey).Result()
|
|
if err != nil {
|
|
t.Fatalf("XInfoGroups for %s: %v", streamKey, err)
|
|
}
|
|
found := false
|
|
for _, g := range groups {
|
|
if g.Name == wp.consumerGroup {
|
|
found = true
|
|
break
|
|
}
|
|
}
|
|
if !found {
|
|
t.Fatalf("consumer group %s not found on stream %s", wp.consumerGroup, streamKey)
|
|
}
|
|
}
|
|
}
|
|
|
|
// TestRedisEnqueuePushFailureCompensatedBySweep verifies that if XADD fails
|
|
// during Enqueue, the job is still delivered via the sweep mechanism.
|
|
func TestRedisEnqueuePushFailureCompensatedBySweep(t *testing.T) {
|
|
db := newWorkerTestDB(t)
|
|
mr, rdb := newMiniRedis(t)
|
|
now := time.Date(2026, 7, 10, 12, 0, 0, 0, time.UTC)
|
|
|
|
var handled atomic.Int32
|
|
wp := newRedisWorkerPool(t, db, rdb,
|
|
WithNow(func() time.Time { return now }),
|
|
WithSweepInterval(50*time.Millisecond),
|
|
WithBlockTimeout(50*time.Millisecond),
|
|
)
|
|
wp.Register("sweep_recovery", func(ctx context.Context, job *model.BackgroundJob) error {
|
|
handled.Add(1)
|
|
return nil
|
|
})
|
|
|
|
if err := wp.Start(); err != nil {
|
|
t.Fatalf("start: %v", err)
|
|
}
|
|
defer func() { require.NoError(t, wp.Stop()) }()
|
|
|
|
// Close miniredis to simulate Redis being down during Enqueue.
|
|
mr.Close()
|
|
|
|
job, err := wp.Enqueue(context.Background(), "sweep_recovery", nil)
|
|
if err != nil {
|
|
t.Fatalf("enqueue should still succeed (DB committed): %v", err)
|
|
}
|
|
// XADD should have failed silently; job is in DB as queued.
|
|
reloaded := loadJob(t, db, job.ID)
|
|
if reloaded.Status != model.BackgroundJobStatusQueued {
|
|
t.Fatalf("expected queued status after failed XADD, got %s", reloaded.Status)
|
|
}
|
|
|
|
// Restart miniredis at the same address won't work (port already freed).
|
|
// Instead, we manually push to Redis to simulate sweep recovery once Redis is back.
|
|
// Use a fresh miniredis on a new port + new client.
|
|
mr2 := miniredis.RunT(t)
|
|
rdb2 := redis.NewClient(&redis.Options{Addr: mr2.Addr()})
|
|
t.Cleanup(func() { rdb2.Close() })
|
|
|
|
// Restart with the recovered client so no worker can use it while it is replaced.
|
|
if err := wp.Stop(); err != nil {
|
|
t.Fatalf("stop before Redis replacement: %v", err)
|
|
}
|
|
wp.rdb = rdb2
|
|
if err := wp.Start(); err != nil {
|
|
t.Fatalf("restart after Redis replacement: %v", err)
|
|
}
|
|
|
|
// Wait for sweep to pick up the job and push it, then worker to process.
|
|
deadline := time.Now().Add(5 * time.Second)
|
|
for time.Now().Before(deadline) {
|
|
if handled.Load() == 1 {
|
|
break
|
|
}
|
|
time.Sleep(20 * time.Millisecond)
|
|
}
|
|
if handled.Load() != 1 {
|
|
t.Fatalf("sweep did not recover job after Redis reconnection")
|
|
}
|
|
}
|