补 dsh waitForTurnEnd / locked 的语义测试(6 例)
全流程逐项验证时唯一没有测试覆盖的一处:两个函数都是 apply() 内的闭包, import 不到,于是把结构原样复刻进 test/turnwait.test.mjs 验语义不变量。 锁住的六条: - turn/end 到了立刻返回,**且注销监听** —— 不注销的话每封邮件泄漏一个监听器 - 别的会话的 turn/end 不该让本会话提前返回(session 身份判据) - 事件永不到来时超时兜底返回,不永久挂起(模型崩了不发 turn/end 的情形) - 超时路径也要注销监听 - locked 严格串行(交错会让 DSH 报 message already pending) - 前一个任务抛错不让后续卡死(release 在 finally 里) - 不同会话不互相串行 dsh 219 → 225。
This commit is contained in:
113
plugins/dsh-mail-bridge/test/turnwait.test.mjs
Normal file
113
plugins/dsh-mail-bridge/test/turnwait.test.mjs
Normal file
@ -0,0 +1,113 @@
|
||||
// waitForTurnEnd / locked 的语义验证(把两者的结构原样复刻,因为它们是
|
||||
// apply() 内的闭包,无法 import)。要验的是:
|
||||
// 1. turn/end 到了立刻返回,并注销监听(不泄漏)
|
||||
// 2. 事件永不到来时 120s 超时兜底返回,而不是永久挂起
|
||||
// 3. locked 严格串行,且一个任务抛错不会让后续任务卡死
|
||||
import { test } from 'node:test';
|
||||
import assert from 'node:assert/strict';
|
||||
|
||||
function makeCtx() {
|
||||
const listeners = new Set();
|
||||
return {
|
||||
on(_evt, fn) { listeners.add(fn); return () => listeners.delete(fn); },
|
||||
emit(session, event) { for (const fn of [...listeners]) fn(session, event); },
|
||||
get listenerCount() { return listeners.size; },
|
||||
};
|
||||
}
|
||||
|
||||
function makeWaiter(ctx) {
|
||||
return async function waitForTurnEnd(agent, timeoutMs = 120_000) {
|
||||
await new Promise((resolve) => {
|
||||
let done = false;
|
||||
const finish = () => { if (!done) { done = true; clearTimeout(timer); dispose?.(); resolve(); } };
|
||||
const timer = setTimeout(finish, timeoutMs);
|
||||
const dispose = ctx.on('session/event', (session, event) => {
|
||||
if (session !== agent.session) return;
|
||||
if (event?.type === 'turn/end') finish();
|
||||
});
|
||||
});
|
||||
};
|
||||
}
|
||||
|
||||
function makeLocked() {
|
||||
const sessionLocks = new Map();
|
||||
return async function locked(id, fn) {
|
||||
const prev = sessionLocks.get(id) ?? Promise.resolve();
|
||||
let release;
|
||||
const next = new Promise((r) => { release = r; });
|
||||
sessionLocks.set(id, next);
|
||||
try {
|
||||
await prev;
|
||||
return await fn();
|
||||
} finally {
|
||||
release();
|
||||
if (sessionLocks.get(id) === next) sessionLocks.delete(id);
|
||||
}
|
||||
};
|
||||
}
|
||||
|
||||
test('turn/end 到了立刻返回,且注销监听不泄漏', async () => {
|
||||
const ctx = makeCtx();
|
||||
const wait = makeWaiter(ctx);
|
||||
const agent = { session: 'S1' };
|
||||
const p = wait(agent, 5000);
|
||||
await new Promise((r) => setTimeout(r, 20));
|
||||
assert.equal(ctx.listenerCount, 1, '等待期间应有一个监听');
|
||||
ctx.emit('S1', { type: 'turn/end' });
|
||||
await p;
|
||||
assert.equal(ctx.listenerCount, 0, '返回后必须注销监听,否则每封邮件泄漏一个');
|
||||
});
|
||||
|
||||
test('别的会话的 turn/end 不该让本会话提前返回', async () => {
|
||||
const ctx = makeCtx();
|
||||
const wait = makeWaiter(ctx);
|
||||
let done = false;
|
||||
const p = wait({ session: 'MINE' }, 400).then(() => { done = true; });
|
||||
ctx.emit('OTHER', { type: 'turn/end' });
|
||||
await new Promise((r) => setTimeout(r, 60));
|
||||
assert.equal(done, false, '别人的事件不算');
|
||||
await p;
|
||||
assert.equal(done, true, '自己的超时仍要兜底');
|
||||
});
|
||||
|
||||
test('事件永不到来时超时兜底返回(不永久挂起)', async () => {
|
||||
const ctx = makeCtx();
|
||||
const wait = makeWaiter(ctx);
|
||||
const t0 = Date.now();
|
||||
await wait({ session: 'GHOST' }, 250);
|
||||
const dt = Date.now() - t0;
|
||||
assert.ok(dt >= 240 && dt < 1500, `应在超时后返回,实际 ${dt}ms`);
|
||||
assert.equal(ctx.listenerCount, 0, '超时路径也要注销监听');
|
||||
});
|
||||
|
||||
test('locked 严格串行', async () => {
|
||||
const locked = makeLocked();
|
||||
const order = [];
|
||||
const mk = (tag, ms) => locked('S', async () => {
|
||||
order.push(`${tag}-start`);
|
||||
await new Promise((r) => setTimeout(r, ms));
|
||||
order.push(`${tag}-end`);
|
||||
});
|
||||
await Promise.all([mk('a', 80), mk('b', 10)]);
|
||||
assert.deepEqual(order, ['a-start', 'a-end', 'b-start', 'b-end'],
|
||||
'同一会话必须串行,不能交错(DSH 会报 message already pending)');
|
||||
});
|
||||
|
||||
test('前一个任务抛错不会让后续永久卡死', async () => {
|
||||
const locked = makeLocked();
|
||||
await assert.rejects(locked('S', async () => { throw new Error('boom'); }));
|
||||
const got = await locked('S', async () => 'ok');
|
||||
assert.equal(got, 'ok', '锁必须在 finally 里释放');
|
||||
});
|
||||
|
||||
test('不同会话不互相串行', async () => {
|
||||
const locked = makeLocked();
|
||||
let peak = 0, cur = 0;
|
||||
const mk = (id) => locked(id, async () => {
|
||||
cur++; peak = Math.max(peak, cur);
|
||||
await new Promise((r) => setTimeout(r, 60));
|
||||
cur--;
|
||||
});
|
||||
await Promise.all([mk('A'), mk('B')]);
|
||||
assert.equal(peak, 2, '不同会话应并发');
|
||||
});
|
||||
Reference in New Issue
Block a user