mqtt.js 4.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175
  1. const _MQTT = uni.requireNativePlugin("zad-socket-mqtt");
  2. class MQTTClient {
  3. constructor() {
  4. this.messageQueue = [];
  5. this.callbacks = new Map();
  6. this.isConnecting = false;
  7. // 绑定方法到当前实例
  8. this.handleMessage = this.handleMessage.bind(this);
  9. }
  10. /**
  11. * MQTT操作方法
  12. * @param {string} method - 操作方法名
  13. * @param {Object} param - 操作参数
  14. * @returns {Promise} 返回Promise对象
  15. */
  16. async execute(method, param = {}) {
  17. // console.log(param,'11111111111111111111111');
  18. return new Promise((resolve, reject) => {
  19. try {
  20. _MQTT.event({
  21. method: method,
  22. param: param
  23. }, (e) => {
  24. this.handleMessage(e);
  25. if (e.code === 200) {
  26. resolve(e);
  27. } else if (e.code === 400) {
  28. reject(new Error(`MQTT ${method} 方法调用失败: ${JSON.stringify(e)}`));
  29. }
  30. });
  31. } catch (error) {
  32. reject(error);
  33. }
  34. });
  35. }
  36. /**
  37. * 连接MQTT服务器
  38. * @param {Object} config - 连接配置
  39. * @param {string} config.url - 服务器地址
  40. * @param {string} config.username - 用户名
  41. * @param {string} config.password - 密码
  42. * @param {string} config.clientId - 客户端ID
  43. * @returns {Promise}
  44. */
  45. connect(config) {
  46. this.isConnecting = true;
  47. return this.execute('connect', config);
  48. }
  49. /**
  50. * 检查连接状态
  51. * @returns {Promise}
  52. */
  53. checkConnection() {
  54. return this.execute('isConnected');
  55. }
  56. /**
  57. * 订阅主题
  58. * @param {string} topic - 主题名称
  59. * @returns {Promise}
  60. */
  61. subscribe(topic) {
  62. return this.execute('subscribe', { topic });
  63. }
  64. /**
  65. * 取消订阅
  66. * @param {string} topic - 主题名称
  67. * @returns {Promise}
  68. */
  69. unsubscribe(topic) {
  70. return this.execute('unsubscribe', { topic });
  71. }
  72. /**
  73. * 发送消息
  74. * @param {string} topic - 主题名称
  75. * @param {string|Object} message - 消息内容
  76. * @returns {Promise}
  77. */
  78. send(topic, message) {
  79. const messageStr = typeof message === 'object' ? JSON.stringify(message) : message;
  80. return this.execute('send', { topic, message: messageStr });
  81. }
  82. /**
  83. * 关闭连接
  84. * @returns {Promise}
  85. */
  86. close() {
  87. this.isConnecting = false;
  88. return this.execute('close');
  89. }
  90. /**
  91. * 处理接收到的消息
  92. * @param {Object} e - 消息事件对象
  93. */
  94. handleMessage(e) {
  95. // 添加到消息队列
  96. this.messageQueue.push(e);
  97. // 触发消息接收回调
  98. if (e.method === 'receive') {
  99. this.triggerCallback('message', e.data);
  100. // console.log('接收到消息:', e.data);
  101. }
  102. // 处理连接状态变化
  103. if (e.method === 'connect' && e.code === 200) {
  104. this.isConnecting = false;
  105. this.triggerCallback('connected', e);
  106. }
  107. // 输出日志
  108. if (e.code === 400) {
  109. // console.error('MQTT操作失败:', e);
  110. } else if (e.code === 200) {
  111. // console.log('MQTT操作成功:', e.method);
  112. }
  113. }
  114. /**
  115. * 注册回调函数
  116. * @param {string} event - 事件名称
  117. * @param {Function} callback - 回调函数
  118. */
  119. on(event, callback) {
  120. this.callbacks.set(event, callback);
  121. }
  122. /**
  123. * 移除回调函数
  124. * @param {string} event - 事件名称
  125. */
  126. off(event) {
  127. this.callbacks.delete(event);
  128. }
  129. /**
  130. * 触发回调函数
  131. * @param {string} event - 事件名称
  132. * @param {*} data - 回调数据
  133. */
  134. triggerCallback(event, data) {
  135. const callback = this.callbacks.get(event);
  136. if (callback && typeof callback === 'function') {
  137. callback(data);
  138. }
  139. }
  140. /**
  141. * 获取消息队列
  142. * @returns {Array}
  143. */
  144. getMessages() {
  145. return this.messageQueue;
  146. }
  147. /**
  148. * 清空消息队列
  149. */
  150. clearMessages() {
  151. this.messageQueue = [];
  152. }
  153. }
  154. // 创建全局实例
  155. const mqttClient = new MQTTClient();
  156. export default mqttClient;