diff --git a/back/schedule/addCron.ts b/back/schedule/addCron.ts index bdb29da0..5ebdf7a5 100644 --- a/back/schedule/addCron.ts +++ b/back/schedule/addCron.ts @@ -68,7 +68,11 @@ const addCron = ( schedule, async () => { Logger.info('[schedule][准备运行任务] 命令: %s', item.command); - await runCron(item.command, item); + await runCron( + item.command, + item, + () => scheduleStacks.get(item.id) === jobs, + ); }, { name: `${item.id}: ${item.name || ''}`, diff --git a/back/shared/interface.ts b/back/shared/interface.ts index 01fd6cb7..2ea99f04 100644 --- a/back/shared/interface.ts +++ b/back/shared/interface.ts @@ -18,6 +18,7 @@ export interface IDependencyFn { export interface ICronFn { (): Promise; cron?: TCron; + isCurrent?: () => boolean; } export interface ISchedule { diff --git a/back/shared/pLimit.ts b/back/shared/pLimit.ts index 47f8b6c8..969a6d58 100644 --- a/back/shared/pLimit.ts +++ b/back/shared/pLimit.ts @@ -50,11 +50,9 @@ class TaskLimit { Buffer.from(tlsConfig.clientCert), ) : credentials.createInsecure(); - this._client = new ApiClient( - `localhost:${config.grpcPort}`, - creds, - { 'grpc.enable_http_proxy': 0 }, - ); + this._client = new ApiClient(`localhost:${config.grpcPort}`, creds, { + 'grpc.enable_http_proxy': 0, + }); } return this._client; } @@ -113,12 +111,15 @@ class TaskLimit { } } - public removeQueuedCron(id: string) { + public removeQueuedCron(id: string, completed?: ICronFn) { if (this.queuedCrons.has(id)) { const runs = this.queuedCrons.get(id); if (runs && runs.length > 0) { - runs.pop(); - this.queuedCrons.set(id, runs); + const remaining = completed + ? runs.filter((run) => run !== completed) + : runs.slice(0, -1); + if (remaining.length) this.queuedCrons.set(id, remaining); + else this.queuedCrons.delete(id); } } } @@ -143,9 +144,15 @@ class TaskLimit { cron: TCron, fn: ICronFn, options?: Partial, + isCurrent: () => boolean = () => true, ): Promise { + if (!isCurrent()) return; fn.cron = cron; - let runs = this.queuedCrons.get(cron.id); + fn.isCurrent = isCurrent; + // Invalidated snapshots must not consume the new revision's repeat budget. + const runs = this.queuedCrons + .get(cron.id) + ?.filter((run) => run.isCurrent?.() !== false); const result = runs?.length ? [...runs, fn] : [fn]; const repeatTimes = this.repeatCronNotifyMap.get(cron.id) || 0; if (result?.length > 5) { diff --git a/back/shared/runCron.ts b/back/shared/runCron.ts index ab5511a8..8a862681 100644 --- a/back/shared/runCron.ts +++ b/back/shared/runCron.ts @@ -8,15 +8,23 @@ import { RunningInstanceModel, InstanceStatus } from '../data/runningInstance'; import dayjs from 'dayjs'; import { observeChildProcess, asError } from './childProcess'; -export function runCron(cmd: string, cron: ICron): Promise { - return taskLimit.runWithCronLimit(cron, async () => { +export function runCron( + cmd: string, + cron: ICron, + isCurrent: () => boolean = () => true, +): Promise { + const execute = async () => { try { + if (!isCurrent()) return; // Check if the cron is already running and stop it (only if multiple instances are not allowed) try { const existingCron = await CrontabModel.findOne({ where: { id: Number(cron.id) }, }); + // A mutation may commit while this queued run is reading the database. + if (!isCurrent() || !existingCron || existingCron.isDisabled) return; + // Default to single instance mode (0) for backward compatibility const allowSingleInstances = existingCron?.allow_multiple_instances === 0; @@ -56,6 +64,8 @@ export function runCron(cmd: string, cron: ICron): Promise { throw error; } + // In particular, invalidation can happen while stopping a previous run. + if (!isCurrent()) return; Logger.info( `[schedule][开始执行任务] 参数 ${JSON.stringify({ ...cron, @@ -94,7 +104,8 @@ export function runCron(cmd: string, cron: ICron): Promise { asError(error).message, ); } finally { - taskLimit.removeQueuedCron(cron.id); + taskLimit.removeQueuedCron(cron.id, execute); } - }); + }; + return taskLimit.runWithCronLimit(cron, execute, undefined, isCurrent); } diff --git a/cli/test/integration/differential.test.cjs b/cli/test/integration/differential.test.cjs index e6bb8b50..76a0e348 100644 --- a/cli/test/integration/differential.test.cjs +++ b/cli/test/integration/differential.test.cjs @@ -31,7 +31,7 @@ test('Shell and CLI agree on isolated script state and account modes', async (t) path.join(root, 'static/build'), ]) await fs.mkdir(directory, { recursive: true }); - for (const name of ['task.sh', 'otask.sh', 'share.sh', 'api.sh', 'env.sh']) + for (const name of ['task.sh', 'task-timeout.sh', 'otask.sh', 'share.sh', 'api.sh', 'env.sh']) await fs.copyFile( path.resolve(__dirname, '../../../shell', name), path.join(context.paths.dir_shell, name), @@ -355,6 +355,7 @@ t() { :; } sleep() { printf DELAYED; } enter_script_workdir() { :; } run_else() { :; } +run_task_command() { "$@"; } which_program=true main "$@" `, diff --git a/cli/test/integration/executablePaths.test.cjs b/cli/test/integration/executablePaths.test.cjs index 41ed5319..8cbf8013 100644 --- a/cli/test/integration/executablePaths.test.cjs +++ b/cli/test/integration/executablePaths.test.cjs @@ -10,7 +10,7 @@ test('legacy Shell and TS resolve executable paths after selecting the working d t.after(() => fs.rmSync(root, { recursive: true, force: true })); for (const dir of ['shell/preload', 'shell/lang', 'data/config', 'data/scripts/bin', 'data/scripts/custom', 'data/log', 'static/build', 'bin']) fs.mkdirSync(path.join(root, dir), { recursive: true }); - for (const name of ['task.sh', 'otask.sh', 'share.sh', 'api.sh', 'env.sh']) + for (const name of ['task.sh', 'task-timeout.sh', 'otask.sh', 'share.sh', 'api.sh', 'env.sh']) fs.copyFileSync(path.resolve(__dirname, '../../../shell', name), path.join(root, 'shell', name)); for (const name of ['zh.sh', 'en.sh']) fs.copyFileSync(path.resolve(__dirname, '../../../shell/lang', name), path.join(root, 'shell/lang', name)); diff --git a/cli/test/integration/stdin.test.cjs b/cli/test/integration/stdin.test.cjs index 6c4100dc..aac8b5a7 100644 --- a/cli/test/integration/stdin.test.cjs +++ b/cli/test/integration/stdin.test.cjs @@ -10,7 +10,7 @@ test('legacy Shell and TS preserve task stdin in each execution mode', t => { t.after(() => fs.rmSync(root, { recursive: true, force: true })); for (const dir of ['shell/preload', 'shell/lang', 'data/config', 'data/scripts', 'data/log', 'static/build', 'bin']) fs.mkdirSync(path.join(root, dir), { recursive: true }); - for (const name of ['task.sh', 'otask.sh', 'share.sh', 'api.sh', 'env.sh']) + for (const name of ['task.sh', 'task-timeout.sh', 'otask.sh', 'share.sh', 'api.sh', 'env.sh']) fs.copyFileSync(path.resolve(__dirname, '../../../shell', name), path.join(root, 'shell', name)); for (const name of ['zh.sh', 'en.sh']) fs.copyFileSync(path.resolve(__dirname, '../../../shell/lang', name), path.join(root, 'shell/lang', name)); diff --git a/docs/development/task-startup.md b/docs/development/task-startup.md index fffc9c74..7546f1c9 100644 --- a/docs/development/task-startup.md +++ b/docs/development/task-startup.md @@ -59,3 +59,20 @@ resolved. Provider initialization, including provider configuration reads, now happens at first use. A script's explicit `import __ql_notify__`, `require(...)`, or independent notification helper still loads that module immediately. + +## Task timeouts + +`task -m DURATION ...` also applies to sourced Shell scripts. The timeout worker +inherits before-hook variables, arrays and functions without serializing them. +At the deadline, the worker's process group receives TERM, followed by KILL after +0.5 seconds, so ordinary child processes cannot continue after the timeout. +Explicit interpreters and other commands use the same mechanism. Processes that +intentionally create a separate session/process group are outside this boundary. +The wrapper records status 124 and runs its after hook; concurrent account mode +retains its existing aggregate status behavior. Zero disables the timeout. + +Cron executions waiting for a concurrency slot retain their schedule revision. +Disabling, deleting or replacing that revision prevents its queued executions +from starting when a slot becomes available. Already running executions retain +the existing stop behavior. Invalidated callbacks do not count against a new +revision's limit of five pending/running executions. diff --git a/shell/otask.sh b/shell/otask.sh index c11041f3..f826bfa5 100755 --- a/shell/otask.sh +++ b/shell/otask.sh @@ -164,7 +164,7 @@ run_normal() { if [[ $isJsOrPythonFile == 'false' ]]; then clear_non_sh_env fi - $timeoutCmd $which_program $file_param "${script_params[@]}" + run_task_command $which_program $file_param "${script_params[@]}" } handle_env_split() { @@ -203,7 +203,7 @@ run_concurrent() { export "${env_param}=${array[$i - 1]}" clear_non_sh_env fi - eval envParam="${env_param}" numParam="${i}" $timeoutCmd $which_program $file_param "${script_params[@]}" &>$single_log_path & + eval envParam="${env_param}" numParam="${i}" run_task_command $which_program $file_param "${script_params[@]}" &>$single_log_path & done wait @@ -243,7 +243,7 @@ run_designated() { enter_script_workdir - envParam="${env_param}" numParam="${num_param}" $timeoutCmd $which_program $file_param "${script_params[@]}" + envParam="${env_param}" numParam="${num_param}" run_task_command $which_program $file_param "${script_params[@]}" } ## 运行其他命令 @@ -299,7 +299,7 @@ run_else() { fi clear_non_sh_env - $timeoutCmd $which_program $file_param "$@" + run_task_command $which_program $file_param "$@" } check_file() { @@ -336,10 +336,6 @@ check_nounset() { main() { if [[ $1 == *.js ]] || [[ $1 == *.mjs ]] || [[ $1 == *.py ]] || [[ $1 == *.pyc ]] || [[ $1 == *.sh ]] || [[ $1 == *.ts ]]; then - if [[ $1 == *.sh ]]; then - timeoutCmd="" - fi - case $# in 1) run_normal "$1" diff --git a/shell/task-timeout.sh b/shell/task-timeout.sh new file mode 100644 index 00000000..a4448d97 --- /dev/null +++ b/shell/task-timeout.sh @@ -0,0 +1,63 @@ +#!/usr/bin/env bash + +# Fork from the current Bash state: an external `bash -c` would lose private +# before-hook variables, arrays and functions used by sourced Shell scripts. +run_task_command() { + local timeout_seconds + if [[ -z ${command_timeout_time:-} ]]; then + "$@" + return $? + fi + timeout_seconds=$(awk -v value="$command_timeout_time" 'BEGIN { + if (value !~ /^([0-9]+([.][0-9]*)?|[.][0-9]+)([eE][+-]?[0-9]+)?[smhd]?$/) exit 1 + unit = substr(value, length(value), 1) + factor = unit == "d" ? 86400 : unit == "h" ? 3600 : unit == "m" ? 60 : 1 + printf "%.9f", (value + 0) * factor + }') || { printf 'Invalid task timeout: %s\n' "$command_timeout_time" >&2; return 125; } + if [[ $timeout_seconds == '0.000000000' ]]; then + "$@" + return $? + fi + ( + timeout_dir=$(mktemp -d "${TMPDIR:-/tmp}/ql-task-timeout.XXXXXXXX") || exit 125 + # Both children get private process groups. The timer and the wrapper must + # survive terminating the task, and descendants must not outlive a timeout. + set -m + ( set +m; "$@" ) <&0 & + timeout_worker=$! + ( + set +m + : > "$timeout_dir/ready" + sleep "$timeout_seconds" || exit 125 + : > "$timeout_dir/expired" + kill -TERM -- "-$timeout_worker" 2>/dev/null || : + sleep 0.5 + kill -KILL -- "-$timeout_worker" 2>/dev/null || : + ) & + timeout_watcher=$! + set +m + trap 'kill -KILL -- "-$timeout_watcher" 2>/dev/null || : + if kill -TERM -- "-$timeout_worker" 2>/dev/null; then sleep 0.1; fi + kill -KILL -- "-$timeout_worker" 2>/dev/null || : + wait "$timeout_worker" 2>/dev/null || : + wait "$timeout_watcher" 2>/dev/null || : + rm -rf -- "$timeout_dir"' EXIT + trap 'exit 130' INT + trap 'exit 143' TERM + trap 'exit 129' HUP + trap 'exit 131' QUIT + # Do not clean up a timer whose process group has not initialized yet. + while [[ ! -f "$timeout_dir/ready" ]]; do + kill -0 "$timeout_watcher" 2>/dev/null || exit 125 + sleep 0.001 + done + wait "$timeout_worker" + timeout_code=$? + if [[ -f "$timeout_dir/expired" ]]; then + # Let the timer finish escalation even if the task leader exited first. + wait "$timeout_watcher" 2>/dev/null || : + exit 124 + fi + exit "$timeout_code" + ) +} diff --git a/shell/task.sh b/shell/task.sh index a4fe900c..92bd620f 100755 --- a/shell/task.sh +++ b/shell/task.sh @@ -3,6 +3,7 @@ dir_shell=$QL_DIR/shell . $dir_shell/share.sh . $dir_shell/api.sh +. $dir_shell/task-timeout.sh trap 'single_hanle SIGINT' INT trap 'single_hanle SIGTERM' TERM @@ -112,11 +113,6 @@ format_params() { mtime_format="%Y-%m-%d %H:%M:%S.%3N" fi timeoutCmd="" - if [[ $command_timeout_time ]]; then - if type timeout &>/dev/null; then - timeoutCmd="timeout --foreground -s 2 -k 10s $command_timeout_time " - fi - fi # params=$(echo "$@" | sed -E 's/([^ ])&([^ ])/\1\\\&\2/g') # 分割 task 内置参数和脚本参数 diff --git a/test/back/execution-lifecycle.test.cjs b/test/back/execution-lifecycle.test.cjs index 3e0b9e91..afa945e6 100644 --- a/test/back/execution-lifecycle.test.cjs +++ b/test/back/execution-lifecycle.test.cjs @@ -36,7 +36,7 @@ test( }, '../loaders/logger': logger, '../data/cron': { - CrontabModel: { findOne: async () => null }, + CrontabModel: { findOne: async () => ({ isDisabled: 0 }) }, CrontabStatus: {}, }, '../data/runningInstance': { diff --git a/test/back/queued-cron-invalidation.test.cjs b/test/back/queued-cron-invalidation.test.cjs new file mode 100644 index 00000000..a08417e4 --- /dev/null +++ b/test/back/queued-cron-invalidation.test.cjs @@ -0,0 +1,98 @@ +const test = require('node:test'); +const assert = require('node:assert/strict'); +const load = require('../helpers/load-security-module.cjs'); +const tick = () => new Promise(setImmediate); +function gate() { let release; const promise = new Promise(r => release = r); return { promise, release }; } + +async function fixture() { + const logger = { info() {}, error() {}, warn() {} }, notices = []; + const limit = load('back/shared/pLimit.ts', { + '../data/system': { AuthDataType: { systemConfig: 'systemConfig' }, SystemModel: { sync: async () => {}, findOne: async () => ({ info: { cronConcurrency: 1 } }) } }, + '../loaders/logger': logger, + '../services/notify': class {}, + '../shared/i18n': { t: s => s, tf: s => s }, + '../config': { grpcPort: 1 }, + '../protos/api': { ApiClient: class { systemNotify(value, cb) { notices.push(value); cb(); } } }, + '../config/grpcCerts': { getGrpcCerts: () => null }, + }).default; + await tick(); + const spawned = [], records = new Map(), jobs = new Map(); + let readGate, killGate; + const { runCron } = load('back/shared/runCron.ts', { + 'cross-spawn': { spawn: command => { spawned.push(command); return {}; } }, + './pLimit': limit, + '../loaders/logger': logger, + '../data/cron': { CrontabModel: { findOne: async ({ where }) => { const row = records.get(String(where.id)); if (readGate) await readGate.promise; return row; }, update: async () => {} }, CrontabStatus: { running: 0, queued: 3, idle: 1 } }, + '../data/runningInstance': { RunningInstanceModel: { update: async () => {} }, InstanceStatus: { running: 0, stopped: 2 } }, + '../config/util': { killTask: async () => { if (killGate) await killGate.promise; } }, + './childProcess': { observeChildProcess: () => ({ completed: Promise.resolve({ code: 0 }) }), asError: e => e }, + }); + const { addCron } = load('back/schedule/addCron.ts', { + '../shared/cronScheduler': { createCronJob: (schedule, callback) => ({ start() {}, cancel() {}, fire: callback }) }, + '../shared/cronSchedule': { isValidCronSchedule: x => x !== 'invalid' }, + './data': { scheduleStacks: jobs }, '../shared/runCron': { runCron }, '../loaders/logger': logger, + '../shared/i18n': { tf: s => s }, + }); + const { delCron } = load('back/schedule/delCron.ts', { './data': { scheduleStacks: jobs }, '../loaders/logger': logger }); + function add(id = '1', command = 'old', replace = false) { + records.set(id, { isDisabled: 0, allow_multiple_instances: 1 }); + addCron({ request: { crons: [{ id, command, schedule: '* * * * * *' }], replace } }, err => assert.ifError(err)); + return jobs.get(id)[0]; + } + function remove(id = '1') { delCron({ request: { ids: [id] } }, err => assert.ifError(err)); } + async function block() { const g = gate(); const done = limit.runWithCronLimit({ id: 'blocker' }, async () => { await g.promise; }); await tick(); return { release: async () => { g.release(); await done; } }; } + return { limit, spawned, records, jobs, add, remove, block, notices, setReadGate: g => readGate = g, setKillGate: g => killGate = g }; +} + +for (const operation of ['disable', 'delete', 'update', 'replace', 'disable-enable']) { + test(`queued scheduled execution is invalidated by ${operation}`, async () => { + const f = await fixture(), old = f.add(), blocker = await f.block(); + const queued = old.fire(); + if (operation === 'disable' || operation === 'disable-enable') { f.records.get('1').isDisabled = 1; f.remove(); } + if (operation === 'delete') { f.records.delete('1'); f.remove(); } + if (operation === 'update') f.add('1', 'new'); + if (operation === 'replace') f.add('1', 'new', true); + if (operation === 'disable-enable') f.add('1', 'old'); + const next = f.jobs.get('1')?.[0].fire(); + await blocker.release(); await queued; await next; + assert.deepEqual(f.spawned, ['update', 'replace'].includes(operation) ? ['new'] : operation === 'disable-enable' ? ['old'] : []); + assert.equal(f.limit.cronLimitPendingCount, 0); + assert.equal(f.limit.cronLimitActiveCount, 0); + }); +} + +test('stale callbacks cannot re-enqueue after cancellation', async () => { + const f = await fixture(), old = f.add(); f.remove(); await old.fire(); + assert.deepEqual(f.spawned, []); +}); + +test('invalidation during an asynchronous database read prevents spawning', async () => { + const f = await fixture(), old = f.add(), pending = gate(); f.setReadGate(pending); + const run = old.fire(); await tick(); f.remove(); pending.release(); await run; + assert.deepEqual(f.spawned, []); +}); + +test('invalidation while replacing a running instance prevents the old spawn', async () => { + const f = await fixture(), old = f.add(), pending = gate(); + f.records.set('1', { isDisabled: 0, allow_multiple_instances: 0, pid: 10, status: 0 }); f.setKillGate(pending); + const run = old.fire(); await tick(); f.remove(); pending.release(); await run; + assert.deepEqual(f.spawned, []); +}); + +test('database removal or disable is respected before scheduler RPC arrives', async () => { + for (const disabled of [true, false]) { + const f = await fixture(), old = f.add(), blocker = await f.block(); const run = old.fire(); + if (disabled) f.records.get('1').isDisabled = 1; else f.records.delete('1'); + await blocker.release(); await run; assert.deepEqual(f.spawned, []); + } +}); + +test('five stale runs do not suppress the new revision or erase its repeat accounting', async () => { + const f = await fixture(), old = f.add(), blocker = await f.block(); + const pending = Array.from({ length: 5 }, () => old.fire()); + const current = f.add('1', 'new'); pending.push(...Array.from({ length: 5 }, () => current.fire())); + assert.equal(f.notices.length, 0); + await blocker.release(); await Promise.all(pending); + assert.deepEqual(f.spawned, Array(5).fill('new')); + await current.fire(); assert.equal(f.spawned.length, 6); +}); diff --git a/test/back/task-shell-lifecycle.test.cjs b/test/back/task-shell-lifecycle.test.cjs index 9aefe072..6f8ddf2a 100644 --- a/test/back/task-shell-lifecycle.test.cjs +++ b/test/back/task-shell-lifecycle.test.cjs @@ -12,6 +12,7 @@ function extract(file, name) { return text.slice(start, text.indexOf('\n}', start) + 2); } const helpers = [ + fs.readFileSync('shell/task-timeout.sh', 'utf8'), ...['handle_task_start', 'handle_task_end', 'run_task_before', 'run_task_after', 'get_env_array', 'clear_env'].map(n => extract('shell/share.sh', n)), ...['run_shell_script', 'define_program', 'format_params'].map(n => extract('shell/task.sh', n)), ].join('\n'); @@ -24,7 +25,7 @@ function fixture(t) { write('env.sh', 'export QA_PANEL="alpha&beta&gamma"\n'); write('before.sh', 'export QA_BEFORE=ready\nqa_function() { printf "HOOK_FUNCTION\\n"; }\n'); write('after.sh', 'printf "AFTER:%s:%s\\n" "$QA_PANEL" "$QA_BEFORE" >> "$QA_ROOT/events"\n'); - const run = (args) => { + const run = (args, timeout = '') => { const r = spawnSync('/bin/bash', ['-c', helpers + ` dir_scripts=$QA_ROOT; dir_shell=$QA_ROOT; dir_dep=$QA_ROOT file_env=$QA_ROOT/env.sh; file_task_before=$QA_ROOT/before.sh; file_task_after=$QA_ROOT/after.sh @@ -40,7 +41,7 @@ function fixture(t) { format_params "$@"; define_program "${'${task_shell_params[@]}'}" . "$QA_TASK_SOURCE" printf 'WRAPPER_FINISHED\\n' >> "$QA_ROOT/events" - `, 'fixture', ...args], { cwd: root, env: { ...process.env, QA_ROOT: root, QA_TASK_SOURCE: taskSource }, encoding: 'utf8', timeout: 10000 }); + `, 'fixture', ...args], { cwd: root, env: { ...process.env, command_timeout_time: timeout, QA_ROOT: root, QA_TASK_SOURCE: taskSource }, encoding: 'utf8', timeout: 10000 }); assert.equal(r.status, 0, r.stdout + r.stderr); return { stdout: r.stdout, events: fs.readFileSync(path.join(root, 'events'), 'utf8').trim().split('\n') }; }; @@ -117,3 +118,77 @@ for (const mode of ['desi', 'conc']) { assert.equal(r.events.filter(x => x.startsWith('AFTER:')).length, 1); }); } + +for (const [label, args] of [ + ['sourced shell', ['probe.sh', '--', 'space value']], + ['explicit bash', ['bash', 'probe.sh', 'space value']], +]) { + test(`${label} timeout stops descendants and finalizes once with 124`, t => { + const f = fixture(t); + f.write('before.sh', 'export QA_BEFORE=ready\nprivate_value=private\nprivate_array=(one "two words")\nqa_function() { printf "HOOK_FUNCTION\\n"; }\n'); + f.write('probe.sh', `[[ "$QA_PANEL" == 'alpha&beta&gamma' && "$QA_BEFORE" == ready && "$1" == 'space value' ]] || exit 91\n${label === 'sourced shell' ? '[[ "$private_value" = private && "${private_array[1]}" = "two words" ]] || exit 92\nqa_function\n' : ''}trap 'echo EXIT_TRAP' EXIT\nsleep 3\necho SHOULD_NOT_RUN\n`); + const before = Date.now(); + const r = f.run(args, '0.2s'); + assert.ok(Date.now() - before < 2500, r.stdout); + assert.doesNotMatch(r.stdout, /SHOULD_NOT_RUN/); + assert.match(r.stdout, /EXIT_TRAP/); + assert.deepEqual(r.events, ['STATUS:0:', 'AFTER:alpha&beta&gamma:ready', 'STATUS:1:124', 'STAT:124', 'WRAPPER_FINISHED']); + }); +} + +test('timeout escalates when a shell and its children ignore TERM', t => { + const f = fixture(t); + f.write('probe.sh', 'trap "" TERM\nsleep 3\necho SHOULD_NOT_RUN\n'); + const before = Date.now(); + const r = f.run(['probe.sh'], '0.1'); + assert.ok(Date.now() - before < 2500); + assert.doesNotMatch(r.stdout, /SHOULD_NOT_RUN/); + assert.ok(r.events.includes('STATUS:1:124')); +}); + +test('completed timed tasks keep their exit code and do not wait for the timer', t => { + const f = fixture(t); + f.write('probe.sh', 'exit 7\n'); + const before = Date.now(); + const r = f.run(['probe.sh'], '1h'); + assert.ok(Date.now() - before < 1500); + assert.ok(r.events.includes('STATUS:1:7')); +}); + +test('zero disables the deadline and invalid durations never start the script', t => { + const f = fixture(t); + f.write('probe.sh', 'echo DID_RUN\n'); + assert.match(f.run(['probe.sh'], '0s').stdout, /DID_RUN/); + f.write('events', ''); + const r = f.run(['probe.sh'], 'wrong'); + assert.doesNotMatch(r.stdout, /DID_RUN/); + assert.ok(r.events.includes('STATUS:1:125')); +}); + +for (const [runtime, filename, body] of [ + ['python3', 'timed.py', 'import time\nprint("STARTED", flush=True)\ntime.sleep(3)\nprint("SHOULD_NOT_RUN", flush=True)\n'], + ['node', 'timed.cjs', 'console.log("STARTED"); setTimeout(() => console.log("SHOULD_NOT_RUN"), 3000);\n'], +]) { + test(`${runtime} timeout stops execution and reports 124`, t => { + const f = fixture(t); + f.write(filename, body); + const r = f.run([runtime, filename], '0.5s'); + assert.match(r.stdout, /STARTED/); + assert.doesNotMatch(r.stdout, /SHOULD_NOT_RUN/); + assert.ok(r.events.includes('STATUS:1:124')); + assert.equal(r.events.filter(x => x.startsWith('AFTER:')).length, 1); + }); +} + +for (const mode of ['desi', 'conc']) { + test(`timed shell ${mode} preserves selected accounts and terminates every worker`, t => { + const f = fixture(t); + f.write('probe.sh', 'echo "ACCOUNT=$QA_PANEL"\nsleep 3\necho SHOULD_NOT_RUN\n'); + const r = f.run(['probe.sh', mode, 'QA_PANEL', '2-3'], '0.2s'); + assert.match(r.stdout, mode === 'desi' ? /ACCOUNT=beta&gamma/ : /ACCOUNT=beta\nACCOUNT=gamma/); + assert.doesNotMatch(r.stdout, /SHOULD_NOT_RUN/); + assert.equal(r.events.filter(x => x.startsWith('AFTER:')).length, 1); + // Concurrent mode retains its existing aggregate wait status behavior. + assert.ok(r.events.includes(`STATUS:1:${mode === 'desi' ? 124 : 0}`)); + }); +}