B2B网络软件

标题: MQTT协议示例AIWROK版 [打印本页]

作者: YYPOST群发软件    时间: 9 小时前
标题: MQTT协议示例AIWROK版
MQTT协议示例AIWROK版
MQTT协议示例AIWROK版 B2B网络软件

MQTT协议示例AIWROK版 B2B网络软件

MQTT协议示例AIWROK版 B2B网络软件

MQTT协议示例AIWROK版 B2B网络软件

  1. /**
  2. * MQTT 协议示例 - AIWROK 版
  3. *
  4. * 使用 Eclipse Paho Java 客户端实现 MQTT v3.1.1 连接
  5. *
  6. * v5 说明:dex 为官方 Paho 1.2.5 源码修补版(2026-08-29 重建):
  7. *   1. LoggerFactory 资源包查找加保护,默认使用 MockLogger(无操作日志)
  8. *      → 修复 dex 环境下 ResourceBundle.getBundle 必抛 MissingResourceException
  9. *   2. NetworkModuleService 加内置协议工厂兜底(tcp/ssl/ws/wss)
  10. *      → 修复 dex 无 META-INF/services 导致 "no NetworkModule installed for scheme tcp"
  11. *   3. 脚本不再需要任何运行时反射补丁
  12. * 支持:发布/订阅模式、QoS 0/1/2、自动重连、遗嘱消息、心跳保活
  13. *
  14. * MQTT vs WebSocket 区别:
  15. *   MQTT 是轻量级发布/订阅消息协议,适合 IoT/移动端
  16. *   比 WebSocket 更省电、更省流量,支持 QoS 消息质量保证
  17. *   支持主题(Topic)订阅,消息路由更灵活
  18. *
  19. * 依赖加载(自动,无需手动操作):
  20. *   dex 文件放在 H5HTML/插件/ 目录,AIWROK 的 loadDex 自动加载
  21. *   loadDex 只从"插件"目录加载,且需要 .dex 格式(不是 .jar)
  22. *
  23. *   dex 文件已随工程打包在:H5HTML/插件/paho-mqtt.dex
  24. *   如果插件目录没有 dex,自动从代码目录的 paho-mqtt-dex-b64.js 还原
  25. *
  26. * 使用方法:
  27. *   1. 修改 CFG 中的 broker 地址、端口、用户名密码
  28. *   2. AIWROK 中直接运行本脚本
  29. *   3. 日志窗口查看连接状态和收发消息
  30. */

  31. // ========== 加载 Paho MQTT Java 库 ==========
  32. // AIWROK loadDex 只从 H5HTML/插件/ 目录加载 .dex 文件
  33. // AIWROK IDE 不会同步插件目录的二进制文件到手机,
  34. // 脚本运行时从内嵌的 base64 数据自动还原 dex 文件到插件目录
  35. var loaded = false;

  36. // 手机端路径
  37. var sdcardPath = "/sdcard";
  38. var projectRoot = sdcardPath + "/auto/H5HTML";
  39. var pluginDir = projectRoot + "/插件";
  40. var codeDir = projectRoot + "/代码";
  41. var dexFileName = "paho-mqtt.dex";
  42. var dexTargetPath = pluginDir + "/" + dexFileName;

  43. // dex 版本判定:以内嵌 base64 解码后的长度为唯一基准
  44. // (这样以后重新打包 dex 不必改任何常量,脚本自己会比对并覆盖旧版)
  45. function fileSize(path) {
  46.     try {
  47.         var f = new java.io.File(path);
  48.         return f.exists() ? Number(f.length()) : -1;
  49.     } catch (e) { return -1; }
  50. }

  51. // 先取内嵌数据:require 失败只记日志,不让脚本崩
  52. var dexBytes = null;
  53. try {
  54.     try { require("paho-mqtt-dex-b64.js"); } catch (reqErr) {
  55.         printl("[WARN] require paho-mqtt-dex-b64.js 失败: " + reqErr.message);
  56.     }
  57.     if (typeof PAHO_MQTT_DEX_B64 !== "undefined" && PAHO_MQTT_DEX_B64) {
  58.         dexBytes = java.util.Base64.getDecoder().decode(PAHO_MQTT_DEX_B64);
  59.         printl("[INFO] 内嵌 dex 数据: " + dexBytes.length + " bytes");
  60.     } else {
  61.         printl("[WARN] 未取到内嵌 base64 数据(PAHO_MQTT_DEX_B64 为空)");
  62.     }
  63. } catch (e) {
  64.     printl("[WARN] 内嵌 dex 解码失败: " + e.message);
  65.     dexBytes = null;
  66. }

  67. var expectedSize = dexBytes ? Number(dexBytes.length) : -1;
  68. var curSize = fileSize(dexTargetPath);
  69. var dexWasThere = curSize > 0;
  70. // 有内嵌数据就按字节数判版本;没有就"有 dex 先用着"
  71. var verOk = (expectedSize > 0) ? (curSize === expectedSize) : dexWasThere;

  72. printl("[INFO] 插件目录 dex: " + (dexWasThere ? ("已存在 " + curSize + " bytes") : "不存在") +
  73.     "," + (verOk ? "版本校验通过" : ("需重装(期望 " + (expectedSize > 0 ? expectedSize : "未知") + " bytes)")));

  74. // 第一步:确保插件目录存在
  75. try { file.mkdir(pluginDir); } catch (e) {}

  76. // 第二步:版本不对就用内嵌数据覆盖还原
  77. var installed = false;
  78. if (!verOk && dexBytes) {
  79.     try {
  80.         var fos = new java.io.FileOutputStream(dexTargetPath);
  81.         fos.write(dexBytes);
  82.         fos.close();
  83.         printl("[INFO] 已还原 " + dexFileName + " 到插件目录");
  84.         installed = true;
  85.     } catch (e) {
  86.         printl("[WARN] 写入 dex 失败: " + e.message);
  87.     }
  88. }

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

  94. // 第三步:用 loadDex 加载
  95. if (dexExists) {
  96.     try {
  97.         rhino.loadDex(dexFileName);
  98.         printl("[INFO] Paho MQTT 加载成功 (loadDex: " + dexFileName + ")");
  99.         loaded = true;
  100.     } catch (e) {
  101.         printl("[WARN] loadDex 失败: " + e.message + ",尝试完整路径...");
  102.         try {
  103.             rhino.loadDex(dexTargetPath);
  104.             printl("[INFO] Paho MQTT 加载成功 (loadDex 完整路径)");
  105.             loaded = true;
  106.         } catch (e2) {
  107.             printl("[WARN] loadDex 完整路径也失败: " + e2.message);
  108.         }
  109.     }
  110. } else {
  111.     try {
  112.         rhino.loadDex(codeDir + "/paho-mqtt.dex");
  113.         printl("[INFO] Paho MQTT 加载成功 (loadDex: 代码目录)");
  114.         loaded = true;
  115.     } catch (e) {
  116.         printl("[WARN] loadDex(代码目录) 失败: " + e.message);
  117.     }
  118. }

  119. if (!loaded) {
  120.     printl("[WARN] 自动加载失败,尝试直接 importClass...");
  121. }

  122. // ========== 导入 Paho MQTT Java 类 ==========
  123. // 设置英文 locale,避免 Paho 查找中文资源包失败
  124. try {
  125.     java.util.Locale.setDefault(java.util.Locale.ENGLISH);
  126. } catch (e) {}

  127. try {
  128.     importClass(Packages.org.eclipse.paho.client.mqttv3.MqttClient);
  129.     importClass(Packages.org.eclipse.paho.client.mqttv3.MqttConnectOptions);
  130.     importClass(Packages.org.eclipse.paho.client.mqttv3.MqttCallback);
  131.     importClass(Packages.org.eclipse.paho.client.mqttv3.MqttMessage);
  132.     importClass(Packages.org.eclipse.paho.client.mqttv3.MqttException);
  133.     importClass(Packages.org.eclipse.paho.client.mqttv3.persist.MemoryPersistence);
  134.     importClass(Packages.org.eclipse.paho.client.mqttv3.IMqttDeliveryToken);
  135.     importClass(Packages.org.eclipse.paho.client.mqttv3.MqttTopic);
  136.     printl("[INFO] Paho MQTT 类导入成功");
  137. } catch (e) {
  138.     printl("[FATAL] Paho MQTT 类导入失败: " + e.message);
  139.     printl("[FATAL] 请确保 paho-mqtt.dex 在 H5HTML/插件/ 目录");
  140.     printl("[FATAL] dex 文件已随工程打包在: H5HTML/插件/paho-mqtt.dex");
  141.     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");
  142.     printl("[FATAL] 转换命令: java -cp dx.jar com.android.dx.command.Main --dex --output=paho-mqtt.dex paho-mqtt.jar");
  143.     exit();
  144. }

  145. // ========== 配置 ==========
  146. var CFG = {
  147.     // MQTT Broker 地址(TCP 方式)
  148.     brokerUrl: "tcp://broker.emqx.io:1883",

  149.     // 客户端 ID(留空自动生成)
  150.     clientId: "",

  151.     // 用户名密码
  152.     username: "",
  153.     password: "",

  154.     // 心跳间隔(秒)
  155.     keepAliveSec: 30,

  156.     // 连接超时(秒)
  157.     connectTimeoutSec: 15,

  158.     // 会话清除标志
  159.     cleanSession: true,

  160.     // 遗嘱消息(LWT)
  161.     willTopic: "",
  162.     willMessage: "",
  163.     willQos: 0,
  164.     willRetained: false,

  165.     // 自动重连
  166.     autoReconnect: true,
  167.     reconnectDelayMs: 3000,
  168.     maxReconnectDelayMs: 30000,

  169.     // 默认 QoS
  170.     defaultQos: 1,

  171.     // 订阅主题列表
  172.     subscribeTopics: [
  173.         { topic: "aiwrok/cmd/+", qos: 1 },
  174.         { topic: "aiwrok/broadcast", qos: 0 },
  175.         { topic: "aiwrok/status/request", qos: 1 }
  176.     ],

  177.     // 发布主题前缀
  178.     pubTopicPrefix: "aiwrok",

  179.     // 调试日志
  180.     debug: true,

  181.     // 脚本运行时长(毫秒),0=无限
  182.     runDuration: 0
  183. };

  184. // ========== 全局状态 ==========
  185. var G = {
  186.     mqttClient: null,
  187.     connected: false,
  188.     running: true,
  189.     startTime: 0,
  190.     msgSentCount: 0,
  191.     msgRecvCount: 0,
  192.     reconnectCount: 0
  193. };

  194. // ========== 时间工具 ==========
  195. function nowMs() {
  196.     return (new Date()).getTime();
  197. }
  198. function nowSec() {
  199.     return Math.floor(nowMs() / 1000);
  200. }

  201. // ========== 设备信息 ==========
  202. function getDeviceInfo() {
  203.     var info = {
  204.         imei: "",
  205.         brand: "",
  206.         model: "",
  207.         android: "",
  208.         screen: ""
  209.     };
  210.     try { info.imei = device.getIMEI() || ""; } catch (e) {}
  211.     try { info.brand = device.getBrand() || ""; } catch (e) {}
  212.     try { info.model = device.getModel() || ""; } catch (e) {}
  213.     try { info.android = java.lang.System.getProperty("os.version") + ""; } catch (e) {}
  214.     try { info.screen = screen.getScreenWidth() + "x" + screen.getScreenHeight(); } catch (e) {}
  215.     return info;
  216. }

  217. // ========== 生成客户端 ID ==========
  218. function generateClientId() {
  219.     var info = getDeviceInfo();
  220.     var ts = nowMs();
  221.     var rand = Math.floor(Math.random() * 10000);
  222.     return "aiwrok_" + (info.imei || info.model || "dev") + "_" + ts + "_" + rand;
  223. }

  224. // ========== 获取发布主题 ==========
  225. function getPubTopic(suffix) {
  226.     return CFG.pubTopicPrefix + "/" + suffix;
  227. }

  228. // ========== 发布消息 ==========
  229. function publish(topic, payload, qos) {
  230.     if (!G.connected || !G.mqttClient) {
  231.         printl("[WARN] 未连接,无法发布: " + topic);
  232.         return false;
  233.     }
  234.     if (qos == null) qos = CFG.defaultQos;

  235.     try {
  236.         var content;
  237.         if (typeof payload === "string") {
  238.             content = payload;
  239.         } else {
  240.             content = JSON.stringify(payload);
  241.         }

  242.         var msg = new MqttMessage();
  243.         msg.setPayload(new java.lang.String(content).getBytes("UTF-8"));
  244.         msg.setQos(qos);

  245.         G.mqttClient.publish(topic, msg);
  246.         G.msgSentCount++;

  247.         printl("[PUB] " + topic + " qos=" + qos + " len=" + content.length);
  248.         return true;
  249.     } catch (e) {
  250.         printl("[ERROR] 发布异常: " + e.message);
  251.         return false;
  252.     }
  253. }

  254. // ========== 处理收到的消息 ==========
  255. function handleMessage(topic, message) {
  256.     G.msgRecvCount++;

  257.     var payload = "";
  258.     try {
  259.         payload = "" + new java.lang.String(message.getPayload(), "UTF-8");
  260.     } catch (e) {
  261.         try { payload = String(message); } catch (e2) { payload = "" + message; }
  262.     }

  263.     printl("[SUB] " + topic + " len=" + payload.length + " | " + payload.slice(0, 200));

  264.     var obj = null;
  265.     try {
  266.         obj = JSON.parse(payload);
  267.     } catch (e) {
  268.         printl("[DEBUG] 非 JSON 消息: " + payload.slice(0, 100));
  269.         return;
  270.     }

  271.     if (topic.indexOf("cmd/") >= 0) {
  272.         handleCmdMessage(topic, obj);
  273.     } else if (topic.indexOf("broadcast") >= 0) {
  274.         printl("[INFO] 广播消息: " + (obj.msg || payload.slice(0, 100)));
  275.     } else if (topic.indexOf("status/request") >= 0) {
  276.         handleStatusRequest(topic, obj);
  277.     } else {
  278.         printl("[DEBUG] 未处理的主题: " + topic);
  279.     }
  280. }

  281. // ========== 处理命令消息 ==========
  282. function handleCmdMessage(topic, obj) {
  283.     var cmd = obj.cmd || obj.type || "";
  284.     printl("[INFO] 收到命令: " + cmd);

  285.     if (cmd === "SCREENSHOT") {
  286.         try {
  287.             var bitmap = screen.screenShotFull();
  288.             if (!bitmap) {
  289.                 publish(getPubTopic("resp/screenshot"), { ok: false, msg: "截图返回空" });
  290.                 return;
  291.             }
  292.             var base64 = "" + bitmap.toBase64();
  293.             try { bitmap.recycle(); } catch (e) {}
  294.             publish(getPubTopic("resp/screenshot"), {
  295.                 ok: true,
  296.                 cmdId: obj.cmdId,
  297.                 base64: base64,
  298.                 length: base64.length
  299.             });
  300.         } catch (e) {
  301.             publish(getPubTopic("resp/screenshot"), { ok: false, msg: e.message });
  302.         }
  303.     } else if (cmd === "CLICK") {
  304.         try {
  305.             var x = parseInt(obj.data.x, 10);
  306.             var y = parseInt(obj.data.y, 10);
  307.             if (isNaN(x) || isNaN(y)) {
  308.                 publish(getPubTopic("resp/click"), { ok: false, msg: "坐标无效" });
  309.                 return;
  310.             }
  311.             if (typeof action !== "undefined") {
  312.                 action.click(x, y);
  313.             } else if (typeof hid !== "undefined") {
  314.                 hid.click(x, y);
  315.             }
  316.             publish(getPubTopic("resp/click"), { ok: true, cmdId: obj.cmdId, x: x, y: y });
  317.         } catch (e) {
  318.             publish(getPubTopic("resp/click"), { ok: false, msg: e.message });
  319.         }
  320.     } else if (cmd === "STATUS") {
  321.         handleStatusRequest(topic, obj);
  322.     } else if (cmd === "HOME") {
  323.         try { if (typeof hid !== "undefined") hid.home(); } catch (e) {}
  324.     } else if (cmd === "BACK") {
  325.         try { if (typeof hid !== "undefined") hid.back(); } catch (e) {}
  326.     } else {
  327.         printl("[DEBUG] 未知命令: " + cmd);
  328.     }
  329. }

  330. // ========== 处理状态查询 ==========
  331. function handleStatusRequest(topic, obj) {
  332.     var info = getDeviceInfo();
  333.     info.ok = true;
  334.     info.cmdId = obj.cmdId;
  335.     info.runtimeSec = nowSec() - G.startTime;
  336.     info.msgSentCount = G.msgSentCount;
  337.     info.msgRecvCount = G.msgRecvCount;
  338.     info.reconnectCount = G.reconnectCount;
  339.     publish(getPubTopic("resp/status"), info);
  340. }

  341. // ========== 订阅主题 ==========
  342. function subscribeTopics() {
  343.     if (!G.connected || !G.mqttClient) return;

  344.     for (var i = 0; i < CFG.subscribeTopics.length; i++) {
  345.         var item = CFG.subscribeTopics[i];
  346.         try {
  347.             G.mqttClient.subscribe(item.topic, item.qos);
  348.             printl("[INFO] 已订阅: " + item.topic + " qos=" + item.qos);
  349.         } catch (e) {
  350.             printl("[ERROR] 订阅失败 " + item.topic + ": " + e.message);
  351.         }
  352.     }
  353. }

  354. // ========== MQTT 回调实现 ==========
  355. function createMqttCallback() {
  356.     return new MqttCallback({
  357.         connectionLost: function (cause) {
  358.             G.connected = false;
  359.             var msg = "unknown";
  360.             try { msg = String(cause.getMessage()); } catch (e) {}
  361.             printl("[WARN] 连接断开: " + msg);
  362.         },

  363.         messageArrived: function (topic, message) {
  364.             try {
  365.                 handleMessage("" + topic, message);
  366.             } catch (e) {
  367.                 printl("[ERROR] 消息处理异常: " + e.message);
  368.             }
  369.         },

  370.         deliveryComplete: function (token) {
  371.             printl("[DEBUG] 消息投递完成");
  372.         }
  373.     });
  374. }

  375. // ========== 连接失败统一诊断打印(分阶段定位 + JS/Java 双堆栈) ==========
  376. function logConnErr(stage, e) {
  377.     var errMsg = "" + e;
  378.     try {
  379.         if (e instanceof MqttException) {
  380.             errMsg = "MQTT Error " + e.getReasonCode() + ": " + e.getMessage();
  381.         }
  382.     } catch (e2) {}
  383.     printl("[ERROR] " + stage + " 失败: " + errMsg);
  384.     // Rhino 的 e.stack 指向脚本行号
  385.     try { printl("[ERROR] " + stage + " JS堆栈: " + (e.stack || "(无)")); } catch (e3) {}
  386.     // e.javaException 是底层 Java Throwable(Rhino 包装对象本身不能直接 printStackTrace)
  387.     try {
  388.         var je = e.javaException;
  389.         if (je) {
  390.             var sw = new java.io.StringWriter();
  391.             var pw = new java.io.PrintWriter(sw);
  392.             je.printStackTrace(pw);
  393.             pw.flush();
  394.             printl("[ERROR] " + stage + " Java堆栈: " + sw.toString());
  395.         }
  396.     } catch (e4) {}
  397. }

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

  402.     if (!CFG.clientId) {
  403.         CFG.clientId = generateClientId();
  404.     }
  405.     printl("[INFO] 客户端 ID: " + CFG.clientId);

  406.     // ---- 阶段1:创建 MqttClient(构造器里有 log.fine,是 NPE 高发点) ----
  407.     try {
  408.         var persistence = new MemoryPersistence();
  409.         var brokerUri = CFG.brokerUrl;
  410.         if (brokerUri.indexOf("://") < 0) {
  411.             brokerUri = "tcp://" + brokerUri;
  412.         }
  413.         G.mqttClient = new MqttClient(brokerUri, CFG.clientId, persistence);
  414.         printl("[INFO] [阶段1/3] MqttClient 创建成功");
  415.     } catch (e1) {
  416.         G.connected = false;
  417.         logConnErr("阶段1 创建MqttClient", e1);
  418.         return false;
  419.     }

  420.     // ---- 阶段2:配置连接选项 ----
  421.     var options = null;
  422.     try {
  423.         options = new MqttConnectOptions();
  424.         options.setKeepAliveInterval(CFG.keepAliveSec);
  425.         options.setConnectionTimeout(CFG.connectTimeoutSec);
  426.         options.setCleanSession(CFG.cleanSession);

  427.         if (CFG.autoReconnect) {
  428.             options.setAutomaticReconnect(true);
  429.             try {
  430.                 options.setMaxReconnectDelay(CFG.maxReconnectDelayMs);
  431.             } catch (e) {}
  432.         }
  433.         if (CFG.username) {
  434.             options.setUserName(CFG.username);
  435.         }
  436.         if (CFG.password) {
  437.             options.setPassword(CFG.password.toCharArray());
  438.         }
  439.         if (CFG.willTopic && CFG.willMessage) {
  440.             options.setWill(
  441.                 CFG.willTopic,
  442.                 new java.lang.String(CFG.willMessage).getBytes("UTF-8"),
  443.                 CFG.willQos,
  444.                 CFG.willRetained
  445.             );
  446.         }
  447.         printl("[INFO] [阶段2/3] 连接选项配置完成");
  448.     } catch (e2) {
  449.         G.connected = false;
  450.         logConnErr("阶段2 配置选项", e2);
  451.         return false;
  452.     }

  453.     // ---- 阶段3:设置回调并执行连接(网络握手在这里) ----
  454.     try {
  455.         G.mqttClient.setCallback(createMqttCallback());
  456.         G.mqttClient.connect(options);

  457.         G.connected = true;
  458.         G.reconnectCount = 0;
  459.         printl("[INFO] [阶段3/3] MQTT 连接成功 — 握手完成");

  460.         try { subscribeTopics(); } catch (eS) { printl("[WARN] 订阅异常: " + eS.message); }
  461.         try {
  462.             publish(getPubTopic("online"), {
  463.                 clientId: CFG.clientId,
  464.                 device: getDeviceInfo(),
  465.                 ts: nowSec()
  466.             }, 0);
  467.         } catch (eP) { printl("[WARN] 上线消息发布异常: " + eP.message); }

  468.         return true;
  469.     } catch (e3) {
  470.         G.connected = false;
  471.         logConnErr("阶段3 执行连接", e3);
  472.         return false;
  473.     }
  474. }

  475. // ========== 断开连接 ==========
  476. function disconnect() {
  477.     if (G.mqttClient) {
  478.         try {
  479.             if (G.connected) {
  480.                 publish(getPubTopic("offline"), {
  481.                     clientId: CFG.clientId,
  482.                     ts: nowSec()
  483.                 }, 0);
  484.             }
  485.             G.mqttClient.disconnect();
  486.             printl("[INFO] 已断开 MQTT 连接");
  487.         } catch (e) {
  488.             printl("[WARN] 断开异常: " + e.message);
  489.         }
  490.     }
  491.     G.connected = false;
  492. }

  493. // ========== 关闭清理 ==========
  494. function cleanup() {
  495.     G.running = false;
  496.     disconnect();
  497.     try {
  498.         if (G.mqttClient) {
  499.             G.mqttClient.close(true);
  500.         }
  501.     } catch (e) {}
  502.     printl("[INFO] 资源已清理");
  503. }

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

  508.         printl("[INFO] ==========================================");
  509.         printl("[INFO]   MQTT 客户端 (Eclipse Paho) 启动");
  510.         printl("[INFO]   Broker: " + CFG.brokerUrl);
  511.         printl("[INFO]   ClientID: " + (CFG.clientId || "自动生成"));
  512.         printl("[INFO]   心跳: " + CFG.keepAliveSec + "s");
  513.         printl("[INFO]   自动重连: " + CFG.autoReconnect);
  514.         printl("[INFO]   CleanSession: " + CFG.cleanSession);
  515.         printl("[INFO]   订阅主题: " + CFG.subscribeTopics.length + " 个");
  516.         printl("[INFO]   设备信息: " + JSON.stringify(getDeviceInfo()));
  517.         printl("[INFO] ==========================================");

  518.         if (!connect()) {
  519.             startReconnectLoop();
  520.         }

  521.         var tick = 0;
  522.         (function keepAlive() {
  523.             try {
  524.                 if (!G.running) {
  525.                     printl("[INFO] 主线程保活退出");
  526.                     return;
  527.                 }
  528.                 tick++;
  529.                 if (tick % 30 === 0) {
  530.                     printl("[INFO] [保活] 已运行=" + (nowSec() - G.startTime) +
  531.                         "s 连接=" + G.connected +
  532.                         " 发送=" + G.msgSentCount +
  533.                         " 接收=" + G.msgRecvCount +
  534.                         " 重连=" + G.reconnectCount);
  535.                 }
  536.             } catch (e) {
  537.                 printl("[ERROR] 保活异常(已兜住): " + e.message);
  538.             }
  539.             try { setTimeout(keepAlive, 1000); }
  540.             catch (e) {
  541.                 try { runTime.setTimeout(keepAlive, 1000); } catch (e2) {}
  542.             }
  543.         })();

  544.         if (CFG.runDuration > 0) {
  545.             sleep.millisecond(CFG.runDuration);
  546.             cleanup();
  547.         }

  548.     } catch (e) {
  549.         printl("[FATAL] 主程序异常: " + e.message + "\n" + (e.stack || ""));
  550.         try { cleanup(); } catch (e2) {}
  551.     }
  552. }

  553. // ========== 手动重连循环 ==========
  554. function startReconnectLoop() {
  555.     (function reconnectStep() {
  556.         try {
  557.             if (!G.running) return;
  558.             if (G.connected) return;

  559.             G.reconnectCount++;
  560.             printl("[INFO] 手动重连尝试 " + G.reconnectCount);

  561.             if (connect()) {
  562.                 printl("[INFO] 手动重连成功");
  563.                 return;
  564.             }

  565.             setTimeout(reconnectStep, CFG.reconnectDelayMs);
  566.         } catch (e) {
  567.             printl("[ERROR] 重连异常(已兜住): " + e.message);
  568.             setTimeout(reconnectStep, CFG.reconnectDelayMs);
  569.         }
  570.     })();
  571. }

  572. // ========== 启动 ==========
  573. main();
复制代码







欢迎光临 B2B网络软件 (http://bbs.niubt.cn/) Powered by Discuz! X3.2