fix(qq): prevent turn feed message spam

This commit is contained in:
xmanrui 2026-08-24 00:48:21 +08:00
parent 6ea048b277
commit ade50210f0
10 changed files with 382 additions and 213 deletions

View file

@ -1,9 +1,21 @@
// QQ markdown 回复投递:长文按代码块/表格边界切分,以 msg_type=2 发送,
// QQ markdown 回复投递:长文尽量按结构边界切分,以 msg_type=2 发送,
// 平台拒绝 markdown 时逐条回退纯文本。
const DEFAULT_CHUNK_LIMIT = 4_500;
const CODE_FENCE_OPEN = /^```/;
const GFM_TABLE_LINE = /^\|.+\|$/;
const PASSIVE_REPLY_LIMIT = Object.freeze({ c2c: 4, group: 5 });
const PARTIAL_REPLY_NOTICE = '回答较长,后续内容未能通过 QQ 完整发送,请回复“继续”。';
function safeSliceIndex(value, limit) {
let index = Math.min(limit, value.length);
const before = value.charCodeAt(index - 1);
const after = value.charCodeAt(index);
if (before >= 0xD800 && before <= 0xDBFF && after >= 0xDC00 && after <= 0xDFFF) {
index -= 1;
}
return Math.max(1, index);
}
/**
* 按换行边界切分 Markdown 文本:
@ -44,8 +56,9 @@ export function chunkMarkdownText(text, limit = DEFAULT_CHUNK_LIMIT) {
}
let remaining = block;
while (remaining.length > bound) {
chunks.push(remaining.slice(0, bound));
remaining = remaining.slice(bound);
const index = safeSliceIndex(remaining, bound);
chunks.push(remaining.slice(0, index));
remaining = remaining.slice(index);
}
current = remaining;
};
@ -65,8 +78,9 @@ export function chunkMarkdownText(text, limit = DEFAULT_CHUNK_LIMIT) {
chunks.push(current);
current = '';
}
chunks.push(remaining.slice(0, bound));
remaining = remaining.slice(bound);
const index = safeSliceIndex(remaining, bound);
chunks.push(remaining.slice(0, index));
remaining = remaining.slice(index);
}
appendBlock(remaining);
};
@ -115,11 +129,32 @@ function nextMsgSeq() {
export async function sendMarkdownReply(bot, target, text, { logger } = {}) {
const chunks = chunkMarkdownText(text);
const results = [];
for (const chunk of chunks) {
const passiveLimit = target?.msgId ? PASSIVE_REPLY_LIMIT[target.scope] : null;
const overflow = passiveLimit !== null && chunks.length > passiveLimit;
const passiveContentCount = overflow ? passiveLimit - 1 : chunks.length;
const proactiveTarget = target?.msgId
? { scope: target.scope, targetId: target.targetId }
: target;
let partialNoticeSent = false;
const sendPartialNotice = async () => {
if (partialNoticeSent || !target?.msgId) return;
partialNoticeSent = true;
try {
results.push(await bot.sendText(target, PARTIAL_REPLY_NOTICE));
} catch (error) {
logger?.warn?.('[dsh-im:qq] unable to send partial reply notice:', error);
}
};
for (const [index, chunk] of chunks.entries()) {
const deliveryTarget = overflow && index >= passiveContentCount
? proactiveTarget
: target;
if (typeof bot?.send === 'function') {
try {
results.push(await bot.send({
target,
target: deliveryTarget,
msgType: 2,
markdown: { content: chunk },
extra: { msg_seq: nextMsgSeq() },
@ -129,7 +164,13 @@ export async function sendMarkdownReply(bot, target, text, { logger } = {}) {
logger?.warn?.('[dsh-im:qq] markdown delivery failed; retrying as plain text:', error);
}
}
results.push(await bot.sendText(target, chunk));
try {
results.push(await bot.sendText(deliveryTarget, chunk));
} catch (error) {
if (results.length === 0) throw error;
await sendPartialNotice();
break;
}
}
return results;
}

View file

@ -592,6 +592,7 @@ 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,16 +645,17 @@ export class QqHarnessBridge {
const content = hasImages
? await promptContentForMessage(promptMessage, { signal: this.#signal })
: undefined;
// 过程流:把一次 Turn 的中间过程逐条推送给用户——
// 模型说明文本、每个 Tool call、失败的工具错误详情,最后是完整回答。
let pendingStepText = null;
const pushNotice = async (text) => {
// QQ C2C keeps one stream bubble. Progress is collected but never submitted:
// some clients reject replacing an already visible stream frame, which would
// otherwise leave a stale progress bubble plus a separate fallback answer.
if (message.kind === 'c2c' && target?.msgId && typeof this.#bot.openStream === 'function') {
try {
await this.#bot.sendText(target, text);
stream = this.#bot.openStream({ target });
} catch (error) {
this.#logger.warn?.('[dsh-im:qq] unable to send a turn progress notice:', error);
this.#logger.warn?.('[dsh-im:qq] unable to start a QQ stream; using markdown fallback:', error);
}
};
}
const toolErrors = [];
let answer;
let artifacts = [];
try {
@ -668,25 +670,13 @@ export class QqHarnessBridge {
timeoutMs: this.#replyTimeoutMs,
signal: this.#signal,
control: { owner: this, key },
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;
}
progressMode: 'all',
onUpdate: (update) => {
if (update.error) {
const label = nonEmptyString(update.toolName)
? `Tool call ${update.toolName}` : 'Tool call';
await pushNotice(`${label}\nError: ${update.error}`);
const text = `${label}\nError: ${update.error}`;
toolErrors.push(text);
}
},
onInteraction: (interaction) => this.#handleInteraction(interaction, {
@ -707,22 +697,42 @@ export class QqHarnessBridge {
}
this.#signal?.throwIfAborted();
const answerText = answerTextForDelivery(answer, artifacts);
// 结束前仍有一段未推送的说明文本(且不是最终回答本身)时补发。
if (nonEmptyString(pendingStepText) && pendingStepText.trim() !== answerText.trim()) {
await pushNotice(pendingStepText.trim());
}
const displayAnswer = answerText;
const displayAnswer = toolErrors.length > 0
? `${answerText}\n\n---\n\n${toolErrors.join('\n\n')}`
: answerText;
let textReceipt = null;
let textSendError = null;
try {
const deliveries = await sendMarkdownReply(this.#bot, target, displayAnswer, {
logger: this.#logger,
});
textReceipt = createDeliveryReceipt({
deliveryId: messageId,
presentation: 'qq-text',
providerMessageIds: deliveries.flatMap((delivery) => providerMessageIdsFor(delivery)),
});
let streamFinished = false;
if (stream) {
try {
await stream.update(displayAnswer);
streamFinished = true;
textReceipt = createDeliveryReceipt({
deliveryId: messageId,
presentation: 'qq-text',
providerMessageIds: providerMessageIdsFor(stream),
});
try {
await stream.complete();
} catch (error) {
this.#logger.warn?.('[dsh-im:qq] QQ stream completion failed after visible final content:', error);
}
} catch (error) {
stream.cancel?.();
this.#logger.warn?.('[dsh-im:qq] QQ stream update failed; using markdown fallback:', error);
}
}
if (!streamFinished) {
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);
@ -741,6 +751,11 @@ export class QqHarnessBridge {
return delivery.receipt;
} catch (error) {
if (error?.code === 'turn-stopped') {
try {
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) {
@ -749,6 +764,11 @@ export class QqHarnessBridge {
await this.#state.markSeen(messageId);
return;
}
try {
stream?.cancel?.();
} catch (streamError) {
this.#logger.warn?.('[dsh-im:qq] unable to cancel a failed QQ stream:', streamError);
}
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);

View file

@ -368,7 +368,7 @@ export class HarnessReplyTracker {
return this.#targetTurn;
}
consume(entries) {
consumeAll(entries) {
const updates = [];
// 同一批轮询内的 text 帧只保留最新累积,其余事件逐帧透出,
// 让消费方能按顺序看到每个工具调用与结果。
@ -461,6 +461,10 @@ export class HarnessReplyTracker {
}
return updates;
}
consume(entries) {
return this.consumeAll(entries).at(-1) ?? null;
}
}
export class HarnessRpcError extends Error {
@ -1081,6 +1085,7 @@ export class HarnessClient {
const timeoutMs = options.timeoutMs ?? 600_000;
const signal = options.signal;
const onUpdate = typeof options.onUpdate === 'function' ? options.onUpdate : null;
const progressMode = options.progressMode === 'all' ? 'all' : 'latest';
const onArtifact = typeof options.onArtifact === 'function' ? options.onArtifact : null;
const onInteraction = typeof options.onInteraction === 'function'
? options.onInteraction
@ -1234,9 +1239,10 @@ export class HarnessClient {
this.#consumeInteractionOwnerships(sessionId, history.events ?? []);
if (!wasActive && ownership.active) ownership.reconnect?.();
}
const updates = tracker.consume(history.events ?? []);
const updates = tracker.consumeAll(history.events ?? []);
if (onUpdate) {
for (const update of updates) {
const visibleUpdates = progressMode === 'all' ? updates : updates.slice(-1);
for (const update of visibleUpdates) {
try {
await onUpdate(update);
} catch (error) {