mirror of
https://github.com/whyour/qinglong.git
synced 2026-08-12 03:10:48 +08:00
Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
bc92ba47f6 | ||
|
|
dfb70abbe3 |
@@ -25,7 +25,7 @@ export class GrpcServerService {
|
||||
const grpcPort = config.grpcPort;
|
||||
const bindAsync = promisify(this.server.bindAsync).bind(this.server);
|
||||
await bindAsync(
|
||||
`0.0.0.0:${grpcPort}`,
|
||||
`[::]:${grpcPort}`,
|
||||
ServerCredentials.createInsecure(),
|
||||
);
|
||||
Logger.debug(`✌️ gRPC service started successfully`);
|
||||
|
||||
@@ -11,7 +11,7 @@ export class HttpServerService {
|
||||
async initialize(expressApp: express.Application, port: number) {
|
||||
try {
|
||||
return new Promise((resolve, reject) => {
|
||||
this.server = expressApp.listen(port, '0.0.0.0', () => {
|
||||
this.server = expressApp.listen(port, '::', () => {
|
||||
Logger.debug(`✌️ HTTP service started successfully`);
|
||||
metricsService.record('http_service_start', 1, {
|
||||
port: port.toString(),
|
||||
|
||||
+2
-28
@@ -14,7 +14,6 @@ import {
|
||||
import config from '../config';
|
||||
import { credentials } from '@grpc/grpc-js';
|
||||
import { ApiClient } from '../protos/api';
|
||||
import { CrontabModel } from '../data/cron';
|
||||
|
||||
class TaskLimit {
|
||||
private dependenyLimit = new PQueue({ concurrency: 1 });
|
||||
@@ -132,38 +131,13 @@ class TaskLimit {
|
||||
let runs = this.queuedCrons.get(cron.id);
|
||||
const result = runs?.length ? [...runs, fn] : [fn];
|
||||
const repeatTimes = this.repeatCronNotifyMap.get(cron.id) || 0;
|
||||
|
||||
// Check instance mode from database to determine queue limit
|
||||
let maxQueueSize = 10; // Default for multi-instance mode (increased from 5)
|
||||
let isSingleInstanceMode = false;
|
||||
try {
|
||||
const cronRecord = await CrontabModel.findOne({
|
||||
where: { id: Number(cron.id) },
|
||||
});
|
||||
|
||||
// Default to single instance mode (0) for backward compatibility
|
||||
// allow_multiple_instances is 1 for multi-instance, 0 or null/undefined for single instance
|
||||
isSingleInstanceMode = cronRecord?.allow_multiple_instances !== 1;
|
||||
|
||||
if (isSingleInstanceMode) {
|
||||
// For single instance mode, allow up to 2 queued tasks
|
||||
// This allows the new task to be queued while the old one is being killed
|
||||
maxQueueSize = 2;
|
||||
}
|
||||
} catch (error) {
|
||||
Logger.error(
|
||||
`[schedule][检查实例模式失败] 任务ID: ${cron.id}, 错误: ${error}`,
|
||||
);
|
||||
}
|
||||
|
||||
if (result?.length > maxQueueSize) {
|
||||
if (result?.length > 5) {
|
||||
if (repeatTimes < 3) {
|
||||
this.repeatCronNotifyMap.set(cron.id, repeatTimes + 1);
|
||||
const modeStr = isSingleInstanceMode ? '单实例' : '多实例';
|
||||
this.client.systemNotify(
|
||||
{
|
||||
title: '任务重复运行',
|
||||
content: `任务:${cron.name}(${modeStr}模式),命令:${cron.command},定时:${cron.schedule},处于运行中的超过 ${maxQueueSize} 个,请检查定时设置`,
|
||||
content: `任务:${cron.name},命令:${cron.command},定时:${cron.schedule},处于运行中的超过 5 个,请检查定时设置`,
|
||||
},
|
||||
(err, res) => {
|
||||
if (err) {
|
||||
|
||||
+3
-27
@@ -15,12 +15,11 @@ export function runCron(cmd: string, cron: ICron): Promise<number | void> {
|
||||
});
|
||||
|
||||
// Default to single instance mode (0) for backward compatibility
|
||||
// allow_multiple_instances is 1 for multi-instance, 0 or null/undefined for single instance
|
||||
const isSingleInstanceMode =
|
||||
existingCron?.allow_multiple_instances !== 1;
|
||||
const allowSingleInstances =
|
||||
existingCron?.allow_multiple_instances === 0;
|
||||
|
||||
if (
|
||||
isSingleInstanceMode &&
|
||||
allowSingleInstances &&
|
||||
existingCron &&
|
||||
existingCron.pid &&
|
||||
(existingCron.status === CrontabStatus.running ||
|
||||
@@ -50,18 +49,6 @@ export function runCron(cmd: string, cron: ICron): Promise<number | void> {
|
||||
);
|
||||
const cp = spawn(cmd, { shell: '/bin/bash' });
|
||||
|
||||
// Update status to running after spawning the process
|
||||
try {
|
||||
await CrontabModel.update(
|
||||
{ status: CrontabStatus.running, pid: cp.pid },
|
||||
{ where: { id: Number(cron.id) } },
|
||||
);
|
||||
} catch (error) {
|
||||
Logger.error(
|
||||
`[schedule][更新任务状态失败] 任务ID: ${cron.id}, 错误: ${error}`,
|
||||
);
|
||||
}
|
||||
|
||||
cp.stderr.on('data', (data) => {
|
||||
Logger.info(
|
||||
'[schedule][执行任务失败] 命令: %s, 错误信息: %j',
|
||||
@@ -79,17 +66,6 @@ export function runCron(cmd: string, cron: ICron): Promise<number | void> {
|
||||
|
||||
cp.on('exit', async (code) => {
|
||||
taskLimit.removeQueuedCron(cron.id);
|
||||
// Update status to idle after task completes
|
||||
try {
|
||||
await CrontabModel.update(
|
||||
{ status: CrontabStatus.idle, pid: undefined },
|
||||
{ where: { id: Number(cron.id) } },
|
||||
);
|
||||
} catch (error) {
|
||||
Logger.error(
|
||||
`[schedule][更新任务状态失败] 任务ID: ${cron.id}, 错误: ${error}`,
|
||||
);
|
||||
}
|
||||
Logger.info(
|
||||
'[schedule][执行任务结束] 参数: %s, 退出码: %j',
|
||||
JSON.stringify({
|
||||
|
||||
Reference in New Issue
Block a user