mirror of
https://github.com/hansjone/dsh-im-ops.git
synced 2026-10-09 03:03:24 +08:00
fix(wecom): separate thinking from answers
This commit is contained in:
parent
2d395860b3
commit
422fd34aa4
2 changed files with 147 additions and 26 deletions
|
|
@ -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
|
||||
? `<think>${thinking}</think>`
|
||||
: `<think>${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);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -22,6 +22,12 @@ const PNG_1X1 = Buffer.from(
|
|||
'iVBORw0KGgoAAAANSUhEUgAAAAEAAAABCAQAAAC1HAwCAAAAC0lEQVR42mNk+A8AAQUBAScY42YAAAAASUVORK5CYII=',
|
||||
'base64',
|
||||
);
|
||||
const INITIAL_THINKING_TEXT = '正在思考中…';
|
||||
const INITIAL_THINKING_STREAM = `<think>${INITIAL_THINKING_TEXT}`;
|
||||
|
||||
function streamedAnswer(answer, thinking = INITIAL_THINKING_TEXT) {
|
||||
return `<think>${thinking}</think>\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: '<think>正在使用网页搜索…', finish: false },
|
||||
{ streamId: 'stream-1', content: '<think>正在整理结果…', 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>正在思考中…<\/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>正在思考中…<\/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>正在思考中…<\/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 () => {
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue