From 7329a8d3d7a01f3f3c106a66882d194c5c7be67b Mon Sep 17 00:00:00 2001 From: ClawSweeper Date: Fri, 21 Aug 2026 21:02:03 -0700 Subject: [PATCH 1/2] fix: Code Mode shell calls stall near yield deadline (#127759) * fix: prevent Code Mode shell waits from stalling Co-authored-by: Tak Hoffman <781889+Takhoffman@users.noreply.github.com> * fix: honor remaining Code Mode deadline --------- Co-authored-by: RoboClaw <309084314+roboclaw-bot@users.noreply.github.com> Co-authored-by: Tak Hoffman <781889+Takhoffman@users.noreply.github.com> --- src/agents/code-mode-bridge.ts | 20 ++++- src/agents/code-mode-execution.ts | 3 + src/agents/code-mode-headless.ts | 1 + src/agents/code-mode-state.ts | 3 + src/agents/code-mode.bridge.test.ts | 116 ++++++++++++++++++++++++++++ 5 files changed, 142 insertions(+), 1 deletion(-) diff --git a/src/agents/code-mode-bridge.ts b/src/agents/code-mode-bridge.ts index c252ac636c8f..ecc9cf917a71 100644 --- a/src/agents/code-mode-bridge.ts +++ b/src/agents/code-mode-bridge.ts @@ -380,6 +380,7 @@ export async function runBridgeRequest(params: { parentToolCallId: string; codeModeRunId: string; maxOutputBytes: number; + remainingMs: number; ctx: ToolSearchToolContext; request: PendingBridgeRequest; signal?: AbortSignal; @@ -439,7 +440,24 @@ export async function runBridgeRequest(params: { if (!binding) { throw new ToolInputError(`Unknown catalog function: ${callableName}.`); } - const called = await params.runtime.callExactId(binding.id, values[1] ?? {}, { + let input = values[1] ?? {}; + if ( + binding.source === "openclaw" && + binding.name === "exec" && + binding.input?.includes("yieldMs") === true && + isRecord(input) && + input.background !== true && + input.yieldMs === undefined + ) { + // The shell's 10s default equals Code Mode's default budget. Yield + // within the remaining shared deadline so late sequential calls can + // still return their process handle and resume the guest inline. + input = { + ...input, + yieldMs: Math.max(1, Math.min(1_000, Math.floor(params.remainingMs / 4))), + }; + } + const called = await params.runtime.callExactId(binding.id, input, { parentToolCallId: params.parentToolCallId, signal: params.signal, onUpdate: params.onUpdate, diff --git a/src/agents/code-mode-execution.ts b/src/agents/code-mode-execution.ts index 10e464293f97..6700a03cbfe4 100644 --- a/src/agents/code-mode-execution.ts +++ b/src/agents/code-mode-execution.ts @@ -334,6 +334,7 @@ async function settleCodeModeResult(params: { namespaceRuntime: params.namespaceRuntime, parentToolCallId: params.parentToolCallId, codeModeRunId: params.codeModeReplayId, + deadlineMs: settleDeadline, activeRunId, ctx: params.ctx, signal: params.signal, @@ -469,6 +470,7 @@ async function settleCodeModeResult(params: { namespaceRuntime: params.namespaceRuntime, parentToolCallId: params.parentToolCallId, codeModeRunId: params.codeModeReplayId, + deadlineMs: settleDeadline, activeRunId, ctx: params.ctx, signal: params.signal, @@ -512,6 +514,7 @@ async function settleCodeModeResult(params: { catalogProjection: params.catalogProjection, namespaceRuntime: params.namespaceRuntime, output, + deadlineMs: settleDeadline, deliveredOutputCount, reservedActiveRunSlot: params.reservedActiveRunSlot, replaySafe: params.replaySafe, diff --git a/src/agents/code-mode-headless.ts b/src/agents/code-mode-headless.ts index 2d264fcf1b6e..785f984bedce 100644 --- a/src/agents/code-mode-headless.ts +++ b/src/agents/code-mode-headless.ts @@ -307,6 +307,7 @@ export async function runCodeModeScriptHeadless(params: { namespaceRuntime, parentToolCallId, codeModeRunId, + deadlineMs: deadline, ctx: params.ctx, signal: abortScope.signal, }), diff --git a/src/agents/code-mode-state.ts b/src/agents/code-mode-state.ts index 865cc129a544..b74436e60308 100644 --- a/src/agents/code-mode-state.ts +++ b/src/agents/code-mode-state.ts @@ -238,6 +238,7 @@ export function snapshotState(params: { catalogProjection: CodeModeCatalogProjection; namespaceRuntime: CodeModeNamespaceRuntime; output: unknown[]; + deadlineMs: number; deliveredOutputCount?: number; reservedActiveRunSlot?: boolean; replaySafe: boolean; @@ -321,6 +322,7 @@ export function createPendingBridgeStates(params: { namespaceRuntime: CodeModeNamespaceRuntime; parentToolCallId: string; codeModeRunId: string; + deadlineMs: number; activeRunId?: string; ctx: ToolSearchToolContext; signal?: AbortSignal; @@ -342,6 +344,7 @@ export function createPendingBridgeStates(params: { parentToolCallId: params.parentToolCallId, codeModeRunId: params.codeModeRunId, maxOutputBytes: params.config.maxOutputBytes, + remainingMs: Math.max(1, params.deadlineMs - Date.now()), ctx: params.ctx, request, signal, diff --git a/src/agents/code-mode.bridge.test.ts b/src/agents/code-mode.bridge.test.ts index 68864953d796..58f0dfeecb1d 100644 --- a/src/agents/code-mode.bridge.test.ts +++ b/src/agents/code-mode.bridge.test.ts @@ -155,6 +155,122 @@ describe("Code Mode bridge settlement and cancellation", () => { expect(testing.activeRuns.size).toBe(0); }); + it("yields nested exec before the Code Mode deadline when continuation args are omitted", async () => { + const catalogRef = createToolSearchCatalogRef(); + const config = { + tools: { codeMode: { enabled: true, timeoutMs: 10_000 } }, + } as never; + const ctx = { + config, + runtimeConfig: config, + sessionId: "session-code-mode", + sessionKey: "agent:main:main", + runId: "run-code-mode", + catalogRef, + }; + const codeModeTools = createCodeModeTools(ctx); + const shell = pluginToolWithExecute("exec", "Run shell", async (_toolCallId, input) => + jsonResult(input), + ); + shell.parameters = Type.Object({ + command: Type.String(), + yieldMs: Type.Optional(Type.Number()), + background: Type.Optional(Type.Boolean()), + }); + applyCodeModeCatalog({ + tools: [...codeModeTools, shell], + config, + sessionId: "session-code-mode", + sessionKey: "agent:main:main", + runId: "run-code-mode", + catalogRef, + }); + + const details = resultDetails( + await expectDefined(codeModeTools[0], "Code Mode exec test invariant").execute( + "code-call-shell-yield", + { + code: `return [ + await exec({ command: "default" }), + await exec({ command: "explicit", yieldMs: 4_000 }), + await exec({ command: "background", background: true }), + ];`, + }, + ), + ); + + expect(details).toMatchObject({ + status: "completed", + value: [ + { command: "default", yieldMs: 1_000 }, + { command: "explicit", yieldMs: 4_000 }, + { command: "background", background: true }, + ], + }); + expect(testing.activeRuns.size).toBe(0); + }); + + it("bounds nested exec yield by the shared remaining deadline", async () => { + vi.useFakeTimers({ toFake: ["setTimeout", "clearTimeout", "Date"] }); + const catalogRef = createToolSearchCatalogRef(); + const config = { + tools: { codeMode: { enabled: true, timeoutMs: 10_000 } }, + } as never; + const ctx = { + config, + runtimeConfig: config, + sessionId: "session-code-mode", + sessionKey: "agent:main:main", + runId: "run-code-mode", + catalogRef, + }; + const codeModeTools = createCodeModeTools(ctx); + const consumeBudget = pluginToolWithExecute( + "fake_consume_budget", + "Consume most of the shared Code Mode deadline", + async () => { + vi.advanceTimersByTime(9_600); + return jsonResult({ consumed: true }); + }, + ); + const shell = pluginToolWithExecute("exec", "Run shell", async (_toolCallId, input) => + jsonResult(input), + ); + shell.parameters = Type.Object({ + command: Type.String(), + yieldMs: Type.Optional(Type.Number()), + background: Type.Optional(Type.Boolean()), + }); + applyCodeModeCatalog({ + tools: [...codeModeTools, consumeBudget, shell], + config, + sessionId: "session-code-mode", + sessionKey: "agent:main:main", + runId: "run-code-mode", + catalogRef, + }); + + const details = resultDetails( + await expectDefined(codeModeTools[0], "Code Mode exec test invariant").execute( + "code-call-late-shell-yield", + { + code: ` + await fake_consume_budget({}); + return await exec({ command: "late" }); + `, + }, + ), + ); + + expect(details).toMatchObject({ + status: "completed", + value: { command: "late", yieldMs: 100 }, + }); + expect(consumeBudget.execute).toHaveBeenCalledOnce(); + expect(shell.execute).toHaveBeenCalledOnce(); + expect(testing.activeRuns.size).toBe(0); + }); + it("supports a guest timer between an action and its observation", async () => { const catalogRef = createToolSearchCatalogRef(); const config = { From 8c1ea2382618e36235d57aaeb86faa989d32d7be Mon Sep 17 00:00:00 2001 From: ClawSweeper Date: Fri, 21 Aug 2026 21:03:58 -0700 Subject: [PATCH 2/2] fix: preserve internal reply images across gateway restart (#127729) Co-authored-by: RoboClaw <309084314+roboclaw-bot@users.noreply.github.com> Co-authored-by: Tak Hoffman <781889+Takhoffman@users.noreply.github.com> --- .../embedded-agent-messaging-extraction.ts | 3 + src/agents/embedded-agent-messaging.types.ts | 1 + .../embedded-agent-runner/run/payloads.ts | 3 + .../run/source-reply-payloads.ts | 3 +- ....internal-source-reply.integration.test.ts | 125 +++++++++++++++ src/auto-reply/reply-payload.ts | 2 + .../reply/dispatch-from-config.transcript.ts | 2 +- .../internal-source-reply-persistence.ts | 149 ++++++++++++++++++ src/infra/outbound/message-action-runner.ts | 36 +++++ 9 files changed, 322 insertions(+), 2 deletions(-) create mode 100644 src/gateway/internal-source-reply-persistence.ts diff --git a/src/agents/embedded-agent-messaging-extraction.ts b/src/agents/embedded-agent-messaging-extraction.ts index 4498dd2bd488..9c39589e0a58 100644 --- a/src/agents/embedded-agent-messaging-extraction.ts +++ b/src/agents/embedded-agent-messaging-extraction.ts @@ -74,6 +74,9 @@ export function extractMessagingToolSourceReplyPayload( if (idempotencyKey) { payload.idempotencyKey = idempotencyKey; } + if (details.sourceReplyTranscriptOwner === true) { + payload.transcriptOwner = true; + } return Object.keys(payload).length > 0 ? payload : undefined; } diff --git a/src/agents/embedded-agent-messaging.types.ts b/src/agents/embedded-agent-messaging.types.ts index 67a1aa5ac30b..a59b73c6acf7 100644 --- a/src/agents/embedded-agent-messaging.types.ts +++ b/src/agents/embedded-agent-messaging.types.ts @@ -29,6 +29,7 @@ export type MessagingToolSourceReplyPayload = Pick< | "text" > & { idempotencyKey?: string; + transcriptOwner?: true; /** Current-source progress (`false`) or completed reply (`true`). */ sourceReplyFinal?: boolean; }; diff --git a/src/agents/embedded-agent-runner/run/payloads.ts b/src/agents/embedded-agent-runner/run/payloads.ts index 2a5941a28e70..5575c7cd822f 100644 --- a/src/agents/embedded-agent-runner/run/payloads.ts +++ b/src/agents/embedded-agent-runner/run/payloads.ts @@ -530,6 +530,9 @@ export function buildEmbeddedRunPayloads(params: { if (item.sourceReplyMirror.idempotencyKey) { sourceReplyTranscriptMirror.idempotencyKey = item.sourceReplyMirror.idempotencyKey; } + if (item.sourceReplyMirror.transcriptOwner) { + sourceReplyTranscriptMirror.transcriptOwner = true; + } setReplyPayloadMetadata(payload, { sourceReplyTranscriptMirror, }); diff --git a/src/agents/embedded-agent-runner/run/source-reply-payloads.ts b/src/agents/embedded-agent-runner/run/source-reply-payloads.ts index d1182733b6e3..8fda0c72237c 100644 --- a/src/agents/embedded-agent-runner/run/source-reply-payloads.ts +++ b/src/agents/embedded-agent-runner/run/source-reply-payloads.ts @@ -23,7 +23,7 @@ type EmbeddedRunReplyItem = { interactive?: ReplyPayload["interactive"]; channelData?: Record; nonTerminalToolErrorWarning?: boolean; - sourceReplyMirror?: { idempotencyKey?: string }; + sourceReplyMirror?: { idempotencyKey?: string; transcriptOwner?: true }; }; /** Builds transcript mirrors and completion evidence for message-tool source replies. */ @@ -70,6 +70,7 @@ export function buildSourceReplyPayloadState(params: { idempotencyKey: payload.idempotencyKey ?? (params.runId ? `${params.runId}:internal-source-reply:${index}` : undefined), + ...(payload.transcriptOwner ? { transcriptOwner: true as const } : {}), }, }, ]; diff --git a/src/agents/tools/message-tool.internal-source-reply.integration.test.ts b/src/agents/tools/message-tool.internal-source-reply.integration.test.ts index 16a7c3bdf4fe..695d3acf57ac 100644 --- a/src/agents/tools/message-tool.internal-source-reply.integration.test.ts +++ b/src/agents/tools/message-tool.internal-source-reply.integration.test.ts @@ -5,6 +5,13 @@ import path from "node:path"; import { describe, expect, it } from "vitest"; import { getReplyPayloadMetadata } from "../../auto-reply/reply-payload.js"; import { buildReplyPayloads } from "../../auto-reply/reply/agent-runner-payloads.js"; +import { mirrorDeliveredReplyToTranscript } from "../../auto-reply/reply/dispatch-from-config.transcript.js"; +import { + loadTranscriptEvents, + replaceSessionEntry, +} from "../../config/sessions/session-accessor.js"; +import { resolveManagedOutgoingMediaArtifactDownload } from "../../gateway/managed-image-attachments.js"; +import { listManagedImageRecordEntries } from "../../gateway/managed-image-record-store.js"; import { withOpenClawTestState } from "../../test-utils/openclaw-test-state.js"; import { extractMessagingToolSourceReplyPayload } from "../embedded-agent-messaging-extraction.js"; import { buildEmbeddedRunPayloads } from "../embedded-agent-runner/run/payloads.js"; @@ -28,6 +35,9 @@ function createCurrentSourceMessageTool(params: { workspaceDir?: string } = {}) }); } +const TINY_PNG_BASE64 = + "iVBORw0KGgoAAAANSUhEUgAAAAEAAAABCAQAAAC1HAwCAAAAC0lEQVR42mP8/x8AAusB9Y9ZQmcAAAAASUVORK5CYII="; + describe("WebChat message tool internal source reply", () => { it("projects a real targetless send and preserves the automatic final reply", async () => { const tool = createCurrentSourceMessageTool(); @@ -132,4 +142,119 @@ describe("WebChat message tool internal source reply", () => { }, ); }); + + it("persists one managed image for overlapping internal source replies", async () => { + await withOpenClawTestState( + { layout: "state-only", prefix: "openclaw-internal-source-reply-" }, + async (state) => { + const stateDir = state.stateDir; + const workspaceDir = state.workspaceDir; + const storePath = path.join(stateDir, "agents", "main", "sessions", "sessions.json"); + const sessionKey = "agent:main:webchat:dm:restart-proof"; + const sessionId = "restart-proof-session"; + const imagePath = path.join(workspaceDir, "restart-proof.png"); + await fs.mkdir(workspaceDir, { recursive: true }); + await fs.writeFile(imagePath, Buffer.from(TINY_PNG_BASE64, "base64")); + + await replaceSessionEntry( + { agentId: "main", sessionKey, storePath }, + { sessionId, chatType: "direct", updatedAt: 1 }, + ); + const config = { + agents: { + entries: { + main: { default: true, workspace: workspaceDir }, + }, + }, + }; + const tool = createMessageTool({ + config, + currentChannelProvider: "webchat", + agentSessionKey: sessionKey, + runSessionKey: sessionKey, + sessionId, + agentId: "main", + runId: "restart-proof-run", + getScopedChannelsCommandSecretTargets: () => ({ targetIds: new Set() }), + resolveCommandSecretRefsViaGateway: async () => ({ + resolvedConfig: config, + diagnostics: [], + targetStatesByPath: {}, + hadUnresolvedTargets: false, + }), + }); + + const sendParams = { + action: "send" as const, + message: "Durable image reply", + media: imagePath, + }; + const [toolResult, overlappingResult] = await Promise.all([ + tool.execute("restart-proof-call", sendParams), + tool.execute("restart-proof-call", sendParams), + ]); + const sourceReply = extractMessagingToolSourceReplyPayload(toolResult); + expect(sourceReply).toMatchObject({ transcriptOwner: true }); + expect(overlappingResult.details).toMatchObject({ + idempotencyKey: sourceReply?.idempotencyKey, + sourceReplyTranscriptOwner: true, + }); + const sourcePayloads = buildEmbeddedRunPayloads({ + assistantTexts: [], + lastAssistant: undefined, + currentAssistant: undefined, + sessionKey, + agentId: "main", + sourceReplyDeliveryMode: "message_tool_only", + messagingToolSourceReplyPayloads: sourceReply ? [sourceReply] : [], + runId: "restart-proof-run", + verboseLevel: "off", + reasoningLevel: "off", + toolResultFormat: "plain", + }); + const mirror = getReplyPayloadMetadata( + sourcePayloads[0] as object, + )?.sourceReplyTranscriptMirror; + expect(mirror).toMatchObject({ transcriptOwner: true }); + await mirrorDeliveredReplyToTranscript({ + metadata: mirror ? { ...mirror, expectedSessionId: sessionId, storePath } : undefined, + cfg: config, + }); + const events = await loadTranscriptEvents({ + agentId: "main", + sessionId, + sessionKey, + storePath, + }); + const assistants = events + .map((event) => (event as { message?: Record }).message) + .filter((message) => message?.role === "assistant"); + expect(assistants).toHaveLength(1); + const assistant = assistants[0]; + const content = Array.isArray(assistant?.content) + ? (assistant.content as Array>) + : []; + const image = content.find((block) => block.type === "image"); + expect(toolResult.details).toMatchObject({ + sourceReplySink: "internal-ui", + idempotencyKey: expect.any(String), + }); + expect(content[0]).toEqual({ type: "text", text: "Durable image reply" }); + expect(image).toMatchObject({ + type: "image", + artifactId: expect.stringMatching(/^artifact_managed_image_/u), + }); + expect(JSON.stringify(assistant)).not.toContain(imagePath); + expect(listManagedImageRecordEntries({ stateDir, sessionKey })).toHaveLength(1); + await expect( + resolveManagedOutgoingMediaArtifactDownload({ + sessionKey, + agentId: "main", + artifactId: String(image?.artifactId), + stateDir, + }), + ).resolves.toMatchObject({ type: "image" }); + }, + ); + }); }); diff --git a/src/auto-reply/reply-payload.ts b/src/auto-reply/reply-payload.ts index 7d04aa090c7e..fd27c6722207 100644 --- a/src/auto-reply/reply-payload.ts +++ b/src/auto-reply/reply-payload.ts @@ -294,6 +294,8 @@ export type ReplyPayloadMetadata = { expectedSessionId?: string; /** Delivery stays live, but neither side may be appended to a transcript. */ transcriptWriteBlocked?: boolean; + /** The visible reply already owns its durable transcript row. */ + transcriptOwner?: boolean; text?: string; mediaUrls?: string[]; idempotencyKey?: string; diff --git a/src/auto-reply/reply/dispatch-from-config.transcript.ts b/src/auto-reply/reply/dispatch-from-config.transcript.ts index b4e76ca55695..7fbc6b64351b 100644 --- a/src/auto-reply/reply/dispatch-from-config.transcript.ts +++ b/src/auto-reply/reply/dispatch-from-config.transcript.ts @@ -30,7 +30,7 @@ export async function mirrorDeliveredReplyToTranscript(params: { cfg: OpenClawConfig; }): Promise { const mirror = params.metadata; - if (!mirror) { + if (!mirror || mirror.transcriptOwner) { return; } try { diff --git a/src/gateway/internal-source-reply-persistence.ts b/src/gateway/internal-source-reply-persistence.ts new file mode 100644 index 000000000000..04804f7ad6e1 --- /dev/null +++ b/src/gateway/internal-source-reply-persistence.ts @@ -0,0 +1,149 @@ +import type { ReplyPayload } from "../auto-reply/reply-payload.js"; +import { appendAssistantMessageToSessionTranscript } from "../config/sessions.js"; +import { resolveSessionStorePathCore } from "../config/sessions/paths.js"; +import { + findTranscriptEvent, + readTranscriptEventMessage, +} from "../config/sessions/session-accessor.sqlite-read.js"; +import { getOwnedSessionTranscriptWriterFence } from "../config/sessions/transcript-write-context.js"; +import type { OpenClawConfig } from "../config/types.openclaw.js"; +import { getAgentScopedMediaLocalRootsForSources } from "../media/local-roots.js"; +import { createKeyedFifoLeaseRegistry } from "../shared/keyed-fifo-lease.js"; +import { isOpenClawDeliveryMirrorAssistantMessage } from "../shared/transcript-only-openclaw-assistant.js"; +import { + attachManagedOutgoingMediaToMessage, + createManagedOutgoingMediaBlocks, +} from "./managed-image-attachments.js"; +import { prepareGatewayInjectedAssistantContent } from "./server-methods/chat-transcript-inject.js"; + +const internalSourceReplyPersistenceLeases = createKeyedFifoLeaseRegistry( + Symbol.for("openclaw.internalSourceReplyPersistenceLeases"), +); + +function collectSourceReplyMediaUrls(payload: ReplyPayload): string[] { + return Array.from( + new Set([...(payload.mediaUrl ? [payload.mediaUrl] : []), ...(payload.mediaUrls ?? [])]), + ).filter((value) => value.trim().length > 0); +} + +async function hasPersistedInternalSourceReply(params: { + cfg: OpenClawConfig; + sessionKey: string; + expectedSessionId?: string; + agentId?: string; + idempotencyKey?: string; +}): Promise { + if (!params.expectedSessionId || !params.idempotencyKey) { + return false; + } + const storePath = resolveSessionStorePathCore(params.cfg.session?.store, { + agentId: params.agentId, + }); + const found = await findTranscriptEvent( + { + agentId: params.agentId, + sessionId: params.expectedSessionId, + sessionKey: params.sessionKey, + storePath, + }, + (event) => { + const message = readTranscriptEventMessage(event); + return ( + message?.idempotencyKey === params.idempotencyKey && + isOpenClawDeliveryMirrorAssistantMessage(message) + ); + }, + ); + return found !== undefined; +} + +function resolveInternalSourceReplyPersistenceLeaseKey(params: { + sessionKey: string; + expectedSessionId?: string; + agentId?: string; + idempotencyKey?: string; +}): string | undefined { + if (!params.idempotencyKey) { + return undefined; + } + return JSON.stringify([ + params.agentId ?? "", + params.sessionKey, + params.expectedSessionId ?? "", + params.idempotencyKey, + ]); +} + +/** Persist the private WebChat source reply before its successful tool result becomes visible. */ +export async function persistInternalSourceReply(params: { + cfg: OpenClawConfig; + sessionKey: string; + expectedSessionId?: string; + agentId?: string; + payload: ReplyPayload; + idempotencyKey?: string; + sourceReplyFinal?: boolean; + toolCallId?: string; + sourceTurnId?: string; +}): Promise { + const leaseKey = resolveInternalSourceReplyPersistenceLeaseKey(params); + const lease = leaseKey ? internalSourceReplyPersistenceLeases.reserve([leaseKey]) : undefined; + await lease?.wait(); + try { + if (await hasPersistedInternalSourceReply(params)) { + return; + } + const mediaUrls = collectSourceReplyMediaUrls(params.payload); + const mediaBlocks = await createManagedOutgoingMediaBlocks({ + sessionKey: params.sessionKey, + agentId: params.agentId, + mediaUrls, + localRoots: getAgentScopedMediaLocalRootsForSources({ + cfg: params.cfg, + agentId: params.agentId, + mediaSources: mediaUrls, + }), + }); + const content: Array> = [ + ...(params.payload.text ? [{ type: "text", text: params.payload.text }] : []), + ...mediaBlocks, + ]; + const writerFence = getOwnedSessionTranscriptWriterFence(); + const appended = await appendAssistantMessageToSessionTranscript({ + agentId: params.agentId, + sessionKey: params.sessionKey, + ...(params.expectedSessionId ? { expectedSessionId: params.expectedSessionId } : {}), + ...(writerFence?.expectedLifecycleRevision !== undefined + ? { expectedLifecycleRevision: writerFence.expectedLifecycleRevision } + : {}), + ...(writerFence ? { expectedWriterRunId: writerFence.expectedWriterRunId } : {}), + content: prepareGatewayInjectedAssistantContent(content), + idempotencyKey: params.idempotencyKey, + ...(params.sourceReplyFinal !== undefined + ? { + deliveryMirror: { + kind: "message-tool-source-reply" as const, + final: params.sourceReplyFinal, + ...(params.toolCallId ? { toolCallId: params.toolCallId } : {}), + ...(params.sourceTurnId ? { sourceTurnId: params.sourceTurnId } : {}), + }, + } + : {}), + config: params.cfg, + }); + if (!appended.ok) { + throw new Error(`Internal source reply persistence failed: ${appended.reason}`); + } + if ( + mediaBlocks.length > 0 && + !attachManagedOutgoingMediaToMessage({ + messageId: appended.messageId, + blocks: mediaBlocks, + }) + ) { + throw new Error("Internal source reply media ownership could not be persisted"); + } + } finally { + lease?.release(); + } +} diff --git a/src/infra/outbound/message-action-runner.ts b/src/infra/outbound/message-action-runner.ts index be6181d77c29..82573d6e87d8 100644 --- a/src/infra/outbound/message-action-runner.ts +++ b/src/infra/outbound/message-action-runner.ts @@ -14,6 +14,7 @@ import { getAgentScopedMediaLocalRoots } from "../../media/local-roots.js"; import { resolveAgentScopedOutboundMediaAccess } from "../../media/read-capability.js"; import { readBooleanParam } from "../../plugin-sdk/boolean-param.js"; import { hasPollCreationParams } from "../../poll-params.js"; +import { createLazyRuntimeModule } from "../../shared/lazy-runtime.js"; import { INTERNAL_MESSAGE_CHANNEL } from "../../utils/message-channel.js"; import { formatErrorMessage } from "../errors.js"; import { throwIfAborted } from "./abort.js"; @@ -49,6 +50,10 @@ import { } from "./outbound-policy.js"; import { getRuntimeVisibleChannelPlugin } from "./runtime-visible-channels.js"; +const loadInternalSourceReplyPersistence = createLazyRuntimeModule( + () => import("../../gateway/internal-source-reply-persistence.js"), +); + export function getToolResult(result: MessageActionResult): AgentToolResult | undefined { return "toolResult" in result ? result.toolResult : undefined; } @@ -294,12 +299,37 @@ async function handleInternalSourceReplySendAction( } const sourceReplyMediaUrls = resolveSendableOutboundReplyParts(sourceReplyPayload).mediaUrls; const sourceReplyMessage = sourceReplyPayload.text ?? sourceReply.message; + const idempotencyKey = normalizeOptionalString(params.idempotencyKey); + let persistedIdempotencyKey: string | undefined; + let persistedTranscriptOwner = false; + if (!dryRun && input.sessionId) { + const sessionKey = input.sourceReplySessionKey ?? input.sessionKey; + if (!sessionKey) { + throw new Error("Internal source reply requires a session key"); + } + const { persistInternalSourceReply } = await loadInternalSourceReplyPersistence(); + await persistInternalSourceReply({ + cfg: input.cfg, + sessionKey, + expectedSessionId: input.sessionId, + agentId: input.agentId ?? resolveSessionAgentId({ sessionKey, config: input.cfg }), + payload: sourceReplyPayload, + idempotencyKey, + sourceReplyFinal: input.sourceReplyFinal, + toolCallId: input.sourceReplyToolCallId, + sourceTurnId: input.messageActionAuthorization?.toolContext?.currentSourceTurnId, + }); + persistedIdempotencyKey = idempotencyKey; + persistedTranscriptOwner = true; + } const payload = { status: "ok", deliveryStatus: dryRun ? "dry_run" : "sent", channel: INTERNAL_MESSAGE_CHANNEL, target: "current-run", sourceReplyDeliveryMode: input.sourceReplyDeliveryMode, + ...(persistedIdempotencyKey ? { idempotencyKey: persistedIdempotencyKey } : {}), + ...(persistedTranscriptOwner ? { sourceReplyTranscriptOwner: true as const } : {}), ...(dryRun ? {} : { sourceReplySink: "internal-ui" as const }), sourceReply: sourceReplyPayload, ...(sourceReplyMessage ? { message: sourceReplyMessage } : {}), @@ -328,6 +358,8 @@ function buildInternalSourceReplyToolResult(payload: { channel: ChannelId; target: string; sourceReplyDeliveryMode?: SourceReplyDeliveryMode; + idempotencyKey?: string; + sourceReplyTranscriptOwner?: true; sourceReplySink?: "internal-ui"; sourceReply: ReplyPayload; message?: string; @@ -340,6 +372,8 @@ function buildInternalSourceReplyToolResult(payload: { channel: ChannelId; target: string; sourceReplyDeliveryMode?: SourceReplyDeliveryMode; + idempotencyKey?: string; + sourceReplyTranscriptOwner?: true; sourceReplySink?: "internal-ui"; sourceReply: ReplyPayload; message?: string; @@ -364,6 +398,8 @@ function buildInternalSourceReplyToolResult(payload: { ...(payload.sourceReplyDeliveryMode ? { sourceReplyDeliveryMode: payload.sourceReplyDeliveryMode } : {}), + ...(payload.idempotencyKey ? { idempotencyKey: payload.idempotencyKey } : {}), + ...(payload.sourceReplyTranscriptOwner ? { sourceReplyTranscriptOwner: true as const } : {}), ...(payload.sourceReplySink ? { sourceReplySink: payload.sourceReplySink } : {}), sourceReply: payload.sourceReply, ...(payload.message ? { message: payload.message } : {}),