package worker import ( "context" "encoding/json" "errors" "fmt" "math" "strconv" "strings" "sync" "time" "github.com/gochat/gochat/internal/model" applogger "github.com/gochat/gochat/pkg/logger" "github.com/redis/go-redis/v9" "gorm.io/gorm" "gorm.io/gorm/clause" ) var ErrWorkerDatabaseRequired = errors.New("worker database is required") const redisPublishTimeout = 2 * time.Second type permanentError struct{ err error } func (e *permanentError) Error() string { return e.err.Error() } func (e *permanentError) Unwrap() error { return e.err } func Permanent(err error) error { if err == nil { return nil } return &permanentError{err: err} } // JobHandler performs one durable background job. type JobHandler func(context.Context, *model.BackgroundJob) error // JobFailureHandler runs after a retrying or dead state is persisted. type JobFailureHandler func(context.Context, *model.BackgroundJob, time.Duration) error type BackoffFunc func(attempt int) time.Duration // WorkerPool manages durable background job processors. // Reference: Chatwoot ActiveJob/Sidekiq workers use persistent jobs, scheduled // execution, retry, and dead-letter visibility for async side effects. type WorkerPool struct { db *gorm.DB rdb redis.UniversalClient handlers map[string]JobHandler failureHandlers map[string]JobFailureHandler queues []string workerID string workerCount int pollInterval time.Duration staleLockTimeout time.Duration backoff BackoffFunc now func() time.Time // Redis stream queue parameters streamPrefix string // default "gochat:jobs" consumerGroup string // default "gochat-workers" blockTimeout time.Duration // XREADGROUP block duration, default 5s sweepInterval time.Duration // compensation sweep interval, default 30s mu sync.RWMutex claimMu sync.RWMutex // serializes lifecycle cancellation with claim transitions ctx context.Context cancel context.CancelFunc jobCtx context.Context jobCancel context.CancelFunc wg sync.WaitGroup } type Option func(*WorkerPool) // NewWorkerPool creates a WorkerPool. The variadic DB preserves older no-arg // construction while allowing Phase 5 durable workers to be wired incrementally. func NewWorkerPool(db ...*gorm.DB) *WorkerPool { wp := &WorkerPool{ handlers: make(map[string]JobHandler), failureHandlers: make(map[string]JobFailureHandler), workerID: fmt.Sprintf("worker-%d", time.Now().UnixNano()), workerCount: 1, pollInterval: 500 * time.Millisecond, staleLockTimeout: 15 * time.Minute, backoff: defaultBackoff, now: time.Now, streamPrefix: "gochat:jobs", consumerGroup: "gochat-workers", blockTimeout: 5 * time.Second, sweepInterval: 30 * time.Second, } if len(db) > 0 { wp.db = db[0] } return wp } func NewWorkerPoolWithOptions(db *gorm.DB, opts ...Option) *WorkerPool { wp := NewWorkerPool(db) for _, opt := range opts { opt(wp) } return wp } func WithWorkerCount(count int) Option { return func(wp *WorkerPool) { if count > 0 { wp.workerCount = count } } } func WithQueues(queues ...string) Option { return func(wp *WorkerPool) { if len(queues) > 0 { wp.queues = append([]string(nil), queues...) } } } func WithPollInterval(interval time.Duration) Option { return func(wp *WorkerPool) { if interval > 0 { wp.pollInterval = interval } } } func WithWorkerID(id string) Option { return func(wp *WorkerPool) { if id != "" { wp.workerID = id } } } func WithStaleLockTimeout(timeout time.Duration) Option { return func(wp *WorkerPool) { if timeout > 0 { wp.staleLockTimeout = timeout } } } func WithBackoff(backoff BackoffFunc) Option { return func(wp *WorkerPool) { if backoff != nil { wp.backoff = backoff } } } func WithNow(now func() time.Time) Option { return func(wp *WorkerPool) { if now != nil { wp.now = now } } } // WithRedisClient configures the WorkerPool to use Redis Streams for job // dispatch instead of DB polling. When set, Enqueue pushes to Redis and // workers consume via XREADGROUP. When nil, the pool falls back to DB polling. func WithRedisClient(rdb redis.UniversalClient) Option { return func(wp *WorkerPool) { wp.rdb = rdb } } // WithStreamPrefix sets the Redis stream key prefix (default "gochat:jobs"). func WithStreamPrefix(prefix string) Option { return func(wp *WorkerPool) { if prefix != "" { wp.streamPrefix = prefix } } } // WithConsumerGroup sets the Redis consumer group name (default "gochat-workers"). func WithConsumerGroup(group string) Option { return func(wp *WorkerPool) { if group != "" { wp.consumerGroup = group } } } // WithBlockTimeout sets the XREADGROUP block duration (default 5s). func WithBlockTimeout(d time.Duration) Option { return func(wp *WorkerPool) { if d > 0 { wp.blockTimeout = d } } } // WithSweepInterval sets the compensation sweep interval (default 30s). func WithSweepInterval(d time.Duration) Option { return func(wp *WorkerPool) { if d > 0 { wp.sweepInterval = d } } } func (wp *WorkerPool) Register(jobType string, handler JobHandler) { wp.mu.Lock() defer wp.mu.Unlock() wp.handlers[jobType] = handler } func (wp *WorkerPool) RegisterFailureHandler(jobType string, handler JobFailureHandler) { wp.mu.Lock() defer wp.mu.Unlock() wp.failureHandlers[jobType] = handler } type EnqueueOption func(*model.BackgroundJob) func WithQueue(queue string) EnqueueOption { return func(job *model.BackgroundJob) { if queue != "" { job.Queue = queue } } } func WithScheduledAt(at time.Time) EnqueueOption { return func(job *model.BackgroundJob) { if !at.IsZero() { job.ScheduledAt = at } } } // WithMaxAttempts sets the absolute claim limit. A claim interrupted by a // process crash still consumes an attempt. func WithMaxAttempts(max int) EnqueueOption { return func(job *model.BackgroundJob) { if max > 0 { job.MaxAttempts = max } } } func WithPriority(priority int) EnqueueOption { return func(job *model.BackgroundJob) { job.Priority = priority } } func WithIdempotencyKey(key string) EnqueueOption { return func(job *model.BackgroundJob) { job.IdempotencyKey = key } } func (wp *WorkerPool) Enqueue(ctx context.Context, jobType string, payload any, opts ...EnqueueOption) (*model.BackgroundJob, error) { if wp.db == nil { return nil, ErrWorkerDatabaseRequired } job, created, err := wp.persistJob(ctx, wp.db, jobType, payload, opts...) if err != nil { return nil, err } if created { wp.Publish(ctx, job) } return job, nil } // EnqueueInTransaction persists a job on the caller's transaction. Call // Publish only after the transaction commits successfully. func (wp *WorkerPool) EnqueueInTransaction(ctx context.Context, tx *gorm.DB, jobType string, payload any, opts ...EnqueueOption) (*model.BackgroundJob, bool, error) { if tx == nil { return nil, false, ErrWorkerDatabaseRequired } return wp.persistJob(ctx, tx, jobType, payload, opts...) } func (wp *WorkerPool) persistJob(ctx context.Context, db *gorm.DB, jobType string, payload any, opts ...EnqueueOption) (*model.BackgroundJob, bool, error) { if jobType == "" { return nil, false, errors.New("job type is required") } payloadBytes, err := marshalPayload(payload) if err != nil { return nil, false, err } job := &model.BackgroundJob{ Queue: model.DefaultBackgroundJobQueue, JobType: jobType, Payload: payloadBytes, Status: model.BackgroundJobStatusQueued, MaxAttempts: 3, ScheduledAt: wp.now(), } for _, opt := range opts { opt(job) } if job.ScheduledAt.IsZero() { job.ScheduledAt = wp.now() } if job.MaxAttempts <= 0 { job.MaxAttempts = 3 } if job.IdempotencyKey != "" { var existing model.BackgroundJob err := db.WithContext(ctx).Where("idempotency_key = ?", job.IdempotencyKey).First(&existing).Error if err == nil { return &existing, false, nil } if !errors.Is(err, gorm.ErrRecordNotFound) { return nil, false, err } } query := db.WithContext(ctx) if job.IdempotencyKey != "" { query = query.Clauses(clause.OnConflict{DoNothing: true}) } result := query.Create(job) if result.Error != nil { return nil, false, result.Error } if job.IdempotencyKey != "" && result.RowsAffected == 0 { var existing model.BackgroundJob if err := db.WithContext(ctx).Where("idempotency_key = ?", job.IdempotencyKey).First(&existing).Error; err != nil { return nil, false, err } return &existing, false, nil } return job, true, nil } func (wp *WorkerPool) Publish(ctx context.Context, job *model.BackgroundJob) { if job == nil { return } // Push to Redis Stream for immediate dispatch if the job is due. // Scheduled (future) jobs are picked up by the sweep goroutine when they mature. // Keep Redis delivery bounded: the durable DB row is already committed, and // a slow/unreachable Redis must not hold the HTTP request until its deadline. // ponytail: 2s cap; the durable sweep retries Redis delivery. if wp.rdb != nil && !job.ScheduledAt.After(wp.now()) { publishCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), redisPublishTimeout) defer cancel() if err := wp.pushToRedis(publishCtx, job); err != nil { // Redis delivery failure does not roll back the DB-committed job. // The sweep compensation mechanism will re-push it on the next cycle. applogger.L().Warnf( "redis XAdd failed for job %d (queue=%s), will be picked up by sweep: %v", job.ID, job.Queue, err, ) } } } func (wp *WorkerPool) Start() error { if wp.db == nil { return nil } wp.mu.Lock() if wp.cancel != nil { wp.mu.Unlock() return nil } wp.ctx, wp.cancel = context.WithCancel(context.Background()) wp.jobCtx, wp.jobCancel = context.WithCancel(context.Background()) ctx, jobCtx := wp.ctx, wp.jobCtx cancel, jobCancel := wp.cancel, wp.jobCancel workerCount := wp.workerCount wp.mu.Unlock() if _, err := wp.RequeueStaleJobs(ctx); err != nil { wp.mu.Lock() wp.cancel = nil wp.ctx = nil wp.jobCancel = nil wp.jobCtx = nil wp.mu.Unlock() cancel() jobCancel() return err } // Ensure Redis consumer groups exist for all configured queues. if wp.rdb != nil { wp.ensureConsumerGroups(ctx) } for i := 0; i < workerCount; i++ { wp.wg.Add(1) go wp.run(ctx, jobCtx, i) } // Start the sweep goroutine for delayed-job delivery and stale-job recovery. if wp.rdb != nil { wp.wg.Add(1) go wp.sweepLoop(ctx) } return nil } func (wp *WorkerPool) Stop() error { wp.mu.Lock() cancel := wp.cancel jobCancel := wp.jobCancel wp.cancel = nil wp.ctx = nil wp.jobCancel = nil wp.jobCtx = nil wp.mu.Unlock() if cancel != nil { cancel() } if jobCancel != nil { jobCancel() } wp.claimMu.Lock() wp.claimMu.Unlock() wp.wg.Wait() return nil } // Shutdown stops claiming new jobs, lets active handlers finish, and only // cancels them if the drain deadline expires. func (wp *WorkerPool) Shutdown(ctx context.Context) error { wp.mu.Lock() cancel := wp.cancel jobCancel := wp.jobCancel wp.cancel = nil wp.ctx = nil wp.jobCancel = nil wp.jobCtx = nil wp.mu.Unlock() if cancel == nil { return nil } cancel() claimsDone := make(chan struct{}) go func() { wp.claimMu.Lock() wp.claimMu.Unlock() close(claimsDone) }() timedOut := false select { case <-claimsDone: case <-ctx.Done(): timedOut = true if jobCancel != nil { jobCancel() } <-claimsDone } done := make(chan struct{}) go func() { wp.wg.Wait() close(done) }() if timedOut { <-done return ctx.Err() } select { case <-done: if jobCancel != nil { jobCancel() } return nil case <-ctx.Done(): if jobCancel != nil { jobCancel() } <-done return ctx.Err() } } func (wp *WorkerPool) ProcessOne(ctx context.Context) (bool, error) { return wp.processOne(nil, ctx) } func (wp *WorkerPool) processOne(lifecycleCtx, jobCtx context.Context) (bool, error) { if wp.db == nil { return false, ErrWorkerDatabaseRequired } if err := jobCtx.Err(); err != nil { return false, err } job, err := wp.claimNext(lifecycleCtx, jobCtx) if err != nil { if errors.Is(err, gorm.ErrRecordNotFound) { return false, nil } return false, err } return true, wp.perform(jobCtx, job) } func (wp *WorkerPool) RequeueStaleJobs(ctx context.Context) (int64, error) { if wp.db == nil { return 0, ErrWorkerDatabaseRequired } now := wp.now() cutoff := now.Add(-wp.staleLockTimeout) updates := map[string]any{ "status": gorm.Expr( "CASE WHEN attempts >= max_attempts THEN ? ELSE ? END", model.BackgroundJobStatusDead, model.BackgroundJobStatusRetrying, ), "failed_at": gorm.Expr("CASE WHEN attempts >= max_attempts THEN ? ELSE failed_at END", now), "locked_at": nil, "locked_by": "", } result := wp.db.WithContext(ctx).Model(&model.BackgroundJob{}). Where("status = ? AND locked_at IS NOT NULL AND locked_at < ?", model.BackgroundJobStatusRunning, cutoff). Updates(updates) return result.RowsAffected, result.Error } func archiveExhaustedRetries(db *gorm.DB, now time.Time) error { return db.Model(&model.BackgroundJob{}). Where("status = ? AND attempts >= max_attempts", model.BackgroundJobStatusRetrying). Updates(map[string]any{ "status": model.BackgroundJobStatusDead, "failed_at": now, }).Error } // run is the per-goroutine consume loop. When Redis is configured it uses // XREADGROUP BLOCK; otherwise it falls back to DB polling. func (wp *WorkerPool) run(ctx, jobCtx context.Context, index int) { defer wp.wg.Done() // Each goroutine gets a unique Redis consumer name so XINFO CONSUMERS // can distinguish them. The DB locked_by field still uses wp.workerID. consumerID := fmt.Sprintf("%s-%d", wp.workerID, index) if wp.rdb == nil { wp.runDBPollLoop(ctx, jobCtx) return } // Build the XREADGROUP streams list: [streamKey1, ">", streamKey2, ">", ...] streams := wp.streamKeys() for { if ctx.Err() != nil { return } results, err := wp.rdb.XReadGroup(ctx, &redis.XReadGroupArgs{ Group: wp.consumerGroup, Consumer: consumerID, Streams: streams, Count: 1, Block: wp.blockTimeout, }).Result() if ctx.Err() != nil { return } if err != nil && !errors.Is(err, redis.Nil) { applogger.L().Errorf("XReadGroup error: %v", err) select { case <-ctx.Done(): return case <-time.After(time.Second): } continue } for _, xstream := range results { for _, msg := range xstream.Messages { if ctx.Err() != nil { return } wp.processRedisMessage(ctx, jobCtx, xstream.Stream, msg) } } } } // runDBPollLoop is the legacy DB-polling loop, used when Redis is unavailable // (e.g. SQLite test mode) or not configured. func (wp *WorkerPool) runDBPollLoop(ctx, jobCtx context.Context) { ticker := time.NewTicker(wp.pollInterval) defer ticker.Stop() for { select { case <-ctx.Done(): return default: } processed, _ := wp.processOne(ctx, jobCtx) if processed { continue } select { case <-ctx.Done(): return case <-ticker.C: } } } // processRedisMessage loads the job from DB, claims it atomically, executes // the handler, and acknowledges the Redis message. The DB claim check // guarantees at-most-once execution even if the same job is XADD'd multiple // times (e.g. by both Enqueue and sweep). func (wp *WorkerPool) processRedisMessage(lifecycleCtx, jobCtx context.Context, stream string, msg redis.XMessage) { if lifecycleCtx.Err() != nil { return } jobIDStr, ok := msg.Values["job_id"] if !ok { wp.ackRedis(jobCtx, stream, msg.ID) return } jobID, err := strconv.ParseUint(fmt.Sprintf("%v", jobIDStr), 10, 64) if err != nil { wp.ackRedis(jobCtx, stream, msg.ID) return } var job model.BackgroundJob if err := wp.db.WithContext(jobCtx).First(&job, jobID).Error; err != nil { if errors.Is(err, gorm.ErrRecordNotFound) { // Job was deleted; ACK and drop. wp.ackRedis(jobCtx, stream, msg.ID) return } // DB error — do not ACK so Redis can redeliver to another consumer. applogger.L().Errorf("load job %d from DB failed: %v", jobID, err) return } now := wp.now() if err := archiveExhaustedRetries(wp.db.WithContext(jobCtx).Where("id = ?", jobID), now); err != nil { applogger.L().Errorf("archive exhausted job %d failed: %v", jobID, err) return } // Claim check: atomically transition queued/retrying → running. // RowsAffected == 0 means another consumer already claimed it or the // job is not yet due; ACK to avoid redelivery loops. wp.claimMu.RLock() if lifecycleCtx.Err() != nil { wp.claimMu.RUnlock() return } result := wp.db.WithContext(jobCtx).Model(&model.BackgroundJob{}). Where("id = ? AND status IN ? AND scheduled_at <= ? AND attempts < max_attempts", jobID, []string{model.BackgroundJobStatusQueued, model.BackgroundJobStatusRetrying}, now, ). Updates(map[string]any{ "status": model.BackgroundJobStatusRunning, "locked_at": &now, "locked_by": wp.workerID, "attempts": gorm.Expr("attempts + 1"), }) wp.claimMu.RUnlock() if result.Error != nil { applogger.L().Errorf("claim job %d failed: %v", jobID, result.Error) return } if result.RowsAffected == 0 { wp.ackRedis(jobCtx, stream, msg.ID) return } // Reload the job with updated attempts/locked_by fields. wp.db.WithContext(jobCtx).First(&job, jobID) if err := wp.perform(jobCtx, &job); err != nil { applogger.L().Errorf("job %d failed: %v", job.ID, err) } wp.ackRedis(jobCtx, stream, msg.ID) } func (wp *WorkerPool) ackRedis(ctx context.Context, stream, msgID string) { if err := wp.rdb.XAck(ctx, stream, wp.consumerGroup, msgID).Err(); err != nil { applogger.L().Warnf("XAck failed for stream=%s msgID=%s: %v", stream, msgID, err) } } // pushToRedis adds a job ID message to the Redis stream for the job's queue. func (wp *WorkerPool) pushToRedis(ctx context.Context, job *model.BackgroundJob) error { streamKey := wp.streamKeyFor(job.Queue) return wp.rdb.XAdd(ctx, &redis.XAddArgs{ Stream: streamKey, Values: map[string]interface{}{ "job_id": strconv.FormatUint(uint64(job.ID), 10), "job_type": job.JobType, "queue": job.Queue, }, }).Err() } func (wp *WorkerPool) streamKeyFor(queue string) string { return fmt.Sprintf("%s:%s", wp.streamPrefix, queue) } // streamKeys returns the go-redis Streams slice: all stream keys first, then all IDs. // The ">" sentinel tells Redis to deliver only messages never seen by this group. func (wp *WorkerPool) streamKeys() []string { queues := wp.allQueues() keys := make([]string, 0, len(queues)*2) for _, q := range queues { keys = append(keys, wp.streamKeyFor(q)) } for range queues { keys = append(keys, ">") } return keys } // allQueues returns the queues this worker listens on. If none are explicitly // configured, it defaults to the full set of known queues. func (wp *WorkerPool) allQueues() []string { if len(wp.queues) > 0 { return wp.queues } return []string{ "default", "high", "medium", "low", "events", "automation", "search", "scheduled_jobs", "deferred", "purgable", } } // ensureConsumerGroups creates the Redis consumer group for each queue's stream. // XGROUP CREATE is idempotent: a BUSYGROUP reply means the group already exists. func (wp *WorkerPool) ensureConsumerGroups(ctx context.Context) { for _, queue := range wp.allQueues() { streamKey := wp.streamKeyFor(queue) err := wp.rdb.XGroupCreateMkStream(ctx, streamKey, wp.consumerGroup, "$").Err() if err != nil && !strings.Contains(err.Error(), "BUSYGROUP") { applogger.L().Warnf("XGroupCreate for stream %s group %s: %v", streamKey, wp.consumerGroup, err) } } } // sweepLoop periodically scans the DB for due jobs that are still in // queued/retrying state and re-pushes them to Redis. It also recovers stale // running jobs. This compensates for Redis delivery failures and handles // delayed (scheduled_at > now) jobs that were not XADD'd at enqueue time. func (wp *WorkerPool) sweepLoop(ctx context.Context) { defer wp.wg.Done() ticker := time.NewTicker(wp.sweepInterval) defer ticker.Stop() for { select { case <-ctx.Done(): return case <-ticker.C: if err := wp.sweepDueJobs(ctx); err != nil { applogger.L().Warnf("sweep due jobs failed: %v", err) } if _, err := wp.RequeueStaleJobs(ctx); err != nil { applogger.L().Warnf("requeue stale jobs failed: %v", err) } } } } func (wp *WorkerPool) sweepDueJobs(ctx context.Context) error { now := wp.now() var jobs []model.BackgroundJob if err := wp.db.WithContext(ctx). Where("status IN ? AND scheduled_at <= ?", []string{model.BackgroundJobStatusQueued, model.BackgroundJobStatusRetrying}, now, ). Limit(100). Find(&jobs).Error; err != nil { return err } for i := range jobs { // Idempotent push: duplicate XADD does not cause duplicate execution // because processRedisMessage has a DB claim check. if err := wp.pushToRedis(ctx, &jobs[i]); err != nil { applogger.L().Warnf("sweep: XAdd failed for job %d: %v", jobs[i].ID, err) } } return nil } func (wp *WorkerPool) claimNext(lifecycleCtx, jobCtx context.Context) (*model.BackgroundJob, error) { wp.claimMu.RLock() defer wp.claimMu.RUnlock() if lifecycleCtx != nil { if err := lifecycleCtx.Err(); err != nil { return nil, err } } if err := archiveExhaustedRetries(wp.db.WithContext(jobCtx), wp.now()); err != nil { return nil, err } var job model.BackgroundJob err := wp.db.WithContext(jobCtx).Transaction(func(tx *gorm.DB) error { query := tx.Where("status IN ? AND scheduled_at <= ? AND attempts < max_attempts", []string{model.BackgroundJobStatusQueued, model.BackgroundJobStatusRetrying}, wp.now()) if len(wp.queues) > 0 { query = query.Where("queue IN ?", wp.queues) } if tx.Dialector != nil && tx.Dialector.Name() == "postgres" { query = query.Clauses(clause.Locking{Strength: "UPDATE", Options: "SKIP LOCKED"}) } if err := query.Order("priority DESC, scheduled_at ASC, id ASC").First(&job).Error; err != nil { return err } lockedAt := wp.now() updates := map[string]any{ "status": model.BackgroundJobStatusRunning, "locked_at": &lockedAt, "locked_by": wp.workerID, "attempts": gorm.Expr("attempts + 1"), } if err := tx.Model(&model.BackgroundJob{}).Where("id = ?", job.ID).Updates(updates).Error; err != nil { return err } return tx.First(&job, job.ID).Error }) if err != nil { return nil, err } return &job, nil } func (wp *WorkerPool) perform(ctx context.Context, job *model.BackgroundJob) error { handler := wp.handlerFor(job.JobType) if handler == nil { return wp.fail(ctx, job, fmt.Errorf("no handler registered for job type %q", job.JobType)) } 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, "locked_at": nil, "locked_by": "", "last_error": "", "finished_at": &finishedAt, } return wp.db.WithContext(ctx).Model(&model.BackgroundJob{}).Where("id = ?", job.ID).Updates(updates).Error } 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, "locked_by": "", "last_error": err.Error(), } var permanent *permanentError retryAfter := time.Duration(0) if job.Attempts >= job.MaxAttempts || errors.As(err, &permanent) { updates["status"] = model.BackgroundJobStatusDead updates["failed_at"] = &now } else { retryAfter = wp.backoff(job.Attempts) updates["status"] = model.BackgroundJobStatusRetrying updates["scheduled_at"] = now.Add(retryAfter) } if updateErr := wp.db.WithContext(ctx).Model(&model.BackgroundJob{}).Where("id = ?", job.ID).Updates(updates).Error; updateErr != nil { return updateErr } job.Status = updates["status"].(string) if job.Status == model.BackgroundJobStatusDead { job.FailedAt = &now } else { job.ScheduledAt = updates["scheduled_at"].(time.Time) } if handler := wp.failureHandlerFor(job.JobType); handler != nil { if handlerErr := handler(ctx, job, retryAfter); handlerErr != nil { return errors.Join(err, handlerErr) } } return err } func (wp *WorkerPool) handlerFor(jobType string) JobHandler { wp.mu.RLock() defer wp.mu.RUnlock() return wp.handlers[jobType] } func (wp *WorkerPool) failureHandlerFor(jobType string) JobFailureHandler { wp.mu.RLock() defer wp.mu.RUnlock() return wp.failureHandlers[jobType] } func marshalPayload(payload any) (json.RawMessage, error) { if payload == nil { return json.RawMessage(`{}`), nil } switch value := payload.(type) { case json.RawMessage: if !json.Valid(value) { return nil, errors.New("job payload must be valid JSON") } return value, nil case []byte: if !json.Valid(value) { return nil, errors.New("job payload must be valid JSON") } return json.RawMessage(value), nil default: payloadBytes, err := json.Marshal(value) if err != nil { return nil, err } return payloadBytes, nil } } func defaultBackoff(attempt int) time.Duration { if attempt <= 0 { attempt = 1 } seconds := math.Pow(2, float64(attempt-1)) if seconds > 300 { seconds = 300 } return time.Duration(seconds) * time.Second }