fix(agents): scope quota failures to auth profiles (#121278)

* fix(agents): scope quota failures to auth profiles

* test: repair provider suspension CI coverage

* test: keep suspension reset fixture internal

* fix(agents): spend cooldown probe only on transient candidates

Consume the one-run cooldown probe only when the candidate’s own unavailable reason is transient, so a billing-disabled pin cannot block a recoverable backup.\n\nFinding from the ClawSweeper review on openclaw/openclaw#121278.

* refactor(sessions): deprecate QuotaSuspension.laneId instead of removing

The shipped plugin-SDK surface deprecation policy requires keeping the inert field until the next surface window.

* fix(agents): extend transient probe policy to plugin-harness auth path

* fix(agents): keep provider overload from cooling auth profiles

* fix(agents): exhaust rotation candidates without cooldown records

Co-authored-by: VACInc <3279061+VACInc@users.noreply.github.com>

* test(agents): align auth rotation mocks with current main

* docs: regenerate plugin SDK API baseline

Co-authored-by: VACInc <3279061+VACInc@users.noreply.github.com>

* style(agents): format session-suspension test after rename resolution

---------

Co-authored-by: VACInc <3279061+VACInc@users.noreply.github.com>
Co-authored-by: Peter Steinberger <steipete@gmail.com>
This commit is contained in:
Vito Cappello
2026-08-12 02:07:04 -04:00
committed by GitHub
parent baefefa815
commit 77d89b2fa8
63 changed files with 819 additions and 1178 deletions
+2 -2
View File
@@ -1,5 +1,5 @@
{
"core": 2295,
"channel": 3716,
"plugin": 4040
"channel": 3582,
"plugin": 3997
}
+4 -4
View File
@@ -1,4 +1,4 @@
b0973756164132b2f14542be9af4d48a927abbcf482ad8746742d7391b9f9d36 config-baseline.json
8c3ffcba19ab9f88fa331d24c1e928ac817a1c5453518ccf5c78edc9995880e8 config-baseline.core.json
ddcf52b6ca3b83d8a72a74e0808abf4b16ab17b5fc28bfca5856373590913388 config-baseline.channel.json
d93639a3d59b9b7ecaa27ff38b844a9ec90ac074c9e53f02930146ed21665c66 config-baseline.plugin.json
894ae65aebdb803a735b50d7ddd3cfc27793888113afbbb05f417a6442e82db9 config-baseline.json
c2fc50e668ab74128c7d2ada9fd5273e13867afb5642a56a2ba76c174781479d config-baseline.core.json
552f5ae69ac13628d754e796bace6e800242d09cacbe762593c17ef3693ba754 config-baseline.channel.json
4bcc2364924c80f38139f0508945b6d28b33b70dd2973f0685a8f221d672ac94 config-baseline.plugin.json
@@ -1 +1 @@
{"contentHash":"3d680ab26bba9b4ccd86bd7734dbf9672920df842ae3446142319a283e711941","entrypoint":"agent-harness-runtime","importSpecifier":"openclaw/plugin-sdk/agent-harness-runtime"}
{"contentHash":"ab9f09ea675935d06a12fff912c1f9109311be8be93a5b8d982a9569d9163359","entrypoint":"agent-harness-runtime","importSpecifier":"openclaw/plugin-sdk/agent-harness-runtime"}
+1 -1
View File
@@ -1 +1 @@
{"contentHash":"08051f15195453d4f35616cdb4f0ff47e974f6c2260f0ee1078e1754f5284899","entrypoint":"agent-harness","importSpecifier":"openclaw/plugin-sdk/agent-harness"}
{"contentHash":"0c4bd22dd46955da487cb750a995732a37ffa6000e94e2381621a8dc58333246","entrypoint":"agent-harness","importSpecifier":"openclaw/plugin-sdk/agent-harness"}
+1 -1
View File
@@ -1 +1 @@
{"contentHash":"0c78a76e24e4afa2f37813ea3202262605159ddc1bba101778fd68036ee0ae70","entrypoint":"channel-core","importSpecifier":"openclaw/plugin-sdk/channel-core"}
{"contentHash":"91cc5b30ad903aea2c6c002f64d21c8236fb141a2f8f8f0ca5217d46c7e79883","entrypoint":"channel-core","importSpecifier":"openclaw/plugin-sdk/channel-core"}
@@ -1 +1 @@
{"contentHash":"3b5041b5034fac21a85f1e6fbe993d79e2168b99abc98279d7e83013de000496","entrypoint":"channel-entry-contract","importSpecifier":"openclaw/plugin-sdk/channel-entry-contract"}
{"contentHash":"2c922f70f8c8155d9be8fb0fb0c6dab4b03f7d1765ad98bab4668a49f984aa5f","entrypoint":"channel-entry-contract","importSpecifier":"openclaw/plugin-sdk/channel-entry-contract"}
@@ -1 +1 @@
{"contentHash":"e9cc066cb5ea878cdd03c53d9114b2ff650a00bb7c74fa31a3fa6be1e9f04668","entrypoint":"channel-inbound","importSpecifier":"openclaw/plugin-sdk/channel-inbound"}
{"contentHash":"d6c0e5bf1e49fc6e003221a3ff6daf951082742fa7e82215572456e4d4eaaf18","entrypoint":"channel-inbound","importSpecifier":"openclaw/plugin-sdk/channel-inbound"}
@@ -1 +1 @@
{"contentHash":"a152541a6f351027b97e0c99a02f27038ea5475e96f4513df1b950465c390242","entrypoint":"channel-message","importSpecifier":"openclaw/plugin-sdk/channel-message"}
{"contentHash":"3c2ea28cb7f23bf09c7ad0cab39751fcefd46055c563d1d2ce226b628f86dc48","entrypoint":"channel-message","importSpecifier":"openclaw/plugin-sdk/channel-message"}
@@ -1 +1 @@
{"contentHash":"d75ce81f95cf3a2382057e90d70f5f8b925cd3d6236a5acce844e3512f5539b6","entrypoint":"channel-outbound","importSpecifier":"openclaw/plugin-sdk/channel-outbound"}
{"contentHash":"a31b1c8f7db35e0e316a59afaa46f27f076b55f11a1dc08dc7630224115038de","entrypoint":"channel-outbound","importSpecifier":"openclaw/plugin-sdk/channel-outbound"}
@@ -1 +1 @@
{"contentHash":"ac15e49c895e670c58a6574e31839a8a04a124f84801eb2e63adc15d1fb59064","entrypoint":"channel-plugin-common","importSpecifier":"openclaw/plugin-sdk/channel-plugin-common"}
{"contentHash":"4edbdb687dfc1780b629e3e45efd1eb183a07d3d08c14d776a2e188b76f7c49b","entrypoint":"channel-plugin-common","importSpecifier":"openclaw/plugin-sdk/channel-plugin-common"}
+1 -1
View File
@@ -1 +1 @@
{"contentHash":"a6436ac71e8ba6dd99d2fb9dee159bf90675a4ed8beec9805109731630a29957","entrypoint":"core","importSpecifier":"openclaw/plugin-sdk/core"}
{"contentHash":"cb3406cb7205502c40a8bbf80521e58fd083620bbed99213d83a804b8deb8264","entrypoint":"core","importSpecifier":"openclaw/plugin-sdk/core"}
+1 -1
View File
@@ -1 +1 @@
{"contentHash":"eb7884340d27b0e901a2431eff2a64fae5a3d018ea6e910725ac7df2d35520b8","entrypoint":"discord","importSpecifier":"openclaw/plugin-sdk/discord"}
{"contentHash":"78ef9651b379d1b673d6c17bc9b636cdaaf2ea4cbe44c546482d53e7b7c8fdd9","entrypoint":"discord","importSpecifier":"openclaw/plugin-sdk/discord"}
@@ -1 +1 @@
{"contentHash":"9d305ddc636bf8fb05767cb6be2bf3320b7801d99b0f24f9b012eae2c28959df","entrypoint":"inbound-reply-dispatch","importSpecifier":"openclaw/plugin-sdk/inbound-reply-dispatch"}
{"contentHash":"077ae11669eaa9cba548cae506db3c11261712b3ffc804c48a28bf159caf124c","entrypoint":"inbound-reply-dispatch","importSpecifier":"openclaw/plugin-sdk/inbound-reply-dispatch"}
@@ -1 +1 @@
{"contentHash":"ebad10b0c7b1838e6cdf68153383fdc57096dd252ccbc24de5691c1f1821f3cd","entrypoint":"meeting-runtime","importSpecifier":"openclaw/plugin-sdk/meeting-runtime"}
{"contentHash":"bd05a541d2ba0e7f335ed0c4d4f187bed1c04953479682616d2b6c416bbd0459","entrypoint":"meeting-runtime","importSpecifier":"openclaw/plugin-sdk/meeting-runtime"}
+1 -1
View File
@@ -1 +1 @@
{"contentHash":"2d9af7dec89f421a59872de6623d51ed84d8ee23204f147d49b5218ae64ddf2c","entrypoint":"plugin-entry","importSpecifier":"openclaw/plugin-sdk/plugin-entry"}
{"contentHash":"f995b2ac4652ccf6a6107700f3a15b9d8e7bd8513466704424b899a831abdb41","entrypoint":"plugin-entry","importSpecifier":"openclaw/plugin-sdk/plugin-entry"}
+1 -1
View File
@@ -1 +1 @@
{"contentHash":"793f98ff0d772f586920d8a75df073deba94321dcba74daf1caea9ba59cca05e","entrypoint":"plugin-runtime","importSpecifier":"openclaw/plugin-sdk/plugin-runtime"}
{"contentHash":"97a6f344d3efa2dd1f7d7620562f22bf632873648a82be69a26cad6a8bd0b762","entrypoint":"plugin-runtime","importSpecifier":"openclaw/plugin-sdk/plugin-runtime"}
@@ -1 +1 @@
{"contentHash":"4841c241310369c97349cc63c956d466a078a1864443e24955ed79031613c806","entrypoint":"provider-catalog-runtime","importSpecifier":"openclaw/plugin-sdk/provider-catalog-runtime"}
{"contentHash":"2d76a60bfd593777339ee478306234773ac9ad626c564c60c95d49f82019ae2c","entrypoint":"provider-catalog-runtime","importSpecifier":"openclaw/plugin-sdk/provider-catalog-runtime"}
+1 -1
View File
@@ -1 +1 @@
{"contentHash":"77a1e0b93fa1a501677e5a571ed5ce44358812407554e5e494488d83bcad34d6","entrypoint":"tool-plugin","importSpecifier":"openclaw/plugin-sdk/tool-plugin"}
{"contentHash":"11c6288303d0e1b2f75ad06f5acac4fbafd94ea98eaa6212d7e0adf87db43607","entrypoint":"tool-plugin","importSpecifier":"openclaw/plugin-sdk/tool-plugin"}
@@ -1 +1 @@
{"contentHash":"03433e97b67ba2d0a95fbb7876710f2f4f4425c66d724f55bb4476797b9f2846","entrypoint":"webhook-ingress","importSpecifier":"openclaw/plugin-sdk/webhook-ingress"}
{"contentHash":"9204480482168d7f4d8f9601466c147d056403b42e3ba7979cae746792c268a8","entrypoint":"webhook-ingress","importSpecifier":"openclaw/plugin-sdk/webhook-ingress"}
+1
View File
@@ -69,6 +69,7 @@ Do not write `type: "aws-sdk"` into the credential store; stored credentials are
- When `auth.order.<provider>` or the auth-store order override is set for a provider, `models status --probe` only probes profile ids that remain in the resolved auth order for that provider. The stored override wins over `auth.order` config.
- A stored profile for that provider that is omitted from the explicit order is not silently tried later. Probe output reports it with `reasonCode: excluded_by_auth_order` and the detail `Excluded by auth.order for this provider.`
- A valid session user pin is an explicit per-session exception: OpenClaw tries that profile first even when it is omitted from the provider order, then uses the ordered same-provider profiles as retry candidates. A cooldown or disabled window applies only to the affected profile; it does not suppress its eligible siblings.
## Probe target resolution
+3 -3
View File
@@ -145,10 +145,10 @@ OpenClaw **pins the automatically chosen auth profile per session** to keep prov
- a compaction completes (compaction count increments)
- the profile is in cooldown/disabled
Manual selection via `/model …@<profileId> -s` sets a **user override**. A valid user pin survives `/new`, `/reset`, session rollover, compaction, and cooldown windows. OpenClaw clears it when the profile disappears, no longer matches the selected provider, or the user selects another explicit profile. `/model default -s` clears the model override while retaining a compatible auth pin and clearing an incompatible one.
Manual selection via `/model …@<profileId> -s` sets a **user override**. A valid user pin survives `/new`, `/reset`, session rollover, compaction, and cooldown windows. It remains the first preference when eligible; while that exact profile is in cooldown or disabled, OpenClaw tries the next eligible same-provider profile without replacing the stored pin. OpenClaw clears the pin when the profile disappears, no longer matches the selected provider, or the user selects another explicit profile. `/model default -s` clears the model override while retaining a compatible auth pin and clearing an incompatible one.
<Note>
Auto-pinned profiles (selected by the session router) are treated as a **preference**: they are tried first, but OpenClaw may rotate to another profile on rate limits/timeouts. When the original profile becomes available again, new runs can prefer it again without changing the selected model or runtime. User-pinned profiles stay locked on eligible same-provider candidates. A retained pin on the configured default can still move through configured model fallbacks; an explicit user model selection remains strict and reports failure instead.
Auto-pinned and user-pinned auth profiles are both retry preferences: OpenClaw tries the selected profile first while it is eligible, then may rotate to another same-provider profile on auth failures, rate limits, billing limits, or timeouts. A user pin stays persisted during that temporary rotation, so new runs prefer it again after its cooldown expires without changing the selected model or runtime. This auth rotation does not loosen model selection: an explicit user provider/model selection remains strict and reports failure after its same-provider auth profiles are exhausted.
</Note>
### OpenAI Codex subscription plus API-key backup
@@ -169,7 +169,7 @@ Use `auth.order.openai` for the user-facing order:
Use `openai:*` for both ChatGPT/Codex OAuth profiles and OpenAI API-key profiles. When the subscription hits a Codex usage limit, OpenClaw records the exact reset time when Codex provides one, tries the next ordered auth profile, and keeps the run inside the Codex harness. Once the reset time passes, the subscription profile is eligible again and the next automatic selection can return to it.
Use a user-pinned profile only when you want to force one account/key for that session. User-pinned profiles are intentionally strict and do not silently jump to another profile.
Use a user-pinned profile to make one account/key the durable first preference for that session. If it becomes unavailable, OpenClaw temporarily rotates through the remaining eligible `auth.order.openai` profiles and returns to the pinned profile after recovery.
## Cooldowns
@@ -28,7 +28,7 @@ const claudeCliProfile = {
function resolveScope(params: {
cfg?: OpenClawConfig;
store?: AuthProfileStore;
userLockedAuthProfileId?: string;
userPinnedAuthProfileId?: string;
}) {
return resolveExternalCliAuthOverlayScopeFromSelection({
provider: "anthropic",
@@ -139,9 +139,10 @@ describe("resolveExternalCliAuthOverlayScopeFromSelection", () => {
expect(resolveScope({ cfg })).toEqual({ ignoreAutoPreferredProfile: false });
});
it("scopes a user lock to the locked profile instead of ambient CLI auth", () => {
it("loads ordered same-provider CLI fallbacks behind a user pin", () => {
const cfg = {
auth: {
order: { anthropic: ["anthropic:claude-cli"] },
profiles: {
"anthropic:api": { provider: "anthropic", mode: "api_key" },
"anthropic:claude-cli": { provider: "claude-cli", mode: "oauth" },
@@ -149,7 +150,8 @@ describe("resolveExternalCliAuthOverlayScopeFromSelection", () => {
},
} satisfies OpenClawConfig;
expect(resolveScope({ cfg, userLockedAuthProfileId: "anthropic:api" })).toEqual({
expect(resolveScope({ cfg, userPinnedAuthProfileId: "anthropic:api" })).toEqual({
providerIds: ["claude-cli"],
ignoreAutoPreferredProfile: false,
});
});
@@ -23,7 +23,7 @@ export function resolveExternalCliAuthOverlayScopeFromSelection(params: {
modelId?: string;
workspaceDir?: string;
store?: AuthProfileStore;
userLockedAuthProfileId?: string;
userPinnedAuthProfileId?: string;
}): {
providerIds?: readonly string[];
ignoreAutoPreferredProfile: boolean;
@@ -33,7 +33,7 @@ export function resolveExternalCliAuthOverlayScopeFromSelection(params: {
cfg: params.cfg,
workspaceDir: params.workspaceDir,
store: params.store,
userLockedAuthProfileId: params.userLockedAuthProfileId,
userPinnedAuthProfileId: params.userPinnedAuthProfileId,
});
const selectedRuntimeProvider =
resolveCliRuntimeExecutionProvider({
@@ -41,7 +41,7 @@ export function resolveExternalCliAuthOverlayScopeFromSelection(params: {
cfg: params.cfg,
agentId: params.agentId,
modelId: params.modelId,
authProfileId: params.userLockedAuthProfileId,
authProfileId: params.userPinnedAuthProfileId,
}) || (params.provider === CLAUDE_CLI_PROVIDER_ID ? CLAUDE_CLI_PROVIDER_ID : undefined);
const selectedProvider =
authScope.selectedProviderId ??
@@ -56,8 +56,8 @@ export function resolveExternalCliAuthOverlayScopeFromSelection(params: {
...(providerIds.length > 0 ? { providerIds } : {}),
ignoreAutoPreferredProfile:
// Claude CLI should not auto-prefer a profile when runtime selection has
// already chosen Claude CLI and the user did not lock a profile.
!params.userLockedAuthProfileId && selectedProvider === CLAUDE_CLI_PROVIDER_ID,
// already chosen Claude CLI and the user did not pin a profile.
!params.userPinnedAuthProfileId && selectedProvider === CLAUDE_CLI_PROVIDER_ID,
};
}
@@ -66,57 +66,29 @@ function resolveExternalCliAuthScopeFromAuthSelection(params: {
cfg?: OpenClawConfig;
workspaceDir?: string;
store?: AuthProfileStore;
userLockedAuthProfileId?: string;
userPinnedAuthProfileId?: string;
}): {
providerIds: string[];
selectedProviderId?: string;
} {
if (params.userLockedAuthProfileId) {
// Locked profile id means discovery should be scoped to that exact profile's
// compatible external CLI provider, if any.
const providerId = resolveExternalCliProviderIdForCompatibleAuthProfile({
...params,
profileId: params.userLockedAuthProfileId,
})?.externalCliProviderId;
return {
providerIds: providerId ? [providerId] : [],
...(providerId ? { selectedProviderId: providerId } : {}),
};
}
const providerIds: string[] = [];
let sawCompatibleOrderedProfile = false;
let selectedProviderId: string | undefined;
for (const profileId of resolveConfiguredAuthProfileOrder(params)) {
const resolved = resolveExternalCliProviderIdForCompatibleAuthProfile({
...params,
profileId,
});
if (!resolved.compatible) {
continue;
}
if (!sawCompatibleOrderedProfile) {
selectedProviderId = resolved.externalCliProviderId;
sawCompatibleOrderedProfile = true;
}
if (resolved.externalCliProviderId) {
providerIds.push(resolved.externalCliProviderId);
}
}
if (sawCompatibleOrderedProfile) {
return {
providerIds: [...new Set(providerIds)],
...(selectedProviderId ? { selectedProviderId } : {}),
};
}
let compatibleProfileCount = 0;
const profileIds = [
const orderedProfileIds = resolveConfiguredAuthProfileOrder(params);
const allProfileIds = [
...new Set([
...Object.keys(params.cfg?.auth?.profiles ?? {}),
...Object.keys(params.store?.profiles ?? {}),
]),
];
const discoveredProfileIds = orderedProfileIds.length > 0 ? orderedProfileIds : allProfileIds;
const profileIds = params.userPinnedAuthProfileId
? [
params.userPinnedAuthProfileId,
...discoveredProfileIds.filter((profileId) => profileId !== params.userPinnedAuthProfileId),
]
: discoveredProfileIds;
let sawCompatibleOrderedProfile = false;
let selectedProviderId: string | undefined;
let compatibleProfileCount = 0;
for (const profileId of profileIds) {
const resolved = resolveExternalCliProviderIdForCompatibleAuthProfile({
...params,
@@ -126,10 +98,21 @@ function resolveExternalCliAuthScopeFromAuthSelection(params: {
continue;
}
compatibleProfileCount += 1;
if (!sawCompatibleOrderedProfile) {
selectedProviderId = resolved.externalCliProviderId;
sawCompatibleOrderedProfile = true;
}
if (resolved.externalCliProviderId) {
providerIds.push(resolved.externalCliProviderId);
}
}
if (params.userPinnedAuthProfileId || orderedProfileIds.length > 0) {
return {
providerIds: [...new Set(providerIds)],
...(selectedProviderId ? { selectedProviderId } : {}),
};
}
const uniqueProviderIds = [...new Set(providerIds)];
return {
providerIds: uniqueProviderIds,
+7 -7
View File
@@ -2074,7 +2074,7 @@ describe("runBtwSideQuestion", () => {
it.each([
{ label: "explicit", source: "user" as const },
{ label: "legacy source-less", source: undefined },
])("keeps $label user-locked static Anthropic auth for BTW", async ({ source }) => {
])("keeps $label user-pinned static Anthropic auth first for BTW", async ({ source }) => {
const staticAuthStore = {
version: 1 as const,
profiles: {
@@ -2085,7 +2085,7 @@ describe("runBtwSideQuestion", () => {
},
},
};
ensureAuthProfileStoreWithoutExternalProfilesMock.mockReturnValueOnce(staticAuthStore);
ensureAuthProfileStoreMock.mockReturnValueOnce(staticAuthStore);
getApiKeyForModelMock.mockResolvedValueOnce({
apiKey: "static-key",
mode: "api-key",
@@ -2112,11 +2112,11 @@ describe("runBtwSideQuestion", () => {
}),
});
expect(ensureAuthProfileStoreMock).not.toHaveBeenCalled();
expect(ensureAuthProfileStoreWithoutExternalProfilesMock).toHaveBeenCalledWith(
DEFAULT_AGENT_DIR,
{ allowKeychainPrompt: false },
);
expect(ensureAuthProfileStoreWithoutExternalProfilesMock).not.toHaveBeenCalled();
expect(ensureAuthProfileStoreMock).toHaveBeenCalledWith(DEFAULT_AGENT_DIR, {
externalCliProviderIds: ["claude-cli"],
allowKeychainPrompt: false,
});
expectRecordFields(mockArg(getApiKeyForModelMock, 0, 0), {
profileId: "anthropic:api",
store: staticAuthStore,
+3 -3
View File
@@ -162,7 +162,7 @@ function resolveBtwAuthProfileStore(params: {
};
}
const userLockedAuthProfileId =
const userPinnedAuthProfileId =
params.authProfileIdSource === "user" ? params.authProfileId : undefined;
let externalCliAuthScope = resolveExternalCliAuthOverlayScopeFromSelection({
provider: params.provider,
@@ -170,7 +170,7 @@ function resolveBtwAuthProfileStore(params: {
agentId: params.agentId,
modelId: params.modelId,
workspaceDir: params.workspaceDir,
userLockedAuthProfileId,
userPinnedAuthProfileId,
});
let store: AuthProfileStore;
if (externalCliAuthScope.providerIds) {
@@ -189,7 +189,7 @@ function resolveBtwAuthProfileStore(params: {
modelId: params.modelId,
workspaceDir: params.workspaceDir,
store,
userLockedAuthProfileId,
userPinnedAuthProfileId,
});
if (externalCliAuthScope.providerIds) {
store = ensureAuthProfileStore(params.agentDir, {
@@ -14,6 +14,7 @@ import {
} from "./auth-profiles.js";
import { ensureAuthProfileStore, saveAuthProfileStore } from "./auth-profiles/store.js";
import type { EmbeddedRunAttemptResult } from "./embedded-agent-runner/run/types.js";
import type { AgentHarness } from "./harness/types.js";
import {
buildEmbeddedRunnerAssistant as buildAssistant,
makeEmbeddedRunnerAttempt as makeAttempt,
@@ -55,26 +56,32 @@ const installRunEmbeddedMocks = () => {
// The model resolver stays deterministic so retry assertions only observe
// profile selection, cooldowns, and provider auth preparation.
vi.doMock("./embedded-agent-runner/model.js", () => ({
resolveModelAsync: async (provider: string, modelId: string) => ({
model: {
id: modelId,
name: modelId,
api: "openai-responses",
provider,
baseUrl:
provider === "github-copilot" ? "https://api.copilot.example" : "https://example.com",
reasoning: false,
input: ["text"],
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0 },
contextWindow: 16_000,
maxTokens: 2048,
},
error: undefined,
authStorage: {
setRuntimeApiKey: vi.fn(),
},
modelRegistry: {},
}),
resolveModelAsync: async (provider: string, modelId: string) => {
const subscriptionModel = modelId === "chatgpt-mock";
return {
model: {
id: modelId,
name: modelId,
api: subscriptionModel ? "openai-chatgpt-responses" : "openai-responses",
provider,
baseUrl: subscriptionModel
? "https://chatgpt.com/backend-api/codex"
: provider === "github-copilot"
? "https://api.copilot.example"
: "https://example.com",
reasoning: false,
input: ["text"],
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0 },
contextWindow: 16_000,
maxTokens: 2048,
},
error: undefined,
authStorage: {
setRuntimeApiKey: vi.fn(),
},
modelRegistry: {},
};
},
}));
installEmbeddedRunnerBackoffE2eMocks({
computeBackoff: (policy, attempt) => computeBackoffMock(policy, attempt),
@@ -103,6 +110,7 @@ let createDiagnosticLogRecordCaptureFn: typeof import("../logging/test-helpers/d
let cleanupLogCapture: (() => void) | undefined;
let resetLoggerFn: typeof import("../logging/logger.js").resetLogger;
let setLoggerOverrideFn: typeof import("../logging/logger.js").setLoggerOverride;
let registerAgentHarnessFn: typeof import("./harness/registry.js").registerAgentHarness;
const originalFetch = globalThis.fetch;
beforeAll(async () => {
@@ -115,6 +123,7 @@ beforeAll(async () => {
await import("../logging/test-helpers/diagnostic-log-capture.js"));
({ resetLogger: resetLoggerFn, setLoggerOverride: setLoggerOverrideFn } =
await import("../logging/logger.js"));
({ registerAgentHarness: registerAgentHarnessFn } = await import("./harness/registry.js"));
});
type RunEmbeddedAgentTestParams = Parameters<typeof runEmbeddedAgent>[0] & {
@@ -308,7 +317,7 @@ const writeCopilotAuthStore = async (agentDir: string, token = "gh-token") => {
);
};
const writeOpenAiCodexAuthStore = async (agentDir: string) => {
const writeOpenAiCodexAuthStore = async (agentDir: string, includeBackup = false) => {
saveAuthProfileStore(
{
version: 1,
@@ -318,7 +327,17 @@ const writeOpenAiCodexAuthStore = async (agentDir: string) => {
provider: "openai",
key: "sk-codex",
},
...(includeBackup
? {
"openai:backup": {
type: "api_key" as const,
provider: "openai",
key: "sk-backup",
},
}
: {}),
},
...(includeBackup ? { order: { openai: ["openai:work", "openai:backup"] } } : {}),
},
agentDir,
);
@@ -378,6 +397,35 @@ const mockPromptErrorThenSuccessfulAttempt = (errorMessage: string) => {
);
};
const mockFailedThenSuccessfulAttemptForModel = (params: {
errorMessage: string;
provider: string;
model: string;
}) => {
runEmbeddedAttemptMock
.mockResolvedValueOnce(
makeErrorAttempt(
{
errorMessage: params.errorMessage,
provider: params.provider,
model: params.model,
},
{ currentAttempt: true },
),
)
.mockResolvedValueOnce(
makeAttempt({
assistantTexts: ["ok"],
lastAssistant: buildAssistant({
provider: params.provider,
model: params.model,
stopReason: "stop",
content: [{ type: "text", text: "ok" }],
}),
}),
);
};
async function runAutoPinnedOpenAiTurn(params: {
agentDir: string;
workspaceDir: string;
@@ -458,10 +506,6 @@ async function runAutoPinnedPromptErrorRotationCase(params: {
});
expect(runEmbeddedAttemptMock).toHaveBeenCalledTimes(2);
await vi.waitFor(async () => {
const usageStats = await readUsageStats(agentDir);
expect(typeof usageStats["openai:p1"]?.cooldownUntil).toBe("number");
});
const usageStats = await readUsageStats(agentDir);
return { usageStats };
});
@@ -479,12 +523,12 @@ function mockSingleSuccessfulAttempt() {
);
}
function mockSingleErrorAttempt(params: {
function mockRepeatedErrorAttempts(params: {
errorMessage: string;
provider?: string;
model?: string;
}) {
runEmbeddedAttemptMock.mockResolvedValueOnce(
runEmbeddedAttemptMock.mockResolvedValue(
makeErrorAttempt(
{
errorMessage: params.errorMessage,
@@ -873,7 +917,7 @@ describe("runEmbeddedAgent auth profile rotation", () => {
runId: "run:overloaded-rotation",
});
expect(typeof usageStats["openai:p2"]?.lastUsed).toBe("number");
expect(typeof usageStats["openai:p1"]?.cooldownUntil).toBe("number");
expect(usageStats["openai:p1"]?.cooldownUntil).toBeUndefined();
expect(computeBackoffMock).not.toHaveBeenCalled();
expect(sleepWithAbortMock).not.toHaveBeenCalled();
});
@@ -911,21 +955,12 @@ describe("runEmbeddedAgent auth profile rotation", () => {
expect(failoverAttributes.providerErrorType).toBe("overloaded_error");
expect(failoverAttributes.rawErrorPreview).toContain('"request_id":"sha256:');
await vi.waitFor(async () => {
await logCapture.flush();
const failureStateUpdate = requireLogRecord(
logCapture.records,
"auth profile failure state updated",
);
const failureStateAttributes = requireRecord(
failureStateUpdate.attributes,
"failure state attributes",
);
expect(failureStateAttributes.event).toBe("auth_profile_failure_state_updated");
expect(failureStateAttributes.runId).toBe("run:overloaded-logging");
expect(failureStateAttributes.profileId).toBe(safeProfileId);
expect(failureStateAttributes.reason).toBe("overloaded");
});
expect(
logCapture.records.some(
(record) =>
requireRecord(record, "log record").message === "auth profile failure state updated",
),
).toBe(false);
});
it("rotates for overloaded prompt failures across auto-pinned profiles", async () => {
@@ -935,7 +970,7 @@ describe("runEmbeddedAgent auth profile rotation", () => {
runId: "run:overloaded-prompt-rotation",
});
expect(typeof usageStats["openai:p2"]?.lastUsed).toBe("number");
expect(typeof usageStats["openai:p1"]?.cooldownUntil).toBe("number");
expect(usageStats["openai:p1"]?.cooldownUntil).toBeUndefined();
expect(computeBackoffMock).not.toHaveBeenCalled();
expect(sleepWithAbortMock).not.toHaveBeenCalled();
});
@@ -1078,51 +1113,45 @@ describe("runEmbeddedAgent auth profile rotation", () => {
});
});
it("surfaces rate limits without rotating for user-pinned profiles", async () => {
it("rotates from a rate-limited user pin to the next same-provider profile", async () => {
await withAgentWorkspace(async ({ agentDir, workspaceDir }) => {
await writeAuthStore(agentDir);
mockSingleErrorAttempt({ errorMessage: "rate limit" });
mockFailedThenSuccessfulAttempt("rate limit");
await expectFailoverError(
runEmbeddedAgentInline({
sessionId: "session:test",
sessionKey: "agent:test:user",
workspaceDir,
agentDir,
config: makeConfig(),
prompt: "hello",
provider: "openai",
model: "mock-1",
authProfileId: "openai:p1",
authProfileIdSource: "user",
timeoutMs: 5_000,
runId: "run:user",
}),
{
profileId: "openai:p1",
reason: "rate_limit",
provider: "openai",
model: "mock-1",
},
);
await runEmbeddedAgentInline({
sessionId: "session:test",
sessionKey: "agent:test:user",
workspaceDir,
agentDir,
config: makeConfig(),
prompt: "hello",
provider: "openai",
model: "mock-1",
authProfileId: "openai:p1",
authProfileIdSource: "user",
timeoutMs: 5_000,
runId: "run:user",
});
expect(runEmbeddedAttemptMock).toHaveBeenCalledTimes(1);
await expectProfileP2UsageUnchanged(agentDir);
expect(runEmbeddedAttemptMock).toHaveBeenCalledTimes(2);
const usageStats = await readUsageStats(agentDir);
expect(typeof usageStats["openai:p1"]?.cooldownUntil).toBe("number");
expect(usageStats["openai:p2"]?.lastUsed).not.toBe(2);
});
});
it("honors user-pinned profiles even when in cooldown", async () => {
const { usageStats } = await runTurnWithCooldownSeed({
it("skips a user-pinned profile while only that profile is in cooldown", async () => {
const { usageStats, now } = await runTurnWithCooldownSeed({
sessionKey: "agent:test:user-cooldown",
runId: "run:user-cooldown",
authProfileId: "openai:p1",
authProfileIdSource: "user",
});
expect(usageStats["openai:p1"]?.cooldownUntil).toBeUndefined();
expect(usageStats["openai:p1"]?.lastUsed).not.toBe(1);
expect(usageStats["openai:p2"]?.lastUsed).toBe(2);
expect(usageStats["openai:p1"]?.cooldownUntil).toBe(now + 60 * 60 * 1000);
expect(usageStats["openai:p1"]?.lastUsed).toBe(1);
expect(usageStats["openai:p2"]?.lastUsed).not.toBe(2);
});
it("honors user-pinned profiles even when stored order excludes them", async () => {
@@ -1188,7 +1217,118 @@ describe("runEmbeddedAgent auth profile rotation", () => {
});
});
it("ignores user-locked profile when provider mismatches", async () => {
it("rotates a user-pinned profile inside the Codex harness", async () => {
await withAgentWorkspace(async ({ agentDir, workspaceDir }) => {
await writeOpenAiCodexAuthStore(agentDir, true);
mockFailedThenSuccessfulAttemptForModel({
errorMessage: "rate limit",
provider: "codex-cli",
model: "gpt-5.4",
});
await runEmbeddedAgentInline({
sessionId: "session:test",
sessionKey: "agent:test:user-auth-alias-rotation",
workspaceDir,
agentDir,
config: makeConfig(),
prompt: "hello",
provider: "codex-cli",
model: "gpt-5.4",
authProfileId: "openai:work",
authProfileIdSource: "user",
timeoutMs: 5_000,
runId: "run:user-auth-alias-rotation",
});
expect(runEmbeddedAttemptMock).toHaveBeenCalledTimes(2);
const firstAttempt = requireRecord(
runEmbeddedAttemptMock.mock.calls.at(0)?.[0],
"first Codex attempt params",
);
const secondAttempt = requireRecord(
runEmbeddedAttemptMock.mock.calls.at(1)?.[0],
"second Codex attempt params",
);
expect(firstAttempt.authProfileId).toBe("openai:work");
expect(firstAttempt.authProfileIdSource).toBe("user");
expect(secondAttempt.authProfileId).toBe("openai:backup");
expect(secondAttempt.authProfileIdSource).toBe("auto");
});
});
it("preserves a transient plugin-harness probe after a billing-disabled user pin", async () => {
await withTimedAgentWorkspace(async ({ agentDir, workspaceDir, now }) => {
saveAuthProfileStore(
{
version: 1,
profiles: {
"openai:pinned": {
type: "token",
provider: "openai",
token: "subscription-pinned",
},
"openai:backup": {
type: "token",
provider: "openai",
token: "subscription-backup",
},
},
order: { openai: ["openai:pinned", "openai:backup"] },
usageStats: {
"openai:pinned": {
disabledUntil: now + 60 * 60 * 1000,
disabledReason: "billing",
},
"openai:backup": {
cooldownUntil: now + 60 * 60 * 1000,
failureCounts: { rate_limit: 1 },
},
},
},
agentDir,
);
const harness: AgentHarness = {
id: "probe-harness",
label: "Probe harness",
authBootstrap: "harness",
supports: (ctx) =>
ctx.requestedRuntime === "probe-harness"
? { supported: true, priority: 100 }
: { supported: false, reason: "test harness requires an explicit runtime" },
runAttempt: async (attemptParams) => await runEmbeddedAttemptMock(attemptParams),
};
registerAgentHarnessFn(harness);
mockSingleSuccessfulAttempt();
await runEmbeddedAgentInline({
sessionId: "session:test",
sessionKey: "agent:test:plugin-harness-mixed-cooldown",
workspaceDir,
agentDir,
config: makeConfig(),
prompt: "hello",
provider: "openai",
model: "chatgpt-mock",
agentHarnessId: "probe-harness",
authProfileId: "openai:pinned",
authProfileIdSource: "user",
allowTransientCooldownProbe: true,
timeoutMs: 5_000,
runId: "run:plugin-harness-mixed-cooldown",
});
expect(runEmbeddedAttemptMock).toHaveBeenCalledOnce();
const attemptParams = requireRecord(
runEmbeddedAttemptMock.mock.calls[0]?.[0],
"plugin harness attempt params",
);
expect(attemptParams.authProfileId).toBe("openai:backup");
expect(attemptParams.authProfileIdSource).toBe("auto");
});
});
it("ignores a user-pinned profile when the provider mismatches", async () => {
await withAgentWorkspace(async ({ agentDir, workspaceDir }) => {
await writeAuthStore(agentDir, { includeAnthropic: true });
@@ -1507,7 +1647,7 @@ describe("runEmbeddedAgent auth profile rotation", () => {
it("uses the active erroring model in billing failover errors", async () => {
await withAgentWorkspace(async ({ agentDir, workspaceDir }) => {
await writeAuthStore(agentDir);
mockSingleErrorAttempt({
mockRepeatedErrorAttempts({
errorMessage: "insufficient credits",
provider: "openai",
model: "mock-rotated",
@@ -1539,7 +1679,7 @@ describe("runEmbeddedAgent auth profile rotation", () => {
expect(errorRecord.model).toBe("mock-rotated");
expect(thrown).toBeInstanceOf(Error);
expect((thrown as Error).message).toContain("openai (mock-rotated) returned a billing error");
expect(runEmbeddedAttemptMock).toHaveBeenCalledTimes(1);
expect(runEmbeddedAttemptMock).toHaveBeenCalledTimes(2);
});
});
@@ -135,11 +135,10 @@ async function runEmbeddedAgentInternal(
// Outer fallback attempts defer session suspension only while another
// candidate remains. Direct and final-candidate runs suspend normally.
const failureSuspension = resolveSessionSuspensionTarget();
const suspendForFailure = (suspensionParams: Omit<SessionSuspensionParams, "laneId">) => {
const suspendForFailure = (suspensionParams: SessionSuspensionParams) => {
const suspension = buildEmbeddedFailureSuspension({
suspension: suspensionParams,
runAgentId: params.agentId,
laneId: globalLane,
});
if (failureSuspension.mode === "defer") {
failureSuspension.defer(suspension);
@@ -119,7 +119,9 @@ function setupCompactionRemovedFallbackAttempt() {
return isCurrentAttemptAssistant(assistant) && assistant.provider === "anthropic";
});
mockedClassifyFailoverReason.mockReturnValue("model_not_found");
mockedRunEmbeddedAttempt.mockResolvedValueOnce(
// The pinned profile may rotate to another same-provider credential before
// the outer model fallback runs, so every credential attempt must fail alike.
mockedRunEmbeddedAttempt.mockResolvedValue(
makeAttemptResult({
assistantTexts: [],
lastAssistant: makeAssistantMessageFixture({
@@ -226,7 +228,7 @@ describe("runEmbeddedAgent cross-provider fallback error handling", () => {
await expect(promise).rejects.toThrow(
`anthropic/test-model: ${COMPACTION_REMOVED_ERROR_MESSAGE}`,
);
expect(mockedIsFailoverAssistantError).toHaveBeenCalledTimes(1);
expect(mockedIsFailoverAssistantError).toHaveBeenCalledTimes(2);
expect(getLastFormattedAssistant()).toMatchObject({
provider: "anthropic",
model: "test-model",
@@ -91,7 +91,7 @@ export async function handleEmbeddedAssistantFailure(input: {
typeof handleAssistantFailover
>[0]["advanceRateLimitAuthProfile"];
traceAttempts: TraceAttempt[];
suspendForFailure: (params: Omit<SessionSuspensionParams, "laneId">) => void;
suspendForFailure: (params: SessionSuspensionParams) => void;
suspensionSessionId: string;
agentDir: string;
isProbeSession: boolean;
@@ -229,7 +229,8 @@ export async function prepareAndDispatchEmbeddedRunAttempt(input: {
model: effectiveModel,
resolvedApiKey: resolvedAttemptApiKey,
authProfileId: runtime.lastProfileId,
authProfileIdSource: lockedProfileId ? "user" : "auto",
authProfileIdSource:
runtime.lastProfileId && runtime.lastProfileId === lockedProfileId ? "user" : "auto",
initialReplayState: input.replayState,
authStorage,
authProfileStore: resolveRunAttemptAuthProfileStore(),
@@ -166,7 +166,10 @@ export async function recoverEmbeddedRunAttempt(input: {
provider: preparedRuntime.provider,
model: preparedRuntime.modelId,
authProfileId: runtime.lastProfileId,
authProfileIdSource: preparedRuntime.lockedProfileId ? "user" : "auto",
authProfileIdSource:
runtime.lastProfileId && runtime.lastProfileId === preparedRuntime.lockedProfileId
? "user"
: "auto",
},
requested: requestedSelection,
});
@@ -101,6 +101,8 @@ function createMutableEmbeddedRunAuthController(params: {
profileCandidates?: string[];
authStore?: AuthProfileStore;
fallbackConfigured?: boolean;
lockedProfileId?: string;
allowTransientCooldownProbe?: boolean;
warn?: (message: string) => void;
prepareModelForAuthProfile?: Parameters<
typeof createEmbeddedRunAuthController
@@ -118,10 +120,11 @@ function createMutableEmbeddedRunAuthController(params: {
} as AuthProfileStore),
authStorage: { setRuntimeApiKey: params.setRuntimeApiKey },
profileCandidates: params.profileCandidates ?? ["default"],
lockedProfileId: params.lockedProfileId,
initialThinkLevel: "medium",
attemptedThinking: new Set(),
fallbackConfigured: params.fallbackConfigured ?? false,
allowTransientCooldownProbe: false,
allowTransientCooldownProbe: params.allowTransientCooldownProbe ?? false,
getProvider: () => "custom-openai",
getModelId: () => "test-model",
getRuntimeModel: () => params.harness.runtimeModel,
@@ -371,6 +374,36 @@ describe("createEmbeddedRunAuthController", () => {
expect(setRuntimeApiKey).toHaveBeenLastCalledWith("custom-openai", "backup-source-key");
});
it("exhausts the remaining auth profile after a non-cooling failure", async () => {
const harness = createMutableAuthControllerHarness();
mocks.getApiKeyForModelCore.mockImplementation(async ({ profileId }) => {
if (profileId === "backup") {
throw new Error("provider overloaded");
}
return {
apiKey: "default-key",
mode: "api-key" as const,
profileId,
source: `profile:${String(profileId)}`,
};
});
mocks.prepareProviderRuntimeAuth.mockResolvedValue(undefined);
const controller = createMutableEmbeddedRunAuthController({
harness,
setRuntimeApiKey: vi.fn(),
profileCandidates: ["default", "backup"],
});
await controller.initializeAuthProfile();
await expect(controller.advanceAuthProfile()).resolves.toBe(false);
await expect(controller.advanceAuthProfile()).resolves.toBe(false);
expect(
mocks.getApiKeyForModelCore.mock.calls.filter(([params]) => params.profileId === "backup"),
).toHaveLength(1);
expect(harness.profileIndex).toBe(2);
});
it("unwraps a sentinel for runtime auth exchange but keeps auth storage opaque", async () => {
const harness = createMutableAuthControllerHarness();
const setRuntimeApiKey = vi.fn<(provider: string, apiKey: string) => void>();
@@ -533,29 +566,84 @@ describe("createEmbeddedRunAuthController", () => {
allowTransientCooldownProbe: true,
});
expect(
resolve(
createStore({
first: { disabledUntil: now + 60_000, disabledReason: "rate_limit" },
}),
),
).toEqual({ allowProbe: false, unavailableReason: null });
expect(
resolve(
createStore({
first: { disabledUntil: now + 60_000, disabledReason: "billing" },
second: { disabledUntil: now + 60_000, disabledReason: "billing" },
}),
),
).toEqual({ allowProbe: false, unavailableReason: "billing" });
expect(
resolve(
createStore({
first: { disabledUntil: now + 60_000, disabledReason: "rate_limit" },
second: { disabledUntil: now + 60_000, disabledReason: "rate_limit" },
}),
),
).toEqual({ allowProbe: true, unavailableReason: "rate_limit" });
const partiallyAvailable = resolve(
createStore({
first: { disabledUntil: now + 60_000, disabledReason: "rate_limit" },
}),
);
expect([...partiallyAvailable.probeProfileIds]).toEqual([]);
expect(partiallyAvailable.unavailableReason).toBeNull();
const billingDisabled = resolve(
createStore({
first: { disabledUntil: now + 60_000, disabledReason: "billing" },
second: { disabledUntil: now + 60_000, disabledReason: "billing" },
}),
);
expect([...billingDisabled.probeProfileIds]).toEqual([]);
expect(billingDisabled.unavailableReason).toBe("billing");
const rateLimited = resolve(
createStore({
first: { disabledUntil: now + 60_000, disabledReason: "rate_limit" },
second: { disabledUntil: now + 60_000, disabledReason: "rate_limit" },
}),
);
expect([...rateLimited.probeProfileIds]).toEqual(["first", "second"]);
expect(rateLimited.unavailableReason).toBe("rate_limit");
const mixedPinnedState = resolveEmbeddedAuthCooldownProbePolicy({
authStore: createStore({
first: { disabledUntil: now + 60_000, disabledReason: "billing" },
second: { disabledUntil: now + 60_000, disabledReason: "rate_limit" },
}),
profileCandidates: ["first", "second"],
lockedProfileId: "first",
modelId: "test-model",
allowTransientCooldownProbe: true,
});
expect([...mixedPinnedState.probeProfileIds]).toEqual(["second"]);
expect(mixedPinnedState.unavailableReason).toBe("rate_limit");
});
it("preserves the transient cooldown probe for a rate-limited backup after a billing-disabled pin", async () => {
const harness = createMutableAuthControllerHarness();
const now = Date.now();
mocks.getApiKeyForModelCore.mockImplementation(async ({ profileId }) => ({
apiKey: `${String(profileId)}-key`,
mode: "api-key" as const,
profileId,
source: `profile:${String(profileId)}`,
}));
mocks.prepareProviderRuntimeAuth.mockResolvedValue(undefined);
const controller = createMutableEmbeddedRunAuthController({
harness,
setRuntimeApiKey: vi.fn(),
profileCandidates: ["pinned", "backup"],
lockedProfileId: "pinned",
allowTransientCooldownProbe: true,
authStore: {
version: 1,
profiles: {
pinned: { type: "api_key", provider: "custom-openai", key: "pinned-key" },
backup: { type: "api_key", provider: "custom-openai", key: "backup-key" },
},
usageStats: {
pinned: { disabledUntil: now + 60_000, disabledReason: "billing" },
backup: { blockedUntil: now + 60_000 },
},
},
});
await controller.initializeAuthProfile();
expect(mocks.getApiKeyForModelCore).toHaveBeenCalledOnce();
expect(mocks.getApiKeyForModelCore).toHaveBeenCalledWith(
expect.objectContaining({ profileId: "backup" }),
);
expect(harness.profileIndex).toBe(1);
expect(harness.lastProfileId).toBe("backup");
});
it("rejects privileged runtime transport overrides on the first auth exchange", async () => {
@@ -64,7 +64,7 @@ export function resolveEmbeddedAuthCooldownProbePolicy(params: {
lockedProfileId?: string;
modelId: string;
allowTransientCooldownProbe: boolean;
}): { allowProbe: boolean; unavailableReason: FailoverReason | null } {
}): { probeProfileIds: ReadonlySet<string>; unavailableReason: FailoverReason | null } {
const autoProfileCandidates = params.profileCandidates.filter(
(candidate): candidate is string =>
typeof candidate === "string" && candidate.length > 0 && candidate !== params.lockedProfileId,
@@ -80,13 +80,24 @@ export function resolveEmbeddedAuthCooldownProbePolicy(params: {
profileIds: autoProfileCandidates,
}) ?? "unknown")
: null;
return {
allowProbe:
params.allowTransientCooldownProbe &&
allAutoProfilesInCooldown &&
shouldUseTransientCooldownProbeSlot(unavailableReason),
unavailableReason,
};
const probeProfileIds = new Set<string>();
if (
params.allowTransientCooldownProbe &&
allAutoProfilesInCooldown &&
shouldUseTransientCooldownProbeSlot(unavailableReason)
) {
for (const candidate of autoProfileCandidates) {
const candidateReason =
resolveProfilesUnavailableReason({
store: params.authStore,
profileIds: [candidate],
}) ?? "unknown";
if (shouldUseTransientCooldownProbeSlot(candidateReason)) {
probeProfileIds.add(candidate);
}
}
}
return { probeProfileIds, unavailableReason };
}
/**
@@ -570,22 +581,20 @@ export function createEmbeddedRunAuthController(params: {
};
const advanceAuthProfile = async (): Promise<boolean> => {
if (params.lockedProfileId) {
return false;
}
let nextIndex = params.getProfileIndex() + 1;
while (nextIndex < params.profileCandidates.length) {
const candidate = params.profileCandidates[nextIndex];
const candidateIndex = nextIndex++;
const candidate = params.profileCandidates[candidateIndex];
// Candidate exhaustion is run-local and never depends on a cooldown write.
params.setProfileIndex(candidateIndex);
if (
candidate &&
isProfileInCooldown(params.authStore, candidate, undefined, params.getModelId())
) {
nextIndex += 1;
continue;
}
try {
await applyApiKeyInfo(candidate, nextIndex);
params.setProfileIndex(nextIndex);
await applyApiKeyInfo(candidate, candidateIndex);
params.setThinkLevel(params.initialThinkLevel);
params.attemptedThinking.clear();
return true;
@@ -593,12 +602,9 @@ export function createEmbeddedRunAuthController(params: {
if (err instanceof SecretSurfaceUnavailableError) {
throw err;
}
if (candidate && candidate === params.lockedProfileId) {
throw err;
}
nextIndex += 1;
}
}
params.setProfileIndex(params.profileCandidates.length);
return false;
};
@@ -617,11 +623,13 @@ export function createEmbeddedRunAuthController(params: {
while (params.getProfileIndex() < params.profileCandidates.length) {
const candidate = params.profileCandidates[params.getProfileIndex()];
const inCooldown =
candidate &&
candidate !== params.lockedProfileId &&
isProfileInCooldown(params.authStore, candidate, undefined, modelId);
candidate && isProfileInCooldown(params.authStore, candidate, undefined, modelId);
if (inCooldown) {
if (cooldownProbePolicy.allowProbe && !didTransientCooldownProbe) {
const canProbeCandidate =
!didTransientCooldownProbe && cooldownProbePolicy.probeProfileIds.has(candidate);
// Spend the single probe slot only on a transiently cooled candidate;
// persistent failures must leave it available for later profiles.
if (canProbeCandidate) {
didTransientCooldownProbe = true;
params.log.warn(
`probing cooldowned auth profile for ${params.getProvider()}/${modelId} due to ${cooldownProbePolicy.unavailableReason ?? "transient"} unavailability`,
@@ -644,9 +652,6 @@ export function createEmbeddedRunAuthController(params: {
if (err instanceof FailoverError || err instanceof SecretSurfaceUnavailableError) {
throw err;
}
if (params.profileCandidates[params.getProfileIndex()] === params.lockedProfileId) {
throwAuthProfileFailover({ allInCooldown: false, error: err });
}
const advanced = await advanceAuthProfile();
if (!advanced) {
throwAuthProfileFailover({ allInCooldown: false, error: err });
@@ -87,7 +87,7 @@ export async function prepareEmbeddedRunAuthPlan(params: {
agentId: runParams.agentId,
modelId: params.modelId,
workspaceDir: params.workspaceDir,
userLockedAuthProfileId:
userPinnedAuthProfileId:
runParams.authProfileIdSource === "user" ? runParams.authProfileId : undefined,
});
let noExternalAuthStore: AuthProfileStore | undefined;
@@ -102,7 +102,7 @@ export async function prepareEmbeddedRunAuthPlan(params: {
modelId: params.modelId,
workspaceDir: params.workspaceDir,
store: noExternalAuthStore,
userLockedAuthProfileId:
userPinnedAuthProfileId:
runParams.authProfileIdSource === "user" ? runParams.authProfileId : undefined,
});
}
@@ -102,6 +102,16 @@ describe("resolveAuthProfileFailureReason", () => {
).toBeNull();
});
it("does not persist provider-scoped overload as auth-profile health", () => {
expect(
resolveAuthProfileFailureReason({
failoverReason: "overloaded",
providerStarted: true,
policy: "shared",
}),
).toBeNull();
});
it("does not persist empty responses as auth-profile health", () => {
expect(
resolveAuthProfileFailureReason({
@@ -29,9 +29,12 @@ export function resolveAuthProfileFailureReason(params: {
if (
params.policy === "local" ||
!params.failoverReason ||
// Provider-scoped overload must not cool one credential (#121341 classification).
// Preserve #121278 credential scoping by rotating without a profile-health write.
params.failoverReason === "overloaded" ||
(params.policy === "local_transient" &&
(params.failoverReason === "overloaded" ||
(params.failoverReason === "rate_limit" && params.transientRateLimit === true))) ||
params.failoverReason === "rate_limit" &&
params.transientRateLimit === true) ||
params.failoverReason === "server_error" ||
params.failoverReason === "tls_certificate" ||
params.failoverReason === "empty_response" ||
@@ -30,6 +30,6 @@ export type PreparedEmbeddedRunInput = {
progressController: ReturnType<typeof createEmbeddedRunProgressController>;
laneController: ReturnType<typeof createEmbeddedRunLaneController>;
lifecycleGeneration: NonNullable<RunEmbeddedAgentParams["lifecycleGeneration"]>;
suspendForFailure: (params: Omit<SessionSuspensionParams, "laneId">) => void;
suspendForFailure: (params: SessionSuspensionParams) => void;
preparedModelRuntime?: PreparedModelRuntimeSnapshot;
};
@@ -17,12 +17,11 @@ describe("buildEmbeddedFailureSuspension", () => {
const suspension = buildEmbeddedFailureSuspension({
suspension: { ...baseSuspension, agentDir: "/state/agents/work/agent" },
runAgentId: "work",
laneId: "main",
});
expect(suspension.agentId).toBe("work");
expect(suspension.agentDir).toBe("/state/agents/work/agent");
expect(suspension.laneId).toBe("main");
expect(suspension).not.toHaveProperty("laneId");
});
it("keeps an explicit caller agent id and tolerates a run without one", () => {
@@ -30,7 +29,6 @@ describe("buildEmbeddedFailureSuspension", () => {
buildEmbeddedFailureSuspension({
suspension: { ...baseSuspension, agentId: "explicit" },
runAgentId: "run-owner",
laneId: "main",
}).agentId,
).toBe("explicit");
@@ -38,7 +36,6 @@ describe("buildEmbeddedFailureSuspension", () => {
buildEmbeddedFailureSuspension({
suspension: baseSuspension,
runAgentId: undefined,
laneId: "main",
}).agentId,
).toBeUndefined();
});
@@ -7,15 +7,13 @@
import type { SessionSuspensionParams } from "../../session-suspension.js";
export function buildEmbeddedFailureSuspension(params: {
suspension: Omit<SessionSuspensionParams, "laneId">;
suspension: SessionSuspensionParams;
runAgentId?: string;
laneId: string;
}): SessionSuspensionParams {
return {
...params.suspension,
// A caller-supplied id wins; the run id only fills the gap so an
// unregistered agentDir cannot fall back to the default agent's store.
agentId: params.suspension.agentId ?? params.runAgentId,
laneId: params.laneId,
};
}
@@ -53,7 +53,7 @@ export async function handleEmbeddedPromptFailure(input: {
suspensionSessionId: string;
runtimeAuthRetry: boolean;
maybeRefreshRuntimeAuthForAuthError: (errorText: string, retry: boolean) => Promise<boolean>;
suspendForFailure: (params: Omit<SessionSuspensionParams, "laneId">) => void;
suspendForFailure: (params: SessionSuspensionParams) => void;
resolveReplayInvalid: () => boolean;
setTerminalLifecycleMeta: NonNullable<EmbeddedRunAttemptResult["setTerminalLifecycleMeta"]>;
buildErrorAgentMeta: () => EmbeddedAgentMeta;
@@ -364,15 +364,25 @@ export async function prepareEmbeddedRunRuntime(input: {
log,
});
authStages?.mark("controller");
const cooldownProbePolicy = resolveEmbeddedAuthCooldownProbePolicy({
authStore: attemptAuthProfileStore,
profileCandidates,
lockedProfileId,
modelId,
allowTransientCooldownProbe: params.allowTransientCooldownProbe === true,
});
let didTransientCooldownProbe = false;
const advancePluginHarnessAuthAttempt = async (): Promise<boolean> => {
if (!pluginHarnessOwnsTransport || lockedProfileId) {
if (!pluginHarnessOwnsTransport) {
return false;
}
let nextIndex = profileIndex + 1;
while (nextIndex < preparedAuthAttempts.length) {
const candidateAttempt = preparedAuthAttempts[nextIndex];
const candidateIndex = nextIndex++;
const candidateAttempt = preparedAuthAttempts[candidateIndex];
// Harness-owned auth shares the controller's run-local exhaustion invariant.
profileIndex = candidateIndex;
if (!candidateAttempt) {
nextIndex += 1;
continue;
}
const candidate = candidateAttempt.profileId;
@@ -380,8 +390,13 @@ export async function prepareEmbeddedRunRuntime(input: {
candidate &&
isProfileInCooldown(attemptAuthProfileStore, candidate, undefined, modelId)
) {
nextIndex += 1;
continue;
if (didTransientCooldownProbe || !cooldownProbePolicy.probeProfileIds.has(candidate)) {
continue;
}
didTransientCooldownProbe = true;
log.warn(
`probing cooldowned auth profile for ${provider}/${modelId} due to ${cooldownProbePolicy.unavailableReason ?? "transient"} unavailability`,
);
}
if (
!canRunPreparedAgentRuntimeAuthAttempt({
@@ -389,22 +404,20 @@ export async function prepareEmbeddedRunRuntime(input: {
priorProfileAttempted: preparedProfileAttempted,
})
) {
profileIndex = preparedAuthAttempts.length;
return false;
}
if (candidateAttempt.plan.modelRoute?.authRequirement === "api-key") {
try {
await authController.applyAuthProfileCandidate(candidate, nextIndex);
profileIndex = nextIndex;
await authController.applyAuthProfileCandidate(candidate, candidateIndex);
thinkLevel = initialThinkLevel;
attemptedThinking.clear();
return true;
} catch {
nextIndex += 1;
continue;
}
}
if (!candidate || candidateAttempt.plan.forwardedAuthProfileId !== candidate) {
nextIndex += 1;
continue;
}
const prepared = await prepareAuthAttempt(candidateAttempt);
@@ -412,12 +425,12 @@ export async function prepareEmbeddedRunRuntime(input: {
apiKeyInfo = null;
runtimeAuthState = null;
prepared.commit();
profileIndex = nextIndex;
lastProfileId = candidate;
thinkLevel = initialThinkLevel;
attemptedThinking.clear();
return true;
}
profileIndex = preparedAuthAttempts.length;
return false;
};
const advanceAttemptAuthProfile = pluginHarnessOwnsAuthBootstrap
@@ -426,21 +439,17 @@ export async function prepareEmbeddedRunRuntime(input: {
if (!pluginHarnessOwnsTransport || pluginHarnessNeedsOpenClawAuthBootstrap) {
await authController.initializeAuthProfile();
} else if (lockedProfileId) {
lastProfileId = lockedProfileId;
} else if (forwardedPluginHarnessProfileId) {
const initialAttempt = preparedAuthAttempts[profileIndex];
const initialProfileInCooldown =
initialAttempt?.kind === "profile" &&
isProfileInCooldown(attemptAuthProfileStore, initialAttempt.profileId, undefined, modelId);
const cooldownProbePolicy = resolveEmbeddedAuthCooldownProbePolicy({
authStore: attemptAuthProfileStore,
profileCandidates,
lockedProfileId,
modelId,
allowTransientCooldownProbe: params.allowTransientCooldownProbe === true,
});
if (initialProfileInCooldown && !cooldownProbePolicy.allowProbe) {
const initialProfileId = initialAttempt?.profileId;
const canProbeInitialProfile =
initialProfileInCooldown &&
initialProfileId !== undefined &&
cooldownProbePolicy.probeProfileIds.has(initialProfileId);
if (initialProfileInCooldown && !canProbeInitialProfile) {
if (!(await advancePluginHarnessAuthAttempt())) {
throw new Error(
`Prepared auth profiles are temporarily unavailable for ${provider}/${modelId}.`,
@@ -448,6 +457,7 @@ export async function prepareEmbeddedRunRuntime(input: {
}
} else {
if (initialProfileInCooldown) {
didTransientCooldownProbe = true;
log.warn(
`probing cooldowned auth profile for ${provider}/${modelId} due to ${cooldownProbePolicy.unavailableReason ?? "transient"} unavailability`,
);
-1
View File
@@ -604,7 +604,6 @@ export function throwFallbackFailureSummary(params: {
agentId: params.agentId,
agentDir: params.agentDir,
sessionId: params.attribution.sessionId,
laneId: params.attribution.lane,
reason: "circuit_open",
failedProvider: params.attempts.at(-1)?.provider ?? "unknown",
failedModel: params.attempts.at(-1)?.model ?? "unknown",
+3 -3
View File
@@ -139,7 +139,7 @@ export const probeThrottleInternals = {
type CooldownDecision =
| { type: "skip"; reason: FailoverReason; error: string }
| { type: "attempt"; reason: FailoverReason; markProbe: boolean }
| { type: "suspend_lanes"; reason: FailoverReason; leaderCandidate?: ModelCandidate };
| { type: "suspend_session"; reason: FailoverReason; leaderCandidate?: ModelCandidate };
export function resolveCooldownDecision(params: {
candidate: ModelCandidate;
@@ -188,7 +188,7 @@ export function resolveCooldownDecision(params: {
return { type: "attempt", reason: inferredReason, markProbe: true };
}
return {
type: "suspend_lanes",
type: "suspend_session",
reason: inferredReason,
leaderCandidate: params.candidate,
};
@@ -199,7 +199,7 @@ export function resolveCooldownDecision(params: {
(!params.isPrimary && shouldUseTransientCooldownProbeSlot(inferredReason));
if (!shouldAttemptDespiteCooldown) {
return {
type: "suspend_lanes",
type: "suspend_session",
reason: inferredReason,
leaderCandidate: params.candidate,
};
+14 -12
View File
@@ -205,8 +205,6 @@ async function runWithModelFallbackInternal<T>(
let exhaustionResult: ModelFallbackExhaustionResult<T> | undefined;
const cooldownProbeUsedProviders = new Set<string>();
const tlsFailedProviders = new Set<string>();
const resolveTerminalSuspensionLane = () =>
deferredSuspension.pending ? deferredSuspension.pending.laneId : params.lane;
const observeDecision = async (decision: ModelFallbackDecisionParams) => {
if (!params.onFallbackStep && !isModelFallbackDecisionLogEnabled()) {
return;
@@ -322,12 +320,19 @@ async function runWithModelFallbackInternal<T>(
profileId: userLockedAuthProfileId,
}).eligible;
if (!candidateHarnessAuth.skipsProviderAuthCooldown) {
candidateAuthProfileIds = authRuntime.resolveAuthProfileOrder({
const orderedProfileIds = authRuntime.resolveAuthProfileOrder({
cfg: params.cfg,
store: authStore,
provider: candidate.provider,
forModel: candidate.model,
});
candidateAuthProfileIds =
userLockedAuthProfileEligible && userLockedAuthProfileId
? [
userLockedAuthProfileId,
...orderedProfileIds.filter((profileId) => profileId !== userLockedAuthProfileId),
]
: orderedProfileIds;
authRuntime.maybeReprobeWhamBlockedProfiles({
store: authStore,
profileIds: candidateAuthProfileIds,
@@ -388,7 +393,7 @@ async function runWithModelFallbackInternal<T>(
(id) => !authRuntime.isProfileInCooldown(authStore, id, undefined, candidate.model),
);
if (profileIds.length > 0 && !isAnyProfileAvailable && !userLockedAuthProfileEligible) {
if (profileIds.length > 0 && !isAnyProfileAvailable) {
// All profiles for this provider are in cooldown.
const now = Date.now();
const probeThrottleKey = resolveProbeThrottleKey(candidate.provider, params.agentDir);
@@ -408,13 +413,12 @@ async function runWithModelFallbackInternal<T>(
? resolveSubscriptionAuthModeForProfiles({ store: authStore, profileIds })
: undefined;
if (decision.type === "suspend_lanes") {
const error = `Provider ${candidate.provider} is in cooldown (suspending lanes)`;
if (decision.type === "suspend_session") {
const error = `Provider ${candidate.provider} is in cooldown`;
pushAttempt(error, decision.reason, { authMode });
// Only lock the lane when no remaining candidates can serve as
// fallbacks. Per-provider cooldown state already prevents
// re-attempting the failed provider on subsequent turns.
// Only record terminal session suspension when no remaining candidate
// can serve the turn. Provider cooldown state prevents repeat probes.
const hasRemainingCandidates = hasRemainingCandidate;
if (params.sessionId) {
emitFailoverEvent({
@@ -426,14 +430,12 @@ async function runWithModelFallbackInternal<T>(
suspended: !hasRemainingCandidates,
});
if (!hasRemainingCandidates) {
const laneId = resolveTerminalSuspensionLane();
deferredSuspension.pending = undefined;
void suspendSession({
cfg: params.cfg,
agentId: params.agentId,
agentDir: params.agentDir,
sessionId: params.sessionId,
laneId,
reason: resolveSessionSuspensionReason(decision.reason),
failedProvider: candidate.provider,
failedModel: candidate.model,
@@ -767,7 +769,7 @@ async function runWithModelFallbackInternal<T>(
cfg: params.cfg,
candidates,
}),
attribution: { sessionId: params.sessionId, lane: resolveTerminalSuspensionLane() },
attribution: { sessionId: params.sessionId, lane: params.lane },
cfg: params.cfg,
agentId: params.agentId,
agentDir: params.agentDir,
+12 -14
View File
@@ -39,7 +39,6 @@ const sessionSuspensionMocks = vi.hoisted(() => ({
onDeferred?.({
cfg: {},
sessionId: "test-session",
laneId: "main",
reason: "quota_exhausted",
failedProvider: "openai",
failedModel: "gpt-4.1-mini",
@@ -313,7 +312,7 @@ describe("runWithModelFallback probe logic", () => {
reason: "rate_limit" | "billing",
) {
expect(decision).toEqual({
type: "suspend_lanes",
type: "suspend_session",
reason,
leaderCandidate: OPENAI_PROBE_CANDIDATE,
});
@@ -837,7 +836,7 @@ describe("runWithModelFallback probe logic", () => {
);
});
it("does not lock lane when fallback candidates remain after suspend_lanes decision", async () => {
it("does not suspend the session when fallback candidates remain", async () => {
const cfg = makeCfg({
agents: {
defaults: {
@@ -870,7 +869,7 @@ describe("runWithModelFallback probe logic", () => {
expect(sessionSuspensionMocks.suspendSession).not.toHaveBeenCalled();
});
it("defers embedded lane suspension only while another candidate remains", async () => {
it("defers embedded session suspension only while another candidate remains", async () => {
const cfg = makeCfg({
agents: {
defaults: {
@@ -964,7 +963,7 @@ describe("runWithModelFallback probe logic", () => {
return [];
});
// Throttle primary probe so billing goes to suspend_lanes
// Throttle primary probe so billing records terminal session suspension.
probeThrottleInternals.lastProbeAttempt.set("openai", NOW - 10_000);
const run = vi.fn().mockResolvedValue("should-not-run");
@@ -981,21 +980,16 @@ describe("runWithModelFallback probe logic", () => {
expect(sessionSuspensionMocks.suspendSession).toHaveBeenCalledWith(
expect.objectContaining({
laneId: undefined,
failedProvider: "anthropic",
}),
);
expect(sessionSuspensionMocks.suspendSession).not.toHaveBeenCalledWith(
expect.objectContaining({ failedProvider: "openai" }),
);
expect(
sessionSuspensionMocks.suspendSession.mock.calls.every(
([params]) => params.laneId === undefined,
),
).toBe(true);
expect(sessionSuspensionMocks.suspendSession.mock.calls[0]?.[0]).not.toHaveProperty("laneId");
});
it("restores a deferred embedded lane when later candidates cannot run", async () => {
it("records the final candidate when later candidates cannot run", async () => {
const cfg = makeCfg({
agents: {
defaults: {
@@ -1029,10 +1023,12 @@ describe("runWithModelFallback probe logic", () => {
expect(run).toHaveBeenCalledOnce();
expect(sessionSuspensionMocks.suspendSession).toHaveBeenCalledWith(
expect.objectContaining({
laneId: "main",
failedProvider: "anthropic",
}),
);
expect(sessionSuspensionMocks.suspendSession.mock.calls.at(-1)?.[0]).not.toHaveProperty(
"laneId",
);
});
it("restores deferred suspension when a later harness precheck fails", async () => {
@@ -1065,9 +1061,11 @@ describe("runWithModelFallback probe logic", () => {
expect(run).toHaveBeenCalledOnce();
expect(sessionSuspensionMocks.suspendSession).toHaveBeenCalledWith(
expect.objectContaining({
laneId: "main",
failedProvider: "openai",
}),
);
expect(sessionSuspensionMocks.suspendSession.mock.calls.at(-1)?.[0]).not.toHaveProperty(
"laneId",
);
});
});
@@ -626,7 +626,7 @@ describe("runWithModelFallback + runEmbeddedAgent failover behavior", () => {
}
});
it("keeps direct embedded-run lane suspension outside the outer fallback loop", async () => {
it("keeps direct embedded-run session suspension outside the outer fallback loop", async () => {
await withAgentWorkspace(async ({ agentDir, workspaceDir }) => {
await writeAuthStore(agentDir);
const sessionId = "session:direct-embedded-suspension";
@@ -652,9 +652,8 @@ describe("runWithModelFallback + runEmbeddedAgent failover behavior", () => {
}),
).rejects.toThrow();
expect(suspendSessionMock).toHaveBeenCalledWith(
expect.objectContaining({ laneId: "direct-lane" }),
);
expect(suspendSessionMock).toHaveBeenCalledOnce();
expect(suspendSessionMock.mock.calls[0]?.[0]).not.toHaveProperty("laneId");
});
});
+30
View File
@@ -3462,6 +3462,36 @@ describe("runWithModelFallback", () => {
expect(store.order?.[provider]).toEqual(orderedProfileIds);
});
it("does not skip a provider when only its user-pinned profile is cooling down", async () => {
const provider = `pinned-cooldown-${crypto.randomUUID()}`;
const pinnedProfileId = `${provider}:pinned`;
const backupProfileId = `${provider}:backup`;
const store: AuthProfileStore = {
version: AUTH_STORE_VERSION,
profiles: {
[pinnedProfileId]: { type: "api_key", provider, key: "pinned-key" },
[backupProfileId]: { type: "api_key", provider, key: "backup-key" },
"fallback:default": { type: "api_key", provider: "fallback", key: "fallback-key" },
},
order: { [provider]: [backupProfileId] },
usageStats: {
[pinnedProfileId]: { cooldownUntil: Date.now() + 60_000 },
},
};
const run = vi.fn().mockResolvedValue("ok");
const result = await runWithStoredAuth({
cfg: makeProviderFallbackCfg(provider),
store,
provider,
run,
userLockedAuthProfileId: pinnedProfileId,
});
expect(result.result).toBe("ok");
expect(run.mock.calls).toEqual([[provider, "m1", { isFinalFallbackAttempt: false }]]);
});
it("discovers an exact external CLI user lock before cooldown admission", async () => {
const provider = "minimax-portal";
const orderedProfileId = "minimax-portal:api";
+54 -8
View File
@@ -356,7 +356,7 @@ describe("prepareAgentRuntimeAuthPlan", () => {
).toThrow(/explicit auth order.*no usable profiles/iu);
});
it("keeps a generic user lock as a singleton despite cooldown", () => {
it("skips a cooldowned user pin and selects the next same-provider profile", () => {
const store = authStore(
{
"xai:p1": apiKeyProfile("xai", "p1-key"),
@@ -376,13 +376,37 @@ describe("prepareAgentRuntimeAuthPlan", () => {
});
expect(plan).toMatchObject({
forwardedAuthProfileId: "xai:p1",
forwardedAuthProfileSource: "user",
forwardedAuthProfileCandidateIds: ["xai:p1"],
forwardedAuthProfileId: "xai:p2",
forwardedAuthProfileSource: "auto",
forwardedAuthProfileCandidateIds: ["xai:p2"],
selectedAuthMode: "api_key",
});
});
it("prepares a user pin first and retains same-provider profile fallbacks", () => {
const prepared = prepareAgentRuntimeAuth({
provider: "xai",
modelId: "grok-4",
env: {},
authProfileStore: authStore(
{
"xai:p1": apiKeyProfile("xai", "p1-key"),
"xai:p2": apiKeyProfile("xai", "p2-key"),
},
{ xai: ["xai:p2", "xai:p1"] },
),
sessionAuthProfileId: "xai:p1",
sessionAuthProfileSource: "user",
});
const profileAttempts = prepared.attempts.filter((attempt) => attempt.kind === "profile");
expect(profileAttempts.map((attempt) => attempt.profileId)).toEqual(["xai:p1", "xai:p2"]);
expect(profileAttempts.map((attempt) => attempt.plan.forwardedAuthProfileSource)).toEqual([
"user",
"auto",
]);
});
it("defers an ambiguous route when native Codex owns auth", () => {
const plan = prepareAgentRuntimeAuthPlan({
...openAIChatGptAuthFixture(),
@@ -778,7 +802,7 @@ describe("prepareAgentRuntimeAuthPlan", () => {
).toThrow(/explicit auth order.*no usable profiles/iu);
});
it("keeps a user-locked profile authoritative and rejects the wrong route class", () => {
it("does not cross to an incompatible auth route for a user pin", () => {
expect(() =>
prepareAgentRuntimeAuthPlan({
...openAIChatGptAuthFixture(),
@@ -800,7 +824,7 @@ describe("prepareAgentRuntimeAuthPlan", () => {
"openai:platform": openAIApiKeyProfile("platform-key"),
}),
}),
).toThrow(/requires subscription authentication/u);
).toThrow(/no route-compatible authentication source/iu);
});
it("lets an explicit provider API key outrank automatic subscription profiles", () => {
@@ -1739,7 +1763,29 @@ describe("prepareAgentRuntimeAuthPlan", () => {
expect(plan.modelRoute).toBeUndefined();
});
it("rejects a user-locked non-OpenAI profile on the virtual Codex provider", () => {
it("keeps same-provider retries behind a user-pinned virtual Codex profile", () => {
const preparation = prepareAgentRuntimeAuth({
...virtualCodexAuthFixture(),
authProfileStore: authStore(
{
"openai:p1": openAITokenProfile("p1-token"),
"openai:p2": openAIApiKeyProfile("p2-key"),
},
{ openai: ["openai:p2", "openai:p1"] },
),
sessionAuthProfileId: "openai:p1",
sessionAuthProfileSource: "user",
});
const profileAttempts = preparation.attempts.filter((attempt) => attempt.kind === "profile");
expect(profileAttempts.map((attempt) => attempt.profileId)).toEqual(["openai:p1", "openai:p2"]);
expect(profileAttempts.map((attempt) => attempt.plan.forwardedAuthProfileSource)).toEqual([
"user",
"auto",
]);
});
it("rejects a user-pinned non-OpenAI profile on the virtual Codex provider", () => {
expect(() =>
prepareAgentRuntimeAuthPlan({
...virtualCodexAuthFixture(),
@@ -1752,7 +1798,7 @@ describe("prepareAgentRuntimeAuthPlan", () => {
).toThrow(/not configured for openai/u);
});
it("rejects unavailable user-locked OpenAI profiles on the virtual Codex provider", () => {
it("rejects unavailable user-pinned OpenAI profiles on the virtual Codex provider", () => {
expect(() =>
prepareAgentRuntimeAuthPlan({
...virtualCodexAuthFixture(),
+33 -20
View File
@@ -114,9 +114,6 @@ export function preparedAgentRuntimeProfileAttemptHasCandidate(params: {
if (params.attempt.kind !== "profile") {
return false;
}
if (params.attempt.plan.forwardedAuthProfileSource === "user") {
return true;
}
const profileIds = params.attempt.plan.forwardedAuthProfileCandidateIds ?? [
params.attempt.profileId,
];
@@ -211,7 +208,7 @@ export function prepareAgentRuntimeAuth(
params: PrepareAgentRuntimeAuthPlanParams,
): PreparedAgentRuntimeAuth {
const requestedProfileId = params.sessionAuthProfileId?.trim() || undefined;
const lockedProfileId =
const userPinnedProfileId =
params.sessionAuthProfileSource === "user" ? requestedProfileId : undefined;
const harnessOwnsOpenAIAuth =
params.harnessId?.trim().toLowerCase() === "codex" ||
@@ -222,38 +219,40 @@ export function prepareAgentRuntimeAuth(
? { id: harnessAuthOwnerId }
: undefined;
const harnessAllowsAuthProfileForwarding = params.allowHarnessAuthProfileForwarding !== false;
if (lockedProfileId && !harnessAllowsAuthProfileForwarding) {
if (userPinnedProfileId && !harnessAllowsAuthProfileForwarding) {
throw new Error(
`Auth profile "${lockedProfileId}" cannot be forwarded to the selected agent harness. Configure that harness's native account instead.`,
`Auth profile "${userPinnedProfileId}" cannot be forwarded to the selected agent harness. Configure that harness's native account instead.`,
);
}
const store = params.authProfileStore;
const authProfileSelectionProvider = harnessOwnsOpenAIAuth ? "openai" : params.provider;
if (lockedProfileId) {
if (userPinnedProfileId) {
const eligibility = store
? resolveAuthProfileEligibility({
cfg: params.config,
store,
provider: authProfileSelectionProvider,
profileId: lockedProfileId,
profileId: userPinnedProfileId,
})
: { eligible: false };
if (!eligibility.eligible) {
throw new Error(
`Auth profile "${lockedProfileId}" is not configured for ${authProfileSelectionProvider}.`,
`Auth profile "${userPinnedProfileId}" is not configured for ${authProfileSelectionProvider}.`,
);
}
}
const configuredProvider = resolveMergedModelProviderConfig(params.config, params.provider);
const configuredAuthMode =
lockedProfileId || !harnessAllowsAuthProfileForwarding ? undefined : configuredProvider?.auth;
userPinnedProfileId || !harnessAllowsAuthProfileForwarding
? undefined
: configuredProvider?.auth;
const configuredAwsSdkAuth = configuredAuthMode === "aws-sdk";
const providerHasApiKeySecretRef =
harnessAllowsAuthProfileForwarding &&
Boolean(coerceSecretRef(configuredProvider?.apiKey, params.config?.secrets?.defaults));
const providerBinding =
harnessAllowsAuthProfileForwarding && !lockedProfileId && store && !configuredAwsSdkAuth
harnessAllowsAuthProfileForwarding && !userPinnedProfileId && store && !configuredAwsSdkAuth
? resolvePreparedProviderEntryApiKeyProfileReference({
config: params.config,
modelId: params.modelId,
@@ -286,8 +285,8 @@ export function prepareAgentRuntimeAuth(
// Explicit auth owns the physical route; apiKey is only its bearer material.
const selectedConfiguredAuthMode =
configuredAuthMode ?? (providerHasDirectMaterial ? "api-key" : undefined);
const selectedProfileId = lockedProfileId ?? boundProfileId;
const automaticOrderResolution =
const selectedProfileId = boundProfileId;
const resolvedAutomaticOrder =
!harnessAllowsAuthProfileForwarding ||
selectedProfileId ||
providerBindingSuppressesProfiles ||
@@ -301,13 +300,25 @@ export function prepareAgentRuntimeAuth(
cfg: params.config,
store,
provider: authProfileSelectionProvider,
preferredProfile: lockedProfileId ? undefined : requestedProfileId,
preferredProfile: requestedProfileId,
forModel: params.modelId,
readinessMode: "read-only",
});
const automaticOrderResolution = userPinnedProfileId
? {
...resolvedAutomaticOrder,
profileIds: [
userPinnedProfileId,
...resolvedAutomaticOrder.profileIds.filter(
(profileId) => profileId !== userPinnedProfileId,
),
],
}
: resolvedAutomaticOrder;
const providerPreferredProfileId =
harnessAllowsAuthProfileForwarding &&
!selectedProfileId &&
!userPinnedProfileId &&
!providerBindingSuppressesProfiles &&
!configuredAwsSdkAuth &&
store
@@ -317,8 +328,8 @@ export function prepareAgentRuntimeAuth(
workspaceDir: params.workspaceDir,
provider: params.provider,
modelId: params.modelId,
preferredProfileId: lockedProfileId ? undefined : requestedProfileId,
lockedProfileId,
preferredProfileId: requestedProfileId,
lockedProfileId: undefined,
profileOrder: automaticOrderResolution.profileIds,
authStore: store,
})
@@ -384,7 +395,7 @@ export function prepareAgentRuntimeAuth(
: selectedConfiguredAuthMode;
const ownership = selectedProfileId
? {
reason: lockedProfileId ? ("user-lock" as const) : ("provider-binding" as const),
reason: "provider-binding" as const,
source: resolveProfile(params, selectedProfileId, { ignoreCooldown: true }),
}
: configuredAwsSdkAuth
@@ -401,7 +412,9 @@ export function prepareAgentRuntimeAuth(
const sourcePlan = buildProviderModelAuthSourcePlan({
...(ownership ? { ownership } : {}),
profiles: resolvedOrderedProfileIds.map((profileId) => resolveProfile(params, profileId)),
...(providerPreferredProfileId ? { preferredProfileId: providerPreferredProfileId } : {}),
...(userPinnedProfileId || providerPreferredProfileId
? { preferredProfileId: userPinnedProfileId ?? providerPreferredProfileId }
: {}),
explicitOrder: automaticOrderResolution.hasExplicitOrder,
...(fallbackDirectSource ? { fallback: fallbackDirectSource } : {}),
allowCooldown: params.allowTransientCooldownProbe,
@@ -445,7 +458,7 @@ export function prepareAgentRuntimeAuth(
(attempt?.kind === "direct" ? attempt.source.mode : selectedConfiguredAuthMode),
sessionAuthProfileId: profile?.profileId,
sessionAuthProfileSource: profile
? sourcePlan.kind === "required" && sourcePlan.reason === "user-lock"
? profile.profileId === userPinnedProfileId
? "user"
: "auto"
: undefined,
@@ -542,7 +555,7 @@ export function prepareAgentRuntimeAuth(
(attempt?.kind === "direct" ? attempt.source.mode : selectedConfiguredAuthMode),
sessionAuthProfileId: profile?.profileId,
sessionAuthProfileSource: profile
? sourcePlan.kind === "required" && sourcePlan.reason === "user-lock"
? profile.profileId === userPinnedProfileId
? "user"
: "auto"
: undefined,
@@ -2,10 +2,6 @@ import "./session-suspension.js";
type SessionSuspensionTestApi = {
resetSessionSuspensionStateForTest(): void;
seedClearedLaneResumeForTest(
laneId: string,
cleared: { resumeConcurrency: number; resumeAtMs: number },
): void;
};
function getTestApi(): SessionSuspensionTestApi {
@@ -21,10 +17,3 @@ function getTestApi(): SessionSuspensionTestApi {
export function resetSessionSuspensionStateForTest(): void {
getTestApi().resetSessionSuspensionStateForTest();
}
export function seedClearedLaneResumeForTest(
laneId: string,
cleared: { resumeConcurrency: number; resumeAtMs: number },
): void {
getTestApi().seedClearedLaneResumeForTest(laneId, cleared);
}
+103 -411
View File
@@ -1,7 +1,8 @@
// Verifies quota suspension persists lane state and auto-resumes safely.
// Verifies quota suspension records recovery state without blocking shared work.
import { afterEach, describe, expect, it, vi } from "vitest";
import { DEFAULT_CRON_MAX_CONCURRENT_RUNS } from "../config/cron-limits.js";
import type { OpenClawConfig } from "../config/types.openclaw.js";
import { enqueueCommandInLane, getCommandLaneSnapshot } from "../process/command-queue.js";
import { resetCommandQueueStateForTest } from "../process/command-queue.test-support.js";
import { CommandLane } from "../process/lanes.js";
import { MAX_TIMER_TIMEOUT_MS } from "../shared/number-coercion.js";
@@ -9,14 +10,8 @@ const sessionAccessorMocks = vi.hoisted(() => ({
patchSessionEntryCore: vi.fn(),
}));
const commandQueueMocks = vi.hoisted(() => ({
setCommandLaneConcurrency: vi.fn(),
}));
vi.mock("../config/sessions/session-accessor.js", () => sessionAccessorMocks);
vi.mock("../process/command-queue.js", () => commandQueueMocks);
const sessionKeyResolverMocks = vi.hoisted(() => ({
resolveStoredSessionKeyForSessionId: vi.fn(() => ({
sessionKey: "session-key",
@@ -26,37 +21,78 @@ const sessionKeyResolverMocks = vi.hoisted(() => ({
vi.mock("./command/session.js", () => sessionKeyResolverMocks);
async function suspendLane(ttlMs: number, cfg: OpenClawConfig, laneId: CommandLane) {
// All cases exercise the public suspendSession path with fixed failure metadata.
async function recordSuspension(ttlMs = 100) {
const { suspendSession } = await import("./session-suspension.js");
await suspendSession({
cfg,
cfg: {} as OpenClawConfig,
sessionId: "session-1",
laneId,
reason: "quota_exhausted",
failedProvider: "anthropic",
failedModel: "claude-opus-4-6",
failedProvider: "openai",
failedModel: "gpt-5.6-sol",
ttlMs,
});
}
describe("session suspension", () => {
afterEach(async () => {
if (vi.isFakeTimers()) {
await vi.runOnlyPendingTimersAsync();
vi.clearAllTimers();
}
vi.useRealTimers();
const { resetSessionSuspensionStateForTest } =
await import("./session-suspension.test-support.js");
resetSessionSuspensionStateForTest();
sessionAccessorMocks.patchSessionEntryCore.mockClear();
commandQueueMocks.setCommandLaneConcurrency.mockClear();
resetCommandQueueStateForTest();
vi.useRealTimers();
vi.restoreAllMocks();
sessionAccessorMocks.patchSessionEntryCore.mockReset();
sessionKeyResolverMocks.resolveStoredSessionKeyForSessionId.mockClear();
});
it("records a bounded recovery marker without pausing the shared main lane", async () => {
vi.useFakeTimers();
vi.setSystemTime(1_000);
sessionAccessorMocks.patchSessionEntryCore.mockImplementation(async (_scope, update) =>
update({}),
);
await recordSuspension(Number.MAX_SAFE_INTEGER);
const buildPatch = sessionAccessorMocks.patchSessionEntryCore.mock.calls[0]?.[1] as (_entry: {
quotaSuspension?: unknown;
}) => {
quotaSuspension?: {
expectedResumeBy?: number;
failedProvider?: string;
failedModel?: string;
state?: string;
};
};
expect(buildPatch({}).quotaSuspension).toEqual(
expect.objectContaining({
expectedResumeBy: 1_000 + MAX_TIMER_TIMEOUT_MS,
failedProvider: "openai",
failedModel: "gpt-5.6-sol",
state: "suspended",
}),
);
expect(getCommandLaneSnapshot(CommandLane.Main).maxConcurrent).toBe(1);
await expect(
enqueueCommandInLane(CommandLane.Main, async () => "unrelated-provider-ok"),
).resolves.toBe("unrelated-provider-ok");
});
it("keeps the shared lane runnable when marker persistence fails", async () => {
sessionAccessorMocks.patchSessionEntryCore.mockRejectedValueOnce(new Error("disk busy"));
await recordSuspension();
await expect(enqueueCommandInLane(CommandLane.Main, async () => "still-runs")).resolves.toBe(
"still-runs",
);
});
it("resolves the session store with the explicit agent id, never the agentDir basename", async () => {
const { suspendSession } = await import("./session-suspension.js");
sessionKeyResolverMocks.resolveStoredSessionKeyForSessionId.mockClear();
sessionAccessorMocks.patchSessionEntryCore.mockImplementation(async (_scope, update) =>
update({}),
);
await suspendSession({
cfg: {} as OpenClawConfig,
@@ -64,11 +100,9 @@ describe("session suspension", () => {
// Default layout: <state>/agents/<id>/agent — basename is always "agent".
agentDir: "/state/agents/work/agent",
sessionId: "session-1",
laneId: CommandLane.Main,
reason: "quota_exhausted",
failedProvider: "anthropic",
failedModel: "claude-opus-4-6",
ttlMs: 1,
failedProvider: "openai",
failedModel: "gpt-5.6-sol",
});
expect(sessionKeyResolverMocks.resolveStoredSessionKeyForSessionId).toHaveBeenCalledWith(
@@ -80,7 +114,9 @@ describe("session suspension", () => {
const { suspendSession } = await import("./session-suspension.js");
const { registerResolvedAgentDir, unregisterResolvedAgentDir } =
await import("./agent-dir-registry.js");
sessionKeyResolverMocks.resolveStoredSessionKeyForSessionId.mockClear();
sessionAccessorMocks.patchSessionEntryCore.mockImplementation(async (_scope, update) =>
update({}),
);
registerResolvedAgentDir({ agentId: "research", agentDir: "/state/agents/research/agent" });
try {
@@ -88,11 +124,9 @@ describe("session suspension", () => {
cfg: {} as OpenClawConfig,
agentDir: "/state/agents/research/agent",
sessionId: "session-2",
laneId: CommandLane.Main,
reason: "quota_exhausted",
failedProvider: "anthropic",
failedModel: "claude-opus-4-6",
ttlMs: 1,
failedProvider: "openai",
failedModel: "gpt-5.6-sol",
});
} finally {
unregisterResolvedAgentDir({
@@ -106,398 +140,55 @@ describe("session suspension", () => {
);
});
it("auto-resumes main lane to configured agent concurrency", async () => {
vi.useFakeTimers();
const cfg = {
agents: { defaults: { maxConcurrent: 4 } },
} as OpenClawConfig;
await suspendLane(100, cfg, CommandLane.Main);
expect(commandQueueMocks.setCommandLaneConcurrency).toHaveBeenCalledWith(CommandLane.Main, 0);
await vi.advanceTimersByTimeAsync(100);
expect(commandQueueMocks.setCommandLaneConcurrency).toHaveBeenLastCalledWith(
CommandLane.Main,
4,
);
});
it("auto-resumes cron lanes to the cron concurrency default", async () => {
vi.useFakeTimers();
await suspendLane(100, {} as OpenClawConfig, CommandLane.CronNested);
expect(commandQueueMocks.setCommandLaneConcurrency).toHaveBeenCalledWith(
CommandLane.CronNested,
0,
);
await vi.advanceTimersByTimeAsync(100);
expect(commandQueueMocks.setCommandLaneConcurrency).toHaveBeenLastCalledWith(
CommandLane.CronNested,
DEFAULT_CRON_MAX_CONCURRENT_RUNS,
);
});
it("auto-resumes hook dispatch to the shared cron concurrency width", async () => {
vi.useFakeTimers();
await suspendLane(100, {} as OpenClawConfig, CommandLane.HookDispatch);
expect(commandQueueMocks.setCommandLaneConcurrency).toHaveBeenCalledWith(
CommandLane.HookDispatch,
0,
);
await vi.advanceTimersByTimeAsync(100);
expect(commandQueueMocks.setCommandLaneConcurrency).toHaveBeenLastCalledWith(
CommandLane.HookDispatch,
DEFAULT_CRON_MAX_CONCURRENT_RUNS,
);
});
it("retargets a suspended hook lane when hooks are disabled before its TTL", async () => {
vi.useFakeTimers();
const { getSuspendedLaneIdsForGatewayPublication, setGatewayLaneResumeConcurrencies } =
it("rolls back a write that finishes after gateway shutdown begins", async () => {
const { fenceSessionSuspensionWritesForGatewayShutdown } =
await import("./session-suspension.js");
await suspendLane(100, {} as OpenClawConfig, CommandLane.HookDispatch);
setGatewayLaneResumeConcurrencies({ [CommandLane.HookDispatch]: 0 });
expect(getSuspendedLaneIdsForGatewayPublication()).toEqual(new Set([CommandLane.HookDispatch]));
commandQueueMocks.setCommandLaneConcurrency.mockClear();
await vi.advanceTimersByTimeAsync(100);
expect(commandQueueMocks.setCommandLaneConcurrency).toHaveBeenCalledExactlyOnceWith(
CommandLane.HookDispatch,
0,
);
});
it("uses hooks-off concurrency when a pending suspension write finishes late", async () => {
vi.useFakeTimers();
const { setGatewayLaneResumeConcurrencies } = await import("./session-suspension.js");
let resolvePatch: (() => void) | undefined;
sessionAccessorMocks.patchSessionEntryCore.mockImplementationOnce(async (_scope, update) => {
await new Promise<void>((resolve) => {
resolvePatch = resolve;
});
return update({});
});
const suspension = suspendLane(100, {} as OpenClawConfig, CommandLane.HookDispatch);
await vi.waitFor(() => {
expect(resolvePatch).toBeTypeOf("function");
});
setGatewayLaneResumeConcurrencies({ [CommandLane.HookDispatch]: 0 });
resolvePatch?.();
await suspension;
expect(commandQueueMocks.setCommandLaneConcurrency).toHaveBeenLastCalledWith(
CommandLane.HookDispatch,
0,
);
commandQueueMocks.setCommandLaneConcurrency.mockClear();
await vi.advanceTimersByTimeAsync(100);
expect(commandQueueMocks.setCommandLaneConcurrency).toHaveBeenCalledExactlyOnceWith(
CommandLane.HookDispatch,
0,
);
});
it("clamps oversized suspension TTLs for timers and persisted resume time", async () => {
// Persisted expectedResumeBy must match the clamped timer, not MAX_SAFE_INTEGER.
vi.useFakeTimers();
vi.setSystemTime(1_000);
const setTimeoutSpy = vi.spyOn(globalThis, "setTimeout");
await suspendLane(Number.MAX_SAFE_INTEGER, {} as OpenClawConfig, CommandLane.Main);
expect(setTimeoutSpy).toHaveBeenCalledWith(expect.any(Function), MAX_TIMER_TIMEOUT_MS);
const buildPatch = sessionAccessorMocks.patchSessionEntryCore.mock.calls[0]?.[1] as (_entry: {
quotaSuspension?: unknown;
}) => {
quotaSuspension?: { expectedResumeBy?: number };
};
const patch = buildPatch({});
expect(patch.quotaSuspension?.expectedResumeBy).toBe(1_000 + MAX_TIMER_TIMEOUT_MS);
});
it("clears pending lane auto-resume timers without pumping queued work during cleanup", async () => {
vi.useFakeTimers();
const { clearSessionSuspensionTimers } = await import("./session-suspension.js");
await suspendLane(
100,
{ agents: { defaults: { maxConcurrent: 3 } } } as OpenClawConfig,
CommandLane.Main,
);
expect(commandQueueMocks.setCommandLaneConcurrency).toHaveBeenCalledWith(CommandLane.Main, 0);
expect(clearSessionSuspensionTimers()).toBe(1);
commandQueueMocks.setCommandLaneConcurrency.mockClear();
await vi.advanceTimersByTimeAsync(100);
expect(commandQueueMocks.setCommandLaneConcurrency).not.toHaveBeenCalled();
expect(clearSessionSuspensionTimers()).toBe(0);
});
it("blocks new suspension timers until gateway startup re-enables them", async () => {
vi.useFakeTimers();
const { clearSessionSuspensionTimers, enableSessionSuspensionTimersForGatewayStart } =
await import("./session-suspension.js");
await suspendLane(100, {} as OpenClawConfig, CommandLane.Nested);
expect(commandQueueMocks.setCommandLaneConcurrency).toHaveBeenCalledWith(CommandLane.Nested, 0);
expect(clearSessionSuspensionTimers()).toBe(1);
commandQueueMocks.setCommandLaneConcurrency.mockClear();
sessionAccessorMocks.patchSessionEntryCore.mockClear();
await suspendLane(100, {} as OpenClawConfig, CommandLane.Nested);
expect(commandQueueMocks.setCommandLaneConcurrency).not.toHaveBeenCalled();
expect(sessionAccessorMocks.patchSessionEntryCore).not.toHaveBeenCalled();
enableSessionSuspensionTimersForGatewayStart();
await suspendLane(100, {} as OpenClawConfig, CommandLane.Nested);
expect(commandQueueMocks.setCommandLaneConcurrency).toHaveBeenCalledWith(CommandLane.Nested, 0);
});
it("restores suspended custom lanes when gateway startup re-enables timers", async () => {
vi.useFakeTimers();
const { clearSessionSuspensionTimers, enableSessionSuspensionTimersForGatewayStart } =
await import("./session-suspension.js");
const customLaneId = "plugin:voice:room-1" as CommandLane;
await suspendLane(100, {} as OpenClawConfig, customLaneId);
expect(commandQueueMocks.setCommandLaneConcurrency).toHaveBeenCalledWith(customLaneId, 0);
expect(clearSessionSuspensionTimers()).toBe(1);
commandQueueMocks.setCommandLaneConcurrency.mockClear();
await vi.advanceTimersByTimeAsync(100);
expect(commandQueueMocks.setCommandLaneConcurrency).not.toHaveBeenCalled();
expect(enableSessionSuspensionTimersForGatewayStart().size).toBe(0);
expect(commandQueueMocks.setCommandLaneConcurrency).toHaveBeenCalledWith(customLaneId, 1);
expect(enableSessionSuspensionTimersForGatewayStart().size).toBe(0);
});
it("reschedules unexpired custom lane suspensions when gateway startup re-enables timers", async () => {
vi.useFakeTimers();
vi.setSystemTime(1_000);
const { clearSessionSuspensionTimers, enableSessionSuspensionTimersForGatewayStart } =
await import("./session-suspension.js");
const customLaneId = "plugin:voice:room-2" as CommandLane;
await suspendLane(100, {} as OpenClawConfig, customLaneId);
expect(clearSessionSuspensionTimers()).toBe(1);
commandQueueMocks.setCommandLaneConcurrency.mockClear();
await vi.advanceTimersByTimeAsync(40);
const suspendedLaneIds = enableSessionSuspensionTimersForGatewayStart();
expect(suspendedLaneIds).toEqual(new Set([customLaneId]));
expect(commandQueueMocks.setCommandLaneConcurrency).toHaveBeenCalledWith(customLaneId, 0);
commandQueueMocks.setCommandLaneConcurrency.mockClear();
await vi.advanceTimersByTimeAsync(59);
expect(commandQueueMocks.setCommandLaneConcurrency).not.toHaveBeenCalled();
await vi.advanceTimersByTimeAsync(1);
expect(commandQueueMocks.setCommandLaneConcurrency).toHaveBeenCalledWith(customLaneId, 1);
});
it("leaves built-in lane restoration to gateway startup concurrency", async () => {
vi.useFakeTimers();
const { clearSessionSuspensionTimers, enableSessionSuspensionTimersForGatewayStart } =
await import("./session-suspension.js");
await suspendLane(
100,
{ agents: { defaults: { maxConcurrent: 3 } } } as OpenClawConfig,
CommandLane.Main,
);
expect(clearSessionSuspensionTimers()).toBe(1);
commandQueueMocks.setCommandLaneConcurrency.mockClear();
expect(enableSessionSuspensionTimersForGatewayStart()).toEqual(new Set([CommandLane.Main]));
expect(commandQueueMocks.setCommandLaneConcurrency).toHaveBeenCalledWith(CommandLane.Main, 0);
});
it("clamps rescheduled cleanup timers after wall-clock rollback", async () => {
vi.useFakeTimers();
vi.setSystemTime(1_000);
const { enableSessionSuspensionTimersForGatewayStart } =
await import("./session-suspension.js");
const { seedClearedLaneResumeForTest } = await import("./session-suspension.test-support.js");
const setTimeoutSpy = vi.spyOn(globalThis, "setTimeout");
const customLaneId = "plugin:voice:room-3";
seedClearedLaneResumeForTest(customLaneId, {
resumeConcurrency: 1,
resumeAtMs: 1_000 + MAX_TIMER_TIMEOUT_MS + 1_000,
});
expect(enableSessionSuspensionTimersForGatewayStart()).toEqual(new Set([customLaneId]));
expect(setTimeoutSpy).toHaveBeenCalledWith(expect.any(Function), MAX_TIMER_TIMEOUT_MS);
});
it("does not throttle lanes when cleanup wins a pending suspension write race", async () => {
vi.useFakeTimers();
vi.setSystemTime(1_000);
const { clearSessionSuspensionTimers } = await import("./session-suspension.js");
const previousQuotaSuspension = {
schemaVersion: 1,
suspendedAt: 500,
reason: "circuit_open",
failedProvider: "openai",
failedModel: "gpt-5.5",
laneId: CommandLane.Main,
expectedResumeBy: 2_000,
state: "suspended",
};
let resolvePatch: (() => void) | undefined;
let writtenQuotaSuspension:
| {
suspendedAt: number;
reason: string;
failedProvider: string;
failedModel: string;
laneId?: string;
}
| undefined;
sessionAccessorMocks.patchSessionEntryCore.mockImplementationOnce(async (_scope, update) => {
await new Promise<void>((resolve) => {
resolvePatch = resolve;
});
const patch = update({ quotaSuspension: previousQuotaSuspension }) as {
quotaSuspension?: typeof writtenQuotaSuspension;
};
writtenQuotaSuspension = patch.quotaSuspension;
return patch;
});
const suspension = suspendLane(100, {} as OpenClawConfig, CommandLane.Main);
await vi.waitFor(() => {
expect(resolvePatch).toBeTypeOf("function");
});
expect(clearSessionSuspensionTimers()).toBe(0);
resolvePatch?.();
await suspension;
expect(commandQueueMocks.setCommandLaneConcurrency).not.toHaveBeenCalled();
expect(writtenQuotaSuspension).toBeUndefined();
expect(sessionAccessorMocks.patchSessionEntryCore).toHaveBeenCalledOnce();
await vi.advanceTimersByTimeAsync(100);
expect(commandQueueMocks.setCommandLaneConcurrency).not.toHaveBeenCalled();
});
it("does not let a pending suspension regain ownership after test state resets", async () => {
let resolvePatch: (() => void) | undefined;
let writtenQuotaSuspension: unknown;
sessionAccessorMocks.patchSessionEntryCore.mockImplementationOnce(async (_scope, update) => {
await new Promise<void>((resolve) => {
resolvePatch = resolve;
});
const patch = update({});
writtenQuotaSuspension = patch?.quotaSuspension;
return patch;
});
const suspension = suspendLane(100, {} as OpenClawConfig, CommandLane.Main);
await vi.waitFor(() => {
expect(resolvePatch).toBeTypeOf("function");
});
const { resetSessionSuspensionStateForTest } =
await import("./session-suspension.test-support.js");
resetSessionSuspensionStateForTest();
resolvePatch?.();
await suspension;
expect(writtenQuotaSuspension).toBeUndefined();
expect(commandQueueMocks.setCommandLaneConcurrency).not.toHaveBeenCalled();
});
it("serializes suspension writes so cleanup cannot leave an intermediate write", async () => {
vi.useFakeTimers();
vi.setSystemTime(1_000);
const { clearSessionSuspensionTimers } = await import("./session-suspension.js");
let storeEntry: {
quotaSuspension?: {
suspendedAt: number;
reason: string;
failedProvider: string;
failedModel: string;
laneId?: string;
};
} = {};
let initialWrites = 0;
let releaseInitialWrites!: () => void;
const initialWritesReleased = new Promise<void>((resolve) => {
releaseInitialWrites = resolve;
});
let releaseWrite: (() => void) | undefined;
let storeEntry: { quotaSuspension?: { suspendedAt: number } } = {};
let writeCount = 0;
sessionAccessorMocks.patchSessionEntryCore.mockImplementation(async (_scope, update) => {
writeCount += 1;
if (writeCount === 1) {
await new Promise<void>((resolve) => {
releaseWrite = resolve;
});
}
const patch = update(storeEntry) as typeof storeEntry | null;
if (patch && "quotaSuspension" in patch) {
storeEntry =
patch.quotaSuspension === undefined ? {} : { quotaSuspension: patch.quotaSuspension };
}
if (initialWrites < 2) {
initialWrites += 1;
await initialWritesReleased;
storeEntry = patch.quotaSuspension ? { quotaSuspension: patch.quotaSuspension } : {};
}
return storeEntry;
});
const first = suspendLane(100, {} as OpenClawConfig, CommandLane.Main);
const second = suspendLane(100, {} as OpenClawConfig, CommandLane.Main);
await vi.waitFor(() => {
expect(initialWrites).toBe(1);
});
expect(clearSessionSuspensionTimers()).toBe(0);
releaseInitialWrites();
await Promise.all([first, second]);
const suspension = recordSuspension();
await vi.waitFor(() => expect(releaseWrite).toBeTypeOf("function"));
fenceSessionSuspensionWritesForGatewayShutdown();
releaseWrite?.();
await suspension;
expect(storeEntry.quotaSuspension).toBeUndefined();
expect(commandQueueMocks.setCommandLaneConcurrency).not.toHaveBeenCalled();
expect(sessionAccessorMocks.patchSessionEntryCore).toHaveBeenCalledTimes(2);
});
it("still throttles the lane when persistence fails while gateway is active", async () => {
vi.useFakeTimers();
sessionAccessorMocks.patchSessionEntryCore.mockRejectedValueOnce(new Error("disk busy"));
await suspendLane(
100,
{ agents: { defaults: { maxConcurrent: 4 } } } as OpenClawConfig,
CommandLane.Main,
it("blocks new state writes until gateway startup re-enables them", async () => {
const {
enableSessionSuspensionWritesForGatewayStart,
fenceSessionSuspensionWritesForGatewayShutdown,
} = await import("./session-suspension.js");
sessionAccessorMocks.patchSessionEntryCore.mockImplementation(async (_scope, update) =>
update({}),
);
expect(commandQueueMocks.setCommandLaneConcurrency).toHaveBeenCalledWith(CommandLane.Main, 0);
await vi.advanceTimersByTimeAsync(100);
expect(commandQueueMocks.setCommandLaneConcurrency).toHaveBeenLastCalledWith(
CommandLane.Main,
4,
);
fenceSessionSuspensionWritesForGatewayShutdown();
await recordSuspension();
expect(sessionAccessorMocks.patchSessionEntryCore).not.toHaveBeenCalled();
enableSessionSuspensionWritesForGatewayStart();
await recordSuspension();
expect(sessionAccessorMocks.patchSessionEntryCore).toHaveBeenCalledOnce();
});
it("defers session suspension only for the outer fallback candidate run", async () => {
it("defers only the outer fallback candidate's marker", async () => {
const { resolveSessionSuspensionTarget, runWithDeferredSessionSuspension } =
await import("./session-suspension.js");
const onDeferred = vi.fn();
@@ -510,16 +201,17 @@ describe("session suspension", () => {
target.defer({
cfg: {},
sessionId: "session-1",
laneId: CommandLane.Main,
reason: "quota_exhausted",
failedProvider: "openai",
failedModel: "gpt-5.5",
failedModel: "gpt-5.6-sol",
});
}
expect(resolveSessionSuspensionTarget()).toEqual({ mode: "suspend" });
}, onDeferred);
expect(onDeferred).toHaveBeenCalledOnce();
expect(onDeferred).toHaveBeenCalledWith(expect.objectContaining({ laneId: CommandLane.Main }));
expect(onDeferred).toHaveBeenCalledExactlyOnceWith(
expect.objectContaining({ sessionId: "session-1", failedProvider: "openai" }),
);
expect(resolveSessionSuspensionTarget()).toEqual({ mode: "suspend" });
});
+9 -199
View File
@@ -1,17 +1,13 @@
/**
* Session suspension and lane auto-resume helpers.
* Session suspension persistence and lifecycle helpers.
*
* Records quota/manual/circuit suspensions and temporarily lowers command-lane concurrency.
* Records quota/manual/circuit suspensions for diagnostics and recovery flows.
*/
import { AsyncLocalStorage } from "node:async_hooks";
import { resolveAgentMaxConcurrent, resolveSubagentMaxConcurrent } from "../config/agent-limits.js";
import { resolveCronMaxConcurrentRuns } from "../config/cron-limits.js";
import { patchSessionEntryCore } from "../config/sessions/session-accessor.js";
import type { QuotaSuspension } from "../config/sessions/types.js";
import type { OpenClawConfig } from "../config/types.openclaw.js";
import { createSubsystemLogger } from "../logging/subsystem.js";
import { setCommandLaneConcurrency } from "../process/command-queue.js";
import { CommandLane } from "../process/lanes.js";
import { resolveGlobalSingleton } from "../shared/global-singleton.js";
import {
resolveExpiresAtMsFromDurationMs,
@@ -23,24 +19,9 @@ import type { FailoverReason } from "./failover/signal.js";
const log = createSubsystemLogger("session-suspension");
const DEFAULT_CUSTOM_LANE_RESUME_CONCURRENCY = 1;
const DEFAULT_QUOTA_SUSPENSION_RESUME_MS = 30 * 60 * 1000; // 30 min
type LaneResumeTimer = {
timer: ReturnType<typeof setTimeout>;
resumeConcurrency: number;
resumeAtMs: number;
};
type ClearedLaneResume = {
resumeConcurrency: number;
resumeAtMs: number;
};
type SessionSuspensionRuntimeState = {
laneResumeTimers: Map<string, LaneResumeTimer>;
clearedLaneResumes: Map<string, ClearedLaneResume>;
gatewayLaneResumeConcurrencies: Map<string, number>;
pendingSuspensionWrites: Map<
string,
{
@@ -65,9 +46,6 @@ function getSessionSuspensionState(): SessionSuspensionRuntimeState {
const state = resolveGlobalSingleton<SessionSuspensionRuntimeState>(
SESSION_SUSPENSION_STATE_KEY,
() => ({
laneResumeTimers: new Map<string, LaneResumeTimer>(),
clearedLaneResumes: new Map<string, ClearedLaneResume>(),
gatewayLaneResumeConcurrencies: new Map<string, number>(),
pendingSuspensionWrites: new Map<
string,
{
@@ -82,12 +60,6 @@ function getSessionSuspensionState(): SessionSuspensionRuntimeState {
cleanupActive: false,
}),
);
if (!state.clearedLaneResumes) {
state.clearedLaneResumes = new Map<string, ClearedLaneResume>();
}
if (!state.gatewayLaneResumeConcurrencies) {
state.gatewayLaneResumeConcurrencies = new Map<string, number>();
}
if (!state.pendingSuspensionWrites) {
state.pendingSuspensionWrites = new Map<
string,
@@ -119,7 +91,6 @@ export type SessionSuspensionParams = {
agentId?: string;
agentDir?: string;
sessionId: string;
laneId?: string;
reason: SessionSuspensionReason;
failedProvider: string;
failedModel: string;
@@ -127,35 +98,6 @@ export type SessionSuspensionParams = {
ttlMs?: number;
};
function resolveLaneResumeConcurrency(cfg: OpenClawConfig | undefined, laneId: string): number {
switch (laneId) {
case "main":
return resolveAgentMaxConcurrent(cfg);
case "subagent":
return resolveSubagentMaxConcurrent(cfg);
case "cron":
case "cron-nested":
case "hook-dispatch":
return resolveCronMaxConcurrentRuns();
default:
return DEFAULT_CUSTOM_LANE_RESUME_CONCURRENCY;
}
}
function isGatewayManagedLane(laneId: string): boolean {
// Lane ids are open strings (plugins mint their own); narrow once so the
// membership check compares within the enum.
const lane = laneId as CommandLane;
return (
lane === CommandLane.Main ||
lane === CommandLane.Subagent ||
lane === CommandLane.Cron ||
lane === CommandLane.CronNested ||
lane === CommandLane.HookDispatch ||
lane === CommandLane.Nested
);
}
export function resolveSessionSuspensionReason(reason: FailoverReason): SessionSuspensionReason {
if (reason === "billing") {
return "manual";
@@ -184,113 +126,16 @@ export function resolveSessionSuspensionTarget(): SessionSuspensionTarget {
return { mode: "defer", defer: (params) => scope.onDeferred?.(params) };
}
function scheduleLaneAutoResume(
laneId: string,
delayMs: number,
resumeConcurrency: number,
opts: { nowMs?: number } = {},
) {
const nowMs = opts.nowMs ?? Date.now();
const state = getSessionSuspensionState();
const existing = state.laneResumeTimers.get(laneId);
if (existing) {
clearTimeout(existing.timer);
}
const canonicalResumeConcurrency = isGatewayManagedLane(laneId)
? (state.gatewayLaneResumeConcurrencies.get(laneId) ?? resumeConcurrency)
: resumeConcurrency;
const entry = {
timer: undefined as unknown as ReturnType<typeof setTimeout>,
resumeConcurrency: canonicalResumeConcurrency,
resumeAtMs: nowMs + delayMs,
};
const timer = setTimeout(() => {
if (state.laneResumeTimers.get(laneId) !== entry) {
return;
}
state.laneResumeTimers.delete(laneId);
setCommandLaneConcurrency(laneId, entry.resumeConcurrency);
log.info("auto-resumed lane after suspension TTL", {
laneId,
delayMs,
resumeConcurrency: entry.resumeConcurrency,
});
}, delayMs);
entry.timer = timer;
if (typeof timer.unref === "function") {
timer.unref();
}
state.laneResumeTimers.set(laneId, entry);
}
export function clearSessionSuspensionTimers(): number {
export function fenceSessionSuspensionWritesForGatewayShutdown(): void {
const state = getSessionSuspensionState();
state.cleanupGeneration += 1;
state.cleanupActive = true;
let cleared = 0;
for (const [laneId, entry] of state.laneResumeTimers) {
clearTimeout(entry.timer);
state.clearedLaneResumes.set(laneId, {
resumeConcurrency: entry.resumeConcurrency,
resumeAtMs: entry.resumeAtMs,
});
cleared += 1;
}
state.laneResumeTimers.clear();
return cleared;
}
export function enableSessionSuspensionTimersForGatewayStart(): Set<string> {
export function enableSessionSuspensionWritesForGatewayStart(): void {
const state = getSessionSuspensionState();
state.cleanupGeneration += 1;
state.cleanupActive = false;
const suspendedLaneIds = new Set<string>();
const nowMs = Date.now();
for (const [laneId, cleared] of state.clearedLaneResumes) {
const remainingMs = resolveTimerTimeoutMs(cleared.resumeAtMs - nowMs, 0, 0);
if (remainingMs > 0) {
setCommandLaneConcurrency(laneId, 0);
scheduleLaneAutoResume(laneId, remainingMs, cleared.resumeConcurrency, { nowMs });
suspendedLaneIds.add(laneId);
continue;
}
if (isGatewayManagedLane(laneId)) {
continue;
}
setCommandLaneConcurrency(laneId, cleared.resumeConcurrency);
}
state.clearedLaneResumes.clear();
return suspendedLaneIds;
}
export function setGatewayLaneResumeConcurrencies(
concurrencies: Readonly<Record<string, number>>,
): void {
// Gateway publication owns the desired post-suspension widths. Record them
// even when no timer exists yet so an asynchronous suspension write that
// finishes after a config reload cannot schedule a stale resume target.
const state = getSessionSuspensionState();
for (const [laneId, rawConcurrency] of Object.entries(concurrencies)) {
if (!isGatewayManagedLane(laneId)) {
continue;
}
const resumeConcurrency = Math.max(0, Math.floor(rawConcurrency));
state.gatewayLaneResumeConcurrencies.set(laneId, resumeConcurrency);
const activeTimer = state.laneResumeTimers.get(laneId);
if (activeTimer) {
activeTimer.resumeConcurrency = resumeConcurrency;
}
const clearedResume = state.clearedLaneResumes.get(laneId);
if (clearedResume) {
clearedResume.resumeConcurrency = resumeConcurrency;
}
}
}
export function getSuspendedLaneIdsForGatewayPublication(): Set<string> {
const state = getSessionSuspensionState();
const suspended = state.cleanupActive ? state.clearedLaneResumes : state.laneResumeTimers;
return new Set(suspended.keys());
}
export async function suspendSession(params: SessionSuspensionParams) {
@@ -358,17 +203,6 @@ async function suspendSessionQueued(params: SessionSuspensionParams, queuedGener
getSessionSuspensionState().pendingSuspensionWrites.delete(pendingWriteKey);
}
};
const throttleLane = () => {
if (!params.laneId) {
return;
}
setCommandLaneConcurrency(params.laneId, 0);
scheduleLaneAutoResume(
params.laneId,
ttlMs,
resolveLaneResumeConcurrency(params.cfg, params.laneId),
);
};
// Assigned at the end of the try; the catch path returns, so every read
// below sees the real patch outcome.
let persistedSuspension: boolean;
@@ -392,7 +226,6 @@ async function suspendSessionQueued(params: SessionSuspensionParams, queuedGener
failedProvider: params.failedProvider,
failedModel: params.failedModel,
summary: params.summary,
laneId: params.laneId,
expectedResumeBy,
state: "suspended",
},
@@ -402,18 +235,11 @@ async function suspendSessionQueued(params: SessionSuspensionParams, queuedGener
);
persistedSuspension = patchedEntry !== null;
} catch (err) {
log.warn("failed to persist quota suspension; applying transient lane throttle", {
log.warn("failed to persist quota suspension", {
sessionId: params.sessionId,
laneId: params.laneId,
error: err instanceof Error ? err.message : String(err),
});
releasePendingWrite();
if (
!getSessionSuspensionState().cleanupActive &&
suspensionGeneration === getSessionSuspensionState().cleanupGeneration
) {
throttleLane();
}
return;
}
@@ -429,8 +255,7 @@ async function suspendSessionQueued(params: SessionSuspensionParams, queuedGener
entry.quotaSuspension?.suspendedAt === now &&
entry.quotaSuspension.reason === params.reason &&
entry.quotaSuspension.failedProvider === params.failedProvider &&
entry.quotaSuspension.failedModel === params.failedModel &&
entry.quotaSuspension.laneId === params.laneId
entry.quotaSuspension.failedModel === params.failedModel
? { quotaSuspension: pendingWrite.previousQuotaSuspension }
: null,
{
@@ -441,7 +266,6 @@ async function suspendSessionQueued(params: SessionSuspensionParams, queuedGener
} catch (err) {
log.warn("failed to clear quota suspension after shutdown cleanup", {
sessionId: params.sessionId,
laneId: params.laneId,
error: err instanceof Error ? err.message : String(err),
});
}
@@ -449,9 +273,6 @@ async function suspendSessionQueued(params: SessionSuspensionParams, queuedGener
return;
}
if (persistedSuspension) {
throttleLane();
}
releasePendingWrite();
}
@@ -460,29 +281,18 @@ function resetSessionSuspensionStateForTest(): void {
// Invalidate in-flight writes before clearing test state. Rewinding to a
// reused generation lets a fire-and-forget suspension regain ownership.
state.cleanupGeneration += 1;
for (const entry of state.laneResumeTimers.values()) {
clearTimeout(entry.timer);
}
state.laneResumeTimers.clear();
state.clearedLaneResumes.clear();
state.gatewayLaneResumeConcurrencies.clear();
state.pendingSuspensionWrites.clear();
state.suspensionWriteChain = Promise.resolve();
state.cleanupActive = false;
}
function seedClearedLaneResumeForTest(
laneId: string,
cleared: { resumeConcurrency: number; resumeAtMs: number },
): void {
const state = getSessionSuspensionState();
state.cleanupActive = true;
state.clearedLaneResumes.set(laneId, cleared);
function isSessionSuspensionWriteCleanupActiveForTest(): boolean {
return getSessionSuspensionState().cleanupActive;
}
if (process.env.VITEST || process.env.NODE_ENV === "test") {
(globalThis as Record<PropertyKey, unknown>)[Symbol.for("openclaw.sessionSuspensionTestApi")] = {
isSessionSuspensionWriteCleanupActiveForTest,
resetSessionSuspensionStateForTest,
seedClearedLaneResumeForTest,
};
}
@@ -186,12 +186,16 @@ export function installEmbeddedRunnerFastRunE2eMocks(
provider?: string;
agentHarnessId?: string;
agentHarnessRuntimeOverride?: string;
}) => ({
id: resolveMockHarnessId(params),
label: "Mock agent harness",
supports: vi.fn(() => ({ supported: false })),
runAttempt: vi.fn(),
});
}) => {
const id = resolveMockHarnessId(params);
return {
id,
label: "Mock agent harness",
...(id === "codex" ? { authBootstrap: "harness" as const } : {}),
supports: vi.fn(() => ({ supported: false })),
runAttempt: vi.fn(),
};
};
vi.doMock("../harness/selection.js", () => ({
agentHarnessBuildsOpenClawTools: vi.fn(
(harnessId: string) => harnessId === "codex" || harnessId === "copilot",
@@ -296,17 +300,21 @@ export function installEmbeddedRunnerFastRunE2eMocks(
: undefined;
const matchingRequestedProfileId =
requestedCredential?.provider === authProvider ? requestedProfileId : undefined;
const lockedProfileId =
const userPinnedProfileId =
params.sessionAuthProfileSource === "user" ? matchingRequestedProfileId : undefined;
const orderedProfileIds = lockedProfileId
? [lockedProfileId]
: resolveAuthProfileOrder({
cfg: params.config,
store,
provider: authProvider,
preferredProfile: matchingRequestedProfileId,
forModel: params.modelId,
});
const resolvedProfileIds = resolveAuthProfileOrder({
cfg: params.config,
store,
provider: authProvider,
preferredProfile: matchingRequestedProfileId,
forModel: params.modelId,
});
const orderedProfileIds = userPinnedProfileId
? [
userPinnedProfileId,
...resolvedProfileIds.filter((profileId) => profileId !== userPinnedProfileId),
]
: resolvedProfileIds;
const profileIds = orderedProfileIds.length > 0 ? orderedProfileIds : [undefined];
const attempts = profileIds.map((profileId, index) => {
const credential = profileId ? store.profiles[profileId] : undefined;
@@ -319,7 +327,7 @@ export function installEmbeddedRunnerFastRunE2eMocks(
? {
forwardedAuthProfileId: profileId,
forwardedAuthProfileSource:
lockedProfileId === profileId ? ("user" as const) : ("auto" as const),
userPinnedProfileId === profileId ? ("user" as const) : ("auto" as const),
forwardedAuthProfileCandidateIds: profileIds
.slice(index)
.filter((candidate): candidate is string => Boolean(candidate)),
-3
View File
@@ -330,8 +330,6 @@ const QUOTA_SUSPENSION_CLEANUP_FACTOR = 2; // entries beyond N*ttl are deleted o
type QuotaSuspensionEntryMaintenanceResult = {
/** Patch to apply to the entry, or null when no TTL transition is due. */
patch: Partial<SessionEntry> | null;
/** Present when the entry transitioned from suspended to resuming. */
resumed?: { laneId?: string };
/** True when the quota-suspension marker should be removed. */
cleared: boolean;
};
@@ -360,7 +358,6 @@ export function resolveQuotaSuspensionEntryMaintenance(params: {
if (suspension.state === "suspended" && params.now >= resumeAtMs) {
return {
patch: { quotaSuspension: { ...suspension, state: "resuming" } },
resumed: { laneId: suspension.laneId },
cleared: false,
};
}
@@ -158,7 +158,6 @@ describe("resolveQuotaSuspensionEntryMaintenance", () => {
reason: "quota_exhausted",
failedProvider: "anthropic",
failedModel: "claude-opus-4-6",
laneId: "main",
},
},
now,
@@ -175,10 +174,8 @@ describe("resolveQuotaSuspensionEntryMaintenance", () => {
reason: "quota_exhausted",
failedProvider: "anthropic",
failedModel: "claude-opus-4-6",
laneId: "main",
},
},
resumed: { laneId: "main" },
cleared: false,
});
});
@@ -196,7 +193,6 @@ describe("resolveQuotaSuspensionEntryMaintenance", () => {
reason: "circuit_open",
failedProvider: "anthropic",
failedModel: "claude-opus-4-6",
laneId: "main",
},
},
now,
+4 -1
View File
@@ -259,7 +259,10 @@ export interface QuotaSuspension {
summary?: string;
/** Opaque pointer to an external snapshot blob (path/key); not the briefing text itself. */
snapshotRef?: string;
/** Lane that was set to concurrency=0 when this suspension was issued. */
/**
* @deprecated Lane suspension was removed; nothing writes this anymore. Kept only to
* hold the shipped SDK surface stable; drop at the next surface window.
*/
laneId?: string;
expectedResumeBy?: number; // Reaper TTL (e.g. 30min)
state: LaneExecutionState; // State machine check for hot-path
+7 -7
View File
@@ -27,9 +27,9 @@ const mocks = vi.hoisted(() => ({
triggerInternalHook: vi.fn<TriggerInternalHookMock>(async (_eventValue) => undefined),
disposeAllBundleLspRuntimes: vi.fn(async () => undefined),
drainRetainedEmbeddingProviders: vi.fn(async () => undefined),
clearSessionSuspensionTimers: vi.fn(() => 0),
disposeAcpSessionManagerInstance: vi.fn(async () => undefined),
getAcpSessionManager: vi.fn(() => ({})),
fenceSessionSuspensionWritesForGatewayShutdown: vi.fn(),
closePluginStateDatabase: vi.fn(async () => undefined),
}));
const WEBSOCKET_CLOSE_GRACE_MS = 1_000;
@@ -91,7 +91,8 @@ vi.mock("./embeddings-http.js", () => ({
}));
vi.mock("../agents/session-suspension.js", () => ({
clearSessionSuspensionTimers: mocks.clearSessionSuspensionTimers,
fenceSessionSuspensionWritesForGatewayShutdown:
mocks.fenceSessionSuspensionWritesForGatewayShutdown,
}));
vi.mock("../acp/control-plane/manager.lifecycle.js", () => ({
@@ -208,11 +209,10 @@ describe("createGatewayCloseHandler", () => {
mocks.disposeAllBundleLspRuntimes.mockResolvedValue(undefined);
mocks.drainRetainedEmbeddingProviders.mockClear();
mocks.drainRetainedEmbeddingProviders.mockResolvedValue(undefined);
mocks.clearSessionSuspensionTimers.mockReset();
mocks.clearSessionSuspensionTimers.mockReturnValue(0);
mocks.disposeAcpSessionManagerInstance.mockReset();
mocks.disposeAcpSessionManagerInstance.mockResolvedValue(undefined);
mocks.getAcpSessionManager.mockClear();
mocks.fenceSessionSuspensionWritesForGatewayShutdown.mockReset();
mocks.closePluginStateDatabase.mockReset();
mocks.closePluginStateDatabase.mockResolvedValue(undefined);
});
@@ -315,7 +315,7 @@ describe("createGatewayCloseHandler", () => {
it("joins an in-flight config reload before mutable runtime teardown", async () => {
const events: string[] = [];
mocks.clearSessionSuspensionTimers.mockImplementation(() => {
mocks.fenceSessionSuspensionWritesForGatewayShutdown.mockImplementation(() => {
events.push("session-suspension-timers");
return 1;
});
@@ -529,7 +529,7 @@ describe("createGatewayCloseHandler", () => {
it("clears session suspension timers before sidecars, plugin services, and channels stop", async () => {
const events: string[] = [];
mocks.clearSessionSuspensionTimers.mockImplementation(() => {
mocks.fenceSessionSuspensionWritesForGatewayShutdown.mockImplementation(() => {
events.push("session-suspension-timers");
return 1;
});
@@ -557,7 +557,7 @@ describe("createGatewayCloseHandler", () => {
await close({ reason: "test shutdown" });
expect(mocks.clearSessionSuspensionTimers).toHaveBeenCalledOnce();
expect(mocks.fenceSessionSuspensionWritesForGatewayShutdown).toHaveBeenCalledOnce();
expect(events).toEqual([
"session-suspension-timers",
"sidecar",
+3 -4
View File
@@ -9,7 +9,7 @@ import { disposeAcpSessionManagerInstance } from "../acp/control-plane/manager.l
import { disposeAllSessionMcpRuntimes } from "../agents/agent-bundle-mcp-tools.js";
import { disposeRegisteredAgentHarnesses } from "../agents/harness/registry.js";
import { createAgentRunRestartAbortError } from "../agents/run-termination.js";
import { clearSessionSuspensionTimers } from "../agents/session-suspension.js";
import { fenceSessionSuspensionWritesForGatewayShutdown } from "../agents/session-suspension.js";
import { type ChannelId, listChannelPlugins } from "../channels/plugins/index.js";
import { createInternalHookEvent, triggerInternalHook } from "../hooks/internal-hooks.js";
import type { HeartbeatRunner } from "../infra/heartbeat-runner.js";
@@ -732,9 +732,8 @@ export function createGatewayCloseHandler(
const measureCloseStep = <T>(name: string, run: () => Promise<T> | T) =>
measureGatewayRestartTrace(`restart.close.${name}`, run, [["reason", reason]]);
try {
// Fence lane auto-resume timers before the first awaited shutdown step;
// later teardown can stall long enough for a TTL callback to mutate queues.
clearSessionSuspensionTimers();
// Fence async session-state writes before the first awaited shutdown step.
fenceSessionSuspensionWritesForGatewayShutdown();
// Debug-level: the signal handler already announced the stop/restart at
// info, and the completion line below reports duration and outcome.
shutdownLog.debug(`shutdown started: ${reason}`);
@@ -319,91 +319,6 @@ describe("cron+hook capacity group", () => {
expect(lateHookStarted).toBe(true);
});
it("clears the group on hooks-off even when the grouped lane is suspended", async () => {
// The teardown path publishes only lanes that are NOT suspended. If every
// grouped member is suspended, the lane map is empty — and a guard that
// skips publication on an empty map would skip the group teardown with it.
// The stale group survives, and its members resume still paying a
// reservation for a hook lane that no longer receives work.
publish(HOOKS_ON);
expect(getCommandLaneSnapshot(CommandLane.CronNested).group).toBe("cron-hooks");
const { seedClearedLaneResumeForTest } =
await import("../agents/session-suspension.test-support.js");
seedClearedLaneResumeForTest(CommandLane.CronNested, {
resumeConcurrency: DEFAULT_CRON_MAX_CONCURRENT_RUNS,
resumeAtMs: Date.now() + 60_000,
});
// gatewayStart consults the cleared-resume map for the suspended set.
applyGatewayLaneConcurrency(resolveGatewayLaneConcurrency(HOOKS_OFF), { gatewayStart: true });
expect(getCommandLaneSnapshot(CommandLane.CronNested).group).toBeUndefined();
});
it("reinstalls the group before suspended lanes resume after hooks are re-enabled", async () => {
vi.useFakeTimers();
vi.setSystemTime(1_000);
publish(HOOKS_ON);
const { seedClearedLaneResumeForTest } =
await import("../agents/session-suspension.test-support.js");
for (const lane of [CommandLane.CronNested, CommandLane.HookDispatch]) {
seedClearedLaneResumeForTest(lane, {
resumeConcurrency: DEFAULT_CRON_MAX_CONCURRENT_RUNS,
resumeAtMs: 1_100,
});
}
applyGatewayLaneConcurrency(resolveGatewayLaneConcurrency(HOOKS_OFF), { gatewayStart: true });
expect(getCommandLaneSnapshot(CommandLane.CronNested).group).toBeUndefined();
// Both lanes remain at zero while their timers are active. Re-enabling
// hooks must still restore membership now; the per-lane resume setters
// deliberately cannot infer or install a missing capacity group later.
publish(HOOKS_ON);
expect(getCommandLaneSnapshot(CommandLane.CronNested)).toMatchObject({
maxConcurrent: 0,
group: "cron-hooks",
});
expect(getCommandLaneSnapshot(CommandLane.HookDispatch)).toMatchObject({
maxConcurrent: 0,
group: "cron-hooks",
reservedForLane: 1,
});
const cronGates = Array.from({ length: DEFAULT_CRON_MAX_CONCURRENT_RUNS }, () => gate());
const hookGates = Array.from({ length: DEFAULT_CRON_MAX_CONCURRENT_RUNS }, () => gate());
const cronRuns = cronGates.map((g) =>
enqueueCommandInLane(CommandLane.CronNested, async () => await g.promise, {
warnAfterMs: 10_000,
}),
);
const hookRuns = hookGates.map((g) =>
enqueueCommandInLane(CommandLane.HookDispatch, async () => await g.promise, {
warnAfterMs: 10_000,
}),
);
await vi.advanceTimersByTimeAsync(100);
expect(getCommandLaneSnapshot(CommandLane.CronNested).groupActive).toBe(
DEFAULT_CRON_MAX_CONCURRENT_RUNS,
);
expect(getCommandLaneSnapshot(CommandLane.HookDispatch).groupActive).toBe(
DEFAULT_CRON_MAX_CONCURRENT_RUNS,
);
expect(
getCommandLaneSnapshot(CommandLane.CronNested).activeCount +
getCommandLaneSnapshot(CommandLane.HookDispatch).activeCount,
).toBe(DEFAULT_CRON_MAX_CONCURRENT_RUNS);
for (const g of [...cronGates, ...hookGates]) {
g.release();
}
await Promise.all([...cronRuns, ...hookRuns]);
});
it("removes the group when hooks are turned off by a config reload", async () => {
publish(HOOKS_ON);
expect(getCommandLaneSnapshot(CommandLane.CronNested).group).toBe("cron-hooks");
-61
View File
@@ -126,65 +126,4 @@ describe("applyGatewayLaneConcurrency", () => {
await nestedRun;
expect(started).toBe(true);
});
it("does not resume cleanup-held built-in lanes during live config publication", async () => {
const { seedClearedLaneResumeForTest } =
await import("../agents/session-suspension.test-support.js");
seedClearedLaneResumeForTest(CommandLane.Main, {
resumeConcurrency: 3,
resumeAtMs: Date.now() + 100,
});
setCommandLaneConcurrency(CommandLane.Main, 0);
applyConfigLaneConcurrency({ agents: { defaults: { maxConcurrent: 3 } } } as OpenClawConfig);
let started = false;
const mainRun = enqueueCommandInLane(
CommandLane.Main,
async () => {
started = true;
},
{ warnAfterMs: 10_000 },
);
await Promise.resolve();
expect(started).toBe(false);
setCommandLaneConcurrency(CommandLane.Main, 1);
await mainRun;
expect(started).toBe(true);
});
it("does not resume an unexpired shared nested lane during gateway startup", async () => {
vi.useFakeTimers();
vi.setSystemTime(1_000);
const { seedClearedLaneResumeForTest } =
await import("../agents/session-suspension.test-support.js");
seedClearedLaneResumeForTest(CommandLane.Nested, {
resumeConcurrency: 1,
resumeAtMs: 1_100,
});
setCommandLaneConcurrency(CommandLane.Nested, 0);
applyConfigLaneConcurrency({} as OpenClawConfig, { gatewayStart: true });
let started = false;
const nestedRun = enqueueCommandInLane(
CommandLane.Nested,
async () => {
started = true;
},
{ warnAfterMs: 10_000 },
);
await Promise.resolve();
expect(started).toBe(false);
await vi.advanceTimersByTimeAsync(99);
expect(started).toBe(false);
await vi.advanceTimersByTimeAsync(1);
await nestedRun;
expect(started).toBe(true);
});
});
+13 -43
View File
@@ -1,8 +1,4 @@
import {
enableSessionSuspensionTimersForGatewayStart,
getSuspendedLaneIdsForGatewayPublication,
setGatewayLaneResumeConcurrencies,
} from "../agents/session-suspension.js";
import { enableSessionSuspensionWritesForGatewayStart } from "../agents/session-suspension.js";
// Gateway command-lane concurrency applier.
// Pushes config-derived agent/cron limits into the process command queue.
import { resolveAgentMaxConcurrent, resolveSubagentMaxConcurrent } from "../config/agent-limits.js";
@@ -51,24 +47,12 @@ export function applyGatewayLaneConcurrency(
concurrency: GatewayLaneConcurrency,
opts: { gatewayStart?: boolean } = {},
): void {
setGatewayLaneResumeConcurrencies({
[CommandLane.Cron]: concurrency.cron,
[CommandLane.CronNested]: concurrency.cron,
[CommandLane.HookDispatch]: concurrency.hookDispatch,
[CommandLane.Main]: concurrency.main,
[CommandLane.Nested]: 1,
[CommandLane.Subagent]: concurrency.subagent,
});
// Lane ids are open strings (plugins mint their own); narrow once so the
// gateway-managed cases compare within the enum.
const suspendedLaneIds: ReadonlySet<string> = opts.gatewayStart
? enableSessionSuspensionTimersForGatewayStart()
: getSuspendedLaneIdsForGatewayPublication();
if (opts.gatewayStart) {
enableSessionSuspensionWritesForGatewayStart();
}
// Resolution is deliberately separate: this commit-edge applier only updates
// live queue state and cannot reject a config midway through publication.
if (!suspendedLaneIds.has(CommandLane.Cron)) {
setCommandLaneConcurrency(CommandLane.Cron, concurrency.cron);
}
setCommandLaneConcurrency(CommandLane.Cron, concurrency.cron);
// `cron-nested` (cron inner agent work) and `hook-dispatch` (external hook
// agent runs) are published as ONE transaction together with the group that
// bounds them. Applying them with the per-lane setter would drain each lane
@@ -81,18 +65,11 @@ export function applyGatewayLaneConcurrency(
// budget while cron immediately expands back to its full width. Retain the
// group without a reservation until a later publication sees no active hook.
const retainInFlightHookBudget = !hooksEnabled && hookSnapshot.activeCount > 0;
const grouped: Record<string, number> = {};
if (!suspendedLaneIds.has(CommandLane.CronNested)) {
grouped[CommandLane.CronNested] = concurrency.cron;
}
if (!suspendedLaneIds.has(CommandLane.HookDispatch)) {
grouped[CommandLane.HookDispatch] = concurrency.hookDispatch;
}
// Publish even when `grouped` is empty. Both lanes can be suspended during
// config reload, but the group still needs its reservation updated or its
// membership cleared before their independent resume timers reopen them.
publishLaneConfiguration({
lanes: grouped,
lanes: {
[CommandLane.CronNested]: concurrency.cron,
[CommandLane.HookDispatch]: concurrency.hookDispatch,
},
// Opt-in. A clean hooks-off publication installs no group and
// `cron-nested` keeps the entire cron budget. During an enabled-to-disabled
// transition, a zero-reservation group may remain while in-flight hooks
@@ -115,17 +92,10 @@ export function applyGatewayLaneConcurrency(
: undefined,
clearGroups: hooksEnabled || retainInFlightHookBudget ? undefined : [CRON_HOOK_LANE_GROUP],
});
if (!suspendedLaneIds.has(CommandLane.Main)) {
setCommandLaneConcurrency(CommandLane.Main, concurrency.main);
}
setCommandLaneConcurrency(CommandLane.Main, concurrency.main);
if (opts.gatewayStart) {
// sessions.send work uses a shared nested lane with no config knob; live
// reload must not resume a currently suspended nested lane before its TTL.
if (!suspendedLaneIds.has(CommandLane.Nested)) {
setCommandLaneConcurrency(CommandLane.Nested, 1);
}
}
if (!suspendedLaneIds.has(CommandLane.Subagent)) {
setCommandLaneConcurrency(CommandLane.Subagent, concurrency.subagent);
// sessions.send work uses a shared nested lane with no config knob.
setCommandLaneConcurrency(CommandLane.Nested, 1);
}
setCommandLaneConcurrency(CommandLane.Subagent, concurrency.subagent);
}
+2 -2
View File
@@ -1,5 +1,5 @@
import { resolveActiveEmbeddedRunSessionId } from "../agents/embedded-agent-runner/run-state.js";
import { clearSessionSuspensionTimers } from "../agents/session-suspension.js";
import { fenceSessionSuspensionWritesForGatewayShutdown } from "../agents/session-suspension.js";
import { getTotalPendingReplies } from "../auto-reply/reply/dispatcher-registry.js";
import { listLoadedChannelPlugins } from "../channels/plugins/registry-loaded.js";
import type { ChannelId } from "../channels/plugins/types.public.js";
@@ -444,7 +444,7 @@ export async function prepareGatewayLifecycle(params: {
return configReloaderStopPromise;
};
const beginClosePrelude = async () => {
clearSessionSuspensionTimers();
fenceSessionSuspensionWritesForGatewayShutdown();
markClosePreludeStarted();
// Owners are fenced synchronously above. Join them before any runtime they
// can publish into is torn down.
+9 -15
View File
@@ -276,7 +276,7 @@ it("clears named plugin runtime slots between files", async () => {
}
});
it("clears session suspension state between files", async () => {
it("clears the session suspension shutdown fence between files", async () => {
const root = await fs.mkdtemp(path.join(os.tmpdir(), "openclaw-session-suspension-runner-"));
try {
const write = (name: string, content: string) =>
@@ -284,9 +284,6 @@ it("clears session suspension state between files", async () => {
const sessionSuspensionPath = JSON.stringify(
path.join(repoRoot, "src", "agents", "session-suspension.ts"),
);
const sessionSuspensionTestSupportPath = JSON.stringify(
path.join(repoRoot, "src", "agents", "session-suspension.test-support.ts"),
);
const sharedVitestConfigPath = JSON.stringify(
path.join(repoRoot, "test", "vitest", "vitest.shared.config.ts"),
);
@@ -298,16 +295,12 @@ it("clears session suspension state between files", async () => {
await write(
"a-seed.test.ts",
[
`import { getSuspendedLaneIdsForGatewayPublication } from ${sessionSuspensionPath};`,
`import { seedClearedLaneResumeForTest } from ${sessionSuspensionTestSupportPath};`,
`import { fenceSessionSuspensionWritesForGatewayShutdown } from ${sessionSuspensionPath};`,
'import { expect, it } from "vitest";',
'const laneId = "plugin:test:session-suspension";',
'it("seeds real process-global suspension state", () => {',
" seedClearedLaneResumeForTest(laneId, {",
" resumeConcurrency: 1,",
" resumeAtMs: Date.now() + 10_000,",
" });",
" expect(getSuspendedLaneIdsForGatewayPublication()).toContain(laneId);",
'const testApi = (globalThis as Record<PropertyKey, { isSessionSuspensionWriteCleanupActiveForTest(): boolean }>)[Symbol.for("openclaw.sessionSuspensionTestApi")];',
'it("seeds the real process-global shutdown fence", () => {',
" fenceSessionSuspensionWritesForGatewayShutdown();",
" expect(testApi?.isSessionSuspensionWriteCleanupActiveForTest()).toBe(true);",
"});",
"",
].join("\n"),
@@ -315,10 +308,11 @@ it("clears session suspension state between files", async () => {
await write(
"b-observe.test.ts",
[
`import { getSuspendedLaneIdsForGatewayPublication } from ${sessionSuspensionPath};`,
`import ${sessionSuspensionPath};`,
'import { expect, it } from "vitest";',
'const testApi = (globalThis as Record<PropertyKey, { isSessionSuspensionWriteCleanupActiveForTest(): boolean }>)[Symbol.for("openclaw.sessionSuspensionTestApi")];',
'it("starts without real suspension state from the previous file", () => {',
" expect(getSuspendedLaneIdsForGatewayPublication()).toEqual(new Set());",
" expect(testApi?.isSessionSuspensionWriteCleanupActiveForTest()).toBe(false);",
"});",
"",
].join("\n"),