fix(agents): keep active turns working through configuration reloads (#154462)

Keep active turns working when configuration or credentials reload during model failover. Retain catalog facts on the matching open admitted-generation lease instead of borrowing a replacement account’s catalog, including same-URL account changes. Preserve native and configured physical-model facts, and retain cron’s carried catalog when hydration has no selected model.

Fixes #154456. Original contribution by @TreyLawrence; maintainer review and corrections requested by @VACInc.

Validated at source head 96c07bc308dc7e78c337382b67062ab7ae4b12cb: required CI/security gates passed; 43 focused catalog/lease/cron tests, 27 restart tests, and the real Gateway reload/failover case passed. ClawSweeper revision 17 has no actionable or security findings. Historical load-dependent fixture and scoped CodeQL limitations remain documented in the PR.

Co-authored-by: VACInc <3279061+VACInc@users.noreply.github.com>
Co-authored-by: Vito Cappello <hixvac@gmail.com>
This commit is contained in:
Trey Lawrence 2026-09-27 18:50:52 -04:00 • committed by GitHub
parent a48bd1c4f3
commit cd11c4bc0d
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
7 changed files with 1009 additions and 41 deletions

View file

@ -2,10 +2,15 @@
import { beforeEach, describe, expect, it, vi } from "vitest";
import { createPluginMetadataSnapshot } from "../config/plugin-auto-enable.test-helpers.js";
import type { OpenClawConfig } from "../config/types.openclaw.js";
import { resolveLegacyInheritedAuthDir } from "./legacy-inherited-auth-dir.js";
import type { ModelCatalogEntry, ModelCatalogSnapshot } from "./model-catalog.types.js";
import { PreparedModelCatalogConfigReplacedError } from "./prepared-model-catalog.errors.js";
import { setPreparedModelFullCatalogAuth } from "./prepared-model-runtime-auth.js";
import { PreparedModelRuntimeOwnerNotPublishedError } from "./prepared-model-runtime.errors.js";
import { withPreparedModelRuntimePluginGenerationScope } from "./prepared-model-runtime-generation-scope.js";
import { capturePreparedModelRuntimeCatalog } from "./prepared-model-runtime.capture.js";
import {
PreparedModelRuntimeOwnerNotPublishedError,
PreparedModelRuntimePublicationSupersededError,
} from "./prepared-model-runtime.errors.js";
import type {
PreparedModelRuntimeInput,
PreparedModelRuntimeSnapshot,
@ -49,6 +54,7 @@ vi.mock("./prepared-model-runtime.scoped-catalog.js", () => ({
function owner(config: OpenClawConfig, entries: ModelCatalogEntry[]): PreparedModelRuntimeSnapshot {
return {
agentDir: "/tmp/model-catalog-passive-test",
inheritedAuthDir: resolveLegacyInheritedAuthDir(config),
activeProjectKeys: [],
catalogOwner: undefined,
config,
@ -70,13 +76,25 @@ function owner(config: OpenClawConfig, entries: ModelCatalogEntry[]): PreparedMo
};
}
const entry: ModelCatalogEntry = {
function withAdmitted<T>(snapshot: PreparedModelRuntimeSnapshot, run: () => T, active = true): T {
return withPreparedModelRuntimePluginGenerationScope(
{
pluginMetadataSnapshot: snapshot.metadataSnapshot,
inlineProviderModels: [],
configuredCatalogEntries: snapshot.modelCatalog.entries,
},
run,
() => (active ? snapshot : undefined),
);
}
const entry = {
provider: "acme",
id: "selected",
name: "Selected",
api: "openai-responses",
baseUrl: "https://provider.invalid/v1",
};
} satisfies ModelCatalogEntry;
describe("loadProviderScopedThinkingCatalog", () => {
beforeEach(() => {
@ -187,17 +205,270 @@ describe("loadProviderScopedThinkingCatalog", () => {
expect(scopedCatalogMock).not.toHaveBeenCalled();
});
it("rejects an owner whose configuration was replaced", async () => {
it("omits replacement facts without a matching admitted generation", async () => {
const config = { skills: { entries: { marker: { enabled: true } } } };
const replaced = { skills: { entries: { marker: { enabled: false } } } };
publishedSnapshotMock.mockReturnValue(owner(replaced, [{ ...entry, reasoning: true }]));
const { loadProviderScopedThinkingCatalog } = await import("./prepared-model-catalog.js");
await expect(
loadProviderScopedThinkingCatalog({ config, provider: entry.provider, model: entry.id }),
).rejects.toBeInstanceOf(PreparedModelCatalogConfigReplacedError);
).resolves.toEqual([]);
expect(scopedCatalogMock).not.toHaveBeenCalled();
});
it.each([
{ name: "same-account refresh", replacementKey: "fixture-account-a", native: false },
{ name: "same-route account switch", replacementKey: "fixture-account-b", native: false },
{ name: "route-free native observation", replacementKey: "fixture-account-b", native: true },
])("retains admitted capabilities across $name", async ({ replacementKey, native }) => {
const config: OpenClawConfig = {
models: {
providers: { acme: { baseUrl: entry.baseUrl, apiKey: "fixture-account-a", models: [] } },
},
};
const replaced: OpenClawConfig = {
...config,
models: {
providers: { acme: { baseUrl: entry.baseUrl, apiKey: replacementKey, models: [] } },
},
skills: { entries: { marker: { enabled: false } } },
};
const admittedEntry: ModelCatalogEntry = native
? {
provider: entry.provider,
id: entry.id,
name: entry.name,
nativeRuntime: "native-one",
reasoning: true,
}
: { ...entry, reasoning: true, input: ["text", "image"] };
const completed = { entries: [admittedEntry], routeVariants: [admittedEntry] };
const readFullModelCatalog = vi.fn(() => completed);
const source = { ...owner(config, [entry]), readFullModelCatalog };
const admitted = capturePreparedModelRuntimeCatalog(source, source);
readFullModelCatalog.mockImplementation(() => {
throw new Error("retired owner");
});
publishedSnapshotMock.mockReturnValue(
owner(replaced, [{ ...entry, reasoning: false, input: ["text"] }]),
);
const { loadProviderScopedThinkingCatalog } = await import("./prepared-model-catalog.js");
const result = await withAdmitted(admitted, () =>
loadProviderScopedThinkingCatalog({
config,
agentDir: admitted.agentDir,
provider: entry.provider,
model: entry.id,
...(native
? { agentRuntime: "native-one" }
: { requiredInputRoute: { api: entry.api, baseUrl: entry.baseUrl } }),
}),
);
expect(result).toEqual([admittedEntry]);
expect(preparedSnapshotMock).not.toHaveBeenCalled();
expect(augmentCatalogMock).not.toHaveBeenCalled();
expect(readFullModelCatalog).toHaveBeenCalledOnce();
});
it.each([
{
label: "same-id physical route",
nativeId: entry.id,
runtime: "openclaw",
publishedSelected: false,
},
{
label: "distinct configured model",
nativeId: "native-only",
runtime: "openclaw",
publishedSelected: false,
},
{ label: "native route", nativeId: entry.id, runtime: "native-one", publishedSelected: false },
{
label: "published physical facts",
nativeId: "native-only",
runtime: "openclaw",
publishedSelected: true,
},
])(
"retains native and physical capture facts: $label",
async ({ nativeId, runtime, publishedSelected }) => {
const config = {};
const configured: ModelCatalogEntry = publishedSelected
? entry
: { ...entry, reasoning: false, input: ["text", "image"] };
const discovered: ModelCatalogEntry = {
...entry,
id: publishedSelected ? entry.id : "discovered",
name: "Discovered",
reasoning: true,
input: ["text", "image"],
};
const native: ModelCatalogEntry = {
provider: entry.provider,
id: nativeId,
name: entry.name,
nativeRuntime: "native-one",
reasoning: true,
input: ["text"],
};
let current = true;
const readFullModelCatalog = vi.fn(() => ({
entries: [native, discovered],
routeVariants: [native, discovered],
}));
const loadNativeModelCatalog = vi.fn(async () => {
throw new Error("A retired owner cannot supply new native observations");
});
const source = {
...owner(config, [configured]),
isCurrent: () => current,
readFullModelCatalog,
loadNativeModelCatalog,
};
const admitted = capturePreparedModelRuntimeCatalog(source, source);
current = false;
const { loadProviderScopedThinkingCatalog } = await import("./prepared-model-catalog.js");
const result = await withAdmitted(admitted, () =>
loadProviderScopedThinkingCatalog({
config,
agentDir: admitted.agentDir,
provider: entry.provider,
model: entry.id,
agentRuntime: runtime,
...(runtime === "openclaw"
? { requiredInputRoute: { api: entry.api, baseUrl: entry.baseUrl } }
: {}),
}),
);
const physical = publishedSelected ? discovered : configured;
expect(result).toContainEqual(runtime === "openclaw" ? physical : native);
expect(result).toContainEqual(discovered);
expect(readFullModelCatalog).toHaveBeenCalledOnce();
expect(loadNativeModelCatalog).not.toHaveBeenCalled();
},
);
it.each(["ready", "retired", "failure"] as const)(
"preserves admitted native observations: %s",
async (outcome) => {
const config = {};
const native: ModelCatalogEntry = {
provider: entry.provider,
id: entry.id,
name: entry.name,
nativeRuntime: "native-one",
reasoning: true,
};
let current = true;
const failure = new Error("native observation unavailable");
const loadNativeModelCatalog = vi.fn(async () => {
if (outcome === "retired") {
current = false;
throw new PreparedModelRuntimePublicationSupersededError("retired during observation");
}
if (outcome === "failure") {
throw failure;
}
return { entries: [native], routeVariants: [native] };
});
const admitted = {
...owner(config, outcome === "retired" ? [native] : [{ ...entry, reasoning: false }]),
isCurrent: () => current,
loadNativeModelCatalog,
};
const { loadProviderScopedThinkingCatalog } = await import("./prepared-model-catalog.js");
const result = withAdmitted(admitted, () =>
loadProviderScopedThinkingCatalog({
config,
agentDir: admitted.agentDir,
provider: entry.provider,
model: entry.id,
agentRuntime: "native-one",
}),
);
if (outcome === "failure") {
await expect(result).rejects.toBe(failure);
} else {
await expect(result).resolves.toEqual([native]);
}
expect(loadNativeModelCatalog).toHaveBeenCalledOnce();
expect(preparedSnapshotMock).not.toHaveBeenCalled();
},
);
it.each([
{
name: "a capability change",
configured: [entry],
firstEntries: [{ ...entry, reasoning: true }],
secondEntries: [{ ...entry, reasoning: false }],
},
{
name: "empty to populated inventory",
configured: [],
firstEntries: [],
secondEntries: [entry],
},
{
name: "populated to empty inventory",
configured: [],
firstEntries: [entry],
secondEntries: [],
},
{
name: "configured facts followed by empty inventory",
configured: [entry],
firstEntries: undefined,
secondEntries: [],
},
])("keeps earlier captures across $name", async ({ configured, firstEntries, secondEntries }) => {
const config = {};
const readFullModelCatalog = vi
.fn<() => ModelCatalogSnapshot | undefined>()
.mockReturnValue(
firstEntries ? { entries: firstEntries, routeVariants: firstEntries } : undefined,
);
const source = { ...owner(config, configured), readFullModelCatalog };
const first = capturePreparedModelRuntimeCatalog(source, source);
readFullModelCatalog.mockReturnValue({ entries: secondEntries, routeVariants: secondEntries });
const second = capturePreparedModelRuntimeCatalog(source, source);
const { loadProviderScopedThinkingCatalog } = await import("./prepared-model-catalog.js");
const read = () =>
loadProviderScopedThinkingCatalog({
config,
agentDir: source.agentDir,
provider: entry.provider,
model: entry.id,
});
await expect(withAdmitted(first, read)).resolves.toEqual(firstEntries ?? configured);
await expect(withAdmitted(second, read)).resolves.toEqual(secondEntries);
});
it.each(["closed lease", "other agent", "other workspace"])(
"cannot borrow facts from a %s",
async (mismatch) => {
const config = {};
const replaced = { skills: { entries: { marker: { enabled: false } } } };
const admitted = owner(config, [{ ...entry, reasoning: true }]);
publishedSnapshotMock.mockReturnValue(owner(replaced, [{ ...entry, reasoning: false }]));
const { loadProviderScopedThinkingCatalog } = await import("./prepared-model-catalog.js");
const result = await withAdmitted(
admitted,
() =>
loadProviderScopedThinkingCatalog({
config,
agentDir: mismatch === "other agent" ? "/tmp/other-agent" : admitted.agentDir,
...(mismatch === "other workspace" ? { workspaceDir: "/tmp/other-workspace" } : {}),
provider: entry.provider,
model: entry.id,
}),
mismatch !== "closed lease",
);
expect(result).toEqual([]);
},
);
it("keeps native harness observations available without a published owner", async () => {
const nativeEntry = { ...entry, nativeRuntime: "test-harness", reasoning: true };
augmentCatalogMock.mockResolvedValue({ entries: [nativeEntry], routeVariants: [nativeEntry] });

View file

@ -20,6 +20,12 @@ import {
loadPreparedModelRuntimeAuth,
bindPreparedModelRuntimeAuth,
} from "./prepared-model-runtime-auth.js";
import {
getPreparedModelRuntimeBorrowedSnapshot,
getPreparedModelRuntimePluginGeneration,
} from "./prepared-model-runtime-generation-scope.js";
import { readCapturedPreparedModelRuntimeCatalog } from "./prepared-model-runtime.capture.js";
import { PreparedModelRuntimePublicationSupersededError } from "./prepared-model-runtime.errors.js";
import { isPreparedModelCatalogFull } from "./prepared-model-runtime.full-catalog.js";
import {
acquireAgentRunPreparedModelRuntime,
@ -384,6 +390,25 @@ async function loadScopedReadOnlyModelCatalog(
);
}
/** Reads only the generation retained by this exact, still-open turn. */
function resolveAdmittedModelCatalogOwner(params: LoadPreparedModelCatalogParams) {
const generation = getPreparedModelRuntimePluginGeneration();
const owner = generation && getPreparedModelRuntimeBorrowedSnapshot(generation);
if (!generation || !owner || owner.metadataSnapshot !== generation.pluginMetadataSnapshot) {
return undefined;
}
const { full, activationFull } = resolveInputs(params);
const matches = [full, activationFull].some(
(input) =>
preparedModelRuntimeConfigsMatch(owner.config, input.config) &&
owner.agentId === input.agentId &&
owner.agentDir === input.agentDir &&
owner.inheritedAuthDir === input.inheritedAuthDir &&
owner.workspaceDir === input.workspaceDir,
);
return matches ? { generation, owner } : undefined;
}
/**
* Missing turn-path capabilities do not authorize another inventory, even without a published
* owner. Native harness observations keep their existing owner.
@ -401,35 +426,71 @@ export async function loadProviderScopedThinkingCatalog(params: {
requiredInputRoute?: Pick<ModelCatalogEntry, "api" | "baseUrl">;
}): Promise<ModelCatalogEntry[]> {
const request = { ...params, readOnly: true };
const publishedOwner = getPreparedModelCatalogOwnerSnapshot(request);
const owner = (await resolveReadOnlyPublishedModelCatalogOwner(request, "exact"))?.snapshot;
const admitted = resolveAdmittedModelCatalogOwner(request);
let snapshot: ModelCatalogSnapshot;
if (owner?.loadNativeModelCatalog && params.agentRuntime && params.agentRuntime !== "openclaw") {
snapshot = await owner.loadNativeModelCatalog({
provider: params.provider,
modelId: params.model,
runtime: params.agentRuntime,
});
if (admitted) {
// Replacement inventory may belong to another account on the very same URL.
// A retained turn reads its captured facts without invoking retired owner callbacks.
snapshot =
readCapturedPreparedModelRuntimeCatalog(admitted.owner) ?? admitted.owner.modelCatalog;
if (
params.agentRuntime &&
params.agentRuntime !== "openclaw" &&
admitted.owner.loadNativeModelCatalog &&
admitted.owner.isCurrent()
) {
try {
snapshot = await admitted.owner.loadNativeModelCatalog({
provider: params.provider,
modelId: params.model,
runtime: params.agentRuntime,
});
} catch (error) {
if (
!(error instanceof PreparedModelRuntimePublicationSupersededError) ||
admitted.owner.isCurrent()
) {
throw error;
}
// Retirement during a native observation leaves only this turn's captured facts.
}
}
} else {
const catalog = owner
? (publishedOwner ? await materializeRequestedModelCatalog(owner, true, undefined) : owner)
.modelCatalog
: { entries: [], routeVariants: [] };
const agentId = params.agentId ?? resolveAmbientOwnerAgentId(params.config);
const { augmentModelCatalogWithAgentHarness } = await import("./harness/model-catalog.js");
snapshot = await augmentModelCatalogWithAgentHarness({
cfg: params.config,
agentId,
agentDir: params.agentDir ?? resolveAgentDir(params.config, agentId),
workspaceDir:
params.workspaceDir ??
resolveAgentWorkspaceDir(params.config, agentId) ??
resolveDefaultAgentWorkspaceDir(),
defaultProvider: params.provider,
defaultModel: `${params.provider}/${params.model}`,
agentRuntime: params.agentRuntime,
snapshot: catalog,
});
const owner = (await resolveReadOnlyPublishedModelCatalogOwner(request, "published"))?.snapshot;
if (owner && !preparedModelRuntimeConfigsMatch(owner.config, params.config)) {
// A caller without a matching admitted generation cannot borrow replacement facts.
return [];
}
if (
owner?.loadNativeModelCatalog &&
params.agentRuntime &&
params.agentRuntime !== "openclaw"
) {
snapshot = await owner.loadNativeModelCatalog({
provider: params.provider,
modelId: params.model,
runtime: params.agentRuntime,
});
} else {
const catalog = owner
? (await materializeRequestedModelCatalog(owner, true, undefined)).modelCatalog
: { entries: [], routeVariants: [] };
const agentId = params.agentId ?? resolveAmbientOwnerAgentId(params.config);
const { augmentModelCatalogWithAgentHarness } = await import("./harness/model-catalog.js");
snapshot = await augmentModelCatalogWithAgentHarness({
cfg: params.config,
agentId,
agentDir: params.agentDir ?? resolveAgentDir(params.config, agentId),
workspaceDir:
params.workspaceDir ??
resolveAgentWorkspaceDir(params.config, agentId) ??
resolveDefaultAgentWorkspaceDir(),
defaultProvider: params.provider,
defaultModel: `${params.provider}/${params.model}`,
agentRuntime: params.agentRuntime,
snapshot: catalog,
});
}
}
let entries = snapshot.entries;
if (params.agentRuntime) {
@ -446,6 +507,9 @@ export async function loadProviderScopedThinkingCatalog(params: {
}
}
entries = normalizeThinkingCatalogProviders(entries);
if (admitted && getPreparedModelRuntimeBorrowedSnapshot(admitted.generation) !== admitted.owner) {
return [];
}
if (params.requiredInputRoute !== undefined) {
const entry = findModelInCatalog(entries, params.provider, params.model);
if (

View file

@ -10,11 +10,21 @@ const catalogCaptures = new WeakMap<
{
models: ReadonlyMap<string, readonly Model[]> | undefined;
catalog: ModelCatalogSnapshot | undefined;
nativeSnapshot: PreparedModelRuntimeSnapshot;
admittedCatalog: ModelCatalogSnapshot;
capturedSnapshot: PreparedModelRuntimeSnapshot;
memo: Map<string, Promise<Model>>;
}
>();
const admittedCatalogs = new WeakMap<PreparedModelRuntimeSnapshot, ModelCatalogSnapshot>();
/** Passive inventory captured at admission; the caller must hold the matching turn lease. */
export function readCapturedPreparedModelRuntimeCatalog(
snapshot: PreparedModelRuntimeSnapshot,
): ModelCatalogSnapshot | undefined {
return admittedCatalogs.get(snapshot);
}
/** Captures published executable and native model facts without changing any open lease. */
export function capturePreparedModelRuntimeCatalog(
snapshot: PreparedModelRuntimeSnapshot,
@ -28,20 +38,37 @@ export function capturePreparedModelRuntimeCatalog(
catalog &&
(catalog.entries.some((entry) => entry.nativeRuntime) ||
catalog.routeVariants.some((entry) => entry.nativeRuntime));
const hasCatalogFacts =
catalog &&
(catalog.entries.length > 0 ||
catalog.routeVariants.length > 0 ||
snapshot.modelCatalog.entries.length > 0 ||
snapshot.modelCatalog.routeVariants.length > 0);
cached = {
models,
catalog,
nativeSnapshot: nativeCatalog
admittedCatalog: nativeCatalog
? mergePreparedNativeCatalog(catalog, {
...catalog,
entries: [...catalog.entries, ...snapshot.modelCatalog.entries],
routeVariants: [...catalog.routeVariants, ...snapshot.modelCatalog.routeVariants],
})
: (catalog ?? snapshot.modelCatalog),
// Preserve bare configured owners when there are no capability rows to capture.
capturedSnapshot: hasCatalogFacts
? Object.freeze({
...snapshot,
modelCatalog: mergePreparedNativeCatalog(catalog, snapshot.modelCatalog),
...(nativeCatalog
? { modelCatalog: mergePreparedNativeCatalog(catalog, snapshot.modelCatalog) }
: {}),
})
: snapshot,
memo: cached && cached.models === models ? cached.memo : new Map(),
};
catalogCaptures.set(snapshot, cached);
}
const capturedNative = cached.nativeSnapshot;
const capturedNative = cached.capturedSnapshot;
admittedCatalogs.set(capturedNative, cached.admittedCatalog);
if (!models?.size) {
if (capturedNative !== snapshot) {
copyPreparedModelRuntimeAuthBindings(snapshot, capturedNative);
@ -61,5 +88,6 @@ export function capturePreparedModelRuntimeCatalog(
},
});
copyPreparedModelRuntimeAuthBindings(snapshot, captured);
admittedCatalogs.set(captured, cached.admittedCatalog);
return captured;
}

View file

@ -193,9 +193,10 @@ describe("prepared model runtime Gateway leases", () => {
runtimePluginSelections: [{ provider: "openai", modelId: "gpt-5.5", runtime: "codex" }],
};
const configured = getPreparedModelRuntimeSnapshot(configuredInput);
const configuredLease = await acquireAgentRunPreparedModelRuntime(configuredInput);
expect(configuredLease.snapshot).toBe(configured);
await configuredLease[Symbol.asyncDispose]();
{
await using configuredLease = await acquireAgentRunPreparedModelRuntime(configuredInput);
expect(configuredLease.snapshot).toBe(configured);
}
for (let index = 0; index < 9; index += 1) {
const lease = await acquireAgentRunPreparedModelRuntime({

View file

@ -89,6 +89,26 @@ describe("resolveCronThinkingSelection scoped hydration", () => {
},
);
it.each([
{ refreshed: [] },
{ refreshed: [{ provider: "other", id: "unrelated", name: "Unrelated", reasoning: true }] },
])("keeps the admitted catalog when hydration has no selected row: %j", async ({ refreshed }) => {
scopedThinkingCatalogMock.mockResolvedValue(refreshed);
const carried = { provider: "openai", id: "gpt-5.6-luna", name: "Selected", reasoning: false };
const { resolveCronThinkingSelection } = await import("./model-selection.js");
const selection = await resolveCronThinkingSelection({
cfg: {},
owner: { ...owner, modelCatalog: { entries: [carried], routeVariants: [] } },
provider: carried.provider,
model: carried.id,
agentRuntime: "codex",
jobThinking: "medium",
});
expect(selection.catalog).toEqual([carried]);
expect(selection.requestedThinkLevel).toBe("medium");
expect(scopedThinkingCatalogMock).toHaveBeenCalledOnce();
});
it("keeps the owner catalog and skips hydration when thinking is off", async () => {
const { resolveCronThinkingSelection } = await import("./model-selection.js");
const selection = await resolveCronThinkingSelection({

View file

@ -1,3 +1,4 @@
import { findModelInCatalog } from "../../agents/model-catalog-lookup.js";
import type { ModelCatalogEntry } from "../../agents/model-catalog.types.js";
import { splitTrailingAuthProfile } from "../../agents/model-ref-profile.js";
import { resolveConfiguredModelPolicyAllow } from "../../agents/model-selection-shared.js";
@ -139,7 +140,7 @@ async function resolveCronThinkingCatalog(params: {
return catalog;
}
// Thinking capability is a per-model fact; never materialize the full live catalog on cron turns.
return normalizeThinkingCatalogProviders(
const refreshed = normalizeThinkingCatalogProviders(
await loadProviderScopedThinkingCatalog({
config: params.owner.config,
provider: params.provider,
@ -150,6 +151,7 @@ async function resolveCronThinkingCatalog(params: {
workspaceDir: params.owner.workspaceDir,
}),
);
return findModelInCatalog(refreshed, params.provider, params.model) ? refreshed : catalog;
}
export async function resolveCronThinkingSelection(params: {

View file

@ -0,0 +1,582 @@
/** Real Gateway reload/failover proof with final-request account and capability checks. */
import fs from "node:fs/promises";
import { createServer, type ServerResponse } from "node:http";
import path from "node:path";
import { performance } from "node:perf_hooks";
import { afterEach, describe, expect, it } from "vitest";
import {
getPreparedModelCatalogOwnerSnapshot,
loadProviderScopedThinkingCatalog,
} from "../src/agents/prepared-model-catalog.js";
import { registerPreparedModelRuntimePublicationListener } from "../src/agents/prepared-model-runtime.publication-events.js";
import {
clearConfigCache,
clearRuntimeConfigSnapshot,
getRuntimeConfig,
} from "../src/config/config.js";
import { clearSessionStoreCacheForTest } from "../src/config/sessions/store-writer-state.js";
import type { ModelDefinitionConfig, ModelProviderConfig } from "../src/config/types.models.js";
import {
disconnectGatewayClient,
startGatewayWithClient,
} from "../src/gateway/test-helpers.e2e.js";
import { captureEnv, setTestEnvValue } from "../src/test-utils/env.js";
import { createSolidPngBuffer } from "./helpers/image-fixtures.js";
import { createDeferred } from "./helpers/promise.js";
import { useAutoCleanupTempDirTracker } from "./helpers/temp-dir.js";
const envKeys = [
"HOME",
"OPENCLAW_STATE_DIR",
"OPENCLAW_CONFIG_PATH",
"OPENCLAW_GATEWAY_TOKEN",
"OPENCLAW_SKIP_CHANNELS",
"OPENCLAW_SKIP_GMAIL_WATCHER",
"OPENCLAW_SKIP_CRON",
"OPENCLAW_SKIP_CANVAS_HOST",
"OPENCLAW_SKIP_BROWSER_CONTROL_SERVER",
"OPENCLAW_SKIP_PROVIDERS",
"OPENCLAW_BUNDLED_PLUGINS_DIR",
"OPENCLAW_DISABLE_BUNDLED_PLUGINS",
] as const;
const PROVIDER_ID = "mock-anthropic";
const PLUGIN_ID = "catalog-reload-proof";
const PRIMARY_MODEL_ID = "claude-opus-5";
const FALLBACK_MODEL_ID = "catalog-only-fallback";
const TOKEN = "pr154462-proof-token";
const REPLY_MARKER = "PR154462_TURN_COMPLETED_AFTER_REPLACEMENT";
const ACCOUNT_A = "fixture-account-a-v1";
const ACCOUNT_A_REFRESHED = "fixture-account-a-v2";
const ACCOUNT_B = "fixture-account-b";
const epoch = performance.now();
function proof(event: string, data: Record<string, unknown> = {}): void {
console.log(
JSON.stringify({
proof: true,
event,
utc: new Date().toISOString(),
ms: Number((performance.now() - epoch).toFixed(1)),
pid: process.pid,
...data,
}),
);
}
function anthropicSse(events: Record<string, unknown>[]): string {
return events
.map((event) => `event: ${String(event.type)}\ndata: ${JSON.stringify(event)}\n\n`)
.join("");
}
/** A plain streamed text answer attributed to the requested model. */
function textTurn(model: string, text: string): string {
return anthropicSse([
{
type: "message_start",
message: {
id: `msg_pr154462_${model}`,
type: "message",
role: "assistant",
model,
content: [],
stop_reason: null,
usage: { input_tokens: 320, output_tokens: 0 },
},
},
{ type: "content_block_start", index: 0, content_block: { type: "text", text: "" } },
{ type: "content_block_delta", index: 0, delta: { type: "text_delta", text } },
{ type: "content_block_stop", index: 0 },
{
type: "message_delta",
delta: { stop_reason: "end_turn", stop_sequence: null },
usage: { output_tokens: 8 },
},
{ type: "message_stop" },
]);
}
function buildMockAnthropicProvider(baseUrl: string) {
const model: ModelDefinitionConfig = {
id: PRIMARY_MODEL_ID,
name: "Mock Claude Opus 5",
api: "anthropic-messages",
reasoning: true,
input: ["text", "image"],
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0 },
contextWindow: 200_000,
maxTokens: 4096,
};
// The fallback model is deliberately NOT declared: with no authored row and no carried
// catalog entry, the failover's resolveRunModelHasVision read must hydrate its
// capabilities through loadProviderScopedThinkingCatalog on the real turn path.
const config: Omit<ModelProviderConfig, "models"> & { models: [ModelDefinitionConfig] } = {
baseUrl,
apiKey: ACCOUNT_A,
api: "anthropic-messages",
models: [model],
};
return {
providerId: PROVIDER_ID,
primaryRef: `${PROVIDER_ID}/${PRIMARY_MODEL_ID}`,
fallbackRef: `${PROVIDER_ID}/${FALLBACK_MODEL_ID}`,
config,
} as const;
}
type ProviderRequest = {
model: string;
credential: "A-original" | "A-refreshed" | "B" | "unknown";
thinking: { type?: string; budget_tokens?: number } | null;
imageCount: number;
};
describe("runtime-config replacement during a turn", () => {
const tempDirs = useAutoCleanupTempDirTracker(afterEach);
it(
"keeps admitted request facts across same-account refresh and same-route account replacement",
{ timeout: 120_000 },
async () => {
const envSnapshot = captureEnv([...envKeys]);
let providerServer: ReturnType<typeof createServer> | undefined;
let gateway: Awaited<ReturnType<typeof startGatewayWithClient>> | undefined;
const requests: ProviderRequest[] = [];
const catalogRequests: ProviderRequest["credential"][] = [];
let catalogModels: ModelDefinitionConfig[] = [];
const identifyCredential = (key: unknown): ProviderRequest["credential"] =>
key === ACCOUNT_A
? "A-original"
: key === ACCOUNT_A_REFRESHED
? "A-refreshed"
: key === ACCOUNT_B
? "B"
: "unknown";
const image = createSolidPngBuffer(8, 8, { r: 40, g: 100, b: 180 }).toString("base64");
const expectRichRequest = (request: ProviderRequest) => {
expect(request.imageCount).toBe(1);
expect(request.thinking?.type).toMatch(/^(enabled|adaptive)$/);
if (request.thinking?.type === "enabled") {
expect(request.thinking.budget_tokens).toBeGreaterThan(0);
}
};
let heldResponse: ServerResponse | undefined;
let onPrimaryHeld: (response: ServerResponse) => void = () => {};
let holdPrimary = false;
let failPrimary = false;
try {
const tempHome = tempDirs.make("openclaw-config-reload-proof-");
const stateDir = path.join(tempHome, ".openclaw");
const workspaceDir = path.join(tempHome, "workspace");
const configPath = path.join(stateDir, "openclaw.json");
const bundledPluginsDir = path.join(tempHome, "bundled-plugins");
await Promise.all([
fs.mkdir(workspaceDir, { recursive: true }),
fs.mkdir(bundledPluginsDir, { recursive: true }),
fs.mkdir(stateDir, { recursive: true }),
]);
for (const [key, value] of Object.entries({
HOME: tempHome,
OPENCLAW_STATE_DIR: stateDir,
OPENCLAW_CONFIG_PATH: configPath,
OPENCLAW_GATEWAY_TOKEN: TOKEN,
OPENCLAW_SKIP_CHANNELS: "1",
OPENCLAW_SKIP_GMAIL_WATCHER: "1",
OPENCLAW_SKIP_CRON: "1",
OPENCLAW_SKIP_CANVAS_HOST: "1",
OPENCLAW_SKIP_BROWSER_CONTROL_SERVER: "1",
OPENCLAW_SKIP_PROVIDERS: "1",
OPENCLAW_BUNDLED_PLUGINS_DIR: bundledPluginsDir,
OPENCLAW_DISABLE_BUNDLED_PLUGINS: "1",
})) {
setTestEnvValue(key, value);
}
const pluginDir = path.join(tempHome, "provider-plugin");
const pluginFile = path.join(pluginDir, "index.mjs");
await fs.mkdir(pluginDir, { recursive: true });
await fs.writeFile(
path.join(pluginDir, "openclaw.plugin.json"),
JSON.stringify({
id: PLUGIN_ID,
providers: [PROVIDER_ID],
providerCatalogEntry: "./provider-discovery.mjs",
modelCatalog: { discovery: { [PROVIDER_ID]: "runtime" }, runtimeAugment: true },
configSchema: { type: "object", properties: {}, additionalProperties: false },
}),
);
await fs.writeFile(
path.join(pluginDir, "provider-discovery.mjs"),
`export default {
id: ${JSON.stringify(PROVIDER_ID)}, label: "Reload proof provider", auth: [],
catalog: { order: "simple", async run(ctx) {
const configured = ctx.config.models.providers[${JSON.stringify(PROVIDER_ID)}];
const auth = ctx.resolveProviderApiKey(${JSON.stringify(PROVIDER_ID)});
const key = auth.discoveryApiKey ?? auth.apiKey;
if (!key) throw new Error("Fixture catalog has no materialized credential");
const response = await fetch(configured.baseUrl + "/models", {
headers: { "x-api-key": key }, signal: ctx.signal,
});
if (!response.ok) throw new Error("Fixture catalog authentication failed");
const result = await response.json();
return { provider: { ...configured, models: result.models } };
} },
};`,
);
await fs.writeFile(
pluginFile,
`import provider from "./provider-discovery.mjs";
export default { id: ${JSON.stringify(PLUGIN_ID)}, register(api) { api.registerProvider(provider); } };`,
);
providerServer = createServer((request, response) => {
const credential = identifyCredential(request.headers["x-api-key"]);
if (credential === "unknown") {
response.writeHead(401, { "content-type": "application/json" });
response.end(
JSON.stringify({
type: "error",
error: { type: "authentication_error", message: "unknown fixture account" },
}),
);
return;
}
if (request.url === "/models") {
catalogRequests.push(credential);
const rich = credential !== "B";
response.writeHead(200, { "content-type": "application/json" });
response.end(
JSON.stringify({
models: catalogModels.map((model) =>
Object.assign({}, model, {
reasoning: rich,
input: rich ? ["text", "image"] : ["text"],
} satisfies Pick<ModelDefinitionConfig, "reasoning" | "input">),
),
}),
);
return;
}
let body = "";
request.setEncoding("utf8");
request.on("data", (chunk) => {
body += chunk;
});
request.on("end", () => {
const parsed = JSON.parse(body) as {
model?: string;
thinking?: ProviderRequest["thinking"];
messages?: Array<{ content?: string | Array<{ type?: string }> }>;
};
const model = parsed.model ?? "";
requests.push({
model,
credential,
thinking: parsed.thinking ?? null,
imageCount: (parsed.messages ?? []).reduce(
(count, message) =>
count +
(Array.isArray(message.content)
? message.content.filter((part) => part.type === "image").length
: 0),
0,
),
});
if (holdPrimary && model === PRIMARY_MODEL_ID) {
holdPrimary = false;
heldResponse = response;
onPrimaryHeld(response);
return;
}
if (failPrimary && model === PRIMARY_MODEL_ID) {
response.writeHead(404, { "content-type": "application/json" });
response.end(
JSON.stringify({
type: "error",
error: { type: "not_found_error", message: `model: ${model}` },
}),
);
return;
}
response.writeHead(200, {
"content-type": "text/event-stream; charset=utf-8",
"cache-control": "no-cache",
});
response.end(textTurn(model, model === FALLBACK_MODEL_ID ? REPLY_MARKER : "warmup-ok"));
});
});
await new Promise<void>((resolve, reject) => {
providerServer!.once("error", reject);
providerServer!.listen(0, "127.0.0.1", resolve);
});
const address = providerServer.address();
if (!address || typeof address === "string") {
throw new Error("loopback provider did not bind");
}
const provider = buildMockAnthropicProvider(`http://127.0.0.1:${address.port}`);
catalogModels = [
...provider.config.models,
{ ...provider.config.models[0], id: FALLBACK_MODEL_ID, name: "Discovered fallback" },
];
gateway = await startGatewayWithClient({
cfg: {
plugins: {
allow: [PLUGIN_ID],
load: { paths: [pluginFile] },
entries: { [PLUGIN_ID]: { enabled: true } },
},
agents: {
defaults: {
workspace: workspaceDir,
skipBootstrap: true,
model: { primary: provider.primaryRef, fallbacks: [provider.fallbackRef] },
thinkingDefault: "low",
},
entries: { main: { default: true } },
},
models: { mode: "merge", providers: { [PROVIDER_ID]: provider.config } },
gateway: { auth: { mode: "token", token: TOKEN } },
},
configPath,
token: TOKEN,
clientDisplayName: "config-reload-proof",
scopes: ["operator.admin", "operator.read", "operator.write"],
hotReloadRecovery: () => ({ status: "emitted" as const }),
});
const client = gateway.client;
const patchConfig = async (models: ModelProviderConfig) => {
const before = await client.request<{ hash: string }>("config.get", {});
const patched = await client.request<{ hash?: string }>(
"config.patch",
{
baseHash: before.hash,
replacePaths: [`models.providers.${PROVIDER_ID}.models[].input`],
raw: JSON.stringify({ models: { providers: { [PROVIDER_ID]: models } } }),
},
{ timeoutMs: 60_000 },
);
expect(patched.hash).toEqual(expect.any(String));
};
const refreshCatalog = async (rich: boolean) => {
const requestsBefore = catalogRequests.length;
const committed = createDeferred<void>();
const checkReady = () => {
const owner = getPreparedModelCatalogOwnerSnapshot({ config: getRuntimeConfig() });
const catalog = owner?.readFullModelCatalog?.() ?? owner?.modelCatalog;
const model = catalog?.entries.find(
(entry) => entry.provider === PROVIDER_ID && entry.id === FALLBACK_MODEL_ID,
);
const ready =
catalogRequests.length > requestsBefore &&
!catalog?.pendingProviders?.includes(PROVIDER_ID) &&
model?.reasoning === rich &&
model.input?.includes("image") === rich;
if (ready) {
committed.resolve();
}
return ready;
};
const unsubscribe = registerPreparedModelRuntimePublicationListener((event) => {
if (event.phase === "catalog-failed") {
committed.reject(event.error);
} else if (event.phase === "catalog-published") {
checkReady();
}
});
try {
await Promise.all([
client
.request(
"models.list",
{ agentId: "main", provider: PROVIDER_ID, refresh: true },
{ timeoutMs: 60_000 },
)
.then(() => {
const ready = checkReady();
const owner = getPreparedModelCatalogOwnerSnapshot({
config: getRuntimeConfig(),
});
const catalog = owner?.readFullModelCatalog?.() ?? owner?.modelCatalog;
proof("catalog_refresh_observed", {
rich,
ready,
requests: [...catalogRequests],
pending: catalog?.pendingProviders ?? [],
entries: catalog?.entries.length ?? 0,
});
if (!ready && !catalog?.pendingProviders?.includes(PROVIDER_ID)) {
throw new Error(
"Catalog refresh settled without the requested provider inventory",
);
}
}),
committed.promise,
]);
} finally {
unsubscribe();
}
const facts = await loadProviderScopedThinkingCatalog({
config: getRuntimeConfig(),
provider: PROVIDER_ID,
model: FALLBACK_MODEL_ID,
});
expect(facts).toContainEqual(
expect.objectContaining({
id: FALLBACK_MODEL_ID,
reasoning: rich,
input: rich ? ["text", "image"] : ["text"],
}),
);
proof("catalog_ready", { credential: catalogRequests.at(-1), rich });
};
const send = async (sessionKey: string, id: string, withImage = false) => {
const started = await client.request<{ status?: string; runId?: string }>("chat.send", {
sessionKey,
message: "Reply with a short answer.",
deliver: false,
idempotencyKey: id,
...(withImage
? {
attachments: [
{
type: "image",
mimeType: "image/png",
fileName: "sample.png",
content: `data:image/png;base64,${image}`,
},
],
}
: {}),
});
expect(started.status).toBe("started");
return started;
};
const waitForRun = async (runId: string | undefined) => {
const result = await client.request<{ status?: string }>(
"agent.wait",
{ runId, timeoutMs: 30_000 },
{ timeoutMs: 60_000 },
);
expect(result).toMatchObject({ status: "ok" });
};
await waitForRun((await send("agent:main:reload-warmup", "reload-warmup")).runId);
await refreshCatalog(true);
const beforeBaseline = requests.length;
failPrimary = true;
await waitForRun((await send("agent:main:reload-baseline", "reload-baseline", true)).runId);
failPrimary = false;
const baseline = requests
.slice(beforeBaseline)
.find((request) => request.model === FALLBACK_MODEL_ID);
if (!baseline) {
throw new Error("no fallback request in the no-reload control");
}
expect(baseline.credential).toBe("A-original");
expectRichRequest(baseline);
const admittedShape = { thinking: baseline.thinking, imageCount: baseline.imageCount };
proof("baseline_request_verified", { credential: baseline.credential, ...admittedShape });
const scenarios = [
{
name: "same-account-refresh",
replacementKey: ACCOUNT_A_REFRESHED,
nextCredential: "A-refreshed",
},
{ name: "same-route-account-switch", replacementKey: ACCOUNT_B, nextCredential: "B" },
] as const;
for (const [scenarioIndex, scenario] of scenarios.entries()) {
if (scenarioIndex > 0) {
await patchConfig(provider.config);
}
await refreshCatalog(true);
const capturedConfig = getRuntimeConfig();
const configuredOwner = getPreparedModelCatalogOwnerSnapshot({ config: capturedConfig });
expect(configuredOwner).toBeDefined();
const carriedFallback = configuredOwner?.modelCatalog.entries.find(
(entry) => entry.provider === PROVIDER_ID && entry.id === FALLBACK_MODEL_ID,
);
expect(carriedFallback?.input?.includes("image") ?? false).toBe(false);
const primaryHeld = new Promise<ServerResponse>((resolve) => {
onPrimaryHeld = resolve;
});
holdPrimary = true;
failPrimary = false;
heldResponse = undefined;
const sessionKey = `agent:main:${scenario.name}`;
const held = await send(sessionKey, scenario.name, true);
const admittedResponse = await primaryHeld;
const richReplacement = scenario.replacementKey !== ACCOUNT_B;
const replacement: ModelProviderConfig = {
...provider.config,
apiKey: scenario.replacementKey,
models: provider.config.models.map((model) =>
Object.assign({}, model, {
reasoning: richReplacement,
input: richReplacement ? ["text", "image"] : ["text"],
} satisfies Pick<ModelDefinitionConfig, "reasoning" | "input">),
),
};
await patchConfig(replacement);
const replacementConfig = getRuntimeConfig();
expect(replacementConfig).not.toBe(capturedConfig);
expect(getPreparedModelCatalogOwnerSnapshot({ config: capturedConfig })).toBeUndefined();
await refreshCatalog(richReplacement);
expect(catalogRequests).toContain(scenario.nextCredential);
const beforeFailover = requests.length;
failPrimary = true;
admittedResponse.writeHead(404, { "content-type": "application/json" });
admittedResponse.end(
JSON.stringify({
type: "error",
error: { type: "not_found_error", message: `model: ${PRIMARY_MODEL_ID}` },
}),
);
await waitForRun(held.runId);
const finalRequest = requests
.slice(beforeFailover)
.find((request) => request.model === FALLBACK_MODEL_ID);
expect(finalRequest).toBeDefined();
expect(finalRequest!.credential).toBe("A-original");
expectRichRequest(finalRequest!);
const shape = { thinking: finalRequest!.thinking, imageCount: finalRequest!.imageCount };
expect(shape).toEqual(admittedShape);
const history = await client.request<{ messages?: unknown[] }>("chat.history", {
sessionKey,
limit: 20,
});
expect(JSON.stringify(history.messages)).toContain(REPLY_MARKER);
failPrimary = false;
const beforeNewTurn = requests.length;
await waitForRun(
(await send(`agent:main:${scenario.name}-new`, `${scenario.name}-new`)).runId,
);
expect(
requests.slice(beforeNewTurn).find((request) => request.model === PRIMARY_MODEL_ID)
?.credential,
).toBe(scenario.nextCredential);
proof("final_request_verified", {
scenario: scenario.name,
admittedCredential: finalRequest!.credential,
nextTurnCredential: scenario.nextCredential,
...shape,
persistedReply: true,
});
}
} finally {
heldResponse?.destroy();
if (gateway) {
await disconnectGatewayClient(gateway.client);
await gateway.server.close();
}
if (providerServer?.listening) {
await new Promise<void>((resolve) => {
providerServer!.close(() => resolve());
});
}
envSnapshot.restore();
clearRuntimeConfigSnapshot();
clearConfigCache();
clearSessionStoreCacheForTest();
}
},
);
});