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]