fix: preserve model failure causes and worker fallback runtime (#152028)

* fix(workers): preserve failure causes and fallback runtime

Keep implicit fallback candidates on the runtime owned by their worker placement. Preserve bounded, redacted provider status and classification, retain failure and usage after oversized partial responses, and show actionable error copy in chat history.

* fix(chat): preserve recorded context overflow errors

* fix: distinguish provider failures in saved chat history

* test: reconcile runtime and SQLite fixtures after main updates
This commit is contained in:
Peter Steinberger 2026-09-18 13:36:53 -07:00 • committed by GitHub
parent 8a7af996c8
commit 3fceb86048
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
38 changed files with 1046 additions and 396 deletions

View file

@ -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

View file

@ -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

View file

@ -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";

View file

@ -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]`);
});

View file

@ -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]",

View file

@ -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);

View file

@ -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({

View file

@ -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,

View file

@ -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;
}

View file

@ -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();
});

View file

@ -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<T extends EmbeddedAgentRunResult> = {
resolveContextEngineHost?: (
provider: string,
model: string,
agentHarnessRuntimeOverride: string | undefined,
) => ContextEngineHostSupport | undefined;
};
behavior: RunEntryBehavior;
@ -172,6 +176,22 @@ async function runEmbeddedAgentEntryInternal<T extends EmbeddedAgentRunResult>(
): Promise<EmbeddedAgentRunEntryResult<T>> {
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<T extends EmbeddedAgentRunResult>(
...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<T extends EmbeddedAgentRunResult>(
const resolvedHost = params.harness.resolveContextEngineHost?.(
candidate.provider,
candidate.model,
agentHarnessRuntimeOverride,
);
const host =
resolvedHost ??
@ -410,6 +431,7 @@ async function runEmbeddedAgentEntryInternal<T extends EmbeddedAgentRunResult>(
};
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

View file

@ -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 <T>(_claim: unknown, run: () => Promise<T>) => 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",
});
}
}
},
);
});

View file

@ -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;
}

View file

@ -47,6 +47,7 @@ type SessionPlacementSandboxParams = {
};
export type SessionPlacementAdmissionProvider = {
resolveRuntimeOverride?: (identity: Omit<LocalTurnPlacementClaim, "runId">) => 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<LocalTurnPlacementClaim, "runId">,
): 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;

View file

@ -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,

View file

@ -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> | void;
run: (
provider: string,
model: string,
options: {
allowTransientCooldownProbe?: boolean;
isFinalFallbackAttempt?: boolean;
modelRoutingProvenance: ModelFallbackAttemptProvenance;
},
) => Promise<EmbeddedAgentRunResult>;
};
export function createMemoryRunEntryMockImplementation(deps: {
runWithModelFallback: (params: ModelFallbackParams) => Promise<unknown>;
ensureSelectedAgentHarnessPlugin: typeof ensureSelectedAgentHarnessPlugin;
}) {
return async (params: Parameters<typeof runEmbeddedAgentEntry<EmbeddedAgentRunResult>>[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<ModelFallbackParams["run"]>[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,
};
};
}

View file

@ -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> | void;
run: (
provider: string,
model: string,
options: {
allowTransientCooldownProbe?: boolean;
isFinalFallbackAttempt?: boolean;
modelRoutingProvenance: ModelFallbackAttemptProvenance;
},
) => Promise<EmbeddedAgentRunResult>;
};
function modelRoutingProvenance(
requestedProvider: string,
requestedModel: string,
@ -453,71 +427,12 @@ describe("runMemoryFlushIfNeeded", () => {
model,
attempts: [],
}));
runEmbeddedAgentEntryMock
.mockReset()
.mockImplementation(
async (params: Parameters<typeof runEmbeddedAgentEntry<EmbeddedAgentRunResult>>[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<ModelFallbackParams["run"]>[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,

View file

@ -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,

View file

@ -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,
};
};
}

View file

@ -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<CronCompletedPromptRun> => {
// 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({

View file

@ -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<string, unknown>): boolean {
return (
isContextOverflowErrorSignal(message.errorCode) ||
isContextOverflowErrorSignal(message.errorType) ||
isContextOverflowErrorSignal(message.errorMessage)
);
}
function getAssistantErrorFallbackText(message: Record<string, unknown>): string {
return (
formatProviderRefusalText(message) ??
@ -162,10 +138,8 @@ function getAssistantErrorFallbackText(message: Record<string, unknown>): 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.

View file

@ -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);
});

View file

@ -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<string, unknown>;
content: Array<Record<string, unknown>>;
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])(

View file

@ -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<ChatEvent, { state: "delta" }> =>

View file

@ -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");

View file

@ -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();

View file

@ -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<WorkerInferenceTerminalOutcome, { type: "error" }>["reason"],
string
>;
function inferenceError(
reason: Extract<WorkerInferenceTerminalOutcome, { type: "error" }>["reason"],
usage?: Usage,
message: string = ERROR_MESSAGES[reason],
): WorkerInferenceTerminalOutcome {
return {
type: "error",
reason,
message,
...(usage ? { usage: structuredClone(usage) } : {}),
};
}
function copyTool(tool: NonNullable<WorkerInferenceContext["tools"]>[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();
}

View file

@ -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<WorkerInferenceTerminalOutcome, { type: "error" }>["reason"],
string
>;
export function inferenceError(
reason: Extract<WorkerInferenceTerminalOutcome, { type: "error" }>["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;

View file

@ -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<WorkerInferenceExecutor>(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<WorkerInferenceExecutor>(async () => ERROR);

View file

@ -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(

View file

@ -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;
}

View file

@ -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({

View file

@ -74,6 +74,16 @@ export function createWorkerSessionTurnPlacementProvider(options: WorkerTurnLaun
workspaceDir: string;
}): Promise<SandboxContext | null>;
} = {
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.

View file

@ -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<typeof import("node:child_process")>();
const { promisify } = await import("node:util");
processMocks.execFile.mockImplementation(actual.execFile);
Object.defineProperty(
processMocks.execFile,

View file

@ -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<typeof import("node:child_process")>();
const { promisify } = await import("node:util");
const execFileSpy = vi.fn(actual.execFile);
Object.defineProperty(
execFileSpy,

View file

@ -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<typeof import("node:child_process")>();
const { promisify } = await import("node:util");
processMocks.execFile.mockImplementation(actual.execFile);
Object.defineProperty(
processMocks.execFile,

View file

@ -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<WorkerInferenceTerminalOutcome, { type: "done" }>["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) => {

View file

@ -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;