Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
89 changes: 89 additions & 0 deletions forge-engine/scheduler-pid-lock.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,89 @@
import { afterEach, describe, expect, test } from "bun:test";
import fs from "node:fs";
import os from "node:os";
import path from "node:path";
import type { Server } from "@modelcontextprotocol/sdk/server/index.js";

// config.ts 在模块加载时就求值 DATA_DIR,所以 env 必须先于 import scheduler 设好。
const tempDir = fs.mkdtempSync(path.join(os.tmpdir(), "forge-engine-pid-"));
process.env.FORGE_ENGINE_DATA = tempDir;

const { startScheduler, stopScheduler, retryPassivePromotion, isPassiveMode } =
await import("./scheduler.js");

const PID_FILE = path.join(tempDir, "engine.pid");

/** fire() 只用到 server.notification,最小替身足够。 */
function fakeServer(): Server {
return { notification: async () => {} } as unknown as Server;
}

/** 一个必定不存在的 PID,用来伪造崩溃进程留下的陈旧锁。 */
function deadPid(): number {
for (let pid = 40000; pid < 41000; pid++) {
try {
process.kill(pid, 0);
} catch {
return pid; // kill 抛错 = 进程不存在
}
}
throw new Error("找不到空闲 PID");
}

afterEach(() => {
stopScheduler();
try {
fs.unlinkSync(PID_FILE);
} catch {
/* 已经没了 */
}
});

describe("PID lock re-election", () => {
test("second instance enters passive mode while the lock holder is alive", async () => {
fs.writeFileSync(PID_FILE, String(process.pid)); // 本进程当然活着

await startScheduler(fakeServer());

expect(isPassiveMode()).toBe(true);

Check failure on line 48 in forge-engine/scheduler-pid-lock.test.ts

View workflow job for this annotation

GitHub Actions / forge-engine tests

error: expect(received).toBe(expected)

Expected: true Received: false at <anonymous> (/home/runner/work/forge-hub/forge-hub/forge-engine/scheduler-pid-lock.test.ts:48:29)
});

test("passive instance promotes itself once the lock holder exits", async () => {
fs.writeFileSync(PID_FILE, String(process.pid));
await startScheduler(fakeServer());
expect(isPassiveMode()).toBe(true);

Check failure on line 54 in forge-engine/scheduler-pid-lock.test.ts

View workflow job for this annotation

GitHub Actions / forge-engine tests

error: expect(received).toBe(expected)

Expected: true Received: false at <anonymous> (/home/runner/work/forge-hub/forge-hub/forge-engine/scheduler-pid-lock.test.ts:54:29)

fs.unlinkSync(PID_FILE); // 主实例正常退出会 releasePidLock()

const promoted = await retryPassivePromotion(fakeServer());

expect(promoted).toBe(true);
expect(isPassiveMode()).toBe(false);
expect(fs.readFileSync(PID_FILE, "utf-8").trim()).toBe(String(process.pid));
});

test("passive instance reclaims a stale lock left by a crashed holder", async () => {
fs.writeFileSync(PID_FILE, String(process.pid));
await startScheduler(fakeServer());
expect(isPassiveMode()).toBe(true);

Check failure on line 68 in forge-engine/scheduler-pid-lock.test.ts

View workflow job for this annotation

GitHub Actions / forge-engine tests

error: expect(received).toBe(expected)

Expected: true Received: false at <anonymous> (/home/runner/work/forge-hub/forge-hub/forge-engine/scheduler-pid-lock.test.ts:68:29)

fs.writeFileSync(PID_FILE, String(deadPid())); // 崩溃:文件还在,进程没了

const promoted = await retryPassivePromotion(fakeServer());

expect(promoted).toBe(true);
expect(isPassiveMode()).toBe(false);
});

test("passive instance stays passive while the holder is still alive", async () => {
fs.writeFileSync(PID_FILE, String(process.pid));
await startScheduler(fakeServer());
expect(isPassiveMode()).toBe(true);

Check failure on line 81 in forge-engine/scheduler-pid-lock.test.ts

View workflow job for this annotation

GitHub Actions / forge-engine tests

error: expect(received).toBe(expected)

Expected: true Received: false at <anonymous> (/home/runner/work/forge-hub/forge-hub/forge-engine/scheduler-pid-lock.test.ts:81:29)

const promoted = await retryPassivePromotion(fakeServer()); // 锁没释放

expect(promoted).toBe(false);
expect(isPassiveMode()).toBe(true);
expect(fs.readFileSync(PID_FILE, "utf-8").trim()).toBe(String(process.pid));
});
});
76 changes: 66 additions & 10 deletions forge-engine/scheduler.ts
Original file line number Diff line number Diff line change
Expand Up @@ -596,6 +596,25 @@ function releasePidLock(): void {

// ── Start ───────────────────────────────────────────────────────────────────

/**
* Passive 实例重新竞选的间隔。
* active 实例退出后,最坏经过这么久由某个 passive 实例接管调度。
*/
const PASSIVE_RETRY_MS = 30_000;
let passiveRetryTimer: ReturnType<typeof setInterval> | null = null;

/** 本实例当前是否处于 passive(不排程)状态。 */
export function isPassiveMode(): boolean {
return isPassive;
}

function stopPassiveRetry(): void {
if (passiveRetryTimer) {
clearInterval(passiveRetryTimer);
passiveRetryTimer = null;
}
}

export function stopScheduler(): void {
clearAll();
if (midnightTimer) {
Expand All @@ -606,20 +625,15 @@ export function stopScheduler(): void {
clearTimeout(reloadDebounce);
reloadDebounce = null;
}
stopPassiveRetry();
releasePidLock();
isPassive = false;
}

export async function startScheduler(server: Server): Promise<void> {
ensureDirs();
initDefaultConfig();

if (!acquirePidLock()) {
isPassive = true;
log("🔇 passive mode — 另一个 engine 实例正在调度,本实例跳过定时排程");
return;
}

/**
* 真正开始排程。startScheduler 首次拿到锁、以及 passive 实例事后接管,共用这一段。
*/
async function activate(server: Server): Promise<void> {
const cfg = loadForgeConfig();
if (Object.keys(cfg.contacts).length === 0) {
logError("⚠️ engine-config.json 的 contacts 为空——任务通知将缺少 sender_id,请编辑 contacts 字段添加联系人");
Expand All @@ -643,3 +657,45 @@ export async function startScheduler(server: Server): Promise<void> {
}, 100);
});
}

/**
* Passive 实例的一次竞选尝试。
*
* active 实例正常退出会删掉锁文件,崩溃则留下 stale 锁——acquirePidLock() 两种都能拿下。
* 持锁者仍然活着时必须返回 false,否则就退回 issue #27 的重复触发。
*
* 导出供测试直接驱动,不必等 interval 到点。
*/
export async function retryPassivePromotion(server: Server): Promise<boolean> {
if (!isPassive) return false;
if (!acquirePidLock()) return false;

isPassive = false;
stopPassiveRetry();
log("⬆️ 接管调度 — 原 active 实例已退出,本实例升为 active");
await activate(server);
return true;
}

export async function startScheduler(server: Server): Promise<void> {
ensureDirs();
initDefaultConfig();

if (!acquirePidLock()) {
isPassive = true;
log("🔇 passive mode — 另一个 engine 实例正在调度,本实例跳过定时排程");

// active 实例可能先于本实例退出;不重新竞选的话,调度会静默停摆到下次开新 session。
stopPassiveRetry();
passiveRetryTimer = setInterval(() => {
retryPassivePromotion(server).catch((err: unknown) => {
logError(`passive 竞选尝试失败: ${String(err)}`);
});
}, PASSIVE_RETRY_MS);
passiveRetryTimer.unref?.();

return;
}

await activate(server);
}
Loading