const path = require('path') const WebSocket = require('ws') const config = require('../../config.json') const Logger = require('../../lib/Logger') const db = require('../DataBase/db.js') const EmailTemplate = require('../Email/emailTemplate') const { ConversationService } = require('../../lib/AIChat/ConversationService') const logger = new Logger(path.join(__dirname, '../../logs/OneBotV11.log'), 'INFO') const BOT_KIND_WORKORDER = 'workorder' const BOT_KIND_AI_CHAT = 'aiChat' function createState(kind) { return { kind, wsClient: null, connectingPromise: null, reconnectTimer: null, wsApiClient: null, connectingApiPromise: null, apiReconnectTimer: null } } const botStates = { [BOT_KIND_WORKORDER]: createState(BOT_KIND_WORKORDER), [BOT_KIND_AI_CHAT]: createState(BOT_KIND_AI_CHAT) } function stringifySafe(data) { try { return JSON.stringify(data) } catch (_) { return '[unserializable]' } } function getRawConfig(kind) { if (kind === BOT_KIND_AI_CHAT) return config.aiChatOnebotv11 || {} return config.onebotv11 || {} } function getOneBotConfig(kind = BOT_KIND_WORKORDER) { const onebot = getRawConfig(kind) const isAiChat = kind === BOT_KIND_AI_CHAT const defaults = isAiChat ? { reverseWsPort: 15701, senderNickname: 'Xiaoyan Assistant Web User', botName: 'Xiaoyan Assistant', botUuid: 'onebot-v11-xiaoyan-chat' } : { reverseWsPort: 15700, senderNickname: 'Ticket System', botName: 'Xiaoyan Assistant', botUuid: 'onebot-v11-xiaoyan-assistant' } return { enabled: onebot.enabled === true, transport: onebot.transport || 'reverse_ws', reverseWsUrl: onebot.reverseWsUrl || '', reverseWsHost: onebot.reverseWsHost || '127.0.0.1', reverseWsPort: Number(onebot.reverseWsPort || defaults.reverseWsPort), reverseWsPath: onebot.reverseWsPath || '', reverseWsToken: onebot.reverseWsToken || '', callbackToken: onebot.callbackToken || '', ticketSenderNickname: onebot.ticketSenderNickname || defaults.senderNickname, chatSenderNickname: onebot.chatSenderNickname || onebot.ticketSenderNickname || defaults.senderNickname, selfId: onebot.selfId || '', botName: onebot.botName || defaults.botName, botUuid: onebot.botUuid || defaults.botUuid } } function getWorkOrderOneBotConfig() { return getOneBotConfig(BOT_KIND_WORKORDER) } function getAiChatOneBotConfig() { return getOneBotConfig(BOT_KIND_AI_CHAT) } function getReverseWsUrl(onebot) { const defaultPath = '/ws' const explicitPath = onebot.reverseWsPath ? (onebot.reverseWsPath.startsWith('/') ? onebot.reverseWsPath : `/${onebot.reverseWsPath}`) : null if (onebot.reverseWsUrl) { try { const u = new URL(onebot.reverseWsUrl) if (explicitPath) { u.pathname = explicitPath return u.href } if (u.pathname && u.pathname !== '/') return u.href u.pathname = defaultPath return u.href } catch (_) { return onebot.reverseWsUrl } } const pathPart = explicitPath || defaultPath return `ws://${onebot.reverseWsHost}:${onebot.reverseWsPort}${pathPart}` } function getState(kind) { return botStates[kind] || botStates[BOT_KIND_WORKORDER] } function getBotLabel(kind) { return kind === BOT_KIND_AI_CHAT ? 'AIChat OneBot v11' : 'WorkOrder OneBot v11' } function scheduleReconnect(kind) { const state = getState(kind) if (state.reconnectTimer) return const onebot = getOneBotConfig(kind) logger.info(`${getBotLabel(kind)} ws reconnect scheduled in 3s`) state.reconnectTimer = setTimeout(() => { state.reconnectTimer = null connectReverseWs(kind).catch(err => { logger.error(`${getBotLabel(kind)} reconnect failed: ${err.message}`) }) }, 3000) if (!onebot.enabled && state.reconnectTimer) { clearTimeout(state.reconnectTimer) state.reconnectTimer = null } } function scheduleApiReconnect(kind) { const state = getState(kind) if (state.apiReconnectTimer) return const onebot = getOneBotConfig(kind) logger.info(`${getBotLabel(kind)} API ws reconnect scheduled in 3s`) state.apiReconnectTimer = setTimeout(() => { state.apiReconnectTimer = null connectApiReverseWs(kind).catch(err => { logger.error(`${getBotLabel(kind)} API ws reconnect failed: ${err.message}`) }) }, 3000) if (!onebot.enabled && state.apiReconnectTimer) { clearTimeout(state.apiReconnectTimer) state.apiReconnectTimer = null } } async function connectReverseWs(kind = BOT_KIND_WORKORDER) { const state = getState(kind) if (state.wsClient && state.wsClient.readyState === WebSocket.OPEN) return state.wsClient if (state.connectingPromise) return state.connectingPromise const onebot = getOneBotConfig(kind) if (!onebot.enabled) throw new Error(`${getBotLabel(kind)} is disabled`) if (onebot.transport !== 'reverse_ws') throw new Error(`${getBotLabel(kind)} only supports reverse_ws`) const wsUrl = getReverseWsUrl(onebot) logger.info(`${getBotLabel(kind)} connecting reverse ws: ${wsUrl}`) if (!onebot.selfId) throw new Error(`${getBotLabel(kind)} missing selfId`) const selfIdStr = String(onebot.selfId).trim() if (!selfIdStr) throw new Error(`${getBotLabel(kind)} selfId is empty`) const headers = { 'User-Agent': `runforge/1.0 (${kind}-onebot-v11-reverse-ws-client)`, 'X-Client-Role': `${kind}:event`, 'X-Self-ID': selfIdStr } if (onebot.reverseWsToken) headers.Authorization = `Bearer ${onebot.reverseWsToken}` state.connectingPromise = new Promise((resolve, reject) => { const client = new WebSocket(wsUrl, { headers }) let opened = false const connTimeout = setTimeout(() => { if (!opened) { try { client.terminate() } catch (_) { } reject(new Error(`${getBotLabel(kind)} connect timeout: ${wsUrl}`)) } }, 5000) client.on('unexpected-response', (_req, res) => { const statusCode = res && res.statusCode const statusMessage = res && res.statusMessage const rh = res && typeof res.getHeaders === 'function' ? res.getHeaders() : {} const rawPreview = res && Array.isArray(res.rawHeaders) ? res.rawHeaders.slice(0, 24) : [] logger.error(`${getBotLabel(kind)} ws unexpected-response: ${statusCode} ${statusMessage}`) logger.error(`${getBotLabel(kind)} ws response headers: ${stringifySafe({ rh, rawPreview })}`) }) client.on('open', () => { clearTimeout(connTimeout) opened = true state.wsClient = client logger.info(`${getBotLabel(kind)} reverse ws connected: ${wsUrl}, selfId=${onebot.selfId}`) client.send(JSON.stringify(buildLifecycleConnectEvent(onebot)), (err) => { if (err) { logger.error(`${getBotLabel(kind)} lifecycle event failed: ${err.message}`) return } logger.info(`${getBotLabel(kind)} lifecycle connect event sent`) }) resolve(client) }) client.on('error', (err) => { clearTimeout(connTimeout) if (!opened) reject(err) logger.error(`${getBotLabel(kind)} ws error: ${err.message}`) }) client.on('close', () => { state.wsClient = null state.connectingPromise = null logger.error(`${getBotLabel(kind)} ws closed`) if (getOneBotConfig(kind).enabled) scheduleReconnect(kind) }) client.on('message', (raw) => { const text = typeof raw === 'string' ? raw : raw.toString('utf8') logger.info(`${getBotLabel(kind)} received: ${text.slice(0, 1000)}`) handleWsIncomingFrame(text, client, kind).catch((err) => { logger.error(`${getBotLabel(kind)} incoming frame failed: ${err.message}`) }) }) }) try { return await state.connectingPromise } finally { state.connectingPromise = null } } async function connectApiReverseWs(kind = BOT_KIND_WORKORDER) { const state = getState(kind) if (state.wsApiClient && state.wsApiClient.readyState === WebSocket.OPEN) return state.wsApiClient if (state.connectingApiPromise) return state.connectingApiPromise const onebot = getOneBotConfig(kind) if (!onebot.enabled) throw new Error(`${getBotLabel(kind)} is disabled`) if (onebot.transport !== 'reverse_ws') throw new Error(`${getBotLabel(kind)} only supports reverse_ws`) const wsUrl = getReverseWsUrl(onebot) if (!onebot.selfId) throw new Error(`${getBotLabel(kind)} API ws missing selfId`) const selfIdStr = String(onebot.selfId).trim() if (!selfIdStr) throw new Error(`${getBotLabel(kind)} API ws selfId is empty`) const headers = { 'User-Agent': `runforge/1.0 (${kind}-onebot-v11-reverse-ws-api-client)`, 'X-Client-Role': `${kind}:api`, 'X-Self-ID': selfIdStr } if (onebot.reverseWsToken) headers.Authorization = `Bearer ${onebot.reverseWsToken}` state.connectingApiPromise = new Promise((resolve, reject) => { const client = new WebSocket(wsUrl, { headers }) let opened = false const connTimeout = setTimeout(() => { if (!opened) { try { client.terminate() } catch (_) { } reject(new Error(`${getBotLabel(kind)} API ws connect timeout: ${wsUrl}`)) } }, 5000) client.on('open', () => { clearTimeout(connTimeout) opened = true state.wsApiClient = client logger.info(`${getBotLabel(kind)} API ws connected: ${wsUrl}`) resolve(client) }) client.on('error', (err) => { clearTimeout(connTimeout) if (!opened) reject(err) logger.error(`${getBotLabel(kind)} API ws error: ${err.message}`) }) client.on('close', () => { state.wsApiClient = null state.connectingApiPromise = null logger.error(`${getBotLabel(kind)} API ws closed`) if (getOneBotConfig(kind).enabled) scheduleApiReconnect(kind) }) client.on('message', (raw) => { const text = typeof raw === 'string' ? raw : raw.toString('utf8') logger.info(`${getBotLabel(kind)} API ws received: ${text.slice(0, 1000)}`) handleWsIncomingFrame(text, client, kind).catch((err) => { logger.error(`${getBotLabel(kind)} API frame failed: ${err.message}`) }) }) }) try { return await state.connectingApiPromise } finally { state.connectingApiPromise = null } } function buildSyntheticPrivateMessageEvent(text, onebot, orderId) { const now = Math.floor(Date.now() / 1000) const selfId = String(onebot.selfId).trim() const senderIdStr = String(orderId || '').trim() if (!senderIdStr || !/^\d+$/.test(senderIdStr) || senderIdStr === selfId) { throw new Error('Invalid orderId for OneBot user_id') } const messageArr = [{ type: 'text', data: { text } }] const messageId = Math.floor(Math.random() * 2000000000) + 1 return { time: now, self_id: selfId, post_type: 'message', message_type: 'private', sub_type: 'friend', message_id: messageId, user_id: senderIdStr, message: messageArr, raw_message: text, font: 0, sender: { user_id: senderIdStr, nickname: onebot.ticketSenderNickname || 'Ticket System', sex: 'unknown', age: 0 } } } function buildSyntheticAiChatMessageEvent(payload, onebot) { const now = Math.floor(Date.now() / 1000) const selfId = String(onebot.selfId).trim() const senderIdStr = String(payload.conversationNo || payload.conversationId || '').trim() if (!senderIdStr || !/^\d+$/.test(senderIdStr) || senderIdStr === selfId) { throw new Error('Invalid conversationId for OneBot user_id') } const lines = [ `[AIChat#${senderIdStr}]`, `user_uuid: ${payload.senderUuid || '-'}`, 'content:', payload.content || '' ] const messageArr = [{ type: 'text', data: { text: lines.join('\n') } }] for (const url of payload.images || []) { messageArr.push({ type: 'image', data: { file: url } }) } const rawImages = (payload.images || []).length > 0 ? `\nimages:\n${(payload.images || []).join('\n')}` : '' const rawMessage = lines.join('\n') + rawImages const messageId = Math.floor(Math.random() * 2000000000) + 1 return { time: now, self_id: selfId, post_type: 'message', message_type: 'private', sub_type: 'friend', message_id: messageId, user_id: senderIdStr, message: messageArr, raw_message: rawMessage, font: 0, sender: { user_id: senderIdStr, nickname: onebot.chatSenderNickname || 'Xiaoyan Assistant Web User', sex: 'unknown', age: 0 } } } function buildLifecycleConnectEvent(onebot) { return { time: Math.floor(Date.now() / 1000), self_id: String(onebot.selfId).trim(), post_type: 'meta_event', meta_event_type: 'lifecycle', sub_type: 'connect' } } async function sendImplementationEvent(eventObj, kind = BOT_KIND_WORKORDER) { const client = await connectReverseWs(kind) const line = stringifySafe(eventObj) logger.info(`${getBotLabel(kind)} reporting event post_type=${eventObj.post_type} message_type=${eventObj.message_type} len=${line.length}`) return new Promise((resolve, reject) => { client.send(line, (err) => { if (err) { logger.error(`${getBotLabel(kind)} event send failed: ${err.message}`) return reject(err) } logger.info(`${getBotLabel(kind)} event written to WebSocket`) resolve(true) }) }) } function decodeMessageText(message) { if (!message) return '' if (typeof message === 'string') return message if (Array.isArray(message)) { return message.map(seg => { if (typeof seg === 'string') return seg if (seg && seg.type === 'text' && seg.data && typeof seg.data.text === 'string') return seg.data.text return '' }).join('') } return String(message) } async function persistAiReplyAsServerReply(orderId, plainText) { const parsedOrderId = Number(orderId) if (!Number.isInteger(parsedOrderId) || parsedOrderId <= 0) return false const rows = await db.query('SELECT msg, state, email FROM work_order WHERE id = ?', [parsedOrderId]) if (!rows || rows.length !== 1 || rows[0].state === 2) return false const now = Date.now() const msg = rows[0].msg || [] msg.push({ time: now, content: plainText, files: [], uuid: 'e4fe0277-0b1a-41a1-b25f-8b6e4cec3281', type: 'ai' }) await db.query('UPDATE work_order SET msg = ?, update_time = ?, state = 3 WHERE id = ?', [msg, now, parsedOrderId]) if (rows[0].email) { await EmailTemplate.orderNewReply(rows[0].email, { id: parsedOrderId, content: plainText, files: [] }) } return true } async function persistAiChatReply(conversationId, plainText) { const parsedConversationNo = String(conversationId || '').trim() if (!/^\d+$/.test(parsedConversationNo)) return false return await ConversationService.addAssistantMessage({ conversationNo: parsedConversationNo, content: plainText }) } function ackOneBotApiRequest(j, client, kind) { if (j.action === undefined || j.echo === undefined || typeof j.echo !== 'object' || j.echo === null || j.echo.seq === undefined) { return false } const resp = { status: 'ok', retcode: 0, data: null, echo: j.echo } try { client.send(JSON.stringify(resp)) logger.info(`${getBotLabel(kind)} API request acked: action=${j.action} seq=${JSON.stringify(j.echo.seq)}`) } catch (err) { logger.error(`${getBotLabel(kind)} API ack failed: ${err.message}`) } return true } async function handleWsIncomingFrame(text, client, kind = BOT_KIND_WORKORDER) { let j try { j = JSON.parse(text) } catch (_) { return } if (j == null || typeof j !== 'object') return if (j.action === 'send_private_msg' && j.params) { const plainText = decodeMessageText(j.params.message) logger.info(`${getBotLabel(kind)} decoded private message: ${plainText}`) if (kind === BOT_KIND_AI_CHAT) { await persistAiChatReply(j.params.user_id, plainText) } else { await persistAiReplyAsServerReply(j.params.user_id, plainText) } } if (ackOneBotApiRequest(j, client, kind)) return if ('post_type' in j) return if (j.echo !== undefined && j.status !== undefined) { if (j.retcode !== 0 && j.retcode !== undefined) { const wording = j.wording || j.message || stringifySafe(j.data) logger.error(`${getBotLabel(kind)} API failed echo=${j.echo} retcode=${j.retcode} ${wording}`) } else { logger.info(`${getBotLabel(kind)} API success echo=${j.echo}`) } } } function buildOrderText(payload) { const { orderId, title, role, content, files = [], senderUuid } = payload const fileText = files.length > 0 ? `\nattachments:\n${files.join('\n')}` : '' return [ `[Ticket#${orderId}] ${title ? title : 'Ticket message updated'}`, `sender: ${role === '用户' ? 'user' : 'system'}`, `user_uuid: ${senderUuid || '-'}`, 'content:', content ].join('\n') + fileText } async function sendOrderMessage(payload) { const onebot = getWorkOrderOneBotConfig() if (!onebot.enabled) return false if (onebot.transport !== 'reverse_ws') { logger.error('WorkOrder OneBot v11 transport error: only reverse_ws is supported') return false } const text = buildOrderText(payload) const eventObj = buildSyntheticPrivateMessageEvent(text, onebot, payload.orderId) await sendImplementationEvent(eventObj, BOT_KIND_WORKORDER) logger.info(`WorkOrder OneBot event reported: orderId=${payload.orderId}, userId=${eventObj.user_id}, messageId=${eventObj.message_id}`) return true } async function sendAiChatMessage(payload) { const onebot = getAiChatOneBotConfig() if (!onebot.enabled) throw new Error('AIChat OneBot v11 is disabled') if (onebot.transport !== 'reverse_ws') { logger.error('AIChat OneBot v11 transport error: only reverse_ws is supported') throw new Error('AIChat OneBot v11 transport error') } const eventObj = buildSyntheticAiChatMessageEvent(payload, onebot) await sendImplementationEvent(eventObj, BOT_KIND_AI_CHAT) logger.info(`AIChat OneBot event reported: conversationId=${payload.conversationId}, conversationNo=${payload.conversationNo}, userId=${eventObj.user_id}, messageId=${eventObj.message_id}`) return true } async function initOneBotKind(kind) { const onebot = getOneBotConfig(kind) if (!onebot.enabled) { logger.info(`${getBotLabel(kind)} disabled, skip ws init`) return false } if (onebot.transport !== 'reverse_ws') { logger.error(`${getBotLabel(kind)} init failed: transport is not reverse_ws`) return false } await connectReverseWs(kind) await connectApiReverseWs(kind) return true } async function initOneBotWs() { const results = await Promise.allSettled([ initOneBotKind(BOT_KIND_WORKORDER), initOneBotKind(BOT_KIND_AI_CHAT) ]) let anyReady = false for (const [idx, result] of results.entries()) { const kind = idx === 0 ? BOT_KIND_WORKORDER : BOT_KIND_AI_CHAT if (result.status === 'fulfilled') { anyReady = anyReady || result.value === true continue } logger.error(`${getBotLabel(kind)} init failed: ${result.reason && result.reason.message ? result.reason.message : result.reason}`) } return anyReady } module.exports = { sendOrderMessage, sendAiChatMessage, getOneBotConfig, getWorkOrderOneBotConfig, getAiChatOneBotConfig, initOneBotWs }