diff --git a/docs/.generated/config-baseline.counts.json b/docs/.generated/config-baseline.counts.json index 3b4290c07751..0c97e43b133d 100644 --- a/docs/.generated/config-baseline.counts.json +++ b/docs/.generated/config-baseline.counts.json @@ -1,5 +1,5 @@ { "core": 2295, - "channel": 3716, - "plugin": 4040 + "channel": 3582, + "plugin": 3997 } diff --git a/docs/.generated/config-baseline.sha256 b/docs/.generated/config-baseline.sha256 index 7fd07ceaacb7..2c486d47721e 100644 --- a/docs/.generated/config-baseline.sha256 +++ b/docs/.generated/config-baseline.sha256 @@ -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 diff --git a/docs/.generated/plugin-sdk-api-baseline/agent-harness-runtime.json b/docs/.generated/plugin-sdk-api-baseline/agent-harness-runtime.json index b99603bd000c..1ba9f2e103c0 100644 --- a/docs/.generated/plugin-sdk-api-baseline/agent-harness-runtime.json +++ b/docs/.generated/plugin-sdk-api-baseline/agent-harness-runtime.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"} diff --git a/docs/.generated/plugin-sdk-api-baseline/agent-harness.json b/docs/.generated/plugin-sdk-api-baseline/agent-harness.json index d94cdcfcf9a0..07ccac7be166 100644 --- a/docs/.generated/plugin-sdk-api-baseline/agent-harness.json +++ b/docs/.generated/plugin-sdk-api-baseline/agent-harness.json @@ -1 +1 @@ -{"contentHash":"08051f15195453d4f35616cdb4f0ff47e974f6c2260f0ee1078e1754f5284899","entrypoint":"agent-harness","importSpecifier":"openclaw/plugin-sdk/agent-harness"} +{"contentHash":"0c4bd22dd46955da487cb750a995732a37ffa6000e94e2381621a8dc58333246","entrypoint":"agent-harness","importSpecifier":"openclaw/plugin-sdk/agent-harness"} diff --git a/docs/.generated/plugin-sdk-api-baseline/channel-core.json b/docs/.generated/plugin-sdk-api-baseline/channel-core.json index 53f568a59f2a..8ed4302b79a5 100644 --- a/docs/.generated/plugin-sdk-api-baseline/channel-core.json +++ b/docs/.generated/plugin-sdk-api-baseline/channel-core.json @@ -1 +1 @@ -{"contentHash":"0c78a76e24e4afa2f37813ea3202262605159ddc1bba101778fd68036ee0ae70","entrypoint":"channel-core","importSpecifier":"openclaw/plugin-sdk/channel-core"} +{"contentHash":"91cc5b30ad903aea2c6c002f64d21c8236fb141a2f8f8f0ca5217d46c7e79883","entrypoint":"channel-core","importSpecifier":"openclaw/plugin-sdk/channel-core"} diff --git a/docs/.generated/plugin-sdk-api-baseline/channel-entry-contract.json b/docs/.generated/plugin-sdk-api-baseline/channel-entry-contract.json index 315c3b9a2c5f..35022086ab59 100644 --- a/docs/.generated/plugin-sdk-api-baseline/channel-entry-contract.json +++ b/docs/.generated/plugin-sdk-api-baseline/channel-entry-contract.json @@ -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"} diff --git a/docs/.generated/plugin-sdk-api-baseline/channel-inbound.json b/docs/.generated/plugin-sdk-api-baseline/channel-inbound.json index 24b57d17e901..59084efcebe3 100644 --- a/docs/.generated/plugin-sdk-api-baseline/channel-inbound.json +++ b/docs/.generated/plugin-sdk-api-baseline/channel-inbound.json @@ -1 +1 @@ -{"contentHash":"e9cc066cb5ea878cdd03c53d9114b2ff650a00bb7c74fa31a3fa6be1e9f04668","entrypoint":"channel-inbound","importSpecifier":"openclaw/plugin-sdk/channel-inbound"} +{"contentHash":"d6c0e5bf1e49fc6e003221a3ff6daf951082742fa7e82215572456e4d4eaaf18","entrypoint":"channel-inbound","importSpecifier":"openclaw/plugin-sdk/channel-inbound"} diff --git a/docs/.generated/plugin-sdk-api-baseline/channel-message.json b/docs/.generated/plugin-sdk-api-baseline/channel-message.json index 36a3df3bd4af..e7c8c61ad8c0 100644 --- a/docs/.generated/plugin-sdk-api-baseline/channel-message.json +++ b/docs/.generated/plugin-sdk-api-baseline/channel-message.json @@ -1 +1 @@ -{"contentHash":"a152541a6f351027b97e0c99a02f27038ea5475e96f4513df1b950465c390242","entrypoint":"channel-message","importSpecifier":"openclaw/plugin-sdk/channel-message"} +{"contentHash":"3c2ea28cb7f23bf09c7ad0cab39751fcefd46055c563d1d2ce226b628f86dc48","entrypoint":"channel-message","importSpecifier":"openclaw/plugin-sdk/channel-message"} diff --git a/docs/.generated/plugin-sdk-api-baseline/channel-outbound.json b/docs/.generated/plugin-sdk-api-baseline/channel-outbound.json index 626938b51ccf..9b2b4d646527 100644 --- a/docs/.generated/plugin-sdk-api-baseline/channel-outbound.json +++ b/docs/.generated/plugin-sdk-api-baseline/channel-outbound.json @@ -1 +1 @@ -{"contentHash":"d75ce81f95cf3a2382057e90d70f5f8b925cd3d6236a5acce844e3512f5539b6","entrypoint":"channel-outbound","importSpecifier":"openclaw/plugin-sdk/channel-outbound"} +{"contentHash":"a31b1c8f7db35e0e316a59afaa46f27f076b55f11a1dc08dc7630224115038de","entrypoint":"channel-outbound","importSpecifier":"openclaw/plugin-sdk/channel-outbound"} diff --git a/docs/.generated/plugin-sdk-api-baseline/channel-plugin-common.json b/docs/.generated/plugin-sdk-api-baseline/channel-plugin-common.json index 6482c0527282..3cde632c6f71 100644 --- a/docs/.generated/plugin-sdk-api-baseline/channel-plugin-common.json +++ b/docs/.generated/plugin-sdk-api-baseline/channel-plugin-common.json @@ -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"} diff --git a/docs/.generated/plugin-sdk-api-baseline/core.json b/docs/.generated/plugin-sdk-api-baseline/core.json index ef56e6b14c73..c2212f81ac9f 100644 --- a/docs/.generated/plugin-sdk-api-baseline/core.json +++ b/docs/.generated/plugin-sdk-api-baseline/core.json @@ -1 +1 @@ -{"contentHash":"a6436ac71e8ba6dd99d2fb9dee159bf90675a4ed8beec9805109731630a29957","entrypoint":"core","importSpecifier":"openclaw/plugin-sdk/core"} +{"contentHash":"cb3406cb7205502c40a8bbf80521e58fd083620bbed99213d83a804b8deb8264","entrypoint":"core","importSpecifier":"openclaw/plugin-sdk/core"} diff --git a/docs/.generated/plugin-sdk-api-baseline/discord.json b/docs/.generated/plugin-sdk-api-baseline/discord.json index 6162a03b53ba..e17e1259574e 100644 --- a/docs/.generated/plugin-sdk-api-baseline/discord.json +++ b/docs/.generated/plugin-sdk-api-baseline/discord.json @@ -1 +1 @@ -{"contentHash":"eb7884340d27b0e901a2431eff2a64fae5a3d018ea6e910725ac7df2d35520b8","entrypoint":"discord","importSpecifier":"openclaw/plugin-sdk/discord"} +{"contentHash":"78ef9651b379d1b673d6c17bc9b636cdaaf2ea4cbe44c546482d53e7b7c8fdd9","entrypoint":"discord","importSpecifier":"openclaw/plugin-sdk/discord"} diff --git a/docs/.generated/plugin-sdk-api-baseline/inbound-reply-dispatch.json b/docs/.generated/plugin-sdk-api-baseline/inbound-reply-dispatch.json index ca7a0d56cfbb..14d5ae222edb 100644 --- a/docs/.generated/plugin-sdk-api-baseline/inbound-reply-dispatch.json +++ b/docs/.generated/plugin-sdk-api-baseline/inbound-reply-dispatch.json @@ -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"} diff --git a/docs/.generated/plugin-sdk-api-baseline/meeting-runtime.json b/docs/.generated/plugin-sdk-api-baseline/meeting-runtime.json index 485fd7d11b53..9b54c33ff1cb 100644 --- a/docs/.generated/plugin-sdk-api-baseline/meeting-runtime.json +++ b/docs/.generated/plugin-sdk-api-baseline/meeting-runtime.json @@ -1 +1 @@ -{"contentHash":"ebad10b0c7b1838e6cdf68153383fdc57096dd252ccbc24de5691c1f1821f3cd","entrypoint":"meeting-runtime","importSpecifier":"openclaw/plugin-sdk/meeting-runtime"} +{"contentHash":"bd05a541d2ba0e7f335ed0c4d4f187bed1c04953479682616d2b6c416bbd0459","entrypoint":"meeting-runtime","importSpecifier":"openclaw/plugin-sdk/meeting-runtime"} diff --git a/docs/.generated/plugin-sdk-api-baseline/plugin-entry.json b/docs/.generated/plugin-sdk-api-baseline/plugin-entry.json index d4439981315b..455a193b261b 100644 --- a/docs/.generated/plugin-sdk-api-baseline/plugin-entry.json +++ b/docs/.generated/plugin-sdk-api-baseline/plugin-entry.json @@ -1 +1 @@ -{"contentHash":"2d9af7dec89f421a59872de6623d51ed84d8ee23204f147d49b5218ae64ddf2c","entrypoint":"plugin-entry","importSpecifier":"openclaw/plugin-sdk/plugin-entry"} +{"contentHash":"f995b2ac4652ccf6a6107700f3a15b9d8e7bd8513466704424b899a831abdb41","entrypoint":"plugin-entry","importSpecifier":"openclaw/plugin-sdk/plugin-entry"} diff --git a/docs/.generated/plugin-sdk-api-baseline/plugin-runtime.json b/docs/.generated/plugin-sdk-api-baseline/plugin-runtime.json index 51a80932b0b8..9a62dcabf752 100644 --- a/docs/.generated/plugin-sdk-api-baseline/plugin-runtime.json +++ b/docs/.generated/plugin-sdk-api-baseline/plugin-runtime.json @@ -1 +1 @@ -{"contentHash":"793f98ff0d772f586920d8a75df073deba94321dcba74daf1caea9ba59cca05e","entrypoint":"plugin-runtime","importSpecifier":"openclaw/plugin-sdk/plugin-runtime"} +{"contentHash":"97a6f344d3efa2dd1f7d7620562f22bf632873648a82be69a26cad6a8bd0b762","entrypoint":"plugin-runtime","importSpecifier":"openclaw/plugin-sdk/plugin-runtime"} diff --git a/docs/.generated/plugin-sdk-api-baseline/provider-catalog-runtime.json b/docs/.generated/plugin-sdk-api-baseline/provider-catalog-runtime.json index 0614f88a5b5a..e9307b9a192f 100644 --- a/docs/.generated/plugin-sdk-api-baseline/provider-catalog-runtime.json +++ b/docs/.generated/plugin-sdk-api-baseline/provider-catalog-runtime.json @@ -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"} diff --git a/docs/.generated/plugin-sdk-api-baseline/tool-plugin.json b/docs/.generated/plugin-sdk-api-baseline/tool-plugin.json index cd40d8e95e15..36ef5bdff544 100644 --- a/docs/.generated/plugin-sdk-api-baseline/tool-plugin.json +++ b/docs/.generated/plugin-sdk-api-baseline/tool-plugin.json @@ -1 +1 @@ -{"contentHash":"77a1e0b93fa1a501677e5a571ed5ce44358812407554e5e494488d83bcad34d6","entrypoint":"tool-plugin","importSpecifier":"openclaw/plugin-sdk/tool-plugin"} +{"contentHash":"11c6288303d0e1b2f75ad06f5acac4fbafd94ea98eaa6212d7e0adf87db43607","entrypoint":"tool-plugin","importSpecifier":"openclaw/plugin-sdk/tool-plugin"} diff --git a/docs/.generated/plugin-sdk-api-baseline/webhook-ingress.json b/docs/.generated/plugin-sdk-api-baseline/webhook-ingress.json index 772f87ea8959..53b3e1f657d1 100644 --- a/docs/.generated/plugin-sdk-api-baseline/webhook-ingress.json +++ b/docs/.generated/plugin-sdk-api-baseline/webhook-ingress.json @@ -1 +1 @@ -{"contentHash":"03433e97b67ba2d0a95fbb7876710f2f4f4425c66d724f55bb4476797b9f2846","entrypoint":"webhook-ingress","importSpecifier":"openclaw/plugin-sdk/webhook-ingress"} +{"contentHash":"9204480482168d7f4d8f9601466c147d056403b42e3ba7979cae746792c268a8","entrypoint":"webhook-ingress","importSpecifier":"openclaw/plugin-sdk/webhook-ingress"} diff --git a/docs/auth-credential-semantics.md b/docs/auth-credential-semantics.md index 8a59a19616de..8b9d3b0891b3 100644 --- a/docs/auth-credential-semantics.md +++ b/docs/auth-credential-semantics.md @@ -69,6 +69,7 @@ Do not write `type: "aws-sdk"` into the credential store; stored credentials are - When `auth.order.` 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 diff --git a/docs/concepts/model-failover.md b/docs/concepts/model-failover.md index b1e300c942ec..8726f9078842 100644 --- a/docs/concepts/model-failover.md +++ b/docs/concepts/model-failover.md @@ -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 …@ -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 …@ -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. -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. ### 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 diff --git a/src/agents/auth-profiles/external-cli-auth-selection.test.ts b/src/agents/auth-profiles/external-cli-auth-selection.test.ts index 1eaff5c9a9c1..8caf6d71b1c2 100644 --- a/src/agents/auth-profiles/external-cli-auth-selection.test.ts +++ b/src/agents/auth-profiles/external-cli-auth-selection.test.ts @@ -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, }); }); diff --git a/src/agents/auth-profiles/external-cli-auth-selection.ts b/src/agents/auth-profiles/external-cli-auth-selection.ts index fd5554425c43..7c65b0776179 100644 --- a/src/agents/auth-profiles/external-cli-auth-selection.ts +++ b/src/agents/auth-profiles/external-cli-auth-selection.ts @@ -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, diff --git a/src/agents/btw.test.ts b/src/agents/btw.test.ts index 938b1be48861..1e7457068b63 100644 --- a/src/agents/btw.test.ts +++ b/src/agents/btw.test.ts @@ -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, diff --git a/src/agents/btw.ts b/src/agents/btw.ts index 7a8270d22338..8e55e7ead8a1 100644 --- a/src/agents/btw.ts +++ b/src/agents/btw.ts @@ -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, { diff --git a/src/agents/embedded-agent-runner.run-embedded-agent.auth-profile-rotation.e2e.test.ts b/src/agents/embedded-agent-runner.run-embedded-agent.auth-profile-rotation.e2e.test.ts index 6ca64f4c003e..6f1d7e034c72 100644 --- a/src/agents/embedded-agent-runner.run-embedded-agent.auth-profile-rotation.e2e.test.ts +++ b/src/agents/embedded-agent-runner.run-embedded-agent.auth-profile-rotation.e2e.test.ts @@ -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[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); }); }); diff --git a/src/agents/embedded-agent-runner/run-orchestrator.ts b/src/agents/embedded-agent-runner/run-orchestrator.ts index 3868f8e5167c..9df2d0e6d305 100644 --- a/src/agents/embedded-agent-runner/run-orchestrator.ts +++ b/src/agents/embedded-agent-runner/run-orchestrator.ts @@ -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) => { + const suspendForFailure = (suspensionParams: SessionSuspensionParams) => { const suspension = buildEmbeddedFailureSuspension({ suspension: suspensionParams, runAgentId: params.agentId, - laneId: globalLane, }); if (failureSuspension.mode === "defer") { failureSuspension.defer(suspension); diff --git a/src/agents/embedded-agent-runner/run.cross-provider-fallback-error-context.test-support.ts b/src/agents/embedded-agent-runner/run.cross-provider-fallback-error-context.test-support.ts index 9f2476688bc8..1d0c56184201 100644 --- a/src/agents/embedded-agent-runner/run.cross-provider-fallback-error-context.test-support.ts +++ b/src/agents/embedded-agent-runner/run.cross-provider-fallback-error-context.test-support.ts @@ -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", diff --git a/src/agents/embedded-agent-runner/run/assistant-failure.ts b/src/agents/embedded-agent-runner/run/assistant-failure.ts index 62bf87475e2e..b9ce09fff744 100644 --- a/src/agents/embedded-agent-runner/run/assistant-failure.ts +++ b/src/agents/embedded-agent-runner/run/assistant-failure.ts @@ -91,7 +91,7 @@ export async function handleEmbeddedAssistantFailure(input: { typeof handleAssistantFailover >[0]["advanceRateLimitAuthProfile"]; traceAttempts: TraceAttempt[]; - suspendForFailure: (params: Omit) => void; + suspendForFailure: (params: SessionSuspensionParams) => void; suspensionSessionId: string; agentDir: string; isProbeSession: boolean; diff --git a/src/agents/embedded-agent-runner/run/attempt-dispatch-preparation.ts b/src/agents/embedded-agent-runner/run/attempt-dispatch-preparation.ts index aaf64a63610d..6adbe33f16d4 100644 --- a/src/agents/embedded-agent-runner/run/attempt-dispatch-preparation.ts +++ b/src/agents/embedded-agent-runner/run/attempt-dispatch-preparation.ts @@ -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(), diff --git a/src/agents/embedded-agent-runner/run/attempt-recovery.ts b/src/agents/embedded-agent-runner/run/attempt-recovery.ts index ae0ddca38193..b5034dcd076e 100644 --- a/src/agents/embedded-agent-runner/run/attempt-recovery.ts +++ b/src/agents/embedded-agent-runner/run/attempt-recovery.ts @@ -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, }); diff --git a/src/agents/embedded-agent-runner/run/auth-controller.test.ts b/src/agents/embedded-agent-runner/run/auth-controller.test.ts index e7bca6227f08..3b99faf9488c 100644 --- a/src/agents/embedded-agent-runner/run/auth-controller.test.ts +++ b/src/agents/embedded-agent-runner/run/auth-controller.test.ts @@ -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 () => { diff --git a/src/agents/embedded-agent-runner/run/auth-controller.ts b/src/agents/embedded-agent-runner/run/auth-controller.ts index 10da023ef537..5009aff128d3 100644 --- a/src/agents/embedded-agent-runner/run/auth-controller.ts +++ b/src/agents/embedded-agent-runner/run/auth-controller.ts @@ -64,7 +64,7 @@ export function resolveEmbeddedAuthCooldownProbePolicy(params: { lockedProfileId?: string; modelId: string; allowTransientCooldownProbe: boolean; -}): { allowProbe: boolean; unavailableReason: FailoverReason | null } { +}): { probeProfileIds: ReadonlySet; 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(); + 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 => { - 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 }); diff --git a/src/agents/embedded-agent-runner/run/auth-plan.ts b/src/agents/embedded-agent-runner/run/auth-plan.ts index 52a6cad8c28f..ac1e20540d3c 100644 --- a/src/agents/embedded-agent-runner/run/auth-plan.ts +++ b/src/agents/embedded-agent-runner/run/auth-plan.ts @@ -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, }); } diff --git a/src/agents/embedded-agent-runner/run/auth-profile-failure-policy.test.ts b/src/agents/embedded-agent-runner/run/auth-profile-failure-policy.test.ts index aeb22ac2af24..ae3ed04ce9f7 100644 --- a/src/agents/embedded-agent-runner/run/auth-profile-failure-policy.test.ts +++ b/src/agents/embedded-agent-runner/run/auth-profile-failure-policy.test.ts @@ -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({ diff --git a/src/agents/embedded-agent-runner/run/auth-profile-failure-policy.ts b/src/agents/embedded-agent-runner/run/auth-profile-failure-policy.ts index 9c2eb97db5d1..77c75477cb77 100644 --- a/src/agents/embedded-agent-runner/run/auth-profile-failure-policy.ts +++ b/src/agents/embedded-agent-runner/run/auth-profile-failure-policy.ts @@ -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" || diff --git a/src/agents/embedded-agent-runner/run/execution-context.ts b/src/agents/embedded-agent-runner/run/execution-context.ts index 1b2ea579f274..56bf62896911 100644 --- a/src/agents/embedded-agent-runner/run/execution-context.ts +++ b/src/agents/embedded-agent-runner/run/execution-context.ts @@ -30,6 +30,6 @@ export type PreparedEmbeddedRunInput = { progressController: ReturnType; laneController: ReturnType; lifecycleGeneration: NonNullable; - suspendForFailure: (params: Omit) => void; + suspendForFailure: (params: SessionSuspensionParams) => void; preparedModelRuntime?: PreparedModelRuntimeSnapshot; }; diff --git a/src/agents/embedded-agent-runner/run/failure-suspension.test.ts b/src/agents/embedded-agent-runner/run/failure-suspension.test.ts index 723476f25413..9cbe5eeb03bc 100644 --- a/src/agents/embedded-agent-runner/run/failure-suspension.test.ts +++ b/src/agents/embedded-agent-runner/run/failure-suspension.test.ts @@ -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(); }); diff --git a/src/agents/embedded-agent-runner/run/failure-suspension.ts b/src/agents/embedded-agent-runner/run/failure-suspension.ts index 03919f836271..54c9a56517b9 100644 --- a/src/agents/embedded-agent-runner/run/failure-suspension.ts +++ b/src/agents/embedded-agent-runner/run/failure-suspension.ts @@ -7,15 +7,13 @@ import type { SessionSuspensionParams } from "../../session-suspension.js"; export function buildEmbeddedFailureSuspension(params: { - suspension: Omit; + 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, }; } diff --git a/src/agents/embedded-agent-runner/run/prompt-failure.ts b/src/agents/embedded-agent-runner/run/prompt-failure.ts index 5a53bc21a0c7..014790526fd8 100644 --- a/src/agents/embedded-agent-runner/run/prompt-failure.ts +++ b/src/agents/embedded-agent-runner/run/prompt-failure.ts @@ -53,7 +53,7 @@ export async function handleEmbeddedPromptFailure(input: { suspensionSessionId: string; runtimeAuthRetry: boolean; maybeRefreshRuntimeAuthForAuthError: (errorText: string, retry: boolean) => Promise; - suspendForFailure: (params: Omit) => void; + suspendForFailure: (params: SessionSuspensionParams) => void; resolveReplayInvalid: () => boolean; setTerminalLifecycleMeta: NonNullable; buildErrorAgentMeta: () => EmbeddedAgentMeta; diff --git a/src/agents/embedded-agent-runner/run/runtime-preparation.ts b/src/agents/embedded-agent-runner/run/runtime-preparation.ts index eea7398c95b6..6cb3063723b2 100644 --- a/src/agents/embedded-agent-runner/run/runtime-preparation.ts +++ b/src/agents/embedded-agent-runner/run/runtime-preparation.ts @@ -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 => { - 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`, ); diff --git a/src/agents/model-fallback-attempt.ts b/src/agents/model-fallback-attempt.ts index 76d91f483ece..d471679aa516 100644 --- a/src/agents/model-fallback-attempt.ts +++ b/src/agents/model-fallback-attempt.ts @@ -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", diff --git a/src/agents/model-fallback-cooldown.ts b/src/agents/model-fallback-cooldown.ts index a8f4a83a0bb4..1e72078271b0 100644 --- a/src/agents/model-fallback-cooldown.ts +++ b/src/agents/model-fallback-cooldown.ts @@ -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, }; diff --git a/src/agents/model-fallback-runner.ts b/src/agents/model-fallback-runner.ts index 33aa552dac2e..32b0b57645fa 100644 --- a/src/agents/model-fallback-runner.ts +++ b/src/agents/model-fallback-runner.ts @@ -205,8 +205,6 @@ async function runWithModelFallbackInternal( let exhaustionResult: ModelFallbackExhaustionResult | undefined; const cooldownProbeUsedProviders = new Set(); const tlsFailedProviders = new Set(); - 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( 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( (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( ? 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( 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( 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, diff --git a/src/agents/model-fallback.probe.test.ts b/src/agents/model-fallback.probe.test.ts index a46084ac6b25..bf02aad85c77 100644 --- a/src/agents/model-fallback.probe.test.ts +++ b/src/agents/model-fallback.probe.test.ts @@ -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", + ); }); }); diff --git a/src/agents/model-fallback.run-embedded.e2e.test.ts b/src/agents/model-fallback.run-embedded.e2e.test.ts index ded09442201e..c80e118d089b 100644 --- a/src/agents/model-fallback.run-embedded.e2e.test.ts +++ b/src/agents/model-fallback.run-embedded.e2e.test.ts @@ -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"); }); }); diff --git a/src/agents/model-fallback.test.ts b/src/agents/model-fallback.test.ts index 703a4a8583c5..98f2e3acac8a 100644 --- a/src/agents/model-fallback.test.ts +++ b/src/agents/model-fallback.test.ts @@ -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"; diff --git a/src/agents/runtime-plan/prepare-auth.test.ts b/src/agents/runtime-plan/prepare-auth.test.ts index d821d212a553..baf84ef7f0ec 100644 --- a/src/agents/runtime-plan/prepare-auth.test.ts +++ b/src/agents/runtime-plan/prepare-auth.test.ts @@ -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(), diff --git a/src/agents/runtime-plan/prepare-auth.ts b/src/agents/runtime-plan/prepare-auth.ts index 08c7cbe4c5bd..9aaa5a488db0 100644 --- a/src/agents/runtime-plan/prepare-auth.ts +++ b/src/agents/runtime-plan/prepare-auth.ts @@ -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, diff --git a/src/agents/session-suspension.test-support.ts b/src/agents/session-suspension.test-support.ts index 0e5c8c7f181f..0c066e729fe5 100644 --- a/src/agents/session-suspension.test-support.ts +++ b/src/agents/session-suspension.test-support.ts @@ -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); -} diff --git a/src/agents/session-suspension.test.ts b/src/agents/session-suspension.test.ts index 3485e536c0d5..ca0d8c3fddbb 100644 --- a/src/agents/session-suspension.test.ts +++ b/src/agents/session-suspension.test.ts @@ -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: /agents//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((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((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((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((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((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" }); }); diff --git a/src/agents/session-suspension.ts b/src/agents/session-suspension.ts index 09beef9dfe0f..eedebdc497ce 100644 --- a/src/agents/session-suspension.ts +++ b/src/agents/session-suspension.ts @@ -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; - resumeConcurrency: number; - resumeAtMs: number; -}; - -type ClearedLaneResume = { - resumeConcurrency: number; - resumeAtMs: number; -}; - type SessionSuspensionRuntimeState = { - laneResumeTimers: Map; - clearedLaneResumes: Map; - gatewayLaneResumeConcurrencies: Map; pendingSuspensionWrites: Map< string, { @@ -65,9 +46,6 @@ function getSessionSuspensionState(): SessionSuspensionRuntimeState { const state = resolveGlobalSingleton( SESSION_SUSPENSION_STATE_KEY, () => ({ - laneResumeTimers: new Map(), - clearedLaneResumes: new Map(), - gatewayLaneResumeConcurrencies: new Map(), pendingSuspensionWrites: new Map< string, { @@ -82,12 +60,6 @@ function getSessionSuspensionState(): SessionSuspensionRuntimeState { cleanupActive: false, }), ); - if (!state.clearedLaneResumes) { - state.clearedLaneResumes = new Map(); - } - if (!state.gatewayLaneResumeConcurrencies) { - state.gatewayLaneResumeConcurrencies = new Map(); - } 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, - 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 { +export function enableSessionSuspensionWritesForGatewayStart(): void { const state = getSessionSuspensionState(); state.cleanupGeneration += 1; state.cleanupActive = false; - const suspendedLaneIds = new Set(); - 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>, -): 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 { - 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)[Symbol.for("openclaw.sessionSuspensionTestApi")] = { + isSessionSuspensionWriteCleanupActiveForTest, resetSessionSuspensionStateForTest, - seedClearedLaneResumeForTest, }; } diff --git a/src/agents/test-helpers/embedded-agent-runner-e2e-mocks.ts b/src/agents/test-helpers/embedded-agent-runner-e2e-mocks.ts index a1ab793c1517..1ec91e37fd6a 100644 --- a/src/agents/test-helpers/embedded-agent-runner-e2e-mocks.ts +++ b/src/agents/test-helpers/embedded-agent-runner-e2e-mocks.ts @@ -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)), diff --git a/src/config/sessions/store-maintenance.ts b/src/config/sessions/store-maintenance.ts index 7f43a159d0e4..faa0688c1f36 100644 --- a/src/config/sessions/store-maintenance.ts +++ b/src/config/sessions/store-maintenance.ts @@ -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 | 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, }; } diff --git a/src/config/sessions/store.pruning.test.ts b/src/config/sessions/store.pruning.test.ts index c18c4154ced9..ec84c33101fe 100644 --- a/src/config/sessions/store.pruning.test.ts +++ b/src/config/sessions/store.pruning.test.ts @@ -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, diff --git a/src/config/sessions/types.ts b/src/config/sessions/types.ts index 186ceb78ae1e..4b9f77ce1313 100644 --- a/src/config/sessions/types.ts +++ b/src/config/sessions/types.ts @@ -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 diff --git a/src/gateway/server-close.test.ts b/src/gateway/server-close.test.ts index 88bc24a5bc62..371e8c138805 100644 --- a/src/gateway/server-close.test.ts +++ b/src/gateway/server-close.test.ts @@ -27,9 +27,9 @@ const mocks = vi.hoisted(() => ({ triggerInternalHook: vi.fn(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", diff --git a/src/gateway/server-close.ts b/src/gateway/server-close.ts index bbafc3c932ae..96796619ffe7 100644 --- a/src/gateway/server-close.ts +++ b/src/gateway/server-close.ts @@ -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 = (name: string, run: () => Promise | 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}`); diff --git a/src/gateway/server-lanes.hook-group.test.ts b/src/gateway/server-lanes.hook-group.test.ts index 3f86449f5a86..e9c279d0c00a 100644 --- a/src/gateway/server-lanes.hook-group.test.ts +++ b/src/gateway/server-lanes.hook-group.test.ts @@ -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"); diff --git a/src/gateway/server-lanes.test.ts b/src/gateway/server-lanes.test.ts index 963fb8b4e8b4..84b1e1dc9e4e 100644 --- a/src/gateway/server-lanes.test.ts +++ b/src/gateway/server-lanes.test.ts @@ -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); - }); }); diff --git a/src/gateway/server-lanes.ts b/src/gateway/server-lanes.ts index aefab012a70c..86fd38fc1a6a 100644 --- a/src/gateway/server-lanes.ts +++ b/src/gateway/server-lanes.ts @@ -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 = 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 = {}; - 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); } diff --git a/src/gateway/server-lifecycle.ts b/src/gateway/server-lifecycle.ts index 2c6884428cb9..8ff49fbac1e7 100644 --- a/src/gateway/server-lifecycle.ts +++ b/src/gateway/server-lifecycle.ts @@ -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. diff --git a/test/non-isolated-runner.test.ts b/test/non-isolated-runner.test.ts index 2a74d71ca2ca..ee411e76ddc4 100644 --- a/test/non-isolated-runner.test.ts +++ b/test/non-isolated-runner.test.ts @@ -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)[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)[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"),