{"_id":"@colu-legacy/osseus-mq","_rev":"1-5a316cea2e404e4171cc28a0b720fb24","name":"@colu-legacy/osseus-mq","dist-tags":{"latest":"0.3.17"},"versions":{"0.3.17":{"name":"@colu-legacy/osseus-mq","version":"0.3.17","description":"Osseus pub/sub, bus, queue and more","scripts":{"build":"npm install","test":"standard"},"main":"index.js","repository":{"type":"git","url":"git+https://github.com/colucom/osseus-mq.git"},"author":{"name":"LiorRabin"},"license":"MIT","devDependencies":{"standard":"^11.0.1"},"dependencies":{"aws-sdk":"^2.825.0","mongoose":"^4.13.21","mongoose-time":"^0.1.0","rhea":"^0.3.9","sprintf-js":"^1.1.2","util":"^0.11.0","uuid":"^3.4.0"},"gitHead":"6243aff512b1e9b81539b409daf085ae30eccaa1","bugs":{"url":"https://github.com/colucom/osseus-mq/issues"},"homepage":"https://github.com/colucom/osseus-mq#readme","_id":"@colu-legacy/osseus-mq@0.3.17","_nodeVersion":"14.15.1","_npmVersion":"7.5.3","dist":{"integrity":"sha512-xNO7cllcR9MUM+CNKFm9V8mYxZiuMtFesoxs+44+6R5qfb0H+1Mjz/vIize/KPtvwSkbzLzgDE8xPr6YjGtlgw==","shasum":"b57f87d2dc6a46e84e3e60870b19249d03b33e8e","tarball":"https://registry.npmjs.org/@colu-legacy/osseus-mq/-/osseus-mq-0.3.17.tgz","fileCount":10,"unpackedSize":47429,"npm-signature":"-----BEGIN PGP SIGNATURE-----\r\nVersion: OpenPGP.js v3.0.13\r\nComment: https://openpgpjs.org\r\n\r\nwsFcBAEBCAAQBQJg4XhFCRA9TVsSAnZWagAAHb4P/RcbyN57Ow73eGWYEXxU\nNi6XYAg79MJEzpuAUWliR4PnI1oc14o2tEdzLeQLu8WuV/liokvNwlfFjG5y\n1fUJ2tjeNHvYfBlFyE3I99zWgUZC8DdL1KSmW72SlcH1epbfMb4A+HRM9k4g\nwyIHoIEsIt7JfyVASPK2d/vgcDCB4mWyFMY1NARAJNY875sNFogW9YRWOEbp\nZMs91dg3/WPOzqKrU1zrvY++uM2jH/Wo1XVOr83mijGgQFYI+lJeY4NNOOJ3\n9FTIEFA1OF1opxfFlaZ/TdyHRAh4IXBbfaRtAfqOD2UqRfjvAan7X5G6URap\n6U5JszkQK1TZY+hD5lb13Y3FJMnr5gLPb9dJ7+fsuGfh1JqRxTs35rIGApdD\nf6A7EmTJyzqn0tXf0CdN5LK/JdNPLJS8euy0i4+ugOl6jGOiqRhKfEvEo43m\n9O+cbixCKcsB0Qzd6/wFzHvNt3FjYxOFLDcAIt8zCKe+nnVXUsKnO/IoGFtu\ntWbnccvfyv4cwslwuRP78sj1vHnVZV4zEmbCcXn0XQpt/6byuIalLqHdcmNF\n5lytHqm5lfgdQlo3FPoM3aA0CiYFnZlYVSuKS40YHSDpg39cIW+R2T8CU8ru\nHSFkrshu6Thrjet5z0liug+29N42uTiWVbGLFKxS5G97kTtNKrR1p31MapTM\nBsMt\r\n=ar9y\r\n-----END PGP SIGNATURE-----\r\n","signatures":[{"keyid":"SHA256:jl3bwswu80PjjokCgh0o2w5c2U4LhQAE57gj9cz1kzA","sig":"MEUCIQDUgnVSze7er1MsELxlJz+htLDkfr2VHlFpr5DbbDEQCwIgYUoULa3+4CyTw2S1fcA862Aj+v9gxMAK96M7aq+9vFQ="}]},"_npmUser":{"name":"romanplt87","email":"roman@colu.com"},"directories":{},"maintainers":[{"name":"romanplt87","email":"roman@colu.com"},{"name":"valery-colu","email":"valery@colu.com"}],"_npmOperationalInternal":{"host":"s3://npm-registry-packages","tmp":"tmp/osseus-mq_0.3.17_1625389124926_0.7752892087256711"},"_hasShrinkwrap":false}},"time":{"created":"2021-07-04T08:58:44.556Z","0.3.17":"2021-07-04T08:58:45.042Z","modified":"2022-04-05T00:23:15.659Z"},"maintainers":[{"name":"romanplt87","email":"roman@colu.com"},{"name":"valery-colu","email":"valery@colu.com"}],"description":"Osseus pub/sub, bus, queue and more","homepage":"https://github.com/colucom/osseus-mq#readme","repository":{"type":"git","url":"git+https://github.com/colucom/osseus-mq.git"},"author":{"name":"LiorRabin"},"bugs":{"url":"https://github.com/colucom/osseus-mq/issues"},"license":"MIT","readme":"[![JavaScript Style Guide](https://cdn.rawgit.com/standard/standard/master/badge.svg)](https://github.com/standard/standard)\n\n# Osseus MQ\n\n[AmazonMQ / ActiveMQ](http://activemq.apache.org/) based osseus queue, topic and bus module\n\n* WARNING current version 0.1.0 is in alpha stage.\n\n## Install\n```bash\n$ npm install osseus-mq\n```\n\n## Usage\n\n#### Configuration\nSee [osseus-config](https://github.com/colucom/osseus-config) for the configuration details.\n\n\n* `OSSEUS_MQ_PROTOCOL` (string, required) - 'AMQP' for AMQP or 'STOMP' for STOMP\n* `OSSEUS_MQ_AMQP_BROKERS` (array, required) - array of brokers\n* `OSSEUS_MQ_AMQP_CONSUMERS_<broker alias>` (array, required) - array of consumers of a certain broker\n* `OSSEUS_MQ_AMQP_PRODUCERS_<broker alias>` (array, required) - array of producers to a certain broker\n\n* Example:\n\n```javascript\n\tOSSEUS_MQ_AMQP_BROKERS: [\n            {\n              'alias': 'cluster',\n              'connect_options': {\n                'username': 'admin',\n                'password': 'admin',\n                'transport': 'ssl', //ONLY for amqps, remove otherwise\n                'connections': [\n                  {'host': 'localhost', 'port': 5672},\n                  {'host': 'localhost', 'port': 5673}\n                ]\n              }\n            }\n        ],\n        OSSEUS_MQ_AMQP_CONSUMERS_CLUSTER: [\n            {\n              name: 'queue://SYSTEM.MESSAGES',\n              alias: 'smsg',\n              options: ['settle'] // settle option for queue:// will redlive untill message is accepted by a client.\n            }\n        ],\n        OSSEUS_MQ_AMQP_PRODUCERS_CLUSTER: [\n            {\n              name: 'queue://SYSTEM.MESSAGES',\n              alias: 'smsg',\n              options: ['settle'] // settle option for queue:// will redlive untill message is accepted by a client.\n            }\n        ]\n```\n##### Credit\nTo limit the number of messages handled at once use credit. With credit only a limited number of messages will be handled and the rest will wait at the broker. Upon each accepted/rejected message a new one would dequeue. \n\n* `OSSEUS_MQ_HANDLE_CREDIT` (boolean) - Wheather to handle credit manually or not.\n* `OSSEUS_MQ_AMQP_INITIAL_CREDIT` (number) - How much credit to start with. If set to `0` no message will be dequeue. Only relevant when **OSSEUS_MQ_HANLDE_CREDIT** is set to `true`. \n* `OSSEUS_MQ_AMQP_CREDIT_LIMIT` (number) - Max credit the queue can reach. Only relevant when **OSSEUS_MQ_HANLDE_CREDIT** is set to `true`. \n\nThese parameters would be set for all queues but can be overriden for each one differently.\n\n* Example:\n\n```javascript\n\t\n        OSSEUS_MQ_AMQP_CONSUMERS_CLUSTER: [\n            {\n              name: 'queue://SYSTEM.MESSAGES',\n              alias: 'smsg',\n              options: ['settle'], // settle option for queue:// will redlive untill message is accepted by a client.\n              credit: {\n        \t\thandle: true,\n        \t\tinitial: 1,\n        \t\tlimit: 1\n      \t\t  }\n              }\n        ],\n        ...\n```\n\n\t\n* `OSSEUS_MQ_DB_USAGE` (string, required) - what database type to use if [send](#send) with persist flag set. Options are:\n\t* `MONGODB`\n\t* `POSTGRES`\n\t* `NONE`\n* `OSSEUS_MQ_DB_CONNECTION_STRING` (string, optional) - connection string for the database that will persist messages.\n\t* Required if `OSSEUS_MQ_DB_USAGE != \"NONE\"` \n* `APPLICATION_NAME` (string, required) - name of the app to use in logging and messaging.\n\n**note: wait for osseus.mq to emit `ready` event before usage**\n \n## Protocols\n### AMQP\n[AMQP](http://www.amqp.org/resources/download) is a specification (not a product) that was created to enable interoperability between multiple integration broker implementations.\nAMQP is a protocol and a fairly complete specification for the most commonly used middleware functionality.\n\n#### Topic / Queue names {#topicnames}\nTopics should be used when a message sould have multiple receivers, order is not an issue and there is no single action that is \"atomic\" that needs to be taken by the consumers.\n\nQueue should be used when only a single consumer should consume the message and report when it is finished with taking action on that message or return the message to the Queue for another consumer.\n\nQueue can have exclusive mode and both of them can have durable and save historical messages.\n\nIn order to use a queue topic name should start with `queue://` and to use a topic with `topic://`, both support wildcards `'*'` and both support complete will end `'>'`, so `topic://A.*.C` will match with wildcard whereas `topic://A.>` will match everything that starts with A. \n\n## Methods\nAll methods are accessable via the `mq` object which will be added to the base `osseus` object.\n\nThe `mq` object is an EventEmitter.\n\n### send (alias, message, options) ⇒ <code>Promise</code> {#send}\nMethod to send data to the named queue or topic, promise will be resolved once message is sent (and persisted if specified), promise will contain the base message object which will have a `requestId` if provided, and the system generated `messageId`\n \n**Kind**: function\n \n| Param             | Type                | Description                                                                     |\n| ------------------| ------------------- | ------------------------------------------------------------------------------- |\n| alias             | <code>string</code> | string alias for the topic / queue                                              |\n| message           | <code>object</code> | JSON to pass as the message (will be available in `<object>.msg`)               |\n| options           | <code>object</code> | options object                                                                  |\n| options.persist   | <code>'async'/'sync'</code> | persist message synchronously or asynchronously, if not set message not persisted |\n| options.requestId | <code>string</code>  | request id to add to the base message obejct                                   |\n| options.topic     | <code>string</code>  | set topic field on message for diffrent message types                          |\n\n### receive (alias, options, callback) ⇒ <code>no return value</code> {#receive}\nMethod to start getting notifications for incomming messages on the queue / topic. seta up a callback for the queue / topic, the callback is optional and user can instead use the `.on('message',(alias, message, done, failed)` method of the object to get the emitted events, though `receive` still needs to be called to set up a listener for the specific queue. queue must call `done` once it is done proccessing the message, calling `failed(string - error)` instead will set the message to error, not calling anything will redelive the messgae to this or any other consumer after the broker sees this instance has dissconnected.\n \n**Kind**: function\n \n| Param               | Type                                                | Description                                     |\n| --------------------| ----------------------------------------------------| ----------------------------------------------- |\n| alias               | <code>string</code>                                 | string alias for the topic / queue              |\n| options             | <code>object</code>                                 | options object                                  |\n| callback (optional) | <code>function(alais, message, done, failed)</code> | called when a message arrives to the queue/topic|\n\n## Statistics\n`osseus-mq` allows getting statistics from queues. Getting statistics is done by sending a message to the broker and waiting for it to respond. The number of messages in the response is as the number of queues requested.\nSince there is no way of telling when the broker is finished answering, gathering the statistics is divided into two parts.\nThe first one is for sending the request to the broker to get statistics for a destination and the second one is to get the last updated statistics received from the broker for the same destination.\n\n### sendDestinationStatisticsRequest (destination) ⇒ <code>no return value</code> \n| Param               | Type                                                | Description                                     |\n| --------------------| ----------------------------------------------------| ----------------------------------------------- |\n| destination               | <code>string</code>                                 | queue              |\n\n### retrieveStatisticsResponse (destination) ⇒ <code>no return value</code> \n| Param               | Type                                                | Description                                     |\n| --------------------| ----------------------------------------------------| ----------------------------------------------- |\n| destination               | <code>string</code>                                 | queue              |\n\nIt's best to wait for a few seconds between the steps to give the broker time to respond. \n\nYou can also use wildcards in the destination. \nExamples:\n\n`TEST.FOO` - Get statistics for `TEST.FOO`\n\n`TEST.>` - Get statistics for all queues start with `TEST`\n\n`>` - Get statistics for all queues.\n\n## Known Issues:\n#### consumer:\n- If setting a retain mode ([ActiveMQ retain](http://activemq.apache.org/subscription-recovery-policy.html)) then messages are expired from a topic based on ttl as well as retain method.\n\n\n#### producer:\n- Topics do not receive expire ([adviosry messages](http://activemq.apache.org/advisory-message.html)) even if ttl expires most of the time, even if there was no consumer or a consumer, should maybe check with delivery count.\n- Queues are automatically retaining all messages based on ttl, while topics fire and forget\n- Settle on disposition don't work, broker sends ack frame always when the message is stored in the broker (both for queues and topics)\n\n\n#### Advisory messages:\n\n##### NoConsumer:\n###### Topic :\n- Works. Message is sent only if publisher connects to a topic and there are no consumers at that time when a consumer goes up or a consumer goes down. Consumer advisory message is sent with `applicationProperties: { consumerCount: int }`\n\n###### Queue :\n- Message is never sent.\n\n##### Consumer:\n###### Topic :\n- Works. Message is sent when a consumer count changes `applicationProperties: { consumerCount: int }`, not if there are 0 consumers when producer starts\n\n###### Queue :\n- Works. Message is sent when a consumer count changes `applicationProperties: { consumerCount: int }`, not if there are 0 consumers when producer starts\n\n\n## License\nCode released under the [MIT License](https://github.com/colucom/osseus-mq/blob/master/LICENSE).\n","readmeFilename":"README.md"}