mirror of
https://github.com/openclaw/openclaw.git
synced 2026-08-25 03:45:46 -06:00
refactor(auto-reply): distill verbose commentary lane wiring
This commit is contained in:
@@ -1038,8 +1038,6 @@ async function processDiscordMessageInner(
|
||||
},
|
||||
onItemEvent: async (payload) => {
|
||||
if (payload.kind === "preamble") {
|
||||
// While the durable verbose commentary lane is active, the ephemeral
|
||||
// draft yields its commentary lines so commentary renders once.
|
||||
if (verboseProgressActive()) {
|
||||
return;
|
||||
}
|
||||
|
||||
@@ -1028,10 +1028,9 @@ export const dispatchTelegramMessage = async ({
|
||||
});
|
||||
let finalAnswerDeliveryStarted = false;
|
||||
let finalAnswerDelivered = false;
|
||||
// While the durable verbose commentary lane is active, the ephemeral draft
|
||||
// yields its commentary lines so commentary is not rendered in both lanes.
|
||||
// Tool/plan lines keep the draft: they have no durable counterpart in
|
||||
// streamed runs, so yielding them would lose information.
|
||||
// While the durable verbose lane is active, the ephemeral draft yields its
|
||||
// commentary lines so they render once. Tool/plan status lines keep the
|
||||
// draft: they have no durable counterpart in streamed runs.
|
||||
let verboseProgressActive: () => boolean = () => false;
|
||||
const pushStreamToolProgress = async (
|
||||
line?: string | ChannelProgressDraftLine,
|
||||
@@ -2041,10 +2040,8 @@ export const dispatchTelegramMessage = async ({
|
||||
}
|
||||
if (segment.lane === "answer" && info.kind === "tool") {
|
||||
if (verboseProgressActive()) {
|
||||
// The durable verbose progress lane owns tool/commentary
|
||||
// payloads for this run: deliver as a real standalone
|
||||
// message instead of diverting into the streaming draft,
|
||||
// which is ephemeral and discarded at final.
|
||||
// Durable lane owns tool payloads: send standalone instead
|
||||
// of diverting into the draft, which is discarded at final.
|
||||
if (
|
||||
await sendPayload(
|
||||
applyTextToPayload(effectivePayload, segment.update.text),
|
||||
|
||||
@@ -224,10 +224,8 @@ function createToolEventBridge(params: {
|
||||
/**
|
||||
* Tracks CLI tool start/result events and renders the same durable tool
|
||||
* summaries the embedded runner emits: a formatToolAggregate line per result
|
||||
* (with args-derived meta captured at start), plus the tool output block when
|
||||
* full verbose output is enabled. The CLI parser emits tool result events, but
|
||||
* until now they were dropped at the bridge, so CLI-backed runs had no durable
|
||||
* tool record under verbose while embedded runs did.
|
||||
* (args-derived meta captured at start), plus the output block under full
|
||||
* verbose. Keeps CLI runs at tool-summary parity with embedded runs.
|
||||
*/
|
||||
export function createCliToolSummaryTracker(params: {
|
||||
detailMode?: "explain" | "raw";
|
||||
|
||||
@@ -2143,13 +2143,9 @@ export async function dispatchReplyFromConfig(
|
||||
return !reply.hasMedia && !hasExecApprovalPayload(payload);
|
||||
};
|
||||
// Durable inter-tool commentary lane: with verbose progress on, preamble
|
||||
// items become standalone progress messages (like tool summaries) instead of
|
||||
// living only in ephemeral channel drafts. The latest text per item id is
|
||||
// buffered so producers that re-emit snapshots for the same item send one
|
||||
// message; the buffer flushes when the producer moves on (a different item,
|
||||
// a tool event, a block reply, or the final reply). The delivery helpers
|
||||
// reference gates declared later in this scope; they only run once the
|
||||
// resolver is executing, well after scope initialization completes.
|
||||
// items become standalone progress messages like tool summaries. The latest
|
||||
// text per item id is buffered (snapshot producers re-emit the same item)
|
||||
// and flushed when the producer moves on, always before the final reply.
|
||||
let pendingCommentaryProgress: { itemId?: string; text: string } | null = null;
|
||||
const deliverCommentaryProgressMessage = async (text: string) => {
|
||||
if (!shouldSendToolSummaries() || shouldSuppressProgressDelivery()) {
|
||||
@@ -2217,8 +2213,7 @@ export async function dispatchReplyFromConfig(
|
||||
}
|
||||
};
|
||||
throwIfFinalDeliveryAborted();
|
||||
// Drain any trailing commentary before the final reply so the durable
|
||||
// progress lane completes ahead of the answer.
|
||||
// Trailing commentary must land ahead of the final answer.
|
||||
await flushPendingCommentaryProgress();
|
||||
throwIfFinalDeliveryAborted();
|
||||
const sourceReplyTranscriptMirror =
|
||||
@@ -2663,31 +2658,34 @@ export async function dispatchReplyFromConfig(
|
||||
// classification in the CLI runners is wired once at run start, so a
|
||||
// mid-run verbose toggle cannot move inter-tool commentary between lanes.
|
||||
const deliverStandaloneCommentaryProgress = shouldEmitVerboseProgress();
|
||||
const coreOwnedOnItemEvent = (() => {
|
||||
const forwardItemEvent = wrapProgressCallback(params.replyOptions?.onItemEvent, {
|
||||
forwardWhenSourceDeliverySuppressed: true,
|
||||
requiresToolSummaryVisibility: true,
|
||||
waitForDirectBlockReplyDelivery: true,
|
||||
onForward: (payload) => {
|
||||
if (hasFailedProgressStatus(payload)) {
|
||||
markVisibleToolErrorProgress();
|
||||
const forwardItemEvent = wrapProgressCallback(params.replyOptions?.onItemEvent, {
|
||||
forwardWhenSourceDeliverySuppressed: true,
|
||||
requiresToolSummaryVisibility: true,
|
||||
waitForDirectBlockReplyDelivery: true,
|
||||
onForward: (payload) => {
|
||||
if (hasFailedProgressStatus(payload)) {
|
||||
markVisibleToolErrorProgress();
|
||||
}
|
||||
},
|
||||
});
|
||||
// Item-event presence gates CLI commentary classification downstream, so
|
||||
// the handler exists exactly when verbose buffers it or a channel consumes it.
|
||||
const onItemEvent =
|
||||
deliverStandaloneCommentaryProgress || forwardItemEvent
|
||||
? async (payload: Parameters<NonNullable<GetReplyOptions["onItemEvent"]>>[0]) => {
|
||||
if (isDispatchOperationAborted()) {
|
||||
return;
|
||||
}
|
||||
if (!forwardItemEvent) {
|
||||
// The wrapped forwarder marks progress itself when present.
|
||||
markProgress();
|
||||
}
|
||||
if (deliverStandaloneCommentaryProgress && payload.kind === "preamble") {
|
||||
await noteCommentaryProgress(payload);
|
||||
}
|
||||
await forwardItemEvent?.(payload);
|
||||
}
|
||||
},
|
||||
});
|
||||
return async (payload: Parameters<NonNullable<GetReplyOptions["onItemEvent"]>>[0]) => {
|
||||
if (isDispatchOperationAborted()) {
|
||||
return;
|
||||
}
|
||||
if (!forwardItemEvent) {
|
||||
// The wrapped forwarder marks progress itself when present.
|
||||
markProgress();
|
||||
}
|
||||
if (payload.kind === "preamble") {
|
||||
await noteCommentaryProgress(payload);
|
||||
}
|
||||
await forwardItemEvent?.(payload);
|
||||
};
|
||||
})();
|
||||
: undefined;
|
||||
// Let draft-rendering channels yield their ephemeral commentary lines while
|
||||
// the durable verbose commentary lane is delivering the same content.
|
||||
params.replyOptions?.onVerboseProgressVisibility?.(
|
||||
@@ -2732,26 +2730,13 @@ export async function dispatchReplyFromConfig(
|
||||
requiresToolSummaryVisibility: true,
|
||||
waitForDirectBlockReplyDelivery: true,
|
||||
onForward: async () => {
|
||||
// A tool is starting: the preceding commentary block is final, so
|
||||
// flush it ahead of the tool's own progress rendering.
|
||||
// Commentary precedes the tool that follows it.
|
||||
await flushPendingCommentaryProgress();
|
||||
},
|
||||
}),
|
||||
onItemEvent: deliverStandaloneCommentaryProgress
|
||||
? coreOwnedOnItemEvent
|
||||
: wrapProgressCallback(params.replyOptions?.onItemEvent, {
|
||||
forwardWhenSourceDeliverySuppressed: true,
|
||||
requiresToolSummaryVisibility: true,
|
||||
waitForDirectBlockReplyDelivery: true,
|
||||
onForward: (payload) => {
|
||||
if (hasFailedProgressStatus(payload)) {
|
||||
markVisibleToolErrorProgress();
|
||||
}
|
||||
},
|
||||
}),
|
||||
commentaryProgressEnabled: deliverStandaloneCommentaryProgress
|
||||
? true
|
||||
: params.replyOptions?.commentaryProgressEnabled,
|
||||
onItemEvent,
|
||||
commentaryProgressEnabled:
|
||||
deliverStandaloneCommentaryProgress || params.replyOptions?.commentaryProgressEnabled,
|
||||
onCommandOutput: wrapProgressCallback(params.replyOptions?.onCommandOutput, {
|
||||
forwardWhenSourceDeliverySuppressed: true,
|
||||
requiresToolSummaryVisibility: true,
|
||||
@@ -2783,8 +2768,7 @@ export async function dispatchReplyFromConfig(
|
||||
return;
|
||||
}
|
||||
markInboundDedupeReplayUnsafe();
|
||||
// The tool finished: any buffered commentary preceded this tool
|
||||
// call, so it must land before the tool summary.
|
||||
// Buffered commentary preceded this tool; land it before the summary.
|
||||
await flushPendingCommentaryProgress();
|
||||
if (!suppressAutomaticSourceDelivery && shouldSendToolSummaries()) {
|
||||
await onToolResultFromReplyOptions?.(payload);
|
||||
@@ -2943,8 +2927,7 @@ export async function dispatchReplyFromConfig(
|
||||
) {
|
||||
markInboundDedupeReplayUnsafe();
|
||||
}
|
||||
// Assistant text is moving on: buffered commentary preceded this
|
||||
// block, so deliver it first to keep the lanes in order.
|
||||
// Buffered commentary preceded this block; deliver it first.
|
||||
await flushPendingCommentaryProgress();
|
||||
if (suppressDelivery) {
|
||||
return;
|
||||
@@ -3097,8 +3080,8 @@ export async function dispatchReplyFromConfig(
|
||||
}
|
||||
|
||||
const replies = replyResult ? (Array.isArray(replyResult) ? replyResult : [replyResult]) : [];
|
||||
// Backstop: a run can end without a visible final reply (silent or
|
||||
// streaming-delivered turns); trailing commentary must still land.
|
||||
// Backstop: silent/streaming-delivered turns end without a visible final
|
||||
// reply; trailing commentary must still land.
|
||||
await flushPendingCommentaryProgress();
|
||||
const beforeAgentRunBlocked = replies.some(
|
||||
(reply) => getReplyPayloadMetadata(reply)?.beforeAgentRunBlocked === true,
|
||||
|
||||
@@ -827,6 +827,29 @@ export function createFollowupRunner(params: {
|
||||
const notifyUserMessagePersisted = () => {
|
||||
queuedUserMessagePersistedAcrossFallback = true;
|
||||
};
|
||||
// Shared by the embedded onToolResult callback and the CLI tool
|
||||
// summary tracker so both runners deliver identical durable summaries.
|
||||
const deliverFollowupToolSummary = (payload: ReplyPayload) =>
|
||||
enqueueProgressDelivery(async () => {
|
||||
if (
|
||||
run.sourceReplyDeliveryMode === "message_tool_only" &&
|
||||
!shouldEmitToolResultProgress()
|
||||
) {
|
||||
return;
|
||||
}
|
||||
await sendFollowupPayloads(
|
||||
[payload],
|
||||
effectiveQueued,
|
||||
{
|
||||
provider,
|
||||
modelId: model,
|
||||
},
|
||||
{ kind: "tool", mirror: false, runId },
|
||||
);
|
||||
if (payload.isError === true) {
|
||||
markVisibleToolErrorProgress();
|
||||
}
|
||||
});
|
||||
try {
|
||||
if (isCliProvider(cliExecutionProvider, runtimeConfig)) {
|
||||
const cliSessionBinding = getCliSessionBinding(
|
||||
@@ -845,27 +868,7 @@ export function createFollowupRunner(params: {
|
||||
detailMode: toolProgressDetail,
|
||||
shouldEmitToolResult: shouldEmitToolResultProgress,
|
||||
shouldEmitToolOutput: shouldEmitToolOutputProgress,
|
||||
deliver: (payload) =>
|
||||
enqueueProgressDelivery(async () => {
|
||||
if (
|
||||
run.sourceReplyDeliveryMode === "message_tool_only" &&
|
||||
!shouldEmitToolResultProgress()
|
||||
) {
|
||||
return;
|
||||
}
|
||||
await sendFollowupPayloads(
|
||||
[payload],
|
||||
effectiveQueued,
|
||||
{
|
||||
provider,
|
||||
modelId: model,
|
||||
},
|
||||
{ kind: "tool", mirror: false, runId },
|
||||
);
|
||||
if (payload.isError === true) {
|
||||
markVisibleToolErrorProgress();
|
||||
}
|
||||
}),
|
||||
deliver: deliverFollowupToolSummary,
|
||||
});
|
||||
const result = await runCliAgentWithLifecycle({
|
||||
runId,
|
||||
@@ -1068,27 +1071,7 @@ export function createFollowupRunner(params: {
|
||||
toolProgressDetail,
|
||||
shouldEmitToolResult: shouldEmitToolResultProgress,
|
||||
shouldEmitToolOutput: shouldEmitToolOutputProgress,
|
||||
onToolResult: (payload) =>
|
||||
enqueueProgressDelivery(async () => {
|
||||
if (
|
||||
run.sourceReplyDeliveryMode === "message_tool_only" &&
|
||||
!shouldEmitToolResultProgress()
|
||||
) {
|
||||
return;
|
||||
}
|
||||
await sendFollowupPayloads(
|
||||
[payload],
|
||||
effectiveQueued,
|
||||
{
|
||||
provider,
|
||||
modelId: model,
|
||||
},
|
||||
{ kind: "tool", mirror: false, runId },
|
||||
);
|
||||
if (payload.isError === true) {
|
||||
markVisibleToolErrorProgress();
|
||||
}
|
||||
}),
|
||||
onToolResult: deliverFollowupToolSummary,
|
||||
onAgentEvent: (evt) =>
|
||||
enqueueProgressDelivery(async () => {
|
||||
await forwardFollowupProgressEvent({
|
||||
|
||||
Reference in New Issue
Block a user