mirror of
https://github.com/whyour/qinglong.git
synced 2026-08-12 03:10:48 +08:00
Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
b33d994316 | ||
|
|
6ede8139ce | ||
|
|
95939bbea5 | ||
|
|
2baf352350 | ||
|
|
827453986b | ||
|
|
15f4bcf363 | ||
|
|
d8ea840266 |
+71
-7
@@ -24,6 +24,7 @@ class Application {
|
|||||||
private grpcServerService?: GrpcServerService;
|
private grpcServerService?: GrpcServerService;
|
||||||
private isShuttingDown = false;
|
private isShuttingDown = false;
|
||||||
private workerMetadataMap = new Map<number, WorkerMetadata>();
|
private workerMetadataMap = new Map<number, WorkerMetadata>();
|
||||||
|
private httpWorker?: Worker;
|
||||||
|
|
||||||
constructor() {
|
constructor() {
|
||||||
this.app = express();
|
this.app = express();
|
||||||
@@ -53,8 +54,19 @@ class Application {
|
|||||||
}
|
}
|
||||||
|
|
||||||
private startMasterProcess() {
|
private startMasterProcess() {
|
||||||
this.forkWorker('http');
|
// Fork gRPC worker first and wait for it to be ready
|
||||||
this.forkWorker('grpc');
|
const grpcWorker = this.forkWorker('grpc');
|
||||||
|
|
||||||
|
// Wait for gRPC worker to signal it's ready before starting HTTP worker
|
||||||
|
this.waitForWorkerReady(grpcWorker, 30000)
|
||||||
|
.then(() => {
|
||||||
|
Logger.info('gRPC worker is ready, starting HTTP worker');
|
||||||
|
this.httpWorker = this.forkWorker('http');
|
||||||
|
})
|
||||||
|
.catch((error) => {
|
||||||
|
Logger.error('Failed to wait for gRPC worker:', error);
|
||||||
|
process.exit(1);
|
||||||
|
});
|
||||||
|
|
||||||
cluster.on('exit', (worker, code, signal) => {
|
cluster.on('exit', (worker, code, signal) => {
|
||||||
const metadata = this.workerMetadataMap.get(worker.id);
|
const metadata = this.workerMetadataMap.get(worker.id);
|
||||||
@@ -64,10 +76,32 @@ class Application {
|
|||||||
`${metadata.serviceType} worker ${worker.process.pid} died (${signal || code
|
`${metadata.serviceType} worker ${worker.process.pid} died (${signal || code
|
||||||
}). Restarting...`,
|
}). Restarting...`,
|
||||||
);
|
);
|
||||||
const newWorker = this.forkWorker(metadata.serviceType);
|
// If gRPC worker died, restart it and wait for it to be ready
|
||||||
Logger.info(
|
if (metadata.serviceType === 'grpc') {
|
||||||
`Restarted ${metadata.serviceType} worker (New PID: ${newWorker.process.pid})`,
|
const newGrpcWorker = this.forkWorker('grpc');
|
||||||
);
|
this.waitForWorkerReady(newGrpcWorker, 30000)
|
||||||
|
.then(() => {
|
||||||
|
Logger.info('gRPC worker restarted and ready');
|
||||||
|
// Re-register cron jobs by notifying the HTTP worker
|
||||||
|
if (this.httpWorker) {
|
||||||
|
try {
|
||||||
|
this.httpWorker.send('reregister-crons');
|
||||||
|
Logger.info('Sent reregister-crons message to HTTP worker');
|
||||||
|
} catch (error) {
|
||||||
|
Logger.error('Failed to send reregister-crons message:', error);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
})
|
||||||
|
.catch((error) => {
|
||||||
|
Logger.error('Failed to restart gRPC worker:', error);
|
||||||
|
process.exit(1);
|
||||||
|
});
|
||||||
|
} else {
|
||||||
|
// For HTTP worker, just restart it
|
||||||
|
const newWorker = this.forkWorker(metadata.serviceType);
|
||||||
|
this.httpWorker = newWorker;
|
||||||
|
Logger.info(`Restarted ${metadata.serviceType} worker (PID: ${newWorker.process.pid})`);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
this.workerMetadataMap.delete(worker.id);
|
this.workerMetadataMap.delete(worker.id);
|
||||||
@@ -77,6 +111,25 @@ class Application {
|
|||||||
this.setupMasterShutdown();
|
this.setupMasterShutdown();
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private waitForWorkerReady(worker: Worker, timeoutMs: number): Promise<void> {
|
||||||
|
return new Promise<void>((resolve, reject) => {
|
||||||
|
const messageHandler = (msg: any) => {
|
||||||
|
if (msg === 'ready') {
|
||||||
|
worker.removeListener('message', messageHandler);
|
||||||
|
clearTimeout(timeoutId);
|
||||||
|
resolve();
|
||||||
|
}
|
||||||
|
};
|
||||||
|
worker.on('message', messageHandler);
|
||||||
|
|
||||||
|
// Timeout after specified milliseconds
|
||||||
|
const timeoutId = setTimeout(() => {
|
||||||
|
worker.removeListener('message', messageHandler);
|
||||||
|
reject(new Error(`Worker failed to start within ${timeoutMs / 1000} seconds`));
|
||||||
|
}, timeoutMs);
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
private forkWorker(serviceType: string): Worker {
|
private forkWorker(serviceType: string): Worker {
|
||||||
const worker = cluster.fork({ SERVICE_TYPE: serviceType });
|
const worker = cluster.fork({ SERVICE_TYPE: serviceType });
|
||||||
|
|
||||||
@@ -206,9 +259,20 @@ class Application {
|
|||||||
}
|
}
|
||||||
|
|
||||||
private setupWorkerShutdown(serviceType: string) {
|
private setupWorkerShutdown(serviceType: string) {
|
||||||
process.on('message', (msg) => {
|
process.on('message', async (msg) => {
|
||||||
if (msg === 'shutdown') {
|
if (msg === 'shutdown') {
|
||||||
this.gracefulShutdown(serviceType);
|
this.gracefulShutdown(serviceType);
|
||||||
|
} else if (msg === 'reregister-crons' && serviceType === 'http') {
|
||||||
|
// Re-register cron jobs when gRPC worker restarts
|
||||||
|
try {
|
||||||
|
Logger.info('Received reregister-crons message, re-registering cron jobs...');
|
||||||
|
const CronService = (await import('./services/cron')).default;
|
||||||
|
const cronService = Container.get(CronService);
|
||||||
|
await cronService.autosave_crontab();
|
||||||
|
Logger.info('Cron jobs re-registered successfully');
|
||||||
|
} catch (error) {
|
||||||
|
Logger.error('Failed to re-register cron jobs:', error);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
});
|
});
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user