mirror of
https://github.com/whyour/qinglong.git
synced 2026-09-28 09:02:12 +08:00
fix: verify task termination and preserve startup configuration
This commit is contained in:
+10
-11
@@ -491,17 +491,6 @@ export function psTree(pid: number): Promise<number[]> {
|
||||
|
||||
export async function killTask(pid: number, waitForExit = false) {
|
||||
const descendants = await psTree(pid);
|
||||
if (!waitForExit) {
|
||||
if (descendants.length) {
|
||||
try {
|
||||
[pid, ...descendants]
|
||||
.reverse()
|
||||
.forEach((target) => process.kill(target, 15));
|
||||
} catch {}
|
||||
} else process.kill(pid, 2);
|
||||
return;
|
||||
}
|
||||
const pids = [...descendants.reverse(), pid];
|
||||
const signal = (target: number, sig: NodeJS.Signals) => {
|
||||
try {
|
||||
process.kill(target, sig);
|
||||
@@ -509,6 +498,16 @@ export async function killTask(pid: number, waitForExit = false) {
|
||||
if (error.code !== 'ESRCH') throw error;
|
||||
}
|
||||
};
|
||||
if (!waitForExit) {
|
||||
if (descendants.length) {
|
||||
// A child may exit after psTree; keep signalling the remaining tree.
|
||||
for (const target of [pid, ...descendants].reverse()) {
|
||||
signal(target, 'SIGTERM');
|
||||
}
|
||||
} else signal(pid, 'SIGINT');
|
||||
return;
|
||||
}
|
||||
const pids = [...descendants.reverse(), pid];
|
||||
for (const target of pids) signal(target, 'SIGTERM');
|
||||
const alive = async (target: number) => {
|
||||
try {
|
||||
|
||||
@@ -58,9 +58,9 @@ export default class ScriptService {
|
||||
taskLimit.removeQueuedCron(relativePath.replace(/ /g, '-'));
|
||||
pid = (await getPid(`${TASK_COMMAND} ${relativePath} now`)) as number;
|
||||
}
|
||||
try {
|
||||
await killTask(pid);
|
||||
} catch (error) {}
|
||||
if (pid) {
|
||||
await killTask(pid, true);
|
||||
}
|
||||
|
||||
return { code: 200 };
|
||||
}
|
||||
|
||||
@@ -322,20 +322,22 @@ export default class SubscriptionService {
|
||||
|
||||
public async stop(ids: number[]) {
|
||||
const docs = await SubscriptionModel.findAll({ where: { id: ids } });
|
||||
let failure: unknown;
|
||||
for (const doc of docs) {
|
||||
if (doc.pid) {
|
||||
try {
|
||||
await killTask(doc.pid);
|
||||
if (doc.pid) {
|
||||
await killTask(doc.pid, true);
|
||||
}
|
||||
await SubscriptionModel.update(
|
||||
{ status: SubscriptionStatus.idle, pid: null } as any,
|
||||
{ where: { id: doc.id, pid: doc.pid ?? null } },
|
||||
);
|
||||
} catch (error) {
|
||||
this.logger.error(error);
|
||||
failure ??= error;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
await SubscriptionModel.update(
|
||||
{ status: SubscriptionStatus.idle, pid: undefined },
|
||||
{ where: { id: ids } },
|
||||
);
|
||||
if (failure) throw failure;
|
||||
}
|
||||
|
||||
private async runSingle(subscriptionId: number) {
|
||||
|
||||
@@ -31,7 +31,7 @@ export function runCron(cmd: string, cron: ICron): Promise<number | void> {
|
||||
Logger.info(
|
||||
`[schedule][停止已运行任务] 任务ID: ${cron.id}, PID: ${existingCron.pid}`,
|
||||
);
|
||||
await killTask(existingCron.pid);
|
||||
await killTask(existingCron.pid, true);
|
||||
// Mark old running instances as stopped
|
||||
const stoppedAt = dayjs().unix();
|
||||
await RunningInstanceModel.update(
|
||||
@@ -53,6 +53,7 @@ export function runCron(cmd: string, cron: ICron): Promise<number | void> {
|
||||
Logger.error(
|
||||
`[schedule][检查已运行任务失败] 任务ID: ${cron.id}, 错误: ${error}`,
|
||||
);
|
||||
throw error;
|
||||
}
|
||||
|
||||
Logger.info(
|
||||
|
||||
+5
-2
@@ -91,8 +91,11 @@ if [[ $command != "reload" ]]; then
|
||||
pip3 install --prefix ${PYTHON_HOME} requests
|
||||
fi
|
||||
|
||||
cd ${QL_DIR}
|
||||
cp -f .env.example .env
|
||||
cd "${QL_DIR}"
|
||||
if [[ ! -e .env ]]; then
|
||||
# Preserve user configuration, including a file created concurrently.
|
||||
(umask 077; set -C; cat .env.example > .env)
|
||||
fi
|
||||
chmod 777 ${QL_DIR}/shell/*.sh
|
||||
|
||||
. ${QL_DIR}/shell/share.sh
|
||||
|
||||
@@ -0,0 +1,77 @@
|
||||
const test = require('node:test');
|
||||
const assert = require('node:assert/strict');
|
||||
const fs = require('node:fs');
|
||||
const ts = require('typescript');
|
||||
|
||||
// Exercise the production function with a deterministic process-tree race.
|
||||
function fixture(descendants, probe = () => {}) {
|
||||
const source = fs.readFileSync('back/config/util.ts', 'utf8');
|
||||
const start = source.indexOf('export async function killTask(');
|
||||
const end = source.indexOf('export async function getPid(', start);
|
||||
const code = ts.transpileModule(source.slice(start, end), {
|
||||
compilerOptions: {
|
||||
target: ts.ScriptTarget.ES2020,
|
||||
module: ts.ModuleKind.CommonJS,
|
||||
},
|
||||
}).outputText;
|
||||
const signals = [];
|
||||
const exports = {};
|
||||
new Function('exports', 'process', 'psTree', 'fs', 'setTimeout', code)(
|
||||
exports,
|
||||
{
|
||||
platform: 'linux',
|
||||
kill: (pid, signal) => {
|
||||
const normalized = { 15: 'SIGTERM', 2: 'SIGINT' }[signal] || signal;
|
||||
signals.push([pid, normalized]);
|
||||
probe(pid, normalized);
|
||||
},
|
||||
},
|
||||
async () => [...descendants],
|
||||
{},
|
||||
setTimeout,
|
||||
);
|
||||
return { killTask: exports.killTask, signals };
|
||||
}
|
||||
|
||||
test('an exited descendant does not prevent stopping its siblings and parent', async () => {
|
||||
const { killTask, signals } = fixture([101, 102], (pid) => {
|
||||
if (pid === 102) throw Object.assign(Error('gone'), { code: 'ESRCH' });
|
||||
});
|
||||
await killTask(100);
|
||||
assert.deepEqual(signals, [
|
||||
[102, 'SIGTERM'],
|
||||
[101, 'SIGTERM'],
|
||||
[100, 'SIGTERM'],
|
||||
]);
|
||||
});
|
||||
|
||||
test('an already exited leaf can be stopped repeatedly', async () => {
|
||||
const { killTask } = fixture([], () => {
|
||||
throw Object.assign(Error('gone'), { code: 'ESRCH' });
|
||||
});
|
||||
await killTask(100);
|
||||
await killTask(100);
|
||||
});
|
||||
|
||||
for (const descendants of [[], [101]]) {
|
||||
test(`signal permission failure is surfaced with ${descendants.length} descendants`, async () => {
|
||||
const denied = Object.assign(Error('denied'), { code: 'EPERM' });
|
||||
const { killTask } = fixture(descendants, () => {
|
||||
throw denied;
|
||||
});
|
||||
await assert.rejects(killTask(100), (error) => error === denied);
|
||||
});
|
||||
}
|
||||
|
||||
test('a leaf retains SIGINT and a process tree retains child-first SIGTERM', async () => {
|
||||
const leaf = fixture([]);
|
||||
await leaf.killTask(100);
|
||||
assert.deepEqual(leaf.signals, [[100, 'SIGINT']]);
|
||||
const tree = fixture([101, 102]);
|
||||
await tree.killTask(100);
|
||||
assert.deepEqual(tree.signals, [
|
||||
[102, 'SIGTERM'],
|
||||
[101, 'SIGTERM'],
|
||||
[100, 'SIGTERM'],
|
||||
]);
|
||||
});
|
||||
@@ -0,0 +1,53 @@
|
||||
const test = require('node:test');
|
||||
const assert = require('node:assert/strict');
|
||||
const fs = require('node:fs');
|
||||
const os = require('node:os');
|
||||
const path = require('node:path');
|
||||
const { spawnSync } = require('node:child_process');
|
||||
|
||||
// Run the actual startup prefix, stopping before service initialization.
|
||||
const source = fs
|
||||
.readFileSync('shell/start.sh', 'utf8')
|
||||
.split('. ${QL_DIR}/shell/share.sh')[0];
|
||||
for (const command of ['start', 'reload']) {
|
||||
for (const existing of [false, true]) {
|
||||
test(`${command} ${
|
||||
existing ? 'preserves existing' : 'initializes missing'
|
||||
} .env`, (t) => {
|
||||
const root = fs.mkdtempSync(path.join(os.tmpdir(), 'ql-env-'));
|
||||
t.after(() => fs.rmSync(root, { recursive: true, force: true }));
|
||||
fs.mkdirSync(path.join(root, 'shell'));
|
||||
fs.writeFileSync(path.join(root, 'shell/test.sh'), '# fixture');
|
||||
fs.writeFileSync(path.join(root, '.env.example'), 'BACK_PORT=5700\n');
|
||||
const custom = 'BACK_PORT=5800\nJWT_SECRET=synthetic-test-only\n';
|
||||
if (existing)
|
||||
fs.writeFileSync(path.join(root, '.env'), custom, { mode: 0o600 });
|
||||
const result = spawnSync(
|
||||
'/bin/bash',
|
||||
[
|
||||
'-c',
|
||||
'npm(){ :; }; pip3(){ :; }; apk(){ :; }; sudo(){ "$@"; }; python3(){ printf 3.11; };\n' +
|
||||
source,
|
||||
'startup-test',
|
||||
command,
|
||||
],
|
||||
{
|
||||
encoding: 'utf8',
|
||||
env: {
|
||||
...process.env,
|
||||
QL_DIR: root,
|
||||
QL_DATA_DIR: root + '/data',
|
||||
QL_OS_TYPE: 'alpine',
|
||||
},
|
||||
},
|
||||
);
|
||||
assert.equal(result.status, 0, result.stderr);
|
||||
assert.equal(
|
||||
fs.readFileSync(path.join(root, '.env'), 'utf8'),
|
||||
existing ? custom : 'BACK_PORT=5700\n',
|
||||
);
|
||||
if (existing)
|
||||
assert.equal(fs.statSync(path.join(root, '.env')).mode & 0o777, 0o600);
|
||||
});
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,158 @@
|
||||
const test = require('node:test');
|
||||
const assert = require('node:assert/strict');
|
||||
const load = require('../helpers/load-security-module.cjs');
|
||||
|
||||
const logger = { info() {}, error() {} };
|
||||
|
||||
function cronFixture(killTask) {
|
||||
const events = [];
|
||||
const { runCron } = load('back/shared/runCron.ts', {
|
||||
'cross-spawn': {
|
||||
spawn: () => {
|
||||
events.push('spawn');
|
||||
return {};
|
||||
},
|
||||
},
|
||||
'./childProcess': {
|
||||
observeChildProcess: () => ({ completed: Promise.resolve({ code: 0 }) }),
|
||||
asError: (e) => e,
|
||||
},
|
||||
'./pLimit': {
|
||||
runWithCronLimit: async (_, fn) => fn(),
|
||||
removeQueuedCron: () => events.push('release'),
|
||||
},
|
||||
'../loaders/logger': logger,
|
||||
'../config/util': { killTask },
|
||||
'../data/cron': {
|
||||
CrontabModel: {
|
||||
findOne: async () => ({
|
||||
pid: 100,
|
||||
status: 0,
|
||||
allow_multiple_instances: 0,
|
||||
}),
|
||||
update: async () => events.push('idle'),
|
||||
},
|
||||
CrontabStatus: { running: 0, idle: 1, queued: 3 },
|
||||
},
|
||||
'../data/runningInstance': {
|
||||
RunningInstanceModel: { update: async () => events.push('stopped') },
|
||||
InstanceStatus: { running: 0, stopped: 2 },
|
||||
},
|
||||
});
|
||||
return { run: () => runCron('ignored', { id: '1' }), events };
|
||||
}
|
||||
|
||||
test('single-instance replacement does not spawn or finalize when termination fails', async () => {
|
||||
const f = cronFixture(async () => {
|
||||
throw Error('EPERM');
|
||||
});
|
||||
await f.run();
|
||||
assert.deepEqual(f.events, ['release']);
|
||||
});
|
||||
|
||||
test('single-instance replacement waits for verified exit before spawning', async () => {
|
||||
let finish, entered, verified;
|
||||
const pending = new Promise((resolve) => {
|
||||
entered = resolve;
|
||||
});
|
||||
const gate = new Promise((resolve) => {
|
||||
finish = resolve;
|
||||
});
|
||||
const f = cronFixture(async (pid, wait) => {
|
||||
assert.equal(pid, 100);
|
||||
verified = wait;
|
||||
entered();
|
||||
await gate;
|
||||
});
|
||||
const running = f.run();
|
||||
await pending;
|
||||
assert.deepEqual(f.events, []);
|
||||
finish();
|
||||
await running;
|
||||
assert.equal(verified, true);
|
||||
assert.deepEqual(f.events, ['stopped', 'idle', 'spawn', 'release']);
|
||||
});
|
||||
|
||||
function subscriptionFixture(killTask) {
|
||||
const updates = [];
|
||||
const Service = load('back/services/subscription.ts', {
|
||||
'../config': {},
|
||||
'../data/cron': {},
|
||||
'../config/const': {},
|
||||
'../data/subscription': {
|
||||
SubscriptionModel: {
|
||||
findAll: async () => [
|
||||
{ id: 1, pid: 101 },
|
||||
{ id: 2, pid: 102 },
|
||||
],
|
||||
update: async (values, query) => updates.push({ values, query }),
|
||||
},
|
||||
SubscriptionStatus: { idle: 1 },
|
||||
},
|
||||
'../config/util': { killTask },
|
||||
'../config/subscription': {},
|
||||
'../shared/i18n': {},
|
||||
'../shared/pLimit': {},
|
||||
'../shared/logReader': {},
|
||||
'../shared/logStreamManager': {},
|
||||
'./schedule': {},
|
||||
'./sock': {},
|
||||
'./sshKey': {},
|
||||
'./cron': {},
|
||||
}).default;
|
||||
return { service: new Service(logger, {}, {}, {}, {}), updates };
|
||||
}
|
||||
|
||||
test('batch subscription stop preserves failed items and reports failure while stopping others', async () => {
|
||||
const calls = [];
|
||||
const denied = Error('EPERM');
|
||||
const f = subscriptionFixture(async (pid, wait) => {
|
||||
calls.push([pid, wait]);
|
||||
if (pid === 101) throw denied;
|
||||
});
|
||||
await assert.rejects(f.service.stop([1, 2]), (e) => e === denied);
|
||||
assert.deepEqual(calls, [
|
||||
[101, true],
|
||||
[102, true],
|
||||
]);
|
||||
assert.equal(f.updates.length, 1);
|
||||
assert.equal(f.updates[0].query.where.id, 2);
|
||||
assert.equal(f.updates[0].values.status, 1);
|
||||
});
|
||||
|
||||
function scriptFixture(killTask, pid = 100) {
|
||||
const Service = load('back/services/script.ts', {
|
||||
'../config': { scriptPath: '/scripts' },
|
||||
'../config/const': { TASK_COMMAND: 'task' },
|
||||
'../config/util': { killTask, getPid: async () => pid },
|
||||
'../shared/pLimit': { removeQueuedCron() {} },
|
||||
'./sock': {},
|
||||
'./cron': {},
|
||||
'./schedule': {},
|
||||
'../shared/fileAccess': {},
|
||||
}).default;
|
||||
return new Service(logger, {}, {}, {});
|
||||
}
|
||||
|
||||
test('script stop propagates failure instead of returning success', async () => {
|
||||
const denied = Error('EPERM');
|
||||
const service = scriptFixture(async () => {
|
||||
throw denied;
|
||||
});
|
||||
await assert.rejects(
|
||||
service.stopScript('/scripts/test.js', 100),
|
||||
(e) => e === denied,
|
||||
);
|
||||
});
|
||||
|
||||
test('script stop verifies exit and treats an absent queued process as already stopped', async () => {
|
||||
const calls = [];
|
||||
const service = scriptFixture(async (...args) => calls.push(args), undefined);
|
||||
assert.equal((await service.stopScript('/scripts/test.js', 100)).code, 200);
|
||||
assert.deepEqual(calls, [[100, true]]);
|
||||
const absent = scriptFixture(async () => assert.fail('no PID to stop'), null);
|
||||
assert.equal(
|
||||
(await absent.stopScript('/scripts/test.js', undefined)).code,
|
||||
200,
|
||||
);
|
||||
});
|
||||
@@ -0,0 +1,66 @@
|
||||
const test = require('node:test');
|
||||
const assert = require('node:assert/strict');
|
||||
const { Sequelize, DataTypes } = require('sequelize');
|
||||
const load = require('../helpers/load-security-module.cjs');
|
||||
|
||||
async function fixture(t, onKill = async () => {}) {
|
||||
const db = new Sequelize({
|
||||
dialect: 'sqlite',
|
||||
storage: ':memory:',
|
||||
logging: false,
|
||||
});
|
||||
t.after(() => db.close());
|
||||
const rows = db.define('Subscription', {
|
||||
pid: DataTypes.INTEGER,
|
||||
status: DataTypes.INTEGER,
|
||||
});
|
||||
await db.sync();
|
||||
await rows.create({ id: 1, pid: 101, status: 0 });
|
||||
const Service = load('back/services/subscription.ts', {
|
||||
'../config': {},
|
||||
'../data/cron': {},
|
||||
'../config/const': {},
|
||||
'../data/subscription': {
|
||||
SubscriptionModel: rows,
|
||||
SubscriptionStatus: { idle: 1 },
|
||||
},
|
||||
'../config/util': { killTask: async () => onKill(rows) },
|
||||
'../config/subscription': {},
|
||||
'../shared/i18n': {},
|
||||
'../shared/pLimit': {},
|
||||
'../shared/logReader': {},
|
||||
'../shared/logStreamManager': {},
|
||||
'./schedule': {},
|
||||
'./sock': {},
|
||||
'./sshKey': {},
|
||||
'./cron': {},
|
||||
}).default;
|
||||
return { rows, service: new Service({ error() {} }, {}, {}, {}, {}) };
|
||||
}
|
||||
|
||||
test('successful subscription stop clears the stored PID', async (t) => {
|
||||
const f = await fixture(t);
|
||||
await f.service.stop([1]);
|
||||
const row = await f.rows.findByPk(1);
|
||||
assert.equal(row.status, 1);
|
||||
assert.equal(row.pid, null);
|
||||
});
|
||||
|
||||
test('subscription stop does not overwrite a replacement process started during termination', async (t) => {
|
||||
const f = await fixture(t, async (rows) => {
|
||||
await rows.update({ pid: 202, status: 0 }, { where: { id: 1 } });
|
||||
});
|
||||
await f.service.stop([1]);
|
||||
const row = await f.rows.findByPk(1);
|
||||
assert.equal(row.status, 0);
|
||||
assert.equal(row.pid, 202);
|
||||
});
|
||||
|
||||
test('subscription stop cancels a queued row without a PID', async (t) => {
|
||||
const f = await fixture(t, () => assert.fail('queued row has no process'));
|
||||
await f.rows.update({ pid: null, status: 3 }, { where: { id: 1 } });
|
||||
await f.service.stop([1]);
|
||||
const row = await f.rows.findByPk(1);
|
||||
assert.equal(row.status, 1);
|
||||
assert.equal(row.pid, null);
|
||||
});
|
||||
Reference in New Issue
Block a user