-
Notifications
You must be signed in to change notification settings - Fork 0
Error Handling
Senna provides robust error handling with automatic retries, configurable backoff, and a dead letter queue for permanently failed jobs.
Return an error from your handler to trigger a retry:
w.Register("send_email", func(ctx context.Context, job *senna.Job) error {
err := sendEmail(job.Args)
if err != nil {
return err // Will be retried
}
return nil // Success
})- Default retry count: 25 (configurable per-client or per-job)
- Retries use exponential backoff
delay = (attempt^4) + 15 + (attempt * 10) seconds
| Attempt | Delay |
|---|---|
| 1 | 26 seconds |
| 2 | 51 seconds |
| 3 | 126 seconds (~2 min) |
| 4 | 311 seconds (~5 min) |
| 5 | 670 seconds (~11 min) |
| 10 | ~2.8 hours |
| 15 | ~7.1 hours |
| 20 | ~18.5 hours |
| 25 | ~43.4 hours |
Per-client default:
c, _ := client.New(&client.Config{
Settings: client.Settings{
DefaultRetry: 10, // Default for all jobs
},
})Per-job override:
c.Enqueue(ctx, "fragile_job", args, client.WithRetry(3))Per-handler override:
w.Register("fragile_job", handler, worker.WithJobMaxRetries(3))Return a standard error for automatic retry with default backoff:
return errors.New("connection failed")
return fmt.Errorf("API error: %w", err)Control exactly when the job should be retried:
return &senna.RetryableError{
Job: job,
Cause: errors.New("rate limited"),
RetryIn: 5 * time.Minute, // Retry in 5 minutes
}Use this when:
- External service tells you when to retry (e.g.,
Retry-Afterheader) - You know the optimal retry timing
- Default backoff isn't appropriate
Immediately move the job to the dead queue (skip remaining retries):
return &senna.MaxRetriesExceededError{
Job: job,
Cause: errors.New("invalid data - will never succeed"),
}Use this for:
- Permanent failures (invalid data, missing resources)
- Errors that won't be fixed by retrying
- Business logic failures
Caution
Use MaxRetriesExceededError carefully. Once a job is moved to the dead queue, it won't be automatically retried. Only use this for errors you're certain won't be fixed by retrying.
// 1s, 2s, 4s, 8s... up to 1 hour max
backoff := senna.ExponentialBackoff(time.Second, time.Hour)
w.Use(senna.RetryMiddleware(5, backoff))backoff := func(attempt int) time.Duration {
// Linear backoff: 1m, 2m, 3m, 4m...
return time.Duration(attempt) * time.Minute
}
w.Use(senna.RetryMiddleware(5, backoff))Add randomness to prevent thundering herd:
backoff := func(attempt int) time.Duration {
base := time.Duration(attempt*attempt) * time.Second
jitter := time.Duration(rand.Intn(1000)) * time.Millisecond
return base + jitter
}Tip
Adding jitter to your backoff strategy prevents multiple failed jobs from retrying at the same time, which can overwhelm external services.
Jobs that exhaust all retries are moved to the dead queue for inspection and possible manual retry.
Note
Dead jobs are preserved for debugging. You can inspect them to understand why they failed and potentially retry them after fixing the underlying issue.
- Exceeded max retry count
- Returned
MaxRetriesExceededError - Handler panicked on every attempt
Dead jobs are stored in a Redis sorted set:
namespace:dead
Score is the Unix timestamp when the job died.
The recovery middleware (automatically added) catches panics and converts them to errors:
w.Register("risky_job", func(ctx context.Context, job *senna.Job) error {
// If this panics, the job will be retried
riskyOperation()
return nil
})If a panic occurs:
- Panic is caught by recovery middleware
- Job is returned to retry queue
- Worker continues processing other jobs
When a job times out or the worker shuts down, the context is cancelled. This is
cooperative cancellation: Senna cannot force-stop a handler that ignores
ctx.Done(). A handler that does not return keeps running and keeps its worker
slot occupied until it finishes.
Check the context around long loops, blocking calls, and external work so the job can be retried or returned to the queue promptly:
w.Register("long_job", func(ctx context.Context, job *senna.Job) error {
for i := 0; i < 1000; i++ {
select {
case <-ctx.Done():
// Save progress for retry
saveProgress(i)
return ctx.Err()
default:
processItem(i)
}
}
return nil
})If a worker crashes, its in-flight jobs become orphans. Senna's reaper process:
- Runs every 30 seconds
- Checks for dead workers (no heartbeat)
- Returns orphaned jobs to their queues
This ensures jobs aren't lost when workers crash.
Important
The reaper only recovers orphaned jobs if workers are still running. If all workers are stopped, orphaned jobs will be recovered when a worker starts again.
w.Register("process", func(ctx context.Context, job *senna.Job) error {
data, err := fetchData(job.Args["id"])
if err != nil {
if isNotFoundError(err) {
// Permanent failure - data doesn't exist
return &senna.MaxRetriesExceededError{
Job: job,
Cause: err,
}
}
// Transient failure - retry later
return err
}
return process(data)
})w.Register("api_call", func(ctx context.Context, job *senna.Job) error {
resp, err := callAPI()
if err != nil {
if isRateLimitError(err) {
retryAfter := parseRetryAfter(err)
return &senna.RetryableError{
Job: job,
Cause: err,
RetryIn: retryAfter,
}
}
return err
}
return nil
})Save progress for long-running jobs:
w.Register("batch_process", func(ctx context.Context, job *senna.Job) error {
startIndex := 0
if idx, ok := job.Args["resume_from"]; ok {
startIndex = int(idx.(float64))
}
items := getItems()
for i := startIndex; i < len(items); i++ {
select {
case <-ctx.Done():
// Update job args for resume
job.Args["resume_from"] = i
return ctx.Err()
default:
if err := processItem(items[i]); err != nil {
job.Args["resume_from"] = i
return err
}
}
}
return nil
})w.Use(senna.LoggingMiddleware(logger))
// Logs all failures with error detailsfunc FailureMetrics() senna.Middleware {
return func(next senna.Handler) senna.Handler {
return func(ctx context.Context, job *senna.Job) error {
err := next(ctx, job)
if err != nil {
metrics.IncrCounter("job.failure",
"type", job.Type,
"error", errorType(err),
)
}
return err
}
}
}func DeathAlertMiddleware() senna.Middleware {
return func(next senna.Handler) senna.Handler {
return func(ctx context.Context, job *senna.Job) error {
err := next(ctx, job)
var maxRetries *senna.MaxRetriesExceededError
if errors.As(err, &maxRetries) {
alertOps(job, err)
}
return err
}
}
}package main
import (
"context"
"errors"
"time"
"github.com/mgomes/senna"
"github.com/mgomes/senna/worker"
)
func main() {
w, _ := worker.New(&worker.Config{
Redis: senna.RedisConfig{Addr: "localhost:6379"},
Namespace: "myapp",
})
w.Register("sync_user", func(ctx context.Context, job *senna.Job) error {
userID := int(job.Args["user_id"].(float64))
user, err := fetchUser(userID)
if err != nil {
if isNotFoundError(err) {
// User doesn't exist - permanent failure
return &senna.MaxRetriesExceededError{
Job: job,
Cause: err,
}
}
// Transient error - retry
return err
}
err = syncToExternalService(user)
if err != nil {
if isRateLimitError(err) {
// Rate limited - retry after delay
return &senna.RetryableError{
Job: job,
Cause: err,
RetryIn: 1 * time.Minute,
}
}
return err
}
return nil
}, worker.WithJobMaxRetries(5))
w.Run(context.Background())
}| Error Type | Behavior |
|---|---|
| Standard error | Retry with default backoff |
RetryableError |
Retry after specified duration |
MaxRetriesExceededError |
Move to dead queue immediately |
| Panic | Recover and retry |
| Context cancelled | Job returned to queue |
| Max retries exceeded | Move to dead queue |
- Configure Middleware for logging and metrics
- Use Rate Limiters to prevent overload
- Group jobs with Batches