| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175 |
- 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;
|