 | |  |  | MQTT协议示例AIWROK版
https://www.yuque.com/aiwork/nba2pr/gw970kaylsvz7eyl
- /**
- * MQTT 协议示例 - AIWROK 版
- *
- * 使用 Eclipse Paho Java 客户端实现 MQTT v3.1.1 连接
- *
- * v5 说明:dex 为官方 Paho 1.2.5 源码修补版(2026-08-29 重建):
- * 1. LoggerFactory 资源包查找加保护,默认使用 MockLogger(无操作日志)
- * → 修复 dex 环境下 ResourceBundle.getBundle 必抛 MissingResourceException
- * 2. NetworkModuleService 加内置协议工厂兜底(tcp/ssl/ws/wss)
- * → 修复 dex 无 META-INF/services 导致 "no NetworkModule installed for scheme tcp"
- * 3. 脚本不再需要任何运行时反射补丁
- * 支持:发布/订阅模式、QoS 0/1/2、自动重连、遗嘱消息、心跳保活
- *
- * MQTT vs WebSocket 区别:
- * MQTT 是轻量级发布/订阅消息协议,适合 IoT/移动端
- * 比 WebSocket 更省电、更省流量,支持 QoS 消息质量保证
- * 支持主题(Topic)订阅,消息路由更灵活
- *
- * 依赖加载(自动,无需手动操作):
- * dex 文件放在 H5HTML/插件/ 目录,AIWROK 的 loadDex 自动加载
- * loadDex 只从"插件"目录加载,且需要 .dex 格式(不是 .jar)
- *
- * dex 文件已随工程打包在:H5HTML/插件/paho-mqtt.dex
- * 如果插件目录没有 dex,自动从代码目录的 paho-mqtt-dex-b64.js 还原
- *
- * 使用方法:
- * 1. 修改 CFG 中的 broker 地址、端口、用户名密码
- * 2. AIWROK 中直接运行本脚本
- * 3. 日志窗口查看连接状态和收发消息
- */
- // ========== 加载 Paho MQTT Java 库 ==========
- // AIWROK loadDex 只从 H5HTML/插件/ 目录加载 .dex 文件
- // AIWROK IDE 不会同步插件目录的二进制文件到手机,
- // 脚本运行时从内嵌的 base64 数据自动还原 dex 文件到插件目录
- var loaded = false;
- // 手机端路径
- var sdcardPath = "/sdcard";
- var projectRoot = sdcardPath + "/auto/H5HTML";
- var pluginDir = projectRoot + "/插件";
- var codeDir = projectRoot + "/代码";
- var dexFileName = "paho-mqtt.dex";
- var dexTargetPath = pluginDir + "/" + dexFileName;
- // dex 版本判定:以内嵌 base64 解码后的长度为唯一基准
- // (这样以后重新打包 dex 不必改任何常量,脚本自己会比对并覆盖旧版)
- function fileSize(path) {
- try {
- var f = new java.io.File(path);
- return f.exists() ? Number(f.length()) : -1;
- } catch (e) { return -1; }
- }
- // 先取内嵌数据:require 失败只记日志,不让脚本崩
- var dexBytes = null;
- try {
- try { require("paho-mqtt-dex-b64.js"); } catch (reqErr) {
- printl("[WARN] require paho-mqtt-dex-b64.js 失败: " + reqErr.message);
- }
- if (typeof PAHO_MQTT_DEX_B64 !== "undefined" && PAHO_MQTT_DEX_B64) {
- dexBytes = java.util.Base64.getDecoder().decode(PAHO_MQTT_DEX_B64);
- printl("[INFO] 内嵌 dex 数据: " + dexBytes.length + " bytes");
- } else {
- printl("[WARN] 未取到内嵌 base64 数据(PAHO_MQTT_DEX_B64 为空)");
- }
- } catch (e) {
- printl("[WARN] 内嵌 dex 解码失败: " + e.message);
- dexBytes = null;
- }
- var expectedSize = dexBytes ? Number(dexBytes.length) : -1;
- var curSize = fileSize(dexTargetPath);
- var dexWasThere = curSize > 0;
- // 有内嵌数据就按字节数判版本;没有就"有 dex 先用着"
- var verOk = (expectedSize > 0) ? (curSize === expectedSize) : dexWasThere;
- printl("[INFO] 插件目录 dex: " + (dexWasThere ? ("已存在 " + curSize + " bytes") : "不存在") +
- "," + (verOk ? "版本校验通过" : ("需重装(期望 " + (expectedSize > 0 ? expectedSize : "未知") + " bytes)")));
- // 第一步:确保插件目录存在
- try { file.mkdir(pluginDir); } catch (e) {}
- // 第二步:版本不对就用内嵌数据覆盖还原
- var installed = false;
- if (!verOk && dexBytes) {
- try {
- var fos = new java.io.FileOutputStream(dexTargetPath);
- fos.write(dexBytes);
- fos.close();
- printl("[INFO] 已还原 " + dexFileName + " 到插件目录");
- installed = true;
- } catch (e) {
- printl("[WARN] 写入 dex 失败: " + e.message);
- }
- }
- // 最终判定能否加载:还原成功,或手机上本来就有 dex(退回使用现有文件)
- var dexExists = installed || dexWasThere;
- if (!installed && !verOk && dexWasThere) {
- printl("[WARN] 内嵌还原未完成,改用手机上现有的 dex 继续加载");
- }
- // 第三步:用 loadDex 加载
- if (dexExists) {
- try {
- rhino.loadDex(dexFileName);
- printl("[INFO] Paho MQTT 加载成功 (loadDex: " + dexFileName + ")");
- loaded = true;
- } catch (e) {
- printl("[WARN] loadDex 失败: " + e.message + ",尝试完整路径...");
- try {
- rhino.loadDex(dexTargetPath);
- printl("[INFO] Paho MQTT 加载成功 (loadDex 完整路径)");
- loaded = true;
- } catch (e2) {
- printl("[WARN] loadDex 完整路径也失败: " + e2.message);
- }
- }
- } else {
- try {
- rhino.loadDex(codeDir + "/paho-mqtt.dex");
- printl("[INFO] Paho MQTT 加载成功 (loadDex: 代码目录)");
- loaded = true;
- } catch (e) {
- printl("[WARN] loadDex(代码目录) 失败: " + e.message);
- }
- }
- if (!loaded) {
- printl("[WARN] 自动加载失败,尝试直接 importClass...");
- }
- // ========== 导入 Paho MQTT Java 类 ==========
- // 设置英文 locale,避免 Paho 查找中文资源包失败
- try {
- java.util.Locale.setDefault(java.util.Locale.ENGLISH);
- } catch (e) {}
- try {
- importClass(Packages.org.eclipse.paho.client.mqttv3.MqttClient);
- importClass(Packages.org.eclipse.paho.client.mqttv3.MqttConnectOptions);
- importClass(Packages.org.eclipse.paho.client.mqttv3.MqttCallback);
- importClass(Packages.org.eclipse.paho.client.mqttv3.MqttMessage);
- importClass(Packages.org.eclipse.paho.client.mqttv3.MqttException);
- importClass(Packages.org.eclipse.paho.client.mqttv3.persist.MemoryPersistence);
- importClass(Packages.org.eclipse.paho.client.mqttv3.IMqttDeliveryToken);
- importClass(Packages.org.eclipse.paho.client.mqttv3.MqttTopic);
- printl("[INFO] Paho MQTT 类导入成功");
- } catch (e) {
- printl("[FATAL] Paho MQTT 类导入失败: " + e.message);
- printl("[FATAL] 请确保 paho-mqtt.dex 在 H5HTML/插件/ 目录");
- printl("[FATAL] dex 文件已随工程打包在: H5HTML/插件/paho-mqtt.dex");
- printl("[FATAL] 如需手动下载 jar 并转换: https://repo1.maven.org/maven2/org/eclipse/paho/org.eclipse.paho.client.mqttv3/1.2.5/org.eclipse.paho.client.mqttv3-1.2.5.jar");
- printl("[FATAL] 转换命令: java -cp dx.jar com.android.dx.command.Main --dex --output=paho-mqtt.dex paho-mqtt.jar");
- exit();
- }
- // ========== 配置 ==========
- var CFG = {
- // MQTT Broker 地址(TCP 方式)
- brokerUrl: "tcp://broker.emqx.io:1883",
- // 客户端 ID(留空自动生成)
- clientId: "",
- // 用户名密码
- username: "",
- password: "",
- // 心跳间隔(秒)
- keepAliveSec: 30,
- // 连接超时(秒)
- connectTimeoutSec: 15,
- // 会话清除标志
- cleanSession: true,
- // 遗嘱消息(LWT)
- willTopic: "",
- willMessage: "",
- willQos: 0,
- willRetained: false,
- // 自动重连
- autoReconnect: true,
- reconnectDelayMs: 3000,
- maxReconnectDelayMs: 30000,
- // 默认 QoS
- defaultQos: 1,
- // 订阅主题列表
- subscribeTopics: [
- { topic: "aiwrok/cmd/+", qos: 1 },
- { topic: "aiwrok/broadcast", qos: 0 },
- { topic: "aiwrok/status/request", qos: 1 }
- ],
- // 发布主题前缀
- pubTopicPrefix: "aiwrok",
- // 调试日志
- debug: true,
- // 脚本运行时长(毫秒),0=无限
- runDuration: 0
- };
- // ========== 全局状态 ==========
- var G = {
- mqttClient: null,
- connected: false,
- running: true,
- startTime: 0,
- msgSentCount: 0,
- msgRecvCount: 0,
- reconnectCount: 0
- };
- // ========== 时间工具 ==========
- function nowMs() {
- return (new Date()).getTime();
- }
- function nowSec() {
- return Math.floor(nowMs() / 1000);
- }
- // ========== 设备信息 ==========
- function getDeviceInfo() {
- var info = {
- imei: "",
- brand: "",
- model: "",
- android: "",
- screen: ""
- };
- try { info.imei = device.getIMEI() || ""; } catch (e) {}
- try { info.brand = device.getBrand() || ""; } catch (e) {}
- try { info.model = device.getModel() || ""; } catch (e) {}
- try { info.android = java.lang.System.getProperty("os.version") + ""; } catch (e) {}
- try { info.screen = screen.getScreenWidth() + "x" + screen.getScreenHeight(); } catch (e) {}
- return info;
- }
- // ========== 生成客户端 ID ==========
- function generateClientId() {
- var info = getDeviceInfo();
- var ts = nowMs();
- var rand = Math.floor(Math.random() * 10000);
- return "aiwrok_" + (info.imei || info.model || "dev") + "_" + ts + "_" + rand;
- }
- // ========== 获取发布主题 ==========
- function getPubTopic(suffix) {
- return CFG.pubTopicPrefix + "/" + suffix;
- }
- // ========== 发布消息 ==========
- function publish(topic, payload, qos) {
- if (!G.connected || !G.mqttClient) {
- printl("[WARN] 未连接,无法发布: " + topic);
- return false;
- }
- if (qos == null) qos = CFG.defaultQos;
- try {
- var content;
- if (typeof payload === "string") {
- content = payload;
- } else {
- content = JSON.stringify(payload);
- }
- var msg = new MqttMessage();
- msg.setPayload(new java.lang.String(content).getBytes("UTF-8"));
- msg.setQos(qos);
- G.mqttClient.publish(topic, msg);
- G.msgSentCount++;
- printl("[PUB] " + topic + " qos=" + qos + " len=" + content.length);
- return true;
- } catch (e) {
- printl("[ERROR] 发布异常: " + e.message);
- return false;
- }
- }
- // ========== 处理收到的消息 ==========
- function handleMessage(topic, message) {
- G.msgRecvCount++;
- var payload = "";
- try {
- payload = "" + new java.lang.String(message.getPayload(), "UTF-8");
- } catch (e) {
- try { payload = String(message); } catch (e2) { payload = "" + message; }
- }
- printl("[SUB] " + topic + " len=" + payload.length + " | " + payload.slice(0, 200));
- var obj = null;
- try {
- obj = JSON.parse(payload);
- } catch (e) {
- printl("[DEBUG] 非 JSON 消息: " + payload.slice(0, 100));
- return;
- }
- if (topic.indexOf("cmd/") >= 0) {
- handleCmdMessage(topic, obj);
- } else if (topic.indexOf("broadcast") >= 0) {
- printl("[INFO] 广播消息: " + (obj.msg || payload.slice(0, 100)));
- } else if (topic.indexOf("status/request") >= 0) {
- handleStatusRequest(topic, obj);
- } else {
- printl("[DEBUG] 未处理的主题: " + topic);
- }
- }
- // ========== 处理命令消息 ==========
- function handleCmdMessage(topic, obj) {
- var cmd = obj.cmd || obj.type || "";
- printl("[INFO] 收到命令: " + cmd);
- if (cmd === "SCREENSHOT") {
- try {
- var bitmap = screen.screenShotFull();
- if (!bitmap) {
- publish(getPubTopic("resp/screenshot"), { ok: false, msg: "截图返回空" });
- return;
- }
- var base64 = "" + bitmap.toBase64();
- try { bitmap.recycle(); } catch (e) {}
- publish(getPubTopic("resp/screenshot"), {
- ok: true,
- cmdId: obj.cmdId,
- base64: base64,
- length: base64.length
- });
- } catch (e) {
- publish(getPubTopic("resp/screenshot"), { ok: false, msg: e.message });
- }
- } else if (cmd === "CLICK") {
- try {
- var x = parseInt(obj.data.x, 10);
- var y = parseInt(obj.data.y, 10);
- if (isNaN(x) || isNaN(y)) {
- publish(getPubTopic("resp/click"), { ok: false, msg: "坐标无效" });
- return;
- }
- if (typeof action !== "undefined") {
- action.click(x, y);
- } else if (typeof hid !== "undefined") {
- hid.click(x, y);
- }
- publish(getPubTopic("resp/click"), { ok: true, cmdId: obj.cmdId, x: x, y: y });
- } catch (e) {
- publish(getPubTopic("resp/click"), { ok: false, msg: e.message });
- }
- } else if (cmd === "STATUS") {
- handleStatusRequest(topic, obj);
- } else if (cmd === "HOME") {
- try { if (typeof hid !== "undefined") hid.home(); } catch (e) {}
- } else if (cmd === "BACK") {
- try { if (typeof hid !== "undefined") hid.back(); } catch (e) {}
- } else {
- printl("[DEBUG] 未知命令: " + cmd);
- }
- }
- // ========== 处理状态查询 ==========
- function handleStatusRequest(topic, obj) {
- var info = getDeviceInfo();
- info.ok = true;
- info.cmdId = obj.cmdId;
- info.runtimeSec = nowSec() - G.startTime;
- info.msgSentCount = G.msgSentCount;
- info.msgRecvCount = G.msgRecvCount;
- info.reconnectCount = G.reconnectCount;
- publish(getPubTopic("resp/status"), info);
- }
- // ========== 订阅主题 ==========
- function subscribeTopics() {
- if (!G.connected || !G.mqttClient) return;
- for (var i = 0; i < CFG.subscribeTopics.length; i++) {
- var item = CFG.subscribeTopics[i];
- try {
- G.mqttClient.subscribe(item.topic, item.qos);
- printl("[INFO] 已订阅: " + item.topic + " qos=" + item.qos);
- } catch (e) {
- printl("[ERROR] 订阅失败 " + item.topic + ": " + e.message);
- }
- }
- }
- // ========== MQTT 回调实现 ==========
- function createMqttCallback() {
- return new MqttCallback({
- connectionLost: function (cause) {
- G.connected = false;
- var msg = "unknown";
- try { msg = String(cause.getMessage()); } catch (e) {}
- printl("[WARN] 连接断开: " + msg);
- },
- messageArrived: function (topic, message) {
- try {
- handleMessage("" + topic, message);
- } catch (e) {
- printl("[ERROR] 消息处理异常: " + e.message);
- }
- },
- deliveryComplete: function (token) {
- printl("[DEBUG] 消息投递完成");
- }
- });
- }
- // ========== 连接失败统一诊断打印(分阶段定位 + JS/Java 双堆栈) ==========
- function logConnErr(stage, e) {
- var errMsg = "" + e;
- try {
- if (e instanceof MqttException) {
- errMsg = "MQTT Error " + e.getReasonCode() + ": " + e.getMessage();
- }
- } catch (e2) {}
- printl("[ERROR] " + stage + " 失败: " + errMsg);
- // Rhino 的 e.stack 指向脚本行号
- try { printl("[ERROR] " + stage + " JS堆栈: " + (e.stack || "(无)")); } catch (e3) {}
- // e.javaException 是底层 Java Throwable(Rhino 包装对象本身不能直接 printStackTrace)
- try {
- var je = e.javaException;
- if (je) {
- var sw = new java.io.StringWriter();
- var pw = new java.io.PrintWriter(sw);
- je.printStackTrace(pw);
- pw.flush();
- printl("[ERROR] " + stage + " Java堆栈: " + sw.toString());
- }
- } catch (e4) {}
- }
- // ========== 建立连接 ==========
- // 拆成 3 个阶段分别 try/catch:哪一阶段抛错一目了然
- function connect() {
- printl("[INFO] 正在连接 MQTT Broker: " + CFG.brokerUrl);
- if (!CFG.clientId) {
- CFG.clientId = generateClientId();
- }
- printl("[INFO] 客户端 ID: " + CFG.clientId);
- // ---- 阶段1:创建 MqttClient(构造器里有 log.fine,是 NPE 高发点) ----
- try {
- var persistence = new MemoryPersistence();
- var brokerUri = CFG.brokerUrl;
- if (brokerUri.indexOf("://") < 0) {
- brokerUri = "tcp://" + brokerUri;
- }
- G.mqttClient = new MqttClient(brokerUri, CFG.clientId, persistence);
- printl("[INFO] [阶段1/3] MqttClient 创建成功");
- } catch (e1) {
- G.connected = false;
- logConnErr("阶段1 创建MqttClient", e1);
- return false;
- }
- // ---- 阶段2:配置连接选项 ----
- var options = null;
- try {
- options = new MqttConnectOptions();
- options.setKeepAliveInterval(CFG.keepAliveSec);
- options.setConnectionTimeout(CFG.connectTimeoutSec);
- options.setCleanSession(CFG.cleanSession);
- if (CFG.autoReconnect) {
- options.setAutomaticReconnect(true);
- try {
- options.setMaxReconnectDelay(CFG.maxReconnectDelayMs);
- } catch (e) {}
- }
- if (CFG.username) {
- options.setUserName(CFG.username);
- }
- if (CFG.password) {
- options.setPassword(CFG.password.toCharArray());
- }
- if (CFG.willTopic && CFG.willMessage) {
- options.setWill(
- CFG.willTopic,
- new java.lang.String(CFG.willMessage).getBytes("UTF-8"),
- CFG.willQos,
- CFG.willRetained
- );
- }
- printl("[INFO] [阶段2/3] 连接选项配置完成");
- } catch (e2) {
- G.connected = false;
- logConnErr("阶段2 配置选项", e2);
- return false;
- }
- // ---- 阶段3:设置回调并执行连接(网络握手在这里) ----
- try {
- G.mqttClient.setCallback(createMqttCallback());
- G.mqttClient.connect(options);
- G.connected = true;
- G.reconnectCount = 0;
- printl("[INFO] [阶段3/3] MQTT 连接成功 — 握手完成");
- try { subscribeTopics(); } catch (eS) { printl("[WARN] 订阅异常: " + eS.message); }
- try {
- publish(getPubTopic("online"), {
- clientId: CFG.clientId,
- device: getDeviceInfo(),
- ts: nowSec()
- }, 0);
- } catch (eP) { printl("[WARN] 上线消息发布异常: " + eP.message); }
- return true;
- } catch (e3) {
- G.connected = false;
- logConnErr("阶段3 执行连接", e3);
- return false;
- }
- }
- // ========== 断开连接 ==========
- function disconnect() {
- if (G.mqttClient) {
- try {
- if (G.connected) {
- publish(getPubTopic("offline"), {
- clientId: CFG.clientId,
- ts: nowSec()
- }, 0);
- }
- G.mqttClient.disconnect();
- printl("[INFO] 已断开 MQTT 连接");
- } catch (e) {
- printl("[WARN] 断开异常: " + e.message);
- }
- }
- G.connected = false;
- }
- // ========== 关闭清理 ==========
- function cleanup() {
- G.running = false;
- disconnect();
- try {
- if (G.mqttClient) {
- G.mqttClient.close(true);
- }
- } catch (e) {}
- printl("[INFO] 资源已清理");
- }
- // ========== 主程序 ==========
- function main() {
- try {
- G.startTime = nowSec();
- printl("[INFO] ==========================================");
- printl("[INFO] MQTT 客户端 (Eclipse Paho) 启动");
- printl("[INFO] Broker: " + CFG.brokerUrl);
- printl("[INFO] ClientID: " + (CFG.clientId || "自动生成"));
- printl("[INFO] 心跳: " + CFG.keepAliveSec + "s");
- printl("[INFO] 自动重连: " + CFG.autoReconnect);
- printl("[INFO] CleanSession: " + CFG.cleanSession);
- printl("[INFO] 订阅主题: " + CFG.subscribeTopics.length + " 个");
- printl("[INFO] 设备信息: " + JSON.stringify(getDeviceInfo()));
- printl("[INFO] ==========================================");
- if (!connect()) {
- startReconnectLoop();
- }
- var tick = 0;
- (function keepAlive() {
- try {
- if (!G.running) {
- printl("[INFO] 主线程保活退出");
- return;
- }
- tick++;
- if (tick % 30 === 0) {
- printl("[INFO] [保活] 已运行=" + (nowSec() - G.startTime) +
- "s 连接=" + G.connected +
- " 发送=" + G.msgSentCount +
- " 接收=" + G.msgRecvCount +
- " 重连=" + G.reconnectCount);
- }
- } catch (e) {
- printl("[ERROR] 保活异常(已兜住): " + e.message);
- }
- try { setTimeout(keepAlive, 1000); }
- catch (e) {
- try { runTime.setTimeout(keepAlive, 1000); } catch (e2) {}
- }
- })();
- if (CFG.runDuration > 0) {
- sleep.millisecond(CFG.runDuration);
- cleanup();
- }
- } catch (e) {
- printl("[FATAL] 主程序异常: " + e.message + "\n" + (e.stack || ""));
- try { cleanup(); } catch (e2) {}
- }
- }
- // ========== 手动重连循环 ==========
- function startReconnectLoop() {
- (function reconnectStep() {
- try {
- if (!G.running) return;
- if (G.connected) return;
- G.reconnectCount++;
- printl("[INFO] 手动重连尝试 " + G.reconnectCount);
- if (connect()) {
- printl("[INFO] 手动重连成功");
- return;
- }
- setTimeout(reconnectStep, CFG.reconnectDelayMs);
- } catch (e) {
- printl("[ERROR] 重连异常(已兜住): " + e.message);
- setTimeout(reconnectStep, CFG.reconnectDelayMs);
- }
- })();
- }
- // ========== 启动 ==========
- main();
复制代码
| |  | |  |
|