diff --git a/plugins/dsh-mail-bridge/test/turnwait.test.mjs b/plugins/dsh-mail-bridge/test/turnwait.test.mjs new file mode 100644 index 0000000..a73df84 --- /dev/null +++ b/plugins/dsh-mail-bridge/test/turnwait.test.mjs @@ -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, '不同会话应并发'); +});