diff --git a/extensions/codex/src/app-server/dynamic-tool-build.ts b/extensions/codex/src/app-server/dynamic-tool-build.ts index ce784a5fe108..7395a92863bc 100644 --- a/extensions/codex/src/app-server/dynamic-tool-build.ts +++ b/extensions/codex/src/app-server/dynamic-tool-build.ts @@ -30,7 +30,6 @@ import { isSystemAgentOnlyCodexDynamicToolAllowlist, normalizeCodexDynamicToolName, } from "./dynamic-tool-profile.js"; -import { addCodexMessageToolOnlyFinalControl } from "./message-tool-final-control.js"; import { resolveCodexNodeExecToolOverrides, resolveCodexNativeExecutionPolicy, @@ -359,13 +358,9 @@ export async function buildDynamicTools(input: DynamicToolBuildParams) { run: buildOpenClawCodingTools, }) : buildOpenClawCodingTools(); - const codexScopedTools = addCodexMessageToolOnlyFinalControl( - allTools, - params.sourceReplyDeliveryMode, - ); toolBuildStages.mark("create-openclaw-coding-tools"); const preNormalizationDiagnostics: RuntimeToolSchemaDiagnostic[] = []; - const readableAllToolProjection = filterProviderNormalizableTools(codexScopedTools); + const readableAllToolProjection = filterProviderNormalizableTools(allTools); preNormalizationDiagnostics.push(...readableAllToolProjection.diagnostics); const webSearchPlan = resolveCodexWebSearchPlan({ config: params.config, diff --git a/extensions/codex/src/app-server/dynamic-tools.test.ts b/extensions/codex/src/app-server/dynamic-tools.test.ts index ce35366cba7e..142cbfa5a522 100644 --- a/extensions/codex/src/app-server/dynamic-tools.test.ts +++ b/extensions/codex/src/app-server/dynamic-tools.test.ts @@ -41,7 +41,6 @@ import { type JsonValue, } from "./protocol.js"; import type { CodexRemoteWorkspaceFileReader } from "./remote-workspace-media.js"; -import { settleCodexSourceReplyFinality } from "./source-reply-finality.js"; const CODEX_OPENCLAW_DYNAMIC_TOOL_NAMESPACE = "openclaw"; @@ -1638,7 +1637,7 @@ describe("createCodexDynamicToolBridge", () => { ]); }); - it("keeps omitted source-reply finality non-terminal until a successful attempt settles", async () => { + it("treats omitted source-reply finality as terminal", async () => { const bridge = createBridgeWithToolResult( "message", textToolResult("Sent.", { messageId: "imessage-6264" }), @@ -1651,19 +1650,15 @@ describe("createCodexDynamicToolBridge", () => { }); expect(result).toEqual(expectInputText("Sent.")); - expect(result.terminate).toBeUndefined(); + expect(result.terminate).toBe(true); expect(bridge.telemetry.didDeliverSourceReplyViaMessageTool).toBe(true); - expect(bridge.telemetry.messagingToolSentTargets.at(-1)).not.toHaveProperty("sourceReplyFinal"); - - expect(settleCodexSourceReplyFinality(bridge.telemetry, true)).toBe(true); - expect(bridge.telemetry.messagingToolSentTargets.at(-1)).toMatchObject({ sourceReplyFinal: true, }); expect(Object.keys(result)).not.toContain("terminate"); }); - it("settles omitted source-reply finality as progress when the attempt fails", async () => { + it("keeps omitted source-reply finality terminal when the tool requests termination", async () => { const bridge = createBridgeWithToolResult( "message", { @@ -1677,70 +1672,13 @@ describe("createCodexDynamicToolBridge", () => { action: "send", message: "visible reply", }); - expect(result.terminate).toBeUndefined(); - expect(settleCodexSourceReplyFinality(bridge.telemetry, false)).toBe(false); + expect(result.terminate).toBe(true); expect(bridge.telemetry.messagingToolSentTargets.at(-1)).toMatchObject({ - sourceReplyFinal: false, + sourceReplyFinal: true, }); }); - it("settles only the latest omitted source reply as final after success", async () => { - const bridge = createBridgeWithToolResult( - "message", - textToolResult("Sent.", { messageId: "imessage-6264" }), - { sourceReplyDeliveryMode: "message_tool_only" }, - ); - - await handleMessageToolCall(bridge, { action: "send", message: "first update" }); - await handleMessageToolCall(bridge, { action: "send", message: "second update" }); - settleCodexSourceReplyFinality(bridge.telemetry, true); - - expect( - bridge.telemetry.messagingToolSentTargets.map((target) => target.sourceReplyFinal), - ).toEqual([false, true]); - }); - - it("does not promote an omitted reply past a later explicit progress reply", async () => { - const bridge = createBridgeWithToolResult( - "message", - textToolResult("Sent.", { messageId: "imessage-6264" }), - { sourceReplyDeliveryMode: "message_tool_only" }, - ); - - await handleMessageToolCall(bridge, { action: "send", message: "first update" }); - await handleMessageToolCall(bridge, { - action: "send", - message: "still working", - final: false, - }); - settleCodexSourceReplyFinality(bridge.telemetry, true); - - expect( - bridge.telemetry.messagingToolSentTargets.map((target) => target.sourceReplyFinal), - ).toEqual([false, false]); - }); - - it("keeps a later explicit final reply authoritative over an omitted reply", async () => { - const bridge = createBridgeWithToolResult( - "message", - textToolResult("Sent.", { messageId: "imessage-6264" }), - { sourceReplyDeliveryMode: "message_tool_only" }, - ); - - await handleMessageToolCall(bridge, { action: "send", message: "first update" }); - await handleMessageToolCall(bridge, { - action: "send", - message: "finished", - final: true, - }); - settleCodexSourceReplyFinality(bridge.telemetry, true); - - expect( - bridge.telemetry.messagingToolSentTargets.map((target) => target.sourceReplyFinal), - ).toEqual([false, true]); - }); - it("honors explicit finality for delivered message-tool-only source replies", async () => { const bridge = createBridgeWithToolResult( "message", @@ -2094,7 +2032,7 @@ describe("createCodexDynamicToolBridge", () => { expect(Object.keys(result)).not.toContain("terminate"); }); - it("defers omitted finality even when the message tool returns legacy termination", async () => { + it("keeps omitted finality terminal when the message tool returns termination", async () => { const bridge = createBridgeWithToolResult( "message", { @@ -2114,12 +2052,8 @@ describe("createCodexDynamicToolBridge", () => { }); expect(result).toEqual(expectInputText("Sent.")); - expect(result.terminate).toBeUndefined(); + expect(result.terminate).toBe(true); expect(bridge.telemetry.didDeliverSourceReplyViaMessageTool).toBe(true); - expect(bridge.telemetry.messagingToolSentTargets.at(-1)).not.toHaveProperty("sourceReplyFinal"); - - settleCodexSourceReplyFinality(bridge.telemetry, true); - expect(bridge.telemetry.messagingToolSentTargets.at(-1)).toMatchObject({ sourceReplyFinal: true, }); @@ -2193,7 +2127,7 @@ describe("createCodexDynamicToolBridge", () => { arguments: { action: "inspect" }, }); - expect(firstResult.terminate).toBeUndefined(); + expect(firstResult.terminate).toBe(true); expect(bridge.telemetry.didSendViaMessagingTool).toBe(true); expect(secondResult).toEqual(expectInputText("No message sent.")); expect(secondResult.terminate).toBeUndefined(); diff --git a/extensions/codex/src/app-server/dynamic-tools.ts b/extensions/codex/src/app-server/dynamic-tools.ts index 2a420becca6b..54f4a335b1cd 100644 --- a/extensions/codex/src/app-server/dynamic-tools.ts +++ b/extensions/codex/src/app-server/dynamic-tools.ts @@ -79,7 +79,6 @@ import { prepareCodexRemoteWorkspaceMessageMedia, type CodexRemoteWorkspaceFileReader, } from "./remote-workspace-media.js"; -import { recordCodexSourceReplyDeliveryIntent } from "./source-reply-finality.js"; import { resolveCodexToolAbortTerminalReason } from "./tool-abort-terminal-reason.js"; type CodexDynamicToolHookContext = NonNullable< @@ -779,17 +778,12 @@ export function createCodexDynamicToolBridge(params: { toolName === "message" && !resultIsError && (rawResult.terminate === true || result.terminate === true); - const hasExplicitFinalControl = typeof executedArgs.final === "boolean"; const confirmedSourceReply = params.hookContext?.sourceReplyDeliveryMode === "message_tool_only" && toolName === "message" && (toolConfirmedSourceReply || deliveredSourceReply || receiptConfirmedSourceReply); - const sourceReplyFinal = confirmedSourceReply - ? hasExplicitFinalControl - ? executedArgs.final === true - : undefined - : undefined; - const sourceReplyRecord = collectToolTelemetry({ + const sourceReplyFinal = confirmedSourceReply ? executedArgs.final !== false : undefined; + collectToolTelemetry({ toolName, args: executedArgs, result, @@ -799,26 +793,19 @@ export function createCodexDynamicToolBridge(params: { messagingTarget: confirmedMessagingTarget, sourceReplyFinal, }); - if (confirmedSourceReply && sourceReplyRecord) { - recordCodexSourceReplyDeliveryIntent(telemetry, { - record: sourceReplyRecord, - final: sourceReplyFinal, - }); - } if (deliveredSourceReply || receiptConfirmedSourceReply || toolConfirmedSourceReply) { telemetry.didDeliverSourceReplyViaMessageTool = true; } - const defersInferredSourceReplyTermination = - confirmedSourceReply && executedArgs.final !== true; + const continuesSourceReplyProgress = confirmedSourceReply && sourceReplyFinal === false; withDynamicToolTermination( response, ((rawResult.terminate === true || result.terminate === true) && - !defersInferredSourceReplyTermination) || + !continuesSourceReplyProgress) || // Yield is an explicit owner-level turn handoff, not termination // inferred from source-reply delivery, so finality does not mask it. isToolResultYield(rawResult) || isToolResultYield(result) || - (confirmedSourceReply && executedArgs.final === true), + (confirmedSourceReply && sourceReplyFinal === true), ); const asyncStarted = isAsyncStartedToolResult(rawResult) || isAsyncStartedToolResult(result); diff --git a/extensions/codex/src/app-server/message-tool-final-control.ts b/extensions/codex/src/app-server/message-tool-final-control.ts deleted file mode 100644 index fb1f0a1c3453..000000000000 --- a/extensions/codex/src/app-server/message-tool-final-control.ts +++ /dev/null @@ -1,52 +0,0 @@ -import type { EmbeddedRunAttemptParams } from "openclaw/plugin-sdk/agent-harness-runtime"; -import { normalizeCodexDynamicToolName } from "./dynamic-tool-profile.js"; - -type MutableDynamicTool = { - name: string; - parameters?: unknown; -}; - -/** - * `final` is a Codex-only control for message-tool-only source delivery. Keep - * it on the projected Codex schema so other agent runtimes never receive an - * API contract they do not implement. - */ -export function addCodexMessageToolOnlyFinalControl( - tools: T[], - sourceReplyDeliveryMode: EmbeddedRunAttemptParams["sourceReplyDeliveryMode"], -): T[] { - if (sourceReplyDeliveryMode !== "message_tool_only") { - return tools; - } - // These tools are attempt-fresh. Mutating preserves their WeakMap ownership - // metadata without exposing a clone helper through the public plugin SDK. - for (const tool of tools) { - if (normalizeCodexDynamicToolName(tool.name) === "message") { - const mutableTool: MutableDynamicTool = tool; - mutableTool.parameters = addCodexMessageToolOnlyFinalParameter(mutableTool.parameters); - } - } - return tools; -} - -function addCodexMessageToolOnlyFinalParameter(parameters: unknown): unknown { - if (!parameters || typeof parameters !== "object" || Array.isArray(parameters)) { - return parameters; - } - const schema = parameters as Record; - const rawProperties = schema.properties; - if (!rawProperties || typeof rawProperties !== "object" || Array.isArray(rawProperties)) { - return parameters; - } - return { - ...schema, - properties: { - ...rawProperties, - final: { - type: "boolean", - description: - "Set false for progress or true to complete the current source reply. If omitted, OpenClaw continues and resolves the latest omitted source reply when the turn ends.", - }, - }, - }; -} diff --git a/extensions/codex/src/app-server/run-attempt-finalize.ts b/extensions/codex/src/app-server/run-attempt-finalize.ts index 27f74873ad5b..154ced4c00a4 100644 --- a/extensions/codex/src/app-server/run-attempt-finalize.ts +++ b/extensions/codex/src/app-server/run-attempt-finalize.ts @@ -31,7 +31,6 @@ import { import type { prepareCodexAttemptTurnRequest } from "./run-attempt-turn-request.js"; import type { CodexAttemptTurnState } from "./run-attempt-turn-state.js"; import { captureCodexSettledTurnFinalizationContext } from "./settled-turn-context.js"; -import { settleCodexSourceReplyFinality } from "./source-reply-finality.js"; import { normalizeCodexTrajectoryError, recordCodexTrajectoryCompletion } from "./trajectory.js"; import { codexTranscriptMirrorRuntime } from "./transcript-mirror.js"; import { @@ -270,10 +269,9 @@ export async function finalizeCodexAttempt( !effectiveTimedOut && (finalPromptError === null || finalPromptError === undefined) && (completedTurnStatus === "completed" || recoveredTurnWatchTimeout || locallyCompletedTurn); - // buildResult retains the bridge's delivery records. Resolve omitted final - // intent only after the authoritative turn outcome is known, before any - // terminal observer consumes the result. - const completedSourceReply = settleCodexSourceReplyFinality(toolBridge.telemetry, turnSucceeded); + const completedSourceReply = toolBridge.telemetry.messagingToolSentTargets.some( + (target) => target.sourceReplyFinal === true, + ); if (completedSourceReply) { // Harness classification only sees assistant/reasoning/plan projections. // A reply delivered entirely through the source message tool is visible diff --git a/extensions/codex/src/app-server/source-reply-finality.ts b/extensions/codex/src/app-server/source-reply-finality.ts deleted file mode 100644 index 5ed91b4f541a..000000000000 --- a/extensions/codex/src/app-server/source-reply-finality.ts +++ /dev/null @@ -1,44 +0,0 @@ -import type { - MessagingToolSend, - MessagingToolSourceReplyPayload, -} from "openclaw/plugin-sdk/agent-harness-runtime"; - -type SourceReplyDeliveryIntent = { - record: MessagingToolSend | MessagingToolSourceReplyPayload; - final: boolean | undefined; -}; - -const sourceReplyDeliveryIntents = new WeakMap(); - -/** Retain source-reply intent until the owning Codex turn has an authoritative outcome. */ -export function recordCodexSourceReplyDeliveryIntent( - owner: object, - intent: SourceReplyDeliveryIntent, -): void { - const intents = sourceReplyDeliveryIntents.get(owner); - if (intents) { - intents.push(intent); - return; - } - sourceReplyDeliveryIntents.set(owner, [intent]); -} - -/** Resolve omitted finality without changing explicit progress or final markers. */ -export function settleCodexSourceReplyFinality(owner: object, turnSucceeded: boolean): boolean { - const intents = sourceReplyDeliveryIntents.get(owner); - if (!intents) { - return false; - } - const lastIntent = intents.at(-1); - for (const intent of intents) { - if (intent.final !== undefined) { - continue; - } - // An omitted marker is progress until the owning turn succeeds. Only the - // latest omitted reply can complete the conversation; a later explicit - // progress/final marker remains authoritative. - intent.record.sourceReplyFinal = turnSucceeded && intent === lastIntent; - } - sourceReplyDeliveryIntents.delete(owner); - return turnSucceeded && intents.some((intent) => intent.record.sourceReplyFinal === true); -} diff --git a/extensions/codex/src/app-server/thread-prompt.ts b/extensions/codex/src/app-server/thread-prompt.ts index bf933f8fd824..bdebcf534d1e 100644 --- a/extensions/codex/src/app-server/thread-prompt.ts +++ b/extensions/codex/src/app-server/thread-prompt.ts @@ -84,7 +84,7 @@ function buildVisibleReplyInstruction( messageToolAvailable: boolean, ): string { if (params.sourceReplyDeliveryMode === "message_tool_only" && messageToolAvailable) { - return "Visible source replies are not automatically delivered for this run. Use `message(action=send)` for user-visible source-channel output. For progress, set `final=false`. When the message is the completed reply to the current source conversation, set `final=true`; OpenClaw stops after confirming delivery. If `final` is omitted, OpenClaw continues and resolves the latest omitted source reply only when the turn ends successfully. Do not repeat visible message content in your final answer."; + return "Visible source replies are not automatically delivered for this run. Use `message(action=send)` for user-visible source-channel output. For progress, set `final=false`. Set `final=true`, or omit it, for the completed reply to the current source conversation; OpenClaw stops after confirming delivery. Do not repeat visible message content in your final answer."; } if (messageToolAvailable) { return "For the current source conversation, reply normally in your final assistant message; OpenClaw will deliver it through the active source conversation. Use `message` for supported non-text actions in the current conversation, such as reacting to its current message. Reserve other `message` actions for explicit out-of-band sends or media/file delivery. Reactions are not delivered automatically."; diff --git a/src/agents/cli-runner.reliability.test.ts b/src/agents/cli-runner.reliability.test.ts index 89ce38b5f3ab..37f43cc85d87 100644 --- a/src/agents/cli-runner.reliability.test.ts +++ b/src/agents/cli-runner.reliability.test.ts @@ -984,6 +984,48 @@ describe("runCliAgent reliability", () => { expect(supervisorSpawnMock).toHaveBeenCalledTimes(1); }); + it("surfaces a CLI failure after a delivered progress reply", async () => { + supervisorSpawnMock.mockClear(); + supervisorSpawnMock.mockImplementationOnce(async (...args: unknown[]) => { + const input = args[0] as Parameters["spawn"]>[0]; + const captureHandle = markMcpLoopbackToolCallStarted({ + captureKey: input.env?.OPENCLAW_MCP_CLI_CAPTURE_KEY ?? "", + toolName: "message", + args: { action: "send", message: "still working", final: false }, + }); + if (!captureHandle) { + throw new Error("Expected message delivery capture"); + } + recordMcpLoopbackToolCallResult({ + captureHandle, + toolName: "message", + args: { action: "send", message: "still working", final: false }, + result: { status: "sent", messageId: "progress-1" }, + outcome: "completed", + }); + markMcpLoopbackToolCallFinished(captureHandle); + return makeManagedRun({ exitCode: 1, durationMs: 150, stderr: "failed after progress" }); + }); + const context = makeClaudePreparedContext({ + sessionKey: "agent:main:telegram:direct:chat123", + runId: "run-progress-failure", + }); + context.mcpDeliveryCapture = true; + context.params.sourceReplyDeliveryMode = "message_tool_only"; + context.params.messageChannel = "telegram"; + context.params.currentChannelId = "chat123"; + + const result = await runPreparedCliAgent(context); + + expect(result.messagingToolSentTargets).toEqual([ + expect.objectContaining({ sourceReplyFinal: false }), + ]); + expect(result.payloads).toEqual([ + { text: "The reply stopped after sending progress. Please try again.", isError: true }, + ]); + expect(supervisorSpawnMock).toHaveBeenCalledTimes(1); + }); + it("clears a soft-resumed binding after confirmed message send followed by failure", async () => { supervisorSpawnMock.mockClear(); supervisorSpawnMock.mockImplementationOnce(async (...args: unknown[]) => { @@ -1153,7 +1195,9 @@ describe("runCliAgent reliability", () => { expect(result.didSendViaMessagingTool).toBe(true); expect(result.didDeliverSourceReplyViaMessageTool).toBe(true); - expect(result.messagingToolSourceReplyPayloads).toEqual([{ text: "sent before failure" }]); + expect(result.messagingToolSourceReplyPayloads).toEqual([ + { text: "sent before failure", sourceReplyFinal: true }, + ]); expect(result.payloads).toEqual([{ text: "sent before failure" }]); expect(getReplyPayloadMetadata(result.payloads?.[0] as object)).toMatchObject({ deliverDespiteSourceReplySuppression: true, diff --git a/src/agents/cli-runner.ts b/src/agents/cli-runner.ts index 3644364d6c4c..7160d6392c37 100644 --- a/src/agents/cli-runner.ts +++ b/src/agents/cli-runner.ts @@ -66,6 +66,7 @@ import { claudeCliSessionTranscriptHasContent as claudeCliSessionTranscriptHasCo import { classifyFailoverReason, isFailoverErrorMessage } from "./embedded-agent-helpers.js"; import type { EmbeddedAgentRunResult } from "./embedded-agent-runner.js"; import { waitForDeferredTurnMaintenanceForSession } from "./embedded-agent-runner/context-engine-maintenance.js"; +import { resolveExplicitFinalSourceReplyDeliveryEvidence } from "./embedded-agent-runner/delivery-evidence.js"; import { resolveAuthProfileFailureReason } from "./embedded-agent-runner/run/auth-profile-failure-policy.js"; import { buildEmbeddedRunPayloads } from "./embedded-agent-runner/run/payloads.js"; import { FailoverError, isFailoverError, resolveFailoverStatus } from "./failover-error.js"; @@ -942,6 +943,7 @@ export async function runPreparedCliAgent( CliOutput, | "didSendViaMessagingTool" | "didDeliverSourceReplyViaMessageTool" + | "messagingToolSentTargets" | "messagingToolSourceReplyPayloads" >, ): ReplyPayload[] => { @@ -955,6 +957,7 @@ export async function runPreparedCliAgent( model: context.modelId, didSendViaMessagingTool: evidence.didSendViaMessagingTool, didDeliverSourceReplyViaMessageTool: evidence.didDeliverSourceReplyViaMessageTool, + messagingToolSentTargets: evidence.messagingToolSentTargets, messagingToolSourceReplyPayloads: evidence.messagingToolSourceReplyPayloads, sourceReplyDeliveryMode: params.sourceReplyDeliveryMode, agentId: params.agentId, @@ -967,6 +970,7 @@ export async function runPreparedCliAgent( CliOutput, | "didSendViaMessagingTool" | "didDeliverSourceReplyViaMessageTool" + | "messagingToolSentTargets" | "messagingToolSourceReplyPayloads" >, ) => { @@ -989,9 +993,15 @@ export async function runPreparedCliAgent( ): EmbeddedAgentRunResult => { const message = formatErrorMessage(error); const { payloads } = resolveCliSourceReplyMirror(evidence); + const visiblePayloads = + payloads.length > 0 + ? payloads + : resolveExplicitFinalSourceReplyDeliveryEvidence(evidence) === false + ? [{ text: "The reply stopped after sending progress. Please try again.", isError: true }] + : undefined; deliveredMessagingSideEffect = true; return { - ...(payloads.length > 0 ? { payloads } : {}), + ...(visiblePayloads ? { payloads: visiblePayloads } : {}), meta: { durationMs: Date.now() - context.started, systemPromptReport: context.systemPromptReport, diff --git a/src/agents/cli-runner/execute-tool-tracking.ts b/src/agents/cli-runner/execute-tool-tracking.ts index cf6907eb19f4..5b5b4e400b4a 100644 --- a/src/agents/cli-runner/execute-tool-tracking.ts +++ b/src/agents/cli-runner/execute-tool-tracking.ts @@ -11,6 +11,7 @@ import type { CliOutput, CliToolUseStartDelta } from "../cli-output-contracts.js import { isDeliveredMessageToolOnlySourceReplyResult, isDeliveredMessagingToolResult, + resolveMessageToolSourceReplyFinal, } from "../embedded-agent-message-tool-source-reply.js"; import { isMessagingTool, @@ -226,6 +227,18 @@ export function createCliToolTracking(context: PreparedCliRunContext) { const toolArgs = params.args ?? {}; const isMessagingSend = isMessagingToolSendAction(params.toolName, toolArgs); const content = isMessagingSend ? extractCliMessagingContent(toolArgs, params.result) : {}; + const deliveredCurrentSourceReply = + isMessagingSend && + isDeliveredMessageToolOnlySourceReplyResult({ + sourceReplyDeliveryMode: context.params.sourceReplyDeliveryMode, + toolName: params.toolName, + args: params.args, + result: params.result, + isError: params.isError, + }); + const sourceReplyFinal = deliveredCurrentSourceReply + ? resolveMessageToolSourceReplyFinal(toolArgs) + : undefined; if (isMessagingSend) { appendUniqueCliMessagingEvidence( messagingToolSentTexts, @@ -237,15 +250,7 @@ export function createCliToolTracking(context: PreparedCliRunContext) { messagingToolSentMediaUrlKeys, content.mediaUrls ?? [], ); - if ( - isDeliveredMessageToolOnlySourceReplyResult({ - sourceReplyDeliveryMode: context.params.sourceReplyDeliveryMode, - toolName: params.toolName, - args: params.args, - result: params.result, - isError: params.isError, - }) - ) { + if (deliveredCurrentSourceReply) { didDeliverSourceReplyViaMessageTool = true; const payload = extractMessagingToolSourceReplyPayload(params.result); if (payload) { @@ -254,7 +259,10 @@ export function createCliToolTracking(context: PreparedCliRunContext) { } // Each internal source-reply send is a distinct delivery, even when // two intentional sends have identical text or media. - messagingToolSourceReplyPayloads.push(payload); + messagingToolSourceReplyPayloads.push({ + ...payload, + ...(sourceReplyFinal !== undefined ? { sourceReplyFinal } : {}), + }); } } } @@ -264,6 +272,7 @@ export function createCliToolTracking(context: PreparedCliRunContext) { const targetWithContent = { ...extractMessagingToolSendResult(params.target, params.result), ...content, + ...(sourceReplyFinal !== undefined ? { sourceReplyFinal } : {}), }; const evidenceKey = buildMessagingToolSendEvidenceKey(targetWithContent); if (messagingToolSentTargetKeys.has(evidenceKey)) { diff --git a/src/agents/cli-runner/execute.supervisor-capture.test.ts b/src/agents/cli-runner/execute.supervisor-capture.test.ts index dcf9a43ee807..ee68689a6511 100644 --- a/src/agents/cli-runner/execute.supervisor-capture.test.ts +++ b/src/agents/cli-runner/execute.supervisor-capture.test.ts @@ -2440,10 +2440,12 @@ describe("executePreparedCliRun supervisor output capture", () => { { text: "implicit reply", mediaUrl: "https://example.com/implicit.png", + sourceReplyFinal: true, }, { text: "implicit reply", mediaUrl: "https://example.com/implicit.png", + sourceReplyFinal: true, }, ]); }); diff --git a/src/agents/embedded-agent-message-tool-source-reply.ts b/src/agents/embedded-agent-message-tool-source-reply.ts index b6d5e26a6867..df7753ee8573 100644 --- a/src/agents/embedded-agent-message-tool-source-reply.ts +++ b/src/agents/embedded-agent-message-tool-source-reply.ts @@ -33,6 +33,11 @@ const PARTIAL_DELIVERY_ENVELOPE_KEYS = [...RESULT_ENVELOPE_KEYS, "error", "cause const SESSIONS_SEND_DELIVERY_STATUSES = new Set(["accepted", "ok"]); const BARE_OK_DELIVERY_STATUS = "ok"; +/** Omission preserves the established one-shot send behavior. */ +export function resolveMessageToolSourceReplyFinal(args: unknown): boolean { + return (asOptionalRecord(args) ?? {}).final !== false; +} + function resultConfirmsCurrentSourceRoute(value: unknown): boolean { return ( (asOptionalRecord(asOptionalRecord(value)?.details) ?? {}).sourceReplyRoute === "current-source" diff --git a/src/agents/embedded-agent-messaging.types.ts b/src/agents/embedded-agent-messaging.types.ts index 4ad1097234a8..67a1aa5ac30b 100644 --- a/src/agents/embedded-agent-messaging.types.ts +++ b/src/agents/embedded-agent-messaging.types.ts @@ -14,7 +14,7 @@ export type MessagingToolSend = { text?: string; mediaUrls?: string[]; hasRichContent?: true; - /** Present only when Codex classified this current-source delivery intent. */ + /** Current-source progress (`false`) or completed reply (`true`). */ sourceReplyFinal?: boolean; }; @@ -29,6 +29,6 @@ export type MessagingToolSourceReplyPayload = Pick< | "text" > & { idempotencyKey?: string; - /** Present only when Codex classified this current-source delivery intent. */ + /** Current-source progress (`false`) or completed reply (`true`). */ sourceReplyFinal?: boolean; }; diff --git a/src/agents/embedded-agent-runner/run/message-tool-terminal.test.ts b/src/agents/embedded-agent-runner/run/message-tool-terminal.test.ts index 4ea2d92d83bb..aab88858c02c 100644 --- a/src/agents/embedded-agent-runner/run/message-tool-terminal.test.ts +++ b/src/agents/embedded-agent-runner/run/message-tool-terminal.test.ts @@ -204,12 +204,13 @@ describe("message-tool-only source replies", () => { ).resolves.toEqual({ content: [{ type: "text", text: "rewritten" }], details: { rewritten: true }, + terminate: true, }); expect(previousAfterToolCall).toHaveBeenCalledTimes(1); expect(onDeliveredSourceReply).toHaveBeenCalledTimes(1); }); - it("records delivery evidence without rewriting the default result", async () => { + it("terminates after a delivered completed source reply", async () => { const agent = {} as unknown as Agent; const onDeliveredSourceReply = vi.fn(); installMessageToolOnlyTerminalHook({ @@ -225,10 +226,27 @@ describe("message-tool-only source replies", () => { args: { action: "send", message: "visible reply" }, }), ), - ).resolves.toBeUndefined(); + ).resolves.toEqual({ terminate: true }); expect(onDeliveredSourceReply).toHaveBeenCalledTimes(1); }); + it("continues after delivered progress", async () => { + const agent = {} as unknown as Agent; + installMessageToolOnlyTerminalHook({ + agent, + sourceReplyDeliveryMode: "message_tool_only", + }); + + await expect( + agent.afterToolCall?.( + createAfterToolCallContext({ + toolName: "message", + args: { action: "send", message: "still working", final: false }, + }), + ), + ).resolves.toBeUndefined(); + }); + it("leaves existing after-tool-call output alone when the send failed", async () => { const previousAfterToolCall = vi.fn(async () => ({ content: [{ type: "text" as const, text: "failed" }], diff --git a/src/agents/embedded-agent-runner/run/message-tool-terminal.ts b/src/agents/embedded-agent-runner/run/message-tool-terminal.ts index f4697baf3bbb..2ba749729a41 100644 --- a/src/agents/embedded-agent-runner/run/message-tool-terminal.ts +++ b/src/agents/embedded-agent-runner/run/message-tool-terminal.ts @@ -2,7 +2,10 @@ import type { SourceReplyDeliveryMode } from "../../../auto-reply/get-reply-opti /** * Detects message-tool-only sends that delivered a visible source reply. */ -import { isDeliveredMessageToolOnlySourceReplyResult } from "../../embedded-agent-message-tool-source-reply.js"; +import { + isDeliveredMessageToolOnlySourceReplyResult, + resolveMessageToolSourceReplyFinal, +} from "../../embedded-agent-message-tool-source-reply.js"; import type { AfterToolCallContext, AfterToolCallResult, Agent } from "../../runtime/index.js"; function argsRecordForToolCall(context: AfterToolCallContext): Record { @@ -55,7 +58,9 @@ export function installMessageToolOnlyTerminalHook(params: { }) ) { params.onDeliveredSourceReply?.(); - return hookResult; + if (resolveMessageToolSourceReplyFinal(argsRecordForToolCall(context))) { + return { ...hookResult, terminate: true }; + } } return hookResult; }; diff --git a/src/agents/embedded-agent-subscribe.handlers.tools.test.ts b/src/agents/embedded-agent-subscribe.handlers.tools.test.ts index 3b60919eca0a..e16e3f0dcecb 100644 --- a/src/agents/embedded-agent-subscribe.handlers.tools.test.ts +++ b/src/agents/embedded-agent-subscribe.handlers.tools.test.ts @@ -3485,6 +3485,7 @@ describe("messaging tool media URL tracking", () => { it("commits internal-ui source replies from successful message sends", async () => { const { ctx } = createTestContext(); + ctx.params.sourceReplyDeliveryMode = "message_tool_only"; const startEvt: ToolExecutionStartEvent = { toolName: "message", @@ -3519,6 +3520,7 @@ describe("messaging tool media URL tracking", () => { mediaUrls: ["file:///tmp/reply.png"], channelData: { source: "tui" }, idempotencyKey: "stable-source-reply", + sourceReplyFinal: true, }, ]); }); diff --git a/src/agents/embedded-agent-subscribe.handlers.tools.ts b/src/agents/embedded-agent-subscribe.handlers.tools.ts index 2225afdd2a18..87ecd80855aa 100644 --- a/src/agents/embedded-agent-subscribe.handlers.tools.ts +++ b/src/agents/embedded-agent-subscribe.handlers.tools.ts @@ -57,6 +57,7 @@ import { isDeliveredMessageToolOnlySourceReplyResult, isDeliveredMessagingToolResult, readMessageToolSourceReplyText, + resolveMessageToolSourceReplyFinal, } from "./embedded-agent-message-tool-source-reply.js"; import { isMessagingTool, @@ -1644,6 +1645,18 @@ export async function handleToolExecutionEnd( didDeliverMessagingResult && isMessagingSend ? [...argumentMediaUrls, ...collectMessagingMediaUrlsFromToolResult(result)] : []; + const deliveredCurrentSourceReply = + didDeliverMessagingResult && + isDeliveredMessageToolOnlySourceReplyResult({ + sourceReplyDeliveryMode: ctx.params.sourceReplyDeliveryMode, + toolName, + args: startArgs, + result, + isError: isToolError, + }); + const sourceReplyFinal = deliveredCurrentSourceReply + ? resolveMessageToolSourceReplyFinal(startArgs) + : undefined; ctx.state.pendingMessagingTexts.delete(toolCallId); ctx.state.pendingMessagingTargets.delete(toolCallId); ctx.state.pendingMessagingMediaUrls.delete(toolCallId); @@ -1661,18 +1674,10 @@ export async function handleToolExecutionEnd( ...(messageText ? { text: messageText } : {}), ...(committedMediaUrls.length > 0 ? { mediaUrls: committedMediaUrls.slice() } : {}), ...(hasRichContent ? { hasRichContent: true as const } : {}), + ...(sourceReplyFinal !== undefined ? { sourceReplyFinal } : {}), }); ctx.trimMessagingToolSent(); } - const deliveredCurrentSourceReply = - didDeliverMessagingResult && - isDeliveredMessageToolOnlySourceReplyResult({ - sourceReplyDeliveryMode: ctx.params.sourceReplyDeliveryMode, - toolName, - args: startArgs, - result, - isError: isToolError, - }); if (deliveredCurrentSourceReply) { ctx.state.messageToolOnlySourceReplyDelivered = true; const sourceReplyText = readMessageToolSourceReplyText(startArgs); @@ -1692,7 +1697,10 @@ export async function handleToolExecutionEnd( } const sourceReplyPayload = extractMessagingToolSourceReplyPayload(result); if (sourceReplyPayload) { - ctx.state.messagingToolSourceReplyPayloads.push(sourceReplyPayload); + ctx.state.messagingToolSourceReplyPayloads.push({ + ...sourceReplyPayload, + ...(sourceReplyFinal !== undefined ? { sourceReplyFinal } : {}), + }); ctx.trimMessagingToolSent(); } } diff --git a/src/agents/system-prompt.ts b/src/agents/system-prompt.ts index 583491941533..98329d6a76b6 100644 --- a/src/agents/system-prompt.ts +++ b/src/agents/system-prompt.ts @@ -601,7 +601,7 @@ function buildMessagingSection(params: { }) { const messageToolOnly = params.sourceReplyDeliveryMode === "message_tool_only"; const visibleReplyInstruction = messageToolOnly - ? "- Current source visible reply MUST use `message(action=send)`; final text is private. Skip tool = user gets nothing. Brief tool-call progress is visible; no hidden instructions/private data/reasoning." + ? "- Current source visible reply MUST use `message(action=send)`; final text is private. Set `final=false` for progress. Set `final=true`, or omit it, for the completed reply. Skip tool = user gets nothing. No hidden instructions/private data/reasoning." : "- Current-session final text normally routes to source. If turn says final private, visible output uses `message(action=send)`."; const messageToolTargetInstruction = params.requireExplicitMessageTarget ? "- `send`: `target` + `message`; target required this turn." diff --git a/src/agents/tools/message-tool-description.ts b/src/agents/tools/message-tool-description.ts index 005f9fa6fe9c..a64a2b09de2a 100644 --- a/src/agents/tools/message-tool-description.ts +++ b/src/agents/tools/message-tool-description.ts @@ -14,7 +14,7 @@ export function appendMessageToolVisibleReplyHint( const targetGuidance = requireExplicitTarget ? "send needs target." : "target defaults current source; set only elsewhere."; - return `${description} This turn visible reply: action="send" + message; ${targetGuidance} Final answer private.`; + return `${description} This turn visible reply: action="send" + message; ${targetGuidance} Set final=false for progress. Set final=true, or omit it, for the completed reply. Final answer private.`; } export function appendMessageToolReadHint( diff --git a/src/agents/tools/message-tool.test.ts b/src/agents/tools/message-tool.test.ts index 8cdb6f313147..5e560ab62fbd 100644 --- a/src/agents/tools/message-tool.test.ts +++ b/src/agents/tools/message-tool.test.ts @@ -524,7 +524,7 @@ describe("message tool gateway timeout", () => { }); }); - it("does not advertise the Codex-only final delivery control", () => { + it("does not advertise source-reply finality on ordinary message tools", () => { expect(getToolProperties(createMessageTool())).not.toHaveProperty("final"); expect(getToolProperties(createMessageTool())).not.toHaveProperty("idempotencyKey"); }); @@ -636,7 +636,16 @@ describe("completion source-reply authority", () => { expect(getActionEnum(properties)).toEqual(["send"]); expect(Object.keys(properties).toSorted()).toEqual( - ["accountId", "action", "channel", "message", "replyTo", "target", "threadId"].toSorted(), + [ + "accountId", + "action", + "channel", + "final", + "message", + "replyTo", + "target", + "threadId", + ].toSorted(), ); expectStringSchema(properties.message, { description: "Text to send to the current source conversation.", @@ -778,7 +787,7 @@ describe("completion source-reply authority", () => { expect(mocks.runMessageAction).not.toHaveBeenCalled(); }); - it("allows Codex final controls and matched canonical source-thread text sends", async () => { + it("allows shared final controls and matched canonical source-thread text sends", async () => { mockSendResult({ channel: "discord", to: "channel:source" }); const tool = createRestrictedTool(); @@ -1045,11 +1054,13 @@ describe("message tool secret scoping", () => { const defaultTool = createMessageTool(); expect(scopedTool.description).toContain('visible reply: action="send" + message'); + expect(getToolProperties(scopedTool).final).toMatchObject({ type: "boolean" }); expect(scopedTool.description).toContain("target defaults current source"); expect(scopedTool.description).toContain("Final answer private"); expect(explicitTargetTool.description).toContain("send needs target"); expect(explicitTargetTool.description).not.toContain("target defaults current source"); expect(defaultTool.description).not.toContain('visible reply: action="send" + message'); + expect(getToolProperties(defaultTool)).not.toHaveProperty("final"); }); it("forwards source reply delivery mode through createOpenClawTools", () => { @@ -1059,6 +1070,7 @@ describe("message tool secret scoping", () => { }).find((candidate) => candidate.name === "message"); expect(tool?.description).toContain('visible reply: action="send" + message'); + expect(getToolProperties(tool!).final).toMatchObject({ type: "boolean" }); }); it("passes source reply delivery mode to the outbound runner", async () => { diff --git a/src/agents/tools/message-tool.ts b/src/agents/tools/message-tool.ts index ddc7c41db5f3..82dce7bd4b44 100644 --- a/src/agents/tools/message-tool.ts +++ b/src/agents/tools/message-tool.ts @@ -9,7 +9,7 @@ import { normalizeOptionalStringifiedId, } from "@openclaw/normalization-core/string-coerce"; import { sortUniqueStrings, uniqueValues } from "@openclaw/normalization-core/string-normalization"; -import { Type, type TSchema } from "typebox"; +import { Type, type TObject, type TSchema } from "typebox"; import { GATEWAY_CLIENT_IDS, GATEWAY_CLIENT_MODES, @@ -1095,6 +1095,22 @@ const SOURCE_REPLY_ONLY_MESSAGE_SCHEMA = Type.Object({ threadId: Type.Optional(Type.String()), }); const SOURCE_REPLY_ONLY_RUNTIME_ARG_NAMES = new Set(["to", "channelId", "final"]); +const SOURCE_REPLY_FINAL_PROPERTY = Type.Optional( + Type.Boolean({ + description: + "Set false for progress. Set true, or omit, for the completed current-source reply.", + }), +); + +function addSourceReplyFinalControl( + schema: T, + sourceReplyDeliveryMode: SourceReplyDeliveryMode | undefined, +): T | TObject { + if (sourceReplyDeliveryMode !== "message_tool_only") { + return schema; + } + return Type.Object({ ...schema.properties, final: SOURCE_REPLY_FINAL_PROPERTY }); +} function enforceSourceReplyOnlyTextDirectives(args: Record): void { if (typeof args.message !== "string" || !args.message.trim()) { @@ -1574,11 +1590,12 @@ export function createMessageTool(options?: MessageToolOptions): AnyAgentTool { const actions = messageToolDiscoveryParams ? resolveMessageToolActionSchemaActions(messageToolDiscoveryParams) : undefined; - const schema = options?.sourceReplyOnly + const baseSchema = options?.sourceReplyOnly ? SOURCE_REPLY_ONLY_MESSAGE_SCHEMA : messageToolDiscoveryParams ? buildMessageToolSchema(messageToolDiscoveryParams, actions ?? []) : MessageToolSchema; + const schema = addSourceReplyFinalControl(baseSchema, sourceReplySinkDeliveryMode); const description = options?.sourceReplyOnly ? appendMessageToolVisibleReplyHint( "Send a message to the current source conversation. Supports actions: send.", diff --git a/src/config/sessions/restart-recovery-receipt.test.ts b/src/config/sessions/restart-recovery-receipt.test.ts index bff08d7a4a6e..776450942bf8 100644 --- a/src/config/sessions/restart-recovery-receipt.test.ts +++ b/src/config/sessions/restart-recovery-receipt.test.ts @@ -58,7 +58,7 @@ describe("restart recovery terminal delivery receipt", () => { await seedClaim(); await beginRestartRecoveryTerminalDelivery(scope()); - await expect(beginRestartRecoveryTerminalDelivery(scope())).resolves.toBe("blocked"); + await expect(beginRestartRecoveryTerminalDelivery(scope())).resolves.toBe("delivery-ambiguous"); }); it.each([undefined, "done" as const])( @@ -103,7 +103,7 @@ describe("restart recovery terminal delivery receipt", () => { }, ); - await expect(beginRestartRecoveryTerminalDelivery(scope())).resolves.toBe("blocked"); + await expect(beginRestartRecoveryTerminalDelivery(scope())).resolves.toBe("already-delivered"); }); it("clears pending only after a proven non-delivery", async () => { diff --git a/src/config/sessions/restart-recovery-receipt.ts b/src/config/sessions/restart-recovery-receipt.ts index ec09072c3a90..a9433c3b16f4 100644 --- a/src/config/sessions/restart-recovery-receipt.ts +++ b/src/config/sessions/restart-recovery-receipt.ts @@ -54,7 +54,7 @@ function loadCurrent(scope: RestartRecoveryTerminalDeliveryScope): SessionEntry /** Persists ambiguity before a terminal external send is allowed to start. */ export async function beginRestartRecoveryTerminalDelivery( scope: RestartRecoveryTerminalDeliveryScope, -): Promise<"started" | "blocked" | "stale" | "not-applicable"> { +): Promise<"started" | "already-delivered" | "delivery-ambiguous" | "stale" | "not-applicable"> { let started = false; const updated = await updateSessionEntry( { sessionKey: scope.sessionKey, storePath: scope.storePath }, @@ -89,7 +89,7 @@ export async function beginRestartRecoveryTerminalDelivery( current?.sessionId === scope.sessionId && hasRestartRecoveryTerminalRun(current, scope.sourceTurnId) ) { - return "blocked"; + return "already-delivered"; } // The gateway already verified a short-lived current-turn capability. Room // events intentionally persist no running recovery state, so only durable @@ -101,7 +101,9 @@ export async function beginRestartRecoveryTerminalDelivery( return "stale"; } if (current.restartRecoveryDeliveryReceiptState || current.restartRecoveryDeliveryToolCallId) { - return "blocked"; + return current.restartRecoveryDeliveryReceiptState === "delivered-terminal" + ? "already-delivered" + : "delivery-ambiguous"; } throw new Error("failed to persist terminal delivery intent"); } diff --git a/src/gateway/server-methods/send.test.ts b/src/gateway/server-methods/send.test.ts index c06917a9040b..d039e17b0da1 100644 --- a/src/gateway/server-methods/send.test.ts +++ b/src/gateway/server-methods/send.test.ts @@ -38,7 +38,9 @@ const mocks = vi.hoisted(() => ({ }), ), beginRestartRecoveryTerminalDelivery: vi.fn< - () => Promise<"started" | "blocked" | "stale" | "not-applicable"> + () => Promise< + "started" | "already-delivered" | "delivery-ambiguous" | "stale" | "not-applicable" + > >(async () => "started"), cancelRestartRecoveryTerminalDelivery: vi.fn(async () => "cleared" as const), completeRestartRecoveryTerminalDelivery: vi.fn(async () => "recorded" as const), @@ -3084,14 +3086,30 @@ describe("gateway send mirroring", () => { idempotencyKey: "idem-shared-source-message-action", }); - await runMessageActionRequest(request("progress"), identity(false)); - await runMessageActionRequest(request("terminal"), identity(true)); + const progress = await runMessageActionRequest(request("progress"), identity(false)); + const terminal = await runMessageActionRequest(request("terminal"), identity(true)); + mocks.beginRestartRecoveryTerminalDelivery.mockResolvedValueOnce("already-delivered"); + const repeatedTerminal = await runMessageActionRequest( + { + ...request("repeated terminal"), + idempotencyKey: "idem-repeated-terminal", + }, + identity(true), + ); expect(mocks.appendAssistantMessageToSessionTranscript.mock.calls).toHaveLength(2); expect(appendTranscriptCall(0)?.idempotencyKey).toBe("idem-shared-source-message-action"); expect(appendTranscriptCall(1)?.idempotencyKey).toBe( "idem-shared-source-message-action:terminal-receipt:channel-user:v1:shared-key", ); + expect(firstRespondCall(progress.respond)[0]).toBe(true); + expect(firstRespondCall(terminal.respond)[0]).toBe(true); + expect(firstRespondCall(repeatedTerminal.respond)[0]).toBe(true); + expect(firstRespondCall(repeatedTerminal.respond)[1]).toMatchObject({ + status: "already_delivered", + delivered: false, + }); + expect(mocks.dispatchChannelMessageAction).toHaveBeenCalledTimes(2); }); it("rejects a terminal source send without tool-call correlation before dispatch", async () => { @@ -3131,8 +3149,8 @@ describe("gateway send mirroring", () => { expect(mocks.appendAssistantMessageToSessionTranscript).toHaveBeenCalledOnce(); }); - it("blocks a repeated terminal send before provider dispatch", async () => { - mocks.beginRestartRecoveryTerminalDelivery.mockResolvedValueOnce("blocked"); + it("returns an already-delivered outcome for a repeated terminal send", async () => { + mocks.beginRestartRecoveryTerminalDelivery.mockResolvedValueOnce("already-delivered"); const { respond } = await runTelegramTerminalAction({ sessionId: "session-duplicate-terminal", idempotencyKey: "idem-duplicate-terminal", @@ -3141,7 +3159,11 @@ describe("gateway send mirroring", () => { message: "duplicate terminal", }); - expect(firstRespondCall(respond)[0]).toBe(false); + expect(firstRespondCall(respond)[0]).toBe(true); + expect(firstRespondCall(respond)[1]).toMatchObject({ + status: "already_delivered", + delivered: false, + }); expect(mocks.dispatchChannelMessageAction).not.toHaveBeenCalled(); }); diff --git a/src/gateway/server-methods/send.ts b/src/gateway/server-methods/send.ts index e34a2f02a81c..8e7715b50068 100644 --- a/src/gateway/server-methods/send.ts +++ b/src/gateway/server-methods/send.ts @@ -1009,10 +1009,19 @@ export const sendHandlers: GatewayRequestHandlers = { ? { sourceReplyFinal: trustedContext.sourceReplyFinal } : {}), }; - const terminalDeliveryReceipt = + const terminalDeliveryStart = trustedContext.sourceReplyFinal === true ? await beginTerminalSourceReplyDelivery(sourceReplyMirror) : undefined; + if (terminalDeliveryStart && "outcome" in terminalDeliveryStart) { + return createGatewayInflightSuccess({ + context, + dedupeKey, + payload: terminalDeliveryStart.result, + channel, + }); + } + const terminalDeliveryReceipt = terminalDeliveryStart; const gatewayClientScopes = client?.connect?.scopes ?? []; const handled = await dispatchChannelMessageAction({ channel, diff --git a/src/infra/outbound/message-action-execution.ts b/src/infra/outbound/message-action-execution.ts index 33912200234e..fa936fc39ee1 100644 --- a/src/infra/outbound/message-action-execution.ts +++ b/src/infra/outbound/message-action-execution.ts @@ -338,9 +338,13 @@ export async function executeGatewayAction(params: { sourceReplyFinal: params.input.sourceReplyFinal, toolCallId: params.input.sourceReplyToolCallId, }; - const terminalDeliveryReceipt = callerOwnsTerminalReceipt + const terminalDeliveryStart = callerOwnsTerminalReceipt ? await beginTerminalSourceReplyDelivery(sourceReplyMirror) : undefined; + if (terminalDeliveryStart && "outcome" in terminalDeliveryStart) { + return params.result(terminalDeliveryStart.result); + } + const terminalDeliveryReceipt = terminalDeliveryStart; let hadUnknownDeliveryOutcome = false; let payload: unknown; try { diff --git a/src/infra/outbound/message-action-gateway-reconciliation.test.ts b/src/infra/outbound/message-action-gateway-reconciliation.test.ts index 66e1b095fa52..804447285928 100644 --- a/src/infra/outbound/message-action-gateway-reconciliation.test.ts +++ b/src/infra/outbound/message-action-gateway-reconciliation.test.ts @@ -490,6 +490,48 @@ describe("runMessageAction plugin dispatch", () => { ); }); + it("returns an ambiguous terminal outcome without remote provider I/O", async () => { + const gatewayPlugin = createGatewayActionPlugin({ + pluginId: "gatewaychat", + label: "Gateway Chat", + blurb: "Gateway Chat ambiguous source reply test plugin.", + actions: ["send"], + messaging: { targetResolver: { looksLikeId: () => true } }, + handleAction: vi.fn(async () => jsonResult({ ok: true, local: true })), + }); + setTestPlugin(gatewayPlugin, "gatewaychat"); + mocks.beginTerminalSourceReplyDelivery.mockResolvedValue({ + outcome: "delivery_ambiguous", + result: { + status: "delivery_ambiguous", + delivered: false, + message: "The completed reply may already have been delivered. Do not retry it.", + }, + }); + + const result = await runMessageAction({ + cfg: { channels: { gatewaychat: { enabled: true } } } as OpenClawConfig, + action: "send", + params: { channel: "gatewaychat", target: "user-123", message: "terminal answer" }, + sourceReplyFinal: true, + sourceReplyToolCallId: "message-call-1", + gateway: { + terminalSourceReplyReceiptOwner: "caller", + clientName: GATEWAY_CLIENT_NAMES.GATEWAY_CLIENT, + mode: GATEWAY_CLIENT_MODES.BACKEND, + }, + dryRun: false, + }); + + expect(extractToolPayload(result)).toMatchObject({ + payload: { + status: "delivery_ambiguous", + delivered: false, + }, + }); + expect(mocks.callGatewayLeastPrivilege).not.toHaveBeenCalled(); + }); + it("cancels caller receipts after confirmed gateway request rejection", async () => { const gatewayPlugin = createGatewayActionPlugin({ pluginId: "gatewaychat", diff --git a/src/infra/outbound/source-reply-mirror.ts b/src/infra/outbound/source-reply-mirror.ts index bc0fb392b7e4..d1ee64db4828 100644 --- a/src/infra/outbound/source-reply-mirror.ts +++ b/src/infra/outbound/source-reply-mirror.ts @@ -43,6 +43,28 @@ type SourceReplyTranscriptMirrorParams = { replyToIsExplicit?: boolean; }; +export type TerminalSourceReplyDeliveryStart = + | TerminalSourceReplyDeliveryReceipt + | { + outcome: "already_delivered" | "delivery_ambiguous"; + result: { status: string; delivered: false; message: string }; + } + | undefined; + +function buildTerminalSourceReplyNoSendResult(outcome: "already_delivered" | "delivery_ambiguous") { + return { + outcome, + result: { + status: outcome, + delivered: false as const, + message: + outcome === "already_delivered" + ? "The completed reply was already delivered. Do not retry it." + : "The completed reply may already have been delivered. Do not retry it.", + }, + }; +} + type MirrorableSourceReplyTranscriptParams = SourceReplyTranscriptMirrorParams & { sessionKey: string; }; @@ -220,7 +242,7 @@ function resolveTerminalSourceReplyDeliveryReceipt( /** Arms the fail-closed state before a terminal source reply can reach a provider. */ export async function beginTerminalSourceReplyDelivery( params: SourceReplyTranscriptMirrorParams, -): Promise { +): Promise { const receipt = resolveTerminalSourceReplyDeliveryReceipt(params); if (!receipt) { return undefined; @@ -229,11 +251,11 @@ export async function beginTerminalSourceReplyDelivery( if (result === "not-applicable") { return undefined; } - if (result === "blocked") { - throw new Error("terminal source reply already has a durable delivery outcome"); + if (result === "already-delivered") { + return buildTerminalSourceReplyNoSendResult("already_delivered"); } - if (result === "stale") { - throw new Error("terminal source reply lost restart recovery ownership"); + if (result === "delivery-ambiguous" || result === "stale") { + return buildTerminalSourceReplyNoSendResult("delivery_ambiguous"); } return receipt; }