This document explains the background job processing system using BullMQ and Redis.
The application uses BullMQ for reliable background processing of trigger actions (Discord, Email, Telegram, Webhooks). This decouples the event poller from HTTP calls, improving reliability and scalability.
Event Poller → Queue (Redis) → Worker Pool → External Services
- Poller (
poller.js): Detects Soroban events and enqueues actions - Queue (
queue.js): Manages job storage and retrieval via Redis - Processor (
processor.js): Worker pool that executes actions with retries - Controller (
queue.controller.js): API endpoints for monitoring
- Guaranteed Delivery: Jobs are persisted in Redis and survive crashes
- Automatic Retries: Failed jobs retry with exponential backoff (3 attempts)
- Concurrency Control: Configurable worker pool size (default: 5)
- Rate Limiting: Built-in rate limiter (10 jobs/second)
- Job Tracking: Monitor job status via API endpoints
- Graceful Shutdown: Workers complete current jobs before stopping
macOS (Homebrew):
brew install redis
brew services start redisUbuntu/Debian:
sudo apt-get update
sudo apt-get install redis-server
sudo systemctl start redisWindows: Download from https://redis.io/download or use Docker:
docker run -d -p 6379:6379 redis:alpineVerify Redis:
redis-cli ping
# Should return: PONG- Install dependencies:
cd backend
npm install- Configure environment variables in
.env:
REDIS_HOST=localhost
REDIS_PORT=6379
REDIS_PASSWORD=
WORKER_CONCURRENCY=5- Start the server (worker starts automatically):
npm start| Variable | Default | Description |
|---|---|---|
REDIS_HOST |
localhost |
Redis server hostname |
REDIS_PORT |
6379 |
Redis server port |
REDIS_PASSWORD |
- | Redis password (optional) |
WORKER_CONCURRENCY |
5 |
Number of concurrent jobs |
Jobs are configured with:
- Attempts: 3 retries on failure
- Backoff: Exponential (2s, 4s, 8s)
- Retention: Completed jobs kept 24h, failed jobs kept 7 days
GET /api/queue/statsResponse:
{
"success": true,
"data": {
"waiting": 5,
"active": 2,
"completed": 150,
"failed": 3,
"delayed": 0,
"total": 160
}
}GET /api/queue/jobs?status=failed&limit=50Query Parameters:
status:waiting,active,completed,failed,delayedlimit: Number of jobs to return (default: 50)
POST /api/queue/cleanRemoves completed jobs older than 24 hours and failed jobs older than 7 days.
POST /api/queue/jobs/{jobId}/retryconst { enqueueAction } = require('./worker/queue');
// When an event is detected
const trigger = {
_id: 'trigger-id',
actionType: 'discord',
actionUrl: 'https://discord.com/api/webhooks/...',
contractId: 'CXXX...',
eventName: 'transfer',
};
const eventPayload = {
from: 'GXXX...',
to: 'GYYY...',
amount: '1000',
};
await enqueueAction(trigger, eventPayload);const { getQueueStats } = require('./worker/queue');
const stats = await getQueueStats();
console.log(`Active jobs: ${stats.active}`);
console.log(`Failed jobs: ${stats.failed}`);To add a web UI for queue monitoring:
- Install Bull Board:
npm install @bull-board/express @bull-board/api- Add to
server.js:
const { createBullBoard } = require('@bull-board/api');
const { BullMQAdapter } = require('@bull-board/api/bullMQAdapter');
const { ExpressAdapter } = require('@bull-board/express');
const { actionQueue } = require('./worker/queue');
const serverAdapter = new ExpressAdapter();
serverAdapter.setBasePath('/admin/queues');
createBullBoard({
queues: [new BullMQAdapter(actionQueue)],
serverAdapter,
});
app.use('/admin/queues', serverAdapter.getRouter());- Access at:
http://localhost:5000/admin/queues
Error: connect ECONNREFUSED 127.0.0.1:6379
Solution: Ensure Redis is running: redis-cli ping
Solution: Check worker logs. Worker may have crashed. Restart server.
Solution: Reduce WORKER_CONCURRENCY or clean old jobs more frequently.
Solution: Check logs for error details. Verify external service credentials (Discord webhook, SMTP, etc.).
- Redis Persistence: Enable AOF or RDB snapshots
- Redis Cluster: Use Redis Cluster for high availability
- Monitoring: Set up alerts for failed job count
- Scaling: Run multiple worker processes on different servers
- Security: Use Redis password and TLS in production
- Increase Concurrency: Higher
WORKER_CONCURRENCYfor more throughput - Adjust Rate Limits: Modify limiter in
processor.js - Job Priority: Set
prioritywhen enqueuing (lower = higher priority) - Batch Processing: Group similar jobs for efficiency
Worker logs include:
- Job enqueued
- Job processing started
- Job completed/failed
- Worker errors
Example:
[INFO] Action enqueued { jobId: 'trigger-123-1234567890', actionType: 'discord' }
[INFO] Processing action job { jobId: 'trigger-123-1234567890', attempt: 1 }
[INFO] Job completed { jobId: 'trigger-123-1234567890' }