const mq = require('./index') const db = require('../DataBase/db') const { publishRunforgeTask } = require('./runforgeTaskMq') const { popDueMessages, stripScheduleMeta, requeueAt, SCHEDULE_KEY } = require('./lepaoAutoScheduleRedis') const mqNames = require('./jkesMqNames') const { prepareAutoRunTask } = require('../jkes/prepareAutoRunTask') let intervalHandle = null const ACCOUNT_SQL = ` SELECT student_num, auto_day, token, target_count, auto_run_distance_min_km, auto_run_distance_max_km, pace_min_sec_per_km, pace_max_sec_per_km FROM lepao_account WHERE student_num = ? AND auto_run = 1 AND state = 1 LIMIT 1 ` async function refreshDelayedStartRunPayload(msg, logger) { const account = msg?.data?.account if (!account || msg?.type !== 'lepao.startRun') { return msg } const rows = await db.query(ACCOUNT_SQL, [account]) if (!rows?.length) { logger.info?.(`[LepaoSchedule] 账号 ${account} 不可用,跳过延迟任务`) return null } const prepared = await prepareAutoRunTask(rows[0]) if (!prepared.run) { logger.info?.(`[LepaoSchedule] ${account} 延迟任务跳过:${prepared.reason}`) return null } return { ...msg, data: { ...msg.data, targetKm: prepared.task.targetKm, autoDoubleSlot: prepared.task.autoDoubleSlot, paceRandomMinSecPerKm: prepared.task.paceRandomMinSecPerKm, paceRandomMaxSecPerKm: prepared.task.paceRandomMaxSecPerKm } } } /** * 定时将 Redis 中已到期的 JKES 乐跑任务写入主任务队列(jkes_runforge_task_queue,支持 mqPrefix) */ function startLepaoSchedulePublisher(options = {}) { const logger = options.logger || console const intervalMs = options.intervalMs ?? 2000 const batch = options.batch ?? 100 if (intervalHandle) return intervalHandle = setInterval(async () => { try { const now = Date.now() const rawList = await popDueMessages(now, batch) if (!rawList.length) return const channel = await mq.getChannel(mqNames.channelScheduleTick) for (const raw of rawList) { let msg try { msg = JSON.parse(raw) } catch (e) { logger.error?.( `[LepaoSchedule] 调度 JSON 无效已丢弃: ${String(raw).slice(0, 120)}` ) continue } try { const refreshed = await refreshDelayedStartRunPayload(msg, logger) if (!refreshed) continue publishRunforgeTask(channel, stripScheduleMeta(refreshed)) } catch (e) { logger.error?.(`[LepaoSchedule] MQ 投递失败,5s 后重试: ${e.message || e}`) const retryAt = now + 5000 try { await requeueAt(retryAt, msg) } catch (re2) { logger.error?.(`[LepaoSchedule] 写回 Redis 失败: ${re2.message || re2}`) } } } } catch (e) { logger.error?.(`[LepaoSchedule] tick 异常: ${e.message || e}`) } }, intervalMs) logger.info?.( `[LepaoSchedule] 已启动 Redis→MQ(JKES)调度(间隔 ${intervalMs}ms,每批最多 ${batch} 条,key=${SCHEDULE_KEY})` ) } module.exports = { startLepaoSchedulePublisher }