refactor: migrate to Tencent Cloud backend
这个提交包含在:
@@ -126,8 +126,8 @@ function buildFrame(type, payload) {
|
||||
if (payload && payload.length > 0) {
|
||||
frame = frame.concat(payload)
|
||||
}
|
||||
var checkBytes = frame.slice(2)
|
||||
frame.push(xorChecksum(checkBytes))
|
||||
var checkBytes = frame.slice(0)
|
||||
frame.push(xorChecksum(checkBytes))
|
||||
return bytesToBuffer(frame)
|
||||
}
|
||||
|
||||
@@ -140,7 +140,7 @@ function parseFrame(buffer) {
|
||||
var type = bytes[3]
|
||||
var payload = bytes.slice(4, 4 + len)
|
||||
var checksum = bytes[4 + len]
|
||||
var expected = xorChecksum(bytes.slice(2, 4 + len))
|
||||
var expected = xorChecksum(bytes.slice(0, 4 + len))
|
||||
if (checksum !== expected) return null
|
||||
return { type: type, payload: payload, seq: payload.length > 0 ? payload[payload.length - 1] : 0 }
|
||||
}
|
||||
|
||||
@@ -0,0 +1,68 @@
|
||||
var ble = require('./ble')
|
||||
var http = require('../utils/request')
|
||||
|
||||
var syncing = false
|
||||
|
||||
function commandPayload(command) {
|
||||
return command.payload || {}
|
||||
}
|
||||
|
||||
function execute(command) {
|
||||
var payload = commandPayload(command)
|
||||
switch (Number(command.opcode)) {
|
||||
case ble.CMD.SET_PARAMS:
|
||||
return ble.setParams({
|
||||
region_mask: payload.region_mask,
|
||||
wavelength: payload.wavelength,
|
||||
brightness: payload.brightness,
|
||||
duration_ms: payload.duration_ms,
|
||||
mode: payload.mode
|
||||
})
|
||||
case ble.CMD.START:
|
||||
return ble.startTreatment(payload.region_mask)
|
||||
case ble.CMD.STOP:
|
||||
return ble.stopTreatment()
|
||||
case ble.CMD.QUERY_STATUS:
|
||||
return ble.queryStatus()
|
||||
default:
|
||||
return Promise.reject({ error_code: 0xFE, error_msg: 'unsupported opcode' })
|
||||
}
|
||||
}
|
||||
|
||||
function report(command, success, result) {
|
||||
return http.post('/api/v1/device/command/result', {
|
||||
command_id: command.seq || command.command_id,
|
||||
seq: command.seq || command.command_id,
|
||||
success: success,
|
||||
opcode: command.opcode,
|
||||
result: result || null
|
||||
}).catch(function () {})
|
||||
}
|
||||
|
||||
function runOne(command) {
|
||||
return execute(command).then(function (result) {
|
||||
return report(command, true, result)
|
||||
}).catch(function (err) {
|
||||
return report(command, false, err)
|
||||
})
|
||||
}
|
||||
|
||||
function sync(deviceId) {
|
||||
if (syncing || !deviceId || !ble.isConnected()) return Promise.resolve()
|
||||
syncing = true
|
||||
return http.get('/api/v1/device/command/pending', { device_id: deviceId }).then(function (data) {
|
||||
var chain = Promise.resolve()
|
||||
var commands = data.commands || []
|
||||
commands.forEach(function (command) {
|
||||
chain = chain.then(function () { return runOne(command) })
|
||||
})
|
||||
return chain
|
||||
}).catch(function () {}).then(function () {
|
||||
syncing = false
|
||||
})
|
||||
}
|
||||
|
||||
module.exports = {
|
||||
sync: sync,
|
||||
execute: execute
|
||||
}
|
||||
@@ -1,172 +0,0 @@
|
||||
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
|
||||
}
|
||||
在新工单中引用
屏蔽一个用户