From 375578cf9b4c12ffc80a5428b67266e35035fc36 Mon Sep 17 00:00:00 2001 From: JianFeeeee Date: Fri, 4 Sep 2026 21:39:50 +0800 Subject: [PATCH] =?UTF-8?q?=E8=A1=A5=20dsh=20waitForTurnEnd=20/=20locked?= =?UTF-8?q?=20=E7=9A=84=E8=AF=AD=E4=B9=89=E6=B5=8B=E8=AF=95=EF=BC=886=20?= =?UTF-8?q?=E4=BE=8B=EF=BC=89?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 全流程逐项验证时唯一没有测试覆盖的一处:两个函数都是 apply() 内的闭包, import 不到,于是把结构原样复刻进 test/turnwait.test.mjs 验语义不变量。 锁住的六条: - turn/end 到了立刻返回,**且注销监听** —— 不注销的话每封邮件泄漏一个监听器 - 别的会话的 turn/end 不该让本会话提前返回(session 身份判据) - 事件永不到来时超时兜底返回,不永久挂起(模型崩了不发 turn/end 的情形) - 超时路径也要注销监听 - locked 严格串行(交错会让 DSH 报 message already pending) - 前一个任务抛错不让后续卡死(release 在 finally 里) - 不同会话不互相串行 dsh 219 → 225。 --- .../dsh-mail-bridge/test/turnwait.test.mjs | 113 ++++++++++++++++++ 1 file changed, 113 insertions(+) create mode 100644 plugins/dsh-mail-bridge/test/turnwait.test.mjs 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, '不同会话应并发'); +});