diff --git a/README.md b/README.md index 19a12e0..287027d 100644 --- a/README.md +++ b/README.md @@ -35,7 +35,7 @@ options object. The jitter applied to the start of querying lookupd instances periodically. * ```lowRdyTimeout: 50```
The timeout in milliseconds for switching between connections when the Reader - maxInFlight is less than the number of connected NSQDs. + maxInFlight is less than the number of connected NSQDs. * ```tls: false```
Use TLS if nsqd has TLS support enabled. * ```tlsVerification: true```
@@ -158,8 +158,9 @@ These methods are available on a Writer object: * `publish(topic, msgs, [callback])`
`topic` is a string. `msgs` is either a string, a `Buffer`, JSON serializable object, a list of strings / `Buffers` / JSON serializable objects. `callback` takes a single `error` argument. -* `deferPublish(topic, msg, timeMs, [callback])`
- `topic` is a string. `msg` is either a string, a `Buffer`, JSON serializable object. `timeMs` is the delay by which the message should be delivered. `callback` takes a single `error` argument. +* `deferPublish(topic, msgs, timeMs, [callback])`
+ `topic` is a string. `msgs` is either a string, a `Buffer`, JSON serializable + object, a list of strings / `Buffers` / JSON serializable objects. `timeMs` is the delay by which the message should be delivered, and `-max-req-timeout` on NSQD must larger than `timeMS`. `callback` takes a single `error` argument. ### Simple example @@ -274,7 +275,7 @@ w.on('ready', () => { w.publish('sample_topic', 'it really tied the room together') w.deferPublish('sample_topic', ['This message gonna arrive 1 sec later.'], 1000) w.publish('sample_topic', [ - 'Uh, excuse me. Mark it zero. Next frame.', + 'Uh, excuse me. Mark it zero. Next frame.', 'Smokey, this is not \'Nam. This is bowling. There are rules.' ]) w.publish('sample_topic', 'Wu?', err => { @@ -356,7 +357,7 @@ w.on('closed', () => { * Bug: Non-fatal nsqd errors would cause RDY count to decrease and never return to normal. This will happen for example when finishing messages that have exceeded their amount of time to process a message. - * + * * **0.7.10** * Properly handles non-string errors * **0.7.9** diff --git a/lib/writer.js b/lib/writer.js index ba68399..193bc00 100644 --- a/lib/writer.js +++ b/lib/writer.js @@ -143,7 +143,7 @@ class Writer extends EventEmitter { * of the messages should either be strings or buffers with the payload encoded. * @param {String} topic - * @param {String|Buffer|Object} msg - A string, a buffer, a + * @param {String|Buffer|Object|Array} msg - A string, a buffer, a * JSON serializable object, or a list of string / buffers / * JSON serializable objects. * @param {Number} timeMs - defer time @@ -152,7 +152,7 @@ class Writer extends EventEmitter { */ deferPublish(topic, msg, timeMs, callback) { let err = this._checkStateValidity() - err = err || this._checkMsgsValidity(msg) + err = err || this._checkMsgsValidity(msgs) err = err || this._checkTimeMsValidity(timeMs) if (err) { @@ -163,12 +163,19 @@ class Writer extends EventEmitter { if (!this.ready) { const onReady = (err) => { if (err) return callback(err) - this.deferPublish(topic, msg, timeMs, callback) + this.deferPublish(topic, msgs, timeMs, callback) } this._callwhenReady(onReady) } - return this.conn.produceMessages(topic, msg, timeMs, callback) + if (!_.isArray(msgs)) { + msgs = [msgs] + } + + // Automatically serialize as JSON if the message isn't a String or a Buffer + msgs = msgs.map(this._serializeMsg) + + return this.conn.produceMessages(topic, msgs, timeMs, callback) } /** diff --git a/test/writer_test.js b/test/writer_test.js index 9dc7c38..77f8982 100644 --- a/test/writer_test.js +++ b/test/writer_test.js @@ -30,7 +30,7 @@ describe('writer', () => { const topic = 'test_topic' const msg = 'hello world!' - writer.publish(topic, msg, 300, () => { + writer.deferPublish(topic, msg, 300, () => { should.equal(writer.conn.produceMessages.calledOnce, true) should.equal(writer.conn.produceMessages.calledWith(topic, [msg]), true) })