diff --git a/docs/concepts/model-failover.md b/docs/concepts/model-failover.md index 8f5c516af4f9..eb4439c39597 100644 --- a/docs/concepts/model-failover.md +++ b/docs/concepts/model-failover.md @@ -50,6 +50,12 @@ policy. OpenClaw does not retry them with thinking disabled. Fallback execution is turn-local. The reply runner persists only fallback notice state so `/status` and transition notices can distinguish the selected model from the model that answered. It does not persist the fallback as the next turn's model selection. +Sessions placed on an OpenClaw cloud worker keep the OpenClaw runtime when +automatic model selection advances to a fallback. The configured provider, +model, and auth-profile fallback rules still apply. Explicit session or +configured runtime choices remain strict: an incompatible runtime reports a +placement error and requires a compatible destination before retrying. + When configured fallback stops because the agent run reaches a final timeout or the idle-timeout cost-runaway breaker returns a terminal error, the `model-fallback/decision` logger records `model_fallback_chain_stopped` with diff --git a/docs/logging.md b/docs/logging.md index cb6b84ef1df4..73ec0134f378 100644 --- a/docs/logging.md +++ b/docs/logging.md @@ -237,6 +237,18 @@ Chat displays recognized request-limit facts, including the allowed and actual number of `cache_control` blocks, in both live failures and saved history. Raw proxy metadata stays in redacted diagnostics rather than the chat message. +Saved failed replies also distinguish rate limits, authentication failures, +provider HTTP errors, and network interruptions. Worker inference preserves +bounded, redacted error details for classification, including when a large +partial response cannot fit in the transcript. Unrecognized errors still use +generic chat copy; inspect the Gateway logs and stored error for diagnosis. + +A worker message-size failure is separate from a model context-window limit. +Retry with a smaller response or continue on the Gateway. If the worker cannot +preserve the model's continuation data, stop or reclaim it before retrying on +the Gateway. Earlier tool actions may already have completed, so check their +results before repeating them. + ### Targeted model transport diagnostics When debugging provider calls, use targeted environment flags instead of raising diff --git a/packages/ai/src/internal/shared.ts b/packages/ai/src/internal/shared.ts index 9a88ab9e5218..a6cb3820965e 100644 --- a/packages/ai/src/internal/shared.ts +++ b/packages/ai/src/internal/shared.ts @@ -4,6 +4,7 @@ export * from "../providers/tool-result-text.js"; export * from "../providers/transform-messages.js"; export * from "../replay-turn-classification.js"; export { createDiagnosticRecord } from "../utils/credential-redaction.js"; +export { projectProviderError } from "../utils/provider-error.js"; export * from "../utils/prompt-cache-stability.js"; export * from "../utils/sanitize-unicode.js"; export * from "../utils/system-prompt-cache-boundary.js"; diff --git a/packages/ai/src/providers/mistral.test.ts b/packages/ai/src/providers/mistral.test.ts index df7fb5dce5c8..bcff042ba0d0 100644 --- a/packages/ai/src/providers/mistral.test.ts +++ b/packages/ai/src/providers/mistral.test.ts @@ -428,7 +428,7 @@ describe("Mistral provider", () => { }, ); - it("preserves Mistral messages while keeping error bodies UTF-16 safe and bounded", async () => { + it("preserves Mistral HTTP status and message while keeping error bodies UTF-16 safe and bounded", async () => { const prefix = "a".repeat(3_999); mistralMockState.streamError = Object.assign(new Error("invalid request"), { statusCode: 400, @@ -437,7 +437,7 @@ describe("Mistral provider", () => { const result = await runMistralFixture(); - expect(result.errorMessage).toBe("invalid request"); + expect(result.errorMessage).toBe("400: invalid request"); expect(result.errorBody).toBe(`${prefix.slice(0, 500)}... [truncated]`); }); diff --git a/packages/ai/src/utils/provider-error.test.ts b/packages/ai/src/utils/provider-error.test.ts index 6f2308dd9c01..c37bc6d84b65 100644 --- a/packages/ai/src/utils/provider-error.test.ts +++ b/packages/ai/src/utils/provider-error.test.ts @@ -132,11 +132,14 @@ describe("projectProviderError", () => { expect(projectProviderError(error).errorMessage).toBe(expected); }); - it("preserves an SDK message that already contains the response body", () => { + it("preserves an SDK response body and its known HTTP status", () => { const body = '{"error":{"message":"permission denied"}}'; const error = Object.assign(new Error(body), { status: 403, body }); - expect(projectProviderError(error).errorMessage).toBe(body); + expect(projectProviderError(error)).toMatchObject({ + errorMessage: `403: ${body}`, + errorBody: body, + }); }); it("preserves a meaningful SDK message alongside its structured body", () => { @@ -208,6 +211,18 @@ describe("projectProviderError", () => { expect(projection.errorBody?.length).toBeLessThanOrEqual(515); }); + it("preserves known HTTP status with a short message and a response body", () => { + const error = Object.assign(new Error("Provider request failed"), { + status: 429, + body: { error: { message: "Try later" } }, + }); + + expect(projectProviderError(error)).toMatchObject({ + errorMessage: "429: Provider request failed", + errorCode: "429", + }); + }); + it("bounds repeated structured diagnostic fragments before extraction", () => { expect(projectProviderError("{}".repeat(8193)).errorMessage).toBe( "[Oversized diagnostic JSON redacted]", diff --git a/packages/ai/src/utils/provider-error.ts b/packages/ai/src/utils/provider-error.ts index 073140e0100b..02c6660eaa36 100644 --- a/packages/ai/src/utils/provider-error.ts +++ b/packages/ai/src/utils/provider-error.ts @@ -88,12 +88,7 @@ function buildProjection(snapshot: unknown, signal?: AbortSignal): ProviderError body ?? stringifyField(snapshot, MAX_ERROR_BODY_LENGTH) ?? "Unknown provider error"); - if ( - status !== undefined && - !body && - originalMessage && - !originalMessage.startsWith(String(status)) - ) { + if (status !== undefined && originalMessage && !originalMessage.startsWith(String(status))) { errorMessage = `${status}: ${errorMessage}`; } const metadata = asOptionalRecord(nestedError?.metadata); diff --git a/src/agents/command/run-embedded-attempt.ts b/src/agents/command/run-embedded-attempt.ts index 465d7a291496..f45466867dfe 100644 --- a/src/agents/command/run-embedded-attempt.ts +++ b/src/agents/command/run-embedded-attempt.ts @@ -384,19 +384,17 @@ export async function runEmbeddedAgentAttempt(params: RunEmbeddedAgentAttemptPar providerOverride === defaultProvider && modelOverride === defaultModel ? configuredDefaultAuthProfileId : undefined; - const agentHarnessRuntimeOverride = resolveSessionRuntimeOverrideForProvider({ - provider: providerOverride, - entry: attemptSessionEntry, - cfg, - }); - const candidateRuntime = resolveEffectiveAgentRuntime({ - cfg, - provider: providerOverride, - modelId: modelOverride, - agentId: sessionAgentId, - sessionKey, - sessionEntry: attemptSessionEntry, - }); + const agentHarnessRuntimeOverride = runOptions.agentHarnessRuntimeOverride; + const candidateRuntime = + agentHarnessRuntimeOverride ?? + resolveEffectiveAgentRuntime({ + cfg, + provider: providerOverride, + modelId: modelOverride, + agentId: sessionAgentId, + sessionKey, + sessionEntry: attemptSessionEntry, + }); const candidateConfiguredThinkLevel = immutableThinkLevel ?? resolveConfiguredThinkingDefault({ diff --git a/src/agents/embedded-agent-runner/run-entry.failures.test-support.ts b/src/agents/embedded-agent-runner/run-entry.failures.test-support.ts index 5d0345d5e5bf..faede59e3525 100644 --- a/src/agents/embedded-agent-runner/run-entry.failures.test-support.ts +++ b/src/agents/embedded-agent-runner/run-entry.failures.test-support.ts @@ -1,5 +1,5 @@ import { expect, it, vi, type Mock } from "vitest"; -import { runEmbeddedAgentEntry } from "./run-entry.js"; +import { runEmbeddedAgentEntry } from "./run-entry.test-harness.js"; import { createDirectHarness, makeResult, diff --git a/src/agents/embedded-agent-runner/run-entry.test-harness.ts b/src/agents/embedded-agent-runner/run-entry.test-harness.ts new file mode 100644 index 000000000000..5a1c325a1756 --- /dev/null +++ b/src/agents/embedded-agent-runner/run-entry.test-harness.ts @@ -0,0 +1,126 @@ +import { beforeEach, expect, vi } from "vitest"; +import type { ContextEngineTurnAttemptFacts } from "../harness/context-engine-turn-attempt.js"; +import { + initialAttemptOptions, + fallbackAttemptOptions, + type FallbackRunnerParams, +} from "./run-entry.test-support.js"; + +const state = vi.hoisted(() => ({ + runWithModelFallback: vi.fn(), + ensureSelectedAgentHarnessPlugin: vi.fn(async (_params: unknown) => undefined), + selectAgentHarness: vi.fn(({ provider }: { provider: string }) => ({ + id: provider === "fallback-provider" ? "fallback-harness" : "primary-harness", + contextEngineHostCapabilities: [], + })), + discardedAttempts: [] as string[], + finalizedAttempts: [] as string[], +})); + +vi.mock("../harness/context-engine-turn-attempt.js", () => ({ + discardContextEngineTurnAttemptIntent: vi.fn( + ({ facts }: { facts: ContextEngineTurnAttemptFacts }) => { + state.discardedAttempts.push(facts.sessionIdUsed); + }, + ), + finalizeAcceptedContextEngineTurn: vi.fn(async ({ facts }) => { + state.finalizedAttempts.push(facts.sessionIdUsed); + }), +})); + +vi.mock("../model-fallback-runner.js", () => ({ + runWithModelFallback: (params: FallbackRunnerParams) => state.runWithModelFallback(params), +})); + +vi.mock("../harness/runtime-plugin.js", () => ({ + ensureSelectedAgentHarnessPlugin: (params: unknown) => + state.ensureSelectedAgentHarnessPlugin(params), +})); + +vi.mock("../harness/selection.js", () => ({ + selectAgentHarness: (params: { provider: string }) => state.selectAgentHarness(params), +})); + +export const { runEmbeddedAgentEntry } = await import("./run-entry.js"); + +export function setupRunEntryTestState() { + beforeEach(() => { + state.discardedAttempts.length = 0; + state.finalizedAttempts.length = 0; + state.ensureSelectedAgentHarnessPlugin.mockReset().mockResolvedValue(undefined); + state.selectAgentHarness + .mockReset() + .mockImplementation(({ provider }: { provider: string }) => ({ + id: provider === "fallback-provider" ? "fallback-harness" : "primary-harness", + contextEngineHostCapabilities: [], + })); + state.runWithModelFallback + .mockReset() + .mockImplementation(async (params: FallbackRunnerParams) => { + await params.prepareCandidateChain?.([ + { + provider: params.provider, + model: params.model, + routeOrigin: "requested", + routeResolution: "raw", + }, + { + provider: "fallback-provider", + model: "fallback-model", + routeOrigin: "configured-fallback", + routeResolution: "raw", + }, + ]); + await params.prepareAgentHarnessRuntime?.({ + provider: params.provider, + model: params.model, + agentHarnessRuntimeOverride: params.resolveAgentHarnessRuntimeOverride?.( + params.provider, + params.model, + ), + }); + const primaryResult = await params.run(params.provider, params.model, { + ...initialAttemptOptions(params), + allowTransientCooldownProbe: true, + }); + const classification = await params.classifyResult?.({ + result: primaryResult, + provider: params.provider, + model: params.model, + attempt: 1, + total: 2, + }); + expect(classification).toBeTruthy(); + const fallbackProvider = "fallback-provider"; + const fallbackModel = "fallback-model"; + await params.prepareAgentHarnessRuntime?.({ + provider: fallbackProvider, + model: fallbackModel, + agentHarnessRuntimeOverride: params.resolveAgentHarnessRuntimeOverride?.( + fallbackProvider, + fallbackModel, + ), + }); + const result = await params.run(fallbackProvider, fallbackModel, { + ...fallbackAttemptOptions(params, "format"), + isFinalFallbackAttempt: true, + }); + return { + outcome: "completed" as const, + result, + provider: fallbackProvider, + model: fallbackModel, + attempts: [ + { + provider: params.provider, + model: params.model, + error: "empty result", + reason: "format" as const, + }, + ], + }; + }); + }); + + return state; +} diff --git a/src/agents/embedded-agent-runner/run-entry.test.ts b/src/agents/embedded-agent-runner/run-entry.test.ts index 6791203eda5d..9c408c5e0b3e 100644 --- a/src/agents/embedded-agent-runner/run-entry.test.ts +++ b/src/agents/embedded-agent-runner/run-entry.test.ts @@ -1,4 +1,4 @@ -import { beforeEach, describe, expect, it, vi } from "vitest"; +import { describe, expect, it, vi } from "vitest"; import type { OpenClawConfig } from "../../config/types.openclaw.js"; import { clearAgentRunContext, @@ -7,134 +7,21 @@ import { resolveProjectedAgentRunModel, registerAgentRunContext, } from "../../infra/agent-run-registry.js"; -import type { ContextEngineTurnAttemptFacts } from "../harness/context-engine-turn-attempt.js"; import { registerRunEntryFailureTests } from "./run-entry.failures.test-support.js"; -import { runEmbeddedAgentEntry } from "./run-entry.js"; +import { runEmbeddedAgentEntry, setupRunEntryTestState } from "./run-entry.test-harness.js"; import { createDirectHarness, makeResult, recordTurnAttempt, initialAttemptOptions, - fallbackAttemptOptions, type FallbackRunnerParams, } from "./run-entry.test-support.js"; -const state = vi.hoisted(() => ({ - runWithModelFallback: vi.fn(), - ensureSelectedAgentHarnessPlugin: vi.fn(async (_params: unknown) => undefined), - selectAgentHarness: vi.fn(({ provider }: { provider: string }) => ({ - id: provider === "fallback-provider" ? "fallback-harness" : "primary-harness", - contextEngineHostCapabilities: [], - })), - discardedAttempts: [] as string[], - finalizedAttempts: [] as string[], -})); - -vi.mock("../harness/context-engine-turn-attempt.js", () => ({ - discardContextEngineTurnAttemptIntent: vi.fn( - ({ facts }: { facts: ContextEngineTurnAttemptFacts }) => { - state.discardedAttempts.push(facts.sessionIdUsed); - }, - ), - finalizeAcceptedContextEngineTurn: vi.fn(async ({ facts }) => { - state.finalizedAttempts.push(facts.sessionIdUsed); - }), -})); - -vi.mock("../model-fallback-runner.js", () => ({ - runWithModelFallback: (params: FallbackRunnerParams) => state.runWithModelFallback(params), -})); - -vi.mock("../harness/runtime-plugin.js", () => ({ - ensureSelectedAgentHarnessPlugin: (params: unknown) => - state.ensureSelectedAgentHarnessPlugin(params), -})); - -vi.mock("../harness/selection.js", () => ({ - selectAgentHarness: (params: { provider: string }) => state.selectAgentHarness(params), -})); +const state = setupRunEntryTestState(); describe("runEmbeddedAgentEntry", () => { registerRunEntryFailureTests(state); - beforeEach(() => { - state.discardedAttempts.length = 0; - state.finalizedAttempts.length = 0; - state.ensureSelectedAgentHarnessPlugin.mockReset().mockResolvedValue(undefined); - state.selectAgentHarness - .mockReset() - .mockImplementation(({ provider }: { provider: string }) => ({ - id: provider === "fallback-provider" ? "fallback-harness" : "primary-harness", - contextEngineHostCapabilities: [], - })); - state.runWithModelFallback - .mockReset() - .mockImplementation(async (params: FallbackRunnerParams) => { - await params.prepareCandidateChain?.([ - { - provider: params.provider, - model: params.model, - routeOrigin: "requested", - routeResolution: "raw", - }, - { - provider: "fallback-provider", - model: "fallback-model", - routeOrigin: "configured-fallback", - routeResolution: "raw", - }, - ]); - await params.prepareAgentHarnessRuntime?.({ - provider: params.provider, - model: params.model, - agentHarnessRuntimeOverride: params.resolveAgentHarnessRuntimeOverride?.( - params.provider, - params.model, - ), - }); - const primaryResult = await params.run(params.provider, params.model, { - ...initialAttemptOptions(params), - allowTransientCooldownProbe: true, - }); - const classification = await params.classifyResult?.({ - result: primaryResult, - provider: params.provider, - model: params.model, - attempt: 1, - total: 2, - }); - expect(classification).toBeTruthy(); - const fallbackProvider = "fallback-provider"; - const fallbackModel = "fallback-model"; - await params.prepareAgentHarnessRuntime?.({ - provider: fallbackProvider, - model: fallbackModel, - agentHarnessRuntimeOverride: params.resolveAgentHarnessRuntimeOverride?.( - fallbackProvider, - fallbackModel, - ), - }); - const result = await params.run(fallbackProvider, fallbackModel, { - ...fallbackAttemptOptions(params, "format"), - isFinalFallbackAttempt: true, - }); - return { - outcome: "completed" as const, - result, - provider: fallbackProvider, - model: fallbackModel, - attempts: [ - { - provider: params.provider, - model: params.model, - error: "empty result", - reason: "format" as const, - }, - ], - }; - }); - }); - it("keeps shared fallback and terminal behavior aligned across entry modes", async ({ onTestFinished, }) => { @@ -337,8 +224,16 @@ describe("runEmbeddedAgentEntry", () => { }), }); - expect(resolveContextEngineHost).toHaveBeenCalledWith("primary-provider", "primary-model"); - expect(resolveContextEngineHost).toHaveBeenCalledWith("fallback-provider", "fallback-model"); + expect(resolveContextEngineHost).toHaveBeenCalledWith( + "primary-provider", + "primary-model", + undefined, + ); + expect(resolveContextEngineHost).toHaveBeenCalledWith( + "fallback-provider", + "fallback-model", + undefined, + ); expect(state.selectAgentHarness).not.toHaveBeenCalled(); }); diff --git a/src/agents/embedded-agent-runner/run-entry.ts b/src/agents/embedded-agent-runner/run-entry.ts index b8850ff89326..a51f0284a629 100644 --- a/src/agents/embedded-agent-runner/run-entry.ts +++ b/src/agents/embedded-agent-runner/run-entry.ts @@ -23,6 +23,7 @@ import { finalizeAcceptedContextEngineTurn, type ContextEngineTurnAttemptFacts, } from "../harness/context-engine-turn-attempt.js"; +import { resolveAgentHarnessPolicy } from "../harness/policy.js"; import { ensureSelectedAgentHarnessPlugin } from "../harness/runtime-plugin.js"; import { selectAgentHarness } from "../harness/selection.js"; import type { ModelFallbackResultClassification } from "../model-fallback-attempt.js"; @@ -36,6 +37,7 @@ import type { import type { ModelManifestNormalizationContext } from "../model-ref-shared.js"; import { settleFailedRequesterRun, settleRequesterRun } from "../requester-run-settlement.js"; import { resolveAgentRunAbortLifecycleFields } from "../run-termination.js"; +import { resolveSessionPlacementRuntimeOverride } from "../session-placement-admission.js"; import { didEmbeddedCyberFailoverTargetCommitWork, EMBEDDED_CYBER_FAILOVER_TRIGGER_CODE, @@ -66,6 +68,7 @@ import type { EmbeddedAgentRunResult } from "./types.js"; export type { EmbeddedAgentRunEntryTerminal } from "./run-entry-terminal.js"; type RunEntryCandidateOptions = { + agentHarnessRuntimeOverride: string | undefined; assistantErrorTranscript: AssistantErrorTranscript; authProfileFailurePolicy?: AuthProfileFailurePolicy; classifyResult: (result: EmbeddedAgentRunResult) => ModelFallbackResultClassification; @@ -135,6 +138,7 @@ type EmbeddedAgentRunEntryParams = { resolveContextEngineHost?: ( provider: string, model: string, + agentHarnessRuntimeOverride: string | undefined, ) => ContextEngineHostSupport | undefined; }; behavior: RunEntryBehavior; @@ -172,6 +176,22 @@ async function runEmbeddedAgentEntryInternal( ): Promise> { const lifecycleGeneration = captureAgentRunLifecycleGeneration(params.identity.runId); const runContext = getAgentRunContext(params.identity.runId); + const placementRuntime = resolveSessionPlacementRuntimeOverride(params.identity); + const resolveRuntimeOverride = (provider: string, model: string) => { + const requestedRuntime = params.harness.resolveRuntimeOverride(provider, model); + if (requestedRuntime || !placementRuntime) { + return requestedRuntime; + } + const policy = resolveAgentHarnessPolicy({ + config: params.selection.cfg, + provider, + modelId: model, + agentId: params.identity.agentId, + sessionKey: params.harness.sessionKey, + }); + // Explicit runtime choices still reach placement's compatibility check. + return policy.runtimeSource === "implicit" ? placementRuntime : undefined; + }; const clearObservedModel = () => { const event = { ...params.identity, @@ -265,11 +285,11 @@ async function runEmbeddedAgentEntryInternal( ...selection, ...params.identity, abortSignal: params.abortSignal, - resolveAgentHarnessRuntimeOverride: params.harness.resolveRuntimeOverride, + resolveAgentHarnessRuntimeOverride: resolveRuntimeOverride, prepareCandidateChain: async (candidates) => { for (const candidate of candidates) { try { - const agentHarnessRuntimeOverride = params.harness.resolveRuntimeOverride( + const agentHarnessRuntimeOverride = resolveRuntimeOverride( candidate.provider, candidate.model, ); @@ -281,6 +301,7 @@ async function runEmbeddedAgentEntryInternal( const resolvedHost = params.harness.resolveContextEngineHost?.( candidate.provider, candidate.model, + agentHarnessRuntimeOverride, ); const host = resolvedHost ?? @@ -410,6 +431,7 @@ async function runEmbeddedAgentEntryInternal( }; try { const result = await params.runCandidate(provider, model, { + agentHarnessRuntimeOverride: resolveRuntimeOverride(provider, model), assistantErrorTranscript, // The original OpenAI refusal proves this turn's credential already // reached the provider. Keep a target-only entitlement rejection from diff --git a/src/agents/embedded-agent-runner/run-entry.worker-placement.test.ts b/src/agents/embedded-agent-runner/run-entry.worker-placement.test.ts new file mode 100644 index 000000000000..d16f717e6d21 --- /dev/null +++ b/src/agents/embedded-agent-runner/run-entry.worker-placement.test.ts @@ -0,0 +1,89 @@ +import { describe, expect, it, onTestFinished, vi } from "vitest"; +import type { OpenClawConfig } from "../../config/types.openclaw.js"; +import { assertSupportedTurn } from "../../gateway/worker-environments/worker-turn-payload.js"; +import { installSessionPlacementAdmissionProvider } from "../session-placement-admission.js"; +import { runEmbeddedAgentEntry, setupRunEntryTestState } from "./run-entry.test-harness.js"; +import { createDirectHarness, makeResult } from "./run-entry.test-support.js"; + +const state = setupRunEntryTestState(); + +describe("runEmbeddedAgentEntry worker placement", () => { + it.each(["implicit", "session", "configured"] as const)( + "keeps worker fallback preparation and execution on its placement runtime (%s selection)", + async (selection) => { + const requestedRuntime = selection === "session" ? "codex" : undefined; + const cfg: OpenClawConfig = + selection === "configured" + ? { + agents: { + defaults: { + models: { "primary-provider/primary-model": { agentRuntime: { id: "codex" } } }, + }, + }, + } + : {}; + const placement = { + assertCompactionSuccessorAllowed: () => {}, + resolveRuntimeOverride: vi.fn(() => "openclaw"), + executeLocalTurn: async (_claim: unknown, run: () => Promise) => run(), + executeTurn: vi.fn(), + }; + onTestFinished(installSessionPlacementAdmissionProvider(placement)); + const entry = runEmbeddedAgentEntry({ + selection: { cfg, provider: "primary-provider", model: "primary-model" }, + identity: { + runId: "run-worker-fallback", + agentId: "main", + sessionId: "session-1", + sessionKey: "agent:main:chat", + }, + harness: { ...createDirectHarness(), resolveRuntimeOverride: () => requestedRuntime }, + behavior: { kind: "command-rpc", hasCommittedSideEffect: () => false }, + sessionOverride: { kind: "preserve" }, + runCandidate: async (provider, model, options) => { + assertSupportedTurn({ + sessionId: "session-1", + sessionFile: "/tmp/worker-fallback.sqlite", + workspaceDir: "/tmp/workspace", + runId: "run-worker-fallback", + prompt: "Continue", + timeoutMs: 1_000, + provider, + model, + agentHarnessRuntimeOverride: + options.agentHarnessRuntimeOverride ?? + (selection === "implicit" + ? options.isFallbackRetry + ? "codex" + : "openclaw" + : undefined), + config: cfg, + }); + return makeResult({ + provider, + model, + ...(options.isFallbackRetry ? {} : { classification: "empty" as const }), + }); + }, + }); + if (selection !== "implicit") { + await expect(entry).rejects.toThrow( + "Cloud worker turns require the OpenClaw runtime, not codex", + ); + } else { + await expect(entry).resolves.toMatchObject({ + outcome: "completed", + model: "fallback-model", + }); + expect(placement.resolveRuntimeOverride).toHaveBeenCalledOnce(); + } + for (const [prepared] of state.ensureSelectedAgentHarnessPlugin.mock.calls) { + if (selection !== "configured") { + expect(prepared).toMatchObject({ + agentHarnessRuntimeOverride: requestedRuntime ?? "openclaw", + }); + } + } + }, + ); +}); diff --git a/src/agents/failover/assistant-request-failure-copy.ts b/src/agents/failover/assistant-request-failure-copy.ts index 11c47948d3dd..7e7b680b3811 100644 --- a/src/agents/failover/assistant-request-failure-copy.ts +++ b/src/agents/failover/assistant-request-failure-copy.ts @@ -1,6 +1,17 @@ +import { normalizeLowercaseStringOrEmpty } from "@openclaw/normalization-core/string-coerce"; import type { GatewayStorageFailure } from "../../infra/sqlite-error-diagnostics.js"; -import { extractErrorHttpStatus, parseApiErrorInfo } from "../../shared/assistant-error-format.js"; -import { isSessionTranscriptValidationErrorMessage } from "./message-patterns.js"; +import { + extractErrorHttpStatus, + formatTransportErrorCopy, + parseApiErrorInfo, +} from "../../shared/assistant-error-format.js"; +import { classifyFailoverSignalCore } from "./classify-core.js"; +import { isContextOverflowErrorFromTables } from "./context-overflow-tables.js"; +import { + isServerErrorMessage, + isSessionTranscriptValidationErrorMessage, +} from "./message-patterns.js"; +import { extractFailoverSignalDetails } from "./signal-details.js"; import type { FailoverReason } from "./signal.js"; export const ERROR_PREFIX_RE = @@ -154,3 +165,60 @@ export function renderAssistantFormatFailureCopy(message: { } return undefined; } + +/** Classify saved error facts without loading providers or publishing their raw diagnostics. */ +export function renderRecordedAssistantFailureCopy(message: { + errorMessage?: unknown; + errorBody?: unknown; + errorCode?: unknown; + errorType?: unknown; +}): string | undefined { + const formatCopy = renderAssistantFormatFailureCopy(message); + if (formatCopy) { + return formatCopy; + } + const raw = typeof message.errorMessage === "string" ? message.errorMessage.trim() : ""; + if (raw === "Worker inference result exceeds the transcript message limit.") { + return "The worker could not save the model response because it exceeded the message size limit. Retry with a smaller response or continue on the Gateway. Earlier actions may have completed; verify their results before continuing."; + } + if ( + raw === + "Cloud worker could not preserve authoritative provider replay. Stop or reclaim the cloud worker, then retry locally." + ) { + return "The worker could not preserve the model's continuation data. Stop or reclaim the worker, then retry on the Gateway. Earlier actions may have completed; verify their results before continuing."; + } + const info = parseApiErrorInfo(raw); + const code = typeof message.errorCode === "string" ? message.errorCode : info?.code; + const status = extractErrorHttpStatus(raw)?.code; + const classification = classifyFailoverSignalCore({ + message: raw, + code, + errorType: typeof message.errorType === "string" ? message.errorType : info?.type, + status, + details: extractFailoverSignalDetails(message.errorBody), + }); + if ( + classification?.kind === "context_overflow" || + [message.errorCode, message.errorType, raw].some( + (value) => + typeof value === "string" && + (normalizeLowercaseStringOrEmpty(value) === "context_overflow" || + isContextOverflowErrorFromTables(value)), + ) + ) { + return "Context overflow: this conversation is too large for the model. Try /compact, use /new to start a fresh session, or retry the command with a tighter output limit."; + } + const classifiedCopy = renderAssistantRequestFailureCopy({ + code, + status, + // The legacy timeout retry bucket also includes explicit server failures. + reason: + classification?.reason === "timeout" && isServerErrorMessage(raw) + ? "server_error" + : classification?.reason, + }); + if (status !== undefined || (classification && classification.reason !== "timeout")) { + return classifiedCopy; + } + return formatTransportErrorCopy([raw, code].filter(Boolean).join(" ")) ?? classifiedCopy; +} diff --git a/src/agents/session-placement-admission.ts b/src/agents/session-placement-admission.ts index 0efd3b88812f..bba61e395f07 100644 --- a/src/agents/session-placement-admission.ts +++ b/src/agents/session-placement-admission.ts @@ -47,6 +47,7 @@ type SessionPlacementSandboxParams = { }; export type SessionPlacementAdmissionProvider = { + resolveRuntimeOverride?: (identity: Omit) => string | undefined; assertCompactionSuccessorAllowed: (params: { currentTarget: SessionTranscriptRuntimeTarget; successorSessionId: string; @@ -87,6 +88,13 @@ export function installSessionPlacementAdmissionProvider( }; } +/** Carries placement-owned runtime selection into candidate preparation and execution. */ +export function resolveSessionPlacementRuntimeOverride( + identity: Omit, +): string | undefined { + return state.provider?.resolveRuntimeOverride?.(identity); +} + /** Captures the exact placement owner, including standalone absence, before awaited work. */ export function captureSessionPlacementCompactionSuccessorAssertion(): SessionPlacementAdmissionProvider["assertCompactionSuccessorAllowed"] { const provider = state.provider; diff --git a/src/auto-reply/reply/agent-runner-fallback-candidate.ts b/src/auto-reply/reply/agent-runner-fallback-candidate.ts index 1a548d369e6d..e820fa1a5c19 100644 --- a/src/auto-reply/reply/agent-runner-fallback-candidate.ts +++ b/src/auto-reply/reply/agent-runner-fallback-candidate.ts @@ -88,14 +88,13 @@ export async function runAgentFallbackCandidates(params: AgentFallbackCycleParam milestone: "before_model_fallback", }); const selection = resolveModelFallbackOptions(params.effectiveRun, params.runtimeConfig); - const resolveCandidateRuntime = (provider: string, model: string) => { + const resolveCandidateRuntime = ( + provider: string, + model: string, + sessionRuntimeOverride: string | undefined, + ) => { const candidateRun = resolveFallbackCandidateRun(params.effectiveRun, provider, model); const activeEntry = params.liveModelSwitchRuntimeEntry ?? turn.getActiveSessionEntry(); - const sessionRuntimeOverride = resolveSessionRuntimeOverrideForProvider({ - provider, - entry: activeEntry, - cfg: params.runtimeConfig, - }); const pinnedHarnessId = resolveSessionPinnedHarnessId(activeEntry); const locksPersistedHarness = pinnedHarnessId !== undefined && pinnedHarnessId === sessionRuntimeOverride; @@ -164,8 +163,8 @@ export async function runAgentFallbackCandidates(params: AgentFallbackCycleParam entry: params.liveModelSwitchRuntimeEntry ?? turn.getActiveSessionEntry(), cfg: params.runtimeConfig, }), - resolveContextEngineHost: (provider, model) => { - const runtime = resolveCandidateRuntime(provider, model); + resolveContextEngineHost: (provider, model, runtimeOverride) => { + const runtime = resolveCandidateRuntime(provider, model, runtimeOverride); if (!runtime.useCliExecution) { return undefined; } @@ -218,7 +217,7 @@ export async function runAgentFallbackCandidates(params: AgentFallbackCycleParam params.state.attemptedRuntimeProvider = provider; params.state.attemptedRuntimeModel = model; const runtime = params.timing.measureSync("fallback_resolve_runtime", () => - resolveCandidateRuntime(provider, model), + resolveCandidateRuntime(provider, model, runOptions.agentHarnessRuntimeOverride), ); const candidateRun = runtime.candidateRun; bindSourceReplyDeliveryRuntime(candidateRun, sourceReplyDeliveryRuntime); @@ -242,6 +241,7 @@ export async function runAgentFallbackCandidates(params: AgentFallbackCycleParam agentId: turn.followupRun.run.agentId, sessionKey: turn.followupRun.run.runtimePolicySessionKey ?? turn.sessionKey, sessionEntry: params.liveModelSwitchRuntimeEntry ?? turn.getActiveSessionEntry(), + agentRuntime: runtime.sessionRuntimeOverride, }); const candidateFastMode = resolveRunFastModeForFallbackCandidate({ run: candidateRun, diff --git a/src/auto-reply/reply/agent-runner-memory.test-support.ts b/src/auto-reply/reply/agent-runner-memory.test-support.ts new file mode 100644 index 000000000000..3bb0bf0c1a20 --- /dev/null +++ b/src/auto-reply/reply/agent-runner-memory.test-support.ts @@ -0,0 +1,98 @@ +import { createAssistantErrorTranscript } from "../../agents/assistant-error-transcript.js"; +import type { runEmbeddedAgentEntry } from "../../agents/embedded-agent-runner/run-entry.js"; +import type { EmbeddedAgentRunResult } from "../../agents/embedded-agent-runner/types.js"; +import type { ensureSelectedAgentHarnessPlugin } from "../../agents/harness/runtime-plugin.js"; +import type { ModelFallbackAttemptProvenance } from "../../agents/model-fallback.types.js"; +import { requireActivePluginRegistry } from "../../plugins/runtime.js"; + +export type ModelFallbackParams = { + provider?: string; + model?: string; + abortSignal?: AbortSignal; + agentId?: string; + sessionId?: string; + sessionKey?: string; + fallbacksOverride?: unknown[]; + requestedRouteResolution?: "raw" | "resolved"; + userLockedAuthProfileId?: string; + resolveAgentHarnessRuntimeOverride?: (provider: string, model: string) => string | undefined; + prepareAgentHarnessRuntime?: (params: { + provider: string; + model: string; + agentHarnessRuntimeOverride?: string; + }) => Promise | void; + run: ( + provider: string, + model: string, + options: { + allowTransientCooldownProbe?: boolean; + isFinalFallbackAttempt?: boolean; + modelRoutingProvenance: ModelFallbackAttemptProvenance; + }, + ) => Promise; +}; + +export function createMemoryRunEntryMockImplementation(deps: { + runWithModelFallback: (params: ModelFallbackParams) => Promise; + ensureSelectedAgentHarnessPlugin: typeof ensureSelectedAgentHarnessPlugin; +}) { + return async (params: Parameters>[0]) => { + const assistantErrorTranscript = createAssistantErrorTranscript({ + runId: params.identity.runId, + }); + const fallbackResult = (await deps.runWithModelFallback({ + ...params.selection, + ...params.identity, + abortSignal: params.abortSignal, + resolveAgentHarnessRuntimeOverride: params.harness.resolveRuntimeOverride, + prepareAgentHarnessRuntime: async ({ + provider, + model, + agentHarnessRuntimeOverride, + }: { + provider: string; + model: string; + agentHarnessRuntimeOverride?: string; + }) => { + await deps.ensureSelectedAgentHarnessPlugin({ + config: params.selection.cfg, + provider, + modelId: model, + agentId: params.identity.agentId, + sessionKey: params.harness.sessionKey, + agentHarnessId: agentHarnessRuntimeOverride, + agentHarnessRuntimeOverride, + workspaceDir: params.harness.workspaceDir, + pluginRegistry: requireActivePluginRegistry(), + }); + }, + run: (provider: string, model: string, options: Parameters[2]) => + params.runCandidate(provider, model, { + agentHarnessRuntimeOverride: params.harness.resolveRuntimeOverride(provider, model), + assistantErrorTranscript, + classifyResult: () => undefined, + allowTransientCooldownProbe: options.allowTransientCooldownProbe, + isFinalFallbackAttempt: options.isFinalFallbackAttempt, + isFallbackRetry: false, + modelRoutingProvenance: options.modelRoutingProvenance, + contextEngineLogicalTurnLease: {} as never, + onContextEngineTurnCandidate: () => {}, + }), + })) as { + outcome?: "completed" | "exhausted"; + result: EmbeddedAgentRunResult; + provider: string; + model: string; + attempts: []; + }; + return { + ...fallbackResult, + outcome: fallbackResult.outcome ?? ("completed" as const), + terminal: { + outcome: { reason: "completed" as const, status: "ok" as const }, + metadata: {}, + }, + settleSessionOverride: async () => undefined, + }; + }; +} diff --git a/src/auto-reply/reply/agent-runner-memory.test.ts b/src/auto-reply/reply/agent-runner-memory.test.ts index 7b89bbc8dad3..98bd1522acd8 100644 --- a/src/auto-reply/reply/agent-runner-memory.test.ts +++ b/src/auto-reply/reply/agent-runner-memory.test.ts @@ -11,12 +11,9 @@ import { type AdmittedRunContext, type PreparedAgentRunAdmission, } from "../../agents/admitted-run-context.js"; -import { createAssistantErrorTranscript } from "../../agents/assistant-error-transcript.js"; import { testing as cliBackendsTesting } from "../../agents/cli-backends.test-support.js"; import { resetContextWindowCacheForTest } from "../../agents/context.js"; import { acceptCompactionSuccessor } from "../../agents/embedded-agent-runner/compaction-successor.js"; -import type { runEmbeddedAgentEntry } from "../../agents/embedded-agent-runner/run-entry.js"; -import type { EmbeddedAgentRunResult } from "../../agents/embedded-agent-runner/types.js"; import type { ModelFallbackAttemptProvenance } from "../../agents/model-fallback.types.js"; import { withSessionCompactionPersistence } from "../../agents/sessions/session-compaction-persistence.js"; import { SessionManager } from "../../agents/sessions/session-manager.js"; @@ -48,6 +45,10 @@ import { runMemoryFlushIfNeeded as runMemoryFlushIfNeededRaw, runSessionCompactionIfNeeded as runSessionCompactionIfNeededRaw, } from "./agent-runner-memory.js"; +import { + createMemoryRunEntryMockImplementation, + type ModelFallbackParams, +} from "./agent-runner-memory.test-support.js"; import { createTestFollowupRun, withTestModelContextTokens, @@ -234,33 +235,6 @@ async function writeTestSessionTranscript(params: { await waitForSessionTranscriptProjection(scope); } -type ModelFallbackParams = { - provider?: string; - model?: string; - abortSignal?: AbortSignal; - agentId?: string; - sessionId?: string; - sessionKey?: string; - fallbacksOverride?: unknown[]; - requestedRouteResolution?: "raw" | "resolved"; - userLockedAuthProfileId?: string; - resolveAgentHarnessRuntimeOverride?: (provider: string, model: string) => string | undefined; - prepareAgentHarnessRuntime?: (params: { - provider: string; - model: string; - agentHarnessRuntimeOverride?: string; - }) => Promise | void; - run: ( - provider: string, - model: string, - options: { - allowTransientCooldownProbe?: boolean; - isFinalFallbackAttempt?: boolean; - modelRoutingProvenance: ModelFallbackAttemptProvenance; - }, - ) => Promise; -}; - function modelRoutingProvenance( requestedProvider: string, requestedModel: string, @@ -453,71 +427,12 @@ describe("runMemoryFlushIfNeeded", () => { model, attempts: [], })); - runEmbeddedAgentEntryMock - .mockReset() - .mockImplementation( - async (params: Parameters>[0]) => { - const assistantErrorTranscript = createAssistantErrorTranscript({ - runId: params.identity.runId, - }); - const fallbackResult = (await runWithModelFallbackMock({ - ...params.selection, - ...params.identity, - abortSignal: params.abortSignal, - resolveAgentHarnessRuntimeOverride: params.harness.resolveRuntimeOverride, - prepareAgentHarnessRuntime: async ({ - provider, - model, - agentHarnessRuntimeOverride, - }: { - provider: string; - model: string; - agentHarnessRuntimeOverride?: string; - }) => { - await ensureSelectedAgentHarnessPluginMock({ - config: params.selection.cfg, - provider, - modelId: model, - agentId: params.identity.agentId, - sessionKey: params.harness.sessionKey, - agentHarnessId: agentHarnessRuntimeOverride, - agentHarnessRuntimeOverride, - workspaceDir: params.harness.workspaceDir, - }); - }, - run: ( - provider: string, - model: string, - options: Parameters[2], - ) => - params.runCandidate(provider, model, { - assistantErrorTranscript, - classifyResult: () => undefined, - allowTransientCooldownProbe: options.allowTransientCooldownProbe, - isFinalFallbackAttempt: options.isFinalFallbackAttempt, - isFallbackRetry: false, - modelRoutingProvenance: options.modelRoutingProvenance, - contextEngineLogicalTurnLease: {} as never, - onContextEngineTurnCandidate: () => {}, - }), - })) as { - outcome?: "completed" | "exhausted"; - result: EmbeddedAgentRunResult; - provider: string; - model: string; - attempts: []; - }; - return { - ...fallbackResult, - outcome: fallbackResult.outcome ?? ("completed" as const), - terminal: { - outcome: { reason: "completed" as const, status: "ok" as const }, - metadata: {}, - }, - settleSessionOverride: async () => undefined, - }; - }, - ); + runEmbeddedAgentEntryMock.mockReset().mockImplementation( + createMemoryRunEntryMockImplementation({ + runWithModelFallback: runWithModelFallbackMock, + ensureSelectedAgentHarnessPlugin: ensureSelectedAgentHarnessPluginMock, + }), + ); compactEmbeddedAgentSessionMock.mockReset().mockResolvedValue({ ok: true, compacted: true, diff --git a/src/auto-reply/reply/agent-runner-memory.ts b/src/auto-reply/reply/agent-runner-memory.ts index adaefc9ea4cd..3d2cbd417147 100644 --- a/src/auto-reply/reply/agent-runner-memory.ts +++ b/src/auto-reply/reply/agent-runner-memory.ts @@ -1692,11 +1692,7 @@ export async function runMemoryFlushIfNeeded(params: { sessionOverride: { kind: "preserve" }, abortSignal: deferredLifecycle.signal, runCandidate: async (provider, model, runOptions) => { - const sessionRuntimeOverride = resolveSessionRuntimeOverrideForProvider({ - provider, - entry: activeSessionEntry, - cfg: params.cfg, - }); + const sessionRuntimeOverride = runOptions.agentHarnessRuntimeOverride; const candidateThinkLevel = resolveRunThinkingLevelForFallbackCandidate({ cfg: params.cfg, provider, diff --git a/src/cron/isolated-agent/run-candidate-runtime.ts b/src/cron/isolated-agent/run-candidate-runtime.ts new file mode 100644 index 000000000000..d01cc25b47a7 --- /dev/null +++ b/src/cron/isolated-agent/run-candidate-runtime.ts @@ -0,0 +1,41 @@ +import { resolveCliRuntimeExecutionProvider } from "../../agents/model-runtime-aliases.js"; +import { isCliProvider } from "./run-execution.runtime.js"; +import { resolveEffectiveAgentRuntime } from "./run.runtime.js"; +import type { CronRunExecutionParams } from "./run.types.js"; + +/** Shares candidate execution policy between harness preparation and dispatch. */ +export function createCronCandidateExecutionResolver( + params: Pick< + CronRunExecutionParams, + "cfgWithAgentDefaults" | "agentId" | "runSessionKey" | "cronSession" + >, +) { + return (provider: string, model: string, sessionRuntimeOverride: string | undefined) => { + const executionProvider = sessionRuntimeOverride + ? isCliProvider(sessionRuntimeOverride, params.cfgWithAgentDefaults) + ? sessionRuntimeOverride + : provider + : (resolveCliRuntimeExecutionProvider({ + provider, + cfg: params.cfgWithAgentDefaults, + agentId: params.agentId, + modelId: model, + }) ?? provider); + const runtime = + sessionRuntimeOverride ?? + resolveEffectiveAgentRuntime({ + cfg: params.cfgWithAgentDefaults, + provider, + modelId: model, + agentId: params.agentId, + sessionKey: params.runSessionKey, + sessionEntry: params.cronSession.sessionEntry, + }); + return { + sessionRuntimeOverride, + executionProvider, + cliExecution: isCliProvider(executionProvider, params.cfgWithAgentDefaults), + runtime, + }; + }; +} diff --git a/src/cron/isolated-agent/run-executor.ts b/src/cron/isolated-agent/run-executor.ts index 08ab73055548..c10b2bb805ea 100644 --- a/src/cron/isolated-agent/run-executor.ts +++ b/src/cron/isolated-agent/run-executor.ts @@ -18,7 +18,6 @@ import { createDeferredEmbeddedRunLifecycleManager } from "../../agents/embedded import type { FastModeAutoProgressState } from "../../agents/fast-mode.js"; import { runAgentHarnessBeforeMessageWriteHook } from "../../agents/harness/hook-helpers.js"; import { findModelInCatalog, modelSupportsInput } from "../../agents/model-catalog-lookup.js"; -import { resolveCliRuntimeExecutionProvider } from "../../agents/model-runtime-aliases.js"; import { resolveConfiguredThinkingDefault } from "../../agents/model-thinking-default.js"; import { rootedAgentRunParams } from "../../agents/rooted-run-params.js"; import { resolveScheduledToolPolicyContext } from "../../agents/scheduled-tool-policy.js"; @@ -53,13 +52,13 @@ import { assertCronRuntimeAuthorityCandidate, prepareCronPromptRunAdmission, } from "./run-admission.js"; +import { createCronCandidateExecutionResolver } from "./run-candidate-runtime.js"; import { appendCronDeliveryInstruction, buildCronDeliveryTargetRuntimeContext, } from "./run-delivery-trace.js"; import { getCliSessionBinding, - isCliProvider, LiveSessionModelSwitchError, logWarn, normalizeVerboseLevel, @@ -75,7 +74,7 @@ import { setCronSessionRuntimeModel, syncCronSessionLiveSelection, } from "./run-session-state.js"; -import { resolveEffectiveAgentRuntime, resolveThinkingSelection } from "./run.runtime.js"; +import { resolveThinkingSelection } from "./run.runtime.js"; import type { AgentTurnPayload, CronCompletedPromptRun, @@ -241,28 +240,7 @@ function createCronPromptExecutor( const currentAttemptCommittedMedia = () => hasNewGeneratedMediaTaskForSessionKey(params.runSessionKey, attemptMediaTaskIds); - const resolveCandidateExecution = (provider: string, model: string) => { - const sessionRuntimeOverride = resolveSessionRuntimeOverrideForProvider({ - provider, - entry: params.cronSession.sessionEntry, - cfg: params.cfgWithAgentDefaults, - }); - const executionProvider = sessionRuntimeOverride - ? isCliProvider(sessionRuntimeOverride, params.cfgWithAgentDefaults) - ? sessionRuntimeOverride - : provider - : (resolveCliRuntimeExecutionProvider({ - provider, - cfg: params.cfgWithAgentDefaults, - agentId: params.agentId, - modelId: model, - }) ?? provider); - return { - sessionRuntimeOverride, - executionProvider, - cliExecution: isCliProvider(executionProvider, params.cfgWithAgentDefaults), - }; - }; + const resolveCandidateExecution = createCronCandidateExecutionResolver(params); return async (promptText: string, runStartedAt: number): Promise => { // A retry can fail during preparation, before any backend start callback. @@ -349,10 +327,18 @@ function createCronPromptExecutor( workspaceDir: params.executionRoot ?? params.workspaceDir, sessionKey: params.runSessionKey, preparation: { kind: "direct" }, - resolveRuntimeOverride: (provider, model) => - resolveCandidateExecution(provider, model).sessionRuntimeOverride, - resolveContextEngineHost: (provider, model) => { - const { executionProvider, cliExecution } = resolveCandidateExecution(provider, model); + resolveRuntimeOverride: (provider) => + resolveSessionRuntimeOverrideForProvider({ + provider, + entry: params.cronSession.sessionEntry, + cfg: params.cfgWithAgentDefaults, + }), + resolveContextEngineHost: (provider, model, runtimeOverride) => { + const { executionProvider, cliExecution } = resolveCandidateExecution( + provider, + model, + runtimeOverride, + ); if (!cliExecution) { return undefined; } @@ -392,16 +378,16 @@ function createCronPromptExecutor( if (params.abortSignal?.aborted) { throw new Error(params.abortReason()); } - const { sessionRuntimeOverride, executionProvider, cliExecution } = - resolveCandidateExecution(providerOverride, modelOverride); - const candidateRuntime = resolveEffectiveAgentRuntime({ - cfg: params.cfgWithAgentDefaults, - provider: providerOverride, - modelId: modelOverride, - agentId: params.agentId, - sessionKey: params.runSessionKey, - sessionEntry: params.cronSession.sessionEntry, - }); + const { + sessionRuntimeOverride, + executionProvider, + cliExecution, + runtime: candidateRuntime, + } = resolveCandidateExecution( + providerOverride, + modelOverride, + runOptions.agentHarnessRuntimeOverride, + ); const candidateConfiguredThinkLevel = params.immutableThinkLevel ?? resolveConfiguredThinkingDefault({ diff --git a/src/gateway/chat-display-projection.core.ts b/src/gateway/chat-display-projection.core.ts index 262ee8210354..fa4462fd3574 100644 --- a/src/gateway/chat-display-projection.core.ts +++ b/src/gateway/chat-display-projection.core.ts @@ -1,15 +1,11 @@ import { STREAM_ERROR_FALLBACK_TEXT } from "@openclaw/ai/internal/shared"; import { GATEWAY_ASSISTANT_ERROR_FALLBACK_TEXT } from "@openclaw/gateway-protocol/gateway-error-details"; import { asOptionalRecord } from "@openclaw/normalization-core/record-coerce"; +import { normalizeOptionalString } from "@openclaw/normalization-core/string-coerce"; import { - normalizeLowercaseStringOrEmpty as normalizeErrorSignal, - normalizeOptionalString, -} from "@openclaw/normalization-core/string-coerce"; -import { - renderAssistantFormatFailureCopy, renderAssistantRequestFailureCopy, + renderRecordedAssistantFailureCopy, } from "../agents/failover/assistant-request-failure-copy.js"; -import { isContextOverflowErrorFromTables } from "../agents/failover/context-overflow-tables.js"; import { readTranscriptSenderIdentity } from "../chat/sender-identity.js"; import { projectAgentHistoryActivity, @@ -135,26 +131,6 @@ type ChatDisplayProjectionResult = { commentaryFallbacksObserved?: true; }; -const GATEWAY_ASSISTANT_CONTEXT_OVERFLOW_FALLBACK_TEXT = - "Context overflow: this conversation is too large for the model. Try /compact, use /new to start a fresh session, or retry the command with a tighter output limit."; - -function isContextOverflowErrorSignal(value: unknown): boolean { - if (typeof value !== "string") { - return false; - } - return ( - normalizeErrorSignal(value) === "context_overflow" || isContextOverflowErrorFromTables(value) - ); -} - -function isContextOverflowAssistantError(message: Record): boolean { - return ( - isContextOverflowErrorSignal(message.errorCode) || - isContextOverflowErrorSignal(message.errorType) || - isContextOverflowErrorSignal(message.errorMessage) - ); -} - function getAssistantErrorFallbackText(message: Record): string { return ( formatProviderRefusalText(message) ?? @@ -162,10 +138,8 @@ function getAssistantErrorFallbackText(message: Record): string storageFailure: classifyGatewayStorageFailure(message), code: typeof message.errorCode === "string" ? message.errorCode : undefined, }) ?? - renderAssistantFormatFailureCopy(message) ?? - (isContextOverflowAssistantError(message) - ? GATEWAY_ASSISTANT_CONTEXT_OVERFLOW_FALLBACK_TEXT - : GATEWAY_ASSISTANT_ERROR_FALLBACK_TEXT) + renderRecordedAssistantFailureCopy(message) ?? + GATEWAY_ASSISTANT_ERROR_FALLBACK_TEXT ); } @@ -212,7 +186,7 @@ function sanitizeAssistantErrorDisplayMessage( const terminalCopy = renderAssistantRequestFailureCopy({ code: typeof message.errorCode === "string" ? message.errorCode : undefined, - }) ?? renderAssistantFormatFailureCopy(message); + }) ?? renderRecordedAssistantFailureCopy(message); if (terminalCopy) { // Apply the normal visibility rules before adding host-owned failure copy. // Put it first in surviving text so phase filtering and display caps retain it. diff --git a/src/gateway/chat-display-projection.provider-owner.test.ts b/src/gateway/chat-display-projection.provider-owner.test.ts index b72dd2bb3d79..efe1aeea5b21 100644 --- a/src/gateway/chat-display-projection.provider-owner.test.ts +++ b/src/gateway/chat-display-projection.provider-owner.test.ts @@ -77,3 +77,85 @@ it("shows the upstream cache limit in persisted history without proxy metadata", expect(projected).not.toHaveProperty("errorMessage"); expect(classifyProviderFailoverSignalWithPlugin).not.toHaveBeenCalled(); }); + +it.each([ + { + errorMessage: '429: {"error":{"type":"rate_limit_error","message":"PRIVATE_CANARY"}}', + expected: + "⚠️ LLM request failed (rate limited, HTTP 429). This is usually temporary — try again shortly.", + }, + { + errorMessage: '503: {"error":{"type":"server_error","message":"DNS PRIVATE_CANARY"}}', + expected: + "⚠️ LLM request failed (provider internal error, HTTP 503). This is usually temporary — try again shortly.", + }, + { + errorMessage: "Internal server error", + expected: + "⚠️ LLM request failed (provider internal error). This is usually temporary — try again shortly.", + }, + { + errorMessage: "Request timed out", + expected: + "⚠️ LLM request failed (request timed out). This is usually temporary — try again shortly.", + }, + { + errorMessage: '{"error":{"code":"invalid_api_key","message":"PRIVATE_CANARY"}}', + expected: + "⚠️ LLM request failed (authentication failed). Re-authenticate the provider and try again.", + }, + { + errorMessage: '429: {"error":{"code":"insufficient_quota","message":"PRIVATE_CANARY"}}', + expected: + "⚠️ LLM request failed (provider billing issue, HTTP 429). Check provider billing and try again.", + }, + { + errorCode: "ECONNRESET", + errorMessage: "PRIVATE_CANARY", + expected: "LLM request failed: network connection was interrupted.", + }, + { + errorMessage: "Worker inference result exceeds the transcript message limit.", + expected: + "The worker could not save the model response because it exceeded the message size limit. Retry with a smaller response or continue on the Gateway. Earlier actions may have completed; verify their results before continuing.", + }, +])( + "retains a safe diagnosis in history and repeated projections: $errorMessage", + ({ expected, ...error }) => { + const raw = { + role: "assistant", + stopReason: "error", + content: [], + ...error, + __openclaw: { runId: "failed-run" }, + }; + const original = structuredClone(raw); + const projected = projectChatDisplayMessage(raw); + expect(projected).toMatchObject({ content: [{ type: "text", text: expected }] }); + expect(projected).not.toHaveProperty("errorMessage"); + expect(projected).not.toHaveProperty("errorCode"); + expect(JSON.stringify(projected)).not.toContain("PRIVATE_CANARY"); + expect(projectChatDisplayMessage(projected)).toEqual(projected); + expect(raw).toEqual(original); + expect(classifyProviderFailoverSignalWithPlugin).not.toHaveBeenCalled(); + }, +); + +it("keeps the failure guidance alongside partial reply text", () => { + const projected = projectChatDisplayMessage({ + role: "assistant", + stopReason: "error", + errorMessage: "429: PRIVATE_CANARY", + content: [{ type: "text", text: "The first step completed." }], + }); + expect(projected).toMatchObject({ + content: [ + { + type: "text", + text: "⚠️ LLM request failed (rate limited, HTTP 429). This is usually temporary — try again shortly.\n\nThe first step completed.", + }, + ], + }); + expect(JSON.stringify(projected)).not.toContain("PRIVATE_CANARY"); + expect(projectChatDisplayMessage(projected)).toEqual(projected); +}); diff --git a/src/gateway/server-methods/server-methods.test.ts b/src/gateway/server-methods/server-methods.test.ts index 4eb39a6e4412..d18e64e78d38 100644 --- a/src/gateway/server-methods/server-methods.test.ts +++ b/src/gateway/server-methods/server-methods.test.ts @@ -1180,12 +1180,15 @@ describe("projectChatDisplayMessages", () => { const safeFailureContent = [ { type: "text", text: "The agent run failed before producing a reply." }, ]; + const networkFailureText = "LLM request failed: network connection error."; + const networkFailureContent = (reply?: string, type = "text") => [ + { type, text: [networkFailureText, reply].filter(Boolean).join("\n\n") }, + ]; const privateError = "private upstream at secret.internal.example failed"; const displayErrorCases: Array<{ name: string; message: Record; content: Array>; - visibleText?: string; }> = [ { name: "projects empty assistant error turns as a generic safe failure", @@ -1193,9 +1196,9 @@ describe("projectChatDisplayMessages", () => { content: safeFailureContent, }, { - name: "projects empty text-block assistant errors as a generic safe failure", + name: "projects empty text-block assistant errors as a safe network failure", message: { content: [{ type: "text", text: "" }], errorMessage: "Connection error." }, - content: safeFailureContent, + content: networkFailureContent(), }, { name: "projects provider refusals before classifying their explanation text", @@ -1223,15 +1226,15 @@ describe("projectChatDisplayMessages", () => { content: [{ type: "output_text", text: "A partial reply before the run failed." }], errorMessage: "Connection error.", }, - content: [{ type: "output_text", text: "A partial reply before the run failed." }], + content: networkFailureContent("A partial reply before the run failed.", "output_text"), }, { - name: "projects thinking-only assistant errors as a generic safe failure", + name: "projects thinking-only assistant errors as a safe network failure", message: { content: [{ type: "thinking", thinking: "private upstream details" }], errorMessage: "Connection error.", }, - content: safeFailureContent, + content: networkFailureContent(), }, { name: "preserves a safe failure for a synthetic sentinel followed only by private thinking", @@ -1261,24 +1264,23 @@ describe("projectChatDisplayMessages", () => { content: safeFailureContent, }, { - name: "projects commentary-phase assistant errors as a visible generic safe failure", + name: "projects commentary-phase assistant errors as a visible safe network failure", message: { phase: "commentary", content: [], text: "private upstream details", errorMessage: "Connection error.", }, - content: safeFailureContent, + content: networkFailureContent(), }, { - name: "leaves legacy top-level assistant error text unchanged", + name: "preserves legacy top-level assistant text with safe network failure details", message: { content: [], text: "A real reply before the run failed.", errorMessage: "Connection error.", }, - content: [], - visibleText: "A real reply before the run failed.", + content: networkFailureContent("A real reply before the run failed."), }, { name: "preserves partial error replies without hidden reasoning or diagnostics", @@ -1366,7 +1368,7 @@ describe("projectChatDisplayMessages", () => { }, ]; - it.each(displayErrorCases)("$name", ({ message, content, visibleText }) => { + it.each(displayErrorCases)("$name", ({ message, content }) => { const result = projectChatDisplayMessages([ { role: "assistant", stopReason: "error", timestamp: 1, ...message }, ]); @@ -1376,7 +1378,6 @@ describe("projectChatDisplayMessages", () => { content, stopReason: "error", timestamp: 1, - ...(visibleText === undefined ? {} : { text: visibleText }), }, ]); expect(JSON.stringify(result)).not.toContain("secret.internal.example"); @@ -1454,7 +1455,7 @@ describe("projectChatDisplayMessages", () => { ["output_text", "NO_REPLY"], ["input_text", ""], ["input_text", "NO_REPLY"], - ])("projects hidden %s assistant errors %j as a generic safe failure", (type, text) => { + ])("projects hidden %s assistant errors %j as a safe network failure", (type, text) => { const result = projectChatDisplayMessages([ { role: "assistant", @@ -1465,9 +1466,7 @@ describe("projectChatDisplayMessages", () => { }, ]); - expect(result[0]?.content).toEqual([ - { type: "text", text: "The agent run failed before producing a reply." }, - ]); + expect(result[0]?.content).toEqual(networkFailureContent()); }); it.each(["NO_REPLY", STREAM_ERROR_FALLBACK_TEXT])( diff --git a/src/gateway/server.chat-recovered-output.test.ts b/src/gateway/server.chat-recovered-output.test.ts index 2cb9e96d416d..a198c9b72c3f 100644 --- a/src/gateway/server.chat-recovered-output.test.ts +++ b/src/gateway/server.chat-recovered-output.test.ts @@ -239,7 +239,9 @@ describe("registered chat.send recovered output over Responses HTTP", () => { ); if (failed) { expect(completed.status).toBe("error"); - expect(messageText(history.messages.at(-1))).toBe(prefix); + expect(messageText(history.messages.at(-1))).toBe( + `⚠️ LLM request failed (provider internal error). This is usually temporary — try again shortly.\n\n${prefix}`, + ); expect(terminal.some((event) => event.state === "error")).toBe(true); const deltas = events.filter( (event): event is Extract => diff --git a/src/gateway/server.chat.gateway-server-chat.test.ts b/src/gateway/server.chat.gateway-server-chat.test.ts index b79db4df5225..aa70e8318450 100644 --- a/src/gateway/server.chat.gateway-server-chat.test.ts +++ b/src/gateway/server.chat.gateway-server-chat.test.ts @@ -1543,6 +1543,9 @@ describe("gateway server chat", () => { }); }); + const contextOverflowCopy = + "Context overflow: this conversation is too large for the model. Try /compact, use /new to start a fresh session, or retry the command with a tighter output limit."; + test.each([ { name: "structured context-overflow code", @@ -1550,7 +1553,7 @@ describe("gateway server chat", () => { errorCode: "context_overflow", errorMessage: "private upstream body: 203557 tokens sent", }, - overflow: true, + expected: contextOverflowCopy, }, { name: "provider request-too-large code", @@ -1558,7 +1561,7 @@ describe("gateway server chat", () => { errorCode: "request_too_large", errorMessage: "private upstream body: 196607 tokens sent", }, - overflow: true, + expected: contextOverflowCopy, }, { name: "provider context-window message", @@ -1566,12 +1569,12 @@ describe("gateway server chat", () => { errorType: "invalid_request_error", errorMessage: "Request size exceeds model context window: 203557 tokens", }, - overflow: true, + expected: contextOverflowCopy, }, { name: "embedded context-overflow message", fields: { errorMessage: "Unhandled stop reason: context_overflow" }, - overflow: true, + expected: contextOverflowCopy, }, { name: "token-per-minute rate limit", @@ -1579,16 +1582,17 @@ describe("gateway server chat", () => { errorCode: "rate_limit_exceeded", errorMessage: "413 request too large: 203557 tokens per minute (TPM)", }, - overflow: false, + expected: + "⚠️ LLM request failed (rate limited, HTTP 413). This is usually temporary — try again shortly.", }, { name: "private upstream failure", fields: { errorMessage: "private upstream at secret.internal.example failed" }, - overflow: false, + expected: "The agent run failed before producing a reply.", }, ])( "chat.history safely displays $name over authenticated WebSocket", - async ({ fields, overflow }) => { + async ({ fields, expected }) => { const historyMessages = await loadChatHistoryWithMessages([ { role: "assistant", @@ -1599,11 +1603,7 @@ describe("gateway server chat", () => { }, ]); - expect(collectHistoryTextValues(historyMessages)).toEqual([ - overflow - ? "Context overflow: this conversation is too large for the model. Try /compact, use /new to start a fresh session, or retry the command with a tighter output limit." - : "The agent run failed before producing a reply.", - ]); + expect(collectHistoryTextValues(historyMessages)).toEqual([expected]); const wirePayload = JSON.stringify(historyMessages); expect(wirePayload).not.toContain("203557"); expect(wirePayload).not.toContain("196607"); diff --git a/src/gateway/worker-environments/inference-runtime.test.ts b/src/gateway/worker-environments/inference-runtime.test.ts index 0e79b1a7b4f1..4d87b5a33d04 100644 --- a/src/gateway/worker-environments/inference-runtime.test.ts +++ b/src/gateway/worker-environments/inference-runtime.test.ts @@ -7,6 +7,7 @@ import type { AssistantMessage } from "../../llm/types.js"; import { createAssistantMessageEventStream } from "../../llm/utils/event-stream.js"; import { createEmptyPluginRegistry } from "../../plugins/registry-empty.js"; import { resetPluginRuntimeStateForTest, setActivePluginRegistry } from "../../plugins/runtime.js"; +import { parseApiErrorInfo } from "../../shared/assistant-error-format.js"; import { isWorkerTranscriptMessageFrameSafe, WORKER_PROVIDER_REPLAY_LOCAL_RETRY_MESSAGE, @@ -74,6 +75,92 @@ describe("worker inference provider runtime", () => { expect(runtime.releaseRuntime).toHaveBeenCalledOnce(); }); + it.each(["insufficient_quota", "invalid_api_key", "context_length_exceeded"])( + "preserves a streamed provider failure identified only by %s", + async (errorCode) => { + const runtime = setup(); + runtime.stream.mockImplementation(() => { + const stream = createAssistantMessageEventStream(); + stream.push({ + type: "error", + reason: "error", + error: { ...finalMessage(), stopReason: "error", errorCode }, + }); + return stream; + }); + + const outcome = await runtime.executor(params(request(), vi.fn())); + + expect(outcome).toMatchObject({ type: "error", reason: "provider-error", usage }); + if (outcome.type !== "error") { + throw new Error("expected provider failure"); + } + expect(parseApiErrorInfo(outcome.message)?.code).toBe(errorCode); + expect(validateWorkerInferenceTerminalOutcome(outcome)).toBe(true); + }, + ); + + it("bounds streamed provider details without losing structured failure facts", async () => { + const runtime = setup(); + const secret = `stream-secret-${"a".repeat(48)}`; + runtime.stream.mockImplementation(() => { + const stream = createAssistantMessageEventStream(); + stream.push({ + type: "error", + reason: "error", + error: { + ...finalMessage(), + stopReason: "error", + errorMessage: `429: Authorization: Bearer ${secret} ${'diagnostic " \\ '.repeat(80)}`, + errorCode: "rate_limit_exceeded", + errorType: "rate_limit_error", + }, + }); + return stream; + }); + + const outcome = await runtime.executor(params(request(), vi.fn())); + + if (outcome.type !== "error") { + throw new Error("expected provider failure"); + } + expect(parseApiErrorInfo(outcome.message)).toMatchObject({ + httpCode: "429", + code: "rate_limit_exceeded", + type: "rate_limit_error", + }); + expect(outcome.message).not.toContain(secret); + expect(outcome.message.length).toBeLessThanOrEqual(256); + expect(validateWorkerInferenceTerminalOutcome(outcome)).toBe(true); + }); + + it.each([ + { name: "short body", status: 503, code: "upstream_unavailable", detail: "Unavailable" }, + { name: "long body", status: 429, code: "insufficient_quota", detail: "x".repeat(520) }, + { name: "bigint diagnostic", status: 429, code: "insufficient_quota", detail: 1n }, + ])("preserves a thrown provider HTTP failure ($name)", async ({ status, code, detail }) => { + const runtime = setup(); + runtime.stream.mockImplementation(() => { + throw Object.assign(new Error("Request rejected"), { + status, + body: { error: { detail, code, message: "Provider request rejected" } }, + }); + }); + + const outcome = await runtime.executor(params(request(), vi.fn())); + + if (outcome.type !== "error") { + throw new Error("expected provider failure"); + } + expect(outcome).toMatchObject({ type: "error", reason: "provider-error" }); + expect(parseApiErrorInfo(outcome.message)).toMatchObject({ + httpCode: String(status), + code, + }); + expect(outcome.message).toContain("Provider request rejected"); + expect(runtime.releaseRuntime).toHaveBeenCalledOnce(); + }); + it("keeps provider construction and execution on the leased generation", async () => { const generationA = createEmptyPluginRegistry(); const generationB = createEmptyPluginRegistry(); diff --git a/src/gateway/worker-environments/inference-runtime.ts b/src/gateway/worker-environments/inference-runtime.ts index 0bb6bb69eb61..26476e0a2035 100644 --- a/src/gateway/worker-environments/inference-runtime.ts +++ b/src/gateway/worker-environments/inference-runtime.ts @@ -5,7 +5,6 @@ import type { WorkerInferenceContext, WorkerInferenceEventParams, WorkerInferenceStartParams, - WorkerInferenceTerminalOutcome, } from "../../../packages/gateway-protocol/src/schema/worker-inference.js"; import { resolveAgentDir, resolveAgentWorkspaceDir } from "../../agents/agent-scope.js"; import { resolveSessionAuthSelection } from "../../agents/auth-profiles/session-override.js"; @@ -62,12 +61,14 @@ import { withPluginRuntimeGenerationScope } from "../../plugins/runtime/generati import { estimateUsageCost, resolveModelCostConfig } from "../../utils/usage-format.js"; import { WORKER_PROVIDER_REPLAY_LOCAL_RETRY_MESSAGE } from "../../worker/transcript-message.js"; import { + ERROR_MESSAGES, + inferenceError, projectWorkerInferenceTerminalMessage, type WorkerInferenceModelIdentity, } from "./inference-terminal-message.js"; import { createWorkerToolCallStream } from "./inference-tool-call-stream.js"; import { resolveWorkerSessionTarget, type ResolvedWorkerSessionTarget } from "./session-target.js"; -import { boundedWorkerError } from "./worker-error.js"; +import { boundedWorkerError, formatWorkerInferenceError } from "./worker-error.js"; type WorkerInferenceStreamEvent = WorkerInferenceEventParams["event"]; export type WorkerInferenceExecutor = import("./inference.js").WorkerInferenceExecutor; @@ -107,31 +108,6 @@ type WorkerInferenceRuntimeDependencies = { recordUsage: (params: WorkerInferenceUsageParams) => void; }; -const ERROR_MESSAGES = { - "model-not-approved": "Model is not approved for this agent.", - "invalid-context": "Inference context is invalid.", - "epoch-mismatch": "Worker run epoch does not match.", - "session-not-attached": "Worker session is not attached.", - "provider-error": "Model provider request failed.", - cancelled: "Inference request was cancelled.", -} as const satisfies Record< - Extract["reason"], - string ->; - -function inferenceError( - reason: Extract["reason"], - usage?: Usage, - message: string = ERROR_MESSAGES[reason], -): WorkerInferenceTerminalOutcome { - return { - type: "error", - reason, - message, - ...(usage ? { usage: structuredClone(usage) } : {}), - }; -} - function copyTool(tool: NonNullable[number]): Tool | undefined { if (!isRecord(tool.parameters) || tool.parameters.type !== "object") { return undefined; @@ -731,6 +707,14 @@ export function createWorkerInferenceExecutor( return inferenceError( event.reason === "aborted" ? "cancelled" : "provider-error", event.error.usage, + event.reason === "aborted" + ? undefined + : formatWorkerInferenceError({ + message: event.error.errorMessage ?? ERROR_MESSAGES["provider-error"], + errorCode: event.error.errorCode, + errorType: event.error.errorType, + errorBody: event.error.errorBody, + }), ); } if (signal.aborted || !params.isCurrent()) { @@ -768,8 +752,12 @@ export function createWorkerInferenceExecutor( } } return inferenceError(signal.aborted ? "cancelled" : "provider-error"); - } catch { - return inferenceError(signal.aborted ? "cancelled" : "provider-error"); + } catch (error) { + return inferenceError( + signal.aborted ? "cancelled" : "provider-error", + undefined, + signal.aborted ? undefined : formatWorkerInferenceError(error), + ); } finally { providerAbort.abort(); } diff --git a/src/gateway/worker-environments/inference-terminal-message.ts b/src/gateway/worker-environments/inference-terminal-message.ts index 6c2a9ff94247..fd740b3f683c 100644 --- a/src/gateway/worker-environments/inference-terminal-message.ts +++ b/src/gateway/worker-environments/inference-terminal-message.ts @@ -1,5 +1,5 @@ import type { WorkerInferenceTerminalOutcome } from "../../../packages/gateway-protocol/src/schema/worker-inference.js"; -import type { AssistantMessage } from "../../llm/types.js"; +import type { AssistantMessage, Usage } from "../../llm/types.js"; import { projectWorkerProviderReplay, type WorkerMessageProjection, @@ -11,6 +11,31 @@ export type WorkerInferenceModelIdentity = { model: string; }; +export const ERROR_MESSAGES = { + "model-not-approved": "Model is not approved for this agent.", + "invalid-context": "Inference context is invalid.", + "epoch-mismatch": "Worker run epoch does not match.", + "session-not-attached": "Worker session is not attached.", + "provider-error": "Model provider request failed.", + cancelled: "Inference request was cancelled.", +} as const satisfies Record< + Extract["reason"], + string +>; + +export function inferenceError( + reason: Extract["reason"], + usage?: Usage, + message: string = ERROR_MESSAGES[reason], +): WorkerInferenceTerminalOutcome { + return { + type: "error", + reason, + message, + ...(usage ? { usage: structuredClone(usage) } : {}), + }; +} + export function projectWorkerInferenceTerminalMessage(params: { message: AssistantMessage; modelIdentity: WorkerInferenceModelIdentity; diff --git a/src/gateway/worker-environments/inference.test.ts b/src/gateway/worker-environments/inference.test.ts index c37374da7de0..6eaefeb8cd8a 100644 --- a/src/gateway/worker-environments/inference.test.ts +++ b/src/gateway/worker-environments/inference.test.ts @@ -5,6 +5,7 @@ import type { WorkerInferenceTerminalOutcome, } from "../../../packages/gateway-protocol/src/schema/worker-inference.js"; import { createDeferred } from "../../../test/helpers/promise.js"; +import { parseApiErrorInfo } from "../../shared/assistant-error-format.js"; import type { WorkerConnectionIdentity } from "./connection-identity.js"; import type { WorkerInferenceStore } from "./inference-store.js"; import { @@ -143,6 +144,34 @@ function makeManager(execute: WorkerInferenceExecutor, store = createMemoryStore } describe("worker inference manager", () => { + it("persists and replays a bounded executor failure without repeating inference", async () => { + const execute = vi.fn(async () => { + throw Object.assign(new Error("Upstream unavailable"), { status: 503, code: "server_error" }); + }); + let stored: WorkerInferenceTerminalOutcome | undefined; + const store = createMemoryStore(); + store.begin = () => (stored ? { kind: "replay", outcome: stored } : { kind: "claimed" }); + store.complete = (input) => (stored = input.outcome); + const instance = makeManager(execute, store); + const sink = createSink(); + accept(instance, { sink: sink.sink }); + await waitForFast(() => expect(terminalFrames(sink.frames)).toHaveLength(1)); + const outcome = terminalFrames(sink.frames)[0]?.payload.outcome; + if (outcome?.type !== "error") { + throw new Error("expected failed inference terminal"); + } + expect(parseApiErrorInfo(outcome.message)).toMatchObject({ + httpCode: "503", + code: "server_error", + message: "Upstream unavailable", + }); + const replay = createSink("replay"); + expect(accept(instance, { sink: replay.sink }).result.status).toBe("replayed"); + expect(terminalFrames(replay.frames)[0]?.payload.outcome).toEqual(outcome); + expect(execute).toHaveBeenCalledOnce(); + await instance.stop(); + }); + it("rejects oversized and concurrent turns", async () => { const store = createMemoryStore(); const execute = vi.fn(async () => ERROR); diff --git a/src/gateway/worker-environments/inference.ts b/src/gateway/worker-environments/inference.ts index 3dc5f73a3bbf..061ebbc2ad49 100644 --- a/src/gateway/worker-environments/inference.ts +++ b/src/gateway/worker-environments/inference.ts @@ -33,6 +33,7 @@ import { serializeWorkerSessionTurnClaim, type WorkerSessionTurnClaim, } from "./placement-record.js"; +import { formatWorkerInferenceError } from "./worker-error.js"; const DEFAULT_REQUEST_MAX_BYTES = WORKER_PROTOCOL_MAX_INFERENCE_PAYLOAD_BYTES; // One active turn plus one provider that ignored abort. This prevents repeated @@ -111,6 +112,7 @@ function trySend( function terminalError( reason: WorkerInferenceErrorReason, outcome?: WorkerInferenceTerminalOutcome, + errorMessage?: string, ): WorkerInferenceTerminalOutcome { const usage = outcome?.type === "done" @@ -138,7 +140,7 @@ function terminalError( return { type: "error", reason, - message, + message: errorMessage ?? message, ...(usage ? { usage } : {}), }; } @@ -364,8 +366,12 @@ export function createWorkerInferenceManager(options: { isCurrent: () => durableFence(entry) === null, ...(config ? { config } : {}), }); - } catch { - outcome = terminalError(entry.abortReason ?? "provider-error"); + } catch (error) { + outcome = terminalError( + entry.abortReason ?? "provider-error", + undefined, + entry.abortReason ? undefined : formatWorkerInferenceError(error), + ); } finish(entry, outcome); }; @@ -378,8 +384,15 @@ export function createWorkerInferenceManager(options: { const operation = runWithGatewayIndependentRootWorkContinuation( () => executeEntry(entry), "worker:dispatch", - ).catch(() => { - finish(entry, terminalError(entry.abortReason ?? "provider-error")); + ).catch((error: unknown) => { + finish( + entry, + terminalError( + entry.abortReason ?? "provider-error", + undefined, + entry.abortReason ? undefined : formatWorkerInferenceError(error), + ), + ); }); operations.set(operation, entry.request.sessionId); void operation.then( diff --git a/src/gateway/worker-environments/worker-error.ts b/src/gateway/worker-environments/worker-error.ts index ee493de01dcb..efac30ff72a1 100644 --- a/src/gateway/worker-environments/worker-error.ts +++ b/src/gateway/worker-environments/worker-error.ts @@ -1,6 +1,11 @@ +import { projectDiagnosticValue } from "@openclaw/ai/diagnostics"; +import { projectProviderError } from "@openclaw/ai/internal/shared"; +import { stableStringify } from "@openclaw/normalization-core"; +import { asOptionalRecord } from "@openclaw/normalization-core/record-coerce"; import { sliceUtf16Safe, truncateUtf16Safe } from "@openclaw/normalization-core/utf16-slice"; import { formatErrorMessage, formatErrorMessageWithCode } from "../../infra/errors.js"; import { redactSensitiveText } from "../../logging/redact.js"; +import { extractErrorHttpStatus, parseApiErrorInfo } from "../../shared/assistant-error-format.js"; function boundWorkerErrorText(text: string, maxChars: number): string { const redacted = redactSensitiveText(text, { mode: "tools" }).replace(/\s+/g, " ").trim(); @@ -25,3 +30,60 @@ export function boundedWorkerError(error: unknown, maxChars = 1_024): string { export function boundedWorkerErrorWithCode(error: unknown, maxChars = 1_024): string { return boundWorkerErrorText(formatErrorMessageWithCode(error), maxChars); } + +/** Preserve provider classification inside the existing bounded inference error text. */ +export function formatWorkerInferenceError(error: unknown): string { + const snapshot = projectDiagnosticValue(error); + const projected = projectProviderError(snapshot); + const record = asOptionalRecord(snapshot); + const response = asOptionalRecord(record?.response); + const status = [ + record?.status, + record?.statusCode, + response?.status, + response?.statusCode, + extractErrorHttpStatus(projected.errorMessage)?.code, + ].find( + (value): value is number => + typeof value === "number" && Number.isInteger(value) && value >= 100 && value <= 599, + ); + const body = + record?.errorBody ?? record?.body ?? response?.body ?? response?.data ?? record?.error; + // The display projection clips bodies; extract classification from the complete + // redacted snapshot so a long body cannot turn billing or context errors into retries. + const info = + parseApiErrorInfo(projected.errorMessage) ?? + parseApiErrorInfo(typeof body === "string" ? body : stableStringify(body)) ?? + parseApiErrorInfo(projected.errorBody); + const code = info?.code ?? projected.errorCode; + const type = projected.errorType ?? info?.type; + if (!code && !type && status === undefined) { + return boundedWorkerError(snapshot, 256); + } + const details = { + ...(code ? { code: boundedWorkerError(code, 64) } : {}), + ...(type ? { type: boundedWorkerError(type, 64) } : {}), + message: boundedWorkerError( + info?.message ?? + extractErrorHttpStatus(projected.errorMessage)?.rest ?? + projected.errorMessage, + 256, + ), + }; + const prefix = status === undefined ? "" : `${status}: `; + const encode = () => `${prefix}${JSON.stringify({ error: details })}`; + let encoded = encode(); + // Clip text before JSON encoding so escaping cannot break the terminal schema + // or destroy the code/type needed by both failover and chat history. + for (const field of ["message", "type", "code"] as const) { + const value = details[field]; + if (encoded.length > 256 && value) { + details[field] = boundWorkerErrorText( + value, + Math.max(0, value.length - (encoded.length - 256)), + ); + encoded = encode(); + } + } + return encoded; +} diff --git a/src/gateway/worker-environments/worker-turn-launcher.test.ts b/src/gateway/worker-environments/worker-turn-launcher.test.ts index d102af09f433..161948b9ff14 100644 --- a/src/gateway/worker-environments/worker-turn-launcher.test.ts +++ b/src/gateway/worker-environments/worker-turn-launcher.test.ts @@ -5,7 +5,10 @@ import { abortAndDrainEmbeddedAgentRun, setActiveEmbeddedRun, } from "../../agents/embedded-agent-runner/runs.js"; -import { installSessionPlacementAdmissionProvider } from "../../agents/session-placement-admission.js"; +import { + installSessionPlacementAdmissionProvider, + resolveSessionPlacementRuntimeOverride, +} from "../../agents/session-placement-admission.js"; import { resolveSessionPlacementForcedTerminalSettlement, resolveSessionPlacementTurnSettlementAssertion, @@ -48,6 +51,40 @@ describe("worker turn launcher local placement", () => { beforeEach(setupWorkerTurnLauncherTest); afterEach(cleanupWorkerTurnLauncherTest); + it.each(["worker-turn", "remote-exec"] as const)( + "uses only the matching %s placement as a runtime default", + (executionMode) => { + const provider = createWorkerSessionTurnPlacementProvider({ + environments: unusedEnvironments(), + placements, + }); + const uninstall = installSessionPlacementAdmissionProvider(provider); + const identity = { sessionId: SESSION_ID, sessionKey: SESSION_KEY, agentId: "main" }; + try { + expect(resolveSessionPlacementRuntimeOverride(identity)).toBeUndefined(); + seedActivePlacement(executionMode); + expect(resolveSessionPlacementRuntimeOverride(identity)).toBe( + executionMode === "worker-turn" ? "openclaw" : undefined, + ); + expect(resolveSessionPlacementRuntimeOverride({ sessionId: SESSION_ID })).toBe( + executionMode === "worker-turn" ? "openclaw" : undefined, + ); + for (const mismatch of [ + { sessionId: "other-session" }, + { sessionKey: "agent:main:other" }, + { agentId: "other-agent" }, + ]) { + expect( + resolveSessionPlacementRuntimeOverride({ ...identity, ...mismatch }), + ).toBeUndefined(); + } + } finally { + uninstall(); + } + expect(resolveSessionPlacementRuntimeOverride(identity)).toBeUndefined(); + }, + ); + it("rejects a transcript target without a session incarnation", () => { expect(() => resolveWorkerTurnTranscriptTarget({ diff --git a/src/gateway/worker-environments/worker-turn-launcher.ts b/src/gateway/worker-environments/worker-turn-launcher.ts index 7325afdab470..56492ec5fa4a 100644 --- a/src/gateway/worker-environments/worker-turn-launcher.ts +++ b/src/gateway/worker-environments/worker-turn-launcher.ts @@ -74,6 +74,16 @@ export function createWorkerSessionTurnPlacementProvider(options: WorkerTurnLaun workspaceDir: string; }): Promise; } = { + resolveRuntimeOverride(identity) { + const placement = options.placements.get(identity.sessionId); + return placement && + placement.state !== "local" && + placement.executionMode === "worker-turn" && + (identity.agentId === undefined || placement.agentId === identity.agentId) && + (identity.sessionKey === undefined || placement.sessionKey === identity.sessionKey) + ? "openclaw" + : undefined; + }, assertCompactionSuccessorAllowed({ currentTarget }) { const placement = options.placements.get(currentTarget.sessionId); // Remote-exec has a local turn claim but still owns remote workspace state. diff --git a/src/infra/sqlite-readonly-location.cancellation.test.ts b/src/infra/sqlite-readonly-location.cancellation.test.ts index a6fe17837dd5..a3e084dbf44c 100644 --- a/src/infra/sqlite-readonly-location.cancellation.test.ts +++ b/src/infra/sqlite-readonly-location.cancellation.test.ts @@ -20,7 +20,6 @@ const processMocks = vi.hoisted(() => ({ vi.mock("node:child_process", async (importOriginal) => { const { promisify } = await import("node:util"); const actual = await importOriginal(); - const { promisify } = await import("node:util"); processMocks.execFile.mockImplementation(actual.execFile); Object.defineProperty( processMocks.execFile, diff --git a/src/infra/sqlite-readonly-worker.test.ts b/src/infra/sqlite-readonly-worker.test.ts index 6f53bc054eee..14fe9e3a386b 100644 --- a/src/infra/sqlite-readonly-worker.test.ts +++ b/src/infra/sqlite-readonly-worker.test.ts @@ -48,7 +48,6 @@ vi.mock("../logging/subsystem.js", async (importOriginal) => { vi.mock("node:child_process", async (importOriginal) => { const { promisify } = await import("node:util"); const actual = await importOriginal(); - const { promisify } = await import("node:util"); const execFileSpy = vi.fn(actual.execFile); Object.defineProperty( execFileSpy, diff --git a/src/infra/sqlite-snapshot-staging.cancellation.test.ts b/src/infra/sqlite-snapshot-staging.cancellation.test.ts index 0d4747080b4a..03e9ab1d8544 100644 --- a/src/infra/sqlite-snapshot-staging.cancellation.test.ts +++ b/src/infra/sqlite-snapshot-staging.cancellation.test.ts @@ -22,7 +22,6 @@ const processMocks = vi.hoisted(() => ({ vi.mock("node:child_process", async (importOriginal) => { const { promisify } = await import("node:util"); const actual = await importOriginal(); - const { promisify } = await import("node:util"); processMocks.execFile.mockImplementation(actual.execFile); Object.defineProperty( processMocks.execFile, diff --git a/src/worker/inference-stream.runtime.test.ts b/src/worker/inference-stream.runtime.test.ts index ce1d4fc94fd6..a9590f621fe8 100644 --- a/src/worker/inference-stream.runtime.test.ts +++ b/src/worker/inference-stream.runtime.test.ts @@ -8,6 +8,7 @@ import { } from "../../packages/gateway-protocol/src/client-info.js"; import { WORKER_PROTOCOL_FEATURES, + WORKER_PROTOCOL_MAX_PAYLOAD_BYTES, WORKER_RPC_SET_VERSION, } from "../../packages/gateway-protocol/src/schema/worker-admission.js"; import { @@ -22,7 +23,10 @@ import type { Usage } from "../llm/types.js"; import { createWorkerInferenceStreamAdapter } from "./inference-stream.runtime.js"; import { createWorkerImageHistory } from "./replay-images.test-support.js"; import { fitWorkerReplayImages } from "./replay-message-window.js"; -import { WORKER_PROVIDER_REPLAY_LOCAL_RETRY_MESSAGE } from "./transcript-message.js"; +import { + isWorkerTranscriptMessageFrameSafe, + WORKER_PROVIDER_REPLAY_LOCAL_RETRY_MESSAGE, +} from "./transcript-message.js"; import { createWorkerConnection } from "./worker-connection.js"; import { WorkerInferenceProxyClient } from "./worker-rpc-clients.js"; @@ -138,6 +142,66 @@ function inferenceRequest(context: WorkerInferenceContext): WorkerInferenceStart }; } +it("preserves the provider failure and usage after an oversized partial response", async () => { + const fixture = createAdapterFixture(); + fixture.start.mockImplementation(async (request, handlers) => { + const events: WorkerInferenceEventParams["event"][] = [ + { type: "text_start", contentIndex: 0 }, + { type: "text_delta", contentIndex: 0, delta: "x".repeat(WORKER_PROTOCOL_MAX_PAYLOAD_BYTES) }, + ]; + events.forEach((event, index) => handlers?.onEvent?.({ ...request, seq: index + 1, event })); + return { type: "error", reason: "provider-error", message: "429: rate limit exceeded", usage }; + }); + + try { + const result = await fixture + .stream({ modelRef, context: { messages: [] }, options: {} }) + .result(); + expect(result).toMatchObject({ + stopReason: "error", + errorMessage: "429: rate limit exceeded", + usage, + content: [], + }); + expect(isWorkerTranscriptMessageFrameSafe(result)).toBe(true); + } finally { + fixture.client.dispose(); + } +}); + +it("rejects an oversized successful reply without truncating it or losing its usage", async () => { + const fixture = createAdapterFixture(); + const reply: Extract["message"] = { + role: "assistant", + content: [{ type: "text", text: "x".repeat(WORKER_PROTOCOL_MAX_PAYLOAD_BYTES) }], + api: "openai-responses", + provider: modelRef.provider, + model: modelRef.model, + stopReason: "stop", + usage, + timestamp: 1, + }; + fixture.start.mockResolvedValue({ type: "done", message: reply }); + + try { + const result = await fixture + .stream({ modelRef, context: { messages: [] }, options: {} }) + .result(); + expect(result).toMatchObject({ + stopReason: "error", + errorMessage: "Worker inference result exceeds the transcript message limit.", + usage, + content: [], + }); + expect(isWorkerTranscriptMessageFrameSafe(result)).toBe(true); + expect(reply.content).toEqual([ + { type: "text", text: "x".repeat(WORKER_PROTOCOL_MAX_PAYLOAD_BYTES) }, + ]); + } finally { + fixture.client.dispose(); + } +}); + it.each([false, true])( "bounds screenshot serialization work while preserving the fitted request (oversized: %s)", async (oversized) => { diff --git a/src/worker/inference-stream.runtime.ts b/src/worker/inference-stream.runtime.ts index 3b7cd028c500..d015819ff18f 100644 --- a/src/worker/inference-stream.runtime.ts +++ b/src/worker/inference-stream.runtime.ts @@ -4,6 +4,7 @@ import { parseTerminalToolCallArguments, type ToolArgumentPreviewSchedule, } from "@openclaw/ai/internal/runtime"; +import { truncateUtf16Safe } from "@openclaw/normalization-core/utf16-slice"; import { WORKER_PROTOCOL_MAX_IDENTIFIER_LENGTH } from "../../packages/gateway-protocol/src/schema/worker-admission.js"; import type { WorkerInferenceContext, @@ -236,8 +237,16 @@ function transcriptSafeErrorMessage( return message; } const replacement = emptyAssistantMessage(modelRef); + replacement.api = message.api; + replacement.provider = message.provider; + replacement.model = message.model; + replacement.timestamp = message.timestamp; replacement.stopReason = message.stopReason === "aborted" ? "aborted" : "error"; - replacement.errorMessage = "Worker inference result exceeds the transcript message limit."; + replacement.errorMessage = truncateUtf16Safe( + message.errorMessage ?? "Worker inference result exceeds the transcript message limit.", + 256, + ); + replacement.usage = structuredClone(message.usage); return replacement; } @@ -405,6 +414,7 @@ export function createWorkerInferenceStreamAdapter( const message = emptyAssistantMessage(adapter.modelRef); message.stopReason = "error"; message.errorMessage = "Worker inference result exceeds the transcript message limit."; + message.usage = structuredClone(outcome.message.usage); stream.push({ type: "error", reason: "error", error: message }); stream.end(); return;