mirror of
https://github.com/whyour/qinglong.git
synced 2026-09-30 04:31:11 +08:00
fix: invalidate queued cron revisions and enforce task process timeouts
This commit is contained in:
@@ -68,7 +68,11 @@ const addCron = (
|
|||||||
schedule,
|
schedule,
|
||||||
async () => {
|
async () => {
|
||||||
Logger.info('[schedule][准备运行任务] 命令: %s', item.command);
|
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 || ''}`,
|
name: `${item.id}: ${item.name || ''}`,
|
||||||
|
|||||||
@@ -18,6 +18,7 @@ export interface IDependencyFn<T> {
|
|||||||
export interface ICronFn<T> {
|
export interface ICronFn<T> {
|
||||||
(): Promise<T>;
|
(): Promise<T>;
|
||||||
cron?: TCron;
|
cron?: TCron;
|
||||||
|
isCurrent?: () => boolean;
|
||||||
}
|
}
|
||||||
|
|
||||||
export interface ISchedule {
|
export interface ISchedule {
|
||||||
|
|||||||
+16
-9
@@ -50,11 +50,9 @@ class TaskLimit {
|
|||||||
Buffer.from(tlsConfig.clientCert),
|
Buffer.from(tlsConfig.clientCert),
|
||||||
)
|
)
|
||||||
: credentials.createInsecure();
|
: credentials.createInsecure();
|
||||||
this._client = new ApiClient(
|
this._client = new ApiClient(`localhost:${config.grpcPort}`, creds, {
|
||||||
`localhost:${config.grpcPort}`,
|
'grpc.enable_http_proxy': 0,
|
||||||
creds,
|
});
|
||||||
{ 'grpc.enable_http_proxy': 0 },
|
|
||||||
);
|
|
||||||
}
|
}
|
||||||
return this._client;
|
return this._client;
|
||||||
}
|
}
|
||||||
@@ -113,12 +111,15 @@ class TaskLimit {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
public removeQueuedCron(id: string) {
|
public removeQueuedCron(id: string, completed?: ICronFn<any>) {
|
||||||
if (this.queuedCrons.has(id)) {
|
if (this.queuedCrons.has(id)) {
|
||||||
const runs = this.queuedCrons.get(id);
|
const runs = this.queuedCrons.get(id);
|
||||||
if (runs && runs.length > 0) {
|
if (runs && runs.length > 0) {
|
||||||
runs.pop();
|
const remaining = completed
|
||||||
this.queuedCrons.set(id, runs);
|
? 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,
|
cron: TCron,
|
||||||
fn: ICronFn<T>,
|
fn: ICronFn<T>,
|
||||||
options?: Partial<QueueAddOptions>,
|
options?: Partial<QueueAddOptions>,
|
||||||
|
isCurrent: () => boolean = () => true,
|
||||||
): Promise<T | void> {
|
): Promise<T | void> {
|
||||||
|
if (!isCurrent()) return;
|
||||||
fn.cron = cron;
|
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 result = runs?.length ? [...runs, fn] : [fn];
|
||||||
const repeatTimes = this.repeatCronNotifyMap.get(cron.id) || 0;
|
const repeatTimes = this.repeatCronNotifyMap.get(cron.id) || 0;
|
||||||
if (result?.length > 5) {
|
if (result?.length > 5) {
|
||||||
|
|||||||
+15
-4
@@ -8,15 +8,23 @@ import { RunningInstanceModel, InstanceStatus } from '../data/runningInstance';
|
|||||||
import dayjs from 'dayjs';
|
import dayjs from 'dayjs';
|
||||||
import { observeChildProcess, asError } from './childProcess';
|
import { observeChildProcess, asError } from './childProcess';
|
||||||
|
|
||||||
export function runCron(cmd: string, cron: ICron): Promise<number | void> {
|
export function runCron(
|
||||||
return taskLimit.runWithCronLimit(cron, async () => {
|
cmd: string,
|
||||||
|
cron: ICron,
|
||||||
|
isCurrent: () => boolean = () => true,
|
||||||
|
): Promise<number | void> {
|
||||||
|
const execute = async () => {
|
||||||
try {
|
try {
|
||||||
|
if (!isCurrent()) return;
|
||||||
// Check if the cron is already running and stop it (only if multiple instances are not allowed)
|
// Check if the cron is already running and stop it (only if multiple instances are not allowed)
|
||||||
try {
|
try {
|
||||||
const existingCron = await CrontabModel.findOne({
|
const existingCron = await CrontabModel.findOne({
|
||||||
where: { id: Number(cron.id) },
|
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
|
// Default to single instance mode (0) for backward compatibility
|
||||||
const allowSingleInstances =
|
const allowSingleInstances =
|
||||||
existingCron?.allow_multiple_instances === 0;
|
existingCron?.allow_multiple_instances === 0;
|
||||||
@@ -56,6 +64,8 @@ export function runCron(cmd: string, cron: ICron): Promise<number | void> {
|
|||||||
throw error;
|
throw error;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// In particular, invalidation can happen while stopping a previous run.
|
||||||
|
if (!isCurrent()) return;
|
||||||
Logger.info(
|
Logger.info(
|
||||||
`[schedule][开始执行任务] 参数 ${JSON.stringify({
|
`[schedule][开始执行任务] 参数 ${JSON.stringify({
|
||||||
...cron,
|
...cron,
|
||||||
@@ -94,7 +104,8 @@ export function runCron(cmd: string, cron: ICron): Promise<number | void> {
|
|||||||
asError(error).message,
|
asError(error).message,
|
||||||
);
|
);
|
||||||
} finally {
|
} finally {
|
||||||
taskLimit.removeQueuedCron(cron.id);
|
taskLimit.removeQueuedCron(cron.id, execute);
|
||||||
}
|
}
|
||||||
});
|
};
|
||||||
|
return taskLimit.runWithCronLimit(cron, execute, undefined, isCurrent);
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -31,7 +31,7 @@ test('Shell and CLI agree on isolated script state and account modes', async (t)
|
|||||||
path.join(root, 'static/build'),
|
path.join(root, 'static/build'),
|
||||||
])
|
])
|
||||||
await fs.mkdir(directory, { recursive: true });
|
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(
|
await fs.copyFile(
|
||||||
path.resolve(__dirname, '../../../shell', name),
|
path.resolve(__dirname, '../../../shell', name),
|
||||||
path.join(context.paths.dir_shell, name),
|
path.join(context.paths.dir_shell, name),
|
||||||
@@ -355,6 +355,7 @@ t() { :; }
|
|||||||
sleep() { printf DELAYED; }
|
sleep() { printf DELAYED; }
|
||||||
enter_script_workdir() { :; }
|
enter_script_workdir() { :; }
|
||||||
run_else() { :; }
|
run_else() { :; }
|
||||||
|
run_task_command() { "$@"; }
|
||||||
which_program=true
|
which_program=true
|
||||||
main "$@"
|
main "$@"
|
||||||
`,
|
`,
|
||||||
|
|||||||
@@ -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 }));
|
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'])
|
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 });
|
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));
|
fs.copyFileSync(path.resolve(__dirname, '../../../shell', name), path.join(root, 'shell', name));
|
||||||
for (const name of ['zh.sh', 'en.sh'])
|
for (const name of ['zh.sh', 'en.sh'])
|
||||||
fs.copyFileSync(path.resolve(__dirname, '../../../shell/lang', name), path.join(root, 'shell/lang', name));
|
fs.copyFileSync(path.resolve(__dirname, '../../../shell/lang', name), path.join(root, 'shell/lang', name));
|
||||||
|
|||||||
@@ -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 }));
|
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'])
|
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 });
|
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));
|
fs.copyFileSync(path.resolve(__dirname, '../../../shell', name), path.join(root, 'shell', name));
|
||||||
for (const name of ['zh.sh', 'en.sh'])
|
for (const name of ['zh.sh', 'en.sh'])
|
||||||
fs.copyFileSync(path.resolve(__dirname, '../../../shell/lang', name), path.join(root, 'shell/lang', name));
|
fs.copyFileSync(path.resolve(__dirname, '../../../shell/lang', name), path.join(root, 'shell/lang', name));
|
||||||
|
|||||||
@@ -59,3 +59,20 @@ resolved.
|
|||||||
Provider initialization, including provider configuration reads, now happens at
|
Provider initialization, including provider configuration reads, now happens at
|
||||||
first use. A script's explicit `import __ql_notify__`, `require(...)`, or independent
|
first use. A script's explicit `import __ql_notify__`, `require(...)`, or independent
|
||||||
notification helper still loads that module immediately.
|
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.
|
||||||
|
|||||||
+4
-8
@@ -164,7 +164,7 @@ run_normal() {
|
|||||||
if [[ $isJsOrPythonFile == 'false' ]]; then
|
if [[ $isJsOrPythonFile == 'false' ]]; then
|
||||||
clear_non_sh_env
|
clear_non_sh_env
|
||||||
fi
|
fi
|
||||||
$timeoutCmd $which_program $file_param "${script_params[@]}"
|
run_task_command $which_program $file_param "${script_params[@]}"
|
||||||
}
|
}
|
||||||
|
|
||||||
handle_env_split() {
|
handle_env_split() {
|
||||||
@@ -203,7 +203,7 @@ run_concurrent() {
|
|||||||
export "${env_param}=${array[$i - 1]}"
|
export "${env_param}=${array[$i - 1]}"
|
||||||
clear_non_sh_env
|
clear_non_sh_env
|
||||||
fi
|
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
|
done
|
||||||
|
|
||||||
wait
|
wait
|
||||||
@@ -243,7 +243,7 @@ run_designated() {
|
|||||||
|
|
||||||
enter_script_workdir
|
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
|
fi
|
||||||
|
|
||||||
clear_non_sh_env
|
clear_non_sh_env
|
||||||
$timeoutCmd $which_program $file_param "$@"
|
run_task_command $which_program $file_param "$@"
|
||||||
}
|
}
|
||||||
|
|
||||||
check_file() {
|
check_file() {
|
||||||
@@ -336,10 +336,6 @@ check_nounset() {
|
|||||||
|
|
||||||
main() {
|
main() {
|
||||||
if [[ $1 == *.js ]] || [[ $1 == *.mjs ]] || [[ $1 == *.py ]] || [[ $1 == *.pyc ]] || [[ $1 == *.sh ]] || [[ $1 == *.ts ]]; then
|
if [[ $1 == *.js ]] || [[ $1 == *.mjs ]] || [[ $1 == *.py ]] || [[ $1 == *.pyc ]] || [[ $1 == *.sh ]] || [[ $1 == *.ts ]]; then
|
||||||
if [[ $1 == *.sh ]]; then
|
|
||||||
timeoutCmd=""
|
|
||||||
fi
|
|
||||||
|
|
||||||
case $# in
|
case $# in
|
||||||
1)
|
1)
|
||||||
run_normal "$1"
|
run_normal "$1"
|
||||||
|
|||||||
@@ -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"
|
||||||
|
)
|
||||||
|
}
|
||||||
+1
-5
@@ -3,6 +3,7 @@
|
|||||||
dir_shell=$QL_DIR/shell
|
dir_shell=$QL_DIR/shell
|
||||||
. $dir_shell/share.sh
|
. $dir_shell/share.sh
|
||||||
. $dir_shell/api.sh
|
. $dir_shell/api.sh
|
||||||
|
. $dir_shell/task-timeout.sh
|
||||||
|
|
||||||
trap 'single_hanle SIGINT' INT
|
trap 'single_hanle SIGINT' INT
|
||||||
trap 'single_hanle SIGTERM' TERM
|
trap 'single_hanle SIGTERM' TERM
|
||||||
@@ -112,11 +113,6 @@ format_params() {
|
|||||||
mtime_format="%Y-%m-%d %H:%M:%S.%3N"
|
mtime_format="%Y-%m-%d %H:%M:%S.%3N"
|
||||||
fi
|
fi
|
||||||
timeoutCmd=""
|
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')
|
# params=$(echo "$@" | sed -E 's/([^ ])&([^ ])/\1\\\&\2/g')
|
||||||
|
|
||||||
# 分割 task 内置参数和脚本参数
|
# 分割 task 内置参数和脚本参数
|
||||||
|
|||||||
@@ -36,7 +36,7 @@ test(
|
|||||||
},
|
},
|
||||||
'../loaders/logger': logger,
|
'../loaders/logger': logger,
|
||||||
'../data/cron': {
|
'../data/cron': {
|
||||||
CrontabModel: { findOne: async () => null },
|
CrontabModel: { findOne: async () => ({ isDisabled: 0 }) },
|
||||||
CrontabStatus: {},
|
CrontabStatus: {},
|
||||||
},
|
},
|
||||||
'../data/runningInstance': {
|
'../data/runningInstance': {
|
||||||
|
|||||||
@@ -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);
|
||||||
|
});
|
||||||
@@ -12,6 +12,7 @@ function extract(file, name) {
|
|||||||
return text.slice(start, text.indexOf('\n}', start) + 2);
|
return text.slice(start, text.indexOf('\n}', start) + 2);
|
||||||
}
|
}
|
||||||
const helpers = [
|
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)),
|
...['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)),
|
...['run_shell_script', 'define_program', 'format_params'].map(n => extract('shell/task.sh', n)),
|
||||||
].join('\n');
|
].join('\n');
|
||||||
@@ -24,7 +25,7 @@ function fixture(t) {
|
|||||||
write('env.sh', 'export QA_PANEL="alpha&beta&gamma"\n');
|
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('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');
|
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 + `
|
const r = spawnSync('/bin/bash', ['-c', helpers + `
|
||||||
dir_scripts=$QA_ROOT; dir_shell=$QA_ROOT; dir_dep=$QA_ROOT
|
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
|
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[@]}'}"
|
format_params "$@"; define_program "${'${task_shell_params[@]}'}"
|
||||||
. "$QA_TASK_SOURCE"
|
. "$QA_TASK_SOURCE"
|
||||||
printf 'WRAPPER_FINISHED\\n' >> "$QA_ROOT/events"
|
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);
|
assert.equal(r.status, 0, r.stdout + r.stderr);
|
||||||
return { stdout: r.stdout, events: fs.readFileSync(path.join(root, 'events'), 'utf8').trim().split('\n') };
|
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);
|
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}`));
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user