Files
MailUI4Agents/plugins/dsh-mail-bridge/test/turnwait.test.mjs
JianFeeeee 375578cf9b 补 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。
2026-09-04 21:39:50 +08:00

114 lines
4.0 KiB
JavaScript
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

// 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, '不同会话应并发');
});