mirror of
https://github.com/hansjone/dsh-im-ops.git
synced 2026-10-09 00:33:20 +08:00
feat(qq): stream the full turn feed as discrete messages
Show what the agent is actually doing during a turn, the way Claude Code does: interim explanation text, every Tool call, and the error detail when a tool fails, each as its own pushed message before the final markdown answer. - HarnessReplyTracker.consume now returns every frame of a polling batch in order (same-batch text frames collapse to the latest) instead of dropping all but the last event, which silently ate tool-call frames that shared a batch with their result - tool frames carry callId; status frames carry the correlated toolName and a flattened error text parsed from the tool/result error payload (message, or name: code) - the QQ bridge holds interim step text until the next tool call (or the end of the turn) pushes it, emits 'Tool call <name>' per call, emits 'Tool call <name>\nError: <reason>' for failures, and skips re-sending interim text that already equals the final answer Other channels keep their existing behaviour: the status frame still reads '正在整理结果…' and they ignore the new optional fields.
This commit is contained in:
parent
6171ec0dc1
commit
e1bc9dd0d0
7 changed files with 325 additions and 190 deletions
311
lib/index.js
311
lib/index.js
File diff suppressed because one or more lines are too long
|
|
@ -595,6 +595,16 @@ export class QqHarnessBridge {
|
|||
const content = hasImages
|
||||
? await promptContentForMessage(promptMessage, { signal: this.#signal })
|
||||
: undefined;
|
||||
// 过程流:把一次 Turn 的中间过程逐条推送给用户——
|
||||
// 模型说明文本、每个 Tool call、失败的工具错误详情,最后是完整回答。
|
||||
let pendingStepText = null;
|
||||
const pushNotice = async (text) => {
|
||||
try {
|
||||
await this.#bot.sendText(target, text);
|
||||
} catch (error) {
|
||||
this.#logger.warn?.('[dsh-im:qq] unable to send a turn progress notice:', error);
|
||||
}
|
||||
};
|
||||
let answer;
|
||||
let artifacts = [];
|
||||
try {
|
||||
|
|
@ -609,14 +619,24 @@ export class QqHarnessBridge {
|
|||
timeoutMs: this.#replyTimeoutMs,
|
||||
signal: this.#signal,
|
||||
control: { owner: this, key },
|
||||
// 中间过程逐条推送:每个工具调用发一条独立消息;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);
|
||||
if (update.type === 'text') {
|
||||
// 当前 step 的累积文本:暂存,等下一个工具或结束时一次性推送。
|
||||
pendingStepText = update.text;
|
||||
return;
|
||||
}
|
||||
if (update.type === 'tool') {
|
||||
if (nonEmptyString(pendingStepText)) {
|
||||
await pushNotice(pendingStepText.trim());
|
||||
}
|
||||
pendingStepText = null;
|
||||
await pushNotice(`Tool call ${update.name}`);
|
||||
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, {
|
||||
|
|
@ -635,7 +655,12 @@ 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 {
|
||||
|
|
|
|||
|
|
@ -251,6 +251,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)
|
||||
|
|
@ -321,6 +336,8 @@ export class HarnessReplyTracker {
|
|||
#latestText = '';
|
||||
#finished = false;
|
||||
#reason = null;
|
||||
#toolNames = new Map();
|
||||
#lastToolName = null;
|
||||
|
||||
constructor({ promptRpcId, afterSeq = -1 }) {
|
||||
this.#promptRpcId = promptRpcId;
|
||||
|
|
@ -348,7 +365,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)
|
||||
|
|
@ -390,7 +419,7 @@ export class HarnessReplyTracker {
|
|||
.trim();
|
||||
if (text && text !== this.#latestText) {
|
||||
this.#latestText = text;
|
||||
update = { type: 'text', text };
|
||||
pushUpdate({ type: 'text', text });
|
||||
}
|
||||
continue;
|
||||
}
|
||||
|
|
@ -399,18 +428,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;
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -1150,12 +1195,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,
|
||||
|
|
|
|||
|
|
@ -713,7 +713,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',
|
||||
|
|
@ -726,7 +726,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([
|
||||
|
|
@ -747,7 +747,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: {
|
||||
|
|
@ -761,7 +761,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: {
|
||||
|
|
@ -778,7 +778,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' });
|
||||
|
|
@ -791,9 +791,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, '先创建再观察:');
|
||||
});
|
||||
|
|
|
|||
|
|
@ -432,8 +432,9 @@ test('QQ pushes the final answer without streaming intermediate text frames', as
|
|||
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 '最终回答';
|
||||
},
|
||||
},
|
||||
|
|
@ -464,9 +465,9 @@ test('QQ pushes one notice per tool call and ignores status frames', async () =>
|
|||
ask: async (_session, _text, { onUpdate }) => {
|
||||
await onUpdate({ type: 'tool', name: 'bash' });
|
||||
// 工具结束后 Harness 客户端会下发 status 帧“正在整理结果…”。
|
||||
await onUpdate({ type: 'status', text: '正在整理结果…' });
|
||||
await onUpdate({ type: 'status', text: '正在整理结果…', toolName: 'bash' });
|
||||
await onUpdate({ type: 'tool', name: 'read' });
|
||||
await onUpdate({ type: 'status', text: '正在整理结果…' });
|
||||
await onUpdate({ type: 'status', text: '正在整理结果…', toolName: 'read' });
|
||||
return '最终回答';
|
||||
},
|
||||
},
|
||||
|
|
@ -480,7 +481,49 @@ test('QQ pushes one notice per tool call and ignores status frames', async () =>
|
|||
});
|
||||
|
||||
await bridge.accept(message({ messageId: 'msg-status-frame' }));
|
||||
assert.deepEqual(sent, ['正在使用bash…', '正在使用read…', '最终回答']);
|
||||
assert.deepEqual(sent, ['Tool call bash', 'Tool call read', '最终回答']);
|
||||
});
|
||||
|
||||
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: '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',
|
||||
'Tool call add_observations\nError: Error calling add_observations. Status code: 404.',
|
||||
'Tool call create_entities',
|
||||
'已存入两套记忆。',
|
||||
]);
|
||||
});
|
||||
|
||||
test('QQ delivers final group answers as markdown messages', async () => {
|
||||
|
|
@ -579,7 +622,7 @@ test('QQ announces a stopped turn after any tool notices already pushed', async
|
|||
|
||||
await bridge.accept(message({ messageId: 'qq-stopped-stream' }));
|
||||
|
||||
assert.deepEqual(sent, ['正在使用bash…', '已停止。']);
|
||||
assert.deepEqual(sent, ['Tool call bash', '已停止。']);
|
||||
assert.equal(loggedErrors, 0);
|
||||
assert.equal(fixture.seen.has('qq-stopped-stream'), true);
|
||||
});
|
||||
|
|
|
|||
|
|
@ -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