From 5a28a491b9a07034f7e23e0f4e2db928bba8a807 Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Tue, 18 Aug 2026 03:17:29 -0700 Subject: [PATCH] fix(media): surface silent attachment failures (#125562) * fix(media): surface silent attachment failures Missing structured images now produce a visible model-facing system note, failed playback transcodes warn once per source identity, and bodyless successful media responses are rejected instead of stored as empty files. * fix(media): persist failed image outcomes * fix(media): skip note lookup without failures --- .../run/attempt-execution-settle.test.ts | 27 +++++++ .../run/attempt-prompt-phase.test.ts | 32 +++++---- .../run/attempt-prompt-phase.ts | 1 - .../attempt-session-runtime-prepare.test.ts | 4 ++ .../run/attempt-session-runtime-prepare.ts | 8 +++ .../run/attempt-settle.ts | 43 +++++++++++ .../run/attempt-setup.ts | 2 + .../run/attempt-stream-settle.test.ts | 71 +++++++++++++++++++ .../run/attempt-stream-settle.ts | 2 + .../run/history-image-prune.test.ts | 39 +++++++++- .../embedded-agent-runner/run/images.ts | 29 +++++++- src/media/fetch.test.ts | 37 +++++++--- src/media/fetch.ts | 30 +++----- src/media/playback-transcode.test.ts | 43 ++++++++++- src/media/playback-transcode.ts | 11 ++- 15 files changed, 332 insertions(+), 47 deletions(-) diff --git a/src/agents/embedded-agent-runner/run/attempt-execution-settle.test.ts b/src/agents/embedded-agent-runner/run/attempt-execution-settle.test.ts index a1beceb441a3..9e736d91b502 100644 --- a/src/agents/embedded-agent-runner/run/attempt-execution-settle.test.ts +++ b/src/agents/embedded-agent-runner/run/attempt-execution-settle.test.ts @@ -123,7 +123,9 @@ function createFixture() { }; const sessionManager = { kind: "session-manager", + appendMessage: vi.fn((message) => messages.push(message)), buildSessionContext: vi.fn(() => ({ messages: [] })), + getSessionTarget: vi.fn(() => undefined), }; const hookRunner = { hasHooks: vi.fn(() => false) }; const cacheTrace = { recordStage: vi.fn() }; @@ -131,6 +133,7 @@ function createFixture() { const toolResultPromptProjectionState = { kind: "tool-result-projection" }; const sessionPromptState = { toolResults: toolResultPromptProjectionState }; const sessionRuntimeState = { + currentTurnImageFailureCount: 0, prePromptMessageCount: 2, promptCache: undefined, systemPromptText: "system prompt", @@ -319,6 +322,7 @@ function createFixture() { order, queueHandle, result, + sessionManager, sessionRuntimeState, state, subscription, @@ -399,6 +403,29 @@ describe("runEmbeddedAttemptSettledPhase", () => { ); }); + it("persists image failure notes after after-turn transcript reconciliation", async () => { + const fixture = createFixture(); + fixture.sessionRuntimeState.currentTurnImageFailureCount = 1; + await runEmbeddedAttemptSettledPhase(fixture.input); + + expect(fixture.sessionManager.appendMessage).toHaveBeenCalledWith( + expect.objectContaining({ + customType: "openclaw.system-note", + display: true, + content: expect.stringMatching(/1.*image contents.*unavailable.*resend.*not claim/is), + }), + ); + expect(mocks.completeResult).toHaveBeenCalledWith( + expect.objectContaining({ + state: expect.objectContaining({ + messagesSnapshot: expect.arrayContaining([ + expect.objectContaining({ customType: "openclaw.system-note", display: true }), + ]), + }), + }), + ); + }); + it("carries a successful hidden target through settlement into the terminal receipt", async () => { const fixture = createFixture(); fixture.input.prepared.toolBase.toolSearchTargetTranscriptProjections.push( diff --git a/src/agents/embedded-agent-runner/run/attempt-prompt-phase.test.ts b/src/agents/embedded-agent-runner/run/attempt-prompt-phase.test.ts index 5a13b3314a8b..fe21d2cdc57b 100644 --- a/src/agents/embedded-agent-runner/run/attempt-prompt-phase.test.ts +++ b/src/agents/embedded-agent-runner/run/attempt-prompt-phase.test.ts @@ -90,6 +90,17 @@ function createFixture() { yieldDetected: false, yieldMessage: null as string | null, }; + const activeSession = { + messages: [], + agent: { + state: { messages: [] }, + streamFn: vi.fn(), + }, + }; + const sessionManager = { + appendCustomEntry: vi.fn(), + getEntries: vi.fn(() => []), + }; let prePromptMessageCount = 1; const setPrePromptMessageCount = vi.fn((count: number) => { @@ -169,18 +180,6 @@ function createFixture() { submissionInput.onSteeringAcknowledged(); }); mocks.handlePromptError.mockResolvedValue({}); - - const activeSession = { - messages: [], - agent: { - state: { messages: [] }, - streamFn: vi.fn(), - }, - }; - const sessionManager = { - appendCustomEntry: vi.fn(), - getEntries: vi.fn(() => []), - }; const input = { attempt: { model: { id: "model-1", provider: "test" }, @@ -213,6 +212,7 @@ function createFixture() { toolResultPromptProjectionState: {}, }, execution: { + mediaOwnerAgentId: "main", effectiveFsWorkspaceOnly: false, effectiveWorkspace: "/tmp/workspace", sandbox: null, @@ -412,6 +412,14 @@ describe("runEmbeddedAttemptPromptPhase", () => { it("releases steering when preflight skips provider submission", async () => { const fixture = createFixture(); const promptError = new Error("preflight rejected"); + mocks.preparePromptExecution.mockResolvedValueOnce({ + images: [], + imageFactIndexes: [], + detectedRefs: [], + failedMediaCount: 1, + loadedCount: 0, + skippedCount: 1, + }); mocks.observePrompt.mockImplementationOnce(() => { fixture.order.push("observe"); return { skipPromptSubmission: true }; diff --git a/src/agents/embedded-agent-runner/run/attempt-prompt-phase.ts b/src/agents/embedded-agent-runner/run/attempt-prompt-phase.ts index 720b3632d170..a32982584bcf 100644 --- a/src/agents/embedded-agent-runner/run/attempt-prompt-phase.ts +++ b/src/agents/embedded-agent-runner/run/attempt-prompt-phase.ts @@ -312,7 +312,6 @@ export async function runEmbeddedAttemptPromptPhase(input: { prompt: promptContext.promptSubmission.prompt, skipPromptSubmission, }); - const reserveTokens = input.getCompactionReserveTokens(); let state: PromptPreflightInput["state"] = { ...input.lifecycle.readState(), diff --git a/src/agents/embedded-agent-runner/run/attempt-session-runtime-prepare.test.ts b/src/agents/embedded-agent-runner/run/attempt-session-runtime-prepare.test.ts index efdd95c79286..6916a21fe0c9 100644 --- a/src/agents/embedded-agent-runner/run/attempt-session-runtime-prepare.test.ts +++ b/src/agents/embedded-agent-runner/run/attempt-session-runtime-prepare.test.ts @@ -254,6 +254,7 @@ describe("prepareEmbeddedAttemptSessionRuntime", () => { }), ); expect(result.state).toEqual({ + currentTurnImageFailureCount: 0, prePromptMessageCount: 2, promptCache: undefined, systemPromptText: "runtime prompt", @@ -283,6 +284,9 @@ describe("prepareEmbeddedAttemptSessionRuntime", () => { expect(guardInput.getPromptCache()).toEqual({ cacheRead: 3 }); expect(guardInput.getPromptCacheRetention()).toBe("long"); expect(guardInput.getSystemPrompt()).toBe("updated prompt"); + guardInput.onCurrentTurnImageFailure(2); + guardInput.onCurrentTurnImageFailure(1); + expect(result.state.currentTurnImageFailureCount).toBe(2); }); it("publishes every cleanup owner before a later transport failure", async () => { diff --git a/src/agents/embedded-agent-runner/run/attempt-session-runtime-prepare.ts b/src/agents/embedded-agent-runner/run/attempt-session-runtime-prepare.ts index 93d5178b9da6..cbb6ffea0046 100644 --- a/src/agents/embedded-agent-runner/run/attempt-session-runtime-prepare.ts +++ b/src/agents/embedded-agent-runner/run/attempt-session-runtime-prepare.ts @@ -26,12 +26,14 @@ type TrajectoryInput = Parameters[0]; type AttemptSessionManager = ReturnType; type SessionSettleTracker = ReturnType; type TrajectoryRecorder = Awaited>; + type ExternalAbortController = Pick< ReturnType, "setActiveSessionAbort" >; type EmbeddedAttemptSessionRuntimeState = { + currentTurnImageFailureCount: number; prePromptMessageCount: number; promptCache: EmbeddedRunAttemptResult["promptCache"]; systemPromptText: string; @@ -109,6 +111,7 @@ export async function prepareEmbeddedAttemptSessionRuntime(input: { preparedSessionManager; const state: EmbeddedAttemptSessionRuntimeState = { + currentTurnImageFailureCount: 0, prePromptMessageCount: 0, promptCache: undefined, systemPromptText: input.initialSystemPrompt, @@ -136,6 +139,9 @@ export async function prepareEmbeddedAttemptSessionRuntime(input: { sessionManager, }); const { activeSession, setActiveSessionSystemPrompt, settingsManager } = preparedAgentSession; + const recordCurrentTurnImageFailure = (count: number) => { + state.currentTurnImageFailureCount = Math.max(state.currentTurnImageFailureCount, count); + }; await attempt.userTurnTranscriptRecorder?.waitForRuntimePersistence(); const boundary = prepareEmbeddedAttemptSessionBoundary({ activeSession, @@ -177,6 +183,7 @@ export async function prepareEmbeddedAttemptSessionRuntime(input: { effectiveWorkspace: input.effectiveWorkspace, getPrePromptMessageCount: () => state.prePromptMessageCount, getPromptCache: () => state.promptCache, + onCurrentTurnImageFailure: recordCurrentTurnImageFailure, getPromptCacheRetention: () => promptCacheRetentionRef.current, getSystemPrompt: () => state.systemPromptText, isOpenAIResponsesApi, @@ -234,6 +241,7 @@ export async function prepareEmbeddedAttemptSessionRuntime(input: { agentDir: input.agentDir, abortSignal: input.transport.abortSignal, getProviderRuntimeHandle: input.transport.getProviderRuntimeHandle, + onCurrentTurnImageFailure: recordCurrentTurnImageFailure, sandboxSessionKey: input.transport.sandboxSessionKey, ...(input.transport.sandbox !== undefined ? { sandbox: input.transport.sandbox } : {}), codeModeControlsEnabled: input.transport.codeModeControlsEnabled, diff --git a/src/agents/embedded-agent-runner/run/attempt-settle.ts b/src/agents/embedded-agent-runner/run/attempt-settle.ts index daeca702a0f4..2e9310a0efaa 100644 --- a/src/agents/embedded-agent-runner/run/attempt-settle.ts +++ b/src/agents/embedded-agent-runner/run/attempt-settle.ts @@ -2,6 +2,7 @@ * Settles prompt dispatch, stream cleanup, and result projection. * It may assume stream runtime preparation and session state are ready. */ +import { asOptionalRecord } from "@openclaw/normalization-core/record-coerce"; import type { AssistantMessage } from "../../../llm/types.js"; import { mergeAgentRunAttemptTerminal, @@ -11,6 +12,7 @@ import { } from "../../agent-run-terminal-outcome.js"; import { sanitizeCompactionReplayMessages } from "../../compaction-replay.js"; import type { AgentMessage } from "../../runtime/index.js"; +import { SessionManager } from "../../sessions/index.js"; import { settleRequesterAfterSessionSpawns } from "../../subagents/registry/subagent-registry.js"; import type { NormalizedUsage } from "../../usage.js"; import { log } from "../logger.js"; @@ -32,6 +34,7 @@ import type { prepareEmbeddedAttemptStream } from "./attempt-stream-prepare.js"; import { settleEmbeddedAttemptStream } from "./attempt-stream-settle.js"; import type { installEmbeddedAttemptStreamGuards } from "./attempt-stream.js"; import type { prepareEmbeddedAttemptTimeout } from "./attempt-timeout-prepare.js"; +import { buildPromptImageFailureNotice } from "./images.js"; import type { EmbeddedRunAttemptParams, EmbeddedRunAttemptResult } from "./types.js"; /** Runs prompt dispatch, stream settlement, cleanup, and result projection. */ @@ -55,6 +58,9 @@ type PreparedStreamRuntime = { timeout: ReturnType; }; +const FAILED_PROMPT_MEDIA_NOTE_TYPE = "openclaw.system-note"; +const FAILED_PROMPT_MEDIA_NOTE_SOURCE = "prompt-image-hydration"; + type StreamCleanupInput = { attempt: EmbeddedRunAttemptParams; clearAttemptTimeoutTimers: () => void; @@ -507,6 +513,43 @@ export async function runEmbeddedAttemptSettledPhase( compactionOccurredThisAttempt: settledStream.compactionOccurredThisAttempt, }, }); + if ( + sessionRuntimeState.currentTurnImageFailureCount > 0 && + !activeSession.messages.some( + (message) => + message.role === "custom" && + message.customType === FAILED_PROMPT_MEDIA_NOTE_TYPE && + asOptionalRecord(message.details)?.source === FAILED_PROMPT_MEDIA_NOTE_SOURCE && + asOptionalRecord(message.details)?.runId === attempt.runId, + ) + ) { + const note = { + role: "custom" as const, + customType: FAILED_PROMPT_MEDIA_NOTE_TYPE, + content: buildPromptImageFailureNotice(sessionRuntimeState.currentTurnImageFailureCount), + display: true, + details: { + source: FAILED_PROMPT_MEDIA_NOTE_SOURCE, + runId: attempt.runId, + failedMediaCount: sessionRuntimeState.currentTurnImageFailureCount, + }, + timestamp: Date.now(), + }; + await input.sessionLock.withOwnedTranscriptWrite(() => { + const target = sessionManager.getSessionTarget(); + if (target) { + SessionManager.appendMessageToTranscript( + target, + note, + attempt.config ? { config: attempt.config } : undefined, + ); + } else { + sessionManager.appendMessage(note); + } + activeSession.agent.state.messages = [...activeSession.messages, note]; + }); + messagesSnapshot = [...messagesSnapshot, note]; + } sessionIdUsed = afterTurn.sessionIdUsed; sessionFileUsed = afterTurn.sessionFileUsed; } finally { diff --git a/src/agents/embedded-agent-runner/run/attempt-setup.ts b/src/agents/embedded-agent-runner/run/attempt-setup.ts index fb95df895270..bb57017e80c5 100644 --- a/src/agents/embedded-agent-runner/run/attempt-setup.ts +++ b/src/agents/embedded-agent-runner/run/attempt-setup.ts @@ -286,6 +286,7 @@ export function installEmbeddedAttemptContextGuards(input: { getPromptCache: () => EmbeddedRunAttemptResult["promptCache"]; getPromptCacheRetention: () => PromptCacheRetention; getSystemPrompt: () => string; + onCurrentTurnImageFailure?: (count: number) => void; isOpenAIResponsesApi: boolean; repairToolUseResultPairing: boolean; sessionAgentId: string; @@ -451,6 +452,7 @@ export function installEmbeddedAttemptContextGuards(input: { input.sandbox?.enabled && input.sandbox.fsBridge ? { root: input.sandbox.workspaceDir, bridge: input.sandbox.fsBridge } : undefined, + onCurrentTurnImageFailure: input.onCurrentTurnImageFailure, }, ); const previousComputerFrameTransform = activeSession.agent.transformContext; diff --git a/src/agents/embedded-agent-runner/run/attempt-stream-settle.test.ts b/src/agents/embedded-agent-runner/run/attempt-stream-settle.test.ts index c2358b4ac687..367ac79d6533 100644 --- a/src/agents/embedded-agent-runner/run/attempt-stream-settle.test.ts +++ b/src/agents/embedded-agent-runner/run/attempt-stream-settle.test.ts @@ -261,4 +261,75 @@ describe("prepareEmbeddedAttemptTransport", () => { await fs.rm(workspaceDir, { recursive: true, force: true }); } }); + + it("records image hydration failures at the provider handoff", async () => { + let providerOptions: ProviderStreamOptions | undefined; + const providerStream = vi.fn((_model, _context, options) => { + providerOptions = options as ProviderStreamOptions; + return {} as never; + }); + bindStreamLlmRuntime(providerStream, { + streamSimple: providerStream, + registry: { getApiProvider: () => undefined }, + } as never); + const session = { + agent: { streamFn: providerStream, transport: "auto" }, + }; + const model = { api: "test-api", provider: "test-provider", id: "test-model-image" }; + const onCurrentTurnImageFailure = vi.fn(); + registerProviderStreamForModel.mockReturnValue(providerStream); + await prepareEmbeddedAttemptTransport({ + attempt: { + config: {}, + model, + modelId: model.id, + provider: model.provider, + runId: "run-native-image-failure", + runtimePlan: { + auth: { forwardedAuthProfileId: undefined }, + transport: { resolveExtraParams: () => ({}) }, + }, + sessionId: "session-native-image-failure", + }, + session, + settingsManager: { + getGlobalSettings: () => ({}), + getProjectSettings: () => ({}), + }, + onCurrentTurnImageFailure, + sessionAgentId: "main", + workspaceDir: "/tmp", + workspaceOnly: false, + agentDir: "/tmp", + abortSignal: new AbortController().signal, + getProviderRuntimeHandle: () => ({ provider: model.provider, modelId: model.id }), + sandboxSessionKey: "agent:main:test", + codeModeControlsEnabled: false, + providerPromptState: { state: {}, effectiveContextTokenBudget: 128_000 }, + } as unknown as PrepareTransportInput); + const message = attachRuntimePromptMediaFacts( + castAgentMessage({ + role: "user", + content: [ + { type: "text", text: "inspect" }, + { type: "image", data: "%%%", mimeType: "image/png" }, + ], + }), + [{ kind: "image" }], + ["inline"], + ); + const context = { systemPrompt: "system", messages: [message], tools: [] }; + + session.agent.streamFn(model as never, context as never, {}); + const provider = await resolveProviderContext(context as never, providerOptions); + + expect(onCurrentTurnImageFailure).toHaveBeenCalledWith(1); + expect(provider.messages[0]?.content).toEqual([ + { type: "text", text: "inspect" }, + { + type: "text", + text: expect.stringMatching(/1.*image contents.*unavailable.*resend.*not claim/is), + }, + ]); + }); }); diff --git a/src/agents/embedded-agent-runner/run/attempt-stream-settle.ts b/src/agents/embedded-agent-runner/run/attempt-stream-settle.ts index eeea3182699d..76d996721a0d 100644 --- a/src/agents/embedded-agent-runner/run/attempt-stream-settle.ts +++ b/src/agents/embedded-agent-runner/run/attempt-stream-settle.ts @@ -452,6 +452,7 @@ export async function prepareEmbeddedAttemptTransport(input: { session: AgentSession; settingsManager: SettingsManager; providerThinkingLevel: ProviderThinkLevel | undefined; + onCurrentTurnImageFailure?: (count: number) => void; sessionAgentId: string; workspaceDir: string; workspaceOnly: boolean; @@ -523,6 +524,7 @@ export async function prepareEmbeddedAttemptTransport(input: { localRoots: input.workspaceOnly ? undefined : getAgentScopedMediaLocalRoots(attempt.config ?? {}, input.sessionAgentId), + onCurrentTurnImageFailure: input.onCurrentTurnImageFailure, sandbox: input.sandbox?.enabled && input.sandbox.fsBridge ? { root: input.sandbox.workspaceDir, bridge: input.sandbox.fsBridge } diff --git a/src/agents/embedded-agent-runner/run/history-image-prune.test.ts b/src/agents/embedded-agent-runner/run/history-image-prune.test.ts index f9e914206f19..2784c0a12cf2 100644 --- a/src/agents/embedded-agent-runner/run/history-image-prune.test.ts +++ b/src/agents/embedded-agent-runner/run/history-image-prune.test.ts @@ -5,7 +5,7 @@ import os from "node:os"; import path from "node:path"; import type { AgentMessage } from "openclaw/plugin-sdk/agent-core"; import type { ImageContent } from "openclaw/plugin-sdk/llm"; -import { describe, expect, it } from "vitest"; +import { describe, expect, it, vi } from "vitest"; import { attachRuntimePromptMediaFacts, readRuntimePromptMediaFacts, @@ -706,6 +706,43 @@ describe("installHistoryImagePruneContextTransform", () => { } }); + it("surfaces a failed active image through the installed context transform", async () => { + const onCurrentTurnImageFailure = vi.fn(); + const message = castAgentMessage({ + role: "user", + content: [ + { type: "text", text: "inspect" }, + { type: "image", data: "%%%", mimeType: "image/png" }, + ], + __openclaw: { + media: [{ kind: "image" }], + mediaImageLayout: { slots: [{ kind: "inline", factIndex: 0 }] }, + }, + }); + const agent: { + transformContext?: (messages: AgentMessage[]) => Promise | AgentMessage[]; + } = {}; + const restore = installHistoryImagePruneContextTransform(agent, { + workspaceDir: "/tmp", + model: { input: ["text", "image"] }, + onCurrentTurnImageFailure, + }); + + try { + const replay = await agent.transformContext?.([message]); + expect(onCurrentTurnImageFailure).toHaveBeenCalledWith(1); + expect(expectArrayMessageContent(replay?.[0], "expected failure notice")).toEqual([ + { type: "text", text: "inspect" }, + { + type: "text", + text: expect.stringMatching(/1.*image contents.*unavailable.*resend.*not claim/is), + }, + ]); + } finally { + restore(); + } + }); + it("strips nested media metadata before old turns can rehydrate", async () => { const workspaceDir = await fs.mkdtemp(path.join(os.tmpdir(), "openclaw-pruned-nested-media-")); const imagePath = path.join(workspaceDir, "old.png"); diff --git a/src/agents/embedded-agent-runner/run/images.ts b/src/agents/embedded-agent-runner/run/images.ts index 6b807cb25a13..d559f6dbafde 100644 --- a/src/agents/embedded-agent-runner/run/images.ts +++ b/src/agents/embedded-agent-runner/run/images.ts @@ -9,6 +9,7 @@ import type { import { formatErrorMessage } from "../../../infra/errors.js"; import { assertNoWindowsNetworkPath, safeFileURLToPath } from "../../../infra/local-file-access.js"; import type { Context, ImageContent, TextContent } from "../../../llm/types.js"; +import { redactSensitiveText } from "../../../logging/redact.js"; import { attachRuntimePromptMediaFacts, isImageMediaFact, @@ -220,6 +221,7 @@ async function loadMediaFromRef( }, ): Promise { options?.signal?.throwIfAborted(); + const redactedRef = redactSensitiveText(ref.raw || ref.resolved); try { let targetPath = ref.resolved; @@ -240,8 +242,8 @@ async function loadMediaFromRef( }); targetPath = resolved.resolved; } catch (err) { - log.debug( - `${options?.label ?? "Native media"}: sandbox validation failed: ${formatErrorMessage(err)}`, + log.warn( + `${options?.label ?? "Native media"}: sandbox validation failed for ${redactedRef}: ${redactSensitiveText(formatErrorMessage(err))}`, ); return null; } @@ -266,7 +268,9 @@ async function loadMediaFromRef( return media; } catch (err) { options?.signal?.throwIfAborted(); - log.debug(`${options?.label ?? "Native media"}: failed to load: ${formatErrorMessage(err)}`); + log.warn( + `${options?.label ?? "Native media"}: failed to load ${redactedRef}: ${redactSensitiveText(formatErrorMessage(err))}`, + ); return null; } } @@ -488,8 +492,13 @@ type PromptMediaOptions = { sandbox?: { root: string; bridge: SandboxFsBridge }; provider?: boolean; signal?: AbortSignal; + onCurrentTurnImageFailure?: (count: number) => void; }; +export function buildPromptImageFailureNotice(count: number): string { + return `System note: ${count} image attachment${count === 1 ? "" : "s"} could not be loaded; their image contents are unavailable. Tell the user and ask them to resend ${count === 1 ? "the image" : "the images"}; do not claim inspection.`; +} + const VIDEO_OMISSION = { unsupported: "(video omitted: provider does not support native video)", unavailable: "(video omitted: source unavailable)", @@ -575,6 +584,7 @@ async function materializePromptMediaMessages( ): Promise { let hydrated: AgentMessage[] | undefined; const videoBudget = { remaining: MAX_VIDEO_BYTES }; + const activeUserIndex = messages.findLastIndex((message) => message.role === "user"); for (const [index, message] of messages.entries()) { if (message.role !== "user") { continue; @@ -613,6 +623,17 @@ async function materializePromptMediaMessages( options, budget: videoBudget, }); + if ( + (options.provider || options.onCurrentTurnImageFailure) && + index === activeUserIndex && + result.failedMediaCount > 0 + ) { + options.onCurrentTurnImageFailure?.(result.failedMediaCount); + projectedContent.push({ + type: "text", + text: buildPromptImageFailureNotice(result.failedMediaCount), + }); + } hydrated ??= messages.slice(); if (options.provider) { hydrated[index] = { @@ -665,6 +686,7 @@ export async function materializeProviderContext(params: { workspaceOnly?: boolean; localRoots?: readonly string[]; sandbox?: { root: string; bridge: SandboxFsBridge }; + onCurrentTurnImageFailure?: (count: number) => void; }): Promise { const messages = await materializePromptMediaMessages(params.context.messages as AgentMessage[], { workspaceDir: params.workspaceDir, @@ -674,6 +696,7 @@ export async function materializeProviderContext(params: { sandbox: params.sandbox, provider: true, signal: params.signal, + onCurrentTurnImageFailure: params.onCurrentTurnImageFailure, }); params.signal?.throwIfAborted(); return messages === params.context.messages diff --git a/src/media/fetch.test.ts b/src/media/fetch.test.ts index 84a4517a4d53..55042b00b2cf 100644 --- a/src/media/fetch.test.ts +++ b/src/media/fetch.test.ts @@ -1,7 +1,9 @@ // Media fetch tests cover remote media download limits and validation. import fs from "node:fs/promises"; +import path from "node:path"; import { MAX_TIMER_TIMEOUT_MS } from "@openclaw/normalization-core/number-coercion"; import { afterAll, beforeAll, beforeEach, describe, expect, it, vi } from "vitest"; +import { hasErrnoCode } from "../infra/errors.js"; import { createTempHomeEnv, type TempHomeEnv } from "../test-utils/temp-home.js"; const fetchWithSsrFGuardMock = vi.hoisted(() => vi.fn()); @@ -1396,15 +1398,34 @@ describe("readRemoteMediaBuffer", () => { }, ); - it("saves bodyless successful responses without unbounded buffering", async () => { - const saved = await saveResponseMedia(new Response(null, { status: 204 }), { - sourceUrl: "https://example.com/empty", - fallbackContentType: "application/octet-stream", - maxBytes: 8, - }); + it("rejects bodyless successful responses without saving an empty file", async () => { + const inboundDir = path.join(tempHome.home, ".openclaw", "media", "inbound"); + const listInboundFiles = async () => { + try { + return (await fs.readdir(inboundDir)).toSorted(); + } catch (error) { + if (hasErrnoCode(error, "ENOENT")) { + return []; + } + throw error; + } + }; + const before = await listInboundFiles(); - expect(saved.size).toBe(0); - await expect(fs.readFile(saved.path)).resolves.toStrictEqual(Buffer.alloc(0)); + await expect( + saveResponseMedia(new Response(null, { status: 204 }), { + sourceUrl: "https://example.com/empty", + fallbackContentType: "application/octet-stream", + maxBytes: 8, + }), + ).rejects.toMatchObject({ + name: "MediaFetchError", + code: "http_error", + status: 204, + message: + "Failed to fetch media from https://example.com/empty: HTTP 204; empty response body", + }); + await expect(listInboundFiles()).resolves.toEqual(before); }); it("uses caller filename hints for MIME detection without preserving storage basenames", async () => { diff --git a/src/media/fetch.ts b/src/media/fetch.ts index 6ac9c1db3a69..4376c6442820 100644 --- a/src/media/fetch.ts +++ b/src/media/fetch.ts @@ -21,7 +21,7 @@ import { retryAsync, type RetryOptions } from "../infra/retry.js"; import { isTransientNetworkError } from "../infra/retryable-network-errors.js"; import { redactSensitiveText } from "../logging/redact.js"; import { buildTimeoutAbortSignal } from "../utils/fetch-timeout.js"; -import { saveMediaBuffer, saveMediaStream, type SavedMedia } from "./store.js"; +import { saveMediaStream, type SavedMedia } from "./store.js"; /** Default remote media fetch cap shared by buffer reads and store writes. */ const DEFAULT_FETCH_MEDIA_MAX_BYTES = MAX_DOCUMENT_BYTES; @@ -393,7 +393,7 @@ async function assertMediaResponseOk(params: { readIdleTimeoutMs?: number; }): Promise { const { res, url, finalUrl, sourceUrl, readIdleTimeoutMs } = params; - if (res.ok) { + if (res.ok && res.body) { return; } const statusText = res.statusText ? ` ${res.statusText}` : ""; @@ -573,23 +573,15 @@ async function saveOkMediaResponse(params: { ? (params.filePathHint ?? fileName) : undefined; try { - const saved = params.res.body - ? await saveMediaStream( - responseBodyChunks(params.res.body, params.readIdleTimeoutMs), - contentType ?? undefined, - params.subdir ?? "inbound", - params.maxBytes, - params.originalFilename, - detectionFilePathHint, - ) - : await saveMediaBuffer( - Buffer.alloc(0), - contentType ?? undefined, - params.subdir ?? "inbound", - params.maxBytes, - params.originalFilename, - detectionFilePathHint, - ); + const body = expectDefined(params.res.body, "media response body"); + const saved = await saveMediaStream( + responseBodyChunks(body, params.readIdleTimeoutMs), + contentType ?? undefined, + params.subdir ?? "inbound", + params.maxBytes, + params.originalFilename, + detectionFilePathHint, + ); return { ...saved, ...(fileName ? { fileName } : {}) }; } catch (err) { if (err instanceof MediaFetchError) { diff --git a/src/media/playback-transcode.test.ts b/src/media/playback-transcode.test.ts index 0cdbba205cd0..ad3ecf9ef26c 100644 --- a/src/media/playback-transcode.test.ts +++ b/src/media/playback-transcode.test.ts @@ -9,11 +9,16 @@ import { waitForPlaybackTranscodeJobsForTest, } from "./playback-transcode.test-support.js"; -const { probePlaybackMediaFileDescriptor, runFfmpeg } = vi.hoisted(() => ({ +const { playbackWarn, probePlaybackMediaFileDescriptor, runFfmpeg } = vi.hoisted(() => ({ + playbackWarn: vi.fn(), probePlaybackMediaFileDescriptor: vi.fn(), runFfmpeg: vi.fn(), })); +vi.mock("../logging/subsystem.js", () => ({ + createSubsystemLogger: (subsystem: string) => ({ subsystem, warn: playbackWarn }), +})); + vi.mock("./ffmpeg-exec.js", () => ({ runFfmpeg, })); @@ -42,6 +47,7 @@ afterAll(async () => { }); beforeEach(() => { + playbackWarn.mockReset(); runFfmpeg.mockReset(); probePlaybackMediaFileDescriptor.mockReset(); probePlaybackMediaFileDescriptor.mockImplementation(async (_fd: number, kind: string) => @@ -759,6 +765,41 @@ describe("resolvePlaybackTranscode", () => { } }); + it("warns once when the same transcode operation fails across cooldown retries", async () => { + const source = await createSource("warn-failed.caf", "caff-source"); + const nowSpy = vi.spyOn(Date, "now").mockReturnValue(1000); + runFfmpeg.mockRejectedValue(new Error("ffmpeg unavailable")); + const params = { + ...source, + mimeType: "audio/x-caf", + kind: "audio" as const, + }; + + try { + await expect(playback.resolvePlaybackTranscode(params)).resolves.toEqual({ + kind: "preparing", + }); + await vi.waitFor(() => expect(playbackWarn).toHaveBeenCalledOnce()); + expect(playbackWarn).toHaveBeenCalledWith( + expect.stringContaining(`${source.sourcePath}: ffmpeg unavailable`), + ); + + nowSpy.mockReturnValue(61_001); + await expect(playback.resolvePlaybackTranscode(params)).resolves.toEqual({ + kind: "preparing", + }); + await vi.waitFor(() => expect(runFfmpeg).toHaveBeenCalledTimes(2)); + await vi.waitFor(async () => { + await expect(playback.resolvePlaybackTranscode(params)).resolves.toEqual({ + kind: "fallback", + }); + }); + expect(playbackWarn).toHaveBeenCalledOnce(); + } finally { + nowSpy.mockRestore(); + } + }); + it("keeps a saturated transcode pool retryable until capacity is available", async () => { const sources = await Promise.all([ createSource("pool-first.mkv"), diff --git a/src/media/playback-transcode.ts b/src/media/playback-transcode.ts index 2b35ff159ba0..6d43e05f3471 100644 --- a/src/media/playback-transcode.ts +++ b/src/media/playback-transcode.ts @@ -4,11 +4,13 @@ import fs, { type FileHandle } from "node:fs/promises"; import path from "node:path"; import { maxBytesForKind, type MediaKind } from "@openclaw/media-core/constants"; import { extensionForMime, normalizeMimeType } from "@openclaw/media-core/mime"; +import { formatErrorMessage } from "../infra/errors.js"; import { fileStore } from "../infra/file-store.js"; import { openLocalFileSafely } from "../infra/fs-safe.js"; import { pruneMapToMaxSize } from "../infra/map-size.js"; import { withTempWorkspace } from "../infra/private-temp-workspace.js"; import { resolvePreferredOpenClawTmpDir } from "../infra/tmp-openclaw-dir.js"; +import { createSubsystemLogger } from "../logging/subsystem.js"; import { getOrCreatePromise } from "../shared/lazy-promise.js"; import { runFfmpeg } from "./ffmpeg-exec.js"; import { probePlaybackMediaFileDescriptor, type PlaybackMediaProbeResult } from "./media-probe.js"; @@ -139,6 +141,7 @@ const playbackJobs = new Map>(); const playbackFailures = new Map(); const playbackInspections = new Map(); const playbackInspectionJobs = new Map>(); +const log = createSubsystemLogger("media/playback"); /** Hashes the immutable source identity used by playback cache file names. */ function createPlaybackTranscodeCacheKey(source: PlaybackSourceIdentity): string { @@ -662,7 +665,6 @@ export async function resolvePlaybackTranscode( if (Date.now() - failedAtMs < PLAYBACK_TRANSCODE_FAILURE_COOLDOWN_MS) { return { kind: "fallback" }; } - playbackFailures.delete(operationKey); } if (playbackJobs.size >= MAX_PLAYBACK_TRANSCODE_JOBS) { return { kind: "preparing" }; @@ -689,8 +691,13 @@ export async function resolvePlaybackTranscode( playbackJobs.delete(operationKey); playbackFailures.delete(operationKey); }, - () => { + (reason: unknown) => { playbackJobs.delete(operationKey); + if (!playbackFailures.has(operationKey)) { + log.warn( + `Playback transcode failed for ${params.sourcePath}: ${formatErrorMessage(reason)}`, + ); + } playbackFailures.delete(operationKey); playbackFailures.set(operationKey, Date.now()); pruneMapToMaxSize(playbackFailures, MAX_PLAYBACK_ENTRIES.failures);