mirror of
https://github.com/hansjone/dsh-im-ops.git
synced 2026-10-09 03:03:24 +08:00
Merge pull request #44 from ferocknew/feat/qqbot-turn-feed
Feat/qqbot turn feed
This commit is contained in:
commit
6ea048b277
9 changed files with 741 additions and 256 deletions
318
lib/index.js
318
lib/index.js
File diff suppressed because one or more lines are too long
135
src/channels/qq/markdown-reply.mjs
Normal file
135
src/channels/qq/markdown-reply.mjs
Normal file
|
|
@ -0,0 +1,135 @@
|
|||
// QQ markdown 回复投递:长文按代码块/表格边界切分,以 msg_type=2 发送,
|
||||
// 平台拒绝 markdown 时逐条回退纯文本。
|
||||
|
||||
const DEFAULT_CHUNK_LIMIT = 4_500;
|
||||
const CODE_FENCE_OPEN = /^```/;
|
||||
const GFM_TABLE_LINE = /^\|.+\|$/;
|
||||
|
||||
/**
|
||||
* 按换行边界切分 Markdown 文本:
|
||||
* - 不在代码块中间断开;
|
||||
* - 不在 GFM 表格中间断开;
|
||||
* - 超长行在 limit 处硬切,避免单行超限无法投递。
|
||||
*/
|
||||
export function chunkMarkdownText(text, limit = DEFAULT_CHUNK_LIMIT) {
|
||||
const value = typeof text === 'string' ? text : '';
|
||||
const bound = Number.isInteger(limit) && limit > 0 ? limit : DEFAULT_CHUNK_LIMIT;
|
||||
if (value.length <= bound) return value ? [value] : [];
|
||||
|
||||
const lines = value.split('\n');
|
||||
const chunks = [];
|
||||
let current = '';
|
||||
let inCodeBlock = false;
|
||||
let tableBuffer = [];
|
||||
|
||||
const appendBlock = (block) => {
|
||||
if (block.length <= bound) {
|
||||
if (!current) {
|
||||
current = block;
|
||||
return;
|
||||
}
|
||||
const candidate = `${current}\n${block}`;
|
||||
if (candidate.length > bound) {
|
||||
chunks.push(current);
|
||||
current = block;
|
||||
} else {
|
||||
current = candidate;
|
||||
}
|
||||
return;
|
||||
}
|
||||
// 超大块:收束当前块后按 bound 硬切,保证每块可投递。
|
||||
if (current) {
|
||||
chunks.push(current);
|
||||
current = '';
|
||||
}
|
||||
let remaining = block;
|
||||
while (remaining.length > bound) {
|
||||
chunks.push(remaining.slice(0, bound));
|
||||
remaining = remaining.slice(bound);
|
||||
}
|
||||
current = remaining;
|
||||
};
|
||||
|
||||
const flushTable = () => {
|
||||
if (tableBuffer.length === 0) return;
|
||||
const block = tableBuffer.join('\n');
|
||||
tableBuffer = [];
|
||||
appendBlock(block);
|
||||
};
|
||||
|
||||
const appendLine = (line) => {
|
||||
let remaining = line;
|
||||
// 超长行先硬切,保证每块不超过 bound。
|
||||
while (remaining.length > bound) {
|
||||
if (current) {
|
||||
chunks.push(current);
|
||||
current = '';
|
||||
}
|
||||
chunks.push(remaining.slice(0, bound));
|
||||
remaining = remaining.slice(bound);
|
||||
}
|
||||
appendBlock(remaining);
|
||||
};
|
||||
|
||||
for (const line of lines) {
|
||||
if (CODE_FENCE_OPEN.test(line)) {
|
||||
flushTable();
|
||||
if (!inCodeBlock && current) {
|
||||
// 代码块开启:先收束当前块,让整个代码块从新块开始。
|
||||
chunks.push(current);
|
||||
current = '';
|
||||
}
|
||||
inCodeBlock = !inCodeBlock;
|
||||
appendLine(line);
|
||||
continue;
|
||||
}
|
||||
if (inCodeBlock) {
|
||||
appendLine(line);
|
||||
continue;
|
||||
}
|
||||
if (GFM_TABLE_LINE.test(line)) {
|
||||
tableBuffer.push(line);
|
||||
continue;
|
||||
}
|
||||
flushTable();
|
||||
appendLine(line);
|
||||
}
|
||||
|
||||
flushTable();
|
||||
if (current) chunks.push(current);
|
||||
return chunks;
|
||||
}
|
||||
|
||||
function nextMsgSeq() {
|
||||
// 与 SDK getNextMsgSeq 相同的随机策略:被动回复同 msg_id 的多条消息
|
||||
// 各自带不同 msg_seq,避免平台去重(错误码 40054005)。
|
||||
const timePart = Date.now() % 100_000_000;
|
||||
const random = Math.floor(Math.random() * 65_536);
|
||||
return (timePart ^ random) % 65_536;
|
||||
}
|
||||
|
||||
/**
|
||||
* 以 markdown(msg_type=2)发送回复;单条被平台拒绝时回退纯文本(msg_type=0)。
|
||||
* 返回每条消息的平台响应,供调用方提取 provider message ids。
|
||||
*/
|
||||
export async function sendMarkdownReply(bot, target, text, { logger } = {}) {
|
||||
const chunks = chunkMarkdownText(text);
|
||||
const results = [];
|
||||
for (const chunk of chunks) {
|
||||
if (typeof bot?.send === 'function') {
|
||||
try {
|
||||
results.push(await bot.send({
|
||||
target,
|
||||
msgType: 2,
|
||||
markdown: { content: chunk },
|
||||
extra: { msg_seq: nextMsgSeq() },
|
||||
}));
|
||||
continue;
|
||||
} catch (error) {
|
||||
logger?.warn?.('[dsh-im:qq] markdown delivery failed; retrying as plain text:', error);
|
||||
}
|
||||
}
|
||||
results.push(await bot.sendText(target, chunk));
|
||||
}
|
||||
return results;
|
||||
}
|
||||
|
|
@ -42,6 +42,7 @@ import {
|
|||
mergeDeliveryReceipts,
|
||||
providerMessageIdsFor,
|
||||
} from '../shared/semantic/delivery.mjs';
|
||||
import { sendMarkdownReply } from './markdown-reply.mjs';
|
||||
|
||||
const INTERACTION_RESOLVED_TEXT = '这个问题已在其他客户端处理,无需再次回答。';
|
||||
const DEFAULT_FILE_UPLOAD_TIMEOUT_MS = 120_000;
|
||||
|
|
@ -591,7 +592,6 @@ export class QqHarnessBridge {
|
|||
const text = promptMessage.content;
|
||||
const hasImages = hasInboundImages(promptMessage);
|
||||
const hasFiles = hasInboundFiles(promptMessage);
|
||||
let stream = null;
|
||||
try {
|
||||
if (!text && !hasImages && !hasFiles) {
|
||||
await this.#bot.sendText(target, '目前支持文字、图片和文件消息。');
|
||||
|
|
@ -644,14 +644,16 @@ export class QqHarnessBridge {
|
|||
const content = hasImages
|
||||
? await promptContentForMessage(promptMessage, { signal: this.#signal })
|
||||
: undefined;
|
||||
let streamFinished = false;
|
||||
if (message.kind === 'c2c' && target?.msgId && typeof this.#bot.openStream === 'function') {
|
||||
// 过程流:把一次 Turn 的中间过程逐条推送给用户——
|
||||
// 模型说明文本、每个 Tool call、失败的工具错误详情,最后是完整回答。
|
||||
let pendingStepText = null;
|
||||
const pushNotice = async (text) => {
|
||||
try {
|
||||
stream = this.#bot.openStream({ target });
|
||||
await this.#bot.sendText(target, text);
|
||||
} catch (error) {
|
||||
this.#logger.warn?.('[dsh-im:qq] unable to start a QQ stream; using a text reply:', error);
|
||||
this.#logger.warn?.('[dsh-im:qq] unable to send a turn progress notice:', error);
|
||||
}
|
||||
}
|
||||
};
|
||||
let answer;
|
||||
let artifacts = [];
|
||||
try {
|
||||
|
|
@ -666,14 +668,27 @@ export class QqHarnessBridge {
|
|||
timeoutMs: this.#replyTimeoutMs,
|
||||
signal: this.#signal,
|
||||
control: { owner: this, key },
|
||||
onUpdate: stream ? async (update) => {
|
||||
const progress = update.type === 'text'
|
||||
? update.text
|
||||
: update.type === 'tool'
|
||||
? `正在使用${update.name}…`
|
||||
: update.text;
|
||||
if (progress) await stream.update(progress);
|
||||
} : undefined,
|
||||
onUpdate: async (update) => {
|
||||
if (update.type === 'text') {
|
||||
// 当前 step 的累积文本:暂存,等下一个工具或结束时一次性推送。
|
||||
pendingStepText = update.text;
|
||||
return;
|
||||
}
|
||||
if (update.type === 'tool') {
|
||||
// 工具调用本身不推送:agent 任务工具调用频繁,逐条推送会刷屏;
|
||||
// 它只作为边界把上一段说明文本定稿推送。失败时随错误详情带出工具名。
|
||||
if (nonEmptyString(pendingStepText)) {
|
||||
await pushNotice(pendingStepText.trim());
|
||||
}
|
||||
pendingStepText = null;
|
||||
return;
|
||||
}
|
||||
if (update.error) {
|
||||
const label = nonEmptyString(update.toolName)
|
||||
? `Tool call ${update.toolName}` : 'Tool call';
|
||||
await pushNotice(`${label}\nError: ${update.error}`);
|
||||
}
|
||||
},
|
||||
onInteraction: (interaction) => this.#handleInteraction(interaction, {
|
||||
key,
|
||||
actor: sender,
|
||||
|
|
@ -691,33 +706,23 @@ export class QqHarnessBridge {
|
|||
]);
|
||||
}
|
||||
this.#signal?.throwIfAborted();
|
||||
const displayAnswer = answerTextForDelivery(answer, artifacts);
|
||||
const answerText = answerTextForDelivery(answer, artifacts);
|
||||
// 结束前仍有一段未推送的说明文本(且不是最终回答本身)时补发。
|
||||
if (nonEmptyString(pendingStepText) && pendingStepText.trim() !== answerText.trim()) {
|
||||
await pushNotice(pendingStepText.trim());
|
||||
}
|
||||
const displayAnswer = answerText;
|
||||
let textReceipt = null;
|
||||
let textSendError = null;
|
||||
try {
|
||||
if (stream) {
|
||||
try {
|
||||
await stream.update(displayAnswer);
|
||||
await stream.complete();
|
||||
streamFinished = true;
|
||||
textReceipt = createDeliveryReceipt({
|
||||
deliveryId: messageId,
|
||||
presentation: 'qq-text',
|
||||
providerMessageIds: providerMessageIdsFor(stream),
|
||||
});
|
||||
} catch (error) {
|
||||
stream.cancel?.();
|
||||
this.#logger.warn?.('[dsh-im:qq] QQ stream finalization failed; using a text reply:', error);
|
||||
}
|
||||
}
|
||||
if (!streamFinished) {
|
||||
const sent = await this.#bot.sendText(target, displayAnswer);
|
||||
textReceipt = createDeliveryReceipt({
|
||||
deliveryId: messageId,
|
||||
presentation: 'qq-text',
|
||||
providerMessageIds: providerMessageIdsFor(sent),
|
||||
});
|
||||
}
|
||||
const deliveries = await sendMarkdownReply(this.#bot, target, displayAnswer, {
|
||||
logger: this.#logger,
|
||||
});
|
||||
textReceipt = createDeliveryReceipt({
|
||||
deliveryId: messageId,
|
||||
presentation: 'qq-text',
|
||||
providerMessageIds: deliveries.flatMap((delivery) => providerMessageIdsFor(delivery)),
|
||||
});
|
||||
} catch (error) {
|
||||
textSendError = error;
|
||||
this.#logger.warn?.('[dsh-im:qq] final text delivery failed; continuing with result files:', error);
|
||||
|
|
@ -736,22 +741,14 @@ export class QqHarnessBridge {
|
|||
return delivery.receipt;
|
||||
} catch (error) {
|
||||
if (error?.code === 'turn-stopped') {
|
||||
if (stream) {
|
||||
try {
|
||||
await stream.cancel?.();
|
||||
} catch (streamError) {
|
||||
this.#logger.warn?.('[dsh-im:qq] unable to cancel a stopped QQ stream:', streamError);
|
||||
}
|
||||
try {
|
||||
await this.#bot.sendText(target, '已停止。');
|
||||
} catch (sendError) {
|
||||
this.#logger.warn?.('[dsh-im:qq] unable to announce a stopped QQ turn:', sendError);
|
||||
}
|
||||
try {
|
||||
await this.#bot.sendText(target, '已停止。');
|
||||
} catch (sendError) {
|
||||
this.#logger.warn?.('[dsh-im:qq] unable to announce a stopped QQ turn:', sendError);
|
||||
}
|
||||
await this.#state.markSeen(messageId);
|
||||
return;
|
||||
}
|
||||
stream?.cancel?.();
|
||||
if (this.#signal?.aborted) return;
|
||||
this.#status.lastError = error?.message ?? String(error);
|
||||
this.#logger.error?.('[dsh-im:qq] failed to process an inbound message:', error);
|
||||
|
|
|
|||
|
|
@ -255,6 +255,21 @@ function assistantMessageText(event) {
|
|||
.trim();
|
||||
}
|
||||
|
||||
function nonEmptyText(value) {
|
||||
return typeof value === 'string' && value.trim() ? value.trim() : null;
|
||||
}
|
||||
|
||||
/** Flatten a tool/result error payload into a displayable one-line reason. */
|
||||
function toolResultErrorText(error) {
|
||||
if (!error || typeof error !== 'object') return null;
|
||||
const message = nonEmptyText(error.message);
|
||||
if (message) return message;
|
||||
const name = nonEmptyText(error.name);
|
||||
const code = nonEmptyText(error.code);
|
||||
if (name || code) return [name ?? 'Error', code].filter(Boolean).join(': ');
|
||||
return null;
|
||||
}
|
||||
|
||||
function consumeInteractionOwnership(ownership, entries) {
|
||||
const ordered = [...entries]
|
||||
.map((entry) => entry?.event ?? entry)
|
||||
|
|
@ -325,6 +340,8 @@ export class HarnessReplyTracker {
|
|||
#latestText = '';
|
||||
#finished = false;
|
||||
#reason = null;
|
||||
#toolNames = new Map();
|
||||
#lastToolName = null;
|
||||
|
||||
constructor({ promptRpcId, afterSeq = -1 }) {
|
||||
this.#promptRpcId = promptRpcId;
|
||||
|
|
@ -352,7 +369,19 @@ export class HarnessReplyTracker {
|
|||
}
|
||||
|
||||
consume(entries) {
|
||||
let update = null;
|
||||
const updates = [];
|
||||
// 同一批轮询内的 text 帧只保留最新累积,其余事件逐帧透出,
|
||||
// 让消费方能按顺序看到每个工具调用与结果。
|
||||
const pushUpdate = (update) => {
|
||||
if (update.type === 'text' && updates.length > 0) {
|
||||
const last = updates[updates.length - 1];
|
||||
if (last.type === 'text') {
|
||||
updates[updates.length - 1] = update;
|
||||
return;
|
||||
}
|
||||
}
|
||||
updates.push(update);
|
||||
};
|
||||
const ordered = [...entries]
|
||||
.map((entry) => entry?.event ?? entry)
|
||||
.filter(Boolean)
|
||||
|
|
@ -394,7 +423,7 @@ export class HarnessReplyTracker {
|
|||
.trim();
|
||||
if (text && text !== this.#latestText) {
|
||||
this.#latestText = text;
|
||||
update = { type: 'text', text };
|
||||
pushUpdate({ type: 'text', text });
|
||||
}
|
||||
continue;
|
||||
}
|
||||
|
|
@ -403,18 +432,34 @@ export class HarnessReplyTracker {
|
|||
const text = assistantMessageText(event);
|
||||
if (text && text !== this.#latestText) {
|
||||
this.#latestText = text;
|
||||
update = { type: 'text', text };
|
||||
pushUpdate({ type: 'text', text });
|
||||
}
|
||||
continue;
|
||||
}
|
||||
|
||||
if (event.type === 'tool/call') {
|
||||
update = { type: 'tool', name: event.data?.name ?? '工具' };
|
||||
const name = nonEmptyText(event.data?.name) ?? '工具';
|
||||
const callId = nonEmptyText(event.data?.callId)
|
||||
?? nonEmptyText(event.data?.subCallId);
|
||||
if (callId) this.#toolNames.set(callId, name);
|
||||
this.#lastToolName = name;
|
||||
pushUpdate({ type: 'tool', name, ...(callId ? { callId } : {}) });
|
||||
} else if (event.type === 'tool/result') {
|
||||
update = { type: 'status', text: '正在整理结果…' };
|
||||
const callId = nonEmptyText(event.data?.message?.source?.callId)
|
||||
?? nonEmptyText(event.data?.callId)
|
||||
?? nonEmptyText(event.data?.subCallId);
|
||||
const toolName = (callId ? this.#toolNames.get(callId) : null)
|
||||
?? this.#lastToolName;
|
||||
const error = toolResultErrorText(event.data?.error);
|
||||
pushUpdate({
|
||||
type: 'status',
|
||||
text: '正在整理结果…',
|
||||
...(toolName ? { toolName } : {}),
|
||||
...(error ? { error } : {}),
|
||||
});
|
||||
}
|
||||
}
|
||||
return update;
|
||||
return updates;
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -1189,12 +1234,14 @@ export class HarnessClient {
|
|||
this.#consumeInteractionOwnerships(sessionId, history.events ?? []);
|
||||
if (!wasActive && ownership.active) ownership.reconnect?.();
|
||||
}
|
||||
const update = tracker.consume(history.events ?? []);
|
||||
if (update && onUpdate) {
|
||||
try {
|
||||
await onUpdate(update);
|
||||
} catch (error) {
|
||||
console.warn(`[${this.#logPrefix}] ignored a progress update failure:`, error.message);
|
||||
const updates = tracker.consume(history.events ?? []);
|
||||
if (onUpdate) {
|
||||
for (const update of updates) {
|
||||
try {
|
||||
await onUpdate(update);
|
||||
} catch (error) {
|
||||
console.warn(`[${this.#logPrefix}] ignored a progress update failure:`, error.message);
|
||||
}
|
||||
}
|
||||
}
|
||||
if (!tracker.finished) continue;
|
||||
|
|
|
|||
|
|
@ -399,7 +399,7 @@ test('reply tracker associates only the Harness turn created by the DingTalk pro
|
|||
data: { turn: 9, step: 0, chunk: { type: 'text-delta', index: 0, text: '钉钉' } },
|
||||
} },
|
||||
]);
|
||||
assert.deepEqual(update, { type: 'text', text: '钉钉' });
|
||||
assert.deepEqual(update, [{ type: 'text', text: '钉钉' }]);
|
||||
tracker.consume([
|
||||
{ event: {
|
||||
seq: 6,
|
||||
|
|
|
|||
|
|
@ -869,7 +869,7 @@ test('HarnessClient delivers an existing file-only Turn directly', async (t) =>
|
|||
test('HarnessReplyTracker correlates the prompt and emits only answer text', () => {
|
||||
const tracker = new HarnessReplyTracker({ promptRpcId: 'prompt-1', afterSeq: 10 });
|
||||
|
||||
assert.equal(tracker.consume([
|
||||
assert.deepEqual(tracker.consume([
|
||||
{ event: { type: 'turn/start', seq: 11, data: { turn: 4 } } },
|
||||
{ event: {
|
||||
type: 'user/message',
|
||||
|
|
@ -882,7 +882,7 @@ test('HarnessReplyTracker correlates the prompt and emits only answer text', ()
|
|||
data: { turn: 4, step: 1, chunk: { type: 'text-delta', index: 0, text: '忽略' } },
|
||||
} },
|
||||
{ event: { type: 'turn/end', seq: 14, data: { turn: 4, reason: { kind: 'completed' } } } },
|
||||
]), null);
|
||||
]), []);
|
||||
assert.equal(tracker.finished, false);
|
||||
|
||||
const first = tracker.consume([
|
||||
|
|
@ -903,7 +903,7 @@ test('HarnessReplyTracker correlates the prompt and emits only answer text', ()
|
|||
data: { turn: 5, step: 1, chunk: { type: 'text-delta', index: 1, text: '深圳' } },
|
||||
} },
|
||||
]);
|
||||
assert.deepEqual(first, { type: 'text', text: '深圳' });
|
||||
assert.deepEqual(first, [{ type: 'text', text: '深圳' }]);
|
||||
|
||||
const second = tracker.consume([
|
||||
{ event: {
|
||||
|
|
@ -917,7 +917,7 @@ test('HarnessReplyTracker correlates the prompt and emits only answer text', ()
|
|||
data: { turn: 5, step: 1, chunk: { type: 'text-delta', index: 1, text: '明天有雨' } },
|
||||
} },
|
||||
]);
|
||||
assert.deepEqual(second, { type: 'text', text: '深圳明天有雨' });
|
||||
assert.deepEqual(second, [{ type: 'text', text: '深圳明天有雨' }]);
|
||||
|
||||
const final = tracker.consume([
|
||||
{ event: {
|
||||
|
|
@ -934,7 +934,7 @@ test('HarnessReplyTracker correlates the prompt and emits only answer text', ()
|
|||
} },
|
||||
{ event: { type: 'turn/end', seq: 21, data: { turn: 5, reason: { kind: 'completed' } } } },
|
||||
]);
|
||||
assert.deepEqual(final, { type: 'text', text: '深圳明天有阵雨。' });
|
||||
assert.deepEqual(final, [{ type: 'text', text: '深圳明天有阵雨。' }]);
|
||||
assert.equal(tracker.finished, true);
|
||||
assert.equal(tracker.answer, '深圳明天有阵雨。');
|
||||
assert.deepEqual(tracker.reason, { kind: 'completed' });
|
||||
|
|
@ -947,9 +947,28 @@ test('HarnessReplyTracker emits tool progress without exposing tool results', ()
|
|||
{ type: 'user/message', seq: 2, data: { source: { rpcId: 'prompt-tool' } } },
|
||||
{ type: 'tool/call', seq: 3, data: { turn: 1, step: 1, name: 'web_search' } },
|
||||
]);
|
||||
assert.deepEqual(update, { type: 'tool', name: 'web_search' });
|
||||
assert.deepEqual(update, [{ type: 'tool', name: 'web_search' }]);
|
||||
|
||||
assert.deepEqual(tracker.consume([
|
||||
{ type: 'tool/result', seq: 4, data: { turn: 1, step: 1, secret: 'not rendered' } },
|
||||
]), { type: 'status', text: '正在整理结果…' });
|
||||
]), [{ type: 'status', text: '正在整理结果…', toolName: 'web_search' }]);
|
||||
});
|
||||
|
||||
test('HarnessReplyTracker keeps every frame of a batched turn in order', () => {
|
||||
const tracker = new HarnessReplyTracker({ promptRpcId: 'prompt-batch' });
|
||||
const updates = tracker.consume([
|
||||
{ type: 'turn/start', seq: 1, data: { turn: 1 } },
|
||||
{ type: 'user/message', seq: 2, data: { source: { rpcId: 'prompt-batch' } } },
|
||||
{ type: 'assistant/chunk', seq: 3, data: { turn: 1, step: 0, chunk: { type: 'text-delta', index: 0, text: '先创建再观察:' } } },
|
||||
{ type: 'tool/call', seq: 4, data: { turn: 1, step: 1, name: 'add_observations' } },
|
||||
{ type: 'tool/result', seq: 5, data: { turn: 1, step: 1, error: { message: 'Status code: 404.' } } },
|
||||
{ type: 'tool/call', seq: 6, data: { turn: 1, step: 2, name: 'create_entities' } },
|
||||
]);
|
||||
assert.deepEqual(updates, [
|
||||
{ type: 'text', text: '先创建再观察:' },
|
||||
{ type: 'tool', name: 'add_observations' },
|
||||
{ type: 'status', text: '正在整理结果…', toolName: 'add_observations', error: 'Status code: 404.' },
|
||||
{ type: 'tool', name: 'create_entities' },
|
||||
]);
|
||||
assert.equal(tracker.answer, '先创建再观察:');
|
||||
});
|
||||
|
|
|
|||
|
|
@ -554,18 +554,12 @@ test('QQ remembers any authorized private inbound as a connection-test target',
|
|||
assert.equal(sent.length, 2);
|
||||
});
|
||||
|
||||
test('QQ private messages stream Harness snapshots and finalize once', async () => {
|
||||
const frames = [];
|
||||
test('QQ pushes the final answer without streaming intermediate text frames', async () => {
|
||||
const sent = [];
|
||||
const seen = new Set();
|
||||
const bridge = new QqHarnessBridge({
|
||||
bot: {
|
||||
sendText: async (_target, text) => sent.push(text),
|
||||
openStream: () => ({
|
||||
update: async (text) => frames.push(text),
|
||||
complete: async () => frames.push('DONE'),
|
||||
cancel() {},
|
||||
}),
|
||||
},
|
||||
ownerUserOpenid: 'owner-openid',
|
||||
harness: {
|
||||
|
|
@ -573,7 +567,9 @@ test('QQ private messages stream Harness snapshots and finalize once', async ()
|
|||
createSession: async () => 'session-new',
|
||||
ensureRunning: async () => true,
|
||||
ask: async (_session, _text, { onUpdate }) => {
|
||||
await onUpdate({ type: 'text', text: '回答中' });
|
||||
// 正文流式帧不逐帧推送;最终回答与已暂存文本相同时不重复补发。
|
||||
await onUpdate({ type: 'text', text: '最终回' });
|
||||
await onUpdate({ type: 'text', text: '最终回答' });
|
||||
return '最终回答';
|
||||
},
|
||||
},
|
||||
|
|
@ -587,26 +583,164 @@ test('QQ private messages stream Harness snapshots and finalize once', async ()
|
|||
});
|
||||
|
||||
await bridge.accept(message());
|
||||
assert.deepEqual(frames, ['回答中', '最终回答', 'DONE']);
|
||||
assert.deepEqual(sent, []);
|
||||
assert.deepEqual(sent, ['最终回答']);
|
||||
assert.equal(seen.has('msg-1'), true);
|
||||
assert.equal(bridge.status.messagesReplied, 1);
|
||||
});
|
||||
|
||||
test('QQ closes an opened progress stream and announces when the Harness turn is stopped', async () => {
|
||||
const fixture = stateFixture([['c2c:owner-openid', 'session-stopped']]);
|
||||
const frames = [];
|
||||
test('QQ pushes one notice per tool call and ignores status frames', async () => {
|
||||
const sent = [];
|
||||
const bridge = new QqHarnessBridge({
|
||||
bot: {
|
||||
sendText: async (_target, text) => sent.push(text),
|
||||
},
|
||||
ownerUserOpenid: 'owner-openid',
|
||||
harness: {
|
||||
sessionExists: async () => true,
|
||||
ask: async (_session, _text, { onUpdate }) => {
|
||||
await onUpdate({ type: 'tool', name: 'bash' });
|
||||
// 工具结束后 Harness 客户端会下发 status 帧“正在整理结果…”。
|
||||
await onUpdate({ type: 'status', text: '正在整理结果…', toolName: 'bash' });
|
||||
await onUpdate({ type: 'tool', name: 'read' });
|
||||
await onUpdate({ type: 'status', text: '正在整理结果…', toolName: 'read' });
|
||||
return '最终回答';
|
||||
},
|
||||
},
|
||||
state: {
|
||||
hasSeen: () => false,
|
||||
markSeen: async () => {},
|
||||
sessionFor: () => 'session-status',
|
||||
setSession: async () => {},
|
||||
clearSession: async () => {},
|
||||
},
|
||||
});
|
||||
|
||||
await bridge.accept(message({ messageId: 'msg-status-frame' }));
|
||||
// 工具调用本身不推送,避免多工具任务刷屏。
|
||||
assert.deepEqual(sent, ['最终回答']);
|
||||
});
|
||||
|
||||
test('QQ pushes a failed tool error and any interim explanation text', async () => {
|
||||
const sent = [];
|
||||
const bridge = new QqHarnessBridge({
|
||||
bot: {
|
||||
sendText: async (_target, text) => sent.push(text),
|
||||
},
|
||||
ownerUserOpenid: 'owner-openid',
|
||||
harness: {
|
||||
sessionExists: async () => true,
|
||||
ask: async (_session, _text, { onUpdate }) => {
|
||||
await onUpdate({ type: 'text', text: '实体不存在,先创建再添加观察:' });
|
||||
await onUpdate({ type: 'tool', name: 'add_observations' });
|
||||
await onUpdate({
|
||||
type: 'status',
|
||||
text: '正在整理结果…',
|
||||
toolName: 'add_observations',
|
||||
error: 'Error calling add_observations. Status code: 404.',
|
||||
});
|
||||
await onUpdate({ type: 'text', text: '改用创建实体的方式:' });
|
||||
await onUpdate({ type: 'tool', name: 'create_entities' });
|
||||
await onUpdate({ type: 'status', text: '正在整理结果…', toolName: 'create_entities' });
|
||||
return '已存入两套记忆。';
|
||||
},
|
||||
},
|
||||
state: {
|
||||
hasSeen: () => false,
|
||||
markSeen: async () => {},
|
||||
sessionFor: () => 'session-tool-error',
|
||||
setSession: async () => {},
|
||||
clearSession: async () => {},
|
||||
},
|
||||
});
|
||||
|
||||
await bridge.accept(message({ messageId: 'msg-tool-error' }));
|
||||
assert.deepEqual(sent, [
|
||||
'实体不存在,先创建再添加观察:',
|
||||
'Tool call add_observations\nError: Error calling add_observations. Status code: 404.',
|
||||
'改用创建实体的方式:',
|
||||
'已存入两套记忆。',
|
||||
]);
|
||||
});
|
||||
|
||||
test('QQ delivers final group answers as markdown messages', async () => {
|
||||
const sentText = [];
|
||||
const markdownCalls = [];
|
||||
const bridge = new QqHarnessBridge({
|
||||
bot: {
|
||||
sendText: async (_target, text) => sentText.push(text),
|
||||
send: async (options) => {
|
||||
markdownCalls.push(options);
|
||||
return { id: `md-${markdownCalls.length}` };
|
||||
},
|
||||
},
|
||||
ownerUserOpenid: 'owner-openid',
|
||||
harness: {
|
||||
sessionExists: async () => true,
|
||||
ask: async () => '## 标题\n\n**加粗**与 `代码`',
|
||||
},
|
||||
state: {
|
||||
hasSeen: () => false,
|
||||
markSeen: async () => {},
|
||||
sessionFor: () => 'session-group',
|
||||
setSession: async () => {},
|
||||
clearSession: async () => {},
|
||||
},
|
||||
});
|
||||
|
||||
await bridge.accept(message({
|
||||
kind: 'group',
|
||||
rawEventType: 'GROUP_AT_MESSAGE_CREATE',
|
||||
groupOpenid: 'group-md',
|
||||
messageId: 'msg-group-md',
|
||||
content: '请回答',
|
||||
replyTarget: { scope: 'group', targetId: 'group-md', msgId: 'msg-group-md' },
|
||||
}));
|
||||
assert.equal(markdownCalls.length, 1);
|
||||
assert.equal(markdownCalls[0].msgType, 2);
|
||||
assert.equal(markdownCalls[0].markdown.content, '## 标题\n\n**加粗**与 `代码`');
|
||||
assert.equal(Number.isInteger(markdownCalls[0].extra.msg_seq), true);
|
||||
assert.deepEqual(sentText, []);
|
||||
assert.equal(bridge.status.messagesReplied, 1);
|
||||
});
|
||||
|
||||
test('QQ falls back to plain text when the platform rejects markdown', async () => {
|
||||
const sentText = [];
|
||||
const bridge = new QqHarnessBridge({
|
||||
bot: {
|
||||
sendText: async (_target, text) => sentText.push(text),
|
||||
send: async () => { throw new Error('markdown rejected'); },
|
||||
},
|
||||
ownerUserOpenid: 'owner-openid',
|
||||
harness: {
|
||||
sessionExists: async () => true,
|
||||
ask: async () => '**回答**内容',
|
||||
},
|
||||
state: {
|
||||
hasSeen: () => false,
|
||||
markSeen: async () => {},
|
||||
sessionFor: () => 'session-fallback',
|
||||
setSession: async () => {},
|
||||
clearSession: async () => {},
|
||||
},
|
||||
logger: { warn() {}, error() {} },
|
||||
});
|
||||
|
||||
await bridge.accept(message({
|
||||
messageId: 'msg-md-fallback',
|
||||
replyTarget: { scope: 'c2c', targetId: 'owner-openid', msgId: 'msg-md-fallback' },
|
||||
}));
|
||||
assert.deepEqual(sentText, ['**回答**内容']);
|
||||
assert.equal(bridge.status.messagesReplied, 1);
|
||||
assert.equal(bridge.status.lastError, null);
|
||||
});
|
||||
|
||||
test('QQ announces a stopped turn after any tool notices already pushed', async () => {
|
||||
const fixture = stateFixture([['c2c:owner-openid', 'session-stopped']]);
|
||||
const sent = [];
|
||||
let cancellations = 0;
|
||||
let loggedErrors = 0;
|
||||
const bridge = new QqHarnessBridge({
|
||||
bot: {
|
||||
sendText: async (_target, text) => sent.push(text),
|
||||
openStream: () => ({
|
||||
update: async (text) => frames.push(text),
|
||||
complete: async () => frames.push('DONE'),
|
||||
cancel: () => { cancellations += 1; },
|
||||
}),
|
||||
},
|
||||
ownerUserOpenid: 'owner-openid',
|
||||
harness: {
|
||||
|
|
@ -624,25 +758,18 @@ test('QQ closes an opened progress stream and announces when the Harness turn is
|
|||
|
||||
await bridge.accept(message({ messageId: 'qq-stopped-stream' }));
|
||||
|
||||
assert.deepEqual(frames, ['正在使用bash…']);
|
||||
assert.equal(cancellations, 1);
|
||||
assert.deepEqual(sent, ['已停止。']);
|
||||
assert.equal(loggedErrors, 0);
|
||||
assert.equal(fixture.seen.has('qq-stopped-stream'), true);
|
||||
});
|
||||
|
||||
test('QQ keeps a stopped turn terminal when stream cleanup and its notice both fail', async () => {
|
||||
test('QQ keeps a stopped turn terminal when its notice cannot be sent', async () => {
|
||||
const fixture = stateFixture([['c2c:owner-openid', 'session-stopped-fallback']]);
|
||||
let warnings = 0;
|
||||
let loggedErrors = 0;
|
||||
const bridge = new QqHarnessBridge({
|
||||
bot: {
|
||||
sendText: async () => { throw new Error('send unavailable'); },
|
||||
openStream: () => ({
|
||||
update: async () => {},
|
||||
complete: async () => {},
|
||||
cancel: () => { throw new Error('cancel unavailable'); },
|
||||
}),
|
||||
},
|
||||
ownerUserOpenid: 'owner-openid',
|
||||
harness: {
|
||||
|
|
@ -662,7 +789,7 @@ test('QQ keeps a stopped turn terminal when stream cleanup and its notice both f
|
|||
|
||||
await bridge.accept(message({ messageId: 'qq-stopped-stream-fallback' }));
|
||||
|
||||
assert.equal(warnings, 2);
|
||||
assert.equal(warnings, 1);
|
||||
assert.equal(loggedErrors, 0);
|
||||
assert.equal(fixture.seen.has('qq-stopped-stream-fallback'), true);
|
||||
});
|
||||
|
|
|
|||
156
test/channels/qq/markdown-reply.test.mjs
Normal file
156
test/channels/qq/markdown-reply.test.mjs
Normal file
|
|
@ -0,0 +1,156 @@
|
|||
import assert from 'node:assert/strict';
|
||||
import test from 'node:test';
|
||||
|
||||
import {
|
||||
chunkMarkdownText,
|
||||
sendMarkdownReply,
|
||||
} from '../../../src/channels/qq/markdown-reply.mjs';
|
||||
|
||||
const target = { scope: 'c2c', targetId: 'user-openid', msgId: 'msg-1' };
|
||||
|
||||
test('chunkMarkdownText keeps short text as a single chunk', () => {
|
||||
assert.deepEqual(chunkMarkdownText('**你好**,世界'), ['**你好**,世界']);
|
||||
});
|
||||
|
||||
test('chunkMarkdownText returns no chunks for empty text', () => {
|
||||
assert.deepEqual(chunkMarkdownText(''), []);
|
||||
assert.deepEqual(chunkMarkdownText(null), []);
|
||||
});
|
||||
|
||||
test('chunkMarkdownText splits long text within the limit', () => {
|
||||
const text = Array.from({ length: 200 }, (_, index) => `第${index}行内容`).join('\n');
|
||||
const chunks = chunkMarkdownText(text, 100);
|
||||
assert.ok(chunks.length > 1);
|
||||
for (const chunk of chunks) {
|
||||
assert.ok(chunk.length <= 100, `chunk exceeds the limit: ${chunk.length}`);
|
||||
}
|
||||
assert.equal(chunks.join('\n'), text);
|
||||
});
|
||||
|
||||
test('chunkMarkdownText does not break inside a code block that fits the limit', () => {
|
||||
const code = '```js\nconsole.log(1);\n```';
|
||||
const text = `${'A'.repeat(80)}\n\n${code}\n\n${'B'.repeat(80)}`;
|
||||
const chunks = chunkMarkdownText(text, 100);
|
||||
const codeChunk = chunks.find((chunk) => chunk.includes('```js'));
|
||||
assert.ok(codeChunk);
|
||||
assert.ok(codeChunk.startsWith(code));
|
||||
});
|
||||
|
||||
test('chunkMarkdownText still chunks an oversized code block within the limit', () => {
|
||||
const code = ['```js', ...Array.from({ length: 30 }, (_, i) => `console.log(${i});`), '```']
|
||||
.join('\n');
|
||||
const chunks = chunkMarkdownText(code, 100);
|
||||
assert.ok(chunks.length > 1);
|
||||
for (const chunk of chunks) {
|
||||
assert.ok(chunk.length <= 100, `chunk exceeds the limit: ${chunk.length}`);
|
||||
}
|
||||
assert.equal(chunks.join('\n'), code);
|
||||
});
|
||||
|
||||
test('chunkMarkdownText keeps a GFM table together', () => {
|
||||
const table = [
|
||||
'| 列一 | 列二 |',
|
||||
'| --- | --- |',
|
||||
'| a | b |',
|
||||
'| c | d |',
|
||||
].join('\n');
|
||||
const chunks = chunkMarkdownText(table, 200);
|
||||
assert.equal(chunks.length, 1);
|
||||
assert.equal(chunks[0], table);
|
||||
});
|
||||
|
||||
test('chunkMarkdownText hard-splits an oversized single line', () => {
|
||||
const line = 'x'.repeat(250);
|
||||
const chunks = chunkMarkdownText(line, 100);
|
||||
assert.deepEqual(chunks, ['x'.repeat(100), 'x'.repeat(100), 'x'.repeat(50)]);
|
||||
});
|
||||
|
||||
test('sendMarkdownReply sends markdown with unique msg_seq per chunk', async () => {
|
||||
const calls = [];
|
||||
const results = await sendMarkdownReply({
|
||||
send: async (options) => {
|
||||
calls.push(options);
|
||||
return { id: `id-${calls.length}` };
|
||||
},
|
||||
sendText: async () => {
|
||||
throw new Error('sendText must not be called when markdown succeeds');
|
||||
},
|
||||
}, target, '# 标题\n\n**加粗**内容');
|
||||
assert.equal(calls.length, 1);
|
||||
assert.equal(calls[0].target, target);
|
||||
assert.equal(calls[0].msgType, 2);
|
||||
assert.equal(calls[0].markdown.content, '# 标题\n\n**加粗**内容');
|
||||
assert.equal(Number.isInteger(calls[0].extra.msg_seq), true);
|
||||
assert.deepEqual(results, [{ id: 'id-1' }]);
|
||||
});
|
||||
|
||||
test('sendMarkdownReply assigns distinct msg_seq values across chunks', async () => {
|
||||
const seqs = [];
|
||||
const bot = {
|
||||
send: async ({ extra }) => {
|
||||
seqs.push(extra.msg_seq);
|
||||
return { id: 'x' };
|
||||
},
|
||||
sendText: async () => { throw new Error('unexpected'); },
|
||||
};
|
||||
const text = Array.from({ length: 1_500 }, (_, index) => `第${index}行`).join('\n');
|
||||
await sendMarkdownReply(bot, target, text);
|
||||
assert.ok(seqs.length > 1);
|
||||
assert.equal(new Set(seqs).size, seqs.length);
|
||||
});
|
||||
|
||||
test('sendMarkdownReply falls back to plain text per chunk on markdown rejection', async () => {
|
||||
const sentText = [];
|
||||
const warnings = [];
|
||||
const results = await sendMarkdownReply({
|
||||
send: async () => {
|
||||
throw new Error('markdown rejected: no permission');
|
||||
},
|
||||
sendText: async (_target, text) => {
|
||||
sentText.push(text);
|
||||
return { id: `text-${sentText.length}` };
|
||||
},
|
||||
}, target, '回答内容', { logger: { warn: (...args) => warnings.push(args) } });
|
||||
assert.deepEqual(sentText, ['回答内容']);
|
||||
assert.deepEqual(results, [{ id: 'text-1' }]);
|
||||
assert.equal(warnings.length, 1);
|
||||
});
|
||||
|
||||
test('sendMarkdownReply uses sendText directly when the bot lacks send()', async () => {
|
||||
const sentText = [];
|
||||
const results = await sendMarkdownReply({
|
||||
sendText: async (_target, text) => {
|
||||
sentText.push(text);
|
||||
return { id: 'plain-1' };
|
||||
},
|
||||
}, target, '没有 send 方法的机器人');
|
||||
assert.deepEqual(sentText, ['没有 send 方法的机器人']);
|
||||
assert.deepEqual(results, [{ id: 'plain-1' }]);
|
||||
});
|
||||
|
||||
test('sendMarkdownReply delivers long answers as multiple markdown chunks', async () => {
|
||||
const markdownChunks = [];
|
||||
const bot = {
|
||||
send: async ({ markdown }) => {
|
||||
markdownChunks.push(markdown.content);
|
||||
return { id: `id-${markdownChunks.length}` };
|
||||
},
|
||||
sendText: async () => { throw new Error('unexpected'); },
|
||||
};
|
||||
const text = Array.from({ length: 600 }, (_, index) => `- 列表项 ${index}`).join('\n');
|
||||
const results = await sendMarkdownReply(bot, target, text);
|
||||
assert.ok(markdownChunks.length > 1);
|
||||
for (const chunk of markdownChunks) {
|
||||
assert.ok(chunk.length <= 4_500);
|
||||
}
|
||||
assert.equal(results.length, markdownChunks.length);
|
||||
assert.equal(markdownChunks.join('\n'), text);
|
||||
});
|
||||
|
||||
test('sendMarkdownReply returns no deliveries for empty text', async () => {
|
||||
const results = await sendMarkdownReply({
|
||||
send: async () => { throw new Error('unexpected'); },
|
||||
sendText: async () => { throw new Error('unexpected'); },
|
||||
}, target, '');
|
||||
assert.deepEqual(results, []);
|
||||
});
|
||||
|
|
@ -513,7 +513,7 @@ test('reply tracker associates only the Harness turn created by the Weixin promp
|
|||
data: { turn: 9, step: 0, chunk: { type: 'text-delta', index: 0, text: '微信' } },
|
||||
} },
|
||||
]);
|
||||
assert.deepEqual(first, { type: 'text', text: '微信' });
|
||||
assert.deepEqual(first, [{ type: 'text', text: '微信' }]);
|
||||
tracker.consume([
|
||||
{ event: {
|
||||
seq: 6,
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue