Workers & Background Jobs
Run deferred and scheduled work on BullMQ over the shared Redis connection instead of blocking requests.
All queueing in plugins and erxes-api-shared is BullMQ on Redis; there is no RabbitMQ/amqp usage. For the standalone background services, see Background Services.
The shared helpers
erxes-api-shared/src/utils/mq-worker.ts wraps BullMQ:
| Helper | What it does |
|---|---|
sendWorkerQueue(service, queueName) | Returns a cached Queue named <service>-<queueName> on the shared redis connection. add() auto-serializes payloads (strips circular references) and applies DEFAULT_JOB_OPTIONS (removeOnComplete: true, removeOnFail: { count: 5000, age: 24h }). |
createMQWorkerWithListeners(service, queueName, processor, redis, onReady, workerOptions?) | Creates a Worker on the same <service>-<queueName> queue with completed/failed/error/ready logging; workerOptions carries concurrency, limiter, etc. |
sendWorkerMessage({ pluginName, queueName, jobName, subdomain, data, defaultValue, timeout, options }) | Request/reply: adds a job to <plugin>-<queue>, waits on QueueEvents for the result, times out after 3000 ms by default; keep the consumer fast for synchronous calls. |
The redis export (utils/redis.ts) is a shared ioredis client built from REDIS_HOST/REDIS_PORT/REDIS_PASSWORD. Reuse it; queue names must match between producer and consumer, so always go through sendWorkerQueue/createMQWorkerWithListeners with the same service and queueName.
The plugin pattern
Plugins keep an initMQWorkers(redis) in src/worker/index.ts and call it from onServerInit in startPlugin. operation_api is the active example:
export const initMQWorkers = async (redis: any) => {
const myQueue = new Queue('operations-daily-cycles-check', {
connection: redis,
defaultJobOptions: { removeOnComplete: true, removeOnFail: true },
});
// cron scheduler: BullMQ upsertJobScheduler
await myQueue.upsertJobScheduler(
'operations-daily-cycles-check',
{ pattern: '0 * * * *', tz: 'UTC' },
{ name: 'operations-daily-cycles-check' },
);
createMQWorkerWithListeners('operations', 'checkCycle', checkCycle, redis, onReady);
createMQWorkerWithListeners('operations', 'daily-cycles-check', dailyCheckCycles, redis, onReady);
};
// src/main.ts
startPlugin({
name: 'operation',
// …
onServerInit: async () => { await initMQWorkers(redis); },
});
The dispatcher job (scheduled) fans out per-tenant work onto a second queue: a job registered via upsertJobScheduler enqueues one job per organization via sendWorkerQueue. In posclient_api the initMQWorkers(redis) call happens in main.ts; in content_api the call exists but is commented out, so check the plugin's own main.ts before assuming its workers run.
Two meta keys also start BullMQ workers for you (see Plugin Metadata & Extensions): payments → <plugin>-payments and importExport → <plugin>-import-processor / <plugin>-export-processor (concurrency and limiter tunable via IMPORT_EXPORT_IMPORT|EXPORT_CONCURRENCY, *_LIMITER_MAX, *_LIMITER_DURATION_MS).
Tenant isolation and delivery state
A job payload carries { subdomain, data }; the worker regenerates models per job with generateModels(subdomain). Never reuse request-scoped models or cache them across jobs. sendWorkerMessage wraps the payload the same way.
content_api's Postiz delivery path (src/modules/cms/postiz/) is the fullest delivery-state example, and it does not use BullMQ at all; startCmsDeliveryWorker runs a polling sweep that:
- atomically claims a
CmsSharesjob by settingleaseUntil(a 120 s lease) and incrementingattempts, - determines the delivery tenant from the job's
subdomain(SaaS) or treats jobs as tenant-bound (self-hosted), - moves the record through a
statemachine (PENDING/QUEUED→ delivered,CANCELLED,UNKNOWN), - backs off with
nextCheck = now + min(attempts * 30000, 300000), capping at 20 attempts.
Follow that shape for new jobs: validate with Zod at the API boundary, persist a delivery record, claim work atomically, make retries idempotent, and show state/attempts in the UI (the PostizDeliveryList component does this for posts).
Picking the right mechanism
| Need | Mechanism |
|---|---|
| Fire-and-forget work off the request path | sendWorkerQueue('<plugin>', '<queue>').add(jobName, { subdomain, ... }) plus a matching createMQWorkerWithListeners consumer |
| Recurring work | queue.upsertJobScheduler(schedulerId, { pattern, tz }, { name }) on a dedicated scheduler queue, fanning out per-tenant jobs |
| Synchronous result from another service | sendWorkerMessage request/reply (3 s default timeout; keep the consumer fast) |
| Platform-driven work | meta.payments / meta.importExport handlers; the shared lib starts the workers |
| Long-running or externally paced delivery | Persisted delivery record + atomic lease + polling sweep (Postiz pattern), not a BullMQ Worker |
How ENABLED_SERVICES relates
Workers that live outside plugins run as separate Nx projects under backend/services/: automations-service and logs-service. pnpm dev:apis maps ENABLED_SERVICES=automations,logs to those project names and serves them next to your <name>_api projects. They run their own initMQWorkers:
automations-service:automations-trigger,automations-action,automations-aiAgentworkers. Plugins enqueue ontoautomations-*queues viasendWorkerQueue('automations', 'trigger')(the same name the service consumes).logs-service: event-log, activity-log, and segment workers, plus the after-process pipeline.
Cross-service queue naming is the rule: <service>-<queue> producer and consumer must agree, so call sendWorkerQueue with the consuming service's name, not your plugin's.