From 7315914ee50de2f7b67e037150152ae08e2a39e6 Mon Sep 17 00:00:00 2001 From: Vincent Koc Date: Sat, 11 Apr 2026 13:30:55 +0100 Subject: [PATCH] fix(acp): suppress commentary relay leakage --- src/agents/acp-spawn-parent-stream.test.ts | 60 +++++++++++++++++++ src/agents/acp-spawn-parent-stream.ts | 14 ++++- ...bedded-subscribe.handlers.messages.test.ts | 2 + ...pi-embedded-subscribe.handlers.messages.ts | 19 +++++- 4 files changed, 92 insertions(+), 3 deletions(-) diff --git a/src/agents/acp-spawn-parent-stream.test.ts b/src/agents/acp-spawn-parent-stream.test.ts index 506e107436..c94758efd7 100644 --- a/src/agents/acp-spawn-parent-stream.test.ts +++ b/src/agents/acp-spawn-parent-stream.test.ts @@ -290,6 +290,66 @@ describe("startAcpSpawnParentStreamRelay", () => { relay.dispose(); }); + it("suppresses commentary-phase assistant relay text", () => { + const relay = startAcpSpawnParentStreamRelay({ + runId: "run-commentary", + parentSessionKey: "agent:main:main", + childSessionKey: "agent:codex:acp:child-commentary", + agentId: "codex", + streamFlushMs: 10, + noOutputNoticeMs: 120_000, + }); + + emitAgentEvent({ + runId: "run-commentary", + stream: "assistant", + data: { + delta: "checking thread context; then post a tight progress reply here.", + phase: "commentary", + }, + }); + vi.advanceTimersByTime(15); + + const texts = collectedTexts(); + expect(texts.some((text) => text.includes("checking thread context"))).toBe(false); + expect(texts.some((text) => text.includes("post a tight progress reply here"))).toBe(false); + relay.dispose(); + }); + + it("still relays final_answer assistant text after suppressed commentary", () => { + const relay = startAcpSpawnParentStreamRelay({ + runId: "run-final", + parentSessionKey: "agent:main:main", + childSessionKey: "agent:codex:acp:child-final", + agentId: "codex", + streamFlushMs: 10, + noOutputNoticeMs: 120_000, + }); + + emitAgentEvent({ + runId: "run-final", + stream: "assistant", + data: { + delta: "checking thread context; then post a tight progress reply here.", + phase: "commentary", + }, + }); + emitAgentEvent({ + runId: "run-final", + stream: "assistant", + data: { + delta: "final answer ready", + phase: "final_answer", + }, + }); + vi.advanceTimersByTime(15); + + const texts = collectedTexts(); + expect(texts.some((text) => text.includes("checking thread context"))).toBe(false); + expect(texts.some((text) => text.includes("codex: final answer ready"))).toBe(true); + relay.dispose(); + }); + it("resolves ACP spawn stream log path from session metadata", () => { readAcpSessionEntryMock.mockReturnValue({ storePath: "/tmp/openclaw/agents/codex/sessions/sessions.json", diff --git a/src/agents/acp-spawn-parent-stream.ts b/src/agents/acp-spawn-parent-stream.ts index e96d91f711..8f8102531a 100644 --- a/src/agents/acp-spawn-parent-stream.ts +++ b/src/agents/acp-spawn-parent-stream.ts @@ -6,6 +6,7 @@ import { onAgentEvent } from "../infra/agent-events.js"; import { requestHeartbeatNow } from "../infra/heartbeat-wake.js"; import { enqueueSystemEvent } from "../infra/system-events.js"; import { scopedHeartbeatWakeOptions } from "../routing/session-key.js"; +import { normalizeAssistantPhase } from "../shared/chat-message-content.js"; import { normalizeOptionalString } from "../shared/string-coerce.js"; import { recordTaskRunProgressByRunId } from "../tasks/task-executor.js"; import { type DeliveryContext } from "../utils/delivery-context.js"; @@ -310,6 +311,9 @@ export function startAcpSpawnParentStreamRelay(params: { if (event.stream === "assistant") { const data = event.data; + const assistantPhase = normalizeAssistantPhase( + (data as { phase?: unknown } | undefined)?.phase, + ); const deltaCandidate = (data as { delta?: unknown } | undefined)?.delta ?? (data as { text?: unknown } | undefined)?.text; @@ -317,7 +321,15 @@ export function startAcpSpawnParentStreamRelay(params: { if (!delta || !delta.trim()) { return; } - logEvent("assistant_delta", { delta }); + logEvent("assistant_delta", { + delta, + ...(assistantPhase ? { phase: assistantPhase } : {}), + }); + + if (assistantPhase === "commentary") { + lastProgressAt = Date.now(); + return; + } if (stallNotified) { stallNotified = false; diff --git a/src/agents/pi-embedded-subscribe.handlers.messages.test.ts b/src/agents/pi-embedded-subscribe.handlers.messages.test.ts index 137b41b232..7075b5f1ad 100644 --- a/src/agents/pi-embedded-subscribe.handlers.messages.test.ts +++ b/src/agents/pi-embedded-subscribe.handlers.messages.test.ts @@ -178,12 +178,14 @@ describe("buildAssistantStreamData", () => { delta: "he", replace: true, mediaUrl: "https://example.com/a.png", + phase: "final_answer", }), ).toEqual({ text: "hello", delta: "he", replace: true, mediaUrls: ["https://example.com/a.png"], + phase: "final_answer", }); }); }); diff --git a/src/agents/pi-embedded-subscribe.handlers.messages.ts b/src/agents/pi-embedded-subscribe.handlers.messages.ts index 2ae52b7b59..8feeebd4b8 100644 --- a/src/agents/pi-embedded-subscribe.handlers.messages.ts +++ b/src/agents/pi-embedded-subscribe.handlers.messages.ts @@ -5,7 +5,10 @@ import { parseReplyDirectives } from "../auto-reply/reply/reply-directives.js"; import { isSilentReplyText, SILENT_REPLY_TOKEN } from "../auto-reply/tokens.js"; import { emitAgentEvent } from "../infra/agent-events.js"; import { createInlineCodeState } from "../markdown/code-spans.js"; -import { resolveAssistantMessagePhase } from "../shared/chat-message-content.js"; +import { + resolveAssistantMessagePhase, + type AssistantPhase, +} from "../shared/chat-message-content.js"; import { normalizeOptionalString } from "../shared/string-coerce.js"; import { isMessagingToolDuplicateNormalized, @@ -165,13 +168,21 @@ export function buildAssistantStreamData(params: { replace?: boolean; mediaUrls?: string[]; mediaUrl?: string; -}): { text: string; delta: string; replace?: true; mediaUrls?: string[] } { + phase?: AssistantPhase; +}): { + text: string; + delta: string; + replace?: true; + mediaUrls?: string[]; + phase?: AssistantPhase; +} { const mediaUrls = resolveSendableOutboundReplyParts(params).mediaUrls; return { text: params.text ?? "", delta: params.delta ?? "", replace: params.replace ? true : undefined, mediaUrls: mediaUrls.length ? mediaUrls : undefined, + phase: params.phase, }; } @@ -212,6 +223,7 @@ export function handleMessageUpdate( ctx.state.deterministicApprovalPromptPending || ctx.state.deterministicApprovalPromptSent; const assistantEvent = evt.assistantMessageEvent; + const assistantPhase = resolveAssistantMessagePhase(msg); const assistantRecord = assistantEvent && typeof assistantEvent === "object" ? (assistantEvent as Record) @@ -388,6 +400,7 @@ export function handleMessageUpdate( delta: deltaText, replace, mediaUrls, + phase: assistantPhase, }); emitAgentEvent({ runId: ctx.params.runId, @@ -440,6 +453,7 @@ export function handleMessageEnd( } const assistantMessage = msg; + const assistantPhase = resolveAssistantMessagePhase(assistantMessage); const suppressVisibleAssistantOutput = shouldSuppressAssistantVisibleOutput(assistantMessage); const suppressDeterministicApprovalOutput = ctx.state.deterministicApprovalPromptPending || ctx.state.deterministicApprovalPromptSent; @@ -512,6 +526,7 @@ export function handleMessageEnd( delta: finalStreamDelta, replace: shouldReplaceFinalStream, mediaUrls, + phase: assistantPhase, }); emitAgentEvent({ runId: ctx.params.runId,