Files
inquiry_robot/aibot-bridge/long_connection.cjs
T
jillion886andCursor d9029dec89 落地空运群内协同与锁舱出站。
群里按 BOT 通道激活当前工单、航线报价和成交跟进;没有 TMS 报价编号时用工单字段补齐仍调锁舱,避免被本地规则拦住。

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-09-16 17:25:42 +08:00

507 lines
14 KiB
JavaScript

/**
* AIBOT 长连接(只传输,无工单状态)。
*
* 本文件职责:
* - WECOM_AIBOT_ENABLED=true 时连 wss://openws.work.weixin.qq.com
* - 订阅后:入站 aibot_msg_callback 转发智能体;出站 /send-group 走 aibot_send_msg
* - 同一 Bot 只能一条长连接;并存期默认关,勿与现网 8813 双连
*
* 禁止:写工单/六态/Graph。密钥只读环境变量,不打日志。
*/
"use strict";
const crypto = require("crypto");
function generateReqId(prefix) {
const timestamp = Date.now();
const random = crypto.randomBytes(4).toString("hex");
return `${prefix}_${timestamp}_${random}`;
}
function loadWs() {
try {
return require("ws");
} catch (_) {
return null;
}
}
/**
* @param {object} opts
* @param {boolean} opts.enabled
* @param {(payloadText: string) => Promise<object>} opts.forwardToAgent
* @param {() => boolean} [opts.shouldStop]
*/
function startLongConnectionLoop(opts) {
const enabled = !!opts.enabled;
const forwardToAgent = opts.forwardToAgent;
const shouldStop = opts.shouldStop || (() => false);
const botId = String(process.env.WECOM_AIBOT_BOT_ID || "").trim();
const secret = String(process.env.WECOM_AIBOT_SECRET || "").trim();
const wsUrl =
process.env.WECOM_AIBOT_WS_URL || "wss://openws.work.weixin.qq.com";
const state = {
ws: null,
connected: false,
subscribed: false,
lastError: "",
pending: new Map(),
lastCallback: new Map(),
heartbeatTimer: null,
reconnectTimer: null,
stop: () => {},
};
function status() {
return {
enabled,
connected: state.connected,
subscribed: state.subscribed,
lastError: state.lastError,
hasCredentials: Boolean(botId && secret),
};
}
function sendWs(payload, reqId) {
return new Promise((resolve, reject) => {
if (!state.ws || state.ws.readyState !== 1 || !state.subscribed) {
const err = new Error("ws_offline");
err.code = "ws_offline";
reject(err);
return;
}
const timer = setTimeout(() => {
state.pending.delete(reqId);
const err = new Error("send_timeout");
err.code = "send_timeout";
reject(err);
}, 8000);
state.pending.set(reqId, { resolve, reject, timer });
state.ws.send(JSON.stringify(payload));
});
}
function sendGroupMessage({ chatId, content }) {
const text = String(content || "");
const callbackReqId = state.lastCallback.get(String(chatId || "").trim()) || "";
if (callbackReqId) {
return sendWs(
{
cmd: "aibot_respond_msg",
headers: { req_id: callbackReqId },
body: {
msgtype: "stream",
stream: { id: crypto.randomUUID(), finish: true, content: text },
},
},
callbackReqId
).finally(() => {
state.lastCallback.delete(String(chatId || "").trim());
});
}
const reqId = generateReqId("aibot_send_msg");
return sendWs(
{
cmd: "aibot_send_msg",
headers: { req_id: reqId },
body: {
chatid: chatId,
chat_type: 2,
msgtype: "markdown",
markdown: { content: text },
},
},
reqId
);
}
/**
* 空运群发报价模板:分片上传后 aibot_send_msg file。
* 单个分片不超过 512KB(编码前)。
*/
async function sendGroupFile({ chatId, filename, fileB64, fileBytes }) {
const buf = Buffer.isBuffer(fileBytes)
? fileBytes
: Buffer.from(String(fileB64 || ""), "base64");
const name = String(filename || "quote.xlsx").trim() || "quote.xlsx";
if (!chatId || !buf.length) {
const err = new Error("file_empty");
err.code = "file_empty";
throw err;
}
const chunkSize = 400 * 1024;
const totalChunks = Math.max(1, Math.ceil(buf.length / chunkSize));
const md5 = crypto.createHash("md5").update(buf).digest("hex");
const initReq = generateReqId("aibot_upload_media_init");
const init = await sendWs(
{
cmd: "aibot_upload_media_init",
headers: { req_id: initReq },
body: {
type: "file",
filename: name,
total_size: buf.length,
total_chunks: totalChunks,
md5,
},
},
initReq
);
const uploadId = String((init.body && init.body.upload_id) || "").trim();
if (!uploadId) {
const err = new Error("upload_init_failed");
err.code = "upload_init_failed";
throw err;
}
for (let i = 0; i < totalChunks; i += 1) {
const slice = buf.subarray(i * chunkSize, Math.min(buf.length, (i + 1) * chunkSize));
const chunkReq = generateReqId("aibot_upload_media_chunk");
await sendWs(
{
cmd: "aibot_upload_media_chunk",
headers: { req_id: chunkReq },
body: {
upload_id: uploadId,
chunk_index: i,
base64_data: slice.toString("base64"),
},
},
chunkReq
);
}
const finishReq = generateReqId("aibot_upload_media_finish");
const finished = await sendWs(
{
cmd: "aibot_upload_media_finish",
headers: { req_id: finishReq },
body: { upload_id: uploadId },
},
finishReq
);
const mediaId = String((finished.body && finished.body.media_id) || "").trim();
if (!mediaId) {
const err = new Error("upload_finish_failed");
err.code = "upload_finish_failed";
throw err;
}
const sendReq = generateReqId("aibot_send_msg");
return sendWs(
{
cmd: "aibot_send_msg",
headers: { req_id: sendReq },
body: {
chatid: chatId,
chat_type: 2,
msgtype: "file",
file: { media_id: mediaId },
},
},
sendReq
);
}
if (!enabled) {
// eslint-disable-next-line no-console
console.log("[aibot-long-connection] 未启用(WECOM_AIBOT_ENABLED=false),不连企微");
return { stop: () => {}, sendGroupMessage, sendGroupFile, status };
}
if (!botId || !secret) {
state.lastError = "missing_bot_credentials";
// eslint-disable-next-line no-console
console.warn(
"[aibot-long-connection] 已启用但缺少 WECOM_AIBOT_BOT_ID / WECOM_AIBOT_SECRET,无法订阅"
);
return { stop: () => {}, sendGroupMessage, sendGroupFile, status };
}
const WebSocket = loadWs();
if (!WebSocket) {
state.lastError = "ws_package_missing";
// eslint-disable-next-line no-console
console.warn("[aibot-long-connection] 未安装 ws 包,请在 aibot-bridge 目录执行 npm install");
return { stop: () => {}, sendGroupMessage, sendGroupFile, status };
}
function startHeartbeat() {
stopHeartbeat();
state.heartbeatTimer = setInterval(() => {
if (!state.ws || state.ws.readyState !== 1) {
return;
}
try {
state.ws.send(
JSON.stringify({
cmd: "ping",
headers: { req_id: generateReqId("ping") },
})
);
} catch (_) {
/* ignore */
}
}, 25000);
}
function stopHeartbeat() {
if (state.heartbeatTimer) {
clearInterval(state.heartbeatTimer);
state.heartbeatTimer = null;
}
}
function scheduleReconnect() {
if (shouldStop() || state.reconnectTimer) {
return;
}
state.reconnectTimer = setTimeout(() => {
state.reconnectTimer = null;
connect();
}, 5000);
}
function settlePending(reqId, errcode, errmsg, frame) {
const wait = state.pending.get(reqId);
if (!wait) {
return;
}
clearTimeout(wait.timer);
state.pending.delete(reqId);
if (Number(errcode) === 0) {
wait.resolve({ ok: true, body: (frame && frame.body) || {}, frame: frame || {} });
return;
}
const err = new Error(errmsg || "send_failed");
err.code = "wecom_error";
err.errcode = errcode;
wait.reject(err);
}
function normalizeInbound(frame) {
const body = frame.body || {};
const from = body.from || {};
const text = (body.text && body.text.content) || "";
return {
sender_id: String(from.userid || "").trim(),
message_id: String(body.msgid || (frame.headers && frame.headers.req_id) || "").trim(),
content: text,
msg_type: String(body.msgtype || "text"),
chat_id: String(body.chatid || "").trim(),
chat_type: body.chattype === "single" ? "c2c" : "group",
raw: {
source: "aibot",
aibot_req_id: frame.headers && frame.headers.req_id,
},
};
}
function handleFrame(raw) {
let frame;
try {
frame = JSON.parse(raw);
} catch (_) {
return;
}
const reqId = frame.headers && frame.headers.req_id;
if (reqId && state.pending.has(reqId)) {
settlePending(reqId, frame.errcode, frame.errmsg, frame);
}
if (frame.cmd === "aibot_event_callback") {
const ev = (frame.body && frame.body.event && frame.body.event.eventtype) || "";
// eslint-disable-next-line no-console
console.warn("[aibot-long-connection] 事件", ev || "-");
if (ev === "disconnected_event") {
state.lastError = "kicked_by_new_connection";
}
return;
}
const bodyHint = frame.body || {};
const looksLikeInbound = Boolean(
frame.cmd === "aibot_msg_callback" ||
bodyHint.msgid ||
bodyHint.chatid ||
(bodyHint.text && bodyHint.text.content)
);
if (looksLikeInbound) {
// eslint-disable-next-line no-console
console.log(
"[aibot-long-connection] 收到入站 cmd=",
frame.cmd || "-",
"msgid=",
bodyHint.msgid || "-",
"chat=",
bodyHint.chatid || "-"
);
const inbound = normalizeInbound(frame);
if (reqId && inbound.chat_id) {
state.lastCallback.set(inbound.chat_id, reqId);
}
const officialBody = Object.assign({}, frame.body || {}, {
aibot_direct: true,
aibot_req_id: reqId || "",
});
const payload = inbound.sender_id
? Object.assign({}, officialBody, {
sender_id: inbound.sender_id,
message_id: inbound.message_id,
content: inbound.content,
chat_id: inbound.chat_id,
chat_type: inbound.chat_type,
})
: officialBody;
if (!payload.sender_id && !payload.from) {
return;
}
Promise.resolve(forwardToAgent(JSON.stringify(payload))).catch((e) => {
// eslint-disable-next-line no-console
console.warn("[aibot-long-connection] 转发入站失败", String(e && e.message ? e.message : e));
});
}
}
function connect() {
if (shouldStop()) {
return;
}
// eslint-disable-next-line no-console
console.warn(
"[aibot-long-connection] 即将连接企微长连接;请确认现网 8813 已停,避免双连"
);
const ws = new WebSocket(wsUrl);
state.ws = ws;
ws.on("open", () => {
state.connected = true;
state.lastError = "";
const reqId = generateReqId("aibot_subscribe");
ws.send(
JSON.stringify({
cmd: "aibot_subscribe",
headers: { req_id: reqId },
body: { bot_id: botId, secret },
})
);
// eslint-disable-next-line no-console
console.log("[aibot-long-connection] 已发订阅 req_id前缀=aibot_subscribe");
});
ws.on("message", (buf) => {
const text = Buffer.isBuffer(buf) ? buf.toString("utf8") : String(buf);
let frame;
try {
frame = JSON.parse(text);
} catch (_) {
// eslint-disable-next-line no-console
console.warn("[aibot-long-connection] 非JSON帧 len=", text.length);
handleFrame(text);
return;
}
// eslint-disable-next-line no-console
const reqPrefix = String((frame.headers && frame.headers.req_id) || "").split("_").slice(0, 2).join("_");
console.log(
"[aibot-long-connection] 帧 cmd=",
frame.cmd || "-",
"req=",
reqPrefix || "-",
"keys=",
Object.keys(frame).join(","),
"errcode=",
frame.errcode,
"errmsg=",
frame.errmsg || ""
);
if (frame.cmd === "ping") {
try {
ws.send(
JSON.stringify({
cmd: "pong",
headers: frame.headers || {},
})
);
} catch (_) {
/* ignore */
}
}
const ackReq = String((frame.headers && frame.headers.req_id) || "");
if (ackReq.startsWith("aibot_subscribe") && Number(frame.errcode) === 0 && !state.subscribed) {
state.subscribed = true;
startHeartbeat();
// eslint-disable-next-line no-console
console.log("[aibot-long-connection] 订阅成功");
} else if (
ackReq.startsWith("aibot_subscribe") &&
frame.errcode !== undefined &&
Number(frame.errcode) !== 0
) {
state.lastError = "subscribe_failed_" + String(frame.errcode);
// eslint-disable-next-line no-console
console.warn("[aibot-long-connection] 订阅失败 errcode=", frame.errcode);
}
handleFrame(text);
});
ws.on("close", (code, reason) => {
stopHeartbeat();
state.connected = false;
state.subscribed = false;
state.ws = null;
// eslint-disable-next-line no-console
console.warn(
"[aibot-long-connection] 连接关闭 code=",
code,
"reason=",
reason ? String(reason).slice(0, 120) : ""
);
if (!shouldStop()) {
scheduleReconnect();
}
});
ws.on("error", (e) => {
state.lastError = "ws_error";
// eslint-disable-next-line no-console
console.warn("[aibot-long-connection] 连接错误", String(e && e.message ? e.message : e));
});
ws.on("ping", () => {
try {
ws.pong();
} catch (_) {
/* ignore */
}
});
}
connect();
state.stop = () => {
stopHeartbeat();
if (state.reconnectTimer) {
clearTimeout(state.reconnectTimer);
state.reconnectTimer = null;
}
if (state.ws) {
try {
state.ws.close();
} catch (_) {
/* ignore */
}
state.ws = null;
}
state.connected = false;
state.subscribed = false;
};
return {
stop: state.stop,
sendGroupMessage,
sendGroupFile,
status,
/**
* 测试注入:模拟 SDK 收到一帧已解析 JSON,只做转发。
* @param {object|string} frame
*/
injectFrameForTest: async (frame) => {
const text = typeof frame === "string" ? frame : JSON.stringify(frame);
return forwardToAgent(text);
},
};
}
module.exports = { startLongConnectionLoop };