From 57da466ecbb5f63b5a1881ae8caa7359c02ec40c Mon Sep 17 00:00:00 2001 From: JC Date: Sun, 17 May 2026 15:59:07 -0700 Subject: [PATCH] Fix Discord verbose tool progress delivery (#80042) Summary: - The PR changes Discord reply delivery, sanitizer, and queued follow-up auto-reply paths so explicit verbose tool-progress payloads are delivered while final assistant replies still use the privacy sanitizer. - Reproducibility: yes. source-level: current main strips tool-looking Discord payload text at the front-chann ... ds compaction events in queued follow-up runs. I did not run a live Discord repro in this read-only review. Automerge notes: - Ran the ClawSweeper repair loop before final review. - Included post-review commit in the final squash: fix: gate queued follow-up progress when verbose is off - Included post-review commit in the final squash: fix: preserve queued verbose progress under preview suppression - Included post-review commit in the final squash: ci: rerun discord verbose progress PR - Included post-review commit in the final squash: fix: preserve Discord verbose progress after rebase - Included post-review commit in the final squash: fix: serialize discord queued progress - Included post-review commit in the final squash: Fix Discord verbose tool progress delivery Validation: - ClawSweeper review passed for head fd845e773ac6d649ee2dd7ff0defb4676b06b3da. - Required merge gates passed before the squash merge. Prepared head SHA: fd845e773ac6d649ee2dd7ff0defb4676b06b3da Review: https://github.com/openclaw/openclaw/pull/80042#issuecomment-4414121881 Co-authored-by: Clawsistant Co-authored-by: anyech Co-authored-by: OpenClaw Assistant Co-authored-by: Shadow Co-authored-by: clawsweeper <274271284+clawsweeper[bot]@users.noreply.github.com> Co-authored-by: clawsweeper[bot] <274271284+clawsweeper[bot]@users.noreply.github.com> Approved-by: thewilloftheshadow Co-authored-by: thewilloftheshadow <35580099+thewilloftheshadow@users.noreply.github.com> --- .../src/monitor/agent-components.dispatch.ts | 3 +- .../src/monitor/message-handler.process.ts | 3 + .../src/monitor/reply-delivery.test.ts | 31 ++ .../discord/src/monitor/reply-delivery.ts | 3 +- .../discord/src/monitor/reply-safety.ts | 6 +- src/auto-reply/reply/agent-runner.ts | 2 + .../reply/dispatch-from-config.test.ts | 37 ++- src/auto-reply/reply/dispatch-from-config.ts | 11 +- src/auto-reply/reply/followup-runner.test.ts | 269 ++++++++++++++++++ src/auto-reply/reply/followup-runner.ts | 209 +++++++++++++- 10 files changed, 555 insertions(+), 19 deletions(-) diff --git a/extensions/discord/src/monitor/agent-components.dispatch.ts b/extensions/discord/src/monitor/agent-components.dispatch.ts index c7ac0b026b67..926d6e612df4 100644 --- a/extensions/discord/src/monitor/agent-components.dispatch.ts +++ b/extensions/discord/src/monitor/agent-components.dispatch.ts @@ -312,7 +312,7 @@ export async function dispatchDiscordComponentEvent(params: { }, }, delivery: { - deliver: async (payload) => { + deliver: async (payload, info) => { const replyToId = replyReference.use(); await deliverDiscordReply({ cfg: ctx.cfg, @@ -333,6 +333,7 @@ export async function dispatchDiscordComponentEvent(params: { tableMode, chunkMode: resolveChunkMode(ctx.cfg, "discord", accountId), mediaLocalRoots, + kind: info.kind, }); replyReference.markSent(); }, diff --git a/extensions/discord/src/monitor/message-handler.process.ts b/extensions/discord/src/monitor/message-handler.process.ts index b4e2e212a9f8..3352001cd0ca 100644 --- a/extensions/discord/src/monitor/message-handler.process.ts +++ b/extensions/discord/src/monitor/message-handler.process.ts @@ -640,6 +640,7 @@ export async function processDiscordMessage( sessionKey: ctxPayload.SessionKey, threadBindings, mediaLocalRoots, + kind: info.kind, }); return true; }, @@ -678,6 +679,7 @@ export async function processDiscordMessage( sessionKey: ctxPayload.SessionKey, threadBindings, mediaLocalRoots, + kind: info.kind, }); return true; }, @@ -716,6 +718,7 @@ export async function processDiscordMessage( sessionKey: ctxPayload.SessionKey, threadBindings, mediaLocalRoots, + kind: info.kind, }); replyReference.markSent(); if (isFinal) { diff --git a/extensions/discord/src/monitor/reply-delivery.test.ts b/extensions/discord/src/monitor/reply-delivery.test.ts index e31cfce487bc..d9bf401a16ae 100644 --- a/extensions/discord/src/monitor/reply-delivery.test.ts +++ b/extensions/discord/src/monitor/reply-delivery.test.ts @@ -119,6 +119,7 @@ describe("deliverDiscordReply", () => { textLimit: 2000, replyToId: "reply-1", replyToMode: "all", + kind: "final", }); const params = firstDeliverParams(); @@ -151,10 +152,30 @@ describe("deliverDiscordReply", () => { runtime, cfg, textLimit: 2000, + kind: "final", }), ).rejects.toThrow("discord final reply produced no delivered message for channel:101"); }); + it("preserves explicit tool progress payloads at the tool delivery boundary", async () => { + await deliverDiscordReply({ + replies: [{ text: "🛠️ Exec: `echo visible`" }], + target: "channel:101", + token: "token", + accountId: "default", + runtime, + cfg, + textLimit: 2000, + kind: "tool", + }); + + expect(sendDurableMessageBatchMock).toHaveBeenCalledWith( + expect.objectContaining({ + payloads: [{ text: "🛠️ Exec: `echo visible`" }], + }), + ); + }); + it("strips internal execution trace lines at the final Discord send boundary", async () => { await deliverDiscordReply({ replies: [ @@ -177,6 +198,7 @@ describe("deliverDiscordReply", () => { runtime, cfg, textLimit: 2000, + kind: "final", }); expect(firstDeliverParams().payloads).toEqual([{ text: "Visible reply." }]); @@ -196,6 +218,7 @@ describe("deliverDiscordReply", () => { runtime, cfg, textLimit: 2000, + kind: "final", }); expect(firstDeliverParams().payloads).toEqual([ @@ -235,6 +258,7 @@ describe("deliverDiscordReply", () => { runtime, cfg, textLimit: 2000, + kind: "final", }); expect(firstDeliverParams().payloads).toEqual([{ channelData, text: undefined }]); @@ -264,6 +288,7 @@ describe("deliverDiscordReply", () => { runtime, cfg, textLimit: 2000, + kind: "final", }); expect(firstDeliverParams().payloads).toEqual([{ presentation, text: undefined }]); @@ -280,6 +305,7 @@ describe("deliverDiscordReply", () => { runtime, cfg, textLimit: 2000, + kind: "final", }); expect(firstDeliverParams().payloads).toEqual([{ text }]); @@ -301,6 +327,7 @@ describe("deliverDiscordReply", () => { runtime, cfg, textLimit: 2000, + kind: "final", }); expect(firstDeliverParams().payloads).toEqual([{ text }]); @@ -334,6 +361,7 @@ describe("deliverDiscordReply", () => { maxLinesPerMessage: 7, tableMode: "off", chunkMode: "newline", + kind: "final", }); expect(firstDeliverParams().cfg).toBe(baseCfg); @@ -363,6 +391,7 @@ describe("deliverDiscordReply", () => { textLimit: 2000, replyToMode: "off", mediaLocalRoots: ["/tmp/openclaw-media"], + kind: "final", }); const params = firstDeliverParams(); @@ -381,6 +410,7 @@ describe("deliverDiscordReply", () => { cfg, textLimit: 2000, replyToId: "reply-1", + kind: "final", }); const deps = firstDeliverParams().deps!; @@ -429,6 +459,7 @@ describe("deliverDiscordReply", () => { replyToId: "reply-1", sessionKey: "agent:main:subagent:child", threadBindings, + kind: "final", }); const params = firstDeliverParams(); diff --git a/extensions/discord/src/monitor/reply-delivery.ts b/extensions/discord/src/monitor/reply-delivery.ts index 15558590ecda..ce79b493472a 100644 --- a/extensions/discord/src/monitor/reply-delivery.ts +++ b/extensions/discord/src/monitor/reply-delivery.ts @@ -172,11 +172,12 @@ export async function deliverDiscordReply(params: { sessionKey?: string; threadBindings?: DiscordThreadBindingLookup; mediaLocalRoots?: readonly string[]; + kind: "tool" | "block" | "final"; }) { void params.runtime; const delivery = resolveDiscordDeliveryOptions(params); - const payloads = sanitizeDiscordFrontChannelReplyPayloads(params.replies); + const payloads = sanitizeDiscordFrontChannelReplyPayloads(params.replies, { kind: params.kind }); if (payloads.length === 0) { return; } diff --git a/extensions/discord/src/monitor/reply-safety.ts b/extensions/discord/src/monitor/reply-safety.ts index 1fa43c5de7e1..cfb840b35861 100644 --- a/extensions/discord/src/monitor/reply-safety.ts +++ b/extensions/discord/src/monitor/reply-safety.ts @@ -71,12 +71,16 @@ export function sanitizeDiscordFrontChannelText(text: string): string { export function sanitizeDiscordFrontChannelReplyPayloads( payloads: readonly ReplyPayload[], + options: { kind?: "tool" | "block" | "final" } = {}, ): ReplyPayload[] { + const preserveVerboseToolProgress = options.kind === "tool"; const safePayloads: ReplyPayload[] = []; for (const payload of payloads) { const safeText = typeof payload.text === "string" - ? sanitizeDiscordFrontChannelText(payload.text) + ? preserveVerboseToolProgress + ? collapseExcessBlankLines(sanitizeAssistantVisibleText(payload.text)).trim() + : sanitizeDiscordFrontChannelText(payload.text) : payload.text; const nextPayload = safeText === payload.text diff --git a/src/auto-reply/reply/agent-runner.ts b/src/auto-reply/reply/agent-runner.ts index 27d3b953c2bf..4b86f5eb08d0 100644 --- a/src/auto-reply/reply/agent-runner.ts +++ b/src/auto-reply/reply/agent-runner.ts @@ -1193,6 +1193,7 @@ export async function runReplyAgent(params: { storePath, defaultModel, agentCfgContextTokens, + toolProgressDetail, }); if (activeRunQueueAction === "drop") { @@ -1415,6 +1416,7 @@ export async function runReplyAgent(params: { storePath, defaultModel, agentCfgContextTokens, + toolProgressDetail, }); let responseUsageLine: string | undefined; diff --git a/src/auto-reply/reply/dispatch-from-config.test.ts b/src/auto-reply/reply/dispatch-from-config.test.ts index 674b13d13485..23cb3fb3954f 100644 --- a/src/auto-reply/reply/dispatch-from-config.test.ts +++ b/src/auto-reply/reply/dispatch-from-config.test.ts @@ -2024,7 +2024,7 @@ describe("dispatchReplyFromConfig", () => { expect(dispatcher.sendFinalReply).not.toHaveBeenCalled(); }); - it("keeps verbose tool summaries suppressed for channel message-tool-only turns", async () => { + it("delivers verbose tool summaries for Discord channel message-tool-only turns", async () => { setNoAbort(); sessionStoreMocks.currentEntry = { sessionId: "s1", @@ -2056,7 +2056,7 @@ describe("dispatchReplyFromConfig", () => { }); expect(result.sourceReplyDeliveryMode).toBe("message_tool_only"); - expect(dispatcher.sendToolResult).not.toHaveBeenCalled(); + expect(dispatcher.sendToolResult).toHaveBeenCalledWith({ text: "🛠️ `pwd (agent)`" }); expect(dispatcher.sendFinalReply).not.toHaveBeenCalled(); }); @@ -4947,6 +4947,39 @@ describe("sendPolicy deny — suppress delivery, not processing (#53328)", () => expect(dispatcher.sendBlockReply).not.toHaveBeenCalled(); }); + it("delivers verbose tool progress in message-tool-only mode", async () => { + setNoAbort(); + sessionStoreMocks.currentEntry = { + sessionId: "s1", + updatedAt: 0, + sendPolicy: "allow", + verboseLevel: "on", + }; + const dispatcher = createDispatcher(); + const replyResolver = vi.fn(async (_ctx: MsgContext, opts?: GetReplyOptions) => { + await opts?.onToolResult?.({ text: "🛠️ Exec: echo post-restart" }); + return { text: "NO_REPLY" } satisfies ReplyPayload; + }); + const ctx = buildTestCtx({ SessionKey: "test:session", ChatType: "channel" }); + + const result = await dispatchReplyFromConfig({ + ctx, + cfg: emptyConfig, + dispatcher, + replyResolver, + replyOptions: { + sourceReplyDeliveryMode: "message_tool_only", + }, + }); + + expect(result.queuedFinal).toBe(false); + expect(result.sourceReplyDeliveryMode).toBe("message_tool_only"); + expect(dispatcher.sendToolResult).toHaveBeenCalledWith( + expect.objectContaining({ text: "🛠️ Exec: echo post-restart" }), + ); + expect(dispatcher.sendFinalReply).not.toHaveBeenCalled(); + }); + it("delivers marked runtime failure notices in message-tool-only mode", async () => { setNoAbort(); sessionStoreMocks.currentEntry = { diff --git a/src/auto-reply/reply/dispatch-from-config.ts b/src/auto-reply/reply/dispatch-from-config.ts index 51e416653bdc..e6c33783ecbc 100644 --- a/src/auto-reply/reply/dispatch-from-config.ts +++ b/src/auto-reply/reply/dispatch-from-config.ts @@ -1046,7 +1046,7 @@ export async function dispatchReplyFromConfig( const shouldSendToolStartStatuses = false; const shouldDeliverVerboseProgressDespiteSourceSuppression = () => suppressAutomaticSourceDelivery && - chatType === "direct" && + sourceReplyDeliveryMode === "message_tool_only" && ctx.InboundEventKind !== "room_event" && !sendPolicyDenied && shouldEmitVerboseProgress() && @@ -1214,7 +1214,7 @@ export async function dispatchReplyFromConfig( return parts.join("\n\n").trim() || "Planning next steps."; }; const maybeSendWorkingStatus = async (label: string): Promise => { - if (suppressDelivery && !shouldDeliverVerboseProgressDespiteSourceSuppression()) { + if (shouldSuppressProgressDelivery()) { return; } const normalizedLabel = normalizeWorkingLabel(label); @@ -1244,7 +1244,7 @@ export async function dispatchReplyFromConfig( steps?: string[]; }): Promise => { if ( - (suppressDelivery && !shouldDeliverVerboseProgressDespiteSourceSuppression()) || + shouldSuppressProgressDelivery() || !shouldEmitVerboseProgress() || !shouldSendVerboseProgressMessages ) { @@ -1349,6 +1349,9 @@ export async function dispatchReplyFromConfig( params.replyOptions?.suppressDefaultToolProgressMessages === true; const shouldSuppressDefaultToolProgressMessages = () => suppressDefaultToolProgressMessages && !shouldEmitVerboseProgress(); + const shouldSuppressProgressDelivery = () => + sendPolicyDenied || + (suppressDelivery && !shouldDeliverVerboseProgressDespiteSourceSuppression()); const onToolResultFromReplyOptions = params.replyOptions?.onToolResult; const onPlanUpdateFromReplyOptions = params.replyOptions?.onPlanUpdate; const onApprovalEventFromReplyOptions = params.replyOptions?.onApprovalEvent; @@ -1420,7 +1423,7 @@ export async function dispatchReplyFromConfig( if (!suppressAutomaticSourceDelivery) { await onToolResultFromReplyOptions?.(payload); } - if (suppressDelivery && !shouldDeliverVerboseProgressDespiteSourceSuppression()) { + if (shouldSuppressProgressDelivery()) { return; } const ttsPayload = await maybeApplyTtsToReplyPayload({ diff --git a/src/auto-reply/reply/followup-runner.test.ts b/src/auto-reply/reply/followup-runner.test.ts index b3566eb6daeb..257f8dc092b4 100644 --- a/src/auto-reply/reply/followup-runner.test.ts +++ b/src/auto-reply/reply/followup-runner.test.ts @@ -1093,6 +1093,275 @@ describe("createFollowupRunner runtime config", () => { }); }); +describe("createFollowupRunner progress forwarding", () => { + it("forwards queued follow-up tool progress and verbose tool result payloads", async () => { + const onToolStart = vi.fn(async () => {}); + const queued = createQueuedRun({ + originatingChannel: "discord", + originatingTo: "channel:C1", + originatingAccountId: "acct-1", + originatingThreadId: "thread-1", + run: { + messageProvider: "discord", + sourceReplyDeliveryMode: "message_tool_only", + verboseLevel: "on", + }, + }); + + runEmbeddedPiAgentMock.mockImplementationOnce( + async (args: { + onAgentEvent?: (evt: { stream: string; data: Record }) => Promise; + onToolResult?: (payload: { text: string }) => Promise; + shouldEmitToolResult?: () => boolean; + shouldEmitToolOutput?: () => boolean; + toolProgressDetail?: "explain" | "raw"; + }) => { + expect(args.shouldEmitToolResult?.()).toBe(true); + expect(args.shouldEmitToolOutput?.()).toBe(false); + expect(args.toolProgressDetail).toBe("raw"); + await args.onAgentEvent?.({ + stream: "tool", + data: { + phase: "start", + name: "exec", + args: { command: "echo queued-progress" }, + }, + }); + await args.onToolResult?.({ text: "🛠️ Exec: echo queued-progress" }); + return { payloads: [], meta: { agentMeta: {} } }; + }, + ); + + const runner = createFollowupRunner({ + opts: { onToolStart }, + typing: createMockTypingController(), + typingMode: "instant", + defaultModel: "claude", + toolProgressDetail: "raw", + }); + + await runner(queued); + + expect(onToolStart).toHaveBeenCalledWith({ + name: "exec", + phase: "start", + args: { command: "echo queued-progress" }, + detailMode: "raw", + }); + expect(routeReplyMock).toHaveBeenCalledTimes(1); + expect(routeReplyMock).toHaveBeenCalledWith( + expect.objectContaining({ + channel: "discord", + to: "channel:C1", + accountId: "acct-1", + threadId: "thread-1", + mirror: false, + payload: expect.objectContaining({ text: "🛠️ Exec: echo queued-progress" }), + }), + ); + }); + + it("drains fire-and-forget queued tool progress before final delivery", async () => { + const queued = createQueuedRun({ + originatingChannel: "discord", + originatingTo: "channel:C1", + originatingAccountId: "acct-1", + originatingThreadId: "thread-1", + run: { + messageProvider: "discord", + verboseLevel: "on", + }, + }); + let releaseProgressRoute: (() => void) | undefined; + const progressRouteStarted = new Promise((resolve) => { + routeReplyMock.mockImplementationOnce( + async () => + await new Promise<{ ok: true }>((release) => { + releaseProgressRoute = () => { + release({ ok: true }); + }; + resolve(); + }), + ); + }); + + runEmbeddedPiAgentMock.mockImplementationOnce( + async (args: { onToolResult?: (payload: { text: string }) => Promise }) => { + void args.onToolResult?.({ text: "🛠️ Exec: echo queued-progress" }); + return { payloads: [{ text: "final reply" }], meta: { agentMeta: {} } }; + }, + ); + + const runner = createFollowupRunner({ + typing: createMockTypingController(), + typingMode: "instant", + defaultModel: "claude", + }); + + const runPromise = runner(queued); + await progressRouteStarted; + await Promise.resolve(); + + expect(routeReplyMock).toHaveBeenCalledTimes(1); + expect(requireMockCallArg(routeReplyMock, 0).payload).toEqual( + expect.objectContaining({ text: "🛠️ Exec: echo queued-progress" }), + ); + expect(requireMockCallArg(routeReplyMock, 0).mirror).toBe(false); + + releaseProgressRoute?.(); + await runPromise; + + expect(routeReplyMock).toHaveBeenCalledTimes(2); + expect(requireMockCallArg(routeReplyMock, 1).payload).toEqual( + expect.objectContaining({ text: "final reply" }), + ); + expect(requireMockCallArg(routeReplyMock, 1).mirror).toBeUndefined(); + }); + + it("preserves queued verbose progress when default tool progress is suppressed", async () => { + const onToolStart = vi.fn(async () => {}); + const onCommandOutput = vi.fn(async () => {}); + const queued = createQueuedRun({ + originatingChannel: "discord", + originatingTo: "channel:C1", + originatingAccountId: "acct-1", + originatingThreadId: "thread-1", + run: { + messageProvider: "discord", + sourceReplyDeliveryMode: "message_tool_only", + verboseLevel: "on", + }, + }); + + runEmbeddedPiAgentMock.mockImplementationOnce( + async (args: { + onAgentEvent?: (evt: { stream: string; data: Record }) => Promise; + onToolResult?: (payload: { text: string }) => Promise; + shouldEmitToolResult?: () => boolean; + shouldEmitToolOutput?: () => boolean; + }) => { + expect(args.shouldEmitToolResult?.()).toBe(true); + expect(args.shouldEmitToolOutput?.()).toBe(false); + await args.onAgentEvent?.({ + stream: "tool", + data: { + phase: "start", + name: "exec", + args: { command: "echo queued-suppressed-preview" }, + }, + }); + await args.onAgentEvent?.({ + stream: "command_output", + data: { phase: "chunk", output: "queued output" }, + }); + await args.onToolResult?.({ text: "🛠️ Exec: echo queued-suppressed-preview" }); + return { payloads: [], meta: { agentMeta: {} } }; + }, + ); + + const runner = createFollowupRunner({ + opts: { suppressDefaultToolProgressMessages: true, onToolStart, onCommandOutput }, + typing: createMockTypingController(), + typingMode: "instant", + defaultModel: "claude", + toolProgressDetail: "raw", + }); + + await runner(queued); + + expect(onToolStart).toHaveBeenCalledWith({ + name: "exec", + phase: "start", + args: { command: "echo queued-suppressed-preview" }, + detailMode: "raw", + }); + expect(onCommandOutput).toHaveBeenCalledWith( + expect.objectContaining({ phase: "chunk", output: "queued output" }), + ); + expect(routeReplyMock).toHaveBeenCalledTimes(1); + expect(routeReplyMock).toHaveBeenCalledWith( + expect.objectContaining({ + channel: "discord", + to: "channel:C1", + accountId: "acct-1", + threadId: "thread-1", + mirror: false, + payload: expect.objectContaining({ text: "🛠️ Exec: echo queued-suppressed-preview" }), + }), + ); + }); + + it("suppresses queued follow-up progress when verbose progress is disabled", async () => { + const storePath = path.join( + await fs.mkdtemp(path.join(tmpdir(), "openclaw-followup-progress-off-")), + "sessions.json", + ); + const sessionEntry: SessionEntry = { + sessionId: "session", + updatedAt: Date.now(), + }; + const sessionStore: Record = { main: sessionEntry }; + const onToolStart = vi.fn(async () => {}); + const onItemEvent = vi.fn(async () => {}); + const onCommandOutput = vi.fn(async () => {}); + registerFollowupTestSessionStore(storePath, sessionStore); + + runEmbeddedPiAgentMock.mockImplementationOnce( + async (args: { + onAgentEvent?: (evt: { stream: string; data: Record }) => Promise; + shouldEmitToolResult?: () => boolean; + shouldEmitToolOutput?: () => boolean; + }) => { + expect(args.shouldEmitToolResult?.()).toBe(false); + expect(args.shouldEmitToolOutput?.()).toBe(false); + await args.onAgentEvent?.({ + stream: "tool", + data: { phase: "start", name: "exec", args: { command: "echo hidden" } }, + }); + await args.onAgentEvent?.({ + stream: "item", + data: { phase: "start", itemId: "item-1", title: "hidden item" }, + }); + await args.onAgentEvent?.({ + stream: "command_output", + data: { phase: "chunk", output: "hidden output" }, + }); + await args.onAgentEvent?.({ + stream: "compaction", + data: { phase: "end", completed: true }, + }); + return { payloads: [{ text: "final" }], meta: { agentMeta: {} } }; + }, + ); + + const runner = createFollowupRunner({ + opts: { onToolStart, onItemEvent, onCommandOutput }, + typing: createMockTypingController(), + typingMode: "instant", + sessionEntry, + sessionStore, + sessionKey: "main", + storePath, + defaultModel: "claude", + }); + + await runner( + createQueuedRun({ + run: { + messageProvider: "discord", + sourceReplyDeliveryMode: "message_tool_only", + verboseLevel: "off", + }, + }), + ); + + expect(onToolStart).not.toHaveBeenCalled(); + expect(onItemEvent).not.toHaveBeenCalled(); + expect(onCommandOutput).not.toHaveBeenCalled(); + expect(sessionStore.main.compactionCount).toBe(1); + }); +}); + describe("createFollowupRunner compaction", () => { it("adds verbose auto-compaction notice and tracks count", async () => { const storePath = path.join( diff --git a/src/auto-reply/reply/followup-runner.ts b/src/auto-reply/reply/followup-runner.ts index e1de0fc8d165..e8339a1aaf53 100644 --- a/src/auto-reply/reply/followup-runner.ts +++ b/src/auto-reply/reply/followup-runner.ts @@ -23,6 +23,7 @@ import { logVerbose } from "../../globals.js"; import { emitAgentEvent, registerAgentRunContext } from "../../infra/agent-events.js"; import { formatErrorMessage } from "../../infra/errors.js"; import { defaultRuntime } from "../../runtime.js"; +import { readStringValue } from "../../shared/string-coerce.js"; import { isInternalMessageChannel } from "../../utils/message-channel.js"; import type { GetReplyOptions, ReplyPayload } from "../types.js"; import { runCliAgentWithLifecycle } from "./agent-runner-cli-dispatch.js"; @@ -53,6 +54,139 @@ import type { TypingController } from "./typing.js"; type EmbeddedAgentRunResult = Awaited>; +type FollowupAgentEvent = { stream: string; data: Record }; + +function readApprovalScopeValue(value: unknown): "turn" | "session" | undefined { + return value === "turn" || value === "session" ? value : undefined; +} + +function filterStringArray(value: unknown): string[] | undefined { + return Array.isArray(value) + ? value.filter((entry): entry is string => typeof entry === "string") + : undefined; +} + +async function forwardFollowupProgressEvent(params: { + evt: FollowupAgentEvent; + opts?: GetReplyOptions; + detailMode?: "explain" | "raw"; + emitChannelProgress?: boolean; + onCompactionComplete?: () => void; +}) { + const { evt, opts } = params; + const emitChannelProgress = params.emitChannelProgress !== false; + if (!emitChannelProgress && evt.stream !== "compaction") { + return; + } + + if (evt.stream === "tool") { + const phase = readStringValue(evt.data.phase) ?? ""; + const name = readStringValue(evt.data.name); + if (phase === "start" || phase === "update") { + await opts?.onToolStart?.({ + name, + phase, + args: + evt.data.args && typeof evt.data.args === "object" + ? (evt.data.args as Record) + : undefined, + detailMode: params.detailMode, + }); + } + } + + const suppressItemChannelProgress = + evt.stream === "item" && + evt.data.suppressChannelProgress === true && + Boolean(opts?.onToolStart); + if (evt.stream === "item" && !suppressItemChannelProgress) { + await opts?.onItemEvent?.({ + itemId: readStringValue(evt.data.itemId), + kind: readStringValue(evt.data.kind), + title: readStringValue(evt.data.title), + name: readStringValue(evt.data.name), + phase: readStringValue(evt.data.phase), + status: readStringValue(evt.data.status), + summary: readStringValue(evt.data.summary), + progressText: readStringValue(evt.data.progressText), + meta: readStringValue(evt.data.meta), + approvalId: readStringValue(evt.data.approvalId), + approvalSlug: readStringValue(evt.data.approvalSlug), + }); + } + + if (evt.stream === "plan") { + await opts?.onPlanUpdate?.({ + phase: readStringValue(evt.data.phase), + title: readStringValue(evt.data.title), + explanation: readStringValue(evt.data.explanation), + steps: filterStringArray(evt.data.steps), + source: readStringValue(evt.data.source), + }); + } + + if (evt.stream === "approval") { + await opts?.onApprovalEvent?.({ + phase: readStringValue(evt.data.phase), + kind: readStringValue(evt.data.kind), + status: readStringValue(evt.data.status), + title: readStringValue(evt.data.title), + itemId: readStringValue(evt.data.itemId), + toolCallId: readStringValue(evt.data.toolCallId), + approvalId: readStringValue(evt.data.approvalId), + approvalSlug: readStringValue(evt.data.approvalSlug), + command: readStringValue(evt.data.command), + host: readStringValue(evt.data.host), + reason: readStringValue(evt.data.reason), + scope: readApprovalScopeValue(evt.data.scope), + message: readStringValue(evt.data.message), + }); + } + + if (evt.stream === "command_output") { + await opts?.onCommandOutput?.({ + itemId: readStringValue(evt.data.itemId), + phase: readStringValue(evt.data.phase), + title: readStringValue(evt.data.title), + toolCallId: readStringValue(evt.data.toolCallId), + name: readStringValue(evt.data.name), + output: readStringValue(evt.data.output), + status: readStringValue(evt.data.status), + exitCode: + typeof evt.data.exitCode === "number" || evt.data.exitCode === null + ? evt.data.exitCode + : undefined, + durationMs: typeof evt.data.durationMs === "number" ? evt.data.durationMs : undefined, + cwd: readStringValue(evt.data.cwd), + }); + } + + if (evt.stream === "patch") { + await opts?.onPatchSummary?.({ + itemId: readStringValue(evt.data.itemId), + phase: readStringValue(evt.data.phase), + title: readStringValue(evt.data.title), + toolCallId: readStringValue(evt.data.toolCallId), + name: readStringValue(evt.data.name), + added: filterStringArray(evt.data.added), + modified: filterStringArray(evt.data.modified), + deleted: filterStringArray(evt.data.deleted), + summary: readStringValue(evt.data.summary), + }); + } + + if (evt.stream === "compaction") { + const phase = readStringValue(evt.data.phase) ?? ""; + if (phase === "start") { + await opts?.onCompactionStart?.(); + } + if (phase === "end" && evt.data?.completed === true) { + params.onCompactionComplete?.(); + await opts?.onCompactionEnd?.(); + } + } +} + export function createFollowupRunner(params: { opts?: GetReplyOptions; typing: TypingController; @@ -63,6 +197,7 @@ export function createFollowupRunner(params: { storePath?: string; defaultModel: string; agentCfgContextTokens?: number; + toolProgressDetail?: "explain" | "raw"; }): (queued: FollowupRun) => Promise { const { opts, @@ -74,6 +209,7 @@ export function createFollowupRunner(params: { storePath, defaultModel, agentCfgContextTokens, + toolProgressDetail, } = params; const typingSignals = createTypingSignaler({ typing, @@ -93,6 +229,7 @@ export function createFollowupRunner(params: { payloads: ReplyPayload[], queued: FollowupRun, resolvedRun: { provider: string; modelId: string }, + options: { mirror?: boolean } = {}, ) => { // Check if we should route to originating channel. const { originatingChannel, originatingTo } = queued; @@ -164,6 +301,7 @@ export function createFollowupRunner(params: { requesterSenderE164: queued.run.senderE164, threadId: queued.originatingThreadId, cfg: runtimeConfig, + mirror: options.mirror, }); if (!result.ok) { const errorMsg = result.error ?? "unknown error"; @@ -226,6 +364,7 @@ export function createFollowupRunner(params: { const queuedImages = queued.images ?? opts?.images; const queuedImageOrder = queued.imageOrder ?? opts?.imageOrder; let replyOperation: ReturnType | undefined; + try { queued.run.config = await resolveQueuedReplyExecutionConfig(queued.run.config, { originatingChannel: queued.originatingChannel, @@ -251,6 +390,30 @@ export function createFollowupRunner(params: { if (run !== effectiveQueued.run) { effectiveQueued = { ...effectiveQueued, run }; } + const shouldEmitVerboseProgress = () => run.verboseLevel !== "off"; + const shouldSuppressDefaultToolProgressMessages = () => + opts?.suppressDefaultToolProgressMessages === true && !shouldEmitVerboseProgress(); + const shouldEmitToolResultProgress = () => + shouldEmitVerboseProgress() && !shouldSuppressDefaultToolProgressMessages(); + const shouldEmitToolOutputProgress = () => + run.verboseLevel === "full" && !shouldSuppressDefaultToolProgressMessages(); + let progressDeliveryChain: Promise = Promise.resolve(); + const pendingProgressDeliveries = new Set>(); + const enqueueProgressDelivery = (deliver: () => Promise) => { + progressDeliveryChain = progressDeliveryChain.then(deliver).catch((err) => { + logVerbose(`followup queue: progress delivery failed: ${formatErrorMessage(err)}`); + }); + const task = progressDeliveryChain.finally(() => { + pendingProgressDeliveries.delete(task); + }); + pendingProgressDeliveries.add(task); + return task; + }; + const drainProgressDeliveries = async () => { + while (pendingProgressDeliveries.size > 0) { + await Promise.all(pendingProgressDeliveries); + } + }; replyOperation = createReplyOperation({ sessionId: run.sessionId, sessionKey: replySessionKey ?? "", @@ -558,16 +721,39 @@ export function createFollowupRunner(params: { bootstrapPromptWarningSignaturesSeen[ bootstrapPromptWarningSignaturesSeen.length - 1 ], - onAgentEvent: (evt) => { - if (evt.stream !== "compaction") { - return; - } - const phase = typeof evt.data.phase === "string" ? evt.data.phase : ""; - const completed = evt.data?.completed === true; - if (phase === "end" && completed) { - attemptCompactionCount += 1; - } - }, + toolProgressDetail, + shouldEmitToolResult: shouldEmitToolResultProgress, + shouldEmitToolOutput: shouldEmitToolOutputProgress, + onToolResult: (payload) => + enqueueProgressDelivery(async () => { + if ( + run.sourceReplyDeliveryMode === "message_tool_only" && + run.verboseLevel === "off" + ) { + return; + } + await sendFollowupPayloads( + [payload], + effectiveQueued, + { + provider, + modelId: model, + }, + { mirror: false }, + ); + }), + onAgentEvent: (evt) => + enqueueProgressDelivery(async () => { + await forwardFollowupProgressEvent({ + evt, + opts, + detailMode: toolProgressDetail, + emitChannelProgress: shouldEmitToolResultProgress(), + onCompactionComplete: () => { + attemptCompactionCount += 1; + }, + }); + }), }); bootstrapPromptWarningSignaturesSeen = resolveBootstrapWarningSignaturesSeen( result.meta?.systemPromptReport, @@ -622,10 +808,13 @@ export function createFollowupRunner(params: { }); pendingDeferredCliTerminal = undefined; } + await drainProgressDeliveries(); defaultRuntime.error?.(`Followup agent failed before reply: ${message}`); return; } + await drainProgressDeliveries(); + const usage = runResult.meta?.agentMeta?.usage; const promptTokens = runResult.meta?.agentMeta?.promptTokens; const modelUsed = runResult.meta?.agentMeta?.model ?? fallbackModel ?? defaultModel;