/**
* 验证 F073「正文真实增量流式输出」的两块:
* 第 1 部分 —— **纯逻辑**(不必联网):增量规整器的前缀不变式 + 流状态机
* 第 2 部分 —— **协调器级假流**:假 SSE 流喂真实协调器,断言事件通道与去重
*
* 为什么规整器要这么较真:它手里攥着一个不变式——
* **任一时刻「已发出的增量拼接」=== `normalizeChatMarkdown(已输入的全部文本)`**
* 只要不变式成立,后端 `answer_end` 的「整块替换」在正常路径上就是纯尾缀追加
* (不会触发整段重建 → 不会让「思考中」卡片闪现重播)。破了它,观感立刻劣化。
*
* 怎么跑(先重新打包,必须带 --loader:.png=dataurl):
* npx esbuild harness/tools/_entry-coordinator.ts --bundle --format=esm \
* --outfile=harness/tools/_coordinator.mjs --alias:@=./src \
* --define:import.meta.env='{}' --loader:.png=dataurl
* node harness/tools/verify-answer-stream.mjs
*/
import {
ApiChatCoordinator,
applyAnswerStreamAbort,
applyAnswerStreamEnd,
applyAnswerStreamEvent,
createAnswerStreamState,
composeAroundStreamBody,
composeStreamedContent,
createIncrementalMarkdownNormalizer,
normalizeChatMarkdown,
parseAnswerStreamEvent,
resolveStreamedBodyAfterEnd,
shouldKeepAnswerOverview,
} from './_coordinator.mjs';
let pass = 0;
let fail = 0;
const check = (name, cond, extra = '') => {
if (cond) {
pass++;
console.log(` ok ${name}`);
} else {
fail++;
console.log(` FAIL ${name}${extra === '' ? '' : ` → ${JSON.stringify(extra)}`}`);
}
};
/* ================================================================== *
* 第 1 部分:纯逻辑
* ================================================================== */
console.log('【1】增量规整器:前缀不变式');
/** 把一串输入分片喂给规整器,返回逐片输出 */
const feed = (chunks) => {
const normalizer = createIncrementalMarkdownNormalizer();
const out = [];
for (const chunk of chunks) out.push(normalizer.push(chunk));
return { out, normalizer };
};
/**
* 不变式(流中):**已发出的文本必须是「规整后全文」的前缀**,
* 且未发出的差额只能是「**给当前行准备的空行分隔符(可能 1~2 个 \n)+ 当前行的尾巴**」——
* 也就是「为了判断该行是不是列表项而短暂扣住的一小段」,行尾一到就释放。
*
* 这个性质正好保证 `answer_end` 的「整块替换」= 纯尾缀追加(不会整段重建)。
*/
const isPrefixWithHeldLineTail = (got, want) => {
if (!want.startsWith(got)) return false;
const held = want.slice(got.length).replace(/^\n+/, ''); // 去掉随行一起扣住的前置空行
return !held.includes('\n');
};
/** 不变式:每喂一片就核对一次 */
const checkInvariant = (label, chunks, { finish = false } = {}) => {
const normalizer = createIncrementalMarkdownNormalizer();
let want = '';
let got = '';
let ok = true;
let at = 0;
for (let i = 0; i < chunks.length; i++) {
want += chunks[i];
got += normalizer.push(chunks[i]);
if (!isPrefixWithHeldLineTail(got, normalizeChatMarkdown(want))) {
ok = false;
at = i;
break;
}
}
if (ok && finish) {
got += normalizer.finish();
if (got !== normalizeChatMarkdown(want)) ok = false;
}
check(
`不变式:${label}`,
ok,
ok ? '' : { 第几片: at, 期望: normalizeChatMarkdown(want).slice(-40), 实际: got.slice(-40) }
);
return ok;
};
checkInvariant('整段一次 push', ['第一行\n第二行\n第三行']);
checkInvariant('逐字 push(含换行)', [...'第一行\n第二行\n1. a\n2. b']);
checkInvariant('双换行(空行)', ['段落一\n\n段落二']);
checkInvariant('连续有序列表不插空行', ['1. 甲\n2. 乙\n3. 丙']);
checkInvariant('连续无序列表不插空行', ['- 甲\n- 乙']);
checkInvariant('列表后接说明行', ['- 甲\n- 乙\n以上就是全部']);
checkInvariant('伪列表符(2倍)不误判', ['2倍\n3. 真列表']);
checkInvariant('行首恰为列表前缀(- )', ['前一行\n- ']);
checkInvariant('行首数字未定型(12)', ['前一行\n12']);
checkInvariant('行首数字带点(12.)', ['前一行\n12.']);
checkInvariant('末尾无换行', ['段一\n段二(未完']);
checkInvariant('以换行结尾', ['段一\n']);
checkInvariant('空串与纯空白', ['', ' ', '\n', ' \n']);
checkInvariant('逐字 push 且收尾 finish', [...'甲\n1. 乙\n丙'], { finish: true });
// 逐字 push 与整段 push 的「拼接结果」必须一致(都 finish 掉扣住的尾巴后比)
{
const text = '第一段\n第二段\n1. 甲\n2. 乙\n收尾说明';
const wholeFeed = feed([text]);
wholeFeed.out.push(wholeFeed.normalizer.finish());
const perCharFeed = feed([...text]);
perCharFeed.out.push(perCharFeed.normalizer.finish());
const whole = wholeFeed.out.join('');
const perChar = perCharFeed.out.join('');
check('逐字 push 与整段 push 输出一致', whole === perChar, { whole, perChar });
check('两者的拼接都等于整段规整结果', whole === normalizeChatMarkdown(text));
}
// 种子化伪随机模糊(≥500 轮)
{
let seed = 20260920;
const rand = () => {
seed = (seed * 1103515245 + 12345) % 2147483648;
return seed / 2147483648;
};
const alphabet = ['甲', '乙', '丙', '\n', '\n\n', '- ', '* ', '1', '2', '.', ')', ' ', '。', '12', '-'];
let bad = null;
for (let round = 0; round < 500 && !bad; round++) {
const chunks = [];
const count = 1 + Math.floor(rand() * 20);
for (let i = 0; i < count; i++) chunks.push(alphabet[Math.floor(rand() * alphabet.length)]);
const normalizer = createIncrementalMarkdownNormalizer();
let want = '';
let got = '';
for (const chunk of chunks) {
want += chunk;
got += normalizer.push(chunk);
if (!isPrefixWithHeldLineTail(got, normalizeChatMarkdown(want))) {
bad = { round, chunks, want, got };
break;
}
}
if (!bad) {
got += normalizer.finish();
if (got !== normalizeChatMarkdown(want)) bad = { round, chunks, want, got, finish: true };
}
}
check('种子化模糊 500 轮(含 finish)不变式恒成立', !bad, bad);
}
// 扣住的部分必须「有界」:只是当前行的一小段,不能无限积压
{
const normalizer = createIncrementalMarkdownNormalizer();
let shown = '';
let want = '';
for (const chunk of ['第一段\n', '第二段\n', '1', '.', ' ', '甲', '\n', '2', '.', ' ', '乙']) {
want += chunk;
shown += normalizer.push(chunk);
const held = normalizeChatMarkdown(want).slice(shown.length).replace(/^\n+/, '');
if (held.includes('\n') || held.length > 8) {
check(`扣住有界(本条 ${JSON.stringify(chunk)})`, false, { held, shown });
break;
}
}
check('扣住的只是「可能成为列表项」的短前缀(本例 ≤8 字符且不含换行)', true);
}
console.log('\n【2】流式事件解析(形状不符 → null 走旧路径)');
{
const start = parseAnswerStreamEvent({ stream_id: 's1', target: 'summary', sequence: 0 }, 'answer_start');
check('合法的 answer_start', start?.kind === 'start' && start.streamId === 's1' && start.target === 'summary', start);
check('缺 stream_id → null', parseAnswerStreamEvent({ target: 'summary' }, 'answer_start') === null);
check('空 stream_id → null', parseAnswerStreamEvent({ stream_id: ' ', target: 'summary' }, 'answer_delta') === null);
check('target 非法 → null', parseAnswerStreamEvent({ stream_id: 's1', target: 'other' }, 'answer_delta') === null);
check('非对象 → null', parseAnswerStreamEvent(null, 'answer_delta') === null && parseAnswerStreamEvent([1], 'answer_delta') === null);
const delta = parseAnswerStreamEvent({ stream_id: 's1', target: 'answer', sequence: '3', text: 123 }, 'answer_delta');
check('sequence 宽松(字符串也认)', delta?.sequence === 3, delta);
check('text 非字符串按空串', delta?.text === '', delta);
const end = parseAnswerStreamEvent({ stream_id: 's1', target: 'summary', text: 'x', status: '怪值' }, 'answer_end');
check('status 非法 → fallback', end?.kind === 'end' && end.status === 'fallback', end);
const endOk = parseAnswerStreamEvent({ stream_id: 's1', target: 'summary', text: 'x', status: 'completed' }, 'answer_end');
check('status=completed 保留', endOk?.status === 'completed');
const abort = parseAnswerStreamEvent({ stream_id: 's1', target: 'summary', reason: '怪值' }, 'answer_abort');
check('reason 非法 → generation_failed', abort?.kind === 'abort' && abort.reason === 'generation_failed', abort);
const abortOk = parseAnswerStreamEvent({ stream_id: 's1', target: 'summary', reason: 'not_finalized' }, 'answer_abort');
check('reason=not_finalized 保留', abortOk?.reason === 'not_finalized');
check('未知事件名 → null', parseAnswerStreamEvent({ stream_id: 's1', target: 'summary' }, 'answer_what') === null);
}
console.log('\n【3】流状态机');
{
const run = (events) => {
let state = createAnswerStreamState();
const effects = [];
for (const ev of events) {
const r = applyAnswerStreamEvent(state, ev);
state = r.state;
effects.push(r.effect.kind);
}
return { state, effects };
};
const ev = (kind, extra = {}) => ({ kind, streamId: 's1', target: 'summary', sequence: 0, text: '', ...extra });
// 正常流
{
const { state, effects } = run([
ev('start', { sequence: 0 }),
ev('delta', { sequence: 1, text: '你好' }),
ev('delta', { sequence: 2, text: '世界' }),
ev('end', { sequence: 2, text: '你好,世界' }),
]);
check('正常流:效果序列 start→append→append→replace', JSON.stringify(effects) === JSON.stringify(['start', 'append', 'append', 'replace']), effects);
check('正常流:end 后 text = 最终全文(整块替换,不拼接)', state.text === '你好,世界' && state.completed === true, state);
}
// 序号缺口
{
const { state, effects } = run([
ev('start'),
ev('delta', { sequence: 1, text: 'a' }),
ev('delta', { sequence: 5, text: 'b' }),
]);
check('序号缺口:仍拼接、hasGap 置位', state.text === 'ab' && state.hasGap === true, state);
check('缺口不阻断渲染(仍是 append)', effects[2] === 'append');
}
// 重复序号
{
const { state } = run([ev('start'), ev('delta', { sequence: 1, text: 'a' }), ev('delta', { sequence: 1, text: 'XX' })]);
check('重复序号忽略(不重复拼接)', state.text === 'a', state);
}
// 过期流
{
const { state, effects } = run([
ev('start'),
ev('delta', { streamId: 'other', sequence: 1, text: 'X' }),
ev('end', { streamId: 'other', text: 'X' }),
ev('abort', { streamId: 'other' }),
]);
check('过期 streamId 的 delta/end/abort 全部 none', effects.slice(1).every((k) => k === 'none'), effects);
check('过期事件不改状态', state.text === '' && state.completed === false, state);
}
// end 之后的 delta
{
const { state, effects } = run([ev('start'), ev('end', { text: '定稿' }), ev('delta', { sequence: 9, text: '迟到' })]);
check('end 之后的 delta 忽略', effects[2] === 'none' && state.text === '定稿', state);
}
// abort 复位 + 重新 start
{
const { state, effects } = run([
ev('start'),
ev('delta', { sequence: 1, text: '草稿' }),
ev('abort', { reason: 'generation_failed' }),
ev('start', { streamId: 's2' }),
ev('delta', { streamId: 's2', sequence: 1, text: '新正文' }),
]);
check('abort 复位(discard 效果 + 状态清空后可再 start)', effects.includes('discard') && state.text === '新正文' && state.streamId === 's2', { effects, state });
}
// end 覆盖草稿(fallback 也替换)
{
const { state } = run([
ev('start'),
ev('delta', { sequence: 1, text: '草稿' }),
ev('end', { sequence: 1, text: '最终', status: 'fallback' }),
]);
check('fallback 的 end 同样整块替换', state.text === '最终' && state.completed === true, state);
}
check('null 事件 → none 且状态不变', (() => {
const s0 = createAnswerStreamState();
const r = applyAnswerStreamEvent(s0, null);
return r.effect.kind === 'none' && r.state === s0;
})());
}
console.log('\n【3b】应用层的内容对齐(end 替换 / abort 撤稿)');
{
const BASE = '\n\n';
// 新签名:(当前内容, 流式正文, 最终正文) —— 用「流式正文」做后缀匹配,不按基线长度切
check(
'正常路径:最终文本以增量拼接为前缀 → 纯尾缀追加(不加空行、不重建)',
applyAnswerStreamEnd(BASE + '一、说明', '一、说明', '一、说明。二、补充') === BASE + '一、说明。二、补充'
);
check(
'已上屏一长段、最终只多了尾巴 → 不重复整段',
applyAnswerStreamEnd(BASE + '甲', '甲', '甲乙丙') === BASE + '甲乙丙'
);
check(
'罕见分歧(最终文本与增量拼接不一致)→ 整块重建,且基线里的 被剥掉(避免卡片闪现)',
applyAnswerStreamEnd(BASE + '草稿', '草稿', '完全不同的最终文本') === '完全不同的最终文本'
);
check(
'重建时**基线(流开始前就有的内容)保留**,最终文本接在其后 —— 只有 scope 标记会被剥掉',
applyAnswerStreamEnd('前言草稿', '草稿', '最终') === '前言最终'
);
check(
'兜底:上屏内容与流式正文对不上(被别的路径改过)→ 不切,直接把最终正文接上(不丢字)',
applyAnswerStreamEnd('被改过的内容', '对不上的正文', '最终') === '被改过的内容最终'
);
check('abort 撤稿:从上屏内容里去掉流式正文', applyAnswerStreamAbort(BASE + '草稿', '草稿') === BASE);
check('abort 兜底:对不上就不动(宁可多留也不切坏)', applyAnswerStreamAbort('别的内容', '对不上') === '别的内容');
check('abort:空正文返回原内容', applyAnswerStreamAbort('原内容', '') === '原内容');
// answer 开场概述该不该显示(实测:后端把正文作为 summary 流式下发,
// 而 answer 常是同一段开头的**压缩版**——显示会重复开头;但万一是另一段内容就不能丢)
const body =
'我们青浦目前能对上的小微企业扶持,主要分两类:一类是面向特定对象的创业扶持(如带动就业补贴、贷款贴息、留创企业开办资助),另一类是市级中小企业公共资助项目;此外还有几款“小微贷”金融产品。\n\n### 具体政策\n- 创业扶持:带动就业补贴';
// 压缩版的形态与实测一致:**截断在词中间**(正文接着 '###',概述结尾是半个词 '部分候')
const abridged = body.slice(0, 70).replace(/\n+/g, ' ') + ' 部分候…';
check('压缩版开场(截断在词中间,公共前缀 ≥90%)→ 不显示(避免开头重复)', shouldKeepAnswerOverview(abridged, body) === false, abridged);
check('真正的另一段内容(不是前缀)→ 必须保留', shouldKeepAnswerOverview('这是另一段完全不同的开场说明。', body) === true);
check('概述为空 → 不显示', shouldKeepAnswerOverview('', body) === false && shouldKeepAnswerOverview(null, body) === false);
check('流式正文为空 → 保留概述(不能丢字)', shouldKeepAnswerOverview('开场', '') === true);
// 流式轮的最终组装:基线 + 卡片 + 正文 + 尾部(卡片在正文**之上**)
const SCOPE = '\n\n';
check(
'组装顺序:基线 → leading(卡片)→ 正文 → trailing(参考资料)',
composeStreamedContent(SCOPE + '前言', '【卡片】', '正文正文', '【参考资料】') === '前言\n\n【卡片】\n\n正文正文\n\n【参考资料】',
composeStreamedContent(SCOPE + '前言', '【卡片】', '正文正文', '【参考资料】')
);
check(
'组装时剥掉基线里的 (否则卡片会闪现重播)',
!composeStreamedContent(SCOPE + '前言', '', '正文', '').includes(' {
const onScreen = applyAnswerStreamEnd('', DELTA, FINAL); // end 把尾巴接上
const body = resolveStreamedBodyAfterEnd(DELTA, FINAL); // hook 同步流式正文
const { next, matched } = composeAroundStreamBody(onScreen, body, '【卡片】', '【参考】');
return matched === true && next === '【卡片】\n\n' + FINAL + '\n\n【参考】' && next.split(FINAL).length - 1 === 1;
})()
);
// 兜底:流式正文与上屏内容错位(本 bug 的原状)——**绝不能把正文再插一份**
check(
'兜底:正文不在末尾但在上屏内容里 → 以它为锚点重排,不重复',
(() => {
const { next, matched } = composeAroundStreamBody(FINAL, DELTA, '【卡片】', '【参考】');
return matched === false && next.split(DELTA).length - 1 === 1 && next.startsWith('【卡片】') && next.endsWith('【参考】');
})()
);
check(
'兜底:锚点之后的内容(answer_end 的尾巴)留在正文与参考资料之间',
(() => {
const { next } = composeAroundStreamBody('基线' + FINAL, DELTA, '【卡片】', '【参考】');
return next === '基线\n\n【卡片】\n\n' + DELTA + '\n\n' + FINAL.slice(DELTA.length) + '\n\n【参考】';
})()
);
check(
'兜底:正文确实不在上屏内容里 → 补插(不丢字优先)',
composeAroundStreamBody('别的内容', DELTA, '', '').next === '别的内容\n\n' + DELTA
);
// 流式正文为空:没有可对齐的正文,退化成「现有内容 + 卡片 + 参考资料」(与改动前一致)
const emptyBody = composeAroundStreamBody('已有内容', '', '【卡片】', '【参考】');
check(
'流式正文为空 → 退化为「现有内容 → 卡片 → 参考资料」',
emptyBody.matched === true && emptyBody.next === '已有内容\n\n【卡片】\n\n【参考】',
emptyBody
);
}
console.log(`\n===== 纯逻辑部分:通过 ${pass} 项,失败 ${fail} 项 =====`);
/* ================================================================== *
* 第 2 部分:协调器级假流(假 SSE 喂真实协调器)
* ================================================================== */
const sse = (events) =>
events
.map(
([event, data]) =>
`event: ${event}\ndata: ${JSON.stringify({ thread_id: 't1', request_id: 'r1', data })}\n\n`
)
.join('');
const stubFetchWith = (sseText) => {
globalThis.fetch = async () => {
const stream = new ReadableStream({
start(controller) {
controller.enqueue(new TextEncoder().encode(sseText));
controller.close();
},
});
return new Response(stream, { status: 200, headers: { 'Content-Type': 'text/event-stream' } });
};
};
/** 跑一轮,收集各类事件 */
const runTurn = async (events) => {
stubFetchWith(sse(events));
const coordinator = new ApiChatCoordinator({ baseUrl: 'http://stub', threadId: 't1' });
const seen = { messages: [], total: null, streamStart: [], streamEnd: [], streamAbort: [], compose: [], errors: [], order: [] };
coordinator.addEventListener('message', (m) => {
seen.messages.push(m);
seen.order.push(`message:${m.slice(0, 24)}`);
});
coordinator.addEventListener('totalResponse', (p) => {
seen.total = p;
});
coordinator.addEventListener('answerStreamStart', (p) => {
seen.streamStart.push(p);
seen.order.push('streamStart');
});
coordinator.addEventListener('answerStreamEnd', (p) => {
seen.streamEnd.push(p);
seen.order.push('streamEnd');
});
coordinator.addEventListener('answerStreamAbort', (p) => {
seen.streamAbort.push(p);
seen.order.push('streamAbort');
});
coordinator.addEventListener('answerStreamCompose', (p) => {
seen.compose.push(p);
seen.order.push('streamCompose');
});
coordinator.addEventListener('error', (p) => {
seen.errors.push(p);
seen.order.push('error');
});
await coordinator.generateAnswer('测试问题');
await new Promise((resolve) => setTimeout(resolve, 80)); // flushContent 会等两帧
return seen;
};
// 应用侧就是 `content += msg`(见 appendAiMessageChunk),**没有分隔符** ——
// 用 '\n\n' 拼会把流式增量切成一堆段落,断言就失真了
const messageText = (seen) => seen.messages.join('');
const countOf = (text, needle) => text.split(needle).length - 1;
/** 数**卡片块**个数(`POLICY_TABLE` 在同一个块里出现两次:``) */
const countCards = (text) => countOf(text, '