|
|
@@ -0,0 +1,175 @@
|
|
|
+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;
|