| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602 |
- 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 || '',
- eventClientRole: onebot.eventClientRole || 'Universal',
- apiClientRole: onebot.apiClientRole || 'Api',
- enableApiWs: onebot.enableApiWs !== false,
- 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 (OneBot-v11-reverse-ws-client)',
- 'X-Client-Role': onebot.eventClientRole,
- '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 (OneBot-v11-reverse-ws-api-client)',
- 'X-Client-Role': onebot.apiClientRole,
- '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)
- if (onebot.enableApiWs) {
- await connectApiReverseWs(kind)
- } else {
- logger.info(`${getBotLabel(kind)} API ws disabled by config`)
- }
- 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
- }
|