From 2d6c2206b35d2a6bacec946826cd6e4d79ee9a7c Mon Sep 17 00:00:00 2001 From: whyour Date: Tue, 22 Sep 2026 22:23:18 +0800 Subject: [PATCH] fix: verify task termination and preserve startup configuration --- back/config/util.ts | 21 ++- back/services/script.ts | 6 +- back/services/subscription.ts | 22 +-- back/shared/runCron.ts | 3 +- shell/start.sh | 7 +- test/back/process-stop-signal.test.cjs | 77 ++++++++++ test/back/start-env-preservation.test.cjs | 53 +++++++ test/back/stop-failure-feedback.test.cjs | 158 +++++++++++++++++++++ test/back/subscription-stop-state.test.cjs | 66 +++++++++ 9 files changed, 386 insertions(+), 27 deletions(-) create mode 100644 test/back/process-stop-signal.test.cjs create mode 100644 test/back/start-env-preservation.test.cjs create mode 100644 test/back/stop-failure-feedback.test.cjs create mode 100644 test/back/subscription-stop-state.test.cjs diff --git a/back/config/util.ts b/back/config/util.ts index 13cad395..dc8bdb8f 100644 --- a/back/config/util.ts +++ b/back/config/util.ts @@ -491,17 +491,6 @@ export function psTree(pid: number): Promise { 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 { diff --git a/back/services/script.ts b/back/services/script.ts index ef4a5e2c..207fd0a1 100644 --- a/back/services/script.ts +++ b/back/services/script.ts @@ -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 }; } diff --git a/back/services/subscription.ts b/back/services/subscription.ts index 513d239f..16fe2b9e 100644 --- a/back/services/subscription.ts +++ b/back/services/subscription.ts @@ -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); - } catch (error) { - this.logger.error(error); + try { + 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) { diff --git a/back/shared/runCron.ts b/back/shared/runCron.ts index 29609e7a..ab5511a8 100644 --- a/back/shared/runCron.ts +++ b/back/shared/runCron.ts @@ -31,7 +31,7 @@ export function runCron(cmd: string, cron: ICron): Promise { 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 { Logger.error( `[schedule][检查已运行任务失败] 任务ID: ${cron.id}, 错误: ${error}`, ); + throw error; } Logger.info( diff --git a/shell/start.sh b/shell/start.sh index da615332..657f7aff 100644 --- a/shell/start.sh +++ b/shell/start.sh @@ -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 diff --git a/test/back/process-stop-signal.test.cjs b/test/back/process-stop-signal.test.cjs new file mode 100644 index 00000000..285aee41 --- /dev/null +++ b/test/back/process-stop-signal.test.cjs @@ -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'], + ]); +}); diff --git a/test/back/start-env-preservation.test.cjs b/test/back/start-env-preservation.test.cjs new file mode 100644 index 00000000..0b3bcc57 --- /dev/null +++ b/test/back/start-env-preservation.test.cjs @@ -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); + }); + } +} diff --git a/test/back/stop-failure-feedback.test.cjs b/test/back/stop-failure-feedback.test.cjs new file mode 100644 index 00000000..5b4290af --- /dev/null +++ b/test/back/stop-failure-feedback.test.cjs @@ -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, + ); +}); diff --git a/test/back/subscription-stop-state.test.cjs b/test/back/subscription-stop-state.test.cjs new file mode 100644 index 00000000..e99fa091 --- /dev/null +++ b/test/back/subscription-stop-state.test.cjs @@ -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); +});