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