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:

HelperWhat 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 CmsShares job by setting leaseUntil (a 120 s lease) and incrementing attempts,
  • determines the delivery tenant from the job's subdomain (SaaS) or treats jobs as tenant-bound (self-hosted),
  • moves the record through a state machine (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

NeedMechanism
Fire-and-forget work off the request pathsendWorkerQueue('<plugin>', '<queue>').add(jobName, { subdomain, ... }) plus a matching createMQWorkerWithListeners consumer
Recurring workqueue.upsertJobScheduler(schedulerId, { pattern, tz }, { name }) on a dedicated scheduler queue, fanning out per-tenant jobs
Synchronous result from another servicesendWorkerMessage request/reply (3 s default timeout; keep the consumer fast)
Platform-driven workmeta.payments / meta.importExport handlers; the shared lib starts the workers
Long-running or externally paced deliveryPersisted 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-aiAgent workers. Plugins enqueue onto automations-* queues via sendWorkerQueue('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.

Was this helpful?