Background Jobs
Grit uses asynq -- a Redis-backed job queue library for Go -- to handle background processing. Send emails, generate thumbnails, clean up expired tokens, or run any async task with automatic retries and priority queues.
Architecture
The job system has two components: a client that enqueues jobs and a worker that processes them. Both connect to the same Redis instance. The worker runs as a goroutine inside the API server -- no separate process needed.
┌──────────────────┐ ┌─────────┐ ┌──────────────────┐│ API Handler │ │ Redis │ │ Worker ││ │ │ │ │ ││ client.Enqueue │────>│ Queue │────>│ handleEmailSend ││ SendEmail(...) │ │ │ │ handleImage... ││ │ │ │ │ handleCleanup │└──────────────────┘ └─────────┘ └──────────────────┘
Job Client
The job client at internal/jobs/client.go provides typed methods for enqueuing each built-in job. It handles JSON serialization of payloads and configures retry policies per job type.
// Task type constantsconst (TypeEmailSend = "email:send"TypeImageProcess = "image:process"TypeTokensCleanup = "tokens:cleanup")// Client wraps asynq.Client for enqueuing background jobs.type Client struct {client *asynq.Client}// NewClient creates a new job queue client connected to Redis.func NewClient(redisURL string) (*Client, error)// Close shuts down the client connection.func (c *Client) Close() error
Enqueue Methods
// Every enqueue method takes a context.Context first and accepts// optional ...EnqueueOption (queue, max retries, delay, idempotency key).// EnqueueSendEmail enqueues an email send job.func (c *Client) EnqueueSendEmail(ctx context.Context,to, subject, template string,data map[string]interface{},opts ...EnqueueOption,) error// EnqueueProcessImage enqueues an image processing job.// (uploadID is the Upload's UUID string.)func (c *Client) EnqueueProcessImage(ctx context.Context,uploadID string,key, mimeType string,opts ...EnqueueOption,) error// EnqueueTokensCleanup enqueues a token cleanup job.func (c *Client) EnqueueTokensCleanup(ctx context.Context) error
Job Payloads
// EmailPayload holds the data for an email send job.type EmailPayload struct {To string `json:"to"`Subject string `json:"subject"`Template string `json:"template"`Data map[string]interface{} `json:"data"`}// ImagePayload holds the data for an image processing job.type ImagePayload struct {UploadID string `json:"upload_id"` // Upload UUIDKey string `json:"key"`MimeType string `json:"mime_type"`}
Worker Setup
The worker is started in main.go alongside the HTTP server. It receives all service dependencies via the WorkerDeps struct, giving handler functions access to the database, mailer, storage, and cache.
// WorkerDeps holds dependencies needed by job handlers.type WorkerDeps struct {DB *gorm.DBMailer *mail.MailerStorage *storage.StorageCache *cache.Cache}// StartWorker starts the asynq worker server in a goroutine.// Returns a stop function and any startup error.func StartWorker(redisURL string, deps WorkerDeps) (func(), error) {redisOpt, _ := asynq.ParseRedisURI(redisURL)srv := asynq.NewServer(redisOpt, asynq.Config{Concurrency: 10,Queues: map[string]int{"default": 6,"critical": 3,"low": 1,},})mux := asynq.NewServeMux()mux.HandleFunc(TypeEmailSend, handleEmailSend(deps))mux.HandleFunc(TypeImageProcess, handleImageProcess(deps))mux.HandleFunc(TypeTokensCleanup, handleTokensCleanup(deps))go func() {srv.Run(mux)}()return func() { srv.Shutdown() }, nil}
Starting the Worker in main.go
// Start the background workerstopWorker, err := jobs.StartWorker(cfg.RedisURL, jobs.WorkerDeps{DB: db,Mailer: mailer,Storage: store,Cache: cacheService,})if err != nil {log.Fatalf("Failed to start worker: %v", err)}defer stopWorker()
Queue Priorities
Jobs are distributed across three priority queues. The worker allocates processing capacity based on these weights: critical gets 30%, default gets 60%, and low gets 10%.
| Queue | Weight | Use Case |
|---|---|---|
| critical | 3 (30%) | Password resets, payment webhooks |
| default | 6 (60%) | Emails, image processing |
| low | 1 (10%) | Cleanup, analytics, reports |
Retry Configuration
Each job type has a configured maximum retry count. When a handler returns an error, asynq automatically retries with exponential backoff. After exhausting all retries, the job is moved to the "archived" (failed) state.
// c.Enqueue(ctx, taskType, payload, ...EnqueueOption) wraps asynq and// applies framework defaults; override per call with an EnqueueOption.// Email: 3 attempts (important to deliver)c.Enqueue(ctx, TypeEmailSend, payload, EnqueueOption{MaxRetries: 3})// Image: 2 attempts (can be re-triggered)c.Enqueue(ctx, TypeImageProcess, payload, EnqueueOption{MaxRetries: 2})// Cleanup: 1 attempt (runs hourly anyway)c.Enqueue(ctx, TypeTokensCleanup, nil, EnqueueOption{MaxRetries: 1})// Custom: specific queue, timeout, and idempotency keyc.Enqueue(ctx, "invoice:generate", payload, EnqueueOption{MaxRetries: 5,Queue: "critical", // default | critical | lowTimeout: 30 * time.Second,IdempotencyKey: "invoice-" + orderID, // dedupes re-enqueues})
Adding Custom Jobs
Generate it: grit generate job InvoiceGenerate writes internal/jobs/invoice_generate.go with the task type, a payload struct, a typed enqueue method and the handler you fill in, and registers the handler with the worker. Add --cron "30 23 * * *" to run it on a schedule as well. A job whose handler is not registered is accepted by Redis and then fails on the worker, far from the line that caused it, which is why the command does both. What it writes is below, for doing it by hand.
// 1. Add task type constantconst TypeInvoiceGenerate = "invoice:generate"// 2. Define payloadtype InvoicePayload struct {OrderID string `json:"order_id"` // Order UUIDUserEmail string `json:"user_email"`}// 3. Add enqueue method to Client (Enqueue marshals the payload for you)func (c *Client) EnqueueGenerateInvoice(ctx context.Context, orderID, email string) error {return c.Enqueue(ctx, TypeInvoiceGenerate, InvoicePayload{OrderID: orderID,UserEmail: email,}, EnqueueOption{MaxRetries: 3})}// 4. Write handler functionfunc handleInvoiceGenerate(deps WorkerDeps) func(ctx context.Context, task *asynq.Task) error {return func(ctx context.Context, task *asynq.Task) error {var payload InvoicePayloadif err := json.Unmarshal(task.Payload(), &payload); err != nil {return err}// Generate PDF, send email, etc.return nil}}// 5. Register in worker mux (in StartWorker)mux.HandleFunc(TypeInvoiceGenerate, handleInvoiceGenerate(deps))
Batches, and running each piece exactly once
Large work fans out: a nightly job queues one task per account rather than doing every account in one long task that a single failure restarts from the top. The risk with fan-out is running a piece twice, when the scheduler fires on two replicas or somebody retries the batch by hand. Give every task an IdempotencyKey that names the piece of work, and the second enqueue of the same key is refused with jobs.ErrDuplicateTask instead of running.
// The nightly task, scheduled with --cron, fans out one task per account.for _, id := range accountIDs {err := client.EnqueueReconcileAccount(ctx, ReconcileAccountPayload{AccountID: id, Day: day},EnqueueOption{IdempotencyKey: "reconcile:" + id + ":" + day})if errors.Is(err, ErrDuplicateTask) {continue // already queued for this account and day}if err != nil {return err // retried with backoff; the keys make the retry safe}}
The key blocks re-enqueues for 24 hours by default; set Window to change it. It guards the queue, not your table, so a handler that writes should still write idempotently: an upsert keyed on the same account and day costs nothing and covers the case where a task ran and then failed before it could report success.
Admin Jobs Dashboard
The admin panel includes a jobs dashboard that shows queue statistics and allows admins to view, retry, and clear jobs. The dashboard uses the asynq Inspector API under the hood.
| Endpoint | Method | Description |
|---|---|---|
| /api/admin/jobs/stats | GET | Queue stats (pending, active, completed, failed) |
| /api/admin/jobs/:status | GET | List jobs by status (active, pending, completed, failed, retry) |
| /api/admin/jobs/:id/retry | POST | Retry a failed job |
| /api/admin/jobs/queue/:queue | DELETE | Clear all completed tasks in a queue |
Concurrency: The default worker concurrency is 10. This means up to 10 jobs can be processed simultaneously. Adjust the Concurrency setting in StartWorker() based on your server resources.
