|
|
@@ -6,6 +6,7 @@ const Redis = require('../../plugin/DataBase/Redis')
|
|
|
const Logger = require('../Logger')
|
|
|
|
|
|
const TASK_STATUS = {
|
|
|
+ PAUSED: 'paused',
|
|
|
PENDING: 'pending',
|
|
|
ASSIGNED: 'assigned',
|
|
|
RUNNING: 'running',
|
|
|
@@ -171,6 +172,38 @@ class TaskScheduler {
|
|
|
return []
|
|
|
}
|
|
|
|
|
|
+ countTaskTargets(courses = [], courseGroups = []) {
|
|
|
+ return this.normalizeArray(courses).length + this.normalizeArray(courseGroups).length
|
|
|
+ }
|
|
|
+
|
|
|
+ getMinTaskInterval(totalCount) {
|
|
|
+ return totalCount >= 5 ? 500 : 200
|
|
|
+ }
|
|
|
+
|
|
|
+ validateTaskCoursesAndInterval(courses, courseGroups, intervalMs) {
|
|
|
+ const normalizedCourses = this.normalizeArray(courses)
|
|
|
+ const normalizedGroups = this.normalizeArray(courseGroups)
|
|
|
+ const totalCount = normalizedCourses.length + normalizedGroups.length
|
|
|
+ if (totalCount <= 0) {
|
|
|
+ throw new Error('至少需要填写一门课程或一个课程分组')
|
|
|
+ }
|
|
|
+ if (totalCount > 10) {
|
|
|
+ throw new Error('每个任务最多选择10门课程或分组')
|
|
|
+ }
|
|
|
+ const minInterval = this.getMinTaskInterval(totalCount)
|
|
|
+ const safeInterval = Number(intervalMs)
|
|
|
+ if (!Number.isFinite(safeInterval) || safeInterval < minInterval || safeInterval > 10000) {
|
|
|
+ throw new Error(`当前共 ${totalCount} 个抢课目标,间隔需在 ${minInterval}-10000ms 之间`)
|
|
|
+ }
|
|
|
+ return {
|
|
|
+ courses: normalizedCourses,
|
|
|
+ courseGroups: normalizedGroups,
|
|
|
+ intervalMs: safeInterval,
|
|
|
+ totalCount,
|
|
|
+ minInterval
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
serializeTask(row, includeSecret = false) {
|
|
|
const result = { ...row }
|
|
|
result.enable_ggxxk = Number(result.enable_ggxxk) === 1
|
|
|
@@ -181,6 +214,12 @@ class TaskScheduler {
|
|
|
result.result_json = JSON.parse(result.result_json)
|
|
|
} catch (_) {}
|
|
|
}
|
|
|
+ result.jx0502zbid = row.batch_jx0502zbid || row.jx0502zbid || ''
|
|
|
+ result.batch_name = row.batch_name || ''
|
|
|
+ if (row.batch_enabled !== undefined) {
|
|
|
+ result.batch_enabled = Number(row.batch_enabled) === 1
|
|
|
+ }
|
|
|
+ delete result.batch_jx0502zbid
|
|
|
if (includeSecret) {
|
|
|
result.password = this.decryptPassword(result.password_enc)
|
|
|
}
|
|
|
@@ -188,6 +227,168 @@ class TaskScheduler {
|
|
|
return result
|
|
|
}
|
|
|
|
|
|
+ taskBatchJoinSql(alias = 't') {
|
|
|
+ return `LEFT JOIN qk_batch b ON b.id = ${alias}.batch_id`
|
|
|
+ }
|
|
|
+
|
|
|
+ taskBatchSelectSql(alias = 't') {
|
|
|
+ return `, b.name AS batch_name, b.jx0502zbid AS batch_jx0502zbid, b.enabled AS batch_enabled`
|
|
|
+ }
|
|
|
+
|
|
|
+ async getBatchById(batchId) {
|
|
|
+ const rows = await db.query('SELECT * FROM qk_batch WHERE id = ?', [batchId])
|
|
|
+ if (!rows || rows.length === 0) {
|
|
|
+ throw new Error('抢课批次不存在')
|
|
|
+ }
|
|
|
+ return rows[0]
|
|
|
+ }
|
|
|
+
|
|
|
+ async requireEnabledBatch(batchId) {
|
|
|
+ const batch = await this.getBatchById(batchId)
|
|
|
+ if (!Number(batch.enabled)) {
|
|
|
+ throw new Error('所选抢课批次未启用')
|
|
|
+ }
|
|
|
+ return batch
|
|
|
+ }
|
|
|
+
|
|
|
+ resolveTaskBatchId(payload = {}) {
|
|
|
+ const batchId = Number(payload.batch_id || payload.batchId)
|
|
|
+ if (!batchId) {
|
|
|
+ throw new Error('请选择抢课批次')
|
|
|
+ }
|
|
|
+ return batchId
|
|
|
+ }
|
|
|
+
|
|
|
+ async releaseAssignedTask(task, now = Date.now()) {
|
|
|
+ if (!task || ![TASK_STATUS.ASSIGNED, TASK_STATUS.RUNNING].includes(task.status)) {
|
|
|
+ return false
|
|
|
+ }
|
|
|
+ if (task.assigned_client_id) {
|
|
|
+ await this.decrementClientSlots(task.assigned_client_id, now)
|
|
|
+ }
|
|
|
+ return true
|
|
|
+ }
|
|
|
+
|
|
|
+ async pauseTasksByBatch(batchId, reason = '批次已停用,任务已收回') {
|
|
|
+ const rows = await db.query(
|
|
|
+ 'SELECT id, status, assigned_client_id FROM qk_task WHERE batch_id = ? AND status IN (?, ?, ?)',
|
|
|
+ [batchId, TASK_STATUS.PENDING, TASK_STATUS.ASSIGNED, TASK_STATUS.RUNNING]
|
|
|
+ )
|
|
|
+ const now = Date.now()
|
|
|
+ let count = 0
|
|
|
+ for (const row of rows || []) {
|
|
|
+ await this.releaseAssignedTask(row, now)
|
|
|
+ const result = await db.query(
|
|
|
+ `UPDATE qk_task SET status = ?, assigned_client_id = NULL, lease_expire_at = NULL,
|
|
|
+ assigned_at = NULL, exclude_client_id = NULL, update_time = ? WHERE id = ? AND status IN (?, ?, ?)`,
|
|
|
+ [TASK_STATUS.PAUSED, now, row.id, TASK_STATUS.PENDING, TASK_STATUS.ASSIGNED, TASK_STATUS.RUNNING]
|
|
|
+ )
|
|
|
+ if (result && result.affectedRows > 0) {
|
|
|
+ count += 1
|
|
|
+ await this.logTask(row.id, row.assigned_client_id, 'batch_paused', reason)
|
|
|
+ }
|
|
|
+ }
|
|
|
+ if (count > 0) {
|
|
|
+ this.logInfo('pauseTasksByBatch', reason, { batchId }, { affected: count })
|
|
|
+ }
|
|
|
+ return count
|
|
|
+ }
|
|
|
+
|
|
|
+ async redispatchTasksByBatch(batchId, reason = '批次 ID 已变更,任务已重新进入分配队列') {
|
|
|
+ const rows = await db.query(
|
|
|
+ 'SELECT id, status, assigned_client_id FROM qk_task WHERE batch_id = ? AND status IN (?, ?, ?)',
|
|
|
+ [batchId, TASK_STATUS.PENDING, TASK_STATUS.ASSIGNED, TASK_STATUS.RUNNING]
|
|
|
+ )
|
|
|
+ const now = Date.now()
|
|
|
+ let count = 0
|
|
|
+ for (const row of rows || []) {
|
|
|
+ if ([TASK_STATUS.ASSIGNED, TASK_STATUS.RUNNING].includes(row.status)) {
|
|
|
+ await this.releaseAssignedTask(row, now)
|
|
|
+ const result = await db.query(
|
|
|
+ `UPDATE qk_task SET status = ?, assigned_client_id = NULL, lease_expire_at = NULL,
|
|
|
+ assigned_at = NULL, exclude_client_id = NULL, update_time = ? WHERE id = ? AND status IN (?, ?)`,
|
|
|
+ [TASK_STATUS.PENDING, now, row.id, TASK_STATUS.ASSIGNED, TASK_STATUS.RUNNING]
|
|
|
+ )
|
|
|
+ if (result && result.affectedRows > 0) {
|
|
|
+ count += 1
|
|
|
+ await this.logTask(row.id, row.assigned_client_id, 'batch_redispatch', reason)
|
|
|
+ }
|
|
|
+ }
|
|
|
+ }
|
|
|
+ if (count > 0) {
|
|
|
+ this.logInfo('redispatchTasksByBatch', reason, { batchId }, { affected: count })
|
|
|
+ }
|
|
|
+ return count
|
|
|
+ }
|
|
|
+
|
|
|
+ serializeBatch(row) {
|
|
|
+ return {
|
|
|
+ id: row.id,
|
|
|
+ name: row.name,
|
|
|
+ jx0502zbid: row.jx0502zbid,
|
|
|
+ enabled: Number(row.enabled) === 1,
|
|
|
+ create_time: row.create_time,
|
|
|
+ update_time: row.update_time
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
+ async listBatches(options = {}) {
|
|
|
+ const enabledOnly = !!options.enabledOnly
|
|
|
+ const where = enabledOnly ? 'WHERE enabled = 1' : ''
|
|
|
+ const rows = await db.query(`SELECT * FROM qk_batch ${where} ORDER BY update_time DESC, id DESC`)
|
|
|
+ return (rows || []).map(row => this.serializeBatch(row))
|
|
|
+ }
|
|
|
+
|
|
|
+ async createBatch(payload = {}) {
|
|
|
+ const name = String(payload.name || '').trim()
|
|
|
+ const jx0502zbid = String(payload.jx0502zbid || '').trim()
|
|
|
+ if (!name || !jx0502zbid) {
|
|
|
+ throw new Error('批次名称和批次 ID 不能为空')
|
|
|
+ }
|
|
|
+ const now = Date.now()
|
|
|
+ const result = await db.query(
|
|
|
+ 'INSERT INTO qk_batch (name, jx0502zbid, enabled, create_time, update_time) VALUES (?, ?, ?, ?, ?)',
|
|
|
+ [name, jx0502zbid, payload.enabled === false ? 0 : 1, now, now]
|
|
|
+ )
|
|
|
+ if (!result || result.affectedRows <= 0) {
|
|
|
+ throw new Error('创建抢课批次失败')
|
|
|
+ }
|
|
|
+ this.logInfo('createBatch', '抢课批次已创建', { batchId: result.insertId }, { name, jx0502zbid })
|
|
|
+ return result.insertId
|
|
|
+ }
|
|
|
+
|
|
|
+ async updateBatch(batchId, payload = {}) {
|
|
|
+ const existing = await this.getBatchById(batchId)
|
|
|
+ const name = payload.name !== undefined ? String(payload.name || '').trim() : existing.name
|
|
|
+ const jx0502zbid = payload.jx0502zbid !== undefined ? String(payload.jx0502zbid || '').trim() : existing.jx0502zbid
|
|
|
+ const enabled = payload.enabled !== undefined ? (payload.enabled ? 1 : 0) : Number(existing.enabled)
|
|
|
+ if (!name || !jx0502zbid) {
|
|
|
+ throw new Error('批次名称和批次 ID 不能为空')
|
|
|
+ }
|
|
|
+ const now = Date.now()
|
|
|
+ const result = await db.query(
|
|
|
+ 'UPDATE qk_batch SET name = ?, jx0502zbid = ?, enabled = ?, update_time = ? WHERE id = ?',
|
|
|
+ [name, jx0502zbid, enabled, now, batchId]
|
|
|
+ )
|
|
|
+ if (!result || result.affectedRows <= 0) {
|
|
|
+ throw new Error('更新抢课批次失败')
|
|
|
+ }
|
|
|
+ const disabledNow = Number(existing.enabled) === 1 && enabled === 0
|
|
|
+ const idChanged = String(existing.jx0502zbid) !== jx0502zbid
|
|
|
+ if (disabledNow) {
|
|
|
+ await this.pauseTasksByBatch(batchId)
|
|
|
+ } else if (idChanged && enabled === 1) {
|
|
|
+ await this.redispatchTasksByBatch(batchId)
|
|
|
+ }
|
|
|
+ this.logInfo('updateBatch', '抢课批次已更新', { batchId }, {
|
|
|
+ name,
|
|
|
+ jx0502zbid,
|
|
|
+ enabled: enabled === 1,
|
|
|
+ disabledNow,
|
|
|
+ idChanged
|
|
|
+ })
|
|
|
+ }
|
|
|
+
|
|
|
async logTask(taskId, clientId, event, message = '', payload = null) {
|
|
|
const sql = 'INSERT INTO qk_task_log (task_id, client_id, event, message, payload_json, create_time) VALUES (?, ?, ?, ?, ?, ?)'
|
|
|
await db.query(sql, [
|
|
|
@@ -286,41 +487,41 @@ class TaskScheduler {
|
|
|
}
|
|
|
|
|
|
async createTask(uuid, payload) {
|
|
|
- const courses = this.normalizeArray(payload.courses || payload.COURSES)
|
|
|
- const courseGroups = this.normalizeArray(payload.course_groups || payload.COURSE_GROUPS)
|
|
|
- const intervalMs = Number(payload.interval_ms || payload.INTERVAL_MS || 500)
|
|
|
- if (!courses.length && !courseGroups.length) {
|
|
|
- throw new Error('至少需要填写一门课程或一个课程分组')
|
|
|
- }
|
|
|
- if (!Number.isFinite(intervalMs) || intervalMs < 200 || intervalMs > 10000) {
|
|
|
- throw new Error('抢课间隔需在 200-10000ms 之间')
|
|
|
- }
|
|
|
+ const validated = this.validateTaskCoursesAndInterval(
|
|
|
+ payload.courses || payload.COURSES,
|
|
|
+ payload.course_groups || payload.COURSE_GROUPS,
|
|
|
+ payload.interval_ms || payload.INTERVAL_MS || 500
|
|
|
+ )
|
|
|
+ const { courses, courseGroups, intervalMs } = validated
|
|
|
+ const batchId = this.resolveTaskBatchId(payload)
|
|
|
+ await this.requireEnabledBatch(batchId)
|
|
|
const time = Date.now()
|
|
|
const sql = `INSERT INTO qk_task
|
|
|
- (create_user, name, jx0502zbid, student_num, password_enc, courses, course_groups, enable_ggxxk, interval_ms, status, create_time, update_time)
|
|
|
- VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`
|
|
|
+ (create_user, name, batch_id, jx0502zbid, student_num, password_enc, courses, course_groups, enable_ggxxk, interval_ms, status, create_time, update_time)
|
|
|
+ VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`
|
|
|
const result = await db.query(sql, [
|
|
|
uuid,
|
|
|
payload.name,
|
|
|
- payload.jx0502zbid || payload.id,
|
|
|
+ batchId,
|
|
|
+ '',
|
|
|
payload.student_num || payload.user,
|
|
|
this.encryptPassword(payload.password || payload.pass),
|
|
|
JSON.stringify(courses),
|
|
|
JSON.stringify(courseGroups),
|
|
|
payload.enable_ggxxk || payload.ENABLE_GGXXK ? 1 : 0,
|
|
|
intervalMs,
|
|
|
- TASK_STATUS.PENDING,
|
|
|
+ TASK_STATUS.PAUSED,
|
|
|
time,
|
|
|
time
|
|
|
])
|
|
|
if (!result || result.affectedRows <= 0) {
|
|
|
throw new Error('创建抢课任务失败')
|
|
|
}
|
|
|
- await this.logTask(result.insertId, null, 'created', '用户提交抢课任务')
|
|
|
+ await this.logTask(result.insertId, null, 'created', '用户提交抢课任务(未开始)')
|
|
|
this.logInfo('createTask', '抢课任务已创建', { taskId: result.insertId, uuid }, {
|
|
|
name: payload.name,
|
|
|
+ batch_id: batchId,
|
|
|
student_num: payload.student_num || payload.user,
|
|
|
- jx0502zbid: payload.jx0502zbid || payload.id,
|
|
|
courses_count: courses.length,
|
|
|
course_groups_count: courseGroups.length,
|
|
|
interval_ms: intervalMs,
|
|
|
@@ -334,38 +535,40 @@ class TaskScheduler {
|
|
|
if (!rows || rows.length === 0) {
|
|
|
throw new Error('任务不存在')
|
|
|
}
|
|
|
- if (![TASK_STATUS.PENDING, TASK_STATUS.FAILED, TASK_STATUS.CANCELLED].includes(rows[0].status)) {
|
|
|
- throw new Error('任务已被客户端领取,暂不能修改')
|
|
|
- }
|
|
|
- const courses = this.normalizeArray(payload.courses || payload.COURSES)
|
|
|
- const courseGroups = this.normalizeArray(payload.course_groups || payload.COURSE_GROUPS)
|
|
|
- const intervalMs = Number(payload.interval_ms || payload.INTERVAL_MS || 500)
|
|
|
- if (!courses.length && !courseGroups.length) {
|
|
|
- throw new Error('至少需要填写一门课程或一个课程分组')
|
|
|
+ if (![TASK_STATUS.PAUSED, TASK_STATUS.FAILED, TASK_STATUS.CANCELLED].includes(rows[0].status)) {
|
|
|
+ throw new Error('任务进行中或已分配,请先暂停后再修改')
|
|
|
}
|
|
|
+ const validated = this.validateTaskCoursesAndInterval(
|
|
|
+ payload.courses || payload.COURSES,
|
|
|
+ payload.course_groups || payload.COURSE_GROUPS,
|
|
|
+ payload.interval_ms || payload.INTERVAL_MS || 500
|
|
|
+ )
|
|
|
+ const { courses, courseGroups, intervalMs } = validated
|
|
|
+ const batchId = this.resolveTaskBatchId(payload)
|
|
|
+ await this.requireEnabledBatch(batchId)
|
|
|
const passwordSql = payload.password || payload.pass ? ', password_enc = ?' : ''
|
|
|
const params = [
|
|
|
payload.name,
|
|
|
- payload.jx0502zbid || payload.id,
|
|
|
+ batchId,
|
|
|
payload.student_num || payload.user,
|
|
|
JSON.stringify(courses),
|
|
|
JSON.stringify(courseGroups),
|
|
|
payload.enable_ggxxk || payload.ENABLE_GGXXK ? 1 : 0,
|
|
|
intervalMs,
|
|
|
- TASK_STATUS.PENDING,
|
|
|
+ TASK_STATUS.PAUSED,
|
|
|
Date.now()
|
|
|
]
|
|
|
if (passwordSql) {
|
|
|
params.splice(3, 0, this.encryptPassword(payload.password || payload.pass))
|
|
|
}
|
|
|
params.push(taskId, uuid)
|
|
|
- const sql = `UPDATE qk_task SET name = ?, jx0502zbid = ?, student_num = ?${passwordSql}, courses = ?, course_groups = ?, enable_ggxxk = ?, interval_ms = ?, status = ?, update_time = ?, assigned_client_id = NULL, lease_expire_at = NULL, result_json = NULL, error_msg = NULL, finished_time = NULL WHERE id = ? AND create_user = ?`
|
|
|
+ const sql = `UPDATE qk_task SET name = ?, batch_id = ?, student_num = ?${passwordSql}, courses = ?, course_groups = ?, enable_ggxxk = ?, 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 = ?`
|
|
|
const result = await db.query(sql, params)
|
|
|
if (!result || result.affectedRows <= 0) {
|
|
|
throw new Error('更新抢课任务失败')
|
|
|
}
|
|
|
await this.logTask(taskId, null, 'updated', '用户更新抢课任务')
|
|
|
- this.logInfo('updateTask', '抢课任务已更新并重新进入待分配队列', { taskId, uuid }, {
|
|
|
+ this.logInfo('updateTask', '抢课任务已更新(未开始)', { taskId, uuid }, {
|
|
|
name: payload.name,
|
|
|
student_num: payload.student_num || payload.user,
|
|
|
courses_count: courses.length,
|
|
|
@@ -376,12 +579,25 @@ class TaskScheduler {
|
|
|
}
|
|
|
|
|
|
async listUserTasks(uuid) {
|
|
|
- const rows = await db.query('SELECT * FROM qk_task WHERE create_user = ? ORDER BY create_time DESC', [uuid])
|
|
|
+ const rows = await db.query(
|
|
|
+ `SELECT t.*${this.taskBatchSelectSql('t')}
|
|
|
+ FROM qk_task t
|
|
|
+ ${this.taskBatchJoinSql('t')}
|
|
|
+ WHERE t.create_user = ?
|
|
|
+ ORDER BY t.create_time DESC`,
|
|
|
+ [uuid]
|
|
|
+ )
|
|
|
return (rows || []).map(row => this.serializeTask(row))
|
|
|
}
|
|
|
|
|
|
async getTaskDetail(uuid, taskId) {
|
|
|
- const rows = await db.query('SELECT * FROM qk_task WHERE id = ? AND create_user = ?', [taskId, uuid])
|
|
|
+ const rows = await db.query(
|
|
|
+ `SELECT t.*${this.taskBatchSelectSql('t')}
|
|
|
+ FROM qk_task t
|
|
|
+ ${this.taskBatchJoinSql('t')}
|
|
|
+ WHERE t.id = ? AND t.create_user = ?`,
|
|
|
+ [taskId, uuid]
|
|
|
+ )
|
|
|
if (!rows || rows.length === 0) {
|
|
|
return null
|
|
|
}
|
|
|
@@ -401,8 +617,8 @@ class TaskScheduler {
|
|
|
|
|
|
async cancelTask(uuid, taskId) {
|
|
|
const result = await db.query(
|
|
|
- 'UPDATE qk_task SET status = ?, update_time = ?, finished_time = ? WHERE id = ? AND create_user = ? AND status IN (?, ?)',
|
|
|
- [TASK_STATUS.CANCELLED, Date.now(), Date.now(), taskId, uuid, TASK_STATUS.PENDING, TASK_STATUS.FAILED]
|
|
|
+ 'UPDATE qk_task SET status = ?, update_time = ?, finished_time = ? WHERE id = ? AND create_user = ? AND status IN (?, ?, ?)',
|
|
|
+ [TASK_STATUS.CANCELLED, Date.now(), Date.now(), taskId, uuid, TASK_STATUS.PAUSED, TASK_STATUS.PENDING, TASK_STATUS.FAILED]
|
|
|
)
|
|
|
if (!result || result.affectedRows <= 0) {
|
|
|
throw new Error('任务不存在或当前状态不可取消')
|
|
|
@@ -411,11 +627,67 @@ class TaskScheduler {
|
|
|
this.logInfo('cancelTask', '用户已取消抢课任务', { taskId, uuid })
|
|
|
}
|
|
|
|
|
|
+ async startTask(uuid, taskId) {
|
|
|
+ const rows = await db.query('SELECT * FROM qk_task WHERE id = ? AND create_user = ?', [taskId, uuid])
|
|
|
+ if (!rows || rows.length === 0) {
|
|
|
+ throw new Error('任务不存在')
|
|
|
+ }
|
|
|
+ const task = rows[0]
|
|
|
+ if (![TASK_STATUS.PAUSED, TASK_STATUS.FAILED].includes(task.status)) {
|
|
|
+ throw new Error('仅未开始或失败的任务可开启')
|
|
|
+ }
|
|
|
+ if (task.batch_id) {
|
|
|
+ await this.requireEnabledBatch(task.batch_id)
|
|
|
+ }
|
|
|
+ const now = Date.now()
|
|
|
+ const result = await db.query(
|
|
|
+ `UPDATE qk_task SET status = ?, 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 = ? AND create_user = ? AND status IN (?, ?)`,
|
|
|
+ [TASK_STATUS.PENDING, now, taskId, uuid, TASK_STATUS.PAUSED, TASK_STATUS.FAILED]
|
|
|
+ )
|
|
|
+ if (!result || result.affectedRows <= 0) {
|
|
|
+ throw new Error('开启抢课任务失败')
|
|
|
+ }
|
|
|
+ await this.logTask(taskId, null, 'started', '用户开启抢课任务,等待分配')
|
|
|
+ this.logInfo('startTask', '用户已开启抢课任务', { taskId, uuid })
|
|
|
+ }
|
|
|
+
|
|
|
+ async pauseTask(uuid, taskId) {
|
|
|
+ const rows = await db.query('SELECT * FROM qk_task WHERE id = ? AND create_user = ?', [taskId, uuid])
|
|
|
+ if (!rows || rows.length === 0) {
|
|
|
+ throw new Error('任务不存在')
|
|
|
+ }
|
|
|
+ const task = rows[0]
|
|
|
+ if (![TASK_STATUS.PENDING, TASK_STATUS.ASSIGNED, TASK_STATUS.RUNNING].includes(task.status)) {
|
|
|
+ throw new Error('当前状态不可暂停')
|
|
|
+ }
|
|
|
+ const now = Date.now()
|
|
|
+ const wasAssigned = [TASK_STATUS.ASSIGNED, TASK_STATUS.RUNNING].includes(task.status)
|
|
|
+ if (wasAssigned && task.assigned_client_id) {
|
|
|
+ await this.decrementClientSlots(task.assigned_client_id, now)
|
|
|
+ }
|
|
|
+ const result = await db.query(
|
|
|
+ `UPDATE qk_task SET status = ?, update_time = ?, assigned_client_id = NULL, lease_expire_at = NULL,
|
|
|
+ assigned_at = NULL, exclude_client_id = NULL
|
|
|
+ WHERE id = ? AND create_user = ? AND status IN (?, ?, ?)`,
|
|
|
+ [TASK_STATUS.PAUSED, now, taskId, uuid, TASK_STATUS.PENDING, TASK_STATUS.ASSIGNED, TASK_STATUS.RUNNING]
|
|
|
+ )
|
|
|
+ if (!result || result.affectedRows <= 0) {
|
|
|
+ throw new Error('暂停抢课任务失败')
|
|
|
+ }
|
|
|
+ await this.logTask(taskId, task.assigned_client_id, 'paused', wasAssigned ? '用户暂停任务,已收回客户端' : '用户暂停抢课任务')
|
|
|
+ this.logInfo('pauseTask', '用户已暂停抢课任务', { taskId, uuid }, {
|
|
|
+ released_from_client: wasAssigned ? task.assigned_client_id : null
|
|
|
+ })
|
|
|
+ }
|
|
|
+
|
|
|
async getAdminTask(taskId) {
|
|
|
const rows = await db.query(
|
|
|
- `SELECT t.*, u.username, u.avatar
|
|
|
+ `SELECT t.*, u.username, u.avatar${this.taskBatchSelectSql('t')}
|
|
|
FROM qk_task t
|
|
|
LEFT JOIN users u ON u.uuid COLLATE utf8mb4_general_ci = t.create_user COLLATE utf8mb4_general_ci
|
|
|
+ ${this.taskBatchJoinSql('t')}
|
|
|
WHERE t.id = ?`,
|
|
|
[taskId]
|
|
|
)
|
|
|
@@ -442,42 +714,40 @@ class TaskScheduler {
|
|
|
if (task.status === TASK_STATUS.SUCCESS) {
|
|
|
throw new Error('已成功的任务不可编辑')
|
|
|
}
|
|
|
- const courses = this.normalizeArray(payload.courses || payload.COURSES)
|
|
|
- const courseGroups = this.normalizeArray(payload.course_groups || payload.COURSE_GROUPS)
|
|
|
- const intervalMs = Number(payload.interval_ms || payload.INTERVAL_MS || task.interval_ms || 500)
|
|
|
- if (!courses.length && !courseGroups.length) {
|
|
|
- throw new Error('至少需要填写一门课程或一个课程分组')
|
|
|
- }
|
|
|
- if (!Number.isFinite(intervalMs) || intervalMs < 200 || intervalMs > 10000) {
|
|
|
- throw new Error('抢课间隔需在 200-10000ms 之间')
|
|
|
- }
|
|
|
+ const validated = this.validateTaskCoursesAndInterval(
|
|
|
+ payload.courses || payload.COURSES,
|
|
|
+ payload.course_groups || payload.COURSE_GROUPS,
|
|
|
+ payload.interval_ms || payload.INTERVAL_MS || task.interval_ms || 500
|
|
|
+ )
|
|
|
+ const { courses, courseGroups, intervalMs } = validated
|
|
|
const now = Date.now()
|
|
|
const wasAssigned = [TASK_STATUS.ASSIGNED, TASK_STATUS.RUNNING].includes(task.status)
|
|
|
if (wasAssigned) {
|
|
|
await this.decrementClientSlots(task.assigned_client_id, now)
|
|
|
}
|
|
|
+ const batchId = this.resolveTaskBatchId(payload)
|
|
|
const passwordSql = payload.password || payload.pass ? ', password_enc = ?' : ''
|
|
|
const params = [
|
|
|
payload.name,
|
|
|
- payload.jx0502zbid || payload.id,
|
|
|
+ batchId,
|
|
|
payload.student_num || payload.user,
|
|
|
JSON.stringify(courses),
|
|
|
JSON.stringify(courseGroups),
|
|
|
payload.enable_ggxxk || payload.ENABLE_GGXXK ? 1 : 0,
|
|
|
intervalMs,
|
|
|
- TASK_STATUS.PENDING,
|
|
|
+ TASK_STATUS.PAUSED,
|
|
|
now
|
|
|
]
|
|
|
if (passwordSql) {
|
|
|
params.splice(3, 0, this.encryptPassword(payload.password || payload.pass))
|
|
|
}
|
|
|
params.push(taskId)
|
|
|
- const sql = `UPDATE qk_task SET name = ?, jx0502zbid = ?, student_num = ?${passwordSql}, courses = ?, course_groups = ?, enable_ggxxk = ?, interval_ms = ?, status = ?, 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 = ?`
|
|
|
+ const sql = `UPDATE qk_task SET name = ?, batch_id = ?, student_num = ?${passwordSql}, courses = ?, course_groups = ?, enable_ggxxk = ?, 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 = ?`
|
|
|
const result = await db.query(sql, params)
|
|
|
if (!result || result.affectedRows <= 0) {
|
|
|
throw new Error('更新抢课任务失败')
|
|
|
}
|
|
|
- await this.logTask(taskId, task.assigned_client_id, 'admin_updated', wasAssigned ? '管理员更新任务并收回重新排队' : '管理员更新抢课任务')
|
|
|
+ await this.logTask(taskId, task.assigned_client_id, 'admin_updated', wasAssigned ? '管理员更新任务并收回(未开始)' : '管理员更新抢课任务')
|
|
|
this.logInfo('adminUpdateTask', '管理员已更新抢课任务', { taskId }, {
|
|
|
name: payload.name,
|
|
|
student_num: payload.student_num || payload.user,
|
|
|
@@ -491,7 +761,7 @@ class TaskScheduler {
|
|
|
throw new Error('任务不存在')
|
|
|
}
|
|
|
const task = rows[0]
|
|
|
- if (![TASK_STATUS.PENDING, TASK_STATUS.ASSIGNED, TASK_STATUS.RUNNING, TASK_STATUS.FAILED].includes(task.status)) {
|
|
|
+ if (![TASK_STATUS.PAUSED, TASK_STATUS.PENDING, TASK_STATUS.ASSIGNED, TASK_STATUS.RUNNING, TASK_STATUS.FAILED].includes(task.status)) {
|
|
|
throw new Error('当前状态不可取消')
|
|
|
}
|
|
|
const now = Date.now()
|
|
|
@@ -500,8 +770,8 @@ class TaskScheduler {
|
|
|
}
|
|
|
const result = await db.query(
|
|
|
`UPDATE qk_task SET status = ?, assigned_client_id = NULL, lease_expire_at = NULL, assigned_at = NULL,
|
|
|
- exclude_client_id = NULL, update_time = ?, finished_time = ? WHERE id = ? AND status IN (?, ?, ?, ?)`,
|
|
|
- [TASK_STATUS.CANCELLED, now, now, taskId, TASK_STATUS.PENDING, TASK_STATUS.ASSIGNED, TASK_STATUS.RUNNING, TASK_STATUS.FAILED]
|
|
|
+ exclude_client_id = NULL, update_time = ?, finished_time = ? WHERE id = ? AND status IN (?, ?, ?, ?, ?)`,
|
|
|
+ [TASK_STATUS.CANCELLED, now, now, taskId, TASK_STATUS.PAUSED, TASK_STATUS.PENDING, TASK_STATUS.ASSIGNED, TASK_STATUS.RUNNING, TASK_STATUS.FAILED]
|
|
|
)
|
|
|
if (!result || result.affectedRows <= 0) {
|
|
|
throw new Error('取消抢课任务失败')
|
|
|
@@ -701,7 +971,11 @@ class TaskScheduler {
|
|
|
const now = Date.now()
|
|
|
const leaseExpireAt = now + this.leaseMs
|
|
|
const rows = await db.query(
|
|
|
- 'SELECT * FROM qk_task WHERE assigned_client_id = ? AND status IN (?, ?) ORDER BY create_time ASC',
|
|
|
+ `SELECT t.*, b.jx0502zbid AS batch_jx0502zbid, b.name AS batch_name, b.enabled AS batch_enabled
|
|
|
+ FROM qk_task t
|
|
|
+ LEFT JOIN qk_batch b ON b.id = t.batch_id
|
|
|
+ WHERE t.assigned_client_id = ? AND t.status IN (?, ?)
|
|
|
+ ORDER BY t.create_time ASC`,
|
|
|
[clientId, TASK_STATUS.ASSIGNED, TASK_STATUS.RUNNING]
|
|
|
)
|
|
|
if (!rows || rows.length === 0) {
|
|
|
@@ -794,10 +1068,13 @@ class TaskScheduler {
|
|
|
try {
|
|
|
await conn.beginTransaction()
|
|
|
const [rows] = await conn.execute(
|
|
|
- `SELECT * FROM qk_task
|
|
|
- WHERE status = ?
|
|
|
- AND (exclude_client_id IS NULL OR exclude_client_id <> ?)
|
|
|
- ORDER BY create_time ASC LIMIT ${safeCount} FOR UPDATE`,
|
|
|
+ `SELECT t.*, b.jx0502zbid AS batch_jx0502zbid, b.name AS batch_name, b.enabled AS batch_enabled
|
|
|
+ FROM qk_task t
|
|
|
+ LEFT JOIN qk_batch b ON b.id = t.batch_id
|
|
|
+ WHERE t.status = ?
|
|
|
+ AND (t.batch_id IS NULL OR b.enabled = 1)
|
|
|
+ AND (t.exclude_client_id IS NULL OR t.exclude_client_id <> ?)
|
|
|
+ ORDER BY t.create_time ASC LIMIT ${safeCount} FOR UPDATE`,
|
|
|
[TASK_STATUS.PENDING, clientId]
|
|
|
)
|
|
|
const now = Date.now()
|
|
|
@@ -912,6 +1189,11 @@ class TaskScheduler {
|
|
|
params.push(`%${filters.username}%`)
|
|
|
countParams.push(`%${filters.username}%`)
|
|
|
}
|
|
|
+ if (filters.batch_id) {
|
|
|
+ where.push('t.batch_id = ?')
|
|
|
+ params.push(Number(filters.batch_id))
|
|
|
+ countParams.push(Number(filters.batch_id))
|
|
|
+ }
|
|
|
const whereSql = where.join(' AND ')
|
|
|
const offset = (current - 1) * pagesize
|
|
|
const countRows = await db.query(
|
|
|
@@ -921,9 +1203,10 @@ class TaskScheduler {
|
|
|
countParams
|
|
|
)
|
|
|
const rows = await db.query(
|
|
|
- `SELECT t.*, u.username, u.avatar
|
|
|
+ `SELECT t.*, u.username, u.avatar${this.taskBatchSelectSql('t')}
|
|
|
FROM qk_task t
|
|
|
LEFT JOIN users u ON u.uuid COLLATE utf8mb4_general_ci = t.create_user COLLATE utf8mb4_general_ci
|
|
|
+ ${this.taskBatchJoinSql('t')}
|
|
|
WHERE ${whereSql}
|
|
|
ORDER BY t.create_time DESC
|
|
|
LIMIT ${pagesize} OFFSET ${offset}`,
|