feat(qq): push progress and answers as separate messages

Replace the in-place streaming reply (openStream replace-mode frames)
with discrete pushed messages, so the chat no longer flickers while
the turn is running:

- each tool call pushes one notice message (正在使用X…); status frames
  such as '正在整理结果…' and per-token text frames are not pushed
- the final answer arrives as one new markdown message
- a stopped turn still announces '已停止。'; no stream session is
  created, so there is nothing to cancel or finalize
This commit is contained in:
jonah.fu 2026-08-23 22:31:19 +08:00
parent 6b176d9053
commit 6171ec0dc1
3 changed files with 101 additions and 160 deletions

File diff suppressed because one or more lines are too long

View file

@ -543,7 +543,6 @@ export class QqHarnessBridge {
const promptMessage = qqInboundMessage(message, { fetchImpl: this.#fetchImpl });
const text = promptMessage.content;
const hasImages = hasInboundImages(promptMessage);
let stream = null;
try {
if (!text && !hasImages) {
await this.#bot.sendText(target, '目前支持文字和图片消息。');
@ -596,14 +595,6 @@ 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') {
try {
stream = this.#bot.openStream({ target });
} catch (error) {
this.#logger.warn?.('[dsh-im:qq] unable to start a QQ stream; using a text reply:', error);
}
}
let answer;
let artifacts = [];
try {
@ -618,16 +609,16 @@ export class QqHarnessBridge {
timeoutMs: this.#replyTimeoutMs,
signal: this.#signal,
control: { owner: this, key },
onUpdate: stream ? async (update) => {
// status 帧(如工具结束后的“正在整理结果…”)不推送:
// 保持上一帧(工具名或已生成文本),避免中间过程反复闪提示。
const progress = update.type === 'text'
? update.text
: update.type === 'tool'
? `正在使用${update.name}…`
: null;
if (progress) await stream.update(progress);
} : undefined,
// 中间过程逐条推送:每个工具调用发一条独立消息;text 帧(正文流)
// 不逐帧推送,最终答案一次性发送;status 帧忽略。
onUpdate: async (update) => {
if (update.type !== 'tool' || !nonEmptyString(update.name)) return;
try {
await this.#bot.sendText(target, `正在使用${update.name}…`);
} catch (error) {
this.#logger.warn?.('[dsh-im:qq] unable to send a tool progress notice:', error);
}
},
onInteraction: (interaction) => this.#handleInteraction(interaction, {
key,
actor: sender,
@ -648,31 +639,14 @@ export class QqHarnessBridge {
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 deliveries = await sendMarkdownReply(this.#bot, target, displayAnswer, {
logger: this.#logger,
});
textReceipt = createDeliveryReceipt({
deliveryId: messageId,
presentation: 'qq-text',
providerMessageIds: deliveries.flatMap((delivery) => providerMessageIdsFor(delivery)),
});
}
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);
@ -691,22 +665,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);

View file

@ -419,18 +419,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: {
@ -438,6 +432,7 @@ 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: '回答中' });
return '最终回答';
},
@ -452,22 +447,16 @@ 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 progress stream keeps the last frame when a tool finishes', async () => {
const frames = [];
test('QQ pushes one notice per tool call and ignores status frames', async () => {
const sent = [];
const bridge = new QqHarnessBridge({
bot: {
sendText: async () => {},
openStream: () => ({
update: async (text) => frames.push(text),
complete: async () => frames.push('DONE'),
cancel() {},
}),
sendText: async (_target, text) => sent.push(text),
},
ownerUserOpenid: 'owner-openid',
harness: {
@ -491,7 +480,7 @@ test('QQ progress stream keeps the last frame when a tool finishes', async () =>
});
await bridge.accept(message({ messageId: 'msg-status-frame' }));
assert.deepEqual(frames, ['正在使用bash…', '正在使用read…', '最终回答', 'DONE']);
assert.deepEqual(sent, ['正在使用bash…', '正在使用read…', '最终回答']);
});
test('QQ delivers final group answers as markdown messages', async () => {
@ -566,20 +555,13 @@ test('QQ falls back to plain text when the platform rejects markdown', async ()
assert.equal(bridge.status.lastError, null);
});
test('QQ closes an opened progress stream and announces when the Harness turn is stopped', async () => {
test('QQ announces a stopped turn after any tool notices already pushed', async () => {
const fixture = stateFixture([['c2c:owner-openid', 'session-stopped']]);
const frames = [];
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: {
@ -597,25 +579,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.deepEqual(sent, ['正在使用bash…', '已停止。']);
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: {
@ -635,7 +610,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);
});