TaskScheduler.js 79 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360136113621363136413651366136713681369137013711372137313741375137613771378137913801381138213831384138513861387138813891390139113921393139413951396139713981399140014011402140314041405140614071408140914101411141214131414141514161417141814191420142114221423142414251426142714281429143014311432143314341435143614371438143914401441144214431444144514461447144814491450145114521453145414551456145714581459146014611462146314641465146614671468146914701471147214731474147514761477147814791480148114821483148414851486148714881489149014911492149314941495149614971498149915001501150215031504150515061507150815091510151115121513151415151516151715181519152015211522152315241525152615271528152915301531153215331534153515361537153815391540154115421543154415451546154715481549155015511552155315541555155615571558155915601561156215631564156515661567156815691570157115721573157415751576157715781579158015811582158315841585158615871588158915901591159215931594159515961597159815991600160116021603160416051606160716081609161016111612161316141615161616171618161916201621162216231624162516261627162816291630163116321633163416351636163716381639164016411642164316441645164616471648164916501651165216531654165516561657165816591660166116621663166416651666166716681669167016711672167316741675167616771678167916801681168216831684168516861687168816891690169116921693169416951696169716981699170017011702170317041705170617071708170917101711171217131714171517161717171817191720172117221723172417251726172717281729173017311732173317341735173617371738173917401741174217431744174517461747174817491750175117521753175417551756175717581759176017611762176317641765176617671768176917701771177217731774177517761777177817791780178117821783178417851786178717881789179017911792179317941795179617971798179918001801180218031804180518061807180818091810181118121813181418151816181718181819182018211822182318241825182618271828182918301831183218331834183518361837183818391840184118421843
  1. const crypto = require('crypto')
  2. const bcryptjs = require('bcryptjs')
  3. const config = require('../../config.json')
  4. const db = require('../../plugin/DataBase/db')
  5. const Redis = require('../../plugin/DataBase/Redis')
  6. const Logger = require('../Logger')
  7. const AccessControl = require('../AccessControl')
  8. const TASK_STATUS = {
  9. PAUSED: 'paused',
  10. PENDING: 'pending',
  11. ASSIGNED: 'assigned',
  12. RUNNING: 'running',
  13. SUCCESS: 'success',
  14. FAILED: 'failed',
  15. CANCELLED: 'cancelled'
  16. }
  17. const REPORT_LOG_EVENTS = new Set([
  18. 'request_result',
  19. 'progress_snapshot',
  20. 'grab_success',
  21. 'grab_fail'
  22. ])
  23. class TaskScheduler {
  24. static _schemaReady = false
  25. constructor(options = {}) {
  26. this.leaseMs = options.leaseMs || config.qk?.leaseMs || 90 * 1000
  27. this.heartbeatTtlSeconds = options.heartbeatTtlSeconds || config.qk?.heartbeatTtlSeconds || 45
  28. this.pullLockTtlSeconds = options.pullLockTtlSeconds || 5
  29. this.staleTaskMs = options.staleTaskMs || config.qk?.staleTaskMs || 60 * 60 * 1000
  30. this.memPerSlotMb = options.memPerSlotMb || config.qk?.memPerSlotMb || 3072
  31. this.memReserveMb = options.memReserveMb || config.qk?.memReserveMb || 1024
  32. this.maxSlotsCap = options.maxSlotsCap || config.qk?.maxSlotsCap || 10
  33. this.slotCapacityFactor = options.slotCapacityFactor || config.qk?.slotCapacityFactor || 1.5
  34. this.logger = options.logger || new Logger()
  35. }
  36. async ensureSchema() {
  37. if (TaskScheduler._schemaReady) {
  38. return
  39. }
  40. try {
  41. const rows = await db.query("SHOW COLUMNS FROM qk_task LIKE 'auto_relogin'")
  42. if (!rows || rows.length === 0) {
  43. await db.query(
  44. 'ALTER TABLE qk_task ADD COLUMN auto_relogin TINYINT(1) NOT NULL DEFAULT 0 AFTER enable_ggxxk'
  45. )
  46. this.logInfo('ensureSchema', '已添加 qk_task.auto_relogin 字段')
  47. }
  48. const jwRows = await db.query("SHOW COLUMNS FROM qk_task LIKE 'jw_account_id'")
  49. if (!jwRows || jwRows.length === 0) {
  50. await db.query(
  51. 'ALTER TABLE qk_task ADD COLUMN jw_account_id INT NULL DEFAULT NULL AFTER student_num'
  52. )
  53. this.logInfo('ensureSchema', '已添加 qk_task.jw_account_id 字段')
  54. }
  55. } catch (err) {
  56. this.logWarn('ensureSchema', '抢课任务表结构检查失败', {}, err)
  57. }
  58. TaskScheduler._schemaReady = true
  59. }
  60. parseAutoRelogin(payload = {}) {
  61. if (payload.auto_relogin !== undefined) {
  62. return payload.auto_relogin === true || Number(payload.auto_relogin) === 1
  63. }
  64. if (payload.AUTO_RELOGIN !== undefined) {
  65. return payload.AUTO_RELOGIN === true || Number(payload.AUTO_RELOGIN) === 1
  66. }
  67. return false
  68. }
  69. extractReportCourseName(payload = {}) {
  70. if (!payload || typeof payload !== 'object') {
  71. return ''
  72. }
  73. if (payload.course_name) {
  74. return String(payload.course_name)
  75. }
  76. if (payload.course) {
  77. return String(payload.course)
  78. }
  79. if (payload.label) {
  80. const label = String(payload.label)
  81. const at = label.indexOf('@')
  82. if (at > 0) {
  83. return label.slice(0, at)
  84. }
  85. return label
  86. }
  87. if (Array.isArray(payload.courses) && payload.courses.length > 0) {
  88. return payload.courses.join('、')
  89. }
  90. return ''
  91. }
  92. formatReportMessage(message = '', payload = {}) {
  93. const courseName = this.extractReportCourseName(payload)
  94. const text = String(message || '').trim()
  95. if (!courseName) {
  96. return text
  97. }
  98. if (!text) {
  99. return `[${courseName}]`
  100. }
  101. if (text.includes(`[${courseName}]`) || text.startsWith(`${courseName}:`)) {
  102. return text
  103. }
  104. return `[${courseName}] ${text}`
  105. }
  106. calculateMaxSlots(profile = {}) {
  107. const freeMb = Math.max(0, Number(profile.free_mem_mb) || 0)
  108. const totalMb = Math.max(0, Number(profile.total_mem_mb) || 0)
  109. const threads = Math.max(1, Number(profile.cpu_threads) || 1)
  110. const allocatableFreeMb = Math.max(0, freeMb - this.memReserveMb)
  111. const allocatableTotalMb = Math.max(0, totalMb - this.memReserveMb)
  112. const byFreeMem = Math.floor(allocatableFreeMb / this.memPerSlotMb)
  113. const byTotalMem = Math.floor(allocatableTotalMb / this.memPerSlotMb)
  114. const byCpu = Math.floor(threads * 0.8)
  115. const fallbackMem = totalMb > 0 ? byTotalMem : 1
  116. const byMem = freeMb > 0 ? Math.min(byFreeMem, byTotalMem) : fallbackMem
  117. const base = Math.max(1, Math.min(byMem, byCpu))
  118. const scaled = Math.max(1, Math.floor(base * this.slotCapacityFactor))
  119. const cap = Math.max(1, Math.floor(this.maxSlotsCap * this.slotCapacityFactor))
  120. return Math.max(1, Math.min(50, cap, scaled))
  121. }
  122. resolveClientMaxSlots(client, payload = {}) {
  123. const fromPayload = this.calculateMaxSlots({
  124. free_mem_mb: payload.free_mem_mb ?? client.free_mem_mb,
  125. total_mem_mb: payload.total_mem_mb ?? client.total_mem_mb,
  126. cpu_threads: payload.cpu_threads ?? client.cpu_threads
  127. })
  128. const reported = Number(payload.max_slots || client.max_slots || 0)
  129. if (!reported) return fromPayload
  130. return Math.max(1, Math.min(reported, fromPayload))
  131. }
  132. getClientAvailableSlots(client, payload = {}) {
  133. const maxSlots = this.resolveClientMaxSlots(client, payload)
  134. const currentSlots = Math.max(0, Number(client.current_slots ?? 0))
  135. return Math.max(0, maxSlots - currentSlots)
  136. }
  137. async countClientActiveTasks(clientId) {
  138. if (!clientId) return 0
  139. const rows = await db.query(
  140. 'SELECT COUNT(*) AS cnt FROM qk_task WHERE assigned_client_id = ? AND status IN (?, ?)',
  141. [clientId, TASK_STATUS.ASSIGNED, TASK_STATUS.RUNNING]
  142. )
  143. return Math.max(0, Number(rows?.[0]?.cnt || 0))
  144. }
  145. async syncClientSlots(clientId, now = Date.now()) {
  146. if (!clientId) return 0
  147. const count = await this.countClientActiveTasks(clientId)
  148. await db.query(
  149. 'UPDATE qk_client SET current_slots = ?, update_time = ? WHERE client_id = ?',
  150. [count, now, clientId]
  151. )
  152. try {
  153. await Redis.set(`qk:client:slots:${clientId}`, String(count), { EX: this.heartbeatTtlSeconds })
  154. } catch (_) {}
  155. return count
  156. }
  157. async findRevokedClientTasks(clientId, runningTaskIds = []) {
  158. const ids = (runningTaskIds || []).map(Number).filter(Boolean)
  159. if (!ids.length) return []
  160. const placeholders = ids.map(() => '?').join(',')
  161. const rows = await db.query(
  162. `SELECT id FROM qk_task WHERE id IN (${placeholders})
  163. AND (assigned_client_id IS NULL OR assigned_client_id <> ? OR status NOT IN (?, ?))`,
  164. [...ids, clientId, TASK_STATUS.ASSIGNED, TASK_STATUS.RUNNING]
  165. )
  166. return (rows || []).map(row => row.id)
  167. }
  168. async getClientAvailableSlotsAsync(client, payload = {}) {
  169. const maxSlots = this.resolveClientMaxSlots(client, payload)
  170. const activeCount = await this.countClientActiveTasks(client.client_id)
  171. return Math.max(0, maxSlots - activeCount)
  172. }
  173. getInFlightTaskStatuses() {
  174. return [TASK_STATUS.ASSIGNED, TASK_STATUS.RUNNING]
  175. }
  176. /** 同一学号在单次拉取中只保留最早创建的一条待分配任务 */
  177. pickPullableTasks(rows, limit) {
  178. const picked = []
  179. const seenStudentNums = new Set()
  180. for (const row of rows || []) {
  181. const studentNum = String(row.student_num || '').trim()
  182. if (!studentNum || seenStudentNums.has(studentNum)) {
  183. continue
  184. }
  185. seenStudentNums.add(studentNum)
  186. picked.push(row)
  187. if (picked.length >= limit) {
  188. break
  189. }
  190. }
  191. return picked
  192. }
  193. async countOnlineClients() {
  194. const now = Date.now()
  195. const threshold = now - this.heartbeatTtlSeconds * 1000
  196. const rows = await db.query(
  197. `SELECT COUNT(*) AS total FROM qk_client
  198. WHERE enabled = 1 AND online = 1
  199. AND last_heartbeat_at IS NOT NULL AND last_heartbeat_at >= ?`,
  200. [threshold]
  201. )
  202. return Number(rows?.[0]?.total || 0)
  203. }
  204. safeStringify(obj) {
  205. const seen = new WeakSet()
  206. return JSON.stringify(obj, (key, value) => {
  207. if (typeof value === 'object' && value !== null) {
  208. if (seen.has(value)) return '[Circular]'
  209. seen.add(value)
  210. }
  211. return value
  212. })
  213. }
  214. sanitizeForLog(payload) {
  215. if (!payload || typeof payload !== 'object') {
  216. return payload
  217. }
  218. const copy = Array.isArray(payload) ? [...payload] : { ...payload }
  219. for (const key of ['password', 'pass', 'password_enc', 'client_secret']) {
  220. if (key in copy) {
  221. copy[key] = '***'
  222. }
  223. }
  224. return copy
  225. }
  226. buildLogPrefix(tag, ctx = {}) {
  227. const parts = ['[QK]', `[${tag}]`]
  228. if (ctx.taskId) parts.push(`[taskId=${ctx.taskId}]`)
  229. if (ctx.clientId) parts.push(`[clientId=${ctx.clientId}]`)
  230. if (ctx.uuid) parts.push(`[uuid=${ctx.uuid}]`)
  231. return parts.join('')
  232. }
  233. logInfo(tag, message, ctx = {}, data = null) {
  234. const prefix = this.buildLogPrefix(tag, ctx)
  235. const suffix = data != null ? ` ${this.safeStringify(this.sanitizeForLog(data))}` : ''
  236. this.logger.info(`${prefix} ${message}${suffix}`)
  237. }
  238. logWarn(tag, message, ctx = {}, data = null) {
  239. const prefix = this.buildLogPrefix(tag, ctx)
  240. const suffix = data != null ? ` ${this.safeStringify(this.sanitizeForLog(data))}` : ''
  241. this.logger.warn(`${prefix} ${message}${suffix}`)
  242. }
  243. logError(tag, message, ctx = {}, err = null) {
  244. const prefix = this.buildLogPrefix(tag, ctx)
  245. const suffix = err ? ` ${err.stack || err}` : ''
  246. this.logger.error(`${prefix} ${message}${suffix}`)
  247. }
  248. getPasswordKey() {
  249. const source = process.env.QK_PASSWORD_KEY || config.qk?.passwordAesKey || config.database?.password || 'runforge-qk-default-key'
  250. return crypto.createHash('sha256').update(String(source)).digest()
  251. }
  252. encryptPassword(password) {
  253. const iv = crypto.randomBytes(16)
  254. const cipher = crypto.createCipheriv('aes-256-cbc', this.getPasswordKey(), iv)
  255. let encrypted = cipher.update(String(password), 'utf8', 'base64')
  256. encrypted += cipher.final('base64')
  257. return `${iv.toString('base64')}:${encrypted}`
  258. }
  259. decryptPassword(encrypted) {
  260. const [ivText, payload] = String(encrypted || '').split(':')
  261. if (!ivText || !payload) {
  262. return ''
  263. }
  264. const decipher = crypto.createDecipheriv('aes-256-cbc', this.getPasswordKey(), Buffer.from(ivText, 'base64'))
  265. let decrypted = decipher.update(payload, 'base64', 'utf8')
  266. decrypted += decipher.final('utf8')
  267. return decrypted
  268. }
  269. normalizeArray(value) {
  270. if (Array.isArray(value)) {
  271. return value.map(item => String(item).trim()).filter(Boolean)
  272. }
  273. if (typeof value === 'string') {
  274. const trimmed = value.trim()
  275. if (!trimmed) {
  276. return []
  277. }
  278. try {
  279. const parsed = JSON.parse(trimmed)
  280. if (Array.isArray(parsed)) {
  281. return parsed.map(item => String(item).trim()).filter(Boolean)
  282. }
  283. } catch (_) {
  284. return trimmed.split(/[\n,,]/).map(item => item.trim()).filter(Boolean)
  285. }
  286. }
  287. return []
  288. }
  289. countTaskTargets(courses = [], courseGroups = []) {
  290. return this.normalizeArray(courses).length + this.normalizeArray(courseGroups).length
  291. }
  292. getMinTaskInterval(totalCount) {
  293. return totalCount >= 5 ? 500 : 200
  294. }
  295. validateTaskCoursesAndInterval(courses, courseGroups, intervalMs) {
  296. const normalizedCourses = this.normalizeArray(courses)
  297. const normalizedGroups = this.normalizeArray(courseGroups)
  298. const totalCount = normalizedCourses.length + normalizedGroups.length
  299. if (totalCount <= 0) {
  300. throw new Error('至少需要填写一门课程或一个课程分组')
  301. }
  302. if (totalCount > 10) {
  303. throw new Error('每个任务最多选择10门课程或分组')
  304. }
  305. const minInterval = this.getMinTaskInterval(totalCount)
  306. const safeInterval = Number(intervalMs)
  307. if (!Number.isFinite(safeInterval) || safeInterval < minInterval || safeInterval > 10000) {
  308. throw new Error(`当前共 ${totalCount} 个抢课目标,间隔需在 ${minInterval}-10000ms 之间`)
  309. }
  310. return {
  311. courses: normalizedCourses,
  312. courseGroups: normalizedGroups,
  313. intervalMs: safeInterval,
  314. totalCount,
  315. minInterval
  316. }
  317. }
  318. async resolveTaskAccount(ownerUuid, payload, { requireAccount = true } = {}) {
  319. const jwAccountId = Number(payload.jw_account_id || payload.jwAccountId)
  320. if (jwAccountId) {
  321. const account = await AccessControl.getVerifiedJwAccount(ownerUuid, jwAccountId)
  322. if (!account) {
  323. throw new Error('统一身份认证账号不存在或未验证')
  324. }
  325. return {
  326. student_num: account.username,
  327. password: account.password,
  328. jw_account_id: account.id
  329. }
  330. }
  331. const studentNum = String(payload.student_num || payload.user || '').trim()
  332. const password = payload.password || payload.pass
  333. if (studentNum && password) {
  334. return {
  335. student_num: studentNum,
  336. password,
  337. jw_account_id: payload.jw_account_id || null
  338. }
  339. }
  340. if (requireAccount) {
  341. throw new Error('请选择统一身份认证账号')
  342. }
  343. return null
  344. }
  345. async inferJwAccountId(ownerUuid, studentNum, currentId) {
  346. if (currentId) {
  347. return Number(currentId)
  348. }
  349. const username = String(studentNum || '').trim()
  350. if (!username) {
  351. return null
  352. }
  353. const rows = await db.query(
  354. 'SELECT id FROM jw_account WHERE create_user = ? AND username = ? AND state = 1 LIMIT 1',
  355. [ownerUuid, username]
  356. )
  357. return rows?.[0]?.id || null
  358. }
  359. serializeTask(row, includeSecret = false) {
  360. const result = { ...row }
  361. result.enable_ggxxk = Number(result.enable_ggxxk) === 1
  362. result.auto_relogin = Number(result.auto_relogin) === 1
  363. result.courses = this.normalizeArray(result.courses)
  364. result.course_groups = this.normalizeArray(result.course_groups)
  365. if (typeof result.result_json === 'string' && result.result_json) {
  366. try {
  367. result.result_json = JSON.parse(result.result_json)
  368. } catch (_) {}
  369. }
  370. result.jx0502zbid = row.batch_jx0502zbid || row.jx0502zbid || ''
  371. result.batch_name = row.batch_name || ''
  372. if (row.batch_enabled !== undefined) {
  373. result.batch_enabled = Number(row.batch_enabled) === 1
  374. }
  375. delete result.batch_jx0502zbid
  376. if (includeSecret) {
  377. result.password = this.decryptPassword(result.password_enc)
  378. }
  379. delete result.password_enc
  380. return result
  381. }
  382. taskBatchJoinSql(alias = 't') {
  383. return `LEFT JOIN qk_batch b ON b.id = ${alias}.batch_id`
  384. }
  385. taskBatchSelectSql(alias = 't') {
  386. return `, b.name AS batch_name, b.jx0502zbid AS batch_jx0502zbid, b.enabled AS batch_enabled`
  387. }
  388. async getBatchById(batchId) {
  389. const rows = await db.query('SELECT * FROM qk_batch WHERE id = ?', [batchId])
  390. if (!rows || rows.length === 0) {
  391. throw new Error('抢课批次不存在')
  392. }
  393. return rows[0]
  394. }
  395. async requireEnabledBatch(batchId) {
  396. const batch = await this.getBatchById(batchId)
  397. if (!Number(batch.enabled)) {
  398. throw new Error('所选抢课批次未启用')
  399. }
  400. return batch
  401. }
  402. resolveTaskBatchId(payload = {}) {
  403. const batchId = Number(payload.batch_id || payload.batchId)
  404. if (!batchId) {
  405. throw new Error('请选择抢课批次')
  406. }
  407. return batchId
  408. }
  409. async releaseAssignedTask(task) {
  410. return !!(task && [TASK_STATUS.ASSIGNED, TASK_STATUS.RUNNING].includes(task.status))
  411. }
  412. async pauseTasksByBatch(batchId, reason = '批次已停用,任务已收回') {
  413. const rows = await db.query(
  414. 'SELECT id, status, assigned_client_id FROM qk_task WHERE batch_id = ? AND status IN (?, ?, ?)',
  415. [batchId, TASK_STATUS.PENDING, TASK_STATUS.ASSIGNED, TASK_STATUS.RUNNING]
  416. )
  417. const now = Date.now()
  418. const affectedClients = new Set()
  419. let count = 0
  420. for (const row of rows || []) {
  421. if (row.assigned_client_id) {
  422. affectedClients.add(row.assigned_client_id)
  423. }
  424. const result = await db.query(
  425. `UPDATE qk_task SET status = ?, assigned_client_id = NULL, lease_expire_at = NULL,
  426. assigned_at = NULL, exclude_client_id = NULL, update_time = ? WHERE id = ? AND status IN (?, ?, ?)`,
  427. [TASK_STATUS.PAUSED, now, row.id, TASK_STATUS.PENDING, TASK_STATUS.ASSIGNED, TASK_STATUS.RUNNING]
  428. )
  429. if (result && result.affectedRows > 0) {
  430. count += 1
  431. await this.logTask(row.id, row.assigned_client_id, 'batch_paused', reason)
  432. }
  433. }
  434. for (const clientId of affectedClients) {
  435. await this.syncClientSlots(clientId, now)
  436. }
  437. if (count > 0) {
  438. this.logInfo('pauseTasksByBatch', reason, { batchId }, { affected: count })
  439. }
  440. return count
  441. }
  442. async redispatchTasksByBatch(batchId, reason = '批次 ID 已变更,任务已重新进入分配队列') {
  443. const rows = await db.query(
  444. 'SELECT id, status, assigned_client_id FROM qk_task WHERE batch_id = ? AND status IN (?, ?, ?)',
  445. [batchId, TASK_STATUS.PENDING, TASK_STATUS.ASSIGNED, TASK_STATUS.RUNNING]
  446. )
  447. const now = Date.now()
  448. let count = 0
  449. for (const row of rows || []) {
  450. if ([TASK_STATUS.ASSIGNED, TASK_STATUS.RUNNING].includes(row.status)) {
  451. await this.releaseAssignedTask(row, now)
  452. const result = await db.query(
  453. `UPDATE qk_task SET status = ?, assigned_client_id = NULL, lease_expire_at = NULL,
  454. assigned_at = NULL, exclude_client_id = NULL, update_time = ? WHERE id = ? AND status IN (?, ?)`,
  455. [TASK_STATUS.PENDING, now, row.id, TASK_STATUS.ASSIGNED, TASK_STATUS.RUNNING]
  456. )
  457. if (result && result.affectedRows > 0) {
  458. count += 1
  459. await this.logTask(row.id, row.assigned_client_id, 'batch_redispatch', reason)
  460. }
  461. }
  462. }
  463. if (count > 0) {
  464. this.logInfo('redispatchTasksByBatch', reason, { batchId }, { affected: count })
  465. }
  466. return count
  467. }
  468. serializeBatch(row) {
  469. return {
  470. id: row.id,
  471. name: row.name,
  472. jx0502zbid: row.jx0502zbid,
  473. enabled: Number(row.enabled) === 1,
  474. create_time: row.create_time,
  475. update_time: row.update_time
  476. }
  477. }
  478. async listBatches(options = {}) {
  479. const enabledOnly = !!options.enabledOnly
  480. const where = enabledOnly ? 'WHERE enabled = 1' : ''
  481. const rows = await db.query(`SELECT * FROM qk_batch ${where} ORDER BY update_time DESC, id DESC`)
  482. return (rows || []).map(row => this.serializeBatch(row))
  483. }
  484. async createBatch(payload = {}) {
  485. const name = String(payload.name || '').trim()
  486. const jx0502zbid = String(payload.jx0502zbid || '').trim()
  487. if (!name || !jx0502zbid) {
  488. throw new Error('批次名称和批次 ID 不能为空')
  489. }
  490. const now = Date.now()
  491. const result = await db.query(
  492. 'INSERT INTO qk_batch (name, jx0502zbid, enabled, create_time, update_time) VALUES (?, ?, ?, ?, ?)',
  493. [name, jx0502zbid, payload.enabled === false ? 0 : 1, now, now]
  494. )
  495. if (!result || result.affectedRows <= 0) {
  496. throw new Error('创建抢课批次失败')
  497. }
  498. this.logInfo('createBatch', '抢课批次已创建', { batchId: result.insertId }, { name, jx0502zbid })
  499. return result.insertId
  500. }
  501. async updateBatch(batchId, payload = {}) {
  502. const existing = await this.getBatchById(batchId)
  503. const name = payload.name !== undefined ? String(payload.name || '').trim() : existing.name
  504. const jx0502zbid = payload.jx0502zbid !== undefined ? String(payload.jx0502zbid || '').trim() : existing.jx0502zbid
  505. const enabled = payload.enabled !== undefined ? (payload.enabled ? 1 : 0) : Number(existing.enabled)
  506. if (!name || !jx0502zbid) {
  507. throw new Error('批次名称和批次 ID 不能为空')
  508. }
  509. const now = Date.now()
  510. const result = await db.query(
  511. 'UPDATE qk_batch SET name = ?, jx0502zbid = ?, enabled = ?, update_time = ? WHERE id = ?',
  512. [name, jx0502zbid, enabled, now, batchId]
  513. )
  514. if (!result || result.affectedRows <= 0) {
  515. throw new Error('更新抢课批次失败')
  516. }
  517. const disabledNow = Number(existing.enabled) === 1 && enabled === 0
  518. const idChanged = String(existing.jx0502zbid) !== jx0502zbid
  519. if (disabledNow) {
  520. await this.pauseTasksByBatch(batchId)
  521. } else if (idChanged && enabled === 1) {
  522. await this.redispatchTasksByBatch(batchId)
  523. }
  524. this.logInfo('updateBatch', '抢课批次已更新', { batchId }, {
  525. name,
  526. jx0502zbid,
  527. enabled: enabled === 1,
  528. disabledNow,
  529. idChanged
  530. })
  531. }
  532. async logTask(taskId, clientId, event, message = '', payload = null) {
  533. const sql = 'INSERT INTO qk_task_log (task_id, client_id, event, message, payload_json, create_time) VALUES (?, ?, ?, ?, ?, ?)'
  534. await db.query(sql, [
  535. taskId,
  536. clientId || null,
  537. event,
  538. message || '',
  539. payload ? JSON.stringify(payload) : null,
  540. Date.now()
  541. ])
  542. this.logInfo('taskLog', message || event, { taskId, clientId }, payload ? { event, ...this.sanitizeForLog(payload) } : { event })
  543. }
  544. serializeReportLog(row) {
  545. const result = { ...row }
  546. if (typeof result.payload_json === 'string' && result.payload_json) {
  547. try {
  548. result.payload_json = JSON.parse(result.payload_json)
  549. } catch (_) {}
  550. }
  551. const payload = result.payload_json && typeof result.payload_json === 'object'
  552. ? result.payload_json
  553. : {}
  554. result.course_name = this.extractReportCourseName(payload)
  555. result.display_message = this.formatReportMessage(result.message, payload)
  556. return result
  557. }
  558. buildReportMessage(payload = {}) {
  559. const base = (() => {
  560. if (payload.message) {
  561. return String(payload.message)
  562. }
  563. if (payload.error_msg) {
  564. return String(payload.error_msg)
  565. }
  566. if (payload.error) {
  567. return String(payload.error)
  568. }
  569. if (payload.success === true) {
  570. return payload.label ? `${payload.label} 成功` : '请求成功'
  571. }
  572. return payload.label ? `${payload.label} 失败` : '请求失败'
  573. })()
  574. return this.formatReportMessage(base, payload)
  575. }
  576. async assertClientTaskAccess(clientId, taskId) {
  577. const rows = await db.query(
  578. 'SELECT id, status, assigned_client_id FROM qk_task WHERE id = ?',
  579. [taskId]
  580. )
  581. if (!rows || rows.length === 0) {
  582. throw new Error('任务不存在')
  583. }
  584. const task = rows[0]
  585. if (task.assigned_client_id !== clientId) {
  586. throw new Error('任务不属于当前客户端')
  587. }
  588. if (![TASK_STATUS.ASSIGNED, TASK_STATUS.RUNNING].includes(task.status)) {
  589. throw new Error('任务当前状态不可上报')
  590. }
  591. return task
  592. }
  593. async reportProgress(clientId, clientSecret, payload = {}) {
  594. const client = await this.authenticateClient(clientId, clientSecret)
  595. if (!client) {
  596. throw new Error('客户端凭证无效')
  597. }
  598. const taskId = Number(payload.task_id || payload.id)
  599. if (!taskId) {
  600. throw new Error('缺少任务 ID')
  601. }
  602. const event = String(payload.event || 'request_result')
  603. if (!REPORT_LOG_EVENTS.has(event) || event === 'grab_success' || event === 'grab_fail') {
  604. throw new Error('不支持的上报类型')
  605. }
  606. const task = await this.assertClientTaskAccess(clientId, taskId)
  607. const message = this.buildReportMessage(payload)
  608. const now = Date.now()
  609. await this.logTask(taskId, clientId, event, message, payload)
  610. const updates = ['update_time = ?']
  611. const params = [now]
  612. if (task.status === TASK_STATUS.ASSIGNED) {
  613. updates.push('status = ?')
  614. params.push(TASK_STATUS.RUNNING)
  615. }
  616. if (payload.success !== true && message) {
  617. updates.push('error_msg = ?')
  618. params.push(message)
  619. }
  620. params.push(taskId, clientId)
  621. await db.query(
  622. `UPDATE qk_task SET ${updates.join(', ')} WHERE id = ? AND assigned_client_id = ?`,
  623. params
  624. )
  625. this.logInfo('reportProgress', '客户端上报抢课进度', { taskId, clientId }, {
  626. event,
  627. success: payload.success === true,
  628. message
  629. })
  630. return { task_id: taskId, event, message }
  631. }
  632. async createTask(uuid, payload) {
  633. await this.ensureSchema()
  634. const validated = this.validateTaskCoursesAndInterval(
  635. payload.courses || payload.COURSES,
  636. payload.course_groups || payload.COURSE_GROUPS,
  637. payload.interval_ms || payload.INTERVAL_MS || 500
  638. )
  639. const { courses, courseGroups, intervalMs } = validated
  640. const batchId = this.resolveTaskBatchId(payload)
  641. await this.requireEnabledBatch(batchId)
  642. const account = await this.resolveTaskAccount(uuid, payload, { requireAccount: true })
  643. const time = Date.now()
  644. const sql = `INSERT INTO qk_task
  645. (create_user, name, batch_id, jx0502zbid, student_num, jw_account_id, password_enc, courses, course_groups, enable_ggxxk, auto_relogin, interval_ms, status, create_time, update_time)
  646. VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`
  647. const result = await db.query(sql, [
  648. uuid,
  649. payload.name,
  650. batchId,
  651. '',
  652. account.student_num,
  653. account.jw_account_id,
  654. this.encryptPassword(account.password),
  655. JSON.stringify(courses),
  656. JSON.stringify(courseGroups),
  657. payload.enable_ggxxk || payload.ENABLE_GGXXK ? 1 : 0,
  658. this.parseAutoRelogin(payload) ? 1 : 0,
  659. intervalMs,
  660. TASK_STATUS.PAUSED,
  661. time,
  662. time
  663. ])
  664. if (!result || result.affectedRows <= 0) {
  665. throw new Error('创建抢课任务失败')
  666. }
  667. await this.logTask(result.insertId, null, 'created', '用户提交抢课任务(未开始)')
  668. this.logInfo('createTask', '抢课任务已创建', { taskId: result.insertId, uuid }, {
  669. name: payload.name,
  670. batch_id: batchId,
  671. student_num: account.student_num,
  672. jw_account_id: account.jw_account_id,
  673. courses_count: courses.length,
  674. course_groups_count: courseGroups.length,
  675. interval_ms: intervalMs,
  676. enable_ggxxk: !!(payload.enable_ggxxk || payload.ENABLE_GGXXK),
  677. auto_relogin: this.parseAutoRelogin(payload)
  678. })
  679. return result.insertId
  680. }
  681. async updateTask(uuid, taskId, payload) {
  682. await this.ensureSchema()
  683. const rows = await db.query('SELECT status FROM qk_task WHERE id = ? AND create_user = ?', [taskId, uuid])
  684. if (!rows || rows.length === 0) {
  685. throw new Error('任务不存在')
  686. }
  687. if (![TASK_STATUS.PAUSED, TASK_STATUS.FAILED, TASK_STATUS.CANCELLED].includes(rows[0].status)) {
  688. throw new Error('任务进行中或已分配,请先暂停后再修改')
  689. }
  690. const validated = this.validateTaskCoursesAndInterval(
  691. payload.courses || payload.COURSES,
  692. payload.course_groups || payload.COURSE_GROUPS,
  693. payload.interval_ms || payload.INTERVAL_MS || 500
  694. )
  695. const { courses, courseGroups, intervalMs } = validated
  696. const batchId = this.resolveTaskBatchId(payload)
  697. await this.requireEnabledBatch(batchId)
  698. const account = await this.resolveTaskAccount(uuid, payload, { requireAccount: true })
  699. const params = [
  700. payload.name,
  701. batchId,
  702. account.student_num,
  703. account.jw_account_id,
  704. this.encryptPassword(account.password),
  705. JSON.stringify(courses),
  706. JSON.stringify(courseGroups),
  707. payload.enable_ggxxk || payload.ENABLE_GGXXK ? 1 : 0,
  708. this.parseAutoRelogin(payload) ? 1 : 0,
  709. intervalMs,
  710. TASK_STATUS.PAUSED,
  711. Date.now()
  712. ]
  713. params.push(taskId, uuid)
  714. const sql = `UPDATE qk_task SET name = ?, batch_id = ?, student_num = ?, jw_account_id = ?, password_enc = ?, courses = ?, course_groups = ?, enable_ggxxk = ?, auto_relogin = ?, interval_ms = ?, status = ?, jx0502zbid = '', update_time = ?, assigned_client_id = NULL, lease_expire_at = NULL, result_json = NULL, error_msg = NULL, finished_time = NULL WHERE id = ? AND create_user = ?`
  715. const result = await db.query(sql, params)
  716. if (!result || result.affectedRows <= 0) {
  717. throw new Error('更新抢课任务失败')
  718. }
  719. await this.logTask(taskId, null, 'updated', '用户更新抢课任务')
  720. this.logInfo('updateTask', '抢课任务已更新(未开始)', { taskId, uuid }, {
  721. name: payload.name,
  722. student_num: account.student_num,
  723. jw_account_id: account.jw_account_id,
  724. courses_count: courses.length,
  725. course_groups_count: courseGroups.length,
  726. interval_ms: intervalMs,
  727. account_changed: true
  728. })
  729. }
  730. async listUserTasks(uuid, filters = {}) {
  731. const where = ['t.create_user = ?']
  732. const params = [uuid]
  733. if (filters.status) {
  734. where.push('t.status = ?')
  735. params.push(filters.status)
  736. }
  737. if (filters.student_num) {
  738. where.push('t.student_num LIKE ?')
  739. params.push(`%${filters.student_num}%`)
  740. }
  741. if (filters.name) {
  742. where.push('t.name LIKE ?')
  743. params.push(`%${filters.name}%`)
  744. }
  745. const courseName = String(filters.course_name || filters.course || '').trim()
  746. if (courseName) {
  747. const keyword = `%${courseName}%`
  748. where.push('(t.courses LIKE ? OR t.course_groups LIKE ?)')
  749. params.push(keyword, keyword)
  750. }
  751. if (filters.batch_id) {
  752. where.push('t.batch_id = ?')
  753. params.push(Number(filters.batch_id))
  754. }
  755. const whereSql = where.join(' AND ')
  756. const rows = await db.query(
  757. `SELECT t.*${this.taskBatchSelectSql('t')}
  758. FROM qk_task t
  759. ${this.taskBatchJoinSql('t')}
  760. WHERE ${whereSql}
  761. ORDER BY t.create_time DESC`,
  762. params
  763. )
  764. return (rows || []).map(row => this.serializeTask(row))
  765. }
  766. async getTaskDetail(uuid, taskId) {
  767. const rows = await db.query(
  768. `SELECT t.*${this.taskBatchSelectSql('t')}
  769. FROM qk_task t
  770. ${this.taskBatchJoinSql('t')}
  771. WHERE t.id = ? AND t.create_user = ?`,
  772. [taskId, uuid]
  773. )
  774. if (!rows || rows.length === 0) {
  775. return null
  776. }
  777. const task = this.serializeTask(rows[0], false)
  778. task.jw_account_id = await this.inferJwAccountId(uuid, task.student_num, task.jw_account_id)
  779. const logs = await db.query('SELECT client_id, event, message, payload_json, create_time FROM qk_task_log WHERE task_id = ? ORDER BY create_time DESC LIMIT 100', [taskId])
  780. return {
  781. task,
  782. logs: (logs || []).map(log => {
  783. if (typeof log.payload_json === 'string' && log.payload_json) {
  784. try {
  785. log.payload_json = JSON.parse(log.payload_json)
  786. } catch (_) {}
  787. }
  788. return log
  789. })
  790. }
  791. }
  792. async cancelTask(uuid, taskId) {
  793. const result = await db.query(
  794. 'UPDATE qk_task SET status = ?, update_time = ?, finished_time = ? WHERE id = ? AND create_user = ? AND status IN (?, ?, ?)',
  795. [TASK_STATUS.CANCELLED, Date.now(), Date.now(), taskId, uuid, TASK_STATUS.PAUSED, TASK_STATUS.PENDING, TASK_STATUS.FAILED]
  796. )
  797. if (!result || result.affectedRows <= 0) {
  798. throw new Error('任务不存在或当前状态不可取消')
  799. }
  800. await this.logTask(taskId, null, 'cancelled', '用户取消抢课任务')
  801. this.logInfo('cancelTask', '用户已取消抢课任务', { taskId, uuid })
  802. }
  803. async startTask(uuid, taskId) {
  804. const rows = await db.query('SELECT * FROM qk_task WHERE id = ? AND create_user = ?', [taskId, uuid])
  805. if (!rows || rows.length === 0) {
  806. throw new Error('任务不存在')
  807. }
  808. const task = rows[0]
  809. if (![TASK_STATUS.PAUSED, TASK_STATUS.FAILED].includes(task.status)) {
  810. throw new Error('仅未开始或失败的任务可开启')
  811. }
  812. if (task.batch_id) {
  813. await this.requireEnabledBatch(task.batch_id)
  814. }
  815. const now = Date.now()
  816. const result = await db.query(
  817. `UPDATE qk_task SET status = ?, update_time = ?, assigned_client_id = NULL, lease_expire_at = NULL,
  818. assigned_at = NULL, exclude_client_id = NULL, result_json = NULL, error_msg = NULL, finished_time = NULL
  819. WHERE id = ? AND create_user = ? AND status IN (?, ?)`,
  820. [TASK_STATUS.PENDING, now, taskId, uuid, TASK_STATUS.PAUSED, TASK_STATUS.FAILED]
  821. )
  822. if (!result || result.affectedRows <= 0) {
  823. throw new Error('开启抢课任务失败')
  824. }
  825. await this.logTask(taskId, null, 'started', '用户开启抢课任务,等待分配')
  826. this.logInfo('startTask', '用户已开启抢课任务', { taskId, uuid })
  827. }
  828. async pauseTask(uuid, taskId) {
  829. const rows = await db.query('SELECT * FROM qk_task WHERE id = ? AND create_user = ?', [taskId, uuid])
  830. if (!rows || rows.length === 0) {
  831. throw new Error('任务不存在')
  832. }
  833. const task = rows[0]
  834. if (![TASK_STATUS.PENDING, TASK_STATUS.ASSIGNED, TASK_STATUS.RUNNING].includes(task.status)) {
  835. throw new Error('当前状态不可暂停')
  836. }
  837. const now = Date.now()
  838. const wasAssigned = [TASK_STATUS.ASSIGNED, TASK_STATUS.RUNNING].includes(task.status)
  839. const releasedClientId = wasAssigned ? task.assigned_client_id : null
  840. const result = await db.query(
  841. `UPDATE qk_task SET status = ?, update_time = ?, assigned_client_id = NULL, lease_expire_at = NULL,
  842. assigned_at = NULL, exclude_client_id = NULL
  843. WHERE id = ? AND create_user = ? AND status IN (?, ?, ?)`,
  844. [TASK_STATUS.PAUSED, now, taskId, uuid, TASK_STATUS.PENDING, TASK_STATUS.ASSIGNED, TASK_STATUS.RUNNING]
  845. )
  846. if (!result || result.affectedRows <= 0) {
  847. throw new Error('暂停抢课任务失败')
  848. }
  849. if (releasedClientId) {
  850. await this.syncClientSlots(releasedClientId, now)
  851. }
  852. await this.logTask(taskId, task.assigned_client_id, 'paused', wasAssigned ? '用户暂停任务,已收回客户端' : '用户暂停抢课任务')
  853. this.logInfo('pauseTask', '用户已暂停抢课任务', { taskId, uuid }, {
  854. released_from_client: wasAssigned ? task.assigned_client_id : null
  855. })
  856. }
  857. async getAdminTask(taskId) {
  858. const rows = await db.query(
  859. `SELECT t.*, u.username, u.avatar${this.taskBatchSelectSql('t')}
  860. FROM qk_task t
  861. LEFT JOIN users u ON u.uuid COLLATE utf8mb4_general_ci = t.create_user COLLATE utf8mb4_general_ci
  862. ${this.taskBatchJoinSql('t')}
  863. WHERE t.id = ?`,
  864. [taskId]
  865. )
  866. if (!rows || rows.length === 0) {
  867. return null
  868. }
  869. const task = this.serializeTask(rows[0], false)
  870. task.jw_account_id = await this.inferJwAccountId(rows[0].create_user, task.student_num, task.jw_account_id)
  871. return task
  872. }
  873. async adminUpdateTask(taskId, payload) {
  874. await this.ensureSchema()
  875. const rows = await db.query('SELECT * FROM qk_task WHERE id = ?', [taskId])
  876. if (!rows || rows.length === 0) {
  877. throw new Error('任务不存在')
  878. }
  879. const task = rows[0]
  880. if (task.status === TASK_STATUS.SUCCESS) {
  881. throw new Error('已成功的任务不可编辑')
  882. }
  883. const validated = this.validateTaskCoursesAndInterval(
  884. payload.courses || payload.COURSES,
  885. payload.course_groups || payload.COURSE_GROUPS,
  886. payload.interval_ms || payload.INTERVAL_MS || task.interval_ms || 500
  887. )
  888. const { courses, courseGroups, intervalMs } = validated
  889. const now = Date.now()
  890. const wasAssigned = [TASK_STATUS.ASSIGNED, TASK_STATUS.RUNNING].includes(task.status)
  891. const releasedClientId = wasAssigned ? task.assigned_client_id : null
  892. const batchId = this.resolveTaskBatchId(payload)
  893. const account = await this.resolveTaskAccount(task.create_user, payload, { requireAccount: true })
  894. const params = [
  895. payload.name,
  896. batchId,
  897. account.student_num,
  898. account.jw_account_id,
  899. this.encryptPassword(account.password),
  900. JSON.stringify(courses),
  901. JSON.stringify(courseGroups),
  902. payload.enable_ggxxk || payload.ENABLE_GGXXK ? 1 : 0,
  903. this.parseAutoRelogin(payload) ? 1 : 0,
  904. intervalMs,
  905. TASK_STATUS.PAUSED,
  906. now,
  907. taskId
  908. ]
  909. const sql = `UPDATE qk_task SET name = ?, batch_id = ?, student_num = ?, jw_account_id = ?, password_enc = ?, courses = ?, course_groups = ?, enable_ggxxk = ?, auto_relogin = ?, interval_ms = ?, status = ?, jx0502zbid = '', update_time = ?, assigned_client_id = NULL, lease_expire_at = NULL, assigned_at = NULL, exclude_client_id = NULL, result_json = NULL, error_msg = NULL, finished_time = NULL WHERE id = ?`
  910. const result = await db.query(sql, params)
  911. if (!result || result.affectedRows <= 0) {
  912. throw new Error('更新抢课任务失败')
  913. }
  914. if (releasedClientId) {
  915. await this.syncClientSlots(releasedClientId, now)
  916. }
  917. await this.logTask(taskId, task.assigned_client_id, 'admin_updated', wasAssigned ? '管理员更新任务并收回(未开始)' : '管理员更新抢课任务')
  918. this.logInfo('adminUpdateTask', '管理员已更新抢课任务', { taskId }, {
  919. name: payload.name,
  920. student_num: account.student_num,
  921. jw_account_id: account.jw_account_id,
  922. released_from_client: wasAssigned ? task.assigned_client_id : null
  923. })
  924. }
  925. async adminCancelTask(taskId) {
  926. const rows = await db.query('SELECT id, status, assigned_client_id FROM qk_task WHERE id = ?', [taskId])
  927. if (!rows || rows.length === 0) {
  928. throw new Error('任务不存在')
  929. }
  930. const task = rows[0]
  931. if (![TASK_STATUS.PAUSED, TASK_STATUS.PENDING, TASK_STATUS.ASSIGNED, TASK_STATUS.RUNNING, TASK_STATUS.FAILED].includes(task.status)) {
  932. throw new Error('当前状态不可取消')
  933. }
  934. const now = Date.now()
  935. const releasedClientId = [TASK_STATUS.ASSIGNED, TASK_STATUS.RUNNING].includes(task.status)
  936. ? task.assigned_client_id
  937. : null
  938. const result = await db.query(
  939. `UPDATE qk_task SET status = ?, assigned_client_id = NULL, lease_expire_at = NULL, assigned_at = NULL,
  940. exclude_client_id = NULL, update_time = ?, finished_time = ? WHERE id = ? AND status IN (?, ?, ?, ?, ?)`,
  941. [TASK_STATUS.CANCELLED, now, now, taskId, TASK_STATUS.PAUSED, TASK_STATUS.PENDING, TASK_STATUS.ASSIGNED, TASK_STATUS.RUNNING, TASK_STATUS.FAILED]
  942. )
  943. if (!result || result.affectedRows <= 0) {
  944. throw new Error('取消抢课任务失败')
  945. }
  946. if (releasedClientId) {
  947. await this.syncClientSlots(releasedClientId, now)
  948. }
  949. await this.logTask(taskId, task.assigned_client_id, 'admin_cancelled', '管理员取消抢课任务')
  950. this.logInfo('adminCancelTask', '管理员已取消抢课任务', { taskId })
  951. }
  952. async adminRetryTask(taskId) {
  953. const rows = await db.query('SELECT id, status FROM qk_task WHERE id = ?', [taskId])
  954. if (!rows || rows.length === 0) {
  955. throw new Error('任务不存在')
  956. }
  957. if (![TASK_STATUS.FAILED, TASK_STATUS.CANCELLED].includes(rows[0].status)) {
  958. throw new Error('仅失败或已取消的任务可重试')
  959. }
  960. const now = Date.now()
  961. const result = await db.query(
  962. `UPDATE qk_task SET status = ?, assigned_client_id = NULL, lease_expire_at = NULL, assigned_at = NULL,
  963. exclude_client_id = NULL, result_json = NULL, error_msg = NULL, finished_time = NULL, update_time = ?
  964. WHERE id = ? AND status IN (?, ?)`,
  965. [TASK_STATUS.PENDING, now, taskId, TASK_STATUS.FAILED, TASK_STATUS.CANCELLED]
  966. )
  967. if (!result || result.affectedRows <= 0) {
  968. throw new Error('重试抢课任务失败')
  969. }
  970. await this.logTask(taskId, null, 'admin_retry', '管理员将任务重新加入待分配队列')
  971. this.logInfo('adminRetryTask', '管理员已重试抢课任务', { taskId })
  972. }
  973. async adminStartTask(taskId) {
  974. const rows = await db.query('SELECT * FROM qk_task WHERE id = ?', [taskId])
  975. if (!rows || rows.length === 0) {
  976. throw new Error('任务不存在')
  977. }
  978. const task = rows[0]
  979. if (![TASK_STATUS.PAUSED, TASK_STATUS.FAILED].includes(task.status)) {
  980. throw new Error('仅未开始或失败的任务可开启')
  981. }
  982. if (task.batch_id) {
  983. await this.requireEnabledBatch(task.batch_id)
  984. }
  985. const now = Date.now()
  986. const result = await db.query(
  987. `UPDATE qk_task SET status = ?, update_time = ?, assigned_client_id = NULL, lease_expire_at = NULL,
  988. assigned_at = NULL, exclude_client_id = NULL, result_json = NULL, error_msg = NULL, finished_time = NULL
  989. WHERE id = ? AND status IN (?, ?)`,
  990. [TASK_STATUS.PENDING, now, taskId, TASK_STATUS.PAUSED, TASK_STATUS.FAILED]
  991. )
  992. if (!result || result.affectedRows <= 0) {
  993. throw new Error('开启抢课任务失败')
  994. }
  995. await this.logTask(taskId, null, 'admin_started', '管理员开启抢课任务,等待分配')
  996. this.logInfo('adminStartTask', '管理员已开启抢课任务', { taskId })
  997. }
  998. async adminPauseTask(taskId) {
  999. const rows = await db.query('SELECT * FROM qk_task WHERE id = ?', [taskId])
  1000. if (!rows || rows.length === 0) {
  1001. throw new Error('任务不存在')
  1002. }
  1003. const task = rows[0]
  1004. if (![TASK_STATUS.PENDING, TASK_STATUS.ASSIGNED, TASK_STATUS.RUNNING].includes(task.status)) {
  1005. throw new Error('当前状态不可暂停')
  1006. }
  1007. const now = Date.now()
  1008. const wasAssigned = [TASK_STATUS.ASSIGNED, TASK_STATUS.RUNNING].includes(task.status)
  1009. const releasedClientId = wasAssigned ? task.assigned_client_id : null
  1010. const result = await db.query(
  1011. `UPDATE qk_task SET status = ?, update_time = ?, assigned_client_id = NULL, lease_expire_at = NULL,
  1012. assigned_at = NULL, exclude_client_id = NULL
  1013. WHERE id = ? AND status IN (?, ?, ?)`,
  1014. [TASK_STATUS.PAUSED, now, taskId, TASK_STATUS.PENDING, TASK_STATUS.ASSIGNED, TASK_STATUS.RUNNING]
  1015. )
  1016. if (!result || result.affectedRows <= 0) {
  1017. throw new Error('暂停抢课任务失败')
  1018. }
  1019. if (releasedClientId) {
  1020. await this.syncClientSlots(releasedClientId, now)
  1021. }
  1022. await this.logTask(taskId, task.assigned_client_id, 'admin_paused', wasAssigned ? '管理员暂停任务,已收回客户端' : '管理员暂停抢课任务')
  1023. this.logInfo('adminPauseTask', '管理员已暂停抢课任务', { taskId }, {
  1024. released_from_client: wasAssigned ? task.assigned_client_id : null
  1025. })
  1026. }
  1027. async authenticateClient(clientId, clientSecret) {
  1028. if (!clientId || !clientSecret) {
  1029. this.logWarn('authClient', '客户端认证失败:缺少凭证', { clientId: clientId || 'unknown' })
  1030. return null
  1031. }
  1032. const rows = await db.query('SELECT * FROM qk_client WHERE client_id = ? AND enabled = 1', [clientId])
  1033. if (!rows || rows.length !== 1) {
  1034. this.logWarn('authClient', '客户端认证失败:客户端不存在或已禁用', { clientId })
  1035. return null
  1036. }
  1037. if (!bcryptjs.compareSync(String(clientSecret), rows[0].client_secret_hash)) {
  1038. this.logWarn('authClient', '客户端认证失败:密钥不匹配', { clientId })
  1039. return null
  1040. }
  1041. return rows[0]
  1042. }
  1043. async enrollOrAuthenticateClient(clientId, clientSecret, payload = {}) {
  1044. const existing = await this.authenticateClient(clientId, clientSecret)
  1045. if (existing) {
  1046. return existing
  1047. }
  1048. if (!clientId || !clientSecret || !String(clientId).startsWith('qk-cli-')) {
  1049. throw new Error('客户端凭证无效')
  1050. }
  1051. const rows = await db.query('SELECT * FROM qk_client WHERE client_id = ?', [clientId])
  1052. if (rows && rows.length > 0) {
  1053. throw new Error('客户端凭证无效')
  1054. }
  1055. const time = Date.now()
  1056. const label = payload.label || payload.hostname || `auto-${clientId}`
  1057. try {
  1058. const result = await db.query(
  1059. 'INSERT INTO qk_client (client_id, client_secret_hash, label, create_time, update_time) VALUES (?, ?, ?, ?, ?)',
  1060. [clientId, bcryptjs.hashSync(String(clientSecret), 10), label, time, time]
  1061. )
  1062. if (!result || result.affectedRows <= 0) {
  1063. throw new Error('客户端自动注册失败')
  1064. }
  1065. } catch (err) {
  1066. if (err?.code === 'ER_DUP_ENTRY') {
  1067. const raced = await this.authenticateClient(clientId, clientSecret)
  1068. if (raced) {
  1069. return raced
  1070. }
  1071. }
  1072. throw err
  1073. }
  1074. const created = await db.query('SELECT * FROM qk_client WHERE client_id = ? AND enabled = 1', [clientId])
  1075. if (!created || created.length !== 1) {
  1076. throw new Error('客户端自动注册失败')
  1077. }
  1078. this.logInfo('enrollClient', '抢课客户端已自动注册', { clientId }, { label })
  1079. return created[0]
  1080. }
  1081. async registerClient(clientId, clientSecret, payload = {}) {
  1082. const client = await this.enrollOrAuthenticateClient(clientId, clientSecret, payload)
  1083. const maxSlots = this.resolveClientMaxSlots(client, payload)
  1084. const time = Date.now()
  1085. await db.query(
  1086. 'UPDATE qk_client SET label = COALESCE(NULLIF(?, \'\'), label), max_slots = ?, hostname = ?, os_username = ?, cpu_model = ?, cpu_threads = ?, total_mem_mb = ?, free_mem_mb = ?, platform = ?, last_heartbeat_at = ?, online = 1, update_time = ? WHERE client_id = ?',
  1087. [
  1088. payload.label || '',
  1089. maxSlots,
  1090. payload.hostname || null,
  1091. payload.os_username || null,
  1092. payload.cpu_model || null,
  1093. payload.cpu_threads || null,
  1094. payload.total_mem_mb || null,
  1095. payload.free_mem_mb || null,
  1096. payload.platform || null,
  1097. time,
  1098. time,
  1099. clientId
  1100. ]
  1101. )
  1102. await Redis.set(`qk:client:hb:${clientId}`, String(time), { EX: this.heartbeatTtlSeconds })
  1103. const currentSlots = await this.syncClientSlots(clientId, time)
  1104. this.logInfo('registerClient', '抢课客户端已注册/上线', { clientId }, {
  1105. label: payload.label || client.label,
  1106. max_slots: maxSlots,
  1107. current_slots: currentSlots,
  1108. hostname: payload.hostname,
  1109. platform: payload.platform,
  1110. cpu_threads: payload.cpu_threads,
  1111. total_mem_mb: payload.total_mem_mb,
  1112. free_mem_mb: payload.free_mem_mb
  1113. })
  1114. return { client_id: clientId, max_slots: maxSlots, current_slots: currentSlots }
  1115. }
  1116. async heartbeat(clientId, clientSecret, payload = {}) {
  1117. const client = await this.authenticateClient(clientId, clientSecret)
  1118. if (!client) {
  1119. throw new Error('客户端凭证无效')
  1120. }
  1121. const runningTasks = Array.isArray(payload.running_tasks) ? payload.running_tasks.map(Number).filter(Boolean) : []
  1122. const revokedTaskIds = await this.findRevokedClientTasks(clientId, runningTasks)
  1123. const activeRunningTasks = runningTasks.filter(id => !revokedTaskIds.includes(Number(id)))
  1124. const maxSlots = this.resolveClientMaxSlots(client, payload)
  1125. const now = Date.now()
  1126. await db.query(
  1127. 'UPDATE qk_client SET max_slots = ?, hostname = ?, os_username = ?, cpu_model = ?, cpu_threads = ?, total_mem_mb = ?, free_mem_mb = ?, platform = ?, last_heartbeat_at = ?, online = 1, update_time = ? WHERE client_id = ?',
  1128. [
  1129. maxSlots,
  1130. payload.hostname || null,
  1131. payload.os_username || null,
  1132. payload.cpu_model || null,
  1133. payload.cpu_threads || null,
  1134. payload.total_mem_mb || null,
  1135. payload.free_mem_mb || null,
  1136. payload.platform || null,
  1137. now,
  1138. now,
  1139. clientId
  1140. ]
  1141. )
  1142. await Redis.set(`qk:client:hb:${clientId}`, String(now), { EX: this.heartbeatTtlSeconds })
  1143. if (activeRunningTasks.length > 0) {
  1144. const leaseExpireAt = now + this.leaseMs
  1145. const placeholders = activeRunningTasks.map(() => '?').join(',')
  1146. await db.query(
  1147. `UPDATE qk_task SET status = ?, lease_expire_at = ?, update_time = ? WHERE assigned_client_id = ? AND id IN (${placeholders}) AND status IN (?, ?)`,
  1148. [TASK_STATUS.RUNNING, leaseExpireAt, now, clientId, ...activeRunningTasks, TASK_STATUS.ASSIGNED, TASK_STATUS.RUNNING]
  1149. )
  1150. this.logInfo('heartbeat', '客户端心跳续租运行中任务', { clientId }, {
  1151. running_tasks: activeRunningTasks,
  1152. revoked_task_ids: revokedTaskIds,
  1153. lease_expire_at: leaseExpireAt,
  1154. max_slots: maxSlots
  1155. })
  1156. } else if (revokedTaskIds.length > 0) {
  1157. this.logInfo('heartbeat', '客户端上报任务已被服务端收回', { clientId }, {
  1158. revoked_task_ids: revokedTaskIds
  1159. })
  1160. }
  1161. await this.releaseOrphanedClientTasks(clientId, activeRunningTasks, now)
  1162. const currentSlots = await this.syncClientSlots(clientId, now)
  1163. return { client_id: clientId, current_slots: currentSlots, max_slots: maxSlots, revoked_task_ids: revokedTaskIds }
  1164. }
  1165. async releaseOrphanedClientTasks(clientId, runningTasks, now = Date.now()) {
  1166. const assigned = await db.query(
  1167. 'SELECT id FROM qk_task WHERE assigned_client_id = ? AND status IN (?, ?)',
  1168. [clientId, TASK_STATUS.ASSIGNED, TASK_STATUS.RUNNING]
  1169. )
  1170. const runningSet = new Set((runningTasks || []).map(Number).filter(Boolean))
  1171. let released = 0
  1172. for (const row of assigned || []) {
  1173. if (runningSet.has(row.id)) continue
  1174. const result = await db.query(
  1175. 'UPDATE qk_task SET status = ?, assigned_client_id = NULL, lease_expire_at = NULL, assigned_at = NULL, update_time = ? WHERE id = ? AND assigned_client_id = ? AND status IN (?, ?)',
  1176. [TASK_STATUS.PENDING, now, row.id, clientId, TASK_STATUS.ASSIGNED, TASK_STATUS.RUNNING]
  1177. )
  1178. if (result && result.affectedRows > 0) {
  1179. released += 1
  1180. await this.logTask(row.id, clientId, 'released', '客户端未继续执行,任务已释放回队列')
  1181. this.logInfo('releaseOrphaned', '释放未在运行的已分配任务', { taskId: row.id, clientId })
  1182. }
  1183. }
  1184. if (released > 0) {
  1185. await this.syncClientSlots(clientId, now)
  1186. }
  1187. return released
  1188. }
  1189. async reclaimTasks(clientId, clientSecret) {
  1190. const client = await this.authenticateClient(clientId, clientSecret)
  1191. if (!client) {
  1192. throw new Error('客户端凭证无效')
  1193. }
  1194. const now = Date.now()
  1195. const leaseExpireAt = now + this.leaseMs
  1196. const rows = await db.query(
  1197. `SELECT t.*, b.jx0502zbid AS batch_jx0502zbid, b.name AS batch_name, b.enabled AS batch_enabled
  1198. FROM qk_task t
  1199. LEFT JOIN qk_batch b ON b.id = t.batch_id
  1200. WHERE t.assigned_client_id = ? AND t.status IN (?, ?)
  1201. ORDER BY t.create_time ASC`,
  1202. [clientId, TASK_STATUS.ASSIGNED, TASK_STATUS.RUNNING]
  1203. )
  1204. if (!rows || rows.length === 0) {
  1205. return []
  1206. }
  1207. for (const row of rows) {
  1208. await db.query(
  1209. 'UPDATE qk_task SET status = ?, lease_expire_at = ?, update_time = ? WHERE id = ? AND assigned_client_id = ? AND status IN (?, ?)',
  1210. [TASK_STATUS.RUNNING, leaseExpireAt, now, row.id, clientId, TASK_STATUS.ASSIGNED, TASK_STATUS.RUNNING]
  1211. )
  1212. await this.logTask(row.id, clientId, 'reclaimed', '客户端重连后回收任务继续执行')
  1213. }
  1214. this.logInfo('reclaimTasks', `客户端回收 ${rows.length} 个进行中任务`, { clientId }, {
  1215. task_ids: rows.map(row => row.id),
  1216. lease_expire_at: leaseExpireAt
  1217. })
  1218. return rows.map(row => this.serializeTask({
  1219. ...row,
  1220. assigned_client_id: clientId,
  1221. lease_expire_at: leaseExpireAt,
  1222. status: TASK_STATUS.RUNNING
  1223. }, true))
  1224. }
  1225. async releaseTasks(clientId, clientSecret, payload = {}) {
  1226. const client = await this.authenticateClient(clientId, clientSecret)
  1227. if (!client) {
  1228. throw new Error('客户端凭证无效')
  1229. }
  1230. const now = Date.now()
  1231. const taskIds = Array.isArray(payload.task_ids)
  1232. ? payload.task_ids.map(Number).filter(Boolean)
  1233. : []
  1234. let rows = []
  1235. if (taskIds.length > 0) {
  1236. const placeholders = taskIds.map(() => '?').join(',')
  1237. rows = await db.query(
  1238. `SELECT id FROM qk_task WHERE assigned_client_id = ? AND id IN (${placeholders}) AND status IN (?, ?)`,
  1239. [clientId, ...taskIds, TASK_STATUS.ASSIGNED, TASK_STATUS.RUNNING]
  1240. )
  1241. } else {
  1242. rows = await db.query(
  1243. 'SELECT id FROM qk_task WHERE assigned_client_id = ? AND status IN (?, ?)',
  1244. [clientId, TASK_STATUS.ASSIGNED, TASK_STATUS.RUNNING]
  1245. )
  1246. }
  1247. let released = 0
  1248. for (const row of rows || []) {
  1249. const result = await db.query(
  1250. 'UPDATE qk_task SET status = ?, assigned_client_id = NULL, lease_expire_at = NULL, assigned_at = NULL, update_time = ? WHERE id = ? AND assigned_client_id = ? AND status IN (?, ?)',
  1251. [TASK_STATUS.PENDING, now, row.id, clientId, TASK_STATUS.ASSIGNED, TASK_STATUS.RUNNING]
  1252. )
  1253. if (result && result.affectedRows > 0) {
  1254. released += 1
  1255. await this.logTask(row.id, clientId, 'released', '客户端主动释放任务')
  1256. }
  1257. }
  1258. if (released > 0) {
  1259. await this.syncClientSlots(clientId, now)
  1260. this.logInfo('releaseTasks', `客户端释放 ${released} 个任务`, { clientId }, {
  1261. task_ids: (rows || []).map(row => row.id)
  1262. })
  1263. }
  1264. return { released, task_ids: (rows || []).map(row => row.id) }
  1265. }
  1266. async pullTasks(clientId, clientSecret, count) {
  1267. const client = await this.authenticateClient(clientId, clientSecret)
  1268. if (!client) {
  1269. throw new Error('客户端凭证无效')
  1270. }
  1271. const availableSlots = await this.getClientAvailableSlotsAsync(client)
  1272. const safeCount = Math.max(0, Math.min(50, Number(count || 0), availableSlots))
  1273. if (safeCount <= 0) {
  1274. this.logWarn('pullTasks', '客户端无可用槽位或拉取数量无效,已忽略', { clientId }, {
  1275. count,
  1276. available_slots: availableSlots,
  1277. max_slots: this.resolveClientMaxSlots(client),
  1278. current_slots: client.current_slots,
  1279. free_mem_mb: client.free_mem_mb
  1280. })
  1281. return []
  1282. }
  1283. const lockKey = `qk:pull:lock:${clientId}`
  1284. const locked = await Redis.set(lockKey, '1', { NX: true, EX: this.pullLockTtlSeconds })
  1285. if (!locked) {
  1286. this.logWarn('pullTasks', '拉取任务被并发锁拦截,本次跳过', { clientId }, { count: safeCount })
  1287. return []
  1288. }
  1289. const conn = await db.connect()
  1290. try {
  1291. await conn.beginTransaction()
  1292. const inFlightStatuses = this.getInFlightTaskStatuses()
  1293. const candidateLimit = Math.min(50, Math.max(safeCount, safeCount * 5))
  1294. const [candidateRows] = await conn.execute(
  1295. `SELECT t.*, b.jx0502zbid AS batch_jx0502zbid, b.name AS batch_name, b.enabled AS batch_enabled
  1296. FROM qk_task t
  1297. LEFT JOIN qk_batch b ON b.id = t.batch_id
  1298. WHERE t.status = ?
  1299. AND (t.batch_id IS NULL OR b.enabled = 1)
  1300. AND (t.exclude_client_id IS NULL OR t.exclude_client_id <> ?)
  1301. AND NOT EXISTS (
  1302. SELECT 1 FROM qk_task active
  1303. WHERE active.student_num = t.student_num
  1304. AND active.status IN (?, ?)
  1305. )
  1306. ORDER BY t.create_time ASC LIMIT ${candidateLimit} FOR UPDATE`,
  1307. [TASK_STATUS.PENDING, clientId, ...inFlightStatuses]
  1308. )
  1309. const rows = this.pickPullableTasks(candidateRows, safeCount)
  1310. const now = Date.now()
  1311. const leaseExpireAt = now + this.leaseMs
  1312. for (const row of rows) {
  1313. await conn.execute(
  1314. 'UPDATE qk_task SET status = ?, assigned_client_id = ?, lease_expire_at = ?, assigned_at = ?, exclude_client_id = NULL, update_time = ? WHERE id = ?',
  1315. [TASK_STATUS.ASSIGNED, clientId, leaseExpireAt, now, now, row.id]
  1316. )
  1317. }
  1318. await conn.commit()
  1319. if (rows.length > 0) {
  1320. await this.syncClientSlots(clientId, now)
  1321. }
  1322. for (const row of rows) {
  1323. await this.logTask(row.id, clientId, 'assigned', '任务已分配给客户端')
  1324. }
  1325. if (rows.length > 0) {
  1326. this.logInfo('pullTasks', `已分配 ${rows.length} 个抢课任务`, { clientId }, {
  1327. task_ids: rows.map(row => row.id),
  1328. lease_expire_at: leaseExpireAt,
  1329. requested_count: safeCount
  1330. })
  1331. }
  1332. return rows.map(row => this.serializeTask({ ...row, assigned_client_id: clientId, lease_expire_at: leaseExpireAt, status: TASK_STATUS.ASSIGNED }, true))
  1333. } catch (err) {
  1334. await conn.rollback()
  1335. this.logError('pullTasks', '拉取并分配任务失败,事务已回滚', { clientId }, err)
  1336. throw err
  1337. } finally {
  1338. await Redis.del(lockKey)
  1339. }
  1340. }
  1341. extractFailureMessage(payload = {}) {
  1342. const result = payload.result && typeof payload.result === 'object' ? payload.result : {}
  1343. const parts = [
  1344. payload.message,
  1345. payload.error_msg,
  1346. payload.error,
  1347. result.message,
  1348. result.error,
  1349. result.error_msg,
  1350. typeof payload.result === 'string' ? payload.result : null
  1351. ]
  1352. return parts.filter(Boolean).map(item => String(item)).join(' ').trim()
  1353. }
  1354. resolveFailureType(payload = {}) {
  1355. const result = payload.result && typeof payload.result === 'object' ? payload.result : {}
  1356. return String(payload.type || payload.event || result.type || result.event || '').trim()
  1357. }
  1358. shouldRequeueTaskOnFailure(task, payload = {}) {
  1359. if (payload.success === true) {
  1360. return false
  1361. }
  1362. if (payload.requeue === true || payload.requeue === 1 || payload.requeue === '1') {
  1363. return true
  1364. }
  1365. const message = this.extractFailureMessage(payload)
  1366. const failureType = this.resolveFailureType(payload)
  1367. const terminalTypes = new Set(['already', 'dajia', 'not_open', 'not_in_time'])
  1368. if (terminalTypes.has(failureType)) {
  1369. return false
  1370. }
  1371. const terminalPatterns = [
  1372. '已选择',
  1373. '冲突',
  1374. '超过',
  1375. '选课不开放',
  1376. '不在选课时间',
  1377. '用户名或密码',
  1378. '未匹配到目标课程',
  1379. '验证码验证失败次数',
  1380. '为避免账号被锁定'
  1381. ]
  1382. if (terminalPatterns.some(pattern => message.includes(pattern))) {
  1383. return false
  1384. }
  1385. // 旧版客户端 task_exit 上报:抢课任务异常结束,code=1
  1386. if (/抢课任务异常结束/.test(message) || /\bcode=1\b/.test(message)) {
  1387. return true
  1388. }
  1389. if (/抢课循环结束但未获得成功结果/.test(message)) {
  1390. return true
  1391. }
  1392. // 客户端被收回槽位、窗口关闭等导致的被动停止,应重新分配而非标为失败
  1393. const exitReason = String(payload.exit_reason || payload.reason || '').trim()
  1394. if (failureType === 'stopped' || exitReason === 'stopped' || message.includes('任务已停止')) {
  1395. return true
  1396. }
  1397. const autoRelogin = Number(task?.auto_relogin) === 1
  1398. if (autoRelogin) {
  1399. // 新版可传 type=login_expired;旧版仅 message / result.message
  1400. if (failureType === 'login_expired') {
  1401. return true
  1402. }
  1403. const loginPatterns = [
  1404. '别处登录',
  1405. '在其他地方登录',
  1406. '账号在别处',
  1407. '登录失效',
  1408. '登录已失效',
  1409. '请重新登录',
  1410. '按登录失效处理',
  1411. '空响应',
  1412. '未登录',
  1413. '登录失败',
  1414. '会话过期',
  1415. '会话失效'
  1416. ]
  1417. if (loginPatterns.some(pattern => message.includes(pattern))) {
  1418. return true
  1419. }
  1420. }
  1421. return false
  1422. }
  1423. async requeueTaskAfterClientFailure(clientId, taskId, payload = {}) {
  1424. const now = Date.now()
  1425. const message = this.buildReportMessage(payload) || '任务异常结束,已重新进入待分配队列'
  1426. const result = await db.query(
  1427. `UPDATE qk_task SET status = ?, assigned_client_id = NULL, lease_expire_at = NULL,
  1428. assigned_at = NULL, result_json = NULL, error_msg = ?, finished_time = NULL, update_time = ?
  1429. WHERE id = ? AND assigned_client_id = ? AND status IN (?, ?)`,
  1430. [TASK_STATUS.PENDING, message, now, taskId, clientId, TASK_STATUS.ASSIGNED, TASK_STATUS.RUNNING]
  1431. )
  1432. if (!result || result.affectedRows <= 0) {
  1433. throw new Error('任务不存在或不属于当前客户端')
  1434. }
  1435. await this.syncClientSlots(clientId, now)
  1436. await this.logTask(taskId, clientId, 'requeued', message, payload.result || payload)
  1437. this.logInfo('reportResult', '抢课任务异常结束,已重新入队', { taskId, clientId }, {
  1438. message,
  1439. requeued: true,
  1440. exit_reason: payload.exit_reason || payload.reason || null
  1441. })
  1442. return { task_id: taskId, status: TASK_STATUS.PENDING, requeued: true }
  1443. }
  1444. async reportResult(clientId, clientSecret, payload = {}) {
  1445. await this.ensureSchema()
  1446. const client = await this.authenticateClient(clientId, clientSecret)
  1447. if (!client) {
  1448. throw new Error('客户端凭证无效')
  1449. }
  1450. const taskId = Number(payload.task_id || payload.id)
  1451. const success = payload.success === true || payload.status === TASK_STATUS.SUCCESS
  1452. const taskRows = await db.query(
  1453. 'SELECT id, auto_relogin FROM qk_task WHERE id = ? AND assigned_client_id = ? AND status IN (?, ?)',
  1454. [taskId, clientId, TASK_STATUS.ASSIGNED, TASK_STATUS.RUNNING]
  1455. )
  1456. if (!taskRows || taskRows.length === 0) {
  1457. this.logWarn('reportResult', '任务结果上报被拒绝:任务不存在或不属于当前客户端', { taskId, clientId }, {
  1458. success,
  1459. status: success ? TASK_STATUS.SUCCESS : TASK_STATUS.FAILED
  1460. })
  1461. throw new Error('任务不存在或不属于当前客户端')
  1462. }
  1463. const task = taskRows[0]
  1464. if (!success && this.shouldRequeueTaskOnFailure(task, payload)) {
  1465. return this.requeueTaskAfterClientFailure(clientId, taskId, payload)
  1466. }
  1467. const status = success ? TASK_STATUS.SUCCESS : TASK_STATUS.FAILED
  1468. const now = Date.now()
  1469. const resultJson = payload.result ? JSON.stringify(payload.result) : JSON.stringify({
  1470. course: payload.course || '',
  1471. message: payload.message || ''
  1472. })
  1473. const result = await db.query(
  1474. 'UPDATE qk_task SET status = ?, result_json = ?, error_msg = ?, update_time = ?, finished_time = ?, lease_expire_at = NULL WHERE id = ? AND assigned_client_id = ? AND status IN (?, ?)',
  1475. [status, resultJson, success ? null : (payload.error_msg || payload.message || '抢课失败'), now, now, taskId, clientId, TASK_STATUS.ASSIGNED, TASK_STATUS.RUNNING]
  1476. )
  1477. if (!result || result.affectedRows <= 0) {
  1478. throw new Error('任务不存在或不属于当前客户端')
  1479. }
  1480. await this.syncClientSlots(clientId, now)
  1481. await db.query(
  1482. `UPDATE qk_client SET total_completed = total_completed + 1, total_success = total_success + ?, update_time = ? WHERE client_id = ?`,
  1483. [success ? 1 : 0, now, clientId]
  1484. )
  1485. await this.logTask(taskId, clientId, success ? 'grab_success' : 'grab_fail', payload.message || payload.error_msg || '', payload.result || payload)
  1486. this.logInfo('reportResult', success ? '抢课成功' : '抢课失败', { taskId, clientId }, {
  1487. success,
  1488. message: payload.message || payload.error_msg || '',
  1489. course: payload.course || payload.result?.course || '',
  1490. result: payload.result || null
  1491. })
  1492. return { task_id: taskId, status }
  1493. }
  1494. async listClients() {
  1495. const rows = await db.query('SELECT id, client_id, label, enabled, max_slots, current_slots, hostname, os_username, cpu_model, cpu_threads, total_mem_mb, free_mem_mb, platform, last_heartbeat_at, online, total_completed, total_success, create_time, update_time FROM qk_client ORDER BY update_time DESC')
  1496. return rows || []
  1497. }
  1498. async deleteClient(clientId) {
  1499. const result = await db.query('UPDATE qk_client SET enabled = 0, online = 0, update_time = ? WHERE client_id = ?', [Date.now(), clientId])
  1500. if (!result || result.affectedRows <= 0) {
  1501. throw new Error('客户端不存在')
  1502. }
  1503. this.logInfo('deleteClient', '抢课客户端已禁用', { clientId })
  1504. }
  1505. async enableClient(clientId) {
  1506. const rows = await db.query('SELECT client_id, enabled FROM qk_client WHERE client_id = ?', [clientId])
  1507. if (!rows || rows.length === 0) {
  1508. throw new Error('客户端不存在')
  1509. }
  1510. if (Number(rows[0].enabled) === 1) {
  1511. throw new Error('客户端已处于启用状态')
  1512. }
  1513. const now = Date.now()
  1514. await db.query('UPDATE qk_client SET enabled = 1, update_time = ? WHERE client_id = ?', [now, clientId])
  1515. this.logInfo('enableClient', '抢课客户端已解禁', { clientId })
  1516. }
  1517. async listAdminTasks(filters = {}) {
  1518. const pagesize = Math.max(1, Math.min(100, Number(filters.pagesize || 20)))
  1519. const current = Math.max(1, Number(filters.current || 1))
  1520. const where = ['1 = 1']
  1521. const params = []
  1522. const countParams = []
  1523. if (filters.status) {
  1524. where.push('t.status = ?')
  1525. params.push(filters.status)
  1526. countParams.push(filters.status)
  1527. }
  1528. if (filters.client_id) {
  1529. where.push('t.assigned_client_id = ?')
  1530. params.push(filters.client_id)
  1531. countParams.push(filters.client_id)
  1532. }
  1533. if (filters.student_num) {
  1534. where.push('t.student_num LIKE ?')
  1535. params.push(`%${filters.student_num}%`)
  1536. countParams.push(`%${filters.student_num}%`)
  1537. }
  1538. if (filters.username) {
  1539. where.push('u.username COLLATE utf8mb4_general_ci LIKE (CONVERT(? USING utf8mb4) COLLATE utf8mb4_general_ci)')
  1540. params.push(`%${filters.username}%`)
  1541. countParams.push(`%${filters.username}%`)
  1542. }
  1543. const courseName = String(filters.course_name || filters.course || '').trim()
  1544. if (courseName) {
  1545. const keyword = `%${courseName}%`
  1546. where.push('(t.courses LIKE ? OR t.course_groups LIKE ?)')
  1547. params.push(keyword, keyword)
  1548. countParams.push(keyword, keyword)
  1549. }
  1550. if (filters.batch_id) {
  1551. where.push('t.batch_id = ?')
  1552. params.push(Number(filters.batch_id))
  1553. countParams.push(Number(filters.batch_id))
  1554. }
  1555. const whereSql = where.join(' AND ')
  1556. const offset = (current - 1) * pagesize
  1557. const countRows = await db.query(
  1558. `SELECT COUNT(*) AS total FROM qk_task t
  1559. LEFT JOIN users u ON u.uuid COLLATE utf8mb4_general_ci = t.create_user COLLATE utf8mb4_general_ci
  1560. WHERE ${whereSql}`,
  1561. countParams
  1562. )
  1563. const rows = await db.query(
  1564. `SELECT t.*, u.username, u.avatar${this.taskBatchSelectSql('t')}
  1565. FROM qk_task t
  1566. LEFT JOIN users u ON u.uuid COLLATE utf8mb4_general_ci = t.create_user COLLATE utf8mb4_general_ci
  1567. ${this.taskBatchJoinSql('t')}
  1568. WHERE ${whereSql}
  1569. ORDER BY t.create_time DESC
  1570. LIMIT ${pagesize} OFFSET ${offset}`,
  1571. params
  1572. )
  1573. return {
  1574. list: (rows || []).map(row => this.serializeTask(row, true)),
  1575. total: countRows?.[0]?.total || 0,
  1576. current,
  1577. pagesize
  1578. }
  1579. }
  1580. async listAdminReports(filters = {}) {
  1581. const pagesize = Math.max(1, Math.min(100, Number(filters.pagesize || 20)))
  1582. const current = Math.max(1, Number(filters.current || 1))
  1583. const where = ['l.event IN (?, ?, ?, ?)']
  1584. const params = ['request_result', 'progress_snapshot', 'grab_success', 'grab_fail']
  1585. const countParams = ['request_result', 'progress_snapshot', 'grab_success', 'grab_fail']
  1586. if (filters.task_id) {
  1587. where.push('l.task_id = ?')
  1588. params.push(Number(filters.task_id))
  1589. countParams.push(Number(filters.task_id))
  1590. }
  1591. if (filters.client_id) {
  1592. where.push('l.client_id = ?')
  1593. params.push(String(filters.client_id))
  1594. countParams.push(String(filters.client_id))
  1595. }
  1596. if (filters.event) {
  1597. where.push('l.event = ?')
  1598. params.push(String(filters.event))
  1599. countParams.push(String(filters.event))
  1600. }
  1601. if (filters.start_time) {
  1602. where.push('l.create_time >= ?')
  1603. params.push(Number(filters.start_time))
  1604. countParams.push(Number(filters.start_time))
  1605. }
  1606. if (filters.end_time) {
  1607. where.push('l.create_time <= ?')
  1608. params.push(Number(filters.end_time))
  1609. countParams.push(Number(filters.end_time))
  1610. }
  1611. if (filters.student_num) {
  1612. where.push('t.student_num LIKE ?')
  1613. params.push(`%${filters.student_num}%`)
  1614. countParams.push(`%${filters.student_num}%`)
  1615. }
  1616. if (filters.name || filters.task_name) {
  1617. where.push('t.name LIKE ?')
  1618. params.push(`%${filters.name || filters.task_name}%`)
  1619. countParams.push(`%${filters.name || filters.task_name}%`)
  1620. }
  1621. if (filters.username) {
  1622. where.push('u.username COLLATE utf8mb4_general_ci LIKE (CONVERT(? USING utf8mb4) COLLATE utf8mb4_general_ci)')
  1623. params.push(`%${filters.username}%`)
  1624. countParams.push(`%${filters.username}%`)
  1625. }
  1626. if (filters.client_label) {
  1627. where.push('c.label LIKE ?')
  1628. params.push(`%${filters.client_label}%`)
  1629. countParams.push(`%${filters.client_label}%`)
  1630. }
  1631. const whereSql = where.join(' AND ')
  1632. const offset = (current - 1) * pagesize
  1633. const joinSql = `
  1634. FROM qk_task_log l
  1635. LEFT JOIN qk_task t ON t.id = l.task_id
  1636. LEFT JOIN qk_client c ON c.client_id = l.client_id
  1637. LEFT JOIN users u ON u.uuid COLLATE utf8mb4_general_ci = t.create_user COLLATE utf8mb4_general_ci`
  1638. const countRows = await db.query(
  1639. `SELECT COUNT(*) AS total ${joinSql} WHERE ${whereSql}`,
  1640. countParams
  1641. )
  1642. const rows = await db.query(
  1643. `SELECT l.id, l.task_id, l.client_id, l.event, l.message, l.payload_json, l.create_time,
  1644. t.name AS task_name, t.student_num, t.status AS task_status,
  1645. c.label AS client_label, u.username, u.avatar
  1646. ${joinSql}
  1647. WHERE ${whereSql}
  1648. ORDER BY l.create_time DESC
  1649. LIMIT ${pagesize} OFFSET ${offset}`,
  1650. params
  1651. )
  1652. return {
  1653. list: (rows || []).map(row => this.serializeReportLog(row)),
  1654. total: countRows?.[0]?.total || 0,
  1655. current,
  1656. pagesize
  1657. }
  1658. }
  1659. async reassignStaleRunningTasks() {
  1660. const onlineCount = await this.countOnlineClients()
  1661. if (onlineCount <= 1) {
  1662. return 0
  1663. }
  1664. const now = Date.now()
  1665. const staleBefore = now - this.staleTaskMs
  1666. const rows = await db.query(
  1667. `SELECT id, assigned_client_id FROM qk_task
  1668. WHERE status IN (?, ?)
  1669. AND assigned_at IS NOT NULL
  1670. AND assigned_at < ?`,
  1671. [TASK_STATUS.ASSIGNED, TASK_STATUS.RUNNING, staleBefore]
  1672. )
  1673. let count = 0
  1674. const affectedClients = new Set()
  1675. for (const row of rows || []) {
  1676. const result = await db.query(
  1677. `UPDATE qk_task
  1678. SET status = ?, assigned_client_id = NULL, lease_expire_at = NULL,
  1679. assigned_at = NULL, exclude_client_id = ?, update_time = ?
  1680. WHERE id = ? AND status IN (?, ?)`,
  1681. [TASK_STATUS.PENDING, row.assigned_client_id, now, row.id, TASK_STATUS.ASSIGNED, TASK_STATUS.RUNNING]
  1682. )
  1683. if (result && result.affectedRows > 0) {
  1684. count += 1
  1685. if (row.assigned_client_id) {
  1686. affectedClients.add(row.assigned_client_id)
  1687. }
  1688. await this.logTask(row.id, row.assigned_client_id, 'reassigned', '任务超过1小时未成功,已收回并等待分配给其他客户端')
  1689. this.logWarn('reassignStale', '长时间未成功任务已收回', {
  1690. taskId: row.id,
  1691. clientId: row.assigned_client_id
  1692. }, { online_clients: onlineCount })
  1693. }
  1694. }
  1695. for (const clientId of affectedClients) {
  1696. await this.syncClientSlots(clientId, now)
  1697. }
  1698. return count
  1699. }
  1700. async requeueExpiredTasks() {
  1701. const now = Date.now()
  1702. const rows = await db.query(
  1703. 'SELECT id, assigned_client_id FROM qk_task WHERE status IN (?, ?) AND lease_expire_at IS NOT NULL AND lease_expire_at < ?',
  1704. [TASK_STATUS.ASSIGNED, TASK_STATUS.RUNNING, now]
  1705. )
  1706. let count = 0
  1707. const affectedClients = new Set()
  1708. for (const row of rows || []) {
  1709. const result = await db.query(
  1710. 'UPDATE qk_task SET status = ?, assigned_client_id = NULL, lease_expire_at = NULL, assigned_at = NULL, update_time = ? WHERE id = ? AND status IN (?, ?)',
  1711. [TASK_STATUS.PENDING, now, row.id, TASK_STATUS.ASSIGNED, TASK_STATUS.RUNNING]
  1712. )
  1713. if (result && result.affectedRows > 0) {
  1714. count++
  1715. if (row.assigned_client_id) {
  1716. affectedClients.add(row.assigned_client_id)
  1717. }
  1718. await this.logTask(row.id, row.assigned_client_id, 'reassigned', '租约过期,任务已重新进入待分配队列')
  1719. this.logWarn('requeueExpired', '租约过期,任务已重新入队', {
  1720. taskId: row.id,
  1721. clientId: row.assigned_client_id
  1722. })
  1723. }
  1724. }
  1725. for (const clientId of affectedClients) {
  1726. await this.syncClientSlots(clientId, now)
  1727. }
  1728. const offlineResult = await db.query(
  1729. 'UPDATE qk_client SET online = 0, current_slots = 0, update_time = ? WHERE online = 1 AND (last_heartbeat_at IS NULL OR last_heartbeat_at < ?)',
  1730. [now, now - this.heartbeatTtlSeconds * 1000]
  1731. )
  1732. const offlineCount = offlineResult?.affectedRows || 0
  1733. const staleCount = await this.reassignStaleRunningTasks()
  1734. if (count > 0 || offlineCount > 0 || staleCount > 0) {
  1735. this.logInfo('requeueExpired', '租约巡检完成', {}, {
  1736. expired_tasks: (rows || []).length,
  1737. requeued: count,
  1738. stale_reassigned: staleCount,
  1739. clients_marked_offline: offlineCount
  1740. })
  1741. }
  1742. return count + staleCount
  1743. }
  1744. }
  1745. module.exports = {
  1746. TaskScheduler,
  1747. TASK_STATUS
  1748. }