Skip to content

Commit 481e560

Browse files
redboltzmcollina
authored andcommitted
Fixed #952. (#953)
* Fixed #952. Conditional flush `outgoing` on close. Scenario: 1. The client connect to the server. 2. The client sends subscribe to the server. 3. The server destroys the client connection before suback sending. 4. The client detect `close` event, then reconnects to the server. At the step4, `outgoing` still stored the callback for subscribe. However, it has never called because server doen't send corresponding suback. The same thing happens on unsubscribe. So I defined subscribe/unsubscribe as volatile. The volatile type of `outgoing` entries should be cleared when `close` from the server is detected. On the contrary, QoS1 and QoS2 publish is not volatile. Because they are resent after reconnection. And then, callback in the `store` is called. This behavior shouldn't be changed. So I added `volatile` flag to `outgoing` element. * Fixed cb accessing code. If `outgoing[mid]` doesn't match, then accessing `outgoing[mid].cb` causes `Cannot read property` error. So added checking code. Fixed outgoing assignment at `_onConnect` function.
1 parent 84ca344 commit 481e560

2 files changed

Lines changed: 96 additions & 21 deletions

File tree

‎lib/client.js‎

Lines changed: 45 additions & 21 deletions
Original file line numberDiff line numberDiff line change
@@ -90,8 +90,19 @@ function sendPacket (client, packet, cb) {
9090
function flush (queue) {
9191
if (queue) {
9292
Object.keys(queue).forEach(function (messageId) {
93-
if (typeof queue[messageId] === 'function') {
94-
queue[messageId](new Error('Connection closed'))
93+
if (typeof queue[messageId].cb === 'function') {
94+
queue[messageId].cb(new Error('Connection closed'))
95+
delete queue[messageId]
96+
}
97+
})
98+
}
99+
}
100+
101+
function flushVolatile (queue) {
102+
if (queue) {
103+
Object.keys(queue).forEach(function (messageId) {
104+
if (queue[messageId].volatile && typeof queue[messageId].cb === 'function') {
105+
queue[messageId].cb(new Error('Connection closed'))
95106
delete queue[messageId]
96107
}
97108
})
@@ -290,6 +301,7 @@ MqttClient.prototype._setupStream = function () {
290301

291302
// Echo stream close
292303
this.stream.on('close', function () {
304+
flushVolatile(that.outgoing)
293305
that.emit('close')
294306
})
295307

@@ -447,7 +459,10 @@ MqttClient.prototype.publish = function (topic, message, opts, callback) {
447459
case 1:
448460
case 2:
449461
// Add to callbacks
450-
this.outgoing[packet.messageId] = callback || nop
462+
this.outgoing[packet.messageId] = {
463+
volatile: false,
464+
cb: callback || nop
465+
}
451466
if (this._storeProcessing) {
452467
this._packetIdsDuringStoreProcessing[packet.messageId] = false
453468
this._storePacket(packet, undefined, opts.cbStorePut)
@@ -606,15 +621,18 @@ MqttClient.prototype.subscribe = function () {
606621
that.messageIdToTopic[packet.messageId] = topics
607622
}
608623

609-
this.outgoing[packet.messageId] = function (err, packet) {
610-
if (!err) {
611-
var granted = packet.granted
612-
for (var i = 0; i < granted.length; i += 1) {
613-
subs[i].qos = granted[i]
624+
this.outgoing[packet.messageId] = {
625+
volatile: true,
626+
cb: function (err, packet) {
627+
if (!err) {
628+
var granted = packet.granted
629+
for (var i = 0; i < granted.length; i += 1) {
630+
subs[i].qos = granted[i]
631+
}
614632
}
615-
}
616633

617-
callback(err, subs)
634+
callback(err, subs)
635+
}
618636
}
619637

620638
this._sendPacket(packet)
@@ -678,7 +696,10 @@ MqttClient.prototype.unsubscribe = function () {
678696
packet.properties = opts.properties
679697
}
680698

681-
this.outgoing[packet.messageId] = callback
699+
this.outgoing[packet.messageId] = {
700+
volatile: true,
701+
cb: callback
702+
}
682703

683704
this._sendPacket(packet)
684705

@@ -772,7 +793,7 @@ MqttClient.prototype.end = function () {
772793
* @example client.removeOutgoingMessage(client.getLastMessageId());
773794
*/
774795
MqttClient.prototype.removeOutgoingMessage = function (mid) {
775-
var cb = this.outgoing[mid]
796+
var cb = this.outgoing[mid] ? this.outgoing[mid].cb : null
776797
delete this.outgoing[mid]
777798
this.outgoingStore.del({messageId: mid}, function () {
778799
cb(new Error('Message removed'))
@@ -957,7 +978,7 @@ MqttClient.prototype._storePacket = function (packet, cb, cbStorePut) {
957978
if (((packet.qos || 0) === 0 && this.queueQoSZero) || packet.cmd !== 'publish') {
958979
this.queue.push({ packet: packet, cb: cb })
959980
} else if (packet.qos > 0) {
960-
cb = this.outgoing[packet.messageId]
981+
cb = this.outgoing[packet.messageId] ? this.outgoing[packet.messageId].cb : null
961982
this.outgoingStore.put(packet, function (err) {
962983
if (err) {
963984
return cb && cb(err)
@@ -1172,7 +1193,7 @@ MqttClient.prototype._handleAck = function (packet) {
11721193
var mid = packet.messageId
11731194
var type = packet.cmd
11741195
var response = null
1175-
var cb = this.outgoing[mid]
1196+
var cb = this.outgoing[mid] ? this.outgoing[mid].cb : null
11761197
var that = this
11771198
var err
11781199

@@ -1395,14 +1416,17 @@ MqttClient.prototype._onConnect = function (packet) {
13951416

13961417
// Avoid unnecessary stream read operations when disconnected
13971418
if (!that.disconnecting && !that.reconnectTimer) {
1398-
cb = that.outgoing[packet.messageId]
1399-
that.outgoing[packet.messageId] = function (err, status) {
1400-
// Ensure that the original callback passed in to publish gets invoked
1401-
if (cb) {
1402-
cb(err, status)
1419+
cb = that.outgoing[packet.messageId] ? that.outgoing[packet.messageId].cb : null
1420+
that.outgoing[packet.messageId] = {
1421+
volatile: false,
1422+
cb: function (err, status) {
1423+
// Ensure that the original callback passed in to publish gets invoked
1424+
if (cb) {
1425+
cb(err, status)
1426+
}
1427+
1428+
storeDeliver()
14031429
}
1404-
1405-
storeDeliver()
14061430
}
14071431
that._packetIdsDuringStoreProcessing[packet.messageId] = true
14081432
that._sendPacket(packet)

‎test/abstract_client.js‎

Lines changed: 51 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2682,6 +2682,57 @@ module.exports = function (server, config) {
26822682
})
26832683
})
26842684

2685+
it('should clear outgoing if close from server', function (done) {
2686+
var reconnect = false
2687+
var client = {}
2688+
var server2 = new Server(function (c) {
2689+
c.on('connect', function (packet) {
2690+
c.connack({returnCode: 0})
2691+
})
2692+
c.on('subscribe', function (packet) {
2693+
if (reconnect) {
2694+
c.suback({
2695+
messageId: packet.messageId,
2696+
granted: packet.subscriptions.map(function (e) {
2697+
return e.qos
2698+
})
2699+
})
2700+
} else {
2701+
c.destroy()
2702+
}
2703+
})
2704+
})
2705+
2706+
server2.listen(port + 50, function () {
2707+
client = mqtt.connect({
2708+
port: port + 50,
2709+
host: 'localhost',
2710+
clean: true,
2711+
clientId: 'cid1',
2712+
reconnectPeriod: 0
2713+
})
2714+
2715+
client.on('connect', function () {
2716+
client.subscribe('test', {qos: 2}, function (e) {
2717+
if (!e) {
2718+
client.end()
2719+
}
2720+
})
2721+
})
2722+
2723+
client.on('close', function () {
2724+
if (reconnect) {
2725+
server2.close()
2726+
done()
2727+
} else {
2728+
Object.keys(client.outgoing).length.should.equal(0)
2729+
reconnect = true
2730+
client.reconnect()
2731+
}
2732+
})
2733+
})
2734+
})
2735+
26852736
it('should resend in-flight QoS 1 publish messages from the client if clean is false', function (done) {
26862737
var reconnect = false
26872738
var client = {}

0 commit comments

Comments
 (0)