{"_rev":"3-ef6b1775130c67a8eca5ca086e2e2a8b","time":{"created":"2023-11-16T21:29:38.188Z","1.0.0":"2023-11-16T21:27:14.620Z","modified":"2023-11-16T21:29:38.932Z","1.0.1":"2023-11-16T21:29:38.632Z"},"_id":"@adwerx/rabbitmq-worker","name":"@adwerx/rabbitmq-worker","dist-tags":{"latest":"1.0.1"},"versions":{"1.0.1":{"name":"@adwerx/rabbitmq-worker","version":"1.0.1","description":"This library provides a high-level API for consuming messages from a rabbitmq exchange with sidekiq-like behavior.","type":"module","exports":{"types":"./dist/index.d.ts","default":"./dist/index.js"},"repository":{"url":"git+https://github.com/AdWerx/rabbitmq_worker_node.git"},"scripts":{"dev":"tsc --watch","build":"tsc","test":"NODE_OPTIONS=--experimental-vm-modules jest","prepublish":"npm run test && npm run build"},"engines":{"node":">=18"},"author":"","license":"ISC","devDependencies":{"@jest/globals":"^29.7.0","@types/amqplib":"^0.10","@types/debug":"^4.1.12","@types/jest":"^29.5.8","@types/node":"^20.9.0","jest":"^29.7.0","ts-jest":"^29.1.1","ts-node":"^10.9.1","ts-node-dev":"^2.0.0","typescript":"^5.2.2"},"peerDependencies":{"amqplib":"^0.10"},"dependencies":{"debug":"^4.3.4","p-limit":"^5.0.0"},"_id":"@adwerx/rabbitmq-worker@1.0.1","gitHead":"0b444f63d9fa8cfe7fe3a2c2840ca7b4ad7974a5","bugs":{"url":"https://github.com/AdWerx/rabbitmq_worker_node/issues"},"homepage":"https://github.com/AdWerx/rabbitmq_worker_node#readme","_nodeVersion":"21.1.0","_npmVersion":"10.2.0","dist":{"integrity":"sha512-XvjH4lxvwTxa3fPlzHdBpvC4njKJk9i1+tU88A/1YIgEpxcuR/p7JNX6y8/1/13bXTC39xjBvIIQ5qJ06DDVtA==","shasum":"32308b7f45e69473d0a8c0f1e974776d4d7cb1f8","tarball":"https://registry.npmjs.org/@adwerx/rabbitmq-worker/-/rabbitmq-worker-1.0.1.tgz","fileCount":14,"unpackedSize":70819,"signatures":[{"keyid":"SHA256:jl3bwswu80PjjokCgh0o2w5c2U4LhQAE57gj9cz1kzA","sig":"MEQCICSQgXthkJSalaohtVZrYMB8OhhG3MOJYq2KBPd+E8PSAiAuw960NFjdkeoKzdsviqI/OMVc1KVpptY43AvcwnYtQQ=="}]},"_npmUser":{"name":"jbielick","email":"jbielick@gmail.com"},"directories":{},"maintainers":[{"name":"jbielick","email":"jbielick@gmail.com"}],"_npmOperationalInternal":{"host":"s3://npm-registry-packages","tmp":"tmp/rabbitmq-worker_1.0.1_1700170178307_0.9371579574037854"},"_hasShrinkwrap":false}},"maintainers":[{"name":"jbielick","email":"jbielick@gmail.com"}],"description":"This library provides a high-level API for consuming messages from a rabbitmq exchange with sidekiq-like behavior.","homepage":"https://github.com/AdWerx/rabbitmq_worker_node#readme","repository":{"url":"git+https://github.com/AdWerx/rabbitmq_worker_node.git"},"bugs":{"url":"https://github.com/AdWerx/rabbitmq_worker_node/issues"},"license":"ISC","readme":"# rabbitmq-worker\n\nThis library provides a high-level API for consuming messages from a rabbitmq exchange with sidekiq-like behavior.\n\n## Installation\n\n```\nnpm install -S @adwerx/rabbitmq-worker amqplib\n```\n\n## Quickstart\n\nCreate a worker instance and use the configuration methods to indicate the exchange to bind to, a _perform function_ to process messages, an error handler, and provide a connection factory function.\n\nCall `start` with an `AbortSignal` to start consuming messages and use the abort signal to stop the worker.\n\n```js\nimport Worker from \"rabbitmq-worker\";\nimport amqplib from \"amqplib\";\n\nWorker.toConnect(() => amqplib.connect(process.env.AMQP_URL));\n\nconst controller = new AbortController();\nprocess.once(\"SIGINT\", () => controller.abort());\n\nawait new Worker()\n  .from(\"events\")\n  .perform(async (data, { fields, properties }) => {\n    // ... do work\n    // data is a *parsed* message\n  })\n  .on(\"error\", console.error(err))\n  .start(controller.signal);\n```\n\n## Configuration\n\n### `.from(exchangeName: string)`\n\nA worker must be configured with an _exchange_ to source messages from. Use `worker.from(exchangeName)` to set the exchange that the worker's queue will be bound to.\n\n### `.bind(queueName: string, routingKey = \"#\", args: StringIndexed = {})`\n\nBy default, starting a worker asserts a _queue_ specifically for this worker. The default queue name is the same as the exchange name. You can indicate the queue name to be used if desired by calling `worker.bind(queueName)`.\n\n### `.perform(fn)`\n\nA worker must be given a _perform function_ which will be called for each message retrieved. Use `worker.perform(fn)` to provide a _perform function_. The _perform function_ can be sync or async.\n\n### `.toConnect(connectionFactory: ConnectionFactoryFunction)`\n\nA worker must have an amqp _connection factory function_. You can provide this globally by setting a connection factory function to the `Worker.connect` static property. It is also possible to provide the connection factory to a single worker instance with `worker.toConnect(factoryFunc)`. The factory function must return (or resolve with) an `amqplib` connection instance. This function will be called any time the worker needs to reconnect to the server.\n\n### `.parse(contentType: ParserFunction | \"json\" | \"none\")`\n\nAll messages are assumed to be JSON data and will be parsed prior to being passed as an argument to your _perform function_. You can override this behavior by using `.parse(\"none\")` to disable parsing or by providing your own parsing function with `.parse((buffer: Buffer) => any)`.\n\n### `.with(options: PartialWorkerOptions = {})`\n\nYou can configure other worker behavior at the class level or instance level. For global default options, use `Worker.defaults = { ... }` and the provided settings will be defaults for all instances. For options specific to one instance, use `worker.with({ ... })` to set options for only that instance.\n\nWorkerOptions is a type that looks like this:\n\n```\ninterface WorkerOptions {\n  tag?: string;                  # the consumerTag the worker will use when connecting\n  concurrency: number;           # how many jobs may be in progress at any given moment\n  retries: number;               # how many times a job may retry before dying\n  parser: ParserFunction;        # a function that is used to parse a message before perform\n  deadJobRetensionMs: number;    # how long should dead jobs remain in the dead queue\n}\n```\n\n## Design\n\nThis library prioritizes a simple convention over configuration. A worker will assert the exchange it intends to bind a queue to, a work queue to hold its messages, a binding between these two, a retry queue to hold messages awaiting retry, a dead queue to hold messages that have exhausted retries, a requeue exchange, and a binding between the requeue exchange and the work queue to automatically requeue messages that have waited their retry period.\n\n### Example 1: Default behavior\n\n```js\nawait new Worker()\n  .from(\"events\")\n  .perform(() => {})\n  .start(signal);\n```\n\nThe worker defined above will create a work queue named `events`. Since no queue name was explicitly provided, a queue name of `events` is assumed. A queue to contain retryable messages is created named `events.retry`. A queue to hold dead messages is created named `events.dead`. An `events.requeue` exchange is created in order to re-queue messages that wait their retry period. The `events` exchange will be bound to the `events` work queue. The `events.requeue` exchange will be bound to the `events` work queue.\n\n### Example 2: Specified queue name\n\n```js\nawait new Worker()\n  .from(\"orders\")\n  .bind(\"provisioning\")\n  .perform(() => {})\n  .start(signal);\n```\n\nThe worker defined above will create a work queue named `provisioning`. A queue to contain retryable messages is created named `provisioning.retry`. A queue to hold dead messages is created named `provisioning.dead`. An `provisioning.requeue` exchange is created in order to re-queue messages that wait their retry period. The `orders` exchange will be bound to the `provisioning` work queue. The `provisioning.requeue` exchange will be bound to the `provisioning` work queue.\n\n### How retries work\n\nJobs retry a default of 25 times with an exponential backoff. When all retries are exhausted, the job is sent to the dead queue.\n\nWhen a message is delivered to the worker,\nand the perform function is successful,\nthe message is acknowledged.\n\nWhen a message is delivered to the worker,\nand an unexpected error occurs,\nand the message is unable to be sent to the retry or dead queue,\nthe message is _not_ acknowledged.\n\nWhen a message is delivered to the worker,\nand an unexpected error occurs,\nand the message is sent to the retry or dead queue,\nthe message is acknowledged.\n\nWhen a message is sent to the retry queue,\nthe message will have an expiration set.\n\nWhen a message in the retry queue expires,\nthe message is dead-lettered to the requeue exchange.\n\nWhen a message's deaths meet or exceed the retry limit,\nthe message is sent to the dead queue with an expiration.\n\nWhen a message in the dead queue expires,\nthe message is discarded.\n\n## Contributing\n\nRun tests with:\n`npm test`\n\nRun build with:\n`npm run build`\n","readmeFilename":"README.md"}