diff --git a/cli.test.ts b/cli.test.ts index 000fa49..49268f3 100644 --- a/cli.test.ts +++ b/cli.test.ts @@ -3,7 +3,7 @@ import fs from "node:fs"; import os from "node:os"; import path from "node:path"; -import { cleanDirContents, HUB_INSTALL_PRESERVE_ENTRIES, syncApiTokenFileAt } from "./cli.js"; +import { cleanDirContents, HUB_INSTALL_PRESERVE_ENTRIES, syncApiTokenFileAt, waitUntil } from "./cli.js"; let tmpDir: string; @@ -75,3 +75,35 @@ describe("syncApiTokenFileAt", () => { expect(fs.existsSync(attackerTarget)).toBe(false); }); }); + +describe("waitUntil", () => { + test("returns true immediately when the predicate already holds", async () => { + let calls = 0; + const ok = await waitUntil(() => { calls += 1; return true; }, 1000, 10); + + expect(ok).toBe(true); + expect(calls).toBe(1); // 已经成立就不该再轮询 + }); + + test("returns true once the predicate flips partway through", async () => { + let n = 0; + const ok = await waitUntil(() => (n += 1) >= 3, 1000, 5); + + expect(ok).toBe(true); + expect(n).toBe(3); + }); + + test("returns false when the predicate never holds before the timeout", async () => { + const ok = await waitUntil(() => false, 60, 10); + + expect(ok).toBe(false); // 超时必须能报失败——否则调用方会把"没等到"当成"成功了" + }); + + test("awaits async predicates", async () => { + let n = 0; + const ok = await waitUntil(async () => (n += 1) >= 2, 1000, 5); + + expect(ok).toBe(true); + expect(n).toBe(2); + }); +}); diff --git a/cli.ts b/cli.ts index a2fd092..3463bef 100755 --- a/cli.ts +++ b/cli.ts @@ -305,9 +305,29 @@ function installCmd(): void { `); } +/** + * 轮询 predicate 直到为真或超时,返回是否在超时前成立。 + * + * 存在的理由:`launchctl bootout` / `bootstrap` 的退出码只说明命令被接受了, + * 不说明进程真的退出或真的起来了。把"发了命令"和"状态真的变了"分开, + * 否则就会出现"报告已重启、实际没换进程"。 + */ +export async function waitUntil( + predicate: () => boolean | Promise, + timeoutMs: number, + intervalMs = 250, +): Promise { + const deadline = Date.now() + timeoutMs; + for (;;) { + if (await predicate()) return true; + if (Date.now() >= deadline) return false; + await new Promise((resolve) => setTimeout(resolve, intervalMs)); + } +} + // ── sync ──────────────────────────────────────────────────────────────────── -function syncCmd(): void { +async function syncCmd(): Promise { log("🔄 Forge Hub sync\n"); // 1. Re-stage package snapshot from current source @@ -336,14 +356,45 @@ function syncCmd(): void { const uid = os.userInfo().uid; const domain = `gui/${uid}`; const label = `${domain}/com.forge-hub`; + const baseUrl = process.env.FORGE_HUB_URL ?? "http://localhost:9900"; + + // 有任何 HTTP 响应就算活着——状态码不重要,能应答就说明端口后面有 Hub + const hubResponding = async (): Promise => { + try { + await fetch(baseUrl, { signal: AbortSignal.timeout(1500) }); + return true; + } catch { + return false; + } + }; + try { execFileSync("launchctl", ["bootout", label], { stdio: "ignore" }); } catch { /* might not be bootstrapped */ } + + // bootout 的退出码只说明命令被接受,不说明进程已经退出。必须等它真的下线: + // 否则紧接着的 bootstrap 会撞上仍占着端口的旧进程,launchd 每 ThrottleInterval + // 崩一次,而这里照样打印"已重启"。 + if (!(await waitUntil(async () => !(await hubResponding()), 15_000, 500))) { + log("⚠️ 旧 Hub 仍在响应,bootout 没能让它下线——它多半不是 launchd 启动的"); + log(" (hub-client 在 Hub 不可达时会 detached spawn 一个,那种进程不是 launchd 的子进程,bootout 管不到)"); + log(" 运行时文件已同步到磁盘,但跑着的仍是旧代码。手动处理:"); + log(" lsof -nP -iTCP:9900 -sTCP:LISTEN # 找出占用者"); + log(` kill && launchctl bootstrap ${domain} ${LAUNCHD_PLIST}`); + return; + } + try { execFileSync("launchctl", ["bootstrap", domain, LAUNCHD_PLIST], { stdio: "inherit" }); - log("✓ Hub 已重启"); } catch { log(`⚠️ 无法重启 Hub。手动执行:launchctl bootout ${label} && launchctl bootstrap ${domain} ${LAUNCHD_PLIST}`); + return; + } + + if (await waitUntil(hubResponding, 20_000, 500)) { + log("✓ Hub 已重启(已验证重新响应)"); + } else { + log("⚠️ bootstrap 已提交,但 Hub 在 20s 内没有恢复响应——查 ~/.forge-hub/hub-stderr.log"); } } @@ -751,7 +802,7 @@ if (import.meta.main) { installCmd(); break; case "sync": - syncCmd(); + await syncCmd(); break; case "uninstall": uninstallCmd(); diff --git a/forge-engine/README.md b/forge-engine/README.md index d4cd448..a0c8522 100644 --- a/forge-engine/README.md +++ b/forge-engine/README.md @@ -7,6 +7,20 @@ > [!IMPORTANT] > Forge Engine 目前是 **experimental / manual setup**。源码、MCP server 和 CLI 都在仓库里,但 **`forge-hub install` 默认不会部署或注册它**。想用的话,按下面步骤单独配置。 +## 运行日志与崩溃可见性 + +`log()` / `logError()` 除了写 stderr,**同时追加到 `/engine.log`**,超过 2MB 转存 `engine.log.1`(只留一份)。 + +**为什么需要文件日志**:stderr 是给 MCP 宿主看的,但宿主通常只在**连接建立那一刻**捕获 stderr —— engine 启动之后打印的任何东西都不再被记录。一旦进程异常退出,什么线索都不剩。 + +实测案例(2026-09-03):engine 在 00:00–03:00 之间消失,宿主日志无断开记录、系统无崩溃报告、无内存压力事件,调度静默停摆 **10 小时 31 分**、漏跑 12 个任务。事后无从判断死因。 + +**崩溃捕获**:`uncaughtException` 与 `unhandledRejection` 会把完整堆栈写进 `engine.log`,随后 `stopScheduler()` 释放 PID 锁并退出(exit 1)。 + +**为什么退出而不是继续跑**:调度器状态可能已损坏,而一个状态损坏的调度器发出的推送比不推送更坏。 + +**为什么不自动重起**:多实例 + 自动重起 + 抢 PID 锁是本项目已经踩过的坑(`cleanOrphans` + 一次性 `acquirePidLock` 的组合)。恢复由**外部**检测触发人工重连,engine 自己不做进程管理。 + ## 架构 ``` diff --git a/forge-engine/config.ts b/forge-engine/config.ts index eeb7906..26d3400 100644 --- a/forge-engine/config.ts +++ b/forge-engine/config.ts @@ -2,6 +2,7 @@ * Forge Engine 路径常量与日志 */ +import fs from "node:fs"; import path from "node:path"; // ── Channel Identity ──────────────────────────────────────────────────────── @@ -25,13 +26,55 @@ export const HANDLERS_DIR = path.resolve(CODE_DIR, "handlers"); export const SCHEDULE_FILE = path.join(DATA_DIR, "engine-schedule.json"); export const ACTION_LOG_FILE = path.join(DATA_DIR, "engine-trigger-log.md"); export const PID_FILE = path.join(DATA_DIR, "engine.pid"); +export const RUNTIME_LOG_FILE = path.join(DATA_DIR, "engine.log"); -// ── Logging (stderr — stdout is MCP stdio) ────────────────────────────────── +// ── Logging ───────────────────────────────────────────────────────────────── +// +// stderr 是给宿主看的(stdout 被 MCP stdio 占用)。但宿主只在**连接建立那一刻** +// 捕获 stderr —— engine 启动之后打印的任何东西都掉进黑洞。 +// +// 后果实测(2026-09-03):engine 在 00:00–03:00 之间死亡,无崩溃报告、无内存压力、 +// 宿主日志里一条断开记录都没有,调度静默停摆 10 小时 31 分、漏跑 12 个任务。 +// 它死之前很可能打印过原因,只是没人听见。 +// +// 所以 log/logError 同时落盘。文件是唯一能在进程死后还留下痕迹的地方。 + +const MAX_LOG_BYTES = 2 * 1024 * 1024; // 2MB 转存一次,只留一份 .1 + +function appendToFile(line: string): void { + try { + fs.mkdirSync(DATA_DIR, { recursive: true }); + // 轮转:日志本身不能变成下一个「不停长大且没人敢删」的文件 + try { + if (fs.statSync(RUNTIME_LOG_FILE).size > MAX_LOG_BYTES) { + fs.renameSync(RUNTIME_LOG_FILE, RUNTIME_LOG_FILE + ".1"); + } + } catch { /* 文件还不存在 —— 首次写入,正常 */ } + fs.appendFileSync(RUNTIME_LOG_FILE, line); + } catch { /* 日志写不进去不能反过来把进程搞死 */ } +} + +function stamp(): string { + return new Date().toISOString().replace("T", " ").slice(0, 19); +} export function log(msg: string) { process.stderr.write(`[engine] ${msg}\n`); + appendToFile(`${stamp()} [engine] ${msg}\n`); } export function logError(msg: string) { process.stderr.write(`[engine] ERROR: ${msg}\n`); + appendToFile(`${stamp()} [engine] ERROR: ${msg}\n`); +} + +/** 致命错误:带完整堆栈落盘。进程即将退出时用。 */ +export function logFatal(kind: string, err: unknown): void { + // err.stack 本身已含 "Name: message" 首行,不要再拼一次 + const detail = err instanceof Error + ? (err.stack ?? `${err.name}: ${err.message}\n(无堆栈)`) + : String(err); + const block = `${stamp()} [engine] FATAL ${kind} · PID ${process.pid}\n${detail}\n`; + process.stderr.write(block); + appendToFile(block); } diff --git a/forge-engine/engine-channel.ts b/forge-engine/engine-channel.ts index 32dab75..66d6844 100644 --- a/forge-engine/engine-channel.ts +++ b/forge-engine/engine-channel.ts @@ -23,6 +23,7 @@ import { SCHEDULE_DIR, log, logError, + logFatal, } from "./config.js"; import { startScheduler, stopScheduler } from "./scheduler.js"; import { resolveTaskTiming } from "./task-timing.js"; @@ -299,7 +300,30 @@ async function main() { process.on("SIGTERM", () => shutdown("SIGTERM")); process.on("SIGINT", () => shutdown("SIGINT")); - log("engine started"); + // ── Crash Visibility ───────────────────────────────────────────────────── + // + // 2026-09-03:engine 在 00:00–03:00 之间消失,宿主日志里没有任何断开记录、 + // 没有崩溃报告、没有内存压力事件,调度静默停摆 10h31m、漏跑 12 个任务。 + // + // 定时器回调里的未捕获异常在 Bun 下会直接带走进程,而 log() 当时只写 stderr, + // 宿主又只在连接建立那一刻捕获 stderr —— **它死之前很可能喊过,只是没人听见。** + // + // 取舍:抓到之后**退出,不带病续跑**。scheduler 的状态可能已经坏了, + // 而一个状态损坏的调度器推送出去的东西,比不推送更坏。 + // 退出后由外部检测(宿主侧的存活检查)发现并提示人工重连 —— 不做自动重起: + // 多实例 + 自动重起 + 抢 PID 锁,正是本项目已经踩过的坑。 + process.on("uncaughtException", (err) => { + logFatal("uncaughtException", err); + stopScheduler(); // 释放 PID 锁,别留一个 stale 锁给下一个实例 + process.exit(1); + }); + process.on("unhandledRejection", (reason) => { + logFatal("unhandledRejection", reason); + stopScheduler(); + process.exit(1); + }); + + log(`engine started · PID ${process.pid}`); } if (import.meta.main) { diff --git a/forge-engine/scheduler-pid-lock.test.ts b/forge-engine/scheduler-pid-lock.test.ts new file mode 100644 index 0000000..85f64b6 --- /dev/null +++ b/forge-engine/scheduler-pid-lock.test.ts @@ -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); + }); + + 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); + + 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); + + 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); + + 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)); + }); +}); diff --git a/forge-engine/scheduler.test.ts b/forge-engine/scheduler.test.ts index 7094723..a19557d 100644 --- a/forge-engine/scheduler.test.ts +++ b/forge-engine/scheduler.test.ts @@ -3,7 +3,14 @@ import fs from "node:fs"; import os from "node:os"; import path from "node:path"; -import { expandRandom, removeScheduleEntryFromFile } from "./scheduler.js"; +import { + expandRandom, + fireKey, + hasFiredToday, + markFired, + removeScheduleEntryFromFile, + shouldFire, +} from "./scheduler.js"; const tempDirs: string[] = []; @@ -92,3 +99,70 @@ describe("scheduler helpers", () => { expect(fs.existsSync(filePath)).toBe(false); }); }); + +describe("shouldFire 时间规则", () => { + // 回归:错过任务检测原本只走 canScheduleToday(仅查 start_date/end_date), + // 未过 shouldFire,导致「每周日」的任务在周六被报成「今天错过」。 + test("weekdays 不匹配当天时返回 false", () => { + const today = new Date().getDay(); + const otherDay = (today + 1) % 7; + expect(shouldFire({ hour: 22, minute: 0, second: 0, sender: "t", weekdays: [otherDay] } as never)).toBe(false); + expect(shouldFire({ hour: 22, minute: 0, second: 0, sender: "t", weekdays: [today] } as never)).toBe(true); + }); + + test("days 不匹配当天日期时返回 false", () => { + const d = new Date().getDate(); + const other = d === 1 ? 2 : 1; + expect(shouldFire({ hour: 9, minute: 0, second: 0, sender: "t", days: [other] } as never)).toBe(false); + expect(shouldFire({ hour: 9, minute: 0, second: 0, sender: "t", days: [d] } as never)).toBe(true); + }); + + test("无时间条件时默认可触发", () => { + expect(shouldFire({ hour: 0, minute: 0, second: 0, sender: "t" } as never)).toBe(true); + }); +}); + +describe("已触发记录(错过检测的第二道闸)", () => { + // 回归:错过判定原本只看「墙上时间 − 排期时刻」落在 2h 窗口内 + shouldFire, + // 从不查该条目今天是否已经触发过。于是任务时刻之后 2 小时内的任何一次重排 + // (启动 / 配置热加载 / 午夜重排在边界上提前几毫秒跑),都会把已跑完的任务 + // 重新报成「错过」。2026-08-22 / 08-23 / 08-28 / 08-29 连续复发。 + const entry = { + hour: 22, minute: 0, second: 0, + sender: "engine", label: "洗澡提醒", origin: "shower.json", + } as never; + + test("fireKey 对同一条目稳定,对不同时刻/来源不同", () => { + expect(fireKey(entry)).toBe(fireKey(entry)); + expect(fireKey({ ...(entry as object), minute: 30 } as never)).not.toBe(fireKey(entry)); + expect(fireKey({ ...(entry as object), origin: "other.json" } as never)).not.toBe(fireKey(entry)); + }); + + test("label 缺省时回退到 sender,不会把两个条目挤成同一个键", () => { + const a = { hour: 9, minute: 0, second: 0, sender: "briefing", origin: "x.json" } as never; + const b = { hour: 9, minute: 0, second: 0, sender: "redline", origin: "x.json" } as never; + expect(fireKey(a)).not.toBe(fireKey(b)); + }); + + test("没有记录时 hasFiredToday 为 false", () => { + expect(hasFiredToday({}, entry, "2026-08-29")).toBe(false); + }); + + test("markFired 之后同日为 true、次日为 false", () => { + const state = markFired({}, entry, "2026-08-29"); + expect(hasFiredToday(state, entry, "2026-08-29")).toBe(true); + expect(hasFiredToday(state, entry, "2026-08-30")).toBe(false); + }); + + test("markFired 清掉非当天的键,状态文件不会无限增长", () => { + let state: Record = { "stale.json|09:00:00|旧任务": "2026-01-01" }; + state = markFired(state, entry, "2026-08-29"); + expect(Object.keys(state)).toEqual([fireKey(entry)]); + }); + + test("同一天重复 markFired 不产生第二个键", () => { + let state = markFired({}, entry, "2026-08-29"); + state = markFired(state, entry, "2026-08-29"); + expect(Object.keys(state)).toHaveLength(1); + }); +}); diff --git a/forge-engine/scheduler.ts b/forge-engine/scheduler.ts index e80e257..0a208a1 100644 --- a/forge-engine/scheduler.ts +++ b/forge-engine/scheduler.ts @@ -147,7 +147,7 @@ function resolveAll(rawEntries: RawScheduleEntry[]): ResolvedEntry[] { /** * 统一时间规则检查。所有条件都满足才触发。 */ -function shouldFire(entry: ResolvedEntry): boolean { +export function shouldFire(entry: ResolvedEntry): boolean { const now = new Date(); if (entry.weekdays?.length && !entry.weekdays.includes(now.getDay())) return false; @@ -276,6 +276,7 @@ async function fire(entry: ResolvedEntry, server: Server): Promise { appendLog(entry, content); updateState(sender); + saveState(FIRES_MODULE, markFired(loadState(FIRES_MODULE), entry)); // Auto-delete one_shot if (entry.one_shot) { @@ -343,13 +344,58 @@ function updateState(sender: string): void { saveState("global", s); } +// ── 已触发记录 ────────────────────────────────────────────────────────────── +// 错过检测原本只看「墙上时间 − 排期时刻」是否落在 MISSED_WINDOW_MS 内(外加 +// shouldFire 的排期规则),从不查该条目今天是否已经触发过。于是任务时刻之后 +// 2 小时内的任何一次重排——启动、配置热加载、午夜重排在边界上提前几毫秒跑到 +// 前一天——都会把已经跑完的任务重新报成「错过」。global.json 只有 last_fire / +// today_count 这类全局标量,回答不了「这一条今天跑没跑」,所以另存一份 per-entry 记录。 + +const FIRES_MODULE = "fires"; + +/** 条目的稳定标识:来源文件 + 时刻 + 名字。名字缺省时回退到 sender。 */ +export function fireKey( + entry: Pick & + Partial>, +): string { + const time = `${pad2(entry.hour)}:${pad2(entry.minute)}:${pad2(entry.second)}`; + return `${entry.origin ?? ""}|${time}|${entry.label ?? entry.sender}`; +} + +/** 该条目今天是否已经触发过。state 由调用方传入,便于一次重排只读一次文件。 */ +export function hasFiredToday( + state: Record, + entry: Parameters[0], + today: string = dateStr(), +): boolean { + return state[fireKey(entry)] === today; +} + +/** 记下该条目今天已触发,并顺手清掉非今天的键(状态文件不随时间增长)。 */ +export function markFired( + state: Record, + entry: Parameters[0], + today: string = dateStr(), +): Record { + const next: Record = {}; + for (const [k, v] of Object.entries(state)) { + if (v === today) next[k] = v; + } + next[fireKey(entry)] = today; + return next; +} + // ── Schedule ──────────────────────────────────────────────────────────────── const MISSED_WINDOW_MS = 2 * 60 * 60 * 1000; +// 刚过点的宽限:排期时刻本身可能正好压在任务时刻上(典型:午夜重排在 00:00 跑, +// 把 00:00 的格子算成"已经过去"),导致该时段任务永远触发不到。 +const FIRE_GRACE_MS = 90 * 1000; function scheduleOrigin(origin: string, entries: ResolvedEntry[], server: Server): number { const now = Date.now(); const today = dateStr(); + const fires = loadState(FIRES_MODULE); const timers: ReturnType[] = []; let count = 0; const missed: { label: string; time: string }[] = []; @@ -363,7 +409,21 @@ function scheduleOrigin(origin: string, entries: ResolvedEntry[], server: Server if (delay > 0) { timers.push(setTimeout(() => fire(entry, server), delay)); count++; - } else if (delay > -MISSED_WINDOW_MS && !entry.one_shot) { + } else if (delay > -FIRE_GRACE_MS && !entry.one_shot && shouldFire(entry)) { + // 刚过点,还在宽限内 → 立即补触发,不算错过。 + // 没有这一支时,排在 00:00 的任务会被午夜重排自己判成已过去,永远跑不到。 + timers.push(setTimeout(() => fire(entry, server), 1000)); + count++; + } else if ( + delay > -MISSED_WINDOW_MS && + !entry.one_shot && + shouldFire(entry) && + !hasFiredToday(fires, entry, today) + ) { + // shouldFire 必须在这里再查一次:canScheduleToday 只看 start_date/end_date, + // 不看 weekdays/days/months,否则"每周日"的任务会在周六被报成今天错过。 + // hasFiredToday 是第二道闸:跑过就不是错过,否则任务时刻后 2h 内的任何一次 + // 重排都会把已完成的任务重报一遍(见上方「已触发记录」段的说明)。 missed.push({ label: entry.label ?? entry.sender, time: timeStr(entry.hour, entry.minute), @@ -596,6 +656,25 @@ function releasePidLock(): void { // ── Start ─────────────────────────────────────────────────────────────────── +/** + * Passive 实例重新竞选的间隔。 + * active 实例退出后,最坏经过这么久由某个 passive 实例接管调度。 + */ +const PASSIVE_RETRY_MS = 30_000; +let passiveRetryTimer: ReturnType | 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) { @@ -606,20 +685,15 @@ export function stopScheduler(): void { clearTimeout(reloadDebounce); reloadDebounce = null; } + stopPassiveRetry(); releasePidLock(); isPassive = false; } -export async function startScheduler(server: Server): Promise { - ensureDirs(); - initDefaultConfig(); - - if (!acquirePidLock()) { - isPassive = true; - log("🔇 passive mode — 另一个 engine 实例正在调度,本实例跳过定时排程"); - return; - } - +/** + * 真正开始排程。startScheduler 首次拿到锁、以及 passive 实例事后接管,共用这一段。 + */ +async function activate(server: Server): Promise { const cfg = loadForgeConfig(); if (Object.keys(cfg.contacts).length === 0) { logError("⚠️ engine-config.json 的 contacts 为空——任务通知将缺少 sender_id,请编辑 contacts 字段添加联系人"); @@ -643,3 +717,45 @@ export async function startScheduler(server: Server): Promise { }, 100); }); } + +/** + * Passive 实例的一次竞选尝试。 + * + * active 实例正常退出会删掉锁文件,崩溃则留下 stale 锁——acquirePidLock() 两种都能拿下。 + * 持锁者仍然活着时必须返回 false,否则就退回 issue #27 的重复触发。 + * + * 导出供测试直接驱动,不必等 interval 到点。 + */ +export async function retryPassivePromotion(server: Server): Promise { + 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 { + 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); +} diff --git a/hub-client/hub-channel.ts b/hub-client/hub-channel.ts index 713f21c..b4c869a 100644 --- a/hub-client/hub-channel.ts +++ b/hub-client/hub-channel.ts @@ -19,7 +19,10 @@ import { import { z } from "zod/v4-mini"; import fs from "node:fs"; import path from "node:path"; -import { spawn } from "node:child_process"; +import { spawn, execFileSync } from "node:child_process"; +import os from "node:os"; + +import { chooseHubStarter, launchdPlistPath, LAUNCHD_LABEL } from "./hub-starter.js"; import { getSessionConfigPaths, isChannelMode, @@ -320,12 +323,37 @@ async function ensureHubRunning(): Promise { // Hub not running, try to start it log("Hub 未运行,尝试自动启动..."); try { - const hubPath = path.join(HUB_DIR, "hub.ts"); - const child = spawn(BUN_BINARY, [hubPath], { - detached: true, - stdio: "ignore", + const spawnDetached = (): void => { + const hubPath = path.join(HUB_DIR, "hub.ts"); + const child = spawn(BUN_BINARY, [hubPath], { + detached: true, + stdio: "ignore", + }); + child.unref(); + }; + + // 装了 launchd plist 就交给 launchd 启动。自己 detached spawn 会造出一个 + // launchd 管不到的孤儿:它占着端口,launchd 的 KeepAlive 每 ThrottleInterval + // 撞一次,而 `forge-hub sync` 的 bootout 也杀不到它——sync 会报"已重启", + // 实际跑的还是旧进程。没装 plist(或非 macOS)时保持原来的直接 spawn。 + const starter = chooseHubStarter({ + platform: process.platform, + plistExists: fs.existsSync(launchdPlistPath(os.homedir())), }); - child.unref(); + + if (starter === "launchd") { + const label = `gui/${process.getuid?.() ?? os.userInfo().uid}/${LAUNCHD_LABEL}`; + try { + execFileSync("launchctl", ["kickstart", label], { stdio: "ignore" }); + log("已请求 launchd 启动 Hub"); + } catch (err) { + // service 没 bootstrap 过时 kickstart 会失败——退回直接启动,至少能用 + logError(`launchctl kickstart 失败,回退到直接启动: ${String(err)}`); + spawnDetached(); + } + } else { + spawnDetached(); + } // Wait for Hub to start for (let i = 0; i < 20; i++) { diff --git a/hub-client/hub-starter.test.ts b/hub-client/hub-starter.test.ts new file mode 100644 index 0000000..7d06f97 --- /dev/null +++ b/hub-client/hub-starter.test.ts @@ -0,0 +1,22 @@ +import { describe, expect, test } from "bun:test"; + +import { chooseHubStarter } from "./hub-starter.js"; + +describe("chooseHubStarter", () => { + test("delegates to launchd on macOS when the plist is installed", () => { + expect(chooseHubStarter({ platform: "darwin", plistExists: true })).toBe("launchd"); + }); + + test("spawns directly on macOS when no plist is installed", () => { + expect(chooseHubStarter({ platform: "darwin", plistExists: false })).toBe("spawn"); + }); + + test("spawns directly on linux even if a plist file happens to be present", () => { + // launchd 是 macOS 专有,别的平台上这个文件没有意义 + expect(chooseHubStarter({ platform: "linux", plistExists: true })).toBe("spawn"); + }); + + test("spawns directly on linux without a plist", () => { + expect(chooseHubStarter({ platform: "linux", plistExists: false })).toBe("spawn"); + }); +}); diff --git a/hub-client/hub-starter.ts b/hub-client/hub-starter.ts new file mode 100644 index 0000000..61dd78d --- /dev/null +++ b/hub-client/hub-starter.ts @@ -0,0 +1,25 @@ +/** + * Hub 启动权归属。 + * + * macOS 上装了 com.forge-hub.plist 时,hub 应该由 launchd 启动。hub-client 若自己 + * `spawn(..., { detached: true })` 再 `unref()`,会造出一个 launchd 管不到的孤儿: + * 它占着 Hub 端口,launchd 每 ThrottleInterval 撞一次、崩一次,而 `cli.ts sync` 的 + * `launchctl bootout` 同样杀不到它(不是 launchd 的子进程)——于是 sync 的"重启" + * 静默失效:磁盘上换成了新代码,跑着的还是旧进程。 + * + * 没装 plist(或非 macOS)时,自己 spawn 仍然是唯一可行的启动方式,保持原行为。 + */ +import path from "node:path"; + +export const LAUNCHD_LABEL = "com.forge-hub"; + +export function launchdPlistPath(home: string): string { + return path.join(home, "Library", "LaunchAgents", `${LAUNCHD_LABEL}.plist`); +} + +export function chooseHubStarter(params: { + platform: string; + plistExists: boolean; +}): "launchd" | "spawn" { + return params.platform === "darwin" && params.plistExists ? "launchd" : "spawn"; +} diff --git a/hub-server/channels/feishu-lark-cli.ts b/hub-server/channels/feishu-lark-cli.ts index a9fad26..d6facfc 100644 --- a/hub-server/channels/feishu-lark-cli.ts +++ b/hub-server/channels/feishu-lark-cli.ts @@ -224,6 +224,73 @@ function startSubscription(): void { }); } +/** 飞书把图片以 "[Image: img_xxx]" / "[图片: img_xxx]" 的形式塞在 content 里。 */ +const IMAGE_PLACEHOLDER = /\[(?:Image|图片):\s*(img_[^\]]+)\]/; + +/** + * 把 content 里**每一个**图片占位符就地换成 "[图片] <本地路径>",其余文字原样保留。 + * + * 原实现用 content.match() 只取第一个 key,随后把整个 content 覆盖成那一张的路径—— + * 一条夹带 N 张图和文字的富文本,到 agent 手里只剩第一张图,文字和其余 N-1 张全部丢失。 + * + * 下载失败的那一张保留原占位符,不牵连同一条消息里其它图片。 + * download 注入,便于测试。 + */ +export async function resolveImagePlaceholders( + content: string, + download: (imageKey: string) => Promise, +): Promise { + // 从不带 g 的常量另建带 g 的副本:共享带 g 的正则会因 lastIndex 残留而漏匹配 + const matches = [...content.matchAll(new RegExp(IMAGE_PLACEHOLDER, "g"))]; + if (matches.length === 0) return content; + + let out = ""; + let cursor = 0; + for (const m of matches) { + const start = m.index ?? 0; + const filePath = await download(m[1]); + out += content.slice(cursor, start) + (filePath ? `[图片] ${filePath}` : m[0]); + cursor = start + m[0].length; + } + return out + content.slice(cursor); +} + +/** 引用回复前缀里保留的被引用正文长度;再长就截断加省略号 */ +const REPLY_EXCERPT_LIMIT = 200; + +/** + * 把「被引用的那条」拼到当前消息前面,agent 才知道用户在回什么。 + * + * 用户在飞书里选中一条消息点「回复」,事件带 reply_to(父消息 id)——此前通道只取 + * sender/chat/type/content 四个字段,reply_to 落地即丢,agent 收到的是一条裸文本, + * 用户只能把原文复制一遍。 + * + * parentText 为 null(取不回父消息)时退化为只带 id,至少让引用关系可追。 + */ +export function formatReplyContext(parentText: string | null, parentId: string, content: string): string { + if (parentText == null) return `[回复 ${parentId}]\n${content}`; + const flat = parentText.replace(/\s+/g, " ").trim(); + const excerpt = flat.length > REPLY_EXCERPT_LIMIT ? `${flat.slice(0, REPLY_EXCERPT_LIMIT)}…` : flat; + return `[回复 ▸ ${excerpt}]\n${content}`; +} + +/** 按 message_id 取一条消息的已渲染正文;任何失败返回 null(调用方决定退化方式) */ +async function fetchMessageText(messageId: string): Promise { + try { + const out = await execFileText(LARK_CLI, [ + "im", "+messages-mget", + "--message-ids", messageId, + "--as", "bot", + ], { timeout: 10000 }); + const parsed = JSON.parse(out) as { ok?: boolean; data?: { messages?: Array<{ content?: unknown }> } }; + const text = parsed.ok ? parsed.data?.messages?.[0]?.content : undefined; + return typeof text === "string" && text ? text : null; + } catch (err) { + hub.logError(`引用消息取回失败 ${messageId}: ${String(err)}`); + return null; + } +} + async function handleMessage(event: Record): Promise { const senderId = (event.sender_id ?? "") as string; const chatId = (event.chat_id ?? "") as string; @@ -268,12 +335,11 @@ async function handleMessage(event: Record): Promise { const messageId = (event.message_id ?? event.id ?? "") as string; let content = (event.content ?? "") as string; - // image key 在 content 里: "[Image: img_v3_xxx]" 或 "[图片: img_v3_xxx]" - const imageKeyMatch = content.match(/\[(?:Image|图片):\s*(img_[^\]]+)\]/); - if (imageKeyMatch && messageId) { - const imageKey = imageKeyMatch[1]; - const filePath = await downloadFeishuMedia(messageId, "image", imageKey); - content = filePath ? `[图片] ${filePath}` : `[图片: ${imageKey}]`; + // 一条富文本可以夹带多张图和文字,逐个就地替换、其余文字原样保留 + if (IMAGE_PLACEHOLDER.test(content) && messageId) { + content = await resolveImagePlaceholders(content, (key) => + downloadFeishuMedia(messageId, "image", key), + ); } else if (msgType === "file" && messageId) { // file key 可能在 content 里: "[File: file_v3_xxx]" const fileKeyMatch = content.match(/\[(?:File|文件):\s*(file_[^\]]+)\]/); @@ -306,6 +372,12 @@ async function handleMessage(event: Record): Promise { if (!content) return; + // 引用回复:把被引用的那条取回来拼在前面,取不到就只带 id + const replyToId = (event.reply_to ?? "") as string; + if (replyToId) { + content = formatReplyContext(await fetchMessageText(replyToId), replyToId, content); + } + const senderDisplay = hub.getNickname(senderId) || senderId; const displayName = isAuthorizedGroup ? `${senderDisplay} @ ${hub.getNickname(chatId)}` @@ -324,6 +396,7 @@ async function handleMessage(event: Record): Promise { chat_id: chatId, message_type: msgType, auth_sender_id: isAuthorizedGroup ? chatId : senderId, + ...(replyToId ? { reply_to: replyToId } : {}), }, }); diff --git a/hub-server/feishu-media-placeholders.test.ts b/hub-server/feishu-media-placeholders.test.ts new file mode 100644 index 0000000..b3d82f4 --- /dev/null +++ b/hub-server/feishu-media-placeholders.test.ts @@ -0,0 +1,79 @@ +import { describe, expect, test } from "bun:test"; + +import { resolveImagePlaceholders } from "./channels/feishu-lark-cli.js"; + +/** 按 key 返回假路径;返回 null 模拟下载失败。 */ +function fakeDownloader(failFor: string[] = []) { + const calls: string[] = []; + const download = async (key: string): Promise => { + calls.push(key); + return failFor.includes(key) ? null : `/media/${key}.png`; + }; + return { download, calls }; +} + +describe("resolveImagePlaceholders", () => { + test("replaces a single placeholder with the downloaded path", async () => { + const { download } = fakeDownloader(); + + const out = await resolveImagePlaceholders("[图片: img_v3_aaa]", download); + + expect(out).toBe("[图片] /media/img_v3_aaa.png"); + }); + + test("replaces every placeholder when one message carries several images", async () => { + const { download, calls } = fakeDownloader(); + + const out = await resolveImagePlaceholders( + "[图片: img_v3_aaa][图片: img_v3_bbb][图片: img_v3_ccc]", + download, + ); + + expect(calls).toEqual(["img_v3_aaa", "img_v3_bbb", "img_v3_ccc"]); + expect(out).toBe( + "[图片] /media/img_v3_aaa.png[图片] /media/img_v3_bbb.png[图片] /media/img_v3_ccc.png", + ); + }); + + test("keeps the surrounding text of a rich-text message", async () => { + const { download } = fakeDownloader(); + + const out = await resolveImagePlaceholders( + "先看这张 [图片: img_v3_aaa] 再看这张 [图片: img_v3_bbb] 完", + download, + ); + + expect(out).toBe( + "先看这张 [图片] /media/img_v3_aaa.png 再看这张 [图片] /media/img_v3_bbb.png 完", + ); + }); + + test("accepts the English [Image: ...] form", async () => { + const { download } = fakeDownloader(); + + const out = await resolveImagePlaceholders("a [Image: img_v3_aaa] b", download); + + expect(out).toBe("a [图片] /media/img_v3_aaa.png b"); + }); + + test("leaves the placeholder in place when the download fails", async () => { + const { download } = fakeDownloader(["img_v3_bbb"]); + + const out = await resolveImagePlaceholders( + "[图片: img_v3_aaa] 中间 [图片: img_v3_bbb]", + download, + ); + + // 失败的那张保留原占位符,成功的那张照常替换——不因为一张失败就丢掉整条消息 + expect(out).toBe("[图片] /media/img_v3_aaa.png 中间 [图片: img_v3_bbb]"); + }); + + test("returns text without placeholders untouched and downloads nothing", async () => { + const { download, calls } = fakeDownloader(); + + const out = await resolveImagePlaceholders("纯文字,没有图", download); + + expect(out).toBe("纯文字,没有图"); + expect(calls).toEqual([]); + }); +}); diff --git a/hub-server/feishu-reply-context.test.ts b/hub-server/feishu-reply-context.test.ts new file mode 100644 index 0000000..16e4c3b --- /dev/null +++ b/hub-server/feishu-reply-context.test.ts @@ -0,0 +1,45 @@ +import { describe, expect, test } from "bun:test"; + +import { formatReplyContext } from "./channels/feishu-lark-cli.js"; + +describe("formatReplyContext", () => { + test("prefixes the quoted parent text so the agent knows what is being replied to", () => { + const out = formatReplyContext("哪两条没有?", "om_parent", "第二条"); + + expect(out).toBe("[回复 ▸ 哪两条没有?]\n第二条"); + }); + + test("collapses newlines inside the quoted text into single spaces", () => { + const out = formatReplyContext("第一行\n\n第二行\t第三行", "om_parent", "好"); + + expect(out).toBe("[回复 ▸ 第一行 第二行 第三行]\n好"); + }); + + test("truncates a long quoted text to 200 chars with an ellipsis", () => { + const parent = "字".repeat(300); + + const out = formatReplyContext(parent, "om_parent", "收到"); + + expect(out).toBe(`[回复 ▸ ${"字".repeat(200)}…]\n收到`); + }); + + test("does not add an ellipsis when the quoted text is exactly at the limit", () => { + const parent = "x".repeat(200); + + const out = formatReplyContext(parent, "om_parent", "ok"); + + expect(out).toBe(`[回复 ▸ ${parent}]\nok`); + }); + + test("falls back to the parent id when the parent text could not be fetched", () => { + const out = formatReplyContext(null, "om_0000000000000000000000000000ffff", "试试"); + + expect(out).toBe("[回复 om_0000000000000000000000000000ffff]\n试试"); + }); + + test("keeps the user's own content untouched, including its newlines", () => { + const out = formatReplyContext("原句", "om_parent", "第一段\n第二段"); + + expect(out).toBe("[回复 ▸ 原句]\n第一段\n第二段"); + }); +});