{"_id":"@appinventiv/bull-mq","_rev":"2-00f83255988d9dbaf5ebcc731ca71fd0","name":"@appinventiv/bull-mq","dist-tags":{"latest":"1.0.1"},"versions":{"1.0.0":{"name":"@appinventiv/bull-mq","version":"1.0.0","author":{"name":"Abhishek Tyagi"},"license":"ISC","_id":"@appinventiv/bull-mq@1.0.0","maintainers":[{"name":"developer-at","email":"abhishektyagi199816@gmail.com"},{"name":"abhishek.tyagi1","email":"abhishek.tyagi1@appinventiv.com"}],"dist":{"shasum":"89a4f5b968172badbfb7c28a5e05335296b0564f","tarball":"https://registry.npmjs.org/@appinventiv/bull-mq/-/bull-mq-1.0.0.tgz","fileCount":52,"integrity":"sha512-pWhffEwt4g1cPLCDgPj2k6ES38m7jM9+/x9DYRM4Sk/S5TsumCccO4/xAhZKz/k8BCwBTyLSF+KU4Ik56rBxUQ==","signatures":[{"sig":"MEUCIQCmTYD8nODHvXuneo7UsjU28Wsd3wdzKaMSlQJfTscTXgIgUYRK7YgIjJJgG26cCtzoYbu/QWroFeAF7c71NDoHd9U=","keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U"}],"unpackedSize":167878},"main":"dist/index.js","types":"dist/index.d.ts","scripts":{"test":"echo \"Error: no test specified\" && exit 1","build":"tsc"},"_npmUser":{"name":"abhishek.tyagi1","email":"abhishek.tyagi1@appinventiv.com"},"_npmVersion":"10.9.3","description":"BullMQ wrapper for Node.js with the same architectural style as `@appinventiv/rabbit-mq`: singleton manager, composable queue/worker/job layers, and Redis injected from your application.","directories":{},"_nodeVersion":"22.19.0","dependencies":{"bullmq":"^5.78.0","ioredis":"^5.9.3"},"_hasShrinkwrap":false,"devDependencies":{"ioredis":"^5.9.3","typescript":"^5.9.3","@types/node":"^25.9.2"},"peerDependencies":{"ioredis":"^5.0.0"},"_npmOperationalInternal":{"tmp":"tmp/bull-mq_1.0.0_1781774691032_0.04237745771362511","host":"s3://npm-registry-packages-npm-production"}},"1.0.1":{"name":"@appinventiv/bull-mq","version":"1.0.1","description":"BullMQ wrapper for Node.js with the same architectural style as `@appinventiv/rabbit-mq`: singleton manager, composable queue/worker/job layers, and Redis injected from your application.","main":"dist/index.js","types":"dist/index.d.ts","scripts":{"build":"tsc","test":"echo \"Error: no test specified\" && exit 1"},"author":{"name":"Abhishek Tyagi"},"license":"ISC","dependencies":{"bullmq":"^5.78.0","ioredis":"^5.9.3"},"peerDependencies":{"ioredis":"^5.0.0"},"devDependencies":{"@types/node":"^25.9.2","ioredis":"^5.9.3","typescript":"^5.9.3"},"_id":"@appinventiv/bull-mq@1.0.1","_nodeVersion":"22.19.0","_npmVersion":"10.9.3","dist":{"integrity":"sha512-wUMmklk3WhJp2kxk1WID0njFCjauzozDje7nn2Fnb851kMl4RrISsDKUuE6lv2qlDpC7KN99lH4sYpkcEZpKzA==","shasum":"e58f22c103a3202038d532aede27af8f5ade8ffc","tarball":"https://registry.npmjs.org/@appinventiv/bull-mq/-/bull-mq-1.0.1.tgz","fileCount":52,"unpackedSize":168641,"signatures":[{"keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U","sig":"MEQCIHuf4Gem07DipjdLObqSDiCuDHefXh6pCU3Xow5bhtN+AiBGSlASqiDLIhMhdlLVbRIC+jEQGl/y13KRdypu++yk7A=="}]},"_npmUser":{"name":"abhishek.tyagi1","email":"abhishek.tyagi1@appinventiv.com"},"directories":{},"maintainers":[{"name":"developer-at","email":"abhishektyagi199816@gmail.com"},{"name":"abhishek.tyagi1","email":"abhishek.tyagi1@appinventiv.com"}],"_npmOperationalInternal":{"host":"s3://npm-registry-packages-npm-production","tmp":"tmp/bull-mq_1.0.1_1781779105884_0.5678740098503823"},"_hasShrinkwrap":false}},"time":{"created":"2026-06-18T09:24:50.900Z","modified":"2026-06-18T10:38:26.171Z","1.0.0":"2026-06-18T09:24:51.217Z","1.0.1":"2026-06-18T10:38:26.035Z"},"author":{"name":"Abhishek Tyagi"},"license":"ISC","description":"BullMQ wrapper for Node.js with the same architectural style as `@appinventiv/rabbit-mq`: singleton manager, composable queue/worker/job layers, and Redis injected from your application.","maintainers":[{"name":"developer-at","email":"abhishektyagi199816@gmail.com"},{"name":"abhishek.tyagi1","email":"abhishek.tyagi1@appinventiv.com"}],"readme":"# @appinventiv/bull-mq\n\nBullMQ wrapper for Node.js with the same architectural style as `@appinventiv/rabbit-mq`: singleton manager, composable queue/worker/job layers, and Redis injected from your application.\n\n## Installation\n\n```bash\nnpm install @appinventiv/bull-mq ioredis\n```\n\n## Features\n\n- Default `bullMQ` singleton with optional custom `BullMQManager` instances\n- **Queue** — pooled queue wrappers per queue name\n- **Worker** — start processors with concurrency\n- **Job** — `add`, `addBulk`, and **scheduler** (`upsertJobScheduler`)\n- Built-in **Redis connection** module (`connection/redis`) with read/write replica support\n- Redis connection passed from outside or via `initFromRedisConfig`\n- TypeScript support\n\n## Prerequisites\n\n- Redis server accessible from your app\n- An `ioredis` client (or `RedisOptions`) you create and own\n\n## Usage\n\n### Basic setup (package Redis connection)\n\n```typescript\nimport { bullMQ, job, worker, queue } from '@appinventiv/bull-mq';\n\nbullMQ.initFromRedisConfig({\n  readReplicaHost: 'localhost',\n  readReplicaPort: 6379,\n  writeReplicaHost: 'localhost',\n  writeReplicaPort: 6379,\n  db: 0\n});\n\nbullMQ.setPrefix('{myapp}'); // optional\nbullMQ.setDefaultJobOptions({ attempts: 3 }); // optional\n\nawait job.add('notifications', 'send-email', { to: 'user@example.com' });\n\nawait worker.startWorker({\n  queueName: 'notifications',\n  processor: async (bullJob) => {\n    console.log('Processing', bullJob.data);\n  },\n  concurrency: 5\n});\n\nawait job.schedule(\n  'notifications',\n  'hourly-digest',\n  { every: 3600000 },\n  { name: 'digest', data: { type: 'hourly' } }\n);\n```\n\n### Basic setup (external ioredis client)\n\n```typescript\nimport { Redis } from 'ioredis';\nimport { bullMQ, job, worker, queue } from '@appinventiv/bull-mq';\n\nconst redis = new Redis({\n  host: 'localhost',\n  port: 6379,\n  maxRetriesPerRequest: null // required for BullMQ blocking workers\n});\n\nbullMQ.setConnection(redis);\nbullMQ.setPrefix('{myapp}'); // optional\nbullMQ.setDefaultJobOptions({ attempts: 3 }); // optional\n\n// Enqueue a job\nawait job.add('notifications', 'send-email', { to: 'user@example.com' });\n\n// Process jobs\nawait worker.startWorker({\n  queueName: 'notifications',\n  processor: async (bullJob) => {\n    console.log('Processing', bullJob.data);\n  },\n  concurrency: 5\n});\n\n// Repeatable / scheduled job\nawait job.schedule(\n  'notifications',\n  'hourly-digest',\n  { every: 3600000 },\n  { name: 'digest', data: { type: 'hourly' } }\n);\n```\n\n### Alternate Redis instance\n\n```typescript\nimport { BullMQManager, job } from '@appinventiv/bull-mq';\n\nconst analyticsRedis = new Redis({\n  host: 'analytics-redis',\n  port: 6379,\n  maxRetriesPerRequest: null\n});\nconst analyticsMq = new BullMQManager();\nanalyticsMq.setConnection(analyticsRedis);\n\nawait job.add('events', 'track', { event: 'click' }, undefined, analyticsMq);\n```\n\n## Scheduled and delayed jobs\n\nBullMQ supports two related patterns:\n\n| Pattern | Method | When to use |\n|---------|--------|-------------|\n| **One-off delayed job** | `job.add(..., { delay })` | Run once after N milliseconds |\n| **Repeating / cron job** | `job.schedule(...)` | Run on an interval or cron pattern (uses `upsertJobScheduler`) |\n\nA **worker must be running** on the same `queueName` for scheduled jobs to be processed.\n\n### One-off delayed job\n\n```typescript\n// Run once after 30 seconds\nawait job.add(\n  'notifications',\n  'send-reminder',\n  { userId: '42', message: 'Complete your profile' },\n  { delay: 30_000 }\n);\n```\n\n### Repeating job — fixed interval (`every`)\n\n```typescript\n// Every hour (milliseconds)\nawait job.schedule(\n  'notifications',           // queueName\n  'hourly-digest',           // schedulerId — unique id for this schedule\n  { every: 3_600_000 },      // repeat — run every 1 hour\n  {\n    name: 'digest',            // jobTemplate.name — job name each run uses\n    data: { type: 'hourly' }, // jobTemplate.data — payload passed to processor\n    opts: { attempts: 3 }      // jobTemplate.opts — per-run job options (optional)\n  }\n);\n```\n\n### Repeating job — cron pattern\n\n```typescript\n// Every day at midnight (server timezone / cron-parser rules)\nawait job.schedule(\n  'reports',\n  'daily-report',\n  { pattern: '0 0 * * *' },\n  {\n    name: 'generate-report',\n    data: { reportType: 'daily' },\n    opts: {\n      removeOnComplete: true,\n      removeOnFail: 50\n    }\n  }\n);\n```\n\n### Cron with immediate first run\n\n```typescript\nawait job.schedule(\n  'sync',\n  'sync-users',\n  {\n    pattern: '0 */6 * * *', // every 6 hours\n    immediately: true       // also enqueue one run right now\n  },\n  { name: 'sync-users', data: { source: 'crm' } }\n);\n```\n\n### Limit how many times a schedule runs\n\n```typescript\nawait job.schedule(\n  'onboarding',\n  'welcome-series',\n  { every: 86_400_000, limit: 7 }, // once per day, max 7 times\n  { name: 'welcome-email', data: { step: 1 } }\n);\n```\n\n### Remove a schedule\n\n```typescript\nawait job.removeSchedule('notifications', 'hourly-digest');\n```\n\n### Full example: schedule + worker\n\n```typescript\nimport { bullMQ, job, worker } from '@appinventiv/bull-mq';\n\nbullMQ.initFromRedisConfig({\n  readReplicaHost: 'localhost',\n  readReplicaPort: 6379,\n  writeReplicaHost: 'localhost',\n  writeReplicaPort: 6379,\n  db: 0\n});\n\n// 1. Register worker first (same queue name as schedule)\nawait worker.startWorker({\n  queueName: 'notifications',\n  processor: async (bullJob) => {\n    if (bullJob.name === 'digest') {\n      await sendDigestEmail(bullJob.data);\n    }\n  },\n  concurrency: 2\n});\n\n// 2. Register repeatable schedule\nawait job.schedule(\n  'notifications',\n  'hourly-digest',\n  { every: 3_600_000 },\n  { name: 'digest', data: { type: 'hourly' } }\n);\n```\n\n---\n\n## Job API parameters\n\n### `job.add(queueName, jobName, data, opts?, bullMq?)`\n\n| Parameter | Type | Required | Description |\n|-----------|------|----------|-------------|\n| `queueName` | `string` | Yes | BullMQ queue name. Jobs land here; worker must listen on the same name. |\n| `jobName` | `string` | Yes | Logical job name (e.g. `'send-email'`). Your processor can branch on `job.name`. |\n| `data` | `any` | Yes | JSON-serializable payload available as `job.data` in the processor. |\n| `opts` | `JobsOptions` | No | Per-job options (see below). |\n| `bullMq` | `BullMQManager` | No | Override Redis manager; defaults to `bullMQ` singleton. |\n\n**Common `opts` (one-off jobs):**\n\n| Option | Type | Description |\n|--------|------|-------------|\n| `delay` | `number` | Milliseconds to wait before the job becomes active. |\n| `attempts` | `number` | Max retry attempts (default `1`). |\n| `backoff` | `number \\| object` | Delay strategy between retries. |\n| `priority` | `number` | Lower number = higher priority (0 is highest). |\n| `jobId` | `string` | Custom job id; must be unique in the queue. |\n| `removeOnComplete` | `boolean \\| number \\| object` | Remove or cap completed jobs in Redis. |\n| `removeOnFail` | `boolean \\| number \\| object` | Remove or cap failed jobs in Redis. |\n| `lifo` | `boolean` | If `true`, add to front of queue (LIFO). |\n\n```typescript\nawait job.add('orders', 'process-order', { orderId: '99' }, {\n  delay: 5_000,\n  attempts: 5,\n  backoff: { type: 'exponential', delay: 1000 },\n  removeOnComplete: 100,\n  removeOnFail: 50\n});\n```\n\n### `job.schedule(queueName, schedulerId, repeat, jobTemplate?, bullMq?)`\n\nWraps BullMQ [`upsertJobScheduler`](https://docs.bullmq.io/guide/job-schedulers). Creates or updates a **repeatable** schedule; each tick enqueues a job from `jobTemplate`.\n\n| Parameter | Type | Required | Description |\n|-----------|------|----------|-------------|\n| `queueName` | `string` | Yes | Queue that receives scheduled job instances. |\n| `schedulerId` | `string` | Yes | Stable id for this schedule. Re-calling `schedule` with the same id **updates** the schedule. |\n| `repeat` | `RepeatOptions` | Yes | When to run (interval or cron). See table below. |\n| `jobTemplate` | `object` | No | Template for each generated job. |\n| `jobTemplate.name` | `string` | No | Job name for each run (processor sees `job.name`). |\n| `jobTemplate.data` | `any` | No | Payload for each run (`job.data`). |\n| `jobTemplate.opts` | `JobSchedulerTemplateOptions` | No | Per-run options (`attempts`, `removeOnComplete`, etc.). Cannot include `delay`, `repeat`, or `jobId`. |\n| `bullMq` | `BullMQManager` | No | Optional manager override. |\n\n**`repeat` options (`RepeatOptions`):**\n\n| Option | Type | Description |\n|--------|------|-------------|\n| `every` | `number` | Repeat every N **milliseconds**. Do not use with `pattern`. |\n| `pattern` | `string` | Cron expression (e.g. `'0 0 * * *'`). Do not use with `every`. |\n| `limit` | `number` | Stop after this many repetitions. |\n| `immediately` | `boolean` | With cron `pattern`, enqueue one run immediately. |\n| `count` | `number` | Start value for repeat iteration count. |\n| `offset` | `number` | Offset in ms affecting next iteration time. |\n| `tz` | `string` | Timezone for cron (via `cron-parser`; see BullMQ docs). |\n\n**Notes:**\n\n- Use **`every`** for simple intervals (e.g. every 5 minutes → `300_000`).\n- Use **`pattern`** for calendar-style schedules (cron).\n- **`schedulerId`** should be unique per schedule on a queue; use it with `job.removeSchedule()` to delete.\n- Changing `repeat` or `jobTemplate` and calling `schedule` again updates the existing scheduler.\n\n### `job.removeSchedule(queueName, schedulerId, bullMq?)`\n\nRemoves a previously registered scheduler. Does not remove jobs already in the queue.\n\n### `job.addBulk(queueName, jobs, bullMq?)`\n\n| `jobs[]` field | Description |\n|----------------|-------------|\n| `name` | Job name |\n| `data` | Payload |\n| `opts` | Optional `JobsOptions` per job |\n\n```typescript\nawait job.addBulk('notifications', [\n  { name: 'send-email', data: { to: 'a@example.com' } },\n  { name: 'send-email', data: { to: 'b@example.com' }, opts: { priority: 1 } }\n]);\n```\n\n### Shutdown\n\n`worker.stopWorkers()` and `queue.disconnectAll()` close BullMQ queue/worker instances.\n\nWhen using `initFromRedisConfig`, also disconnect the package Redis module:\n\n```typescript\nawait worker.stopWorkers();\nawait queue.disconnectAll();\nawait bullMQ.disconnect();\nawait redisConnection.disconnect();\n```\n\nWhen using an external ioredis client:\n\n```typescript\nawait worker.stopWorkers();\nawait queue.disconnectAll();\nawait bullMQ.disconnect();\nawait redis.quit();\n```\n\n## API overview\n\n| Export | Role |\n|--------|------|\n| `bullMQ` | Default manager — `initFromRedisConfig()` or `setConnection(redis)` |\n| `redisConnection` | Redis read/write clients (`getConnection()` → write client for BullMQ) |\n| `RedisConnection` | Custom Redis connection instance |\n| `BullMQManager` | Custom manager for a second Redis |\n| `queue` | Pooled `BullMQQueue` instances |\n| `worker` | Register workers |\n| `job` | Add / schedule jobs |\n| `BullMQQueue`, `BullMQWorker`, `BullMQJob` | Low-level wrappers |\n\n### BullMQManager\n\n- `initFromRedisConfig(config)` — create Redis via `redisConnection` and bind write client\n- `setConnection(connection)` — use an existing ioredis client\n- `setPrefix(prefix)` — optional BullMQ key prefix\n- `setDefaultJobOptions(opts)` — merged into queue defaults\n- `isConfigured()` — whether connection was set\n- `getQueueOptions()` / `getWorkerOptions()` — used internally by wrappers\n\n### RedisConnection (`redisConnection`)\n\n- `initiateRedisConnection(config)` — create read/write clients\n- `getWriteClient()` / `getConnection()` — write client for BullMQ (created with `maxRetriesPerRequest: null`)\n- `getReadClient()` — read replica client\n- `disconnect()` — quit both clients\n\n**Note:** The write client sets `maxRetriesPerRequest: null` automatically. BullMQ workers rely on blocking Redis commands; a finite retry limit causes ioredis to fail those waits. If you use `bullMQ.setConnection()` with your own ioredis instance, set `maxRetriesPerRequest: null` on that client too.\n\n### JobService (`job`)\n\n- `add(queueName, jobName, data, opts?, bullMq?)` — enqueue a one-off (or delayed) job\n- `addBulk(queueName, jobs, bullMq?)` — enqueue many jobs at once\n- `schedule(queueName, schedulerId, repeat, jobTemplate?, bullMq?)` — repeatable / cron schedule\n- `removeSchedule(queueName, schedulerId, bullMq?)` — remove a scheduler\n- `getJob(queueName, jobId, bullMq?)` — fetch a job by id\n\nSee [Scheduled and delayed jobs](#scheduled-and-delayed-jobs) and [Job API parameters](#job-api-parameters) for examples and field descriptions.\n\n### WorkerService (`worker`)\n\n- `startWorker({ queueName, processor, concurrency?, workerOptions?, bullMq? })`\n- `stopWorkers()` / `disconnectWorkers()`\n\n### QueueService (`queue`)\n\n- `getQueue(queueName, bullMq?)`\n- `disconnectAll()`\n\n## Comparison with rabbit-mq\n\n| rabbit-mq | bull-mq |\n|-----------|---------|\n| `rabbitMQ.setConfig({ url })` | `bullMQ.initFromRedisConfig(config)` or `bullMQ.setConnection(redis)` |\n| `producer.produce(queue, msg)` | `job.add(queueName, jobName, data)` |\n| `consumer.consume({ queue, onMessage })` | `worker.startWorker({ queueName, processor })` |\n| N/A | `job.schedule(...)` |\n\n## License\n\nISC\n","readmeFilename":"README.md"}