mirror of
https://github.com/whyour/qinglong.git
synced 2026-09-28 09:02:12 +08:00
fix: isolate invalid cron rules during scheduler recovery
This commit is contained in:
@@ -5,6 +5,7 @@ const ts = require('typescript');
|
||||
const load = require('../helpers/load-security-module.cjs');
|
||||
const { SchedulerReadiness } = require('../../back/shared/schedulerReadiness');
|
||||
const { AddCronRequest } = require('../../back/protos/cron');
|
||||
const { getInvalidCronSchedules } = require('../../back/shared/cronSchedule');
|
||||
|
||||
function fixture(client) {
|
||||
const source = fs.readFileSync('back/services/cron.ts', 'utf8');
|
||||
@@ -15,11 +16,12 @@ function fixture(client) {
|
||||
{ compilerOptions: { target: ts.ScriptTarget.ES2020 } },
|
||||
).outputText;
|
||||
const module = { exports: {} };
|
||||
new Function('module', 'isDemoEnv', 'cronClient', 'withSchedulerMutation', js)(
|
||||
new Function('module', 'isDemoEnv', 'cronClient', 'withSchedulerMutation', 'getInvalidCronSchedules', js)(
|
||||
module,
|
||||
() => false,
|
||||
client,
|
||||
(fn) => fn(),
|
||||
getInvalidCronSchedules,
|
||||
);
|
||||
const service = new module.exports();
|
||||
service.setCrontab = async () => {};
|
||||
@@ -101,7 +103,7 @@ test('invalid replacement leaves the previous schedule intact', async () => {
|
||||
{
|
||||
request: {
|
||||
replace: true,
|
||||
crons: [{ id: 'bad', schedule: '?', extra_schedules: [] }],
|
||||
crons: [{ id: 'bad', schedule: 'not a cron', extra_schedules: [] }],
|
||||
},
|
||||
},
|
||||
(err) => (err ? reject(err) : resolve()),
|
||||
@@ -111,6 +113,56 @@ test('invalid replacement leaves the previous schedule intact', async () => {
|
||||
assert.deepEqual([...stacks.keys()], ['old']);
|
||||
});
|
||||
|
||||
test('invalid persisted schedules cannot block recovery, valid jobs or HTTP health', async (t) => {
|
||||
const stacks = new Map();
|
||||
const warnings = [];
|
||||
const { addCron } = load('back/schedule/addCron.ts', {
|
||||
'./data': { scheduleStacks: stacks },
|
||||
'../shared/runCron': {},
|
||||
'../loaders/logger': { info() {}, warn() {} },
|
||||
'../shared/i18n': { tf: require('node:util').format },
|
||||
});
|
||||
const client = { addCron: (crons, replace) => new Promise((resolve, reject) => {
|
||||
addCron({ request: { crons, replace } }, (error) => error ? reject(error) : resolve());
|
||||
}) };
|
||||
const service = fixture(client);
|
||||
let snapshot;
|
||||
service.setCrontab = async (tabs) => { snapshot = tabs.data; };
|
||||
service.logger.warn = (...args) => warnings.push(require('node:util').format(...args));
|
||||
let rows = [
|
||||
{ id: 'good', name: 'valid', schedule: '0 0/30 * * * ?', isDisabled: 0 },
|
||||
{ id: 'bad-main', name: 'bad main', schedule: 'not a cron', isDisabled: 0 },
|
||||
{ id: 'bad-extra', name: 'bad extra', schedule: '* * * * *', extra_schedules: [{ schedule: '0 70 * * * ?' }], isDisabled: 0 },
|
||||
{ id: 'unsupported', name: 'unsupported', schedule: 'H * * * *', isDisabled: 0 },
|
||||
{ id: 'disabled', name: 'disabled', schedule: 'bad', isDisabled: 1 },
|
||||
];
|
||||
service.crontabs = async () => ({ data: rows });
|
||||
const state = new SchedulerReadiness(async () => {}, 60000);
|
||||
t.after(() => {
|
||||
clearTimeout(state.retry);
|
||||
for (const jobs of stacks.values()) for (const job of jobs) job.cancel();
|
||||
});
|
||||
state.configure(() => service.autosave_crontab(true));
|
||||
assert.equal(await state.recover(), true);
|
||||
assert.deepEqual([...stacks.keys()], ['good']);
|
||||
assert.deepEqual(snapshot.map((x) => x.id), ['good', 'disabled']);
|
||||
assert.equal(rows.length, 5, 'persisted rows are preserved for editing');
|
||||
assert.equal(rows[1].isDisabled, 0);
|
||||
assert.equal(warnings.length, 3);
|
||||
assert.match(warnings.join('\n'), /bad-main.*bad main.*not a cron/);
|
||||
assert.match(warnings.join('\n'), /bad-extra.*0 70/);
|
||||
const { HealthService } = load('back/services/health.ts', {
|
||||
typedi: { Service: () => (x) => x }, './http': {},
|
||||
'../schedule/client': { readiness: state }, '../loaders/logger': { error() {} },
|
||||
});
|
||||
const health = new HealthService({ getServer: () => ({}) });
|
||||
assert.equal((await health.check()).status, 'ok');
|
||||
rows = rows.filter((x) => x.id !== 'good');
|
||||
assert.equal(await state.recover(), true, 'even all-invalid snapshots can recover');
|
||||
assert.equal(stacks.size, 0);
|
||||
assert.equal((await health.check()).status, 'ok');
|
||||
});
|
||||
|
||||
test('pre-RPC channel failures invalidate readiness and return 503 without executing or replaying writes', async () => {
|
||||
for (const method of ['addCron', 'delCron']) {
|
||||
let invalidations = 0,
|
||||
|
||||
Reference in New Issue
Block a user