From 422fd34aa476b04b5291f2acc0520cc3c5d8034a Mon Sep 17 00:00:00 2001 From: xmanrui <841206367@qq.com> Date: Wed, 26 Aug 2026 17:31:20 +0800 Subject: [PATCH] fix(wecom): separate thinking from answers --- src/channels/wecom/wecom-bridge.mjs | 59 +++++++++++--- test/channels/wecom/bridge.test.mjs | 114 ++++++++++++++++++++++++---- 2 files changed, 147 insertions(+), 26 deletions(-) diff --git a/src/channels/wecom/wecom-bridge.mjs b/src/channels/wecom/wecom-bridge.mjs index c2c0454..6f88d78 100644 --- a/src/channels/wecom/wecom-bridge.mjs +++ b/src/channels/wecom/wecom-bridge.mjs @@ -296,12 +296,23 @@ function splitUtf8(text, maxBytes = MAX_REPLY_BYTES) { return chunks; } -function progressText(update) { - if (update?.type === 'text') return update.text; +function thinkingProgressText(update) { if (update?.type === 'tool') return t('正在使用{name}…', { name: update.name }); return update?.text; } +function streamContent(thinkingText, answerText = '', { finish = false } = {}) { + const thinking = String(thinkingText ?? '') + .replace(/<\/?think>/gi, '') + .trim(); + const answer = String(answerText ?? '').trim(); + if (!thinking) return answer; + const thinkBlock = finish || answer + ? `${thinking}` + : `${thinking}`; + return answer ? `${thinkBlock}\n${answer}` : thinkBlock; +} + function artifactFailureText(fileName, error) { const name = String(fileName ?? t('结果文件')).replace(/[\r\n]+/g, ' ').trim() || t('结果文件'); @@ -867,6 +878,8 @@ export class WecomHarnessBridge { const key = conversationKey(frame); let streamId = null; let streamStarted = false; + let streamThinkingText = t('正在思考中…'); + let streamAnswerText = ''; let batchSettled = batchSubmission === null; let promptRecorded = false; try { @@ -920,7 +933,12 @@ export class WecomHarnessBridge { streamId = this.#generateReqId('stream'); try { - await this.#client.replyStream(frame, streamId, t('正在思考中…'), false); + await this.#client.replyStream( + frame, + streamId, + streamContent(streamThinkingText), + false, + ); streamStarted = true; } catch (error) { this.#logger.warn?.('[dsh-im:wecom] unable to start a stream; using an active reply:', error); @@ -945,8 +963,15 @@ export class WecomHarnessBridge { control: { owner: this, key }, onUpdate: streamStarted && typeof this.#client.replyStreamNonBlocking === 'function' ? async (update) => { - const progress = splitUtf8(progressText(update))[0]; - if (progress) await this.#client.replyStreamNonBlocking(frame, streamId, progress, false); + if (update?.type === 'text') { + streamAnswerText = update.text; + } else { + streamThinkingText = thinkingProgressText(update) || streamThinkingText; + } + const preview = splitUtf8( + streamContent(streamThinkingText, streamAnswerText), + )[0]; + if (preview) await this.#client.replyStreamNonBlocking(frame, streamId, preview, false); } : undefined, onInteraction: (interaction) => this.#handleInteraction(interaction, { @@ -966,18 +991,20 @@ export class WecomHarnessBridge { this.#signal?.throwIfAborted(); const displayAnswer = answerTextForDelivery(answer, artifacts); - const chunks = splitUtf8(displayAnswer); + const streamChunks = splitUtf8( + streamContent(streamThinkingText, displayAnswer, { finish: true }), + ); let finalSent = false; let textReceipt = null; let textSendError = null; try { - if (streamStarted && chunks.length > 0) { + if (streamStarted && streamChunks.length > 0) { try { const providerMessageIds = []; - const streamed = await this.#client.replyStream(frame, streamId, chunks[0], true); + const streamed = await this.#client.replyStream(frame, streamId, streamChunks[0], true); const streamedMessageId = providerMessageId(streamed); if (streamedMessageId) providerMessageIds.push(streamedMessageId); - for (const chunk of chunks.slice(1)) { + for (const chunk of streamChunks.slice(1)) { const sent = await this.#client.sendMessage( chatId, { msgtype: 'markdown', markdown: { content: chunk } }, @@ -1040,7 +1067,12 @@ export class WecomHarnessBridge { } if (error?.code === 'turn-stopped') { if (streamStarted && streamId) { - await this.#client.replyStream(frame, streamId, t('已停止。'), true) + await this.#client.replyStream( + frame, + streamId, + streamContent(streamThinkingText, t('已停止。'), { finish: true }), + true, + ) .catch(() => undefined); } if (!promptRecorded) await this.#state.markSeen(messageId); @@ -1063,7 +1095,12 @@ export class WecomHarnessBridge { : errorText; try { if (streamStarted && streamId) { - await this.#client.replyStream(frame, streamId, visibleError, true); + await this.#client.replyStream( + frame, + streamId, + streamContent(streamThinkingText, visibleError, { finish: true }), + true, + ); } else { await this.#sendImmediate(frame, chatId, visibleError); } diff --git a/test/channels/wecom/bridge.test.mjs b/test/channels/wecom/bridge.test.mjs index 790e3bc..07a253d 100644 --- a/test/channels/wecom/bridge.test.mjs +++ b/test/channels/wecom/bridge.test.mjs @@ -22,6 +22,12 @@ const PNG_1X1 = Buffer.from( 'iVBORw0KGgoAAAANSUhEUgAAAAEAAAABCAQAAAC1HAwCAAAAC0lEQVR42mNk+A8AAQUBAScY42YAAAAASUVORK5CYII=', 'base64', ); +const INITIAL_THINKING_TEXT = '正在思考中…'; +const INITIAL_THINKING_STREAM = `${INITIAL_THINKING_TEXT}`; + +function streamedAnswer(answer, thinking = INITIAL_THINKING_TEXT) { + return `${thinking}\n${answer}`; +} test('Enterprise WeChat native image adapter uploads and sends an image media message', async () => { const calls = []; @@ -491,6 +497,7 @@ test('Enterprise WeChat messages stream Harness progress and finalize once', asy ensureRunning: async () => true, ask: async (_session, _text, { onUpdate }) => { await onUpdate({ type: 'tool', name: '网页搜索' }); + await onUpdate({ type: 'status', text: '正在整理结果…' }); await onUpdate({ type: 'text', text: '回答中' }); return '最终回答'; }, @@ -500,16 +507,89 @@ test('Enterprise WeChat messages stream Harness progress and finalize once', asy await bridge.accept(frame()); assert.deepEqual(replies, [ - { streamId: 'stream-1', content: '正在思考中…', finish: false }, - { streamId: 'stream-1', content: '正在使用网页搜索…', finish: false }, - { streamId: 'stream-1', content: '回答中', finish: false }, - { streamId: 'stream-1', content: '最终回答', finish: true }, + { streamId: 'stream-1', content: INITIAL_THINKING_STREAM, finish: false }, + { streamId: 'stream-1', content: '正在使用网页搜索…', finish: false }, + { streamId: 'stream-1', content: '正在整理结果…', finish: false }, + { + streamId: 'stream-1', + content: streamedAnswer('回答中', '正在整理结果…'), + finish: false, + }, + { + streamId: 'stream-1', + content: streamedAnswer('最终回答', '正在整理结果…'), + finish: true, + }, ]); assert.deepEqual(active, []); assert.equal(store.seen.has('msg-1'), true); assert.equal(bridge.status.messagesReplied, 1); }); +test('Enterprise WeChat keeps the thinking block closed when a long answer is split', async () => { + const replies = []; + const active = []; + const answer = '结'.repeat(7_000); + const bridge = new WecomHarnessBridge({ + client: { + replyStream: async (_frame, streamId, content, finish) => { + replies.push({ streamId, content, finish }); + }, + replyStreamNonBlocking: async () => {}, + sendMessage: async (chatId, body) => active.push({ chatId, body }), + }, + generateStreamId: () => 'stream-long', + harness: { + sessionExists: async () => true, + ask: async () => answer, + }, + state: state(), + }); + + await bridge.accept(frame({ msgid: 'wecom-long-answer' })); + + const final = replies.find(({ finish }) => finish); + assert.ok(final); + assert.match(final.content, /^正在思考中…<\/think>\n/); + assert.ok(Buffer.byteLength(final.content) <= 18_000); + assert.ok(active.length > 0); + assert.equal( + [final.content, ...active.map(({ body }) => body.markdown.content)].join(''), + streamedAnswer(answer), + ); + assert.equal(active.every(({ body }) => Buffer.byteLength(body.markdown.content) <= 18_000), true); +}); + +test('Enterprise WeChat falls back to a plain active reply when the stream cannot start', async () => { + const active = []; + let streamAttempts = 0; + const bridge = new WecomHarnessBridge({ + client: { + replyStream: async () => { + streamAttempts += 1; + throw new Error('stream unavailable'); + }, + replyStreamNonBlocking: async () => assert.fail('an unopened stream cannot be updated'), + sendMessage: async (chatId, body) => active.push({ chatId, body }), + }, + generateStreamId: () => 'stream-unavailable', + harness: { + sessionExists: async () => true, + ask: async () => '最终回答', + }, + state: state(), + logger: { warn() {} }, + }); + + await bridge.accept(frame({ msgid: 'wecom-stream-unavailable' })); + + assert.equal(streamAttempts, 1); + assert.deepEqual(active, [{ + chatId: 'member-1', + body: { msgtype: 'markdown', markdown: { content: '最终回答' } }, + }]); +}); + test('Enterprise WeChat downloads an image with its AES key and submits structured content once', async () => { const transport = testClient(); const downloads = []; @@ -556,7 +636,7 @@ test('Enterprise WeChat downloads an image with its AES key and submits structur }, ], }]); - assert.equal(transport.streamed.at(-1).content, '图片识别完成'); + assert.equal(transport.streamed.at(-1).content, streamedAnswer('图片识别完成')); }); test('Enterprise WeChat preserves mixed-message text and image order', async () => { @@ -664,7 +744,7 @@ test('Enterprise WeChat bridge hands its prefetched native file to the current H name: 'file', loaded: { data: bytes, name: '企微报告.docx' }, }]); - assert.equal(transport.streamed.at(-1).content, '文件已收到'); + assert.equal(transport.streamed.at(-1).content, streamedAnswer('文件已收到')); }); test('Enterprise WeChat starts image download before an earlier conversation turn finishes', async () => { @@ -766,7 +846,7 @@ test('Enterprise WeChat bounds prefetched image memory while a conversation is q assert.equal(downloads.length, 4); assert.equal(prompts.length, 5); assert.equal(transport.streamed.some(({ content }) => ( - content.startsWith('当前待处理图片较多,请稍后重新发送。') + content.includes('当前待处理图片较多,请稍后重新发送。') && /错误码:INPUT_INVALID;参考号:MF-[A-F0-9]{8}$/.test(content) )), true); }); @@ -819,10 +899,11 @@ test('Enterprise WeChat finalizes an existing progress stream when Harness fails await bridge.accept(frame()); assert.deepEqual(replies[0], { - streamId: 'stream-failure', content: '正在思考中…', finish: false, + streamId: 'stream-failure', content: INITIAL_THINKING_STREAM, finish: false, }); assert.equal(replies[1].streamId, 'stream-failure'); assert.equal(replies[1].finish, true); + assert.match(replies[1].content, /^正在思考中…<\/think>\n/); assert.match(replies[1].content, /任务未完成,暂时无法确定原因/); assert.match(replies[1].content, /错误码:INTERNAL_UNKNOWN;参考号:MF-[A-F0-9]{8}$/); assert.equal(store.seen.has('msg-1'), true); @@ -950,7 +1031,7 @@ test('an Enterprise WeChat answer bypasses the original conversation queue', asy answer: { answers: [{ id: 'environment', selected: ['测试环境'] }] }, }, }]); - assert.equal(transport.streamed.at(-1).content, '你选择了:测试环境'); + assert.equal(transport.streamed.at(-1).content, streamedAnswer('你选择了:测试环境')); assert.equal(transport.streamed.at(-1).finish, true); }); @@ -1018,7 +1099,7 @@ test('an answer waits for the first Enterprise WeChat question delivery acknowle assert.deepEqual(responses[0].value.answer.answers, [ { id: 'environment', selected: ['测试环境'] }, ]); - assert.equal(streamed.at(-1).content, '首问回答完成'); + assert.equal(streamed.at(-1).content, streamedAnswer('首问回答完成')); }); test('pending Enterprise WeChat questions stay isolated by conversation', async () => { @@ -1791,7 +1872,7 @@ test('Enterprise WeChat returns the authoritative receipt and one safe notice wh const receipt = await bridge.accept(frame({ msgid: 'wecom-all-fail' })); - assert.deepEqual(finalStreamTexts, ['文字结果']); + assert.deepEqual(finalStreamTexts, [streamedAnswer('文字结果')]); assert.equal(attemptedActiveTexts.length, 2, 'must not append a generic error after the safe notice'); assert.equal(visibleActiveTexts.length, 1); assert.match(visibleActiveTexts[0], /暂时无法读取或准备发送.*仍可访问/); @@ -1846,8 +1927,11 @@ test('Enterprise WeChat keeps the generic error when no answer or file failure n await bridge.accept(frame({ msgid: 'wecom-no-visible-failure' })); assert.equal(attemptedActiveTexts.length, 2); - assert.equal(finalStreamTexts[0], '文字结果'); - assert.match(finalStreamTexts[1], /^回复发送结果未能确认/); + assert.equal(finalStreamTexts[0], streamedAnswer('文字结果')); + assert.match( + finalStreamTexts[1], + /^正在思考中…<\/think>\n回复发送结果未能确认/, + ); assert.match(finalStreamTexts[1], /错误码:CHANNEL_DELIVERY_UNCERTAIN;参考号:MF-[A-F0-9]{8}$/); }); @@ -1956,7 +2040,7 @@ test('Enterprise WeChat uses a neutral final text for a file-only Turn', async ( await bridge.accept(frame({ msgid: 'wecom-file-only' })); - assert.deepEqual(finalTexts, ['结果文件已生成。']); + assert.deepEqual(finalTexts, [streamedAnswer('结果文件已生成。')]); assert.deepEqual(files, [{ chatId: 'member-1', type: 'file', mediaId: 'media-only' }]); }); @@ -2021,7 +2105,7 @@ test('Enterprise WeChat batch input collects ten texts and submits one ordered H assert.match(prompts[0], /\[消息 1\]\n企微内容 1/); assert.match(prompts[0], /\[消息 10\]\n企微内容 10/); assert.doesNotMatch(prompts[0], /不会收录/); - assert.equal(transport.streamed.at(-1).content, '批量完成'); + assert.equal(transport.streamed.at(-1).content, streamedAnswer('批量完成')); }); test('Enterprise WeChat rejects batch commands in groups without invoking Harness', async () => {