fix(acp): suppress commentary relay leakage

This commit is contained in:
Vincent Koc
2026-04-11 13:30:55 +01:00
parent 43bd5545f8
commit 7315914ee5
4 changed files with 92 additions and 3 deletions
@@ -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",
+13 -1
View File
@@ -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;
@@ -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",
});
});
});
@@ -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<string, unknown>)
@@ -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,