Server.js 7.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207
  1. const express = require('express')
  2. const cors = require('cors')
  3. const path = require('path')
  4. const fs = require('fs')
  5. const config = require('../config.json')
  6. const Logger = require('./Logger')
  7. const MySQL = require('../plugin/DataBase/MySQL')
  8. const Worker = require('./Lepao/Worker')
  9. const mq = require('../plugin/mq')
  10. const { mq: mqName } = require('../plugin/mq/mqPrefix')
  11. const { startLepaoSchedulePublisher } = require('../plugin/mq/lepaoSchedulePublisher')
  12. const OneBotV11 = require('../plugin/OneBot/OneBotV11')
  13. const WeixinBindingService = require('./AIChat/WeixinBindingService')
  14. const AccessControl = require('./AccessControl')
  15. const LeaseWatcher = require('./QK/LeaseWatcher')
  16. const { TaskScheduler } = require('./QK/TaskScheduler')
  17. const {
  18. resolveServerRole,
  19. shouldServeApi,
  20. shouldRunLepaoWorker,
  21. shouldRunOrderPaymentWorker
  22. } = require('./serverRole')
  23. const { TASK_QUEUE } = require('../plugin/mq/runforgeTaskMq')
  24. const { PREFIX } = require('../plugin/mq/mqPrefix')
  25. class SERVER {
  26. constructor() {
  27. this.serverRole = resolveServerRole()
  28. this.port = config.port || 3000
  29. this.apiDirectory = path.join(__dirname, '../apis')
  30. const logName = this.serverRole === 'worker' ? 'WorkerServer.log' : 'Server.log'
  31. this.logger = new Logger(path.join(__dirname, '../logs', logName), 'INFO')
  32. this.db = new MySQL()
  33. if (shouldServeApi(this.serverRole)) {
  34. this.app = express()
  35. this.app.use(express.json())
  36. this.app.use(cors())
  37. this.app.use('/uploads', express.static('./uploads'))
  38. this.app.use('/models', express.static('./models'))
  39. this.loadAPIs(this.apiDirectory)
  40. }
  41. }
  42. async initDB() {
  43. try {
  44. this.logger.info('正在测试数据库连接')
  45. await this.db.connect()
  46. await this.db.close()
  47. } catch (error) {
  48. this.logger.error(`数据库连接失败: ${error.stack}`)
  49. process.exit(1)
  50. }
  51. }
  52. async startLepaoWorker() {
  53. const worker = new Worker()
  54. await worker.start()
  55. this.logger.info('RunForge Worker 已启动,正在监听 MQ 任务...')
  56. startLepaoSchedulePublisher({
  57. logger: this.logger,
  58. intervalMs: config.rabbitmq?.lepaoScheduleTickMs ?? 2000,
  59. batch: config.rabbitmq?.lepaoScheduleBatch ?? 100
  60. })
  61. }
  62. async initMQ() {
  63. try {
  64. await mq.init()
  65. const ch = await mq.getChannel('health')
  66. await ch.assertQueue(mqName('mq_health_check'), { durable: false })
  67. this.logger.info('✅ RabbitMQ 初始化 & 测试成功')
  68. if (shouldRunLepaoWorker(this.serverRole)) {
  69. try {
  70. await this.startLepaoWorker()
  71. } catch (err) {
  72. console.error('RunForge Worker 启动失败:', err)
  73. process.exit(1)
  74. }
  75. } else {
  76. this.logger.info(
  77. `serverRole=api,跳过乐跑 Worker;乐跑任务将投递到 MQ 队列「${TASK_QUEUE}」(mqPrefix=${JSON.stringify(PREFIX)})`
  78. )
  79. }
  80. if (shouldRunOrderPaymentWorker(this.serverRole)) {
  81. const { startOrderPaymentWorker } = require('../plugin/mq/orderPaymentWorker')
  82. await startOrderPaymentWorker(this.logger)
  83. } else if (shouldServeApi(this.serverRole)) {
  84. this.logger.info('serverRole=api,跳过订单支付 MQ 消费者(由 all 进程消费)')
  85. }
  86. if (shouldRunLepaoWorker(this.serverRole)) {
  87. WeixinBindingService.startPolling()
  88. this.logger.info('WeChat AIChat binding poller started')
  89. }
  90. } catch (e) {
  91. this.logger.error('❌ RabbitMQ 初始化失败')
  92. process.exit(1)
  93. }
  94. }
  95. async initOneBot() {
  96. try {
  97. const ok = await OneBotV11.initOneBotWs()
  98. if (ok) {
  99. this.logger.info('OneBot v11 ws 初始化成功,已开始监听消息')
  100. } else {
  101. this.logger.info('OneBot v11 ws 未初始化(可能未启用)')
  102. }
  103. } catch (err) {
  104. this.logger.error(`OneBot v11 ws 初始化失败: ${err.message}`)
  105. }
  106. }
  107. async initAccessControlSchema() {
  108. await AccessControl.ensurePermissionSchema()
  109. }
  110. loadAPIs(directory) {
  111. const items = fs.readdirSync(directory)
  112. items.forEach(item => {
  113. const itemPath = path.join(directory, item)
  114. const stats = fs.statSync(itemPath)
  115. if (stats.isDirectory()) {
  116. this.loadAPIs(itemPath)
  117. } else if (stats.isFile() && itemPath.endsWith('.js')) {
  118. this.loadAPIFile(itemPath)
  119. }
  120. })
  121. }
  122. loadAPIFile(filePath) {
  123. try {
  124. const APIClass = require(filePath)
  125. for (const key in APIClass) {
  126. if (APIClass.hasOwnProperty(key)) {
  127. const apiInstance = new APIClass[key]()
  128. apiInstance.setupRoute()
  129. this.app.use('/', apiInstance.getRouter())
  130. this.logger.info(`已加载API:${apiInstance.path} 类型:${apiInstance.method}`)
  131. }
  132. }
  133. } catch (error) {
  134. this.logger.error(`加载API文件失败: ${filePath},错误: ${error.stack}`)
  135. }
  136. }
  137. listenHttp() {
  138. this.app.listen(this.port, () => {
  139. this.logger.info(`==========服务器正在 ${this.port} 端口上运行 (role=${this.serverRole})==========`)
  140. })
  141. }
  142. async startApiServices() {
  143. try {
  144. await this.initAccessControlSchema()
  145. } catch (err) {
  146. this.logger.error(`权限模型初始化异常: ${err.message}`)
  147. }
  148. try {
  149. await this.initOneBot()
  150. } catch (err) {
  151. this.logger.error(`OneBot 初始化异常: ${err.message}`)
  152. }
  153. this.qkLeaseWatcher = new LeaseWatcher({
  154. logger: this.logger,
  155. intervalMs: config.qk?.leaseWatcherIntervalMs || 30 * 1000,
  156. scheduler: new TaskScheduler({ logger: this.logger })
  157. })
  158. this.qkLeaseWatcher.start()
  159. this.listenHttp()
  160. }
  161. start() {
  162. this.logger.info(
  163. `============正在启动 (role=${this.serverRole}, lepaoQueue=${TASK_QUEUE}, mqPrefix=${JSON.stringify(PREFIX)})============`
  164. )
  165. this.initDB()
  166. .then(() => this.initMQ())
  167. .then(async () => {
  168. if (!shouldServeApi(this.serverRole)) {
  169. this.logger.info('serverRole=worker,未启动 HTTP API,进程将保持运行以消费乐跑任务')
  170. return
  171. }
  172. await this.startApiServices()
  173. })
  174. .catch((err) => {
  175. this.logger.error(`启动失败: ${err.message || err}`)
  176. if (err.stack) this.logger.error(err.stack)
  177. process.exit(1)
  178. })
  179. }
  180. }
  181. module.exports = SERVER