WeixinBindingService.js 15 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384
  1. const crypto = require('crypto')
  2. const db = require('../../plugin/DataBase/db')
  3. const Logger = require('../Logger')
  4. const config = require('../../config.json')
  5. const { ConversationService, normalizeText } = require('./ConversationService')
  6. const { consumeQuota } = require('./QuotaService')
  7. const WeixinBotClient = require('./WeixinBotClient')
  8. const OneBotV11 = require('../../plugin/OneBot/OneBotV11')
  9. const { getUserVipInfo } = require('../VipService')
  10. const STATE_PENDING = 0
  11. const STATE_ACTIVE = 1
  12. const STATE_DISABLED = 2
  13. const STATE_EXPIRED = 3
  14. const WECHAT_EXPIRE_MS = 24 * 60 * 60 * 1000
  15. const WECHAT_REMIND_MS = 23 * 60 * 60 * 1000
  16. const logger = new Logger()
  17. let pollTimer = null
  18. let running = false
  19. function keyBuffer() {
  20. const raw = config.weixinBot?.tokenAesKey || config.qk?.passwordAesKey || 'runforge-weixin-token-key'
  21. return crypto.createHash('sha256').update(String(raw)).digest()
  22. }
  23. function encryptToken(text) {
  24. if (!text) return ''
  25. const iv = crypto.randomBytes(12)
  26. const cipher = crypto.createCipheriv('aes-256-gcm', keyBuffer(), iv)
  27. const encrypted = Buffer.concat([cipher.update(String(text), 'utf8'), cipher.final()])
  28. const tag = cipher.getAuthTag()
  29. return `${iv.toString('base64')}.${tag.toString('base64')}.${encrypted.toString('base64')}`
  30. }
  31. function decryptToken(value) {
  32. if (!value) return ''
  33. const [ivRaw, tagRaw, encryptedRaw] = String(value).split('.')
  34. const decipher = crypto.createDecipheriv('aes-256-gcm', keyBuffer(), Buffer.from(ivRaw, 'base64'))
  35. decipher.setAuthTag(Buffer.from(tagRaw, 'base64'))
  36. return Buffer.concat([
  37. decipher.update(Buffer.from(encryptedRaw, 'base64')),
  38. decipher.final()
  39. ]).toString('utf8')
  40. }
  41. async function assertVip(uuid) {
  42. const vip = await getUserVipInfo(uuid)
  43. return vip.vip
  44. }
  45. function sanitizeBinding(row) {
  46. if (!row) return null
  47. return {
  48. user_uuid: row.user_uuid,
  49. state: Number(row.state),
  50. qrcode_url: row.qrcode_url,
  51. wechat_last_user_id: row.wechat_last_user_id ? 'bound' : '',
  52. last_user_message_time: Number(row.last_user_message_time || 0),
  53. expire_remind_time: Number(row.expire_remind_time || 0),
  54. last_error: row.last_error || '',
  55. bind_time: Number(row.bind_time || 0),
  56. create_time: Number(row.create_time || 0),
  57. update_time: Number(row.update_time || 0)
  58. }
  59. }
  60. async function getBindingByUser(uuid, includeSecret = false) {
  61. const rows = await db.query(
  62. `SELECT * FROM weixin_bot_binding WHERE user_uuid = ? LIMIT 1`,
  63. [uuid]
  64. )
  65. const row = rows?.[0] || null
  66. return includeSecret ? row : sanitizeBinding(row)
  67. }
  68. async function createBinding(uuid) {
  69. if (!await assertVip(uuid)) return { ok: false, msg: '微信绑定仅限 VIP 用户使用' }
  70. let qr
  71. try {
  72. qr = await WeixinBotClient.getQrCode()
  73. } catch (err) {
  74. logger.error(`Weixin get qrcode failed: ${err.stack || err.message || err}`)
  75. return { ok: false, msg: err.userMessage || '获取微信二维码失败' }
  76. }
  77. const qrcode = qr.qrcode || qr.data?.qrcode || ''
  78. let qrcodeUrl = qr.url || qr.qrcode_url || qr.qrcode_img_content || qr.data?.url || qr.data?.qrcode_url || qr.data?.qrcode_img_content || ''
  79. if (qrcodeUrl && !/^https?:\/\//i.test(qrcodeUrl) && !String(qrcodeUrl).startsWith('data:image/')) {
  80. qrcodeUrl = `data:image/png;base64,${qrcodeUrl}`
  81. }
  82. if (!qrcode) return { ok: false, msg: '获取微信二维码失败' }
  83. const now = Date.now()
  84. await db.query(
  85. `INSERT INTO weixin_bot_binding
  86. (user_uuid, state, qrcode, qrcode_url, bot_token_enc, base_url, get_updates_buf, last_error, create_time, update_time)
  87. VALUES (?, ?, ?, ?, '', '', '', '', ?, ?)
  88. ON DUPLICATE KEY UPDATE state = VALUES(state), qrcode = VALUES(qrcode), qrcode_url = VALUES(qrcode_url),
  89. bot_token_enc = '', base_url = '', get_updates_buf = '', last_error = '', update_time = VALUES(update_time)`,
  90. [uuid, STATE_PENDING, qrcode, qrcodeUrl, now, now]
  91. )
  92. return { ok: true, data: await getBindingByUser(uuid) }
  93. }
  94. async function refreshQrStatus(uuid) {
  95. const row = await getBindingByUser(uuid, true)
  96. if (!row) return null
  97. if (Number(row.state) !== STATE_PENDING || !row.qrcode) return sanitizeBinding(row)
  98. const status = await WeixinBotClient.getQrCodeStatus(row.qrcode)
  99. const confirmed = status.status === 'confirmed' || status.data?.status === 'confirmed'
  100. if (!confirmed) return sanitizeBinding(row)
  101. const botToken = status.bot_token || status.data?.bot_token
  102. const baseUrl = status.baseurl || status.base_url || status.data?.baseurl || status.data?.base_url || WeixinBotClient.DEFAULT_BASE_URL
  103. if (!botToken) return sanitizeBinding(row)
  104. const now = Date.now()
  105. await db.query(
  106. `UPDATE weixin_bot_binding
  107. SET state = ?, bot_token_enc = ?, base_url = ?, bind_time = ?, last_user_message_time = ?, last_error = '', update_time = ?
  108. WHERE user_uuid = ?`,
  109. [STATE_ACTIVE, encryptToken(botToken), baseUrl, now, now, now, uuid]
  110. )
  111. return getBindingByUser(uuid)
  112. }
  113. async function disableBinding(uuid) {
  114. const now = Date.now()
  115. await db.query(
  116. 'UPDATE weixin_bot_binding SET state = ?, update_time = ? WHERE user_uuid = ?',
  117. [STATE_DISABLED, now, uuid]
  118. )
  119. return true
  120. }
  121. function extractInbound(msg) {
  122. const items = Array.isArray(msg?.item_list) ? msg.item_list : []
  123. const texts = []
  124. const images = []
  125. for (const item of items) {
  126. if (Number(item.type) === 1 && item.text_item?.text) texts.push(String(item.text_item.text))
  127. if (Number(item.type) === 2) {
  128. const image = item.image_item || item.image || {}
  129. const url = image.url || image.cdn_url || image.download_url || image.file_url || ''
  130. if (url) images.push(String(url))
  131. else texts.push('[微信图片]')
  132. }
  133. }
  134. return {
  135. text: normalizeText(texts.join('\n'), 2000),
  136. images,
  137. fromUserId: msg?.from_user_id || '',
  138. toUserId: msg?.to_user_id || '',
  139. contextToken: msg?.context_token || ''
  140. }
  141. }
  142. async function getOrCreateWechatConversation(uuid) {
  143. return ConversationService.getOrCreateChannelConversation({
  144. uuid,
  145. channel: 'wechat',
  146. title: '微信对话'
  147. })
  148. }
  149. async function handleInboundMessage(binding, msg) {
  150. const inbound = extractInbound(msg)
  151. if (!inbound.fromUserId || (!inbound.text && inbound.images.length === 0)) return
  152. const now = Date.now()
  153. await db.query(
  154. `UPDATE weixin_bot_binding
  155. SET wechat_bot_user_id = ?, wechat_last_user_id = ?, last_context_token = ?, last_user_message_time = ?, expire_remind_time = 0, update_time = ?
  156. WHERE user_uuid = ?`,
  157. [inbound.toUserId, inbound.fromUserId, inbound.contextToken, now, now, binding.user_uuid]
  158. )
  159. const quota = await consumeQuota({ uuid: binding.user_uuid, channel: 'wechat' })
  160. if (!quota.allowed) {
  161. await sendTextToUser(binding.user_uuid, quota.message || '今日与小妍助理聊天次数已达上限,请明天再试')
  162. return
  163. }
  164. const conversation = await getOrCreateWechatConversation(binding.user_uuid)
  165. const saved = await ConversationService.addUserMessage({
  166. uuid: binding.user_uuid,
  167. conversationId: conversation.id,
  168. content: inbound.text,
  169. images: inbound.images,
  170. channel: 'wechat'
  171. })
  172. if (!saved || saved.missingContent) return
  173. try {
  174. await OneBotV11.sendAiChatMessage({
  175. conversationId: saved.conversationId,
  176. conversationNo: saved.conversationNo,
  177. senderUuid: binding.user_uuid,
  178. content: saved.content,
  179. images: saved.images,
  180. channel: 'wechat'
  181. })
  182. } catch (err) {
  183. logger.error(`WeChat AIChat OneBot forward failed: ${err.stack || err}`)
  184. await ConversationService.addSystemMessage({
  185. conversationId: saved.conversationId,
  186. content: '消息已保存,但暂时无法连接小妍助理。',
  187. status: 'error',
  188. errorMsg: err.message || 'OneBot send failed',
  189. channel: 'wechat'
  190. })
  191. await sendTextToUser(binding.user_uuid, '消息已保存,但暂时无法连接小妍助理,请稍后再试。')
  192. }
  193. }
  194. async function pollOneBinding(row) {
  195. const now = Date.now()
  196. if (!await assertVip(row.user_uuid)) {
  197. await db.query('UPDATE weixin_bot_binding SET state = ?, update_time = ?, last_error = ? WHERE user_uuid = ?', [STATE_DISABLED, now, 'VIP expired', row.user_uuid])
  198. return
  199. }
  200. if (Number(row.last_user_message_time || 0) > 0 && now - Number(row.last_user_message_time) >= WECHAT_EXPIRE_MS) {
  201. await db.query('UPDATE weixin_bot_binding SET state = ?, update_time = ? WHERE user_uuid = ?', [STATE_EXPIRED, now, row.user_uuid])
  202. return
  203. }
  204. const botToken = decryptToken(row.bot_token_enc)
  205. const result = await WeixinBotClient.getUpdates({
  206. baseUrl: row.base_url,
  207. botToken,
  208. getUpdatesBuf: row.get_updates_buf || ''
  209. })
  210. const nextBuf = result.get_updates_buf ?? result.data?.get_updates_buf ?? row.get_updates_buf ?? ''
  211. await db.query(
  212. 'UPDATE weixin_bot_binding SET get_updates_buf = ?, last_poll_time = ?, last_error = ?, update_time = ? WHERE user_uuid = ?',
  213. [nextBuf, Date.now(), '', Date.now(), row.user_uuid]
  214. )
  215. const msgs = result.msgs || result.data?.msgs || []
  216. for (const msg of msgs) {
  217. if (Number(msg.message_type) === 1) await handleInboundMessage(row, msg)
  218. }
  219. }
  220. async function pollActiveBindings() {
  221. if (running) return
  222. running = true
  223. try {
  224. const rows = await db.query(
  225. `SELECT * FROM weixin_bot_binding
  226. WHERE state = ?
  227. ORDER BY last_poll_time ASC
  228. LIMIT ?`,
  229. [STATE_ACTIVE, String(Number(config.weixinBot?.pollBatchSize || 5))]
  230. ) || []
  231. for (const row of rows) {
  232. try {
  233. await pollOneBinding(row)
  234. } catch (err) {
  235. logger.error(`WeChat binding poll failed user=${row.user_uuid}: ${err.stack || err}`)
  236. await db.query(
  237. 'UPDATE weixin_bot_binding SET last_error = ?, last_poll_time = ?, update_time = ? WHERE user_uuid = ?',
  238. [String(err.message || err).slice(0, 255), Date.now(), Date.now(), row.user_uuid]
  239. )
  240. }
  241. }
  242. await sendExpireReminders()
  243. } finally {
  244. running = false
  245. }
  246. }
  247. async function sendExpireReminders() {
  248. const threshold = Date.now() - WECHAT_REMIND_MS
  249. const rows = await db.query(
  250. `SELECT user_uuid FROM weixin_bot_binding
  251. WHERE state = ? AND last_user_message_time > 0 AND last_user_message_time <= ? AND expire_remind_time = 0
  252. LIMIT 30`,
  253. [STATE_ACTIVE, threshold]
  254. ) || []
  255. for (const row of rows) {
  256. try {
  257. await sendTextToUser(row.user_uuid, '微信登录即将过期,请在 1 小时内向小妍助理发送任意消息以保持登录。')
  258. await db.query('UPDATE weixin_bot_binding SET expire_remind_time = ?, update_time = ? WHERE user_uuid = ?', [Date.now(), Date.now(), row.user_uuid])
  259. } catch (err) {
  260. logger.error(`WeChat expire reminder failed user=${row.user_uuid}: ${err.stack || err}`)
  261. }
  262. }
  263. }
  264. async function sendTextToUser(uuid, text) {
  265. if (!await assertVip(uuid)) return false
  266. const row = await getBindingByUser(uuid, true)
  267. if (!row || Number(row.state) !== STATE_ACTIVE) return false
  268. if (!row.wechat_last_user_id || !row.last_context_token) return false
  269. try {
  270. await WeixinBotClient.sendText({
  271. baseUrl: row.base_url,
  272. botToken: decryptToken(row.bot_token_enc),
  273. toUserId: row.wechat_last_user_id,
  274. contextToken: row.last_context_token,
  275. text
  276. })
  277. await db.query(
  278. 'UPDATE weixin_bot_binding SET last_error = ?, update_time = ? WHERE user_uuid = ?',
  279. ['', Date.now(), uuid]
  280. )
  281. } catch (err) {
  282. await db.query(
  283. 'UPDATE weixin_bot_binding SET last_error = ?, update_time = ? WHERE user_uuid = ?',
  284. [String(err.message || err).slice(0, 255), Date.now(), uuid]
  285. )
  286. throw err
  287. }
  288. return true
  289. }
  290. async function listAdminBindings({ user_uuid, username, state, vip, expired, current = 1, pagesize = 20 }) {
  291. const c = Math.max(1, Number(current || 1))
  292. const p = Math.min(100, Math.max(1, Number(pagesize || 20)))
  293. const where = ['1=1']
  294. const params = []
  295. const countParams = []
  296. if (user_uuid) {
  297. where.push('b.user_uuid LIKE ?')
  298. params.push(`%${user_uuid}%`)
  299. countParams.push(`%${user_uuid}%`)
  300. }
  301. if (username) {
  302. where.push('u.username LIKE ?')
  303. params.push(`%${username}%`)
  304. countParams.push(`%${username}%`)
  305. }
  306. if (state !== undefined && state !== null && state !== '' && Number(state) !== -1) {
  307. where.push('b.state = ?')
  308. params.push(Number(state))
  309. countParams.push(Number(state))
  310. }
  311. if (String(vip) === '1') where.push('(u.vip_expire_time > UNIX_TIMESTAMP(CURRENT_TIMESTAMP(3)) * 1000 OR (u.vip = 1 AND u.vip_expire_time = 0))')
  312. if (String(vip) === '0') where.push('NOT (u.vip_expire_time > UNIX_TIMESTAMP(CURRENT_TIMESTAMP(3)) * 1000 OR (u.vip = 1 AND u.vip_expire_time = 0))')
  313. if (String(expired) === '1') where.push('(b.state = 3 OR (b.last_user_message_time > 0 AND b.last_user_message_time <= UNIX_TIMESTAMP(CURRENT_TIMESTAMP(3)) * 1000 - 86400000))')
  314. const whereSql = where.join(' AND ')
  315. const rows = await db.query(
  316. `SELECT b.user_uuid, b.state, b.wechat_last_user_id, b.last_user_message_time, b.expire_remind_time,
  317. b.last_error, b.last_poll_time, b.bind_time, b.create_time, b.update_time,
  318. u.username, u.vip, u.vip_expire_time
  319. FROM weixin_bot_binding b
  320. LEFT JOIN users u ON u.uuid = b.user_uuid
  321. WHERE ${whereSql}
  322. ORDER BY b.update_time DESC
  323. LIMIT ? OFFSET ?`,
  324. [...params, String(p), String((c - 1) * p)]
  325. )
  326. const countRows = await db.query(
  327. `SELECT COUNT(*) AS total FROM weixin_bot_binding b LEFT JOIN users u ON u.uuid = b.user_uuid WHERE ${whereSql}`,
  328. countParams
  329. )
  330. return {
  331. data: (rows || []).map(row => ({
  332. ...row,
  333. wechat_last_user_id: row.wechat_last_user_id ? 'bound' : '',
  334. token: undefined,
  335. bot_token_enc: undefined
  336. })),
  337. pagination: { current: c, pagesize: p, total: Number(countRows?.[0]?.total || 0) }
  338. }
  339. }
  340. function startPolling() {
  341. if (config.weixinBot?.enabled !== true) return
  342. if (pollTimer) return
  343. const interval = Math.max(5000, Number(config.weixinBot?.pollIntervalMs || 5000))
  344. pollTimer = setInterval(() => {
  345. pollActiveBindings().catch(err => logger.error(`WeChat poll loop failed: ${err.stack || err}`))
  346. }, interval)
  347. pollActiveBindings().catch(err => logger.error(`WeChat initial poll failed: ${err.stack || err}`))
  348. }
  349. module.exports = {
  350. STATE_PENDING,
  351. STATE_ACTIVE,
  352. STATE_DISABLED,
  353. STATE_EXPIRED,
  354. getBindingByUser,
  355. createBinding,
  356. refreshQrStatus,
  357. disableBinding,
  358. sendTextToUser,
  359. listAdminBindings,
  360. startPolling
  361. }