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 TestWorkerPoolDeadLettersJobWhenFinalClaimGoesStale(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)) job, err := wp.Enqueue(context.Background(), "crashing_job", nil, WithMaxAttempts(3)) require.NoError(t, err) for attempt := 1; attempt <= job.MaxAttempts; attempt++ { claimed, claimErr := wp.claimNext(nil, context.Background()) require.NoError(t, claimErr) require.Equal(t, attempt, claimed.Attempts) now = now.Add(2 * time.Minute) count, requeueErr := wp.RequeueStaleJobs(context.Background()) require.NoError(t, requeueErr) require.EqualValues(t, 1, count) } reloaded := loadJob(t, db, job.ID) require.Equal(t, model.BackgroundJobStatusDead, reloaded.Status) require.Equal(t, job.MaxAttempts, reloaded.Attempts) require.NotNil(t, reloaded.FailedAt) processed, err := wp.ProcessOne(context.Background()) require.NoError(t, err) require.False(t, processed, "an exhausted stale job must not receive a fourth claim") } func TestWorkerPoolArchivesPersistedExhaustedRetry(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 })) var handled atomic.Int32 wp.Register("legacy_retry", func(context.Context, *model.BackgroundJob) error { handled.Add(1) return nil }) job, err := wp.Enqueue(context.Background(), "legacy_retry", nil, WithMaxAttempts(3)) require.NoError(t, err) require.NoError(t, db.Model(job).Updates(map[string]any{ "status": model.BackgroundJobStatusRetrying, "attempts": job.MaxAttempts, }).Error) processed, err := wp.ProcessOne(context.Background()) require.NoError(t, err) require.False(t, processed) require.Zero(t, handled.Load()) reloaded := loadJob(t, db, job.ID) require.Equal(t, model.BackgroundJobStatusDead, reloaded.Status) require.Equal(t, job.MaxAttempts, reloaded.Attempts) require.NotNil(t, reloaded.FailedAt) require.Equal(t, now, *reloaded.FailedAt) failedAt := *reloaded.FailedAt now = now.Add(time.Hour) processed, err = wp.ProcessOne(context.Background()) require.NoError(t, err) require.False(t, processed) require.Equal(t, failedAt, *loadJob(t, db, job.ID).FailedAt) } 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") } }