mirror of
https://github.com/openclaw/openclaw.git
synced 2026-10-03 09:39:25 +00:00
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:
parent
8a7af996c8
commit
3fceb86048
38 changed files with 1046 additions and 396 deletions
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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";
|
||||
|
|
|
|||
|
|
@ -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]`);
|
||||
});
|
||||
|
||||
|
|
|
|||
|
|
@ -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]",
|
||||
|
|
|
|||
|
|
@ -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);
|
||||
|
|
|
|||
|
|
@ -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({
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
126
src/agents/embedded-agent-runner/run-entry.test-harness.ts
Normal file
126
src/agents/embedded-agent-runner/run-entry.test-harness.ts
Normal 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;
|
||||
}
|
||||
|
|
@ -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();
|
||||
});
|
||||
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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",
|
||||
});
|
||||
}
|
||||
}
|
||||
},
|
||||
);
|
||||
});
|
||||
|
|
@ -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;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
98
src/auto-reply/reply/agent-runner-memory.test-support.ts
Normal file
98
src/auto-reply/reply/agent-runner-memory.test-support.ts
Normal 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,
|
||||
};
|
||||
};
|
||||
}
|
||||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
41
src/cron/isolated-agent/run-candidate-runtime.ts
Normal file
41
src/cron/isolated-agent/run-candidate-runtime.ts
Normal 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,
|
||||
};
|
||||
};
|
||||
}
|
||||
|
|
@ -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({
|
||||
|
|
|
|||
|
|
@ -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.
|
||||
|
|
|
|||
|
|
@ -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);
|
||||
});
|
||||
|
|
|
|||
|
|
@ -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])(
|
||||
|
|
|
|||
|
|
@ -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" }> =>
|
||||
|
|
|
|||
|
|
@ -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");
|
||||
|
|
|
|||
|
|
@ -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();
|
||||
|
|
|
|||
|
|
@ -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();
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
|
|
|
|||
|
|
@ -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);
|
||||
|
|
|
|||
|
|
@ -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(
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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({
|
||||
|
|
|
|||
|
|
@ -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.
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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) => {
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue