Commit 215f3196 by luoqi

merge: main → test(反向合)—— 僵尸锁周期回收

parents 5d970f15 3e86060c
Pipeline #3671 failed in 0 seconds
......@@ -40,9 +40,28 @@ import { schedulerDisabled } from './scheduler-switch';
/**
* 僵尸锁的最小年龄。比「一轮同步的正常耗时」留足余量 —— 生产实测单轮摄入 28~52 分钟,
* 取 3 小时:真崩溃留下的锁必然远超此值,而正常在跑的绝不会。
*
* 🔴 2026-09-04 测试机事故:**这条阈值单独用在启动回收上,对 deploy 孤儿锁 100% 失效。**
* deploy 孤儿锁的年龄 = T_重启 − T_起跑,而那轮在重启时**尚未跑完** ⇒ 年龄恒 < 一轮耗时;
* 本阈值又被刻意取成 > 一轮耗时。两式相与 ⇒ 任何 deploy 期间在飞的 sync 留下的锁,
* 在新进程启动那一刻必然"太年轻",必然逃过回收 —— 不是概率,是结构性必然。
* 实测:锁 14:15:00 起,容器 14:30:24 重建,回收在 14:30:31 跑了、查到 0 行、静默 return;
* 锁当时 15.5 分钟大,离 180 分钟阈值差 164.5 分钟。此后每轮 cron 全被并发锁拦下,
* 游标 24.7 小时不动直到人工重启。
* ⇒ 修法不是调阈值(调多小都会被"一轮耗时"从上方封住),而是**给回收增加时机**:
* 见 reapStaleRunningLocks 的周期态调用。
*/
const REAP_MIN_AGE_MS = 3 * 60 * 60 * 1000;
/**
* 周期回收只碰 **cron 增量自己**留下的锁(triggered_by = `sync:<host>:<runId>`)。
*
* ⛔ 不能扩到全部前缀:`full:` 是手工全量 / 按诊所补摄,建国门实测 import 段就跑了 2h25m,
* 离 3 小时阈值只剩 35 分钟;`patient_refresh:` 是客服点「刷新」的前台请求。
* 周期回收每轮都跑,把它们纳入 = 把 2026-08-30 那个「误杀正在跑的活儿」换个入口放回来。
*/
const REAP_PERIODIC_PREFIX = 'sync:';
@Injectable()
export class SyncIncrementalSchedulerService implements OnModuleInit {
private readonly logger = new Logger(SyncIncrementalSchedulerService.name);
......@@ -97,7 +116,7 @@ export class SyncIncrementalSchedulerService implements OnModuleInit {
// 启动即回收上一个进程留下的僵尸锁(部署/崩溃重启时,running 行来不及标终态,
// 会永久占住 sync_logs 的 partial UNIQUE(host_id) WHERE status='running',
// 导致之后每次 cron 增量都被并发锁 skip → 游标不动 → DW 滞后告警)。
await this.reapStaleRunningLocks();
await this.reapStaleRunningLocks(PROCESS_STARTED_AT);
const globalCron = process.env.PAC_INCREMENTAL_CRON;
const dataDir = this.dataDir();
......@@ -171,23 +190,47 @@ export class SyncIncrementalSchedulerService implements OnModuleInit {
* @param processStartedAt 回收判据的时间界(默认本进程启动时刻)。显式传入是为了让判据可测 ——
* 测试不必靠真实时钟凑时间差(高负载下事件循环被拖慢会越过阈值边界,产生间歇性假失败)。
*/
private async reapStaleRunningLocks(processStartedAt: Date = PROCESS_STARTED_AT): Promise<void> {
private staleCutoff(processStartedAt?: Date): Date {
const ageCutoff = new Date(Date.now() - REAP_MIN_AGE_MS);
// 周期态(不传 processStartedAt):**只用年龄闸**。
// 本进程是长驻的,带上 processStartedAt 会让 min() 恒取它 → 对"本进程启动之后
// 才产生的孤儿锁"永远为 no-op,周期回收等于没加。
if (!processStartedAt) return ageCutoff;
// 启动态:双条件取更早的那个界(保住 2026-08-30 的 CLI 误杀修复)。
return ageCutoff < processStartedAt ? ageCutoff : processStartedAt;
}
/**
* ⛔ **这里不能写默认参数** `= PROCESS_STARTED_AT`。JS 的默认值对显式传入的 `undefined`
* 同样生效,于是周期态调用 `reap(undefined, opts)` 会被悄悄补成启动态 —— 判据变回
* min(age, 进程启动),对长驻进程恒 no-op,周期回收等于没加、还不报错。
* 2026-09-04 写这个修复时就踩了一次,是 sync-lock-reap.spec 的远期基准用例咬出来的。
* 启动态由调用方显式传 PROCESS_STARTED_AT。
*/
private async reapStaleRunningLocks(
processStartedAt?: Date,
opts?: { triggerPrefix?: string; reason?: string },
): Promise<number> {
try {
// 双条件取更早的那个界:既要早于本进程启动,又要已经躺够 REAP_MIN_AGE_MS。
const ageCutoff = new Date(Date.now() - REAP_MIN_AGE_MS);
const cutoff = ageCutoff < processStartedAt ? ageCutoff : processStartedAt;
const cutoff = this.staleCutoff(processStartedAt);
const where = {
status: SyncStatus.RUNNING,
startedAt: { lt: cutoff },
...(opts?.triggerPrefix ? { triggeredBy: { startsWith: opts.triggerPrefix } } : {}),
};
const stale = await this.prisma.syncLog.findMany({
where: { status: SyncStatus.RUNNING, startedAt: { lt: cutoff } },
where,
select: { id: true, hostId: true, startedAt: true, triggeredBy: true },
});
if (stale.length === 0) return;
if (stale.length === 0) return 0;
const { count } = await this.prisma.syncLog.updateMany({
where: { status: SyncStatus.RUNNING, startedAt: { lt: cutoff } },
where,
data: {
status: SyncStatus.FAILED,
endedAt: new Date(),
errorMessage: '进程重启前未完成,启动时回收僵尸锁(reapStaleRunningLocks)',
errorMessage:
opts?.reason ?? '进程重启前未完成,启动时回收僵尸锁(reapStaleRunningLocks)',
},
});
for (const s of stale) {
......@@ -196,10 +239,14 @@ export class SyncIncrementalSchedulerService implements OnModuleInit {
`trigger=${s.triggeredBy}`,
);
}
this.logger.log(`sync-incremental: 启动回收 ${count} 个僵尸同步锁`);
this.logger.log(
`sync-incremental: 回收 ${count} 个僵尸同步锁(${processStartedAt ? '启动态' : '周期态'})`,
);
return count;
} catch (err) {
// 回收失败不阻断启动(cron 仍注册);最坏退回人工清理,不比现状糟
this.logger.error(`sync-incremental: 僵尸锁回收失败: ${(err as Error).message}`);
return 0;
}
}
......@@ -215,6 +262,18 @@ export class SyncIncrementalSchedulerService implements OnModuleInit {
}
this.runningHosts.add(host);
try {
// ⭐ 2026-09-04:**回收的第二次机会**。启动态回收挡不住 deploy 孤儿锁
// (见 REAP_MIN_AGE_MS 注释:那种锁在重启那一刻必然"太年轻"),而回收原本
// 只有 onModuleInit 一次机会 → 逃过就永久卡死,实测把测试机的增量停了 24.7 小时。
// 这里每轮开跑前再收一次:判据只用年龄闸(周期态),且**只碰本 host 的 cron 增量锁**。
// 走到这一行说明 runningHosts 里没有本 host —— 即本进程没有在跑它,
// 那么一把躺满 3 小时的 `sync:<host>:` 锁必定是死进程留下的。
// 最坏情况:锁比阈值年轻 → 本轮照旧被并发锁拦下,下轮(2 小时后)再收,
// 自愈时延上限 = REAP_MIN_AGE_MS + 一个 cron 间隔,不再需要人工重启。
await this.reapStaleRunningLocks(undefined, {
triggerPrefix: `${REAP_PERIODIC_PREFIX}${host}:`,
reason: '上一轮未正常终结,下一轮开跑前回收僵尸锁(周期态 reapStaleRunningLocks)',
});
await this.runOne(path.join(this.dataDir(), host));
} catch (err) {
if (err instanceof SyncAlreadyRunningError) {
......
......@@ -45,8 +45,14 @@ describe('闸门覆盖范围 — 只关写侧', () => {
test('⭐ 增量摄入(会改数据)受闸控,且连僵尸锁回收一并跳过', () => {
const s = read('sync-incremental.scheduler.ts');
expect(s).toContain('schedulerDisabled()');
// 闸判定必须在 reapStaleRunningLocks 之前 —— 共库场景下回收会误清老实例正在跑的锁
expect(s.indexOf('schedulerDisabled()')).toBeLessThan(s.indexOf('await this.reapStaleRunningLocks()'));
// 闸判定必须在**第一次** reapStaleRunningLocks 调用之前 —— 共库场景下回收会误清老实例正在跑的锁。
// ⚠️ 2026-09-04:调用点从 `reapStaleRunningLocks()` 变成带参数的两处(启动态传
// PROCESS_STARTED_AT,周期态在 runHostSafe 里传 undefined+前缀),写死字面量会假失败。
// 周期态那处是**传递性受闸**的:闸命中时 onModuleInit 提前 return,cron 根本不注册,
// runHostSafe 无从被调用 —— 所以这里只需锁住"闸在第一次回收之前"。
const firstReap = s.indexOf('this.reapStaleRunningLocks(');
expect(firstReap).toBeGreaterThan(-1);
expect(s.indexOf('schedulerDisabled()')).toBeLessThan(firstReap);
});
test('⭐ persona stale-scan(会 enqueue 重算)受闸控', () => {
......
......@@ -57,6 +57,10 @@ describe('runHostSafe 防重入', () => {
const { calls, run, release } = makeService();
const a = run('jvs-dw');
const b = run('friday');
// ⚠️ 2026-09-04:runHostSafe 在进 runOne 之前多了一次 `await 周期回收僵尸锁`,
// 于是 runOne 不再在同一微任务里被调到 —— 断言前必须把微任务队列放干,
// 否则 calls 还是空的(这是真实的行为变化,不是测试写错)。
await new Promise((r) => setImmediate(r));
// runOne 收到的是 path.join(dataDir, host) 的完整路径,按后缀断言
expect(calls.map((d) => d.split('/').pop())).toEqual(['jvs-dw', 'friday']);
release();
......
......@@ -25,14 +25,16 @@ const before = (ms: number) => new Date(PROCESS_STARTED_AT.getTime() - ms);
const after = (ms: number) => new Date(PROCESS_STARTED_AT.getTime() + ms);
function makeService(rows: Array<{ id: string; hostId: string; startedAt: Date; triggeredBy: string }>) {
const findMany = jest.fn().mockImplementation(async ({ where }) => {
const lt: Date = where.startedAt.lt;
return rows.filter((r) => r.startedAt < lt);
});
const updateMany = jest.fn().mockImplementation(async ({ where }) => {
const lt: Date = where.startedAt.lt;
return { count: rows.filter((r) => r.startedAt < lt).length };
});
// 假 prisma 必须**同时**认 startedAt.lt 和 triggeredBy.startsWith ——
// 只认前者的话,周期回收的前缀过滤(不许碰 full: / patient_refresh:)就测不出来。
const match = (where: { startedAt: { lt: Date }; triggeredBy?: { startsWith: string } }) =>
rows.filter(
(r) =>
r.startedAt < where.startedAt.lt &&
(!where.triggeredBy || r.triggeredBy.startsWith(where.triggeredBy.startsWith)),
);
const findMany = jest.fn().mockImplementation(async ({ where }) => match(where));
const updateMany = jest.fn().mockImplementation(async ({ where }) => ({ count: match(where).length }));
const prisma = { syncLog: { findMany, updateMany } };
const svc = new SyncIncrementalSchedulerService(
prisma as never,
......@@ -48,9 +50,20 @@ function makeService(rows: Array<{ id: string; hostId: string; startedAt: Date;
// 触达 private 方法(纯编排,无需走 onModuleInit 以免顺带注册 cron);时间界显式注入
const reap = (svc: SyncIncrementalSchedulerService, processStartedAt = PROCESS_STARTED_AT) =>
(
svc as unknown as { reapStaleRunningLocks: (at: Date) => Promise<void> }
svc as unknown as { reapStaleRunningLocks: (at: Date) => Promise<unknown> }
).reapStaleRunningLocks(processStartedAt);
/// 周期态回收:不传 processStartedAt(只用年龄闸),可带前缀过滤
const reapPeriodic = (svc: SyncIncrementalSchedulerService, triggerPrefix?: string) =>
(
svc as unknown as {
reapStaleRunningLocks: (
at: Date | undefined,
opts?: { triggerPrefix?: string; reason?: string },
) => Promise<unknown>;
}
).reapStaleRunningLocks(undefined, triggerPrefix ? { triggerPrefix } : undefined);
describe('reapStaleRunningLocks', () => {
test('回收进程启动前的 running 行 → 标 failed', async () => {
const { svc, updateMany } = makeService([
......@@ -102,6 +115,100 @@ describe('reapStaleRunningLocks', () => {
const svc = new SyncIncrementalSchedulerService(
prisma as never, {} as never, {} as never, {} as never, {} as never, {} as never,
);
await expect(reap(svc)).resolves.toBeUndefined();
// 契约是「不抛」;返回值 2026-09-04 从 void 改成「回收了几条」,失败路径返回 0。
await expect(reap(svc)).resolves.toBe(0);
});
});
/**
* 🔴 2026-09-04 测试机事故的回归 —— 以及**为什么上面那组用例全绿却没拦住它**。
*
* 事故:锁 14:15:00 起跑,容器 14:30:24 被 deploy 拆掉,新进程 14:30:31 跑启动回收 →
* 查到 0 行静默返回(锁才 15.5 分钟大,离 REAP_MIN_AGE_MS=180 分钟差 164.5 分钟)。
* 此后每轮 cron 被并发锁拦下,游标 24.7 小时不动,直到人工重启。
*
* 上面那组为什么测不到:基准 `PROCESS_STARTED_AT` 是一个**过去**的固定日期,
* 而 `ageCutoff = 真实now − 3h` 永远晚于它 ⇒ `min()` 恒取 processStartedAt,
* **REAP_MIN_AGE_MS 一行都没被执行到**。更糟的是「where 用 lt」那条直接断言
* `lt === PROCESS_STARTED_AT`,把**加年龄闸之前**的行为钉成契约 —— 生产已坏、用例全绿。
* ⇒ 这里改用 fake timers 钉住 Date.now(),让年龄闸真正进入判据。
*/
describe('reapStaleRunningLocks — deploy 孤儿锁(年龄闸的盲区)', () => {
/// 重启时刻。锁比它早 15 分钟 —— 复刻事故里的 14:15 起跑 / 14:30 重启。
const RESTART = new Date('2026-09-03T06:30:24.000Z');
const LOCK_AT = new Date(RESTART.getTime() - 15 * 60_000);
const THREE_H = 3 * 60 * 60 * 1000;
/**
* ⚠️ 周期态的用例必须用**远期**基准,不能沿用上面的 2026 日期 —— 否则测不出东西。
* 模块里的 PROCESS_STARTED_AT 是**真实**的模块加载时刻。若假时钟停在过去,
* `ageCutoff = 假now − 3h` 就恒早于它,`min()` 两边选出同一个值,
* "周期态有没有混进 processStartedAt"这件事在判据上**不可观测**。
* 实测:把周期态改回 min(age, PROCESS_STARTED_AT),用 2026 基准时 11 条用例全绿。
* 钉到 2099 后,真实模块加载时刻必然更早 → min() 会选它 → 判据能分辨 → 用例有牙。
*/
const FAR = new Date('2099-01-01T00:00:00.000Z');
const FAR_LOCK = new Date(FAR.getTime() - THREE_H - 60_000); // 躺满 3h 零 1 分
beforeEach(() => jest.useFakeTimers().setSystemTime(RESTART));
afterEach(() => jest.useRealTimers());
const orphan = () => [
{ id: 'deploy-orphan', hostId: 'h1', startedAt: LOCK_AT, triggeredBy: 'sync:jvs-dw:r1' },
];
const farOrphan = () => [
{ id: 'deploy-orphan', hostId: 'h1', startedAt: FAR_LOCK, triggeredBy: 'sync:jvs-dw:r1' },
];
test('🔴 复现:启动态回收**收不掉** deploy 孤儿锁(锁太年轻,这是结构性必然不是运气)', async () => {
const { svc, updateMany } = makeService(orphan());
await reap(svc, RESTART); // 启动态:cutoff = min(now-3h, RESTART) = now-3h
expect(updateMany).not.toHaveBeenCalled(); // ← 事故当场的行为,原样钉住
});
test('⭐ 周期态:锁躺满 3 小时后,下一轮 cron 开跑前把它收掉(自愈,不必人工重启)', async () => {
const { svc, updateMany } = makeService(farOrphan());
jest.setSystemTime(FAR);
await reapPeriodic(svc, 'sync:jvs-dw:');
expect(updateMany).toHaveBeenCalledTimes(1);
expect(updateMany.mock.calls[0][0].data.status).toBe('failed');
});
test('⭐ 周期态判据里**不许**混进 processStartedAt —— 混进去对长驻进程就是 no-op', async () => {
// 长驻服务:进程启动早于锁。若周期态仍取 min(age, processStart) → cutoff=processStart,
// 锁比它晚 → 永远收不掉,周期回收等于没加。这条就是钉这个。
const { svc, findMany } = makeService(farOrphan());
jest.setSystemTime(FAR);
await reapPeriodic(svc, 'sync:jvs-dw:');
const lt: Date = findMany.mock.calls[0][0].where.startedAt.lt;
expect(lt.getTime()).toBe(Date.now() - THREE_H); // 纯年龄闸
});
test('⛔ 周期态只碰 sync: 前缀 —— full: 补摄(实测 2h25m)和 patient_refresh: 绝不碰', async () => {
const { svc, updateMany } = makeService([
{ id: 'cron', hostId: 'h1', startedAt: FAR_LOCK, triggeredBy: 'sync:jvs-dw:r1' },
{ id: 'backfill', hostId: 'h1', startedAt: FAR_LOCK, triggeredBy: 'full:jvs-dw:r2' },
{ id: 'refresh', hostId: 'h1', startedAt: FAR_LOCK, triggeredBy: 'patient_refresh:P9:r3' },
]);
jest.setSystemTime(FAR);
await reapPeriodic(svc, 'sync:jvs-dw:');
expect(updateMany).toHaveBeenCalledTimes(1);
expect(updateMany.mock.calls[0][0].where.triggeredBy).toEqual({ startsWith: 'sync:jvs-dw:' });
});
test('⛔ 周期态也认前缀里的 host —— 别的 host 的 cron 锁不碰', async () => {
const { svc, updateMany } = makeService([
{ id: 'other-host', hostId: 'h2', startedAt: FAR_LOCK, triggeredBy: 'sync:friday:r9' },
]);
jest.setSystemTime(FAR);
await reapPeriodic(svc, 'sync:jvs-dw:');
expect(updateMany).not.toHaveBeenCalled();
});
test('⛔ 周期态不许收还没躺满 3 小时的锁(正在跑的那轮)', async () => {
const { svc, updateMany } = makeService(farOrphan());
jest.setSystemTime(new Date(FAR_LOCK.getTime() + THREE_H - 60_000)); // 差 1 分钟
await reapPeriodic(svc, 'sync:jvs-dw:');
expect(updateMany).not.toHaveBeenCalled();
});
});
Markdown is supported
0% or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment