const _MQTT = uni.requireNativePlugin("zad-socket-mqtt"); class MQTTClient { constructor() { this.messageQueue = []; this.callbacks = new Map(); this.isConnecting = false; // 绑定方法到当前实例 this.handleMessage = this.handleMessage.bind(this); } /** * MQTT操作方法 * @param {string} method - 操作方法名 * @param {Object} param - 操作参数 * @returns {Promise} 返回Promise对象 */ async execute(method, param = {}) { // console.log(param,'11111111111111111111111'); return new Promise((resolve, reject) => { try { _MQTT.event({ method: method, param: param }, (e) => { this.handleMessage(e); if (e.code === 200) { resolve(e); } else if (e.code === 400) { reject(new Error(`MQTT ${method} 方法调用失败: ${JSON.stringify(e)}`)); } }); } catch (error) { reject(error); } }); } /** * 连接MQTT服务器 * @param {Object} config - 连接配置 * @param {string} config.url - 服务器地址 * @param {string} config.username - 用户名 * @param {string} config.password - 密码 * @param {string} config.clientId - 客户端ID * @returns {Promise} */ connect(config) { this.isConnecting = true; return this.execute('connect', config); } /** * 检查连接状态 * @returns {Promise} */ checkConnection() { return this.execute('isConnected'); } /** * 订阅主题 * @param {string} topic - 主题名称 * @returns {Promise} */ subscribe(topic) { return this.execute('subscribe', { topic }); } /** * 取消订阅 * @param {string} topic - 主题名称 * @returns {Promise} */ unsubscribe(topic) { return this.execute('unsubscribe', { topic }); } /** * 发送消息 * @param {string} topic - 主题名称 * @param {string|Object} message - 消息内容 * @returns {Promise} */ send(topic, message) { const messageStr = typeof message === 'object' ? JSON.stringify(message) : message; return this.execute('send', { topic, message: messageStr }); } /** * 关闭连接 * @returns {Promise} */ close() { this.isConnecting = false; return this.execute('close'); } /** * 处理接收到的消息 * @param {Object} e - 消息事件对象 */ handleMessage(e) { // 添加到消息队列 this.messageQueue.push(e); // 触发消息接收回调 if (e.method === 'receive') { this.triggerCallback('message', e.data); // console.log('接收到消息:', e.data); } // 处理连接状态变化 if (e.method === 'connect' && e.code === 200) { this.isConnecting = false; this.triggerCallback('connected', e); } // 输出日志 if (e.code === 400) { // console.error('MQTT操作失败:', e); } else if (e.code === 200) { // console.log('MQTT操作成功:', e.method); } } /** * 注册回调函数 * @param {string} event - 事件名称 * @param {Function} callback - 回调函数 */ on(event, callback) { this.callbacks.set(event, callback); } /** * 移除回调函数 * @param {string} event - 事件名称 */ off(event) { this.callbacks.delete(event); } /** * 触发回调函数 * @param {string} event - 事件名称 * @param {*} data - 回调数据 */ triggerCallback(event, data) { const callback = this.callbacks.get(event); if (callback && typeof callback === 'function') { callback(data); } } /** * 获取消息队列 * @returns {Array} */ getMessages() { return this.messageQueue; } /** * 清空消息队列 */ clearMessages() { this.messageQueue = []; } } // 创建全局实例 const mqttClient = new MQTTClient(); export default mqttClient;