var request = require('../utils/request') var MQTT_BROKER = 'wxs://iotcloud.tencent.com/socket/mqtt' var _socketTask = null var _connected = false var _subscriptions = {} var _listeners = {} var _reconnectTimer = null var _heartbeatTimer = null var _userId = null function on(event, callback) { if (!_listeners[event]) _listeners[event] = [] _listeners[event].push(callback) } function off(event, callback) { if (!_listeners[event]) return if (callback) { _listeners[event] = _listeners[event].filter(function (cb) { return cb !== callback }) } else { _listeners[event] = [] } } function emit(event, data) { if (!_listeners[event]) return _listeners[event].forEach(function (cb) { try { cb(data) } catch (e) { console.error('mqtt emit error:', e) } }) } function connect(userId) { _userId = userId var clientId = 'user_' + userId var token = wx.getStorageSync('token') _socketTask = wx.connectSocket({ url: MQTT_BROKER, header: { 'Authorization': 'Bearer ' + token }, success: function () { console.log('mqtt connecting...') }, fail: function () { emit('error', { msg: 'MQTT连接失败' }) scheduleReconnect() } }) _socketTask.onOpen(function () { _connected = true startHeartbeat() subscribeUserTopics() emit('connected') }) _socketTask.onMessage(function (res) { handleMessage(res.data) }) _socketTask.onClose(function () { _connected = false stopHeartbeat() emit('disconnected') scheduleReconnect() }) _socketTask.onError(function () { _connected = false emit('error', { msg: 'MQTT连接异常' }) }) } function subscribeUserTopics() { if (!_userId) return subscribe('users/' + _userId + '/subscription') subscribe('users/' + _userId + '/devices') subscribe('users/' + _userId + '/notification') } function subscribe(topic) { _subscriptions[topic] = true } function unsubscribe(topic) { delete _subscriptions[topic] } function handleMessage(data) { try { var msg = JSON.parse(data) var topic = msg.topic if (topic && _subscriptions[topic]) { if (topic.indexOf('/subscription') !== -1) { emit('subscription_change', msg.payload || msg) } else if (topic.indexOf('/devices') !== -1) { emit('devices_change', msg.payload || msg) } else if (topic.indexOf('/notification') !== -1) { emit('notification', msg.payload || msg) } else { emit('message', msg.payload || msg) } } } catch (e) { console.error('mqtt parse error:', e) } } function publish(topic, payload) { if (!_connected || !_socketTask) return _socketTask.send({ data: JSON.stringify({ topic: topic, payload: payload }) }) } function startHeartbeat() { stopHeartbeat() _heartbeatTimer = setInterval(function () { if (_connected && _socketTask) { _socketTask.send({ data: JSON.stringify({ type: 'ping' }) }) } }, 30000) } function stopHeartbeat() { if (_heartbeatTimer) { clearInterval(_heartbeatTimer) _heartbeatTimer = null } } function scheduleReconnect() { if (_reconnectTimer) return _reconnectTimer = setTimeout(function () { _reconnectTimer = null if (_userId) connect(_userId) }, 5000) } function disconnect() { if (_reconnectTimer) { clearTimeout(_reconnectTimer) _reconnectTimer = null } stopHeartbeat() if (_socketTask) { _socketTask.close({}) _socketTask = null } _connected = false _subscriptions = {} _listeners = {} _userId = null } function isConnected() { return _connected } module.exports = { connect: connect, disconnect: disconnect, subscribe: subscribe, unsubscribe: unsubscribe, publish: publish, on: on, off: off, isConnected: isConnected }