YYPOST群发软件 发表于 9 小时前

MQTT协议示例AIWROK版

MQTT协议示例AIWROK版








/**
* 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(" 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(" 内嵌 dex 数据: " + dexBytes.length + " bytes");
    } else {
      printl(" 未取到内嵌 base64 数据(PAHO_MQTT_DEX_B64 为空)");
    }
} catch (e) {
    printl(" 内嵌 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(" 插件目录 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(" 已还原 " + dexFileName + " 到插件目录");
      installed = true;
    } catch (e) {
      printl(" 写入 dex 失败: " + e.message);
    }
}

// 最终判定能否加载:还原成功,或手机上本来就有 dex(退回使用现有文件)
var dexExists = installed || dexWasThere;
if (!installed && !verOk && dexWasThere) {
    printl(" 内嵌还原未完成,改用手机上现有的 dex 继续加载");
}

// 第三步:用 loadDex 加载
if (dexExists) {
    try {
      rhino.loadDex(dexFileName);
      printl(" Paho MQTT 加载成功 (loadDex: " + dexFileName + ")");
      loaded = true;
    } catch (e) {
      printl(" loadDex 失败: " + e.message + ",尝试完整路径...");
      try {
            rhino.loadDex(dexTargetPath);
            printl(" Paho MQTT 加载成功 (loadDex 完整路径)");
            loaded = true;
      } catch (e2) {
            printl(" loadDex 完整路径也失败: " + e2.message);
      }
    }
} else {
    try {
      rhino.loadDex(codeDir + "/paho-mqtt.dex");
      printl(" Paho MQTT 加载成功 (loadDex: 代码目录)");
      loaded = true;
    } catch (e) {
      printl(" loadDex(代码目录) 失败: " + e.message);
    }
}

if (!loaded) {
    printl(" 自动加载失败,尝试直接 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(" Paho MQTT 类导入成功");
} catch (e) {
    printl(" Paho MQTT 类导入失败: " + e.message);
    printl(" 请确保 paho-mqtt.dex 在 H5HTML/插件/ 目录");
    printl(" dex 文件已随工程打包在: H5HTML/插件/paho-mqtt.dex");
    printl(" 如需手动下载 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(" 转换命令: 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(" 未连接,无法发布: " + 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(" " + topic + " qos=" + qos + " len=" + content.length);
      return true;
    } catch (e) {
      printl(" 发布异常: " + 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(" " + topic + " len=" + payload.length + " | " + payload.slice(0, 200));

    var obj = null;
    try {
      obj = JSON.parse(payload);
    } catch (e) {
      printl(" 非 JSON 消息: " + payload.slice(0, 100));
      return;
    }

    if (topic.indexOf("cmd/") >= 0) {
      handleCmdMessage(topic, obj);
    } else if (topic.indexOf("broadcast") >= 0) {
      printl(" 广播消息: " + (obj.msg || payload.slice(0, 100)));
    } else if (topic.indexOf("status/request") >= 0) {
      handleStatusRequest(topic, obj);
    } else {
      printl(" 未处理的主题: " + topic);
    }
}

// ========== 处理命令消息 ==========
function handleCmdMessage(topic, obj) {
    var cmd = obj.cmd || obj.type || "";
    printl(" 收到命令: " + 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(" 未知命令: " + 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;
      try {
            G.mqttClient.subscribe(item.topic, item.qos);
            printl(" 已订阅: " + item.topic + " qos=" + item.qos);
      } catch (e) {
            printl(" 订阅失败 " + 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(" 连接断开: " + msg);
      },

      messageArrived: function (topic, message) {
            try {
                handleMessage("" + topic, message);
            } catch (e) {
                printl(" 消息处理异常: " + e.message);
            }
      },

      deliveryComplete: function (token) {
            printl(" 消息投递完成");
      }
    });
}

// ========== 连接失败统一诊断打印(分阶段定位 + JS/Java 双堆栈) ==========
function logConnErr(stage, e) {
    var errMsg = "" + e;
    try {
      if (e instanceof MqttException) {
            errMsg = "MQTT Error " + e.getReasonCode() + ": " + e.getMessage();
      }
    } catch (e2) {}
    printl(" " + stage + " 失败: " + errMsg);
    // Rhino 的 e.stack 指向脚本行号
    try { printl(" " + 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(" " + stage + " Java堆栈: " + sw.toString());
      }
    } catch (e4) {}
}

// ========== 建立连接 ==========
// 拆成 3 个阶段分别 try/catch:哪一阶段抛错一目了然
function connect() {
    printl(" 正在连接 MQTT Broker: " + CFG.brokerUrl);

    if (!CFG.clientId) {
      CFG.clientId = generateClientId();
    }
    printl(" 客户端 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(" [阶段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(" [阶段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(" [阶段3/3] MQTT 连接成功 — 握手完成");

      try { subscribeTopics(); } catch (eS) { printl(" 订阅异常: " + eS.message); }
      try {
            publish(getPubTopic("online"), {
                clientId: CFG.clientId,
                device: getDeviceInfo(),
                ts: nowSec()
            }, 0);
      } catch (eP) { printl(" 上线消息发布异常: " + 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(" 已断开 MQTT 连接");
      } catch (e) {
            printl(" 断开异常: " + e.message);
      }
    }
    G.connected = false;
}

// ========== 关闭清理 ==========
function cleanup() {
    G.running = false;
    disconnect();
    try {
      if (G.mqttClient) {
            G.mqttClient.close(true);
      }
    } catch (e) {}
    printl(" 资源已清理");
}

// ========== 主程序 ==========
function main() {
    try {
      G.startTime = nowSec();

      printl(" ==========================================");
      printl("   MQTT 客户端 (Eclipse Paho) 启动");
      printl("   Broker: " + CFG.brokerUrl);
      printl("   ClientID: " + (CFG.clientId || "自动生成"));
      printl("   心跳: " + CFG.keepAliveSec + "s");
      printl("   自动重连: " + CFG.autoReconnect);
      printl("   CleanSession: " + CFG.cleanSession);
      printl("   订阅主题: " + CFG.subscribeTopics.length + " 个");
      printl("   设备信息: " + JSON.stringify(getDeviceInfo()));
      printl(" ==========================================");

      if (!connect()) {
            startReconnectLoop();
      }

      var tick = 0;
      (function keepAlive() {
            try {
                if (!G.running) {
                  printl(" 主线程保活退出");
                  return;
                }
                tick++;
                if (tick % 30 === 0) {
                  printl(" [保活] 已运行=" + (nowSec() - G.startTime) +
                        "s 连接=" + G.connected +
                        " 发送=" + G.msgSentCount +
                        " 接收=" + G.msgRecvCount +
                        " 重连=" + G.reconnectCount);
                }
            } catch (e) {
                printl(" 保活异常(已兜住): " + 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(" 主程序异常: " + 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(" 手动重连尝试 " + G.reconnectCount);

            if (connect()) {
                printl(" 手动重连成功");
                return;
            }

            setTimeout(reconnectStep, CFG.reconnectDelayMs);
      } catch (e) {
            printl(" 重连异常(已兜住): " + e.message);
            setTimeout(reconnectStep, CFG.reconnectDelayMs);
      }
    })();
}

// ========== 启动 ==========
main();


页: [1]
查看完整版本: MQTT协议示例AIWROK版