diff --git a/src/agents/embedded-agent-runner.ts b/src/agents/embedded-agent-runner.ts index ecc62ba93bdd..dcebae0b734e 100644 --- a/src/agents/embedded-agent-runner.ts +++ b/src/agents/embedded-agent-runner.ts @@ -8,7 +8,9 @@ export { runEmbeddedAgent } from "./embedded-agent-runner/run.js"; export { abortAndDrainEmbeddedAgentRun, abortEmbeddedAgentRun, + isEmbeddedAgentRunAbortableForCompaction, isEmbeddedAgentRunActive, + isEmbeddedAgentRunHandleActive, isEmbeddedAgentRunStreaming, queueEmbeddedAgentMessage, queueEmbeddedAgentMessageWithOutcome, diff --git a/src/agents/embedded-agent-runner/runs.test.ts b/src/agents/embedded-agent-runner/runs.test.ts index 6c1b73011543..45247a5f02b0 100644 --- a/src/agents/embedded-agent-runner/runs.test.ts +++ b/src/agents/embedded-agent-runner/runs.test.ts @@ -24,6 +24,7 @@ import { clearEmbeddedRunAbandonment, consumeEmbeddedRunModelSwitch, getActiveEmbeddedRunSnapshot, + isEmbeddedAgentRunAbortableForCompaction, isEmbeddedAgentRunHandleActive, isEmbeddedRunAbandoned, formatEmbeddedAgentQueueFailureSummary, @@ -91,6 +92,20 @@ describe("embedded-agent runner run registry", () => { expect(abortNormal).not.toHaveBeenCalled(); }); + it("keeps queued reply operations out of compact abort checks", () => { + const operation = createReplyOperation({ + sessionKey: "agent:main:main", + sessionId: "session-reply-run", + resetTriggered: false, + }); + + expect(isEmbeddedAgentRunAbortableForCompaction("session-reply-run")).toBe(false); + + operation.setPhase("running"); + + expect(isEmbeddedAgentRunAbortableForCompaction("session-reply-run")).toBe(true); + }); + it("aborts every active run in all mode", () => { const abortA = vi.fn(); const abortB = vi.fn(); diff --git a/src/agents/embedded-agent-runner/runs.ts b/src/agents/embedded-agent-runner/runs.ts index 1abb0d9e1130..3869a9f67326 100644 --- a/src/agents/embedded-agent-runner/runs.ts +++ b/src/agents/embedded-agent-runner/runs.ts @@ -7,6 +7,7 @@ import { abortReplyRunBySessionId, forceClearReplyRunBySessionId, isReplyRunActiveForSessionId, + isReplyRunAbortableForCompaction, isReplyRunStreamingForSessionId, queueReplyRunMessage, resolveActiveReplyRunSessionId, @@ -520,6 +521,14 @@ export function isEmbeddedAgentRunHandleActive(sessionId: string): boolean { return active; } +export function isEmbeddedAgentRunAbortableForCompaction(sessionId: string): boolean { + const active = ACTIVE_EMBEDDED_RUNS.has(sessionId) || isReplyRunAbortableForCompaction(sessionId); + if (active) { + diag.debug(`run compact abort check: sessionId=${sessionId} active=true`); + } + return active; +} + export function isEmbeddedAgentRunStreaming(sessionId: string): boolean { const handle = ACTIVE_EMBEDDED_RUNS.get(sessionId); if (!handle) { diff --git a/src/agents/embedded-agent.ts b/src/agents/embedded-agent.ts index 65aa6972aa4d..845aab58e948 100644 --- a/src/agents/embedded-agent.ts +++ b/src/agents/embedded-agent.ts @@ -10,7 +10,9 @@ export { abortAndDrainEmbeddedAgentRun, abortEmbeddedAgentRun, compactEmbeddedAgentSession, + isEmbeddedAgentRunAbortableForCompaction, isEmbeddedAgentRunActive, + isEmbeddedAgentRunHandleActive, isEmbeddedAgentRunStreaming, queueEmbeddedAgentMessage, queueEmbeddedAgentMessageWithOutcome, diff --git a/src/auto-reply/reply/commands-compact.runtime.ts b/src/auto-reply/reply/commands-compact.runtime.ts index 996a0b460783..064b308e4395 100644 --- a/src/auto-reply/reply/commands-compact.runtime.ts +++ b/src/auto-reply/reply/commands-compact.runtime.ts @@ -2,7 +2,7 @@ export { abortEmbeddedAgentRun, compactEmbeddedAgentSession, - isEmbeddedAgentRunActive, + isEmbeddedAgentRunAbortableForCompaction, waitForEmbeddedAgentRunEnd, } from "../../agents/embedded-agent.js"; export { diff --git a/src/auto-reply/reply/commands-compact.test.ts b/src/auto-reply/reply/commands-compact.test.ts index 8ef1143fba28..f0d18c16d249 100644 --- a/src/auto-reply/reply/commands-compact.test.ts +++ b/src/auto-reply/reply/commands-compact.test.ts @@ -14,7 +14,7 @@ vi.mock("./commands-compact.runtime.js", () => ({ formatContextUsageShort: vi.fn(() => "Context 12.1k"), formatTokenCount: vi.fn((value: number) => `${value}`), incrementCompactionCount: vi.fn(), - isEmbeddedAgentRunActive: vi.fn().mockReturnValue(false), + isEmbeddedAgentRunAbortableForCompaction: vi.fn().mockReturnValue(false), resolveFreshSessionTotalTokens: vi.fn(() => 12_345), resolveSessionFilePath: vi.fn(() => "/tmp/session.json"), resolveSessionFilePathOptions: vi.fn(() => ({})), @@ -22,10 +22,13 @@ vi.mock("./commands-compact.runtime.js", () => ({ })); const { + abortEmbeddedAgentRun, compactEmbeddedAgentSession, formatContextUsageShort, incrementCompactionCount, + isEmbeddedAgentRunAbortableForCompaction, resolveSessionFilePathOptions, + waitForEmbeddedAgentRunEnd, } = await import("./commands-compact.runtime.js"); const { handleCompactCommand } = await import("./commands-compact.js"); @@ -191,6 +194,62 @@ describe("handleCompactCommand", () => { expect(call.senderE164).toBe("+15551234567"); expect(call.agentDir).toBe("/tmp/openclaw-agent-compact"); expect(call.authProfileId).toBe("github-copilot:work"); + expect(vi.mocked(abortEmbeddedAgentRun)).not.toHaveBeenCalled(); + expect(vi.mocked(waitForEmbeddedAgentRunEnd)).not.toHaveBeenCalled(); + }); + + it("does not abort the command reply run before compacting", async () => { + vi.mocked(isEmbeddedAgentRunAbortableForCompaction).mockReturnValueOnce(false); + vi.mocked(compactEmbeddedAgentSession).mockResolvedValueOnce({ + ok: true, + compacted: false, + }); + + const result = await handleCompactCommand( + { + ...buildCompactParams("/compact", { + commands: { text: true }, + channels: { whatsapp: { allowFrom: ["*"] } }, + } as OpenClawConfig), + sessionEntry: { + sessionId: "session-1", + updatedAt: Date.now(), + }, + } as HandleCommandsParams, + true, + ); + + expect(result?.shouldContinue).toBe(false); + expect(vi.mocked(isEmbeddedAgentRunAbortableForCompaction)).toHaveBeenCalledWith("session-1"); + expect(vi.mocked(abortEmbeddedAgentRun)).not.toHaveBeenCalled(); + expect(vi.mocked(waitForEmbeddedAgentRunEnd)).not.toHaveBeenCalled(); + expect(vi.mocked(compactEmbeddedAgentSession)).toHaveBeenCalledOnce(); + }); + + it("aborts an active embedded run before compacting", async () => { + vi.mocked(isEmbeddedAgentRunAbortableForCompaction).mockReturnValueOnce(true); + vi.mocked(compactEmbeddedAgentSession).mockResolvedValueOnce({ + ok: true, + compacted: false, + }); + + await handleCompactCommand( + { + ...buildCompactParams("/compact", { + commands: { text: true }, + channels: { whatsapp: { allowFrom: ["*"] } }, + } as OpenClawConfig), + sessionEntry: { + sessionId: "session-1", + updatedAt: Date.now(), + }, + } as HandleCommandsParams, + true, + ); + + expect(vi.mocked(abortEmbeddedAgentRun)).toHaveBeenCalledWith("session-1"); + expect(vi.mocked(waitForEmbeddedAgentRunEnd)).toHaveBeenCalledWith("session-1", 15_000); + expect(vi.mocked(compactEmbeddedAgentSession)).toHaveBeenCalledOnce(); }); it("treats already-under-target manual compaction as skipped", async () => { diff --git a/src/auto-reply/reply/commands-compact.ts b/src/auto-reply/reply/commands-compact.ts index 41b2988c797e..70ada78321df 100644 --- a/src/auto-reply/reply/commands-compact.ts +++ b/src/auto-reply/reply/commands-compact.ts @@ -217,7 +217,7 @@ export const handleCompactCommand: CommandHandler = async (params) => { } const runtime = await loadCompactRuntime(); const sessionId = targetSessionEntry.sessionId; - if (runtime.isEmbeddedAgentRunActive(sessionId)) { + if (runtime.isEmbeddedAgentRunAbortableForCompaction(sessionId)) { runtime.abortEmbeddedAgentRun(sessionId); await runtime.waitForEmbeddedAgentRunEnd(sessionId, 15_000); } diff --git a/src/auto-reply/reply/reply-run-registry.test.ts b/src/auto-reply/reply/reply-run-registry.test.ts index 8b2987b1eccc..d8a13facb2aa 100644 --- a/src/auto-reply/reply/reply-run-registry.test.ts +++ b/src/auto-reply/reply/reply-run-registry.test.ts @@ -11,6 +11,7 @@ import { createReplyOperation, forceClearReplyRunBySessionId, isReplyRunActiveForSessionId, + isReplyRunAbortableForCompaction, queueReplyRunMessage, replyRunRegistry, resolveActiveReplyRunSessionId, @@ -58,6 +59,21 @@ describe("reply run registry", () => { } }); + it("treats queued reply operations as non-abortable for compaction", () => { + const operation = createReplyOperation({ + sessionKey: "agent:main:main", + sessionId: "session-compact", + resetTriggered: false, + }); + + expect(isReplyRunActiveForSessionId("session-compact")).toBe(true); + expect(isReplyRunAbortableForCompaction("session-compact")).toBe(false); + + operation.setPhase("running"); + + expect(isReplyRunAbortableForCompaction("session-compact")).toBe(true); + }); + it("mirrors active reply operations into diagnostic work state", () => { const operation = createReplyOperation({ sessionKey: "agent:main:telegram:direct:chat-1", diff --git a/src/auto-reply/reply/reply-run-registry.ts b/src/auto-reply/reply/reply-run-registry.ts index 85410c85cb23..ceceae6ea5de 100644 --- a/src/auto-reply/reply/reply-run-registry.ts +++ b/src/auto-reply/reply/reply-run-registry.ts @@ -518,6 +518,11 @@ export function isReplyRunActiveForSessionId(sessionId: string): boolean { return resolveReplyRunForCurrentSessionId(sessionId) !== undefined; } +export function isReplyRunAbortableForCompaction(sessionId: string): boolean { + const operation = resolveReplyRunForCurrentSessionId(sessionId); + return Boolean(operation && operation.phase !== "queued"); +} + export function isReplyRunStreamingForSessionId(sessionId: string): boolean { const operation = resolveReplyRunForCurrentSessionId(sessionId); if (!operation || operation.phase !== "running") { diff --git a/src/gateway/server-methods/chat.directive-tags.test.ts b/src/gateway/server-methods/chat.directive-tags.test.ts index 49fcc8a587e4..e0b88b586b93 100644 --- a/src/gateway/server-methods/chat.directive-tags.test.ts +++ b/src/gateway/server-methods/chat.directive-tags.test.ts @@ -57,6 +57,7 @@ const mockState = vi.hoisted(() => ({ replyToId?: string; replyToCurrent?: boolean; isReasoning?: boolean; + isStatusNotice?: boolean; isError?: boolean; }; }>, @@ -1490,6 +1491,38 @@ describe("chat directive tag stripping for non-streaming final payloads", () => ]); }); + it("broadcasts agent-run status notices without source reply mirrors", async () => { + createTranscriptFixture("openclaw-chat-send-agent-status-notice-"); + mockState.triggerAgentRunStart = true; + mockState.dispatchedReplies = [ + { + kind: "final", + payload: { + text: "⚙️ Codex compaction started • Context 2k/200k", + isStatusNotice: true, + }, + }, + ]; + const respond = vi.fn(); + const context = createChatContext(); + + const broadcast = await runNonStreamingChatSend({ + context, + respond, + idempotencyKey: "idem-agent-status-notice", + message: "/compact", + }); + + expect(broadcast).toMatchObject({ + runId: "idem-agent-status-notice", + sessionKey: "main", + state: "final", + }); + expect(extractFirstTextBlock(broadcast)).toBe("⚙️ Codex compaction started • Context 2k/200k"); + const assistantEntries = await readActiveAssistantTranscriptMessages(); + expect(assistantEntries).toStrictEqual([]); + }); + it("does not duplicate media-bearing internal-ui source replies in the transcript", async () => { await withTranscriptFixtureState( "openclaw-chat-send-agent-source-reply-media-", @@ -2175,6 +2208,48 @@ describe("chat directive tag stripping for non-streaming final payloads", () => }); }); + it("broadcasts returned agent errors after status notices", async () => { + createTranscriptFixture("openclaw-chat-send-agent-status-notice-error-"); + const errorMessage = "LLM idle timeout (120s): no response from model"; + mockState.triggerAgentRunStart = true; + mockState.dispatchedReplies = [ + { + kind: "final", + payload: { + text: "⚙️ Codex compaction started • Context 2k/200k", + isStatusNotice: true, + }, + }, + { + kind: "final", + payload: { + text: errorMessage, + isError: true, + }, + }, + ]; + const respond = vi.fn(); + const context = createChatContext(); + + const broadcast = await runNonStreamingChatSend({ + context, + respond, + idempotencyKey: "idem-agent-status-notice-error", + message: "/compact", + }); + + expect(broadcast).toMatchObject({ + runId: "idem-agent-status-notice-error", + sessionKey: "main", + state: "error", + errorMessage, + }); + const finalBroadcasts = ( + context.broadcast as unknown as ReturnType + ).mock.calls.filter(([, payload]) => (payload as { state?: unknown })?.state === "final"); + expect(finalBroadcasts).toStrictEqual([]); + }); + it("broadcasts returned agent-run error payloads after an agent starts", async () => { createTranscriptFixture("openclaw-chat-send-agent-returned-error-"); const errorMessage = "LLM idle timeout (120s): no response from model"; diff --git a/src/gateway/server-methods/chat.ts b/src/gateway/server-methods/chat.ts index 21ce01ddf328..ebccf1b1803e 100644 --- a/src/gateway/server-methods/chat.ts +++ b/src/gateway/server-methods/chat.ts @@ -46,7 +46,11 @@ import { resolveAgentTimeoutMs } from "../../agents/timeout.js"; import type { ModelCatalogEntry } from "../../agents/model-catalog.types.js"; import { modelCatalogBrowseRequiresFullDiscovery } from "../../agents/model-catalog-browse.js"; import { dispatchInboundMessage } from "../../auto-reply/dispatch.js"; -import { getReplyPayloadMetadata, type ReplyPayload } from "../../auto-reply/reply-payload.js"; +import { + getReplyPayloadMetadata, + isReplyPayloadStatusNotice, + type ReplyPayload, +} from "../../auto-reply/reply-payload.js"; import { createReplyDispatcher } from "../../auto-reply/reply/reply-dispatcher.js"; import { stageSandboxMedia } from "../../auto-reply/reply/stage-sandbox-media.js"; import type { MsgContext, TemplateContext } from "../../auto-reply/templating.js"; @@ -4043,17 +4047,25 @@ export const chatHandlers: GatewayRequestHandlers = { }); } } else { - const sourceReplyPayloads = deliveredReplies + const hasReturnedAgentErrorPayloads = returnedAgentErrorPayloads.length > 0; + const agentRunReplyPayloads = deliveredReplies .filter((entryEntry) => entryEntry.kind === "final") .map((entryResult) => entryResult.payload) - .filter(isSourceReplyTranscriptMirrorPayload); - if (sourceReplyPayloads.length > 0) { + .filter( + (payload) => + isSourceReplyTranscriptMirrorPayload(payload) || + (!hasReturnedAgentErrorPayloads && isReplyPayloadStatusNotice(payload)), + ); + if (agentRunReplyPayloads.length > 0) { + const hasSourceReplyTranscriptMirror = agentRunReplyPayloads.some( + isSourceReplyTranscriptMirrorPayload, + ); const finalPayloads = await normalizeWebchatReplyMediaPathsForDisplay({ cfg, sessionKey, agentId, accountId, - payloads: sourceReplyPayloads, + payloads: agentRunReplyPayloads, }); const { storePath: latestStorePath, entry: latestEntry } = loadSessionEntry( sessionKey, @@ -4100,11 +4112,11 @@ export const chatHandlers: GatewayRequestHandlers = { }, }); const combinedAssistantContent = - sourceReplyPayloads.length === 1 + agentRunReplyPayloads.length === 1 ? await buildReplyAssistantContent(finalPayloads) : undefined; const combinedMediaMessage = - sourceReplyPayloads.length === 1 + agentRunReplyPayloads.length === 1 ? await buildReplyMediaMessage(finalPayloads) : undefined; type SourceReplyContentState = { @@ -4115,17 +4127,17 @@ export const chatHandlers: GatewayRequestHandlers = { }; const sourceReplyContentStates: SourceReplyContentState[] = []; const sourceReplyBroadcastContent: AssistantDisplayContentBlock[] = []; - for (const [replyIndex] of sourceReplyPayloads.entries()) { + for (const [replyIndex] of agentRunReplyPayloads.entries()) { const finalPayload = finalPayloads[replyIndex]; if (!finalPayload) { continue; } const replyAssistantContent = - sourceReplyPayloads.length === 1 + agentRunReplyPayloads.length === 1 ? combinedAssistantContent : await buildReplyAssistantContent([finalPayload]); const replyMediaMessage = - sourceReplyPayloads.length === 1 + agentRunReplyPayloads.length === 1 ? combinedMediaMessage : await buildReplyMediaMessage([finalPayload]); const replyBroadcastContent = hasAssistantDisplayMediaContent( @@ -4163,7 +4175,10 @@ export const chatHandlers: GatewayRequestHandlers = { >["sourceReplyTranscriptMirror"]; state: SourceReplyContentState; }> = []; - for (const [replyIndex, sourceReplyPayload] of sourceReplyPayloads.entries()) { + for (const [ + replyIndex, + sourceReplyPayload, + ] of agentRunReplyPayloads.entries()) { const state = sourceReplyContentStates[replyIndex]; if (!state || !hasAssistantDisplayMediaContent(state.persistedContent)) { continue; @@ -4210,7 +4225,7 @@ export const chatHandlers: GatewayRequestHandlers = { for (const [ replyIndex, sourceReplyPayload, - ] of sourceReplyPayloads.entries()) { + ] of agentRunReplyPayloads.entries()) { if (!sourceReplyContentStates[replyIndex]) { continue; } @@ -4352,7 +4367,7 @@ export const chatHandlers: GatewayRequestHandlers = { agentId, message, }); - broadcastedSourceReplyFinal = true; + broadcastedSourceReplyFinal = hasSourceReplyTranscriptMirror; } } }