fix: Gateway 후속 프로필 큐 정체를 자동 복구
PM2 Node client 세션을 직렬화하고 callback timeout 뒤 다음 orchestrator poll이 계속 진행되게 한다. 운영 판별과 제한 복구 경계 및 회귀 테스트를 함께 기록한다.
This commit is contained in:
@@ -8,29 +8,44 @@ import {
|
||||
type ProcessDefinition,
|
||||
} from './processManager.js';
|
||||
|
||||
type Pm2Module = typeof Pm2;
|
||||
export interface Pm2Client {
|
||||
connect(callback: (error?: Error) => void): void;
|
||||
disconnect(): void;
|
||||
list(callback: (error: Error | null, list?: Pm2.ProcessDescription[]) => void): void;
|
||||
start(options: Pm2.StartOptions, callback: (error?: Error) => void): void;
|
||||
stop(name: string, callback: (error?: Error) => void): void;
|
||||
delete(name: string, callback: (error?: Error) => void): void;
|
||||
}
|
||||
|
||||
export interface Pm2ProcessManagerOptions {
|
||||
loadPm2?: () => Pm2Client;
|
||||
connectTimeoutMs?: number;
|
||||
listTimeoutMs?: number;
|
||||
mutationTimeoutMs?: number;
|
||||
}
|
||||
|
||||
const require = createRequire(import.meta.url);
|
||||
|
||||
const loadPm2 = (): Pm2Module => require('pm2') as Pm2Module;
|
||||
const loadPm2 = (): Pm2Client => require('pm2') as Pm2Client;
|
||||
const DEFAULT_PM2_CONNECT_TIMEOUT_MS = 5_000;
|
||||
const DEFAULT_PM2_LIST_TIMEOUT_MS = 5_000;
|
||||
const DEFAULT_PM2_MUTATION_TIMEOUT_MS = 30_000;
|
||||
|
||||
const withPm2 = async <T>(handler: (pm2: Pm2Module) => Promise<T>): Promise<T> => {
|
||||
const pm2 = loadPm2();
|
||||
await new Promise<void>((resolve, reject) => {
|
||||
pm2.connect((error) => {
|
||||
if (error) {
|
||||
const withTimeout = <T>(promise: Promise<T>, timeoutMs: number, label: string): Promise<T> =>
|
||||
new Promise<T>((resolve, reject) => {
|
||||
const timer = setTimeout(() => reject(new Error(`${label} timed out after ${timeoutMs}ms.`)), timeoutMs);
|
||||
timer.unref();
|
||||
promise.then(
|
||||
(value) => {
|
||||
clearTimeout(timer);
|
||||
resolve(value);
|
||||
},
|
||||
(error: unknown) => {
|
||||
clearTimeout(timer);
|
||||
reject(error);
|
||||
return;
|
||||
}
|
||||
resolve();
|
||||
});
|
||||
);
|
||||
});
|
||||
try {
|
||||
return await handler(pm2);
|
||||
} finally {
|
||||
pm2.disconnect();
|
||||
}
|
||||
};
|
||||
|
||||
export const buildPm2StartOptions = (definition: ProcessDefinition) => ({
|
||||
name: definition.name,
|
||||
@@ -47,8 +62,52 @@ export const buildPm2StartOptions = (definition: ProcessDefinition) => ({
|
||||
});
|
||||
|
||||
export class Pm2ProcessManager implements ProcessManager {
|
||||
private readonly loadPm2: () => Pm2Client;
|
||||
private readonly connectTimeoutMs: number;
|
||||
private readonly listTimeoutMs: number;
|
||||
private readonly mutationTimeoutMs: number;
|
||||
private sessionTail: Promise<void> = Promise.resolve();
|
||||
|
||||
constructor(options: Pm2ProcessManagerOptions = {}) {
|
||||
this.loadPm2 = options.loadPm2 ?? loadPm2;
|
||||
this.connectTimeoutMs = options.connectTimeoutMs ?? DEFAULT_PM2_CONNECT_TIMEOUT_MS;
|
||||
this.listTimeoutMs = options.listTimeoutMs ?? DEFAULT_PM2_LIST_TIMEOUT_MS;
|
||||
this.mutationTimeoutMs = options.mutationTimeoutMs ?? DEFAULT_PM2_MUTATION_TIMEOUT_MS;
|
||||
}
|
||||
|
||||
private withPm2<T>(label: string, timeoutMs: number, handler: (pm2: Pm2Client) => Promise<T>): Promise<T> {
|
||||
const task = this.sessionTail.then(async () => {
|
||||
const pm2 = this.loadPm2();
|
||||
try {
|
||||
await withTimeout(
|
||||
new Promise<void>((resolve, reject) => {
|
||||
pm2.connect((error) => {
|
||||
if (error) {
|
||||
reject(error);
|
||||
return;
|
||||
}
|
||||
resolve();
|
||||
});
|
||||
}),
|
||||
this.connectTimeoutMs,
|
||||
'PM2 connect'
|
||||
);
|
||||
return await withTimeout(handler(pm2), timeoutMs, label);
|
||||
} finally {
|
||||
pm2.disconnect();
|
||||
}
|
||||
});
|
||||
this.sessionTail = task.then(
|
||||
() => undefined,
|
||||
() => undefined
|
||||
);
|
||||
return task;
|
||||
}
|
||||
|
||||
async list(): Promise<ManagedProcessInfo[]> {
|
||||
return withPm2(
|
||||
return this.withPm2(
|
||||
'PM2 list',
|
||||
this.listTimeoutMs,
|
||||
(pm2) =>
|
||||
new Promise<ManagedProcessInfo[]>((resolve, reject) => {
|
||||
pm2.list((error, list) => {
|
||||
@@ -72,7 +131,9 @@ export class Pm2ProcessManager implements ProcessManager {
|
||||
}
|
||||
|
||||
async start(definition: ProcessDefinition): Promise<void> {
|
||||
await withPm2(
|
||||
await this.withPm2(
|
||||
`PM2 start ${definition.name}`,
|
||||
this.mutationTimeoutMs,
|
||||
(pm2) =>
|
||||
new Promise<void>((resolve, reject) => {
|
||||
pm2.list((listError, list) => {
|
||||
@@ -100,7 +161,9 @@ export class Pm2ProcessManager implements ProcessManager {
|
||||
}
|
||||
|
||||
async stop(name: string): Promise<void> {
|
||||
await withPm2(
|
||||
await this.withPm2(
|
||||
`PM2 stop ${name}`,
|
||||
this.mutationTimeoutMs,
|
||||
(pm2) =>
|
||||
new Promise<void>((resolve, reject) => {
|
||||
pm2.stop(name, (error) => {
|
||||
@@ -115,7 +178,9 @@ export class Pm2ProcessManager implements ProcessManager {
|
||||
}
|
||||
|
||||
async delete(name: string): Promise<void> {
|
||||
await withPm2(
|
||||
await this.withPm2(
|
||||
`PM2 delete ${name}`,
|
||||
this.mutationTimeoutMs,
|
||||
(pm2) =>
|
||||
new Promise<void>((resolve, reject) => {
|
||||
pm2.delete(name, (error) => {
|
||||
|
||||
@@ -1,6 +1,10 @@
|
||||
import { describe, expect, it } from 'vitest';
|
||||
import { describe, expect, it, vi } from 'vitest';
|
||||
|
||||
import { buildPm2StartOptions } from '../src/orchestrator/pm2ProcessManager.js';
|
||||
import {
|
||||
buildPm2StartOptions,
|
||||
Pm2ProcessManager,
|
||||
type Pm2Client,
|
||||
} from '../src/orchestrator/pm2ProcessManager.js';
|
||||
|
||||
describe('buildPm2StartOptions', () => {
|
||||
it('enforces bounded restart policy and strips inherited PM2 identity at the PM2 boundary', () => {
|
||||
@@ -56,3 +60,118 @@ describe('buildPm2StartOptions', () => {
|
||||
expect(options.env).not.toHaveProperty('args');
|
||||
});
|
||||
});
|
||||
|
||||
describe('Pm2ProcessManager session recovery', () => {
|
||||
it('serializes concurrent PM2 sessions so one disconnect cannot interrupt another request', async () => {
|
||||
const events: string[] = [];
|
||||
let listCall = 0;
|
||||
let releaseFirstList: (() => void) | undefined;
|
||||
const pm2 = {
|
||||
connect(callback: Parameters<Pm2Client['connect']>[0]) {
|
||||
events.push('connect');
|
||||
callback();
|
||||
},
|
||||
disconnect() {
|
||||
events.push('disconnect');
|
||||
},
|
||||
list(callback: Parameters<Pm2Client['list']>[0]) {
|
||||
listCall += 1;
|
||||
const currentCall = listCall;
|
||||
events.push(`list:${currentCall}`);
|
||||
if (currentCall === 1) {
|
||||
releaseFirstList = () => callback(null, []);
|
||||
return;
|
||||
}
|
||||
callback(null, []);
|
||||
},
|
||||
start() {
|
||||
throw new Error('unused');
|
||||
},
|
||||
stop() {
|
||||
throw new Error('unused');
|
||||
},
|
||||
delete() {
|
||||
throw new Error('unused');
|
||||
},
|
||||
} satisfies Pm2Client;
|
||||
const manager = new Pm2ProcessManager({ loadPm2: () => pm2 });
|
||||
|
||||
const first = manager.list();
|
||||
const second = manager.list();
|
||||
await vi.waitFor(() => expect(events).toEqual(['connect', 'list:1']));
|
||||
|
||||
releaseFirstList?.();
|
||||
await expect(Promise.all([first, second])).resolves.toEqual([[], []]);
|
||||
expect(events).toEqual(['connect', 'list:1', 'disconnect', 'connect', 'list:2', 'disconnect']);
|
||||
});
|
||||
|
||||
it('times out a lost PM2 callback and lets the next queued session proceed', async () => {
|
||||
let listCall = 0;
|
||||
const pm2 = {
|
||||
connect(callback: Parameters<Pm2Client['connect']>[0]) {
|
||||
callback();
|
||||
},
|
||||
disconnect() {},
|
||||
list(callback: Parameters<Pm2Client['list']>[0]) {
|
||||
listCall += 1;
|
||||
if (listCall === 1) {
|
||||
return;
|
||||
}
|
||||
callback(null, []);
|
||||
},
|
||||
start() {
|
||||
throw new Error('unused');
|
||||
},
|
||||
stop() {
|
||||
throw new Error('unused');
|
||||
},
|
||||
delete() {
|
||||
throw new Error('unused');
|
||||
},
|
||||
} satisfies Pm2Client;
|
||||
const manager = new Pm2ProcessManager({
|
||||
loadPm2: () => pm2,
|
||||
listTimeoutMs: 10,
|
||||
});
|
||||
|
||||
await expect(manager.list()).rejects.toThrow('PM2 list timed out after 10ms.');
|
||||
await expect(manager.list()).resolves.toEqual([]);
|
||||
});
|
||||
|
||||
it('disconnects a timed-out PM2 connection before releasing the serialized session', async () => {
|
||||
let connectCall = 0;
|
||||
let disconnectCall = 0;
|
||||
const pm2 = {
|
||||
connect(callback: Parameters<Pm2Client['connect']>[0]) {
|
||||
connectCall += 1;
|
||||
if (connectCall > 1) {
|
||||
callback();
|
||||
}
|
||||
},
|
||||
disconnect() {
|
||||
disconnectCall += 1;
|
||||
},
|
||||
list(callback: Parameters<Pm2Client['list']>[0]) {
|
||||
callback(null, []);
|
||||
},
|
||||
start() {
|
||||
throw new Error('unused');
|
||||
},
|
||||
stop() {
|
||||
throw new Error('unused');
|
||||
},
|
||||
delete() {
|
||||
throw new Error('unused');
|
||||
},
|
||||
} satisfies Pm2Client;
|
||||
const manager = new Pm2ProcessManager({
|
||||
loadPm2: () => pm2,
|
||||
connectTimeoutMs: 10,
|
||||
});
|
||||
|
||||
await expect(manager.list()).rejects.toThrow('PM2 connect timed out after 10ms.');
|
||||
expect(disconnectCall).toBe(1);
|
||||
await expect(manager.list()).resolves.toEqual([]);
|
||||
expect(disconnectCall).toBe(2);
|
||||
});
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user