verify-answer-stream.mjs 27 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564
  1. /**
  2. * 验证 F073「正文真实增量流式输出」的两块:
  3. * 第 1 部分 —— **纯逻辑**(不必联网):增量规整器的前缀不变式 + 流状态机
  4. * 第 2 部分 —— **协调器级假流**:假 SSE 流喂真实协调器,断言事件通道与去重
  5. *
  6. * 为什么规整器要这么较真:它手里攥着一个不变式——
  7. * **任一时刻「已发出的增量拼接」=== `normalizeChatMarkdown(已输入的全部文本)`**
  8. * 只要不变式成立,后端 `answer_end` 的「整块替换」在正常路径上就是纯尾缀追加
  9. * (不会触发整段重建 → 不会让「思考中」卡片闪现重播)。破了它,观感立刻劣化。
  10. *
  11. * 怎么跑(先重新打包,必须带 --loader:.png=dataurl):
  12. * npx esbuild harness/tools/_entry-coordinator.ts --bundle --format=esm \
  13. * --outfile=harness/tools/_coordinator.mjs --alias:@=./src \
  14. * --define:import.meta.env='{}' --loader:.png=dataurl
  15. * node harness/tools/verify-answer-stream.mjs
  16. */
  17. import {
  18. ApiChatCoordinator,
  19. applyAnswerStreamAbort,
  20. applyAnswerStreamEnd,
  21. applyAnswerStreamEvent,
  22. createAnswerStreamState,
  23. createIncrementalMarkdownNormalizer,
  24. normalizeChatMarkdown,
  25. parseAnswerStreamEvent,
  26. shouldKeepAnswerOverview,
  27. } from './_coordinator.mjs';
  28. let pass = 0;
  29. let fail = 0;
  30. const check = (name, cond, extra = '') => {
  31. if (cond) {
  32. pass++;
  33. console.log(` ok ${name}`);
  34. } else {
  35. fail++;
  36. console.log(` FAIL ${name}${extra === '' ? '' : ` → ${JSON.stringify(extra)}`}`);
  37. }
  38. };
  39. /* ================================================================== *
  40. * 第 1 部分:纯逻辑
  41. * ================================================================== */
  42. console.log('【1】增量规整器:前缀不变式');
  43. /** 把一串输入分片喂给规整器,返回逐片输出 */
  44. const feed = (chunks) => {
  45. const normalizer = createIncrementalMarkdownNormalizer();
  46. const out = [];
  47. for (const chunk of chunks) out.push(normalizer.push(chunk));
  48. return { out, normalizer };
  49. };
  50. /**
  51. * 不变式(流中):**已发出的文本必须是「规整后全文」的前缀**,
  52. * 且未发出的差额只能是「**给当前行准备的空行分隔符(可能 1~2 个 \n)+ 当前行的尾巴**」——
  53. * 也就是「为了判断该行是不是列表项而短暂扣住的一小段」,行尾一到就释放。
  54. *
  55. * 这个性质正好保证 `answer_end` 的「整块替换」= 纯尾缀追加(不会整段重建)。
  56. */
  57. const isPrefixWithHeldLineTail = (got, want) => {
  58. if (!want.startsWith(got)) return false;
  59. const held = want.slice(got.length).replace(/^\n+/, ''); // 去掉随行一起扣住的前置空行
  60. return !held.includes('\n');
  61. };
  62. /** 不变式:每喂一片就核对一次 */
  63. const checkInvariant = (label, chunks, { finish = false } = {}) => {
  64. const normalizer = createIncrementalMarkdownNormalizer();
  65. let want = '';
  66. let got = '';
  67. let ok = true;
  68. let at = 0;
  69. for (let i = 0; i < chunks.length; i++) {
  70. want += chunks[i];
  71. got += normalizer.push(chunks[i]);
  72. if (!isPrefixWithHeldLineTail(got, normalizeChatMarkdown(want))) {
  73. ok = false;
  74. at = i;
  75. break;
  76. }
  77. }
  78. if (ok && finish) {
  79. got += normalizer.finish();
  80. if (got !== normalizeChatMarkdown(want)) ok = false;
  81. }
  82. check(
  83. `不变式:${label}`,
  84. ok,
  85. ok ? '' : { 第几片: at, 期望: normalizeChatMarkdown(want).slice(-40), 实际: got.slice(-40) }
  86. );
  87. return ok;
  88. };
  89. checkInvariant('整段一次 push', ['第一行\n第二行\n第三行']);
  90. checkInvariant('逐字 push(含换行)', [...'第一行\n第二行\n1. a\n2. b']);
  91. checkInvariant('双换行(空行)', ['段落一\n\n段落二']);
  92. checkInvariant('连续有序列表不插空行', ['1. 甲\n2. 乙\n3. 丙']);
  93. checkInvariant('连续无序列表不插空行', ['- 甲\n- 乙']);
  94. checkInvariant('列表后接说明行', ['- 甲\n- 乙\n以上就是全部']);
  95. checkInvariant('伪列表符(2倍)不误判', ['2倍\n3. 真列表']);
  96. checkInvariant('行首恰为列表前缀(- )', ['前一行\n- ']);
  97. checkInvariant('行首数字未定型(12)', ['前一行\n12']);
  98. checkInvariant('行首数字带点(12.)', ['前一行\n12.']);
  99. checkInvariant('末尾无换行', ['段一\n段二(未完']);
  100. checkInvariant('以换行结尾', ['段一\n']);
  101. checkInvariant('空串与纯空白', ['', ' ', '\n', ' \n']);
  102. checkInvariant('逐字 push 且收尾 finish', [...'甲\n1. 乙\n丙'], { finish: true });
  103. // 逐字 push 与整段 push 的「拼接结果」必须一致(都 finish 掉扣住的尾巴后比)
  104. {
  105. const text = '第一段\n第二段\n1. 甲\n2. 乙\n收尾说明';
  106. const wholeFeed = feed([text]);
  107. wholeFeed.out.push(wholeFeed.normalizer.finish());
  108. const perCharFeed = feed([...text]);
  109. perCharFeed.out.push(perCharFeed.normalizer.finish());
  110. const whole = wholeFeed.out.join('');
  111. const perChar = perCharFeed.out.join('');
  112. check('逐字 push 与整段 push 输出一致', whole === perChar, { whole, perChar });
  113. check('两者的拼接都等于整段规整结果', whole === normalizeChatMarkdown(text));
  114. }
  115. // 种子化伪随机模糊(≥500 轮)
  116. {
  117. let seed = 20260920;
  118. const rand = () => {
  119. seed = (seed * 1103515245 + 12345) % 2147483648;
  120. return seed / 2147483648;
  121. };
  122. const alphabet = ['甲', '乙', '丙', '\n', '\n\n', '- ', '* ', '1', '2', '.', ')', ' ', '。', '12', '-'];
  123. let bad = null;
  124. for (let round = 0; round < 500 && !bad; round++) {
  125. const chunks = [];
  126. const count = 1 + Math.floor(rand() * 20);
  127. for (let i = 0; i < count; i++) chunks.push(alphabet[Math.floor(rand() * alphabet.length)]);
  128. const normalizer = createIncrementalMarkdownNormalizer();
  129. let want = '';
  130. let got = '';
  131. for (const chunk of chunks) {
  132. want += chunk;
  133. got += normalizer.push(chunk);
  134. if (!isPrefixWithHeldLineTail(got, normalizeChatMarkdown(want))) {
  135. bad = { round, chunks, want, got };
  136. break;
  137. }
  138. }
  139. if (!bad) {
  140. got += normalizer.finish();
  141. if (got !== normalizeChatMarkdown(want)) bad = { round, chunks, want, got, finish: true };
  142. }
  143. }
  144. check('种子化模糊 500 轮(含 finish)不变式恒成立', !bad, bad);
  145. }
  146. // 扣住的部分必须「有界」:只是当前行的一小段,不能无限积压
  147. {
  148. const normalizer = createIncrementalMarkdownNormalizer();
  149. let shown = '';
  150. let want = '';
  151. for (const chunk of ['第一段\n', '第二段\n', '1', '.', ' ', '甲', '\n', '2', '.', ' ', '乙']) {
  152. want += chunk;
  153. shown += normalizer.push(chunk);
  154. const held = normalizeChatMarkdown(want).slice(shown.length).replace(/^\n+/, '');
  155. if (held.includes('\n') || held.length > 8) {
  156. check(`扣住有界(本条 ${JSON.stringify(chunk)})`, false, { held, shown });
  157. break;
  158. }
  159. }
  160. check('扣住的只是「可能成为列表项」的短前缀(本例 ≤8 字符且不含换行)', true);
  161. }
  162. console.log('\n【2】流式事件解析(形状不符 → null 走旧路径)');
  163. {
  164. const start = parseAnswerStreamEvent({ stream_id: 's1', target: 'summary', sequence: 0 }, 'answer_start');
  165. check('合法的 answer_start', start?.kind === 'start' && start.streamId === 's1' && start.target === 'summary', start);
  166. check('缺 stream_id → null', parseAnswerStreamEvent({ target: 'summary' }, 'answer_start') === null);
  167. check('空 stream_id → null', parseAnswerStreamEvent({ stream_id: ' ', target: 'summary' }, 'answer_delta') === null);
  168. check('target 非法 → null', parseAnswerStreamEvent({ stream_id: 's1', target: 'other' }, 'answer_delta') === null);
  169. check('非对象 → null', parseAnswerStreamEvent(null, 'answer_delta') === null && parseAnswerStreamEvent([1], 'answer_delta') === null);
  170. const delta = parseAnswerStreamEvent({ stream_id: 's1', target: 'answer', sequence: '3', text: 123 }, 'answer_delta');
  171. check('sequence 宽松(字符串也认)', delta?.sequence === 3, delta);
  172. check('text 非字符串按空串', delta?.text === '', delta);
  173. const end = parseAnswerStreamEvent({ stream_id: 's1', target: 'summary', text: 'x', status: '怪值' }, 'answer_end');
  174. check('status 非法 → fallback', end?.kind === 'end' && end.status === 'fallback', end);
  175. const endOk = parseAnswerStreamEvent({ stream_id: 's1', target: 'summary', text: 'x', status: 'completed' }, 'answer_end');
  176. check('status=completed 保留', endOk?.status === 'completed');
  177. const abort = parseAnswerStreamEvent({ stream_id: 's1', target: 'summary', reason: '怪值' }, 'answer_abort');
  178. check('reason 非法 → generation_failed', abort?.kind === 'abort' && abort.reason === 'generation_failed', abort);
  179. const abortOk = parseAnswerStreamEvent({ stream_id: 's1', target: 'summary', reason: 'not_finalized' }, 'answer_abort');
  180. check('reason=not_finalized 保留', abortOk?.reason === 'not_finalized');
  181. check('未知事件名 → null', parseAnswerStreamEvent({ stream_id: 's1', target: 'summary' }, 'answer_what') === null);
  182. }
  183. console.log('\n【3】流状态机');
  184. {
  185. const run = (events) => {
  186. let state = createAnswerStreamState();
  187. const effects = [];
  188. for (const ev of events) {
  189. const r = applyAnswerStreamEvent(state, ev);
  190. state = r.state;
  191. effects.push(r.effect.kind);
  192. }
  193. return { state, effects };
  194. };
  195. const ev = (kind, extra = {}) => ({ kind, streamId: 's1', target: 'summary', sequence: 0, text: '', ...extra });
  196. // 正常流
  197. {
  198. const { state, effects } = run([
  199. ev('start', { sequence: 0 }),
  200. ev('delta', { sequence: 1, text: '你好' }),
  201. ev('delta', { sequence: 2, text: '世界' }),
  202. ev('end', { sequence: 2, text: '你好,世界' }),
  203. ]);
  204. check('正常流:效果序列 start→append→append→replace', JSON.stringify(effects) === JSON.stringify(['start', 'append', 'append', 'replace']), effects);
  205. check('正常流:end 后 text = 最终全文(整块替换,不拼接)', state.text === '你好,世界' && state.completed === true, state);
  206. }
  207. // 序号缺口
  208. {
  209. const { state, effects } = run([
  210. ev('start'),
  211. ev('delta', { sequence: 1, text: 'a' }),
  212. ev('delta', { sequence: 5, text: 'b' }),
  213. ]);
  214. check('序号缺口:仍拼接、hasGap 置位', state.text === 'ab' && state.hasGap === true, state);
  215. check('缺口不阻断渲染(仍是 append)', effects[2] === 'append');
  216. }
  217. // 重复序号
  218. {
  219. const { state } = run([ev('start'), ev('delta', { sequence: 1, text: 'a' }), ev('delta', { sequence: 1, text: 'XX' })]);
  220. check('重复序号忽略(不重复拼接)', state.text === 'a', state);
  221. }
  222. // 过期流
  223. {
  224. const { state, effects } = run([
  225. ev('start'),
  226. ev('delta', { streamId: 'other', sequence: 1, text: 'X' }),
  227. ev('end', { streamId: 'other', text: 'X' }),
  228. ev('abort', { streamId: 'other' }),
  229. ]);
  230. check('过期 streamId 的 delta/end/abort 全部 none', effects.slice(1).every((k) => k === 'none'), effects);
  231. check('过期事件不改状态', state.text === '' && state.completed === false, state);
  232. }
  233. // end 之后的 delta
  234. {
  235. const { state, effects } = run([ev('start'), ev('end', { text: '定稿' }), ev('delta', { sequence: 9, text: '迟到' })]);
  236. check('end 之后的 delta 忽略', effects[2] === 'none' && state.text === '定稿', state);
  237. }
  238. // abort 复位 + 重新 start
  239. {
  240. const { state, effects } = run([
  241. ev('start'),
  242. ev('delta', { sequence: 1, text: '草稿' }),
  243. ev('abort', { reason: 'generation_failed' }),
  244. ev('start', { streamId: 's2' }),
  245. ev('delta', { streamId: 's2', sequence: 1, text: '新正文' }),
  246. ]);
  247. check('abort 复位(discard 效果 + 状态清空后可再 start)', effects.includes('discard') && state.text === '新正文' && state.streamId === 's2', { effects, state });
  248. }
  249. // end 覆盖草稿(fallback 也替换)
  250. {
  251. const { state } = run([
  252. ev('start'),
  253. ev('delta', { sequence: 1, text: '草稿' }),
  254. ev('end', { sequence: 1, text: '最终', status: 'fallback' }),
  255. ]);
  256. check('fallback 的 end 同样整块替换', state.text === '最终' && state.completed === true, state);
  257. }
  258. check('null 事件 → none 且状态不变', (() => {
  259. const s0 = createAnswerStreamState();
  260. const r = applyAnswerStreamEvent(s0, null);
  261. return r.effect.kind === 'none' && r.state === s0;
  262. })());
  263. }
  264. console.log('\n【3b】应用层的内容对齐(end 替换 / abort 撤稿)');
  265. {
  266. const BASE = '<scope title="思考中">\n</scope>\n';
  267. check(
  268. '正常路径:最终文本是已上屏内容的前缀延伸 → 纯尾缀追加(不加空行、不重建)',
  269. applyAnswerStreamEnd(BASE + '一、说明', BASE, '一、说明。二、补充') === BASE + '一、说明。二、补充'
  270. );
  271. check(
  272. '已上屏一长段、最终只多了尾巴 → 不重复整段',
  273. applyAnswerStreamEnd(BASE + '甲', BASE, '甲乙丙') === BASE + '甲乙丙'
  274. );
  275. check(
  276. '罕见分歧(最终文本与增量拼接不一致)→ 整块重建,且基线里的 <scope> 被剥掉(避免卡片闪现)',
  277. applyAnswerStreamEnd(BASE + '草稿', BASE, '完全不同的最终文本') === '完全不同的最终文本'
  278. );
  279. check(
  280. '重建时**基线(流开始前就有的内容)保留**,最终文本接在其后 —— 只有 scope 标记会被剥掉',
  281. applyAnswerStreamEnd('前言', '前言', '最终') === '前言最终'
  282. );
  283. check('abort 撤稿:回到流开始前的内容', applyAnswerStreamAbort(BASE) === BASE);
  284. check('abort 撤稿:空基线返回空串', applyAnswerStreamAbort('') === '');
  285. // answer 开场概述该不该显示(实测:后端把正文作为 summary 流式下发,
  286. // 而 answer 常是同一段开头的**压缩版**——显示会重复开头;但万一是另一段内容就不能丢)
  287. const body =
  288. '我们青浦目前能对上的小微企业扶持,主要分两类:一类是面向特定对象的创业扶持(如带动就业补贴、贷款贴息、留创企业开办资助),另一类是市级中小企业公共资助项目;此外还有几款“小微贷”金融产品。\n\n### 具体政策\n- 创业扶持:带动就业补贴';
  289. // 压缩版的形态与实测一致:**截断在词中间**(正文接着 '###',概述结尾是半个词 '部分候')
  290. const abridged = body.slice(0, 70).replace(/\n+/g, ' ') + ' 部分候…';
  291. check('压缩版开场(截断在词中间,公共前缀 ≥90%)→ 不显示(避免开头重复)', shouldKeepAnswerOverview(abridged, body) === false, abridged);
  292. check('真正的另一段内容(不是前缀)→ 必须保留', shouldKeepAnswerOverview('这是另一段完全不同的开场说明。', body) === true);
  293. check('概述为空 → 不显示', shouldKeepAnswerOverview('', body) === false && shouldKeepAnswerOverview(null, body) === false);
  294. check('流式正文为空 → 保留概述(不能丢字)', shouldKeepAnswerOverview('开场', '') === true);
  295. }
  296. console.log(`\n===== 纯逻辑部分:通过 ${pass} 项,失败 ${fail} 项 =====`);
  297. /* ================================================================== *
  298. * 第 2 部分:协调器级假流(假 SSE 喂真实协调器)
  299. * ================================================================== */
  300. const sse = (events) =>
  301. events
  302. .map(
  303. ([event, data]) =>
  304. `event: ${event}\ndata: ${JSON.stringify({ thread_id: 't1', request_id: 'r1', data })}\n\n`
  305. )
  306. .join('');
  307. const stubFetchWith = (sseText) => {
  308. globalThis.fetch = async () => {
  309. const stream = new ReadableStream({
  310. start(controller) {
  311. controller.enqueue(new TextEncoder().encode(sseText));
  312. controller.close();
  313. },
  314. });
  315. return new Response(stream, { status: 200, headers: { 'Content-Type': 'text/event-stream' } });
  316. };
  317. };
  318. /** 跑一轮,收集各类事件 */
  319. const runTurn = async (events) => {
  320. stubFetchWith(sse(events));
  321. const coordinator = new ApiChatCoordinator({ baseUrl: 'http://stub', threadId: 't1' });
  322. const seen = { messages: [], total: null, streamStart: [], streamEnd: [], streamAbort: [], errors: [], order: [] };
  323. coordinator.addEventListener('message', (m) => {
  324. seen.messages.push(m);
  325. seen.order.push(`message:${m.slice(0, 24)}`);
  326. });
  327. coordinator.addEventListener('totalResponse', (p) => {
  328. seen.total = p;
  329. });
  330. coordinator.addEventListener('answerStreamStart', (p) => {
  331. seen.streamStart.push(p);
  332. seen.order.push('streamStart');
  333. });
  334. coordinator.addEventListener('answerStreamEnd', (p) => {
  335. seen.streamEnd.push(p);
  336. seen.order.push('streamEnd');
  337. });
  338. coordinator.addEventListener('answerStreamAbort', (p) => {
  339. seen.streamAbort.push(p);
  340. seen.order.push('streamAbort');
  341. });
  342. coordinator.addEventListener('error', (p) => {
  343. seen.errors.push(p);
  344. seen.order.push('error');
  345. });
  346. await coordinator.generateAnswer('测试问题');
  347. await new Promise((resolve) => setTimeout(resolve, 80)); // flushContent 会等两帧
  348. return seen;
  349. };
  350. // 应用侧就是 `content += msg`(见 appendAiMessageChunk),**没有分隔符** ——
  351. // 用 '\n\n' 拼会把流式增量切成一堆段落,断言就失真了
  352. const messageText = (seen) => seen.messages.join('');
  353. const countOf = (text, needle) => text.split(needle).length - 1;
  354. const DELTA_TEXTS = ['一、', '第一段说明', '。\n\n', '二、', '第二段说明', '。'];
  355. const FINAL_TEXT = '一、第一段说明。\n\n二、第二段说明。';
  356. const MARKER = '第一段说明';
  357. console.log('\n【4】协调器假流:happy path(带卡片轮形态,含流中 heartbeat)');
  358. {
  359. const seen = await runTurn([
  360. ['accepted', { message: '收到' }],
  361. ['progress', { stage: 'retrieve', message: '正在检索……' }],
  362. ['answer_start', { stream_id: 's1', target: 'summary', sequence: 0, provisional: true }],
  363. ...DELTA_TEXTS.map((text, i) => ['answer_delta', { stream_id: 's1', target: 'summary', sequence: i + 1, text }]),
  364. ['heartbeat', { status: 'processing', message: '仍在处理……', elapsed_seconds: 20 }],
  365. ['answer', { text: FINAL_TEXT }],
  366. ['source', { id: 's1', title: '测试政策', type: 'policy', source: { url: 'https://example.com/p1' } }],
  367. // 有 item 才有 POLICY_TABLE 卡片块(buildPolicyTableContent 是按 items 建的)
  368. ['item', { source_id: 's1', title: '测试政策', card: { name: { text: '测试政策' } } }],
  369. ['answer_end', { stream_id: 's1', target: 'summary', sequence: DELTA_TEXTS.length, text: FINAL_TEXT, operation: 'replace', status: 'completed', provisional: false }],
  370. ['summary', { text: FINAL_TEXT, source_ids: ['s1'], notices: [] }],
  371. ['result', { response: FINAL_TEXT, recommendation: {}, errors: [] }],
  372. ['done', { status: 'completed' }],
  373. ]);
  374. const text = messageText(seen);
  375. if (process.env.DBG) console.log('[dbg] text=' + JSON.stringify(text));
  376. check('answerStreamStart 发出一次', seen.streamStart.length === 1, seen.streamStart);
  377. check('answerStreamEnd 发出一次', seen.streamEnd.length === 1, seen.streamEnd);
  378. check('无 abort', seen.streamAbort.length === 0, seen.streamAbort);
  379. check('scope 卡片只发了一次(流中 heartbeat 不再发)', countOf(text, '<scope') === 1, countOf(text, '<scope'));
  380. check('answer_end 的文本是规整后的全文', seen.streamEnd[0]?.text === normalizeChatMarkdown(FINAL_TEXT), seen.streamEnd[0]?.text);
  381. check(
  382. '正文在 message 通道**只出现一次**(完整 answer/summary 事件没有渲染成第二份)',
  383. countOf(text, MARKER) === 1,
  384. countOf(text, MARKER)
  385. );
  386. check(
  387. '流式增量拼接 === 规整后全文(前缀不变式在真链路成立)',
  388. text.includes(normalizeChatMarkdown(FINAL_TEXT).split('\n')[0]),
  389. normalizeChatMarkdown(FINAL_TEXT).slice(0, 30)
  390. );
  391. check('totalResponse.answer 以最终全文开头', String(seen.total?.answer || '').startsWith('一、第一段说明。'), String(seen.total?.answer || '').slice(0, 40));
  392. check('totalResponse 里含政策卡片与参考资料', String(seen.total?.answer || '').includes('POLICY_TABLE') && String(seen.total?.answer || '').includes('<ref_links>'));
  393. check('卡片/参考资料是在流结束之后补发的', text.indexOf('POLICY_TABLE') > text.indexOf(MARKER));
  394. }
  395. console.log('\n【5】协调器假流:纯正文轮(无卡片、无 answer 事件)不丢正文');
  396. {
  397. const seen = await runTurn([
  398. ['accepted', { message: '收到' }],
  399. ['answer_start', { stream_id: 's2', target: 'summary', sequence: 0 }],
  400. ...DELTA_TEXTS.map((text, i) => ['answer_delta', { stream_id: 's2', target: 'summary', sequence: i + 1, text }]),
  401. ['answer_end', { stream_id: 's2', target: 'summary', sequence: DELTA_TEXTS.length, text: FINAL_TEXT, status: 'completed' }],
  402. ['summary', { text: FINAL_TEXT }],
  403. ['result', { response: FINAL_TEXT }],
  404. ['done', { status: 'completed' }],
  405. ]);
  406. check('totalResponse.answer === 最终全文(没有因为"没有卡片"而丢正文)', String(seen.total?.answer || '').trim() === normalizeChatMarkdown(FINAL_TEXT).trim(), String(seen.total?.answer || ''));
  407. }
  408. console.log('\n【6】协调器假流:answer_end 与增量拼接不同 → 整块替换生效');
  409. {
  410. const seen = await runTurn([
  411. ['answer_start', { stream_id: 's3', target: 'summary', sequence: 0 }],
  412. ['answer_delta', { stream_id: 's3', target: 'summary', sequence: 1, text: '草稿内容' }],
  413. ['answer_end', { stream_id: 's3', target: 'summary', sequence: 1, text: '定稿内容(后端最终校验改过)', status: 'completed' }],
  414. ['summary', { text: '定稿内容(后端最终校验改过)' }],
  415. ['result', { response: '定稿内容(后端最终校验改过)' }],
  416. ['done', { status: 'completed' }],
  417. ]);
  418. check('End.text = 规整后的最终文本(不是增量拼接)', seen.streamEnd[0]?.text === '定稿内容(后端最终校验改过)', seen.streamEnd[0]?.text);
  419. check('totalResponse.answer 用的是最终文本', String(seen.total?.answer || '').includes('定稿内容'), String(seen.total?.answer || ''));
  420. }
  421. console.log('\n【7】协调器假流:后端 error → 撤草稿(abort 先于 error)');
  422. {
  423. const seen = await runTurn([
  424. ['answer_start', { stream_id: 's4', target: 'summary', sequence: 0 }],
  425. ['answer_delta', { stream_id: 's4', target: 'summary', sequence: 1, text: '半截草稿' }],
  426. ['error', { code: 'processing_failed' }],
  427. ]);
  428. const abortAt = seen.order.indexOf('streamAbort');
  429. const errorAt = seen.order.indexOf('error');
  430. check('abort 在 error **之前**发出(否则 hook 已摘任务、abort 会被守卫丢弃)', abortAt > -1 && abortAt < errorAt, seen.order);
  431. check('abort 的 reason = generation_failed', seen.streamAbort[0]?.reason === 'generation_failed', seen.streamAbort[0]);
  432. check('草稿文本未进入 totalResponse', !String(seen.total?.answer || '').includes('半截草稿'), String(seen.total?.answer || ''));
  433. }
  434. console.log('\n【8】协调器假流:断流(EOF 无 done)→ 撤草稿 + 兼容回退');
  435. {
  436. const seen = await runTurn([
  437. ['answer_start', { stream_id: 's5', target: 'summary', sequence: 0 }],
  438. ['answer_delta', { stream_id: 's5', target: 'summary', sequence: 1, text: '半截草稿' }],
  439. ['answer', { text: '完整事件里的正文(断流前的兼容副本)' }],
  440. ]);
  441. check('断流时发出 abort(reason=not_finalized)', seen.streamAbort[0]?.reason === 'not_finalized', seen.streamAbort);
  442. // ⚠️ 分层:协调器**只发信号**,真正把草稿从界面上撤掉的是应用层(hook 调 applyAnswerStreamAbort)。
  443. // 断流路径也不发 totalResponse(既有行为:内容走 message 通道,hook 的 close 监听负责补写 DMS)。
  444. check('兼容路径的完整 answer 事件仍进了 message 通道', messageText(seen).includes('完整事件里的正文'), messageText(seen).slice(0, 80));
  445. check('报的是连接断开', seen.errors.some((e) => e?.code === 'disconnected'), seen.errors);
  446. }
  447. console.log('\n【9】协调器假流:旧后端(无任何流式事件)行为不变');
  448. {
  449. const seen = await runTurn([
  450. ['accepted', { message: '收到' }],
  451. ['progress', { stage: 'retrieve', message: '正在检索……' }],
  452. ['answer', { text: FINAL_TEXT }],
  453. ['source', { id: 's1', title: '测试政策', type: 'policy', source: { url: 'https://example.com/p1' } }],
  454. ['summary', { text: FINAL_TEXT }],
  455. ['result', { response: FINAL_TEXT }],
  456. ['done', { status: 'completed' }],
  457. ]);
  458. const text = messageText(seen);
  459. check('没有流事件、没有 abort', seen.streamStart.length === 0 && seen.streamAbort.length === 0);
  460. check('scope 卡片正常', countOf(text, '<scope') === 1);
  461. check('正文按原路径整段输出(规整后)', text.includes(normalizeChatMarkdown(FINAL_TEXT)));
  462. check('顺序仍是 正文 → 卡片 → 参考资料', text.indexOf('正文') === -1 || text.indexOf('第一段说明') < text.indexOf('POLICY_TABLE') || text.indexOf('POLICY_TABLE') === -1);
  463. }
  464. console.log('\n【10】协调器:流式轮之后跑旧轮,状态不残留');
  465. {
  466. stubFetchWith(
  467. sse([
  468. ['answer_start', { stream_id: 's6', target: 'summary', sequence: 0 }],
  469. ['answer_delta', { stream_id: 's6', target: 'summary', sequence: 1, text: '流式正文' }],
  470. ['answer_end', { stream_id: 's6', target: 'summary', sequence: 1, text: '流式正文', status: 'completed' }],
  471. ['done', { status: 'completed' }],
  472. ])
  473. );
  474. const coordinator = new ApiChatCoordinator({ baseUrl: 'http://stub', threadId: 't1' });
  475. const rounds = [];
  476. coordinator.addEventListener('message', (m) => rounds.push(m));
  477. await coordinator.generateAnswer('第一轮');
  478. await new Promise((r) => setTimeout(r, 80));
  479. const cutoff = rounds.length; // 第一轮结束时的分界
  480. stubFetchWith(
  481. sse([
  482. ['progress', { stage: 'retrieve', message: '正在检索……' }],
  483. ['answer', { text: '第二轮正文' }],
  484. ['done', { status: 'completed' }],
  485. ])
  486. );
  487. await coordinator.generateAnswer('第二轮');
  488. await new Promise((r) => setTimeout(r, 80));
  489. const first = rounds.slice(0, cutoff).join('');
  490. const second = rounds.slice(cutoff).join('');
  491. check('第一轮是流式正文', first.includes('流式正文'), first.slice(0, 60));
  492. check('第二轮回到旧路径:有 scope、正文是整段输出', countOf(second, '<scope') === 1 && second.includes('第二轮正文'), second.slice(0, 80));
  493. }
  494. console.log('\n【11】两个协调器并行流式:互不串台');
  495. {
  496. stubFetchWith(
  497. sse([
  498. ['answer_start', { stream_id: 'A', target: 'summary', sequence: 0 }],
  499. ['answer_delta', { stream_id: 'A', target: 'summary', sequence: 1, text: '甲会话正文' }],
  500. ['answer_end', { stream_id: 'A', target: 'summary', sequence: 1, text: '甲会话正文', status: 'completed' }],
  501. ['done', { status: 'completed' }],
  502. ])
  503. );
  504. const a = new ApiChatCoordinator({ baseUrl: 'http://stub', threadId: 'A' });
  505. const b = new ApiChatCoordinator({ baseUrl: 'http://stub', threadId: 'B' });
  506. const gotA = [];
  507. const gotB = [];
  508. a.addEventListener('message', (m) => gotA.push(m));
  509. b.addEventListener('message', (m) => gotB.push(m));
  510. b.addEventListener('answerStreamAbort', () => gotB.push('ABORT'));
  511. await a.generateAnswer('问 A');
  512. await new Promise((r) => setTimeout(r, 80));
  513. check('B 未收到任何事件(含 abort)', gotB.length === 0, gotB);
  514. check('A 正常收到正文', gotA.join('').includes('甲会话正文'), gotA);
  515. }
  516. console.log(`\n===== 总计:通过 ${pass} 项,失败 ${fail} 项 =====`);
  517. process.exit(fail ? 1 : 0);