mirror of
https://github.com/whyour/qinglong.git
synced 2026-10-01 05:13:50 +08:00
refactor: migrate cron scheduling to node-cron (#3079)
* refactor: migrate cron scheduling to node-cron * feat: support annually midnight and minutely cron macros * fix: harden cron recovery and system scheduler compatibility
This commit is contained in:
@@ -3,7 +3,7 @@ import { Container } from 'typedi';
|
||||
import { Logger } from 'winston';
|
||||
import SubscriptionService from '../services/subscription';
|
||||
import { celebrate, Joi } from 'celebrate';
|
||||
import CronExpressionParser from 'cron-parser';
|
||||
import { isValidCronSchedule } from '../shared/cronSchedule';
|
||||
const route = Router();
|
||||
|
||||
export default (app: Router) => {
|
||||
@@ -58,10 +58,7 @@ export default (app: Router) => {
|
||||
async (req: Request, res: Response, next: NextFunction) => {
|
||||
const logger: Logger = Container.get('logger');
|
||||
try {
|
||||
if (
|
||||
!req.body.schedule ||
|
||||
CronExpressionParser.parse(req.body.schedule).hasNext()
|
||||
) {
|
||||
if (!req.body.schedule || isValidCronSchedule(req.body.schedule)) {
|
||||
const subscriptionService = Container.get(SubscriptionService);
|
||||
const data = await subscriptionService.create(req.body);
|
||||
return res.send({ code: 200, data });
|
||||
@@ -204,11 +201,7 @@ export default (app: Router) => {
|
||||
async (req: Request, res: Response, next: NextFunction) => {
|
||||
const logger: Logger = Container.get('logger');
|
||||
try {
|
||||
if (
|
||||
!req.body.schedule ||
|
||||
typeof req.body.schedule === 'object' ||
|
||||
CronExpressionParser.parse(req.body.schedule).hasNext()
|
||||
) {
|
||||
if (!req.body.schedule || isValidCronSchedule(req.body.schedule)) {
|
||||
const subscriptionService = Container.get(SubscriptionService);
|
||||
const data = await subscriptionService.update(req.body);
|
||||
return res.send({ code: 200, data });
|
||||
|
||||
+48
-73
@@ -1,6 +1,6 @@
|
||||
import { ServerUnaryCall, sendUnaryData, status } from '@grpc/grpc-js';
|
||||
import { AddCronRequest, AddCronResponse } from '../protos/cron';
|
||||
import nodeSchedule from 'node-schedule';
|
||||
import { createCronJob, CronJob } from '../shared/cronScheduler';
|
||||
import { isValidCronSchedule } from '../shared/cronSchedule';
|
||||
import { scheduleStacks } from './data';
|
||||
import { runCron } from '../shared/runCron';
|
||||
@@ -24,11 +24,7 @@ const addCron = (
|
||||
|
||||
if (!isValidCronField(schedule)) {
|
||||
validationErrors.push(
|
||||
tf(
|
||||
'任务ID %s: 无效的 cron 表达式 "%s"',
|
||||
String(id),
|
||||
schedule,
|
||||
),
|
||||
tf('任务ID %s: 无效的 cron 表达式 "%s"', String(id), schedule),
|
||||
);
|
||||
}
|
||||
|
||||
@@ -56,79 +52,58 @@ const addCron = (
|
||||
return;
|
||||
}
|
||||
|
||||
// Recovery replaces the whole snapshot, including deletions and disabled jobs.
|
||||
// Validation above must finish before touching the previous schedule.
|
||||
if (call.request.replace) {
|
||||
for (const jobs of scheduleStacks.values()) {
|
||||
for (const job of jobs) job?.cancel();
|
||||
// Prepare stopped jobs first: construction errors must preserve the old snapshot.
|
||||
const prepared = new Map<string, CronJob[]>();
|
||||
try {
|
||||
for (const item of call.request.crons) {
|
||||
const jobs: CronJob[] = [];
|
||||
prepared.get(item.id)?.forEach((job) => job.cancel());
|
||||
prepared.set(item.id, jobs);
|
||||
for (const schedule of [
|
||||
item.schedule,
|
||||
...(item.extra_schedules || []).map((x) => x.schedule),
|
||||
]) {
|
||||
jobs.push(
|
||||
createCronJob(
|
||||
schedule,
|
||||
async () => {
|
||||
Logger.info('[schedule][准备运行任务] 命令: %s', item.command);
|
||||
await runCron(item.command, item);
|
||||
},
|
||||
{
|
||||
name: `${item.id}: ${item.name || ''}`,
|
||||
logger: Logger,
|
||||
start: false,
|
||||
},
|
||||
),
|
||||
);
|
||||
}
|
||||
}
|
||||
scheduleStacks.clear();
|
||||
} catch (error) {
|
||||
for (const jobs of prepared.values()) jobs.forEach((job) => job.cancel());
|
||||
const err: any = new Error(
|
||||
error instanceof Error ? error.message : String(error),
|
||||
);
|
||||
err.code = status.INVALID_ARGUMENT;
|
||||
err.details = err.message;
|
||||
callback(err, null);
|
||||
return;
|
||||
}
|
||||
|
||||
// ===== 第二遍:注册所有任务 =====
|
||||
for (const item of call.request.crons) {
|
||||
const { id, schedule, command, extra_schedules, name } = item;
|
||||
|
||||
// 取消该 id 已有的旧任务
|
||||
if (scheduleStacks.has(id)) {
|
||||
scheduleStacks.get(id)?.forEach((x) => x.cancel());
|
||||
}
|
||||
|
||||
if (call.request.replace) {
|
||||
for (const jobs of scheduleStacks.values())
|
||||
jobs.forEach((job) => job?.cancel());
|
||||
scheduleStacks.clear();
|
||||
}
|
||||
for (const [id, jobs] of prepared) {
|
||||
scheduleStacks.get(id)?.forEach((job) => job.cancel());
|
||||
scheduleStacks.set(id, jobs);
|
||||
jobs.forEach((job) => job.start());
|
||||
Logger.info(
|
||||
'[schedule][创建定时任务] 任务ID: %s, 名称: %s, cron: %s, 执行命令: %s',
|
||||
'[schedule][创建定时任务] 任务ID: %s, 规则数: %s',
|
||||
id,
|
||||
name,
|
||||
schedule,
|
||||
command,
|
||||
jobs.length,
|
||||
);
|
||||
|
||||
if (extra_schedules?.length) {
|
||||
extra_schedules.forEach((x) => {
|
||||
Logger.info(
|
||||
'[schedule][创建定时任务] 任务ID: %s, 名称: %s, cron: %s, 执行命令: %s',
|
||||
id,
|
||||
name,
|
||||
x.schedule,
|
||||
command,
|
||||
);
|
||||
});
|
||||
}
|
||||
|
||||
const mainJob = nodeSchedule.scheduleJob(id, schedule, async () => {
|
||||
Logger.info(`[schedule][准备运行任务] 命令: ${command}`);
|
||||
runCron(command, item);
|
||||
});
|
||||
|
||||
if (!mainJob) {
|
||||
Logger.warn(
|
||||
'[schedule][创建定时任务] scheduleJob 返回 null(不符合预期,已通过预校验): 任务ID: %s, cron: %s',
|
||||
id,
|
||||
schedule,
|
||||
);
|
||||
}
|
||||
|
||||
const extraJobs = extra_schedules?.length
|
||||
? extra_schedules.map((x) => {
|
||||
const job = nodeSchedule.scheduleJob(id, x.schedule, async () => {
|
||||
Logger.info(`[schedule][准备运行任务] 命令: ${command}`);
|
||||
runCron(command, item);
|
||||
});
|
||||
if (!job) {
|
||||
Logger.warn(
|
||||
'[schedule][创建定时任务] scheduleJob 返回 null(不符合预期,已通过预校验): 任务ID: %s, cron: %s',
|
||||
id,
|
||||
x.schedule,
|
||||
);
|
||||
}
|
||||
return job;
|
||||
})
|
||||
: [];
|
||||
|
||||
// 过滤 null(兜底保护,正常情况下预校验已拦截)
|
||||
const jobs = [mainJob, ...extraJobs].filter((x) => x != null);
|
||||
if (jobs.length > 0) {
|
||||
scheduleStacks.set(id, jobs);
|
||||
}
|
||||
}
|
||||
|
||||
callback(null, null);
|
||||
|
||||
@@ -67,11 +67,17 @@ class Client {
|
||||
replace = false
|
||||
): Promise<AddCronResponse> {
|
||||
await this.waitForReady(2000);
|
||||
// Include every rule in a recovery snapshot, but bound a stuck worker.
|
||||
const ruleCount = request.reduce(
|
||||
(count, item) => count + 1 + (item.extra_schedules?.length || 0),
|
||||
0,
|
||||
);
|
||||
const registrationTimeoutMs = Math.min(120000, 5000 + ruleCount * 5);
|
||||
return new Promise((resolve, reject) => {
|
||||
this.client.addCron(
|
||||
{ crons: request, replace },
|
||||
new Metadata(),
|
||||
{ deadline: Date.now() + 5000 },
|
||||
{ deadline: Date.now() + registrationTimeoutMs },
|
||||
(err, res) => {
|
||||
if (err) {
|
||||
if (err.code === status.UNAVAILABLE || err.code === status.DEADLINE_EXCEEDED) {
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
import nodeSchedule from 'node-schedule';
|
||||
import { CronJob } from '../shared/cronScheduler';
|
||||
import { ToadScheduler } from 'toad-scheduler';
|
||||
|
||||
export const scheduleStacks = new Map<string, nodeSchedule.Job[]>();
|
||||
export const scheduleStacks = new Map<string, CronJob[]>();
|
||||
|
||||
export const intervalSchedule = new ToadScheduler();
|
||||
|
||||
@@ -13,7 +13,7 @@ const delCron = (
|
||||
'[schedule][取消定时任务] 任务ID: %s',
|
||||
id,
|
||||
);
|
||||
// 过滤掉 nodeSchedule.scheduleJob() 对无效表达式返回的 null,
|
||||
// 防御性过滤历史调度栈中的空任务,
|
||||
// 否则对 null 调 cancel() 会让整个取消流程抛出 UNKNOWN 错误,
|
||||
// 进而导致 HTTP 端的 remove() 跳过 setCrontab(),造成 crontab.list 残留。
|
||||
scheduleStacks.get(id)?.filter((x) => x != null).forEach((x) => {
|
||||
|
||||
+13
-6
@@ -1,5 +1,5 @@
|
||||
import { randomUUID } from 'crypto';
|
||||
import { getInvalidCronSchedules } from '../shared/cronSchedule';
|
||||
import { getInvalidCronSchedules, isValidCronSchedule } from '../shared/cronSchedule';
|
||||
import {
|
||||
withSchedulerMutation,
|
||||
schedulerRegistrationError,
|
||||
@@ -14,7 +14,6 @@ import {
|
||||
} from '../data/runningInstance';
|
||||
import { exec, execSync } from 'child_process';
|
||||
import fs from 'fs/promises';
|
||||
import CronExpressionParser from 'cron-parser';
|
||||
import {
|
||||
getFileContentByName,
|
||||
fileExist,
|
||||
@@ -47,12 +46,20 @@ export default class CronService {
|
||||
|
||||
private isNodeCron(cron: Crontab) {
|
||||
const { schedule, extra_schedules } = cron;
|
||||
const fields = schedule?.trim().split(/\s+/) || [];
|
||||
// System crontab only receives portable numeric five-field expressions.
|
||||
// Extended syntax, macros and legacy shorthand belong to node-schedule.
|
||||
// Extended syntax, macros and legacy shorthand use the Node scheduler.
|
||||
return (
|
||||
schedule?.trim().split(/\s+/).length !== 5 ||
|
||||
fields.length !== 5 ||
|
||||
/[^\d\s*,/\-]/.test(schedule || '') ||
|
||||
Boolean(extra_schedules?.length)
|
||||
Boolean(extra_schedules?.length) ||
|
||||
// BusyBox treats N/step as a single value and doesn't support Sunday=7.
|
||||
fields.some((field) => field.split(',').some((part) => /^\d+\//.test(part))) ||
|
||||
fields[4].includes('7') ||
|
||||
// BusyBox steps DOM from zero; its full-range day/week wildcard
|
||||
// handling also differs from the legacy Node parser's OR semantics.
|
||||
fields[2].includes('/') ||
|
||||
(fields[2] !== '*' && fields[4] !== '*')
|
||||
);
|
||||
}
|
||||
|
||||
@@ -1057,7 +1064,7 @@ export default class CronService {
|
||||
if (
|
||||
command &&
|
||||
schedule &&
|
||||
CronExpressionParser.parse(schedule).hasNext()
|
||||
isValidCronSchedule(schedule)
|
||||
) {
|
||||
const name = namePrefix + '_' + index;
|
||||
|
||||
|
||||
+33
-14
@@ -1,6 +1,6 @@
|
||||
import { Service, Inject } from 'typedi';
|
||||
import winston from 'winston';
|
||||
import nodeSchedule from 'node-schedule';
|
||||
import { createCronJob, CronJob } from '../shared/cronScheduler';
|
||||
import { ChildProcessWithoutNullStreams } from 'child_process';
|
||||
import {
|
||||
ToadScheduler,
|
||||
@@ -11,7 +11,11 @@ import {
|
||||
import dayjs from 'dayjs';
|
||||
import taskLimit from '../shared/pLimit';
|
||||
import { spawn } from 'cross-spawn';
|
||||
import { observeChildProcess, asError, ProcessResult } from '../shared/childProcess';
|
||||
import {
|
||||
observeChildProcess,
|
||||
asError,
|
||||
ProcessResult,
|
||||
} from '../shared/childProcess';
|
||||
|
||||
export interface ScheduleTaskType {
|
||||
id?: number;
|
||||
@@ -38,7 +42,7 @@ export interface TaskCallbacks {
|
||||
|
||||
@Service()
|
||||
export default class ScheduleService {
|
||||
private scheduleStacks = new Map<string, nodeSchedule.Job>();
|
||||
private scheduleStacks = new Map<string, CronJob>();
|
||||
|
||||
private intervalSchedule = new ToadScheduler();
|
||||
|
||||
@@ -159,18 +163,33 @@ export default class ScheduleService {
|
||||
command,
|
||||
);
|
||||
|
||||
this.scheduleStacks.set(
|
||||
_id,
|
||||
nodeSchedule.scheduleJob(_id, schedule, async () => {
|
||||
this.runTask(command, callbacks, {
|
||||
name,
|
||||
if (schedule) {
|
||||
try {
|
||||
const job = createCronJob(
|
||||
schedule,
|
||||
command,
|
||||
id: _id,
|
||||
runOrigin,
|
||||
});
|
||||
}),
|
||||
);
|
||||
() =>
|
||||
this.runTask(command, callbacks, {
|
||||
name,
|
||||
schedule,
|
||||
command,
|
||||
id: _id,
|
||||
runOrigin,
|
||||
}),
|
||||
{ name: `${_id}: ${name || ''}`, logger: this.logger, start: false },
|
||||
);
|
||||
this.scheduleStacks.get(_id)?.cancel();
|
||||
this.scheduleStacks.set(_id, job);
|
||||
job.start();
|
||||
} catch (error) {
|
||||
// A persisted invalid subscription must not become an unhandled startup rejection.
|
||||
this.logger.warn(
|
||||
'[panel][跳过无效定时任务] 任务ID: %s, cron: %s, 错误: %s',
|
||||
_id,
|
||||
schedule,
|
||||
asError(error).message,
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
if (runImmediately) {
|
||||
this.runTask(command, callbacks, {
|
||||
|
||||
@@ -1,12 +1,93 @@
|
||||
import cronParser from 'cron-parser-v4';
|
||||
import CronExpressionParser from 'cron-parser';
|
||||
import { validate } from 'node-cron';
|
||||
import { ScheduleType } from '../interface/schedule';
|
||||
|
||||
// Keep validation aligned with node-schedule's locked cron-parser version.
|
||||
// cron-parser 5 accepts syntax (e.g. H) that the scheduler cannot execute.
|
||||
const aliases: Record<string, string> = {
|
||||
'@yearly': '0 0 0 1 1 *',
|
||||
'@annually': '0 0 0 1 1 *',
|
||||
'@monthly': '0 0 0 1 * *',
|
||||
'@weekly': '0 0 0 * * 0',
|
||||
'@daily': '0 0 0 * * *',
|
||||
'@midnight': '0 0 0 * * *',
|
||||
'@hourly': '0 0 * * * *',
|
||||
'@minutely': '0 * * * * *',
|
||||
};
|
||||
|
||||
export interface CronSchedule {
|
||||
patterns: string[];
|
||||
dayCount: number;
|
||||
nthDay: number;
|
||||
}
|
||||
|
||||
// Validation and registration repeatedly see identical rules in large snapshots.
|
||||
// Cache only immutable syntax, never dates or running tasks, and bound user input.
|
||||
const scheduleCache = new Map<string, CronSchedule>();
|
||||
|
||||
// Canonicalize legacy shorthand, aliases, names and numeric-start steps before
|
||||
// node-cron sees them. Extend legacy aliases explicitly; reject H and bare /N.
|
||||
export function parseCronSchedule(schedule: unknown): CronSchedule {
|
||||
if (typeof schedule !== 'string' || !schedule.trim()) {
|
||||
throw new Error('Invalid cron schedule');
|
||||
}
|
||||
let source = schedule.trim();
|
||||
const cacheKey = source;
|
||||
const cached = scheduleCache.get(cacheKey);
|
||||
if (cached) return cached;
|
||||
if (source.startsWith('@')) {
|
||||
if (!Object.hasOwn(aliases, source)) throw new Error('Invalid cron alias');
|
||||
source = aliases[source];
|
||||
}
|
||||
let parts = source.split(/\s+/);
|
||||
if (
|
||||
parts.length > 6 ||
|
||||
parts.some((part) => /(^|,)\//.test(part) || /H(?:\(|\/|$)/i.test(part))
|
||||
) {
|
||||
throw new Error('Unsupported cron syntax');
|
||||
}
|
||||
parts = [
|
||||
...['0', '*', '*', '*', '*', '*'].slice(0, 6 - parts.length),
|
||||
...parts,
|
||||
];
|
||||
if (
|
||||
parts.some(
|
||||
(part, index) => index !== 3 && index !== 5 && part.includes('?'),
|
||||
)
|
||||
) {
|
||||
throw new Error('Question mark is only valid in day fields');
|
||||
}
|
||||
const expression = CronExpressionParser.parse(parts.join(' '));
|
||||
if (!expression.hasNext())
|
||||
throw new Error('Cron schedule has no next execution');
|
||||
const normalized = expression.stringify(true).replace(/\?/g, '*').split(' ');
|
||||
const dayCount = expression.fields.dayOfMonth.values.length;
|
||||
const weekCount = expression.fields.dayOfWeek.values.length;
|
||||
const patterns =
|
||||
dayCount < 31 && weekCount < 8
|
||||
? [
|
||||
[...normalized.slice(0, 5), '*'].join(' '),
|
||||
[...normalized.slice(0, 3), '*', ...normalized.slice(4)].join(' '),
|
||||
]
|
||||
: [normalized.join(' ')];
|
||||
if (!patterns.every((pattern) => validate(pattern))) {
|
||||
throw new Error('Unsupported cron schedule');
|
||||
}
|
||||
const parsed = {
|
||||
patterns,
|
||||
dayCount,
|
||||
nthDay: expression.fields.dayOfWeek.nthDay,
|
||||
};
|
||||
Object.freeze(patterns);
|
||||
Object.freeze(parsed);
|
||||
if (scheduleCache.size >= 512)
|
||||
scheduleCache.delete(scheduleCache.keys().next().value);
|
||||
scheduleCache.set(cacheKey, parsed);
|
||||
return parsed;
|
||||
}
|
||||
|
||||
export function isValidCronSchedule(schedule: unknown): schedule is string {
|
||||
if (typeof schedule !== 'string' || !schedule.trim()) return false;
|
||||
try {
|
||||
return cronParser.parseExpression(schedule).hasNext();
|
||||
parseCronSchedule(schedule);
|
||||
return true;
|
||||
} catch {
|
||||
return false;
|
||||
}
|
||||
@@ -22,6 +103,8 @@ export function getInvalidCronSchedules(cron: {
|
||||
) {
|
||||
return [];
|
||||
}
|
||||
return [cron.schedule, ...(cron.extra_schedules || []).map((x) => x.schedule)]
|
||||
.filter((schedule) => !isValidCronSchedule(schedule));
|
||||
return [
|
||||
cron.schedule,
|
||||
...(cron.extra_schedules || []).map((x) => x.schedule),
|
||||
].filter((schedule) => !isValidCronSchedule(schedule));
|
||||
}
|
||||
|
||||
@@ -0,0 +1,89 @@
|
||||
import { createTask, ScheduledTask } from 'node-cron';
|
||||
import { parseCronSchedule } from './cronSchedule';
|
||||
|
||||
interface SchedulerLogger {
|
||||
warn(message: string, ...args: unknown[]): unknown;
|
||||
error(message: string, ...args: unknown[]): unknown;
|
||||
}
|
||||
|
||||
export interface CronJob {
|
||||
start(): void;
|
||||
cancel(): void;
|
||||
}
|
||||
|
||||
// node-cron owns timers and calendar calculation; Qinglong owns execution
|
||||
// concurrency and process lifetime. Missed in-process slots join that same queue.
|
||||
export function createCronJob(
|
||||
schedule: string,
|
||||
callback: (date: Date) => unknown | Promise<unknown>,
|
||||
options: { name: string; logger: SchedulerLogger; start?: boolean },
|
||||
): CronJob {
|
||||
const parsed = parseCronSchedule(schedule);
|
||||
const tasks: ScheduledTask[] = [];
|
||||
let cancelled = false;
|
||||
let started = false;
|
||||
const accepts = (index: number, date: Date) => {
|
||||
if (tasks.length === 1) return true;
|
||||
// Preserve cron-parser 4's day/weekday OR semantics, including its
|
||||
// month-dependent full-day range and the global nth-week constraint.
|
||||
if (parsed.nthDay && Math.ceil(date.getDate() / 7) !== parsed.nthDay)
|
||||
return false;
|
||||
if (index === 1) return !tasks[0].match(date);
|
||||
const monthDays = [31, 29, 31, 30, 31, 30, 31, 31, 30, 31, 30, 31];
|
||||
return parsed.dayCount < monthDays[date.getMonth()] || tasks[1].match(date);
|
||||
};
|
||||
const dispatch = async (index: number, date: Date, missed: boolean) => {
|
||||
if (cancelled || !started || !accepts(index, date)) return;
|
||||
if (missed) {
|
||||
options.logger.warn(
|
||||
'[schedule][补执行迟到任务] 任务: %s, 计划时间: %s',
|
||||
options.name,
|
||||
date.toISOString(),
|
||||
);
|
||||
}
|
||||
try {
|
||||
await callback(date);
|
||||
} catch (error) {
|
||||
options.logger.error(
|
||||
'[schedule][定时回调失败] 任务: %s, 错误: %s',
|
||||
options.name,
|
||||
error instanceof Error ? error.message : String(error),
|
||||
);
|
||||
}
|
||||
};
|
||||
try {
|
||||
for (const [index, pattern] of parsed.patterns.entries()) {
|
||||
const task = createTask(
|
||||
pattern,
|
||||
(context) => dispatch(index, context.date, false),
|
||||
{
|
||||
noOverlap: false,
|
||||
missedExecutionTolerance: 5000,
|
||||
},
|
||||
);
|
||||
tasks.push(task);
|
||||
task.on('execution:missed', (context) =>
|
||||
dispatch(index, context.date, true),
|
||||
);
|
||||
// Fail before an existing schedule is cancelled, including impossible dates.
|
||||
task.getNextRuns(1);
|
||||
}
|
||||
} catch (error) {
|
||||
tasks.forEach((task) => task.destroy());
|
||||
throw error;
|
||||
}
|
||||
const job: CronJob = {
|
||||
start() {
|
||||
if (started || cancelled) return;
|
||||
started = true;
|
||||
tasks.forEach((task) => task.start());
|
||||
},
|
||||
cancel() {
|
||||
if (cancelled) return;
|
||||
cancelled = true;
|
||||
tasks.forEach((task) => task.destroy());
|
||||
},
|
||||
};
|
||||
if (options.start !== false) job.start();
|
||||
return job;
|
||||
}
|
||||
Reference in New Issue
Block a user