From 4967cf2e56e513e231e989c54618915ec553f30f Mon Sep 17 00:00:00 2001 From: Theo Browne Date: Fri, 25 Sep 2026 14:08:51 -0700 Subject: [PATCH 01/16] perf(server): Claude instances that share a home share one capabilities probe Each Claude instance had its own private capabilities cache, so 30 instances on 9 homes ran 30 SDK probes, and they all started at once. Add a server-wide ClaudeProbeCache service. It owns one cache keyed by the full probe input (binary path, home path, cwd, and the instance env vars). The lookup gets only that input, so the probe cannot read anything the key leaves out. A 3-permit gate limits SDK probes that run at once across all instances; the probe's own timeout stays inside the gate. Successful results keep 5 min. Failed probes keep only 30 s, because a failure now marks every sibling on that home as unverified. invalidateCaches and the reset-credit re-probe drop the entry for that input only. Behavior change: sibling instances on one home now see one shared probe result, so their usage values can be up to 5 min old for a sibling. Codex, the version check, and refreshAll concurrency are unchanged. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../src/provider/Drivers/ClaudeDriver.ts | 49 +++----- .../src/provider/Drivers/ClaudeHome.test.ts | 15 +-- .../server/src/provider/Drivers/ClaudeHome.ts | 11 -- .../provider/Drivers/ClaudeProbeCache.test.ts | 117 ++++++++++++++++++ .../src/provider/Drivers/ClaudeProbeCache.ts | 75 +++++++++++ .../src/provider/Layers/ClaudeProvider.ts | 4 +- .../ProviderInstanceRegistryLive.test.ts | 3 + .../provider/Layers/ProviderRegistry.test.ts | 5 + apps/server/src/provider/ProviderDriver.ts | 5 +- apps/server/src/server.ts | 8 +- 10 files changed, 231 insertions(+), 61 deletions(-) create mode 100644 apps/server/src/provider/Drivers/ClaudeProbeCache.test.ts create mode 100644 apps/server/src/provider/Drivers/ClaudeProbeCache.ts diff --git a/apps/server/src/provider/Drivers/ClaudeDriver.ts b/apps/server/src/provider/Drivers/ClaudeDriver.ts index b0d7ec6d3ba6..95aa3e9190ad 100644 --- a/apps/server/src/provider/Drivers/ClaudeDriver.ts +++ b/apps/server/src/provider/Drivers/ClaudeDriver.ts @@ -7,14 +7,13 @@ * * Unlike Codex, the Claude snapshot probe may invoke a secondary probe * (`probeClaudeCapabilities`) to read Anthropic account + slash-command - * metadata. That probe is per-instance and keyed by binary + resolved HOME so - * two concurrent Claude instances don't cross-contaminate account metadata. + * metadata. That probe goes through the server-wide `ClaudeProbeCache`, keyed + * on the full probe input, so instances on one home share one probe and + * instances on different homes never see each other's account. * * @module provider/Drivers/ClaudeDriver */ import { ClaudeSettings, ProviderDriverKind } from "@t3tools/contracts"; -import * as Cache from "effect/Cache"; -import * as Duration from "effect/Duration"; import * as Crypto from "effect/Crypto"; import * as Effect from "effect/Effect"; import * as FileSystem from "effect/FileSystem"; @@ -33,11 +32,7 @@ import { makeClaudeAdapter } from "../Layers/ClaudeAdapter.ts"; import { makeClaudeScopedLimitNames } from "../Layers/claudeUsageLimits.ts"; import * as ClaudeResetCredits from "../Layers/claudeResetCredits.ts"; import * as ResetCreditCoordinator from "../Layers/resetCreditCoordinator.ts"; -import { - checkClaudeProviderStatus, - makePendingClaudeProvider, - probeClaudeCapabilities, -} from "../Layers/ClaudeProvider.ts"; +import { checkClaudeProviderStatus, makePendingClaudeProvider } from "../Layers/ClaudeProvider.ts"; import { ProviderEventLoggers } from "../Layers/ProviderEventLoggers.ts"; import { resolveClaudeModelCatalog } from "../ClaudeModelCatalog.ts"; import { makeManagedServerProvider } from "../makeManagedServerProvider.ts"; @@ -61,16 +56,12 @@ import { makeProviderSnapshotSettingsSource, type ProviderSnapshotSettings, } from "../providerUpdateSettings.ts"; -import { - makeClaudeCapabilitiesCacheKey, - makeClaudeContinuationGroupKey, - resolveClaudeHomePath, -} from "./ClaudeHome.ts"; +import { makeClaudeContinuationGroupKey, resolveClaudeHomePath } from "./ClaudeHome.ts"; +import * as ClaudeProbeCache from "./ClaudeProbeCache.ts"; import { discoverClaudeSkills } from "./ClaudeSkills.ts"; const decodeClaudeSettings = Schema.decodeSync(ClaudeSettings); const DRIVER_KIND = ProviderDriverKind.make("claudeAgent"); -const CAPABILITIES_PROBE_TTL = Duration.minutes(5); function isClaudeNativeCommandPath(commandPath: string): boolean { const normalized = normalizeCommandPath(commandPath); @@ -93,6 +84,7 @@ const UPDATE = makePackageManagedProviderMaintenanceResolver({ export type ClaudeDriverEnv = | BackgroundPolicy.BackgroundPolicy | ChildProcessSpawner.ChildProcessSpawner + | ClaudeProbeCache.ClaudeProbeCache | ResetCreditCoordinator.ResetCreditCoordinator | Crypto.Crypto | FileSystem.FileSystem @@ -119,6 +111,7 @@ export const ClaudeDriver: ProviderDriver = { const { cwd } = yield* ServerConfig; const httpClient = yield* HttpClient.HttpClient; const resetCreditCoordinator = yield* ResetCreditCoordinator.ResetCreditCoordinator; + const probeCache = yield* ClaudeProbeCache.ClaudeProbeCache; const serverSettings = yield* ServerSettingsService; const eventLoggers = yield* ProviderEventLoggers; const modelManifest = yield* ModelManifest.ModelManifest; @@ -178,21 +171,13 @@ export const ClaudeDriver: ProviderDriver = { modelCatalog, ); - // Per-instance capabilities cache: keyed on binary + resolved HOME so - // account-specific probes never share auth metadata across instances. - const capabilitiesProbeCache = yield* Cache.make({ - capacity: 1, - timeToLive: CAPABILITIES_PROBE_TTL, - lookup: () => - probeClaudeCapabilities(effectiveConfig, processEnv, cwd).pipe( - Effect.provideService(Path.Path, path), - ), - }); - const capabilitiesCacheKey = yield* makeClaudeCapabilitiesCacheKey( - effectiveConfig, + // Shared with every instance that has the same probe input. + const probeInput = { + binaryPath: effectiveConfig.binaryPath, + homePath: effectiveConfig.homePath, cwd, - processEnv, - ); + environment, + } satisfies ClaudeProbeCache.ClaudeProbeInput; // Start the TTL-gated refresh without delaying provider readiness. The // next check observes a remote manifest after the background fetch lands. @@ -202,7 +187,7 @@ export const ClaudeDriver: ProviderDriver = { Effect.flatMap((manifest) => checkClaudeProviderStatus( effectiveConfig, - () => Cache.get(capabilitiesProbeCache, capabilitiesCacheKey), + () => probeCache.capabilities(probeInput), processEnv, cwd, resolveClaudeModelCatalog(manifest), @@ -312,7 +297,7 @@ export const ClaudeDriver: ProviderDriver = { Effect.tap((outcome) => Effect.gen(function* () { const before = (yield* snapshot.getSnapshot).usageLimits?.checkedAt; - yield* Cache.invalidateAll(capabilitiesProbeCache); + yield* probeCache.invalidate(probeInput); const refreshed = yield* snapshot.refresh; const after = refreshed.usageLimits?.checkedAt; if ( @@ -343,7 +328,7 @@ export const ClaudeDriver: ProviderDriver = { accentColor, enabled, snapshot, - invalidateCaches: Cache.invalidateAll(capabilitiesProbeCache), + invalidateCaches: probeCache.invalidate(probeInput), snapshotForCwd, adapter, textGeneration, diff --git a/apps/server/src/provider/Drivers/ClaudeHome.test.ts b/apps/server/src/provider/Drivers/ClaudeHome.test.ts index fb1caf753f86..bdfefa3fc504 100644 --- a/apps/server/src/provider/Drivers/ClaudeHome.test.ts +++ b/apps/server/src/provider/Drivers/ClaudeHome.test.ts @@ -7,7 +7,6 @@ import * as Path from "effect/Path"; import { claudeSignedOutMessage, - makeClaudeCapabilitiesCacheKey, makeClaudeContinuationGroupKey, makeClaudeEnvironment, resolveClaudeHomePath, @@ -32,7 +31,7 @@ it.layer(NodeServices.layer)("ClaudeHome", (it) => { }), ); - it.effect("resolves configured Claude HOME and stamps continuation/cache keys with it", () => + it.effect("resolves configured Claude HOME and stamps continuation keys with it", () => Effect.gen(function* () { const path = yield* Path.Path; const homePath = "~/.claude-work"; @@ -41,9 +40,6 @@ it.layer(NodeServices.layer)("ClaudeHome", (it) => { expect(yield* resolveClaudeHomePath({ homePath })).toBe(resolved); expect((yield* makeClaudeEnvironment({ homePath })).CLAUDE_CONFIG_DIR).toBe(resolved); expect(yield* makeClaudeContinuationGroupKey({ homePath })).toBe(`claude:home:${resolved}`); - expect(yield* makeClaudeCapabilitiesCacheKey({ binaryPath: "claude", homePath })).toBe( - `claude\0${resolved}\0`, - ); }), ); @@ -75,14 +71,5 @@ it.layer(NodeServices.layer)("ClaudeHome", (it) => { expect(message).not.toContain("CLAUDE_CONFIG_DIR="); expect(message).toContain("then start a new thread"); }); - - it.effect("separates capability probes by cwd", () => - Effect.gen(function* () { - const config = { binaryPath: "claude", homePath: "" }; - const first = yield* makeClaudeCapabilitiesCacheKey(config, "/repo-a"); - const second = yield* makeClaudeCapabilitiesCacheKey(config, "/repo-b"); - expect(first).not.toBe(second); - }), - ); }); }); diff --git a/apps/server/src/provider/Drivers/ClaudeHome.ts b/apps/server/src/provider/Drivers/ClaudeHome.ts index 70699ca669e8..02fb86b2f0ba 100644 --- a/apps/server/src/provider/Drivers/ClaudeHome.ts +++ b/apps/server/src/provider/Drivers/ClaudeHome.ts @@ -63,17 +63,6 @@ export const makeClaudeContinuationGroupKey = Effect.fn("makeClaudeContinuationG }, ); -export const makeClaudeCapabilitiesCacheKey = Effect.fn("makeClaudeCapabilitiesCacheKey")( - function* ( - config: Pick, - cwd?: string, - environment?: NodeJS.ProcessEnv, - ): Effect.fn.Return { - const resolvedHomePath = yield* resolveClaudeHomePath(config, environment); - return `${config.binaryPath}\0${resolvedHomePath}\0${cwd ?? ""}`; - }, -); - /** * Describe the spawned CLI's environment separately from the login command so * paths remain literal on every shell, including relative inherited values. diff --git a/apps/server/src/provider/Drivers/ClaudeProbeCache.test.ts b/apps/server/src/provider/Drivers/ClaudeProbeCache.test.ts new file mode 100644 index 000000000000..7f0753322bc3 --- /dev/null +++ b/apps/server/src/provider/Drivers/ClaudeProbeCache.test.ts @@ -0,0 +1,117 @@ +import * as ClaudeSdk from "@anthropic-ai/claude-agent-sdk"; +import * as NodeServices from "@effect/platform-node/NodeServices"; +import { assert, it } from "@effect/vitest"; +import * as Effect from "effect/Effect"; +import * as Layer from "effect/Layer"; +import * as TestClock from "effect/testing/TestClock"; +import { vi } from "vite-plus/test"; + +import * as ClaudeProbeCache from "./ClaudeProbeCache.ts"; + +vi.mock("@anthropic-ai/claude-agent-sdk", { spy: true }); + +const testLayer = ClaudeProbeCache.layer.pipe(Layer.provide(NodeServices.layer)); + +// A fresh object per call, so sharing depends on equal inputs, not identity. +const input = ( + homePath: string, + environment: ClaudeProbeCache.ClaudeProbeInput["environment"] = [], +): ClaudeProbeCache.ClaudeProbeInput => ({ + binaryPath: "claude", + homePath, + cwd: "/repo", + environment, +}); + +const failedQuery = () => + ({ + initializationResult: () => Promise.reject(new Error("not logged in")), + }) as ReturnType; + +// Stands in for the SDK. Each probe reports the Claude home it was started +// with as the account email. +const mockSdk = Effect.gen(function* () { + const query = vi.spyOn(ClaudeSdk, "query").mockImplementation( + ({ options }) => + ({ + initializationResult: async () => ({ + account: { email: options?.env?.CLAUDE_CONFIG_DIR ?? "" }, + commands: [{ name: "review", description: "Review changes", argumentHint: "" }], + }), + usage_EXPERIMENTAL_MAY_CHANGE_DO_NOT_RELY_ON_THIS_API_YET: () => + Promise.reject(new Error("usage unavailable")), + }) as ReturnType, + ); + yield* Effect.addFinalizer(() => Effect.sync(() => query.mockRestore())); + return query; +}); + +it.effect("instances with the same probe input share one probe", () => + Effect.gen(function* () { + const query = yield* mockSdk; + const cache = yield* ClaudeProbeCache.ClaudeProbeCache; + + const [first, second] = yield* Effect.all( + [cache.capabilities(input("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/homes/work")), cache.capabilities(input("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/homes/work"))], + { concurrency: "unbounded" }, + ); + const later = yield* cache.capabilities(input("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/homes/work")); + + assert.equal(query.mock.calls.length, 1); + assert.match(first?.email ?? "", /work$/); + assert.deepEqual(second, first); + assert.deepEqual(later, first); + }).pipe(Effect.scoped, Effect.provide(testLayer)), +); + +it.effect("different homes or instance env vars run separate probes", () => + Effect.gen(function* () { + const query = yield* mockSdk; + const cache = yield* ClaudeProbeCache.ClaudeProbeCache; + + const work = yield* cache.capabilities(input("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/homes/work")); + const personal = yield* cache.capabilities(input("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/homes/personal")); + yield* cache.capabilities( + input("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/homes/work", [{ name: "ANTHROPIC_API_KEY", value: "sk-test", sensitive: true }]), + ); + + assert.equal(query.mock.calls.length, 3); + assert.match(work?.email ?? "", /work$/); + assert.match(personal?.email ?? "", /personal$/); + }).pipe(Effect.scoped, Effect.provide(testLayer)), +); + +it.effect("retries a failed probe after a short TTL and keeps a success longer", () => + Effect.gen(function* () { + const query = yield* mockSdk; + query.mockImplementationOnce(failedQuery); + const cache = yield* ClaudeProbeCache.ClaudeProbeCache; + + assert.equal(yield* cache.capabilities(input("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/homes/work")), undefined); + assert.equal(yield* cache.capabilities(input("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/homes/work")), undefined); + assert.equal(query.mock.calls.length, 1); + + yield* TestClock.adjust("30 seconds"); + assert.match((yield* cache.capabilities(input("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/homes/work")))?.email ?? "", /work$/); + assert.equal(query.mock.calls.length, 2); + + yield* TestClock.adjust("1 minute"); + yield* cache.capabilities(input("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/homes/work")); + assert.equal(query.mock.calls.length, 2); + }).pipe(Effect.scoped, Effect.provide(testLayer)), +); + +it.effect("invalidate re-probes only that input", () => + Effect.gen(function* () { + const query = yield* mockSdk; + const cache = yield* ClaudeProbeCache.ClaudeProbeCache; + + yield* cache.capabilities(input("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/homes/work")); + yield* cache.capabilities(input("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/homes/personal")); + yield* cache.invalidate(input("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/homes/work")); + yield* cache.capabilities(input("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/homes/work")); + yield* cache.capabilities(input("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/homes/personal")); + + assert.equal(query.mock.calls.length, 3); + }).pipe(Effect.scoped, Effect.provide(testLayer)), +); diff --git a/apps/server/src/provider/Drivers/ClaudeProbeCache.ts b/apps/server/src/provider/Drivers/ClaudeProbeCache.ts new file mode 100644 index 000000000000..a5603485dbd3 --- /dev/null +++ b/apps/server/src/provider/Drivers/ClaudeProbeCache.ts @@ -0,0 +1,75 @@ +/** + * One server-wide cache for the Claude capabilities probe. Claude instances + * whose probe inputs match (same binary, home, cwd and instance env vars) + * read the same account, so they share one cached result instead of each + * starting its own SDK session. A small gate limits how many SDK probes run + * at once across all instances. + * + * @module provider/Drivers/ClaudeProbeCache + */ +import type { ClaudeSettings, ProviderInstanceEnvironment } from "@t3tools/contracts"; +import * as Cache from "effect/Cache"; +import * as Context from "effect/Context"; +import * as Duration from "effect/Duration"; +import * as Effect from "effect/Effect"; +import * as Exit from "effect/Exit"; +import * as Layer from "effect/Layer"; +import * as Semaphore from "effect/Semaphore"; + +import { type ClaudeCapabilitiesProbe, probeClaudeCapabilities } from "../Layers/ClaudeProvider.ts"; +import { mergeProviderInstanceEnvironment } from "../ProviderInstanceEnvironment.ts"; + +/** + * Everything the probe reads, and also the cache key. The lookup gets only + * this value, so the probe cannot depend on an input the key leaves out. + * Keys compare structurally. + */ +export type ClaudeProbeInput = Pick & { + readonly cwd: string; + readonly environment: ProviderInstanceEnvironment; +}; + +const PROBE_TTL = Duration.minutes(5); +// A failed probe marks every instance on that home as unverified, so retry soon. +const FAILED_PROBE_TTL = Duration.seconds(30); +const MAX_CONCURRENT_PROBES = 3; + +export class ClaudeProbeCache extends Context.Service< + ClaudeProbeCache, + { + /** Cached probe result for `input`, or `undefined` when the probe failed. */ + readonly capabilities: ( + input: ClaudeProbeInput, + ) => Effect.Effect; + /** Drop the result for `input`, so the next read probes again. */ + readonly invalidate: (input: ClaudeProbeInput) => Effect.Effect; + } +>()("t3/provider/Drivers/ClaudeProbeCache") {} + +export const layer = Layer.effect( + ClaudeProbeCache, + Effect.gen(function* () { + const gate = yield* Semaphore.make(MAX_CONCURRENT_PROBES); + // The probe keeps its own timeout inside the gate, so time spent waiting + // for a permit cannot turn into a false failure. + const cache = yield* Cache.makeWith( + (input: ClaudeProbeInput) => + gate.withPermits(1)( + probeClaudeCapabilities( + input, + mergeProviderInstanceEnvironment(input.environment), + input.cwd, + ), + ), + { + capacity: 64, + timeToLive: (exit) => + Exit.isSuccess(exit) && exit.value !== undefined ? PROBE_TTL : FAILED_PROBE_TTL, + }, + ); + return { + capabilities: (input) => Cache.get(cache, input), + invalidate: (input) => Cache.invalidate(cache, input), + } satisfies ClaudeProbeCache["Service"]; + }), +); diff --git a/apps/server/src/provider/Layers/ClaudeProvider.ts b/apps/server/src/provider/Layers/ClaudeProvider.ts index db06557c7c8e..b296a9fd1cda 100644 --- a/apps/server/src/provider/Layers/ClaudeProvider.ts +++ b/apps/server/src/provider/Layers/ClaudeProvider.ts @@ -225,7 +225,7 @@ function nonEmptyProbeString(value: string): string | undefined { return candidate ? candidate : undefined; } -type ClaudeCapabilitiesProbe = { +export type ClaudeCapabilitiesProbe = { readonly email: string | undefined; readonly subscriptionType: string | undefined; readonly tokenSource: string | undefined; @@ -330,7 +330,7 @@ function waitForAbortSignal(signal: AbortSignal): Promise { * subscription type information. */ const probeClaudeCapabilities = ( - claudeSettings: ClaudeSettings, + claudeSettings: Pick, environment?: NodeJS.ProcessEnv, cwd?: string, ) => { diff --git a/apps/server/src/provider/Layers/ProviderInstanceRegistryLive.test.ts b/apps/server/src/provider/Layers/ProviderInstanceRegistryLive.test.ts index c33b1dfc690a..3cd3835ce570 100644 --- a/apps/server/src/provider/Layers/ProviderInstanceRegistryLive.test.ts +++ b/apps/server/src/provider/Layers/ProviderInstanceRegistryLive.test.ts @@ -57,6 +57,7 @@ import { OpenCodeDriver } from "../Drivers/OpenCodeDriver.ts"; import * as ModelManifest from "../ModelManifest.ts"; import { OpenCodeRuntimeLive } from "../opencodeRuntime.ts"; import * as ResetCreditCoordinator from "./resetCreditCoordinator.ts"; +import * as ClaudeProbeCache from "../Drivers/ClaudeProbeCache.ts"; import { NoOpProviderEventLoggers, ProviderEventLoggers } from "./ProviderEventLoggers.ts"; import { makeProviderInstanceRegistry } from "./ProviderInstanceRegistryLive.ts"; @@ -245,6 +246,7 @@ describe("ProviderInstanceRegistryLive — multi-instance codex slice", () => { const testLayer = ServerConfig.layerTest(process.cwd(), { prefix: "provider-instance-registry-test", }).pipe( + Layer.provideMerge(ClaudeProbeCache.layer), Layer.provideMerge(NodeServices.layer), Layer.provideMerge(BackgroundPolicyAlwaysRunLayer), Layer.provideMerge(ServerSettingsService.layerTest()), @@ -601,6 +603,7 @@ describe("ProviderInstanceRegistryLive — all drivers slice", () => { prefix: "provider-instance-registry-all-drivers-test", }), ), + Layer.provideMerge(ClaudeProbeCache.layer), Layer.provideMerge(infraLayer), Layer.provideMerge(BackgroundPolicyAlwaysRunLayer), Layer.provideMerge(ServerSettingsService.layerTest()), diff --git a/apps/server/src/provider/Layers/ProviderRegistry.test.ts b/apps/server/src/provider/Layers/ProviderRegistry.test.ts index ee1a23573277..abc2cd42fdb7 100644 --- a/apps/server/src/provider/Layers/ProviderRegistry.test.ts +++ b/apps/server/src/provider/Layers/ProviderRegistry.test.ts @@ -40,6 +40,7 @@ import { AntigravityInstallation } from "../AntigravityInstallation.ts"; import * as ModelManifest from "../ModelManifest.ts"; import { applyProviderCompatibility } from "../providerCompatibility.ts"; import * as ResetCreditCoordinator from "./resetCreditCoordinator.ts"; +import * as ClaudeProbeCache from "../Drivers/ClaudeProbeCache.ts"; import * as OpenCodeRuntime from "../opencodeRuntime.ts"; import * as ProviderEventLoggers from "./ProviderEventLoggers.ts"; import { ProviderInstanceRegistryHydrationLive } from "./ProviderInstanceRegistryHydration.ts"; @@ -2322,6 +2323,7 @@ it.layer(Layer.mergeAll(NodeServices.layer, ServerSettingsModule.layerTest(), Te ), Layer.provideMerge(ModelManifest.layerTest), Layer.provideMerge(ResetCreditCoordinator.layerTest), + Layer.provideMerge(ClaudeProbeCache.layer), Layer.provideMerge(OpenCodeRuntime.OpenCodeRuntimeLive), Layer.provideMerge(BackgroundPolicyAlwaysRunLayer), // NO spawner mock — `ChildProcessSpawner` is supplied by the @@ -2421,6 +2423,7 @@ it.layer(Layer.mergeAll(NodeServices.layer, ServerSettingsModule.layerTest(), Te ), Layer.provideMerge(ModelManifest.layerTest), Layer.provideMerge(ResetCreditCoordinator.layerTest), + Layer.provideMerge(ClaudeProbeCache.layer), Layer.provideMerge(OpenCodeRuntime.OpenCodeRuntimeLive), Layer.updateService(ChildProcessSpawner.ChildProcessSpawner, (spawner) => ChildProcessSpawner.make((command) => { @@ -2537,6 +2540,7 @@ it.layer(Layer.mergeAll(NodeServices.layer, ServerSettingsModule.layerTest(), Te ), Layer.provideMerge(ModelManifest.layerTest), Layer.provideMerge(ResetCreditCoordinator.layerTest), + Layer.provideMerge(ClaudeProbeCache.layer), Layer.provideMerge(OpenCodeRuntime.OpenCodeRuntimeLive), Layer.provideMerge(NodeServices.layer), Layer.provideMerge(BackgroundPolicyAlwaysRunLayer), @@ -2599,6 +2603,7 @@ it.layer(Layer.mergeAll(NodeServices.layer, ServerSettingsModule.layerTest(), Te ), Layer.provideMerge(ModelManifest.layerTest), Layer.provideMerge(ResetCreditCoordinator.layerTest), + Layer.provideMerge(ClaudeProbeCache.layer), Layer.provideMerge(OpenCodeRuntime.OpenCodeRuntimeLive), Layer.provideMerge(BackgroundPolicyAlwaysRunLayer), Layer.provideMerge( diff --git a/apps/server/src/provider/ProviderDriver.ts b/apps/server/src/provider/ProviderDriver.ts index 059a2508e40f..6f77a047d509 100644 --- a/apps/server/src/provider/ProviderDriver.ts +++ b/apps/server/src/provider/ProviderDriver.ts @@ -131,7 +131,10 @@ export interface ProviderDriverCreateInput { * `create` is responsible for *all* per-instance state — process handles, * pubsub topics, refs, file watchers — and must release them when its * scope closes. Two calls to `create` with different `instanceId` / - * `config` MUST yield instances with no shared mutable state. + * `config` MUST yield instances with no shared mutable state. State that + * belongs to an account rather than an instance lives in a server-wide + * service from `R` instead (see `ResetCreditCoordinator`, `ClaudeProbeCache`), + * keyed so instances share it only when their inputs match. */ export interface ProviderDriver { readonly driverKind: ProviderDriverKind; diff --git a/apps/server/src/server.ts b/apps/server/src/server.ts index 88c08fd3bdee..c72a91a0c8c9 100644 --- a/apps/server/src/server.ts +++ b/apps/server/src/server.ts @@ -52,6 +52,7 @@ import * as ProviderSessionRuntime from "./persistence/ProviderSessionRuntime.ts import { ProviderAdapterRegistryLive } from "./provider/Layers/ProviderAdapterRegistry.ts"; import * as ModelManifest from "./provider/ModelManifest.ts"; import * as ResetCreditCoordinator from "./provider/Layers/resetCreditCoordinator.ts"; +import * as ClaudeProbeCache from "./provider/Drivers/ClaudeProbeCache.ts"; import * as ProviderEventLoggers from "./provider/Layers/ProviderEventLoggers.ts"; import { ProviderServiceLive } from "./provider/Layers/ProviderService.ts"; import { ProviderAuthServiceLive } from "./provider/Layers/ProviderAuthService.ts"; @@ -542,7 +543,12 @@ const RuntimeCoreDependenciesLive = ReactorLayerLive.pipe( // from the repo's `model-manifest.json` on `main` and applied by the // Codex/Claude drivers. Layer.provideMerge( - Layer.mergeAll(ProviderEventLoggers.layer, ModelManifest.layer, ResetCreditCoordinator.layer), + Layer.mergeAll( + ProviderEventLoggers.layer, + ModelManifest.layer, + ResetCreditCoordinator.layer, + ClaudeProbeCache.layer, + ), ), // `OpenCodeDriver.create()` yields `OpenCodeRuntime`; previously the old // `ProviderRegistryLive` pulled `OpenCodeRuntimeLive` in for itself, but From de08e34f5a1c4a60dc9e5acdba774185dac7d810 Mon Sep 17 00:00:00 2001 From: Theo Browne Date: Fri, 25 Sep 2026 14:15:56 -0700 Subject: [PATCH 02/16] test(server): cover the Claude probe gate and fix the driver header Add a test that the shared probe cache runs at most 3 SDK probes at once. It fails when the gate is widened. Reword the ClaudeDriver header: instances share one probe only when the whole probe input matches, not just the home. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../src/provider/Drivers/ClaudeDriver.ts | 4 +- .../provider/Drivers/ClaudeProbeCache.test.ts | 65 +++++++++++++------ 2 files changed, 47 insertions(+), 22 deletions(-) diff --git a/apps/server/src/provider/Drivers/ClaudeDriver.ts b/apps/server/src/provider/Drivers/ClaudeDriver.ts index 95aa3e9190ad..871f3d5bac1e 100644 --- a/apps/server/src/provider/Drivers/ClaudeDriver.ts +++ b/apps/server/src/provider/Drivers/ClaudeDriver.ts @@ -8,8 +8,8 @@ * Unlike Codex, the Claude snapshot probe may invoke a secondary probe * (`probeClaudeCapabilities`) to read Anthropic account + slash-command * metadata. That probe goes through the server-wide `ClaudeProbeCache`, keyed - * on the full probe input, so instances on one home share one probe and - * instances on different homes never see each other's account. + * on the full probe input, so instances with the same probe input share one + * probe and instances on different homes never see each other's account. * * @module provider/Drivers/ClaudeDriver */ diff --git a/apps/server/src/provider/Drivers/ClaudeProbeCache.test.ts b/apps/server/src/provider/Drivers/ClaudeProbeCache.test.ts index 7f0753322bc3..fa2a60e5d85a 100644 --- a/apps/server/src/provider/Drivers/ClaudeProbeCache.test.ts +++ b/apps/server/src/provider/Drivers/ClaudeProbeCache.test.ts @@ -2,6 +2,7 @@ import * as ClaudeSdk from "@anthropic-ai/claude-agent-sdk"; import * as NodeServices from "@effect/platform-node/NodeServices"; import { assert, it } from "@effect/vitest"; import * as Effect from "effect/Effect"; +import * as Fiber from "effect/Fiber"; import * as Layer from "effect/Layer"; import * as TestClock from "effect/testing/TestClock"; import { vi } from "vite-plus/test"; @@ -29,26 +30,30 @@ const failedQuery = () => }) as ReturnType; // Stands in for the SDK. Each probe reports the Claude home it was started -// with as the account email. -const mockSdk = Effect.gen(function* () { - const query = vi.spyOn(ClaudeSdk, "query").mockImplementation( - ({ options }) => - ({ - initializationResult: async () => ({ - account: { email: options?.env?.CLAUDE_CONFIG_DIR ?? "" }, - commands: [{ name: "review", description: "Review changes", argumentHint: "" }], - }), - usage_EXPERIMENTAL_MAY_CHANGE_DO_NOT_RELY_ON_THIS_API_YET: () => - Promise.reject(new Error("usage unavailable")), - }) as ReturnType, - ); - yield* Effect.addFinalizer(() => Effect.sync(() => query.mockRestore())); - return query; -}); +// with as the account email. Probes finish once `ready` resolves. +const mockSdk = (ready: Promise = Promise.resolve()) => + Effect.gen(function* () { + const query = vi.spyOn(ClaudeSdk, "query").mockImplementation( + ({ options }) => + ({ + initializationResult: async () => { + await ready; + return { + account: { email: options?.env?.CLAUDE_CONFIG_DIR ?? "" }, + commands: [{ name: "review", description: "Review changes", argumentHint: "" }], + }; + }, + usage_EXPERIMENTAL_MAY_CHANGE_DO_NOT_RELY_ON_THIS_API_YET: () => + Promise.reject(new Error("usage unavailable")), + }) as ReturnType, + ); + yield* Effect.addFinalizer(() => Effect.sync(() => query.mockRestore())); + return query; + }); it.effect("instances with the same probe input share one probe", () => Effect.gen(function* () { - const query = yield* mockSdk; + const query = yield* mockSdk(); const cache = yield* ClaudeProbeCache.ClaudeProbeCache; const [first, second] = yield* Effect.all( @@ -66,7 +71,7 @@ it.effect("instances with the same probe input share one probe", () => it.effect("different homes or instance env vars run separate probes", () => Effect.gen(function* () { - const query = yield* mockSdk; + const query = yield* mockSdk(); const cache = yield* ClaudeProbeCache.ClaudeProbeCache; const work = yield* cache.capabilities(input("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/homes/work")); @@ -83,7 +88,7 @@ it.effect("different homes or instance env vars run separate probes", () => it.effect("retries a failed probe after a short TTL and keeps a success longer", () => Effect.gen(function* () { - const query = yield* mockSdk; + const query = yield* mockSdk(); query.mockImplementationOnce(failedQuery); const cache = yield* ClaudeProbeCache.ClaudeProbeCache; @@ -103,7 +108,7 @@ it.effect("retries a failed probe after a short TTL and keeps a success longer", it.effect("invalidate re-probes only that input", () => Effect.gen(function* () { - const query = yield* mockSdk; + const query = yield* mockSdk(); const cache = yield* ClaudeProbeCache.ClaudeProbeCache; yield* cache.capabilities(input("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/homes/work")); @@ -115,3 +120,23 @@ it.effect("invalidate re-probes only that input", () => assert.equal(query.mock.calls.length, 3); }).pipe(Effect.scoped, Effect.provide(testLayer)), ); + +it.effect("runs at most 3 SDK probes at once", () => + Effect.gen(function* () { + const { promise: released, resolve: release } = Promise.withResolvers(); + const query = yield* mockSdk(released); + const cache = yield* ClaudeProbeCache.ClaudeProbeCache; + const homes = ["a", "b", "c", "d", "e"].map((name) => `/homes/${name}`); + + const probes = yield* Effect.forEach(homes, (home) => cache.capabilities(input(home)), { + concurrency: "unbounded", + }).pipe(Effect.forkChild); + yield* Effect.yieldNow; + assert.equal(query.mock.calls.length, 3); + + release(); + const results = yield* Fiber.join(probes); + assert.equal(query.mock.calls.length, 5); + assert.isTrue(results.every((result) => result !== undefined)); + }).pipe(Effect.scoped, Effect.provide(testLayer)), +); From cd071286951d72e6dfb8f9cd236d3e6b1d96021d Mon Sep 17 00:00:00 2001 From: Theo Browne Date: Fri, 25 Sep 2026 14:32:32 -0700 Subject: [PATCH 03/16] fix(server): Claude probe cache no longer thrashes past 64 instances The cache capacity was 64. Effect's Cache evicts the least recently used key, and each refresh reads the keys in the same order. So with 65 or more probe inputs, every read evicted a key the next refresh needed, and every refresh re-probed every instance. Raise the cap to 1024 and add a test that 100 inputs stay cached across two refreshes. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../provider/Drivers/ClaudeProbeCache.test.ts | 16 ++++++++++++++++ .../src/provider/Drivers/ClaudeProbeCache.ts | 7 ++++++- 2 files changed, 22 insertions(+), 1 deletion(-) diff --git a/apps/server/src/provider/Drivers/ClaudeProbeCache.test.ts b/apps/server/src/provider/Drivers/ClaudeProbeCache.test.ts index fa2a60e5d85a..d2e71d209469 100644 --- a/apps/server/src/provider/Drivers/ClaudeProbeCache.test.ts +++ b/apps/server/src/provider/Drivers/ClaudeProbeCache.test.ts @@ -121,6 +121,22 @@ it.effect("invalidate re-probes only that input", () => }).pipe(Effect.scoped, Effect.provide(testLayer)), ); +it.effect("keeps every input cached across refreshes when there are many instances", () => + Effect.gen(function* () { + const query = yield* mockSdk(); + const cache = yield* ClaudeProbeCache.ClaudeProbeCache; + const homes = Array.from({ length: 100 }, (_, index) => `/homes/${index}`); + const refresh = Effect.forEach(homes, (home) => cache.capabilities(input(home)), { + concurrency: "unbounded", + }); + + yield* refresh; + yield* refresh; + + assert.equal(query.mock.calls.length, homes.length); + }).pipe(Effect.scoped, Effect.provide(testLayer)), +); + it.effect("runs at most 3 SDK probes at once", () => Effect.gen(function* () { const { promise: released, resolve: release } = Promise.withResolvers(); diff --git a/apps/server/src/provider/Drivers/ClaudeProbeCache.ts b/apps/server/src/provider/Drivers/ClaudeProbeCache.ts index a5603485dbd3..dc7e6dab29d5 100644 --- a/apps/server/src/provider/Drivers/ClaudeProbeCache.ts +++ b/apps/server/src/provider/Drivers/ClaudeProbeCache.ts @@ -33,6 +33,11 @@ const PROBE_TTL = Duration.minutes(5); // A failed probe marks every instance on that home as unverified, so retry soon. const FAILED_PROBE_TTL = Duration.seconds(30); const MAX_CONCURRENT_PROBES = 3; +// Keep this far above any real instance count. The cache evicts the least +// recently used key, and every refresh reads the keys in the same order, so a +// cap below the live key count makes each refresh re-probe every instance. +// Stale keys only come from instance edits, and entries are small. +const MAX_CACHED_PROBES = 1024; export class ClaudeProbeCache extends Context.Service< ClaudeProbeCache, @@ -62,7 +67,7 @@ export const layer = Layer.effect( ), ), { - capacity: 64, + capacity: MAX_CACHED_PROBES, timeToLive: (exit) => Exit.isSuccess(exit) && exit.value !== undefined ? PROBE_TTL : FAILED_PROBE_TTL, }, From 8691ad78f50a97f1b5f8839b6d8b3301b5e9742d Mon Sep 17 00:00:00 2001 From: Theo Browne Date: Fri, 25 Sep 2026 14:36:44 -0700 Subject: [PATCH 04/16] refactor(server): export make for the Claude probe cache Co-Authored-By: Claude Opus 5.5 (1M context) --- .../src/provider/Drivers/ClaudeProbeCache.ts | 51 +++++++++---------- 1 file changed, 25 insertions(+), 26 deletions(-) diff --git a/apps/server/src/provider/Drivers/ClaudeProbeCache.ts b/apps/server/src/provider/Drivers/ClaudeProbeCache.ts index dc7e6dab29d5..17b6d7c298bc 100644 --- a/apps/server/src/provider/Drivers/ClaudeProbeCache.ts +++ b/apps/server/src/provider/Drivers/ClaudeProbeCache.ts @@ -51,30 +51,29 @@ export class ClaudeProbeCache extends Context.Service< } >()("t3/provider/Drivers/ClaudeProbeCache") {} -export const layer = Layer.effect( - ClaudeProbeCache, - Effect.gen(function* () { - const gate = yield* Semaphore.make(MAX_CONCURRENT_PROBES); - // The probe keeps its own timeout inside the gate, so time spent waiting - // for a permit cannot turn into a false failure. - const cache = yield* Cache.makeWith( - (input: ClaudeProbeInput) => - gate.withPermits(1)( - probeClaudeCapabilities( - input, - mergeProviderInstanceEnvironment(input.environment), - input.cwd, - ), +export const make = Effect.gen(function* () { + const gate = yield* Semaphore.make(MAX_CONCURRENT_PROBES); + // The probe keeps its own timeout inside the gate, so time spent waiting + // for a permit cannot turn into a false failure. + const cache = yield* Cache.makeWith( + (input: ClaudeProbeInput) => + gate.withPermits(1)( + probeClaudeCapabilities( + input, + mergeProviderInstanceEnvironment(input.environment), + input.cwd, ), - { - capacity: MAX_CACHED_PROBES, - timeToLive: (exit) => - Exit.isSuccess(exit) && exit.value !== undefined ? PROBE_TTL : FAILED_PROBE_TTL, - }, - ); - return { - capabilities: (input) => Cache.get(cache, input), - invalidate: (input) => Cache.invalidate(cache, input), - } satisfies ClaudeProbeCache["Service"]; - }), -); + ), + { + capacity: MAX_CACHED_PROBES, + timeToLive: (exit) => + Exit.isSuccess(exit) && exit.value !== undefined ? PROBE_TTL : FAILED_PROBE_TTL, + }, + ); + return { + capabilities: (input) => Cache.get(cache, input), + invalidate: (input) => Cache.invalidate(cache, input), + } satisfies ClaudeProbeCache["Service"]; +}); + +export const layer = Layer.effect(ClaudeProbeCache, make); From e84169fe5121ff59952786fa3b8e12c71ee81cbd Mon Sep 17 00:00:00 2001 From: Theo Browne Date: Fri, 25 Sep 2026 14:44:27 -0700 Subject: [PATCH 05/16] fix(server): keep ClaudeProbeCache make module-private Knip flags the exported make as unused. The layer is the only consumer. Co-Authored-By: Claude Opus 5.5 (1M context) --- apps/server/src/provider/Drivers/ClaudeProbeCache.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/apps/server/src/provider/Drivers/ClaudeProbeCache.ts b/apps/server/src/provider/Drivers/ClaudeProbeCache.ts index 17b6d7c298bc..11ccda3dc129 100644 --- a/apps/server/src/provider/Drivers/ClaudeProbeCache.ts +++ b/apps/server/src/provider/Drivers/ClaudeProbeCache.ts @@ -51,7 +51,7 @@ export class ClaudeProbeCache extends Context.Service< } >()("t3/provider/Drivers/ClaudeProbeCache") {} -export const make = Effect.gen(function* () { +const make = Effect.gen(function* () { const gate = yield* Semaphore.make(MAX_CONCURRENT_PROBES); // The probe keeps its own timeout inside the gate, so time spent waiting // for a permit cannot turn into a false failure. From e996211969127299e5b98657086a7d9654006dc3 Mon Sep 17 00:00:00 2001 From: Theo Browne Date: Fri, 25 Sep 2026 18:09:44 -0700 Subject: [PATCH 06/16] refactor(server): export ClaudeProbeCache make like other services Mark make @public, as resetCreditCoordinator and other service modules do, so knip and the service conventions agree. Say the probe key holds every instance input, not everything the probe reads, and group the test import with the other driver imports. Co-Authored-By: Claude Opus 5.5 (1M context) --- apps/server/src/provider/Drivers/ClaudeProbeCache.ts | 9 +++++---- .../provider/Layers/ProviderInstanceRegistryLive.test.ts | 2 +- 2 files changed, 6 insertions(+), 5 deletions(-) diff --git a/apps/server/src/provider/Drivers/ClaudeProbeCache.ts b/apps/server/src/provider/Drivers/ClaudeProbeCache.ts index 11ccda3dc129..d8e111b94b9e 100644 --- a/apps/server/src/provider/Drivers/ClaudeProbeCache.ts +++ b/apps/server/src/provider/Drivers/ClaudeProbeCache.ts @@ -20,9 +20,9 @@ import { type ClaudeCapabilitiesProbe, probeClaudeCapabilities } from "../Layers import { mergeProviderInstanceEnvironment } from "../ProviderInstanceEnvironment.ts"; /** - * Everything the probe reads, and also the cache key. The lookup gets only - * this value, so the probe cannot depend on an input the key leaves out. - * Keys compare structurally. + * Every instance input the probe reads, and also the cache key. The lookup + * gets only this value, so the probe cannot depend on an input the key leaves + * out. Keys compare structurally. */ export type ClaudeProbeInput = Pick & { readonly cwd: string; @@ -51,7 +51,8 @@ export class ClaudeProbeCache extends Context.Service< } >()("t3/provider/Drivers/ClaudeProbeCache") {} -const make = Effect.gen(function* () { +/** @public Service construction is part of the canonical Effect module API. */ +export const make = Effect.gen(function* () { const gate = yield* Semaphore.make(MAX_CONCURRENT_PROBES); // The probe keeps its own timeout inside the gate, so time spent waiting // for a permit cannot turn into a false failure. diff --git a/apps/server/src/provider/Layers/ProviderInstanceRegistryLive.test.ts b/apps/server/src/provider/Layers/ProviderInstanceRegistryLive.test.ts index 3cd3835ce570..0ed9ad315591 100644 --- a/apps/server/src/provider/Layers/ProviderInstanceRegistryLive.test.ts +++ b/apps/server/src/provider/Layers/ProviderInstanceRegistryLive.test.ts @@ -50,6 +50,7 @@ import { ServerConfig } from "../../config.ts"; import { expandHomePath } from "../../pathExpansion.ts"; import { ServerSettingsService } from "../../serverSettings.ts"; import { ClaudeDriver } from "../Drivers/ClaudeDriver.ts"; +import * as ClaudeProbeCache from "../Drivers/ClaudeProbeCache.ts"; import { CodexDriver } from "../Drivers/CodexDriver.ts"; import { CursorDriver } from "../Drivers/CursorDriver.ts"; import { GrokDriver } from "../Drivers/GrokDriver.ts"; @@ -57,7 +58,6 @@ import { OpenCodeDriver } from "../Drivers/OpenCodeDriver.ts"; import * as ModelManifest from "../ModelManifest.ts"; import { OpenCodeRuntimeLive } from "../opencodeRuntime.ts"; import * as ResetCreditCoordinator from "./resetCreditCoordinator.ts"; -import * as ClaudeProbeCache from "../Drivers/ClaudeProbeCache.ts"; import { NoOpProviderEventLoggers, ProviderEventLoggers } from "./ProviderEventLoggers.ts"; import { makeProviderInstanceRegistry } from "./ProviderInstanceRegistryLive.ts"; From 9fb00395ec06233e5d23b1ebcbec05f9b1c4078d Mon Sep 17 00:00:00 2001 From: Theo Browne Date: Fri, 25 Sep 2026 18:26:08 -0700 Subject: [PATCH 07/16] fix(server): Claude probe cache retries a failed usage read after 30 seconds Co-Authored-By: Claude Opus 5.5 (1M context) --- .../provider/Drivers/ClaudeProbeCache.test.ts | 27 ++++++++++++++----- .../src/provider/Drivers/ClaudeProbeCache.ts | 5 ++-- 2 files changed, 24 insertions(+), 8 deletions(-) diff --git a/apps/server/src/provider/Drivers/ClaudeProbeCache.test.ts b/apps/server/src/provider/Drivers/ClaudeProbeCache.test.ts index d2e71d209469..0365d6e04149 100644 --- a/apps/server/src/provider/Drivers/ClaudeProbeCache.test.ts +++ b/apps/server/src/provider/Drivers/ClaudeProbeCache.test.ts @@ -29,6 +29,13 @@ const failedQuery = () => initializationResult: () => Promise.reject(new Error("not logged in")), }) as ReturnType; +const failedUsageQuery = () => + ({ + initializationResult: async () => ({}), + usage_EXPERIMENTAL_MAY_CHANGE_DO_NOT_RELY_ON_THIS_API_YET: () => + Promise.reject(new Error("usage unavailable")), + }) as ReturnType; + // Stands in for the SDK. Each probe reports the Claude home it was started // with as the account email. Probes finish once `ready` resolves. const mockSdk = (ready: Promise = Promise.resolve()) => @@ -43,8 +50,10 @@ const mockSdk = (ready: Promise = Promise.resolve()) => commands: [{ name: "review", description: "Review changes", argumentHint: "" }], }; }, - usage_EXPERIMENTAL_MAY_CHANGE_DO_NOT_RELY_ON_THIS_API_YET: () => - Promise.reject(new Error("usage unavailable")), + usage_EXPERIMENTAL_MAY_CHANGE_DO_NOT_RELY_ON_THIS_API_YET: async () => ({ + rate_limits_available: false, + rate_limits: null, + }), }) as ReturnType, ); yield* Effect.addFinalizer(() => Effect.sync(() => query.mockRestore())); @@ -86,10 +95,10 @@ it.effect("different homes or instance env vars run separate probes", () => }).pipe(Effect.scoped, Effect.provide(testLayer)), ); -it.effect("retries a failed probe after a short TTL and keeps a success longer", () => +it.effect("retries a failed probe or usage read after a short TTL and keeps a success longer", () => Effect.gen(function* () { const query = yield* mockSdk(); - query.mockImplementationOnce(failedQuery); + query.mockImplementationOnce(failedQuery).mockImplementationOnce(failedUsageQuery); const cache = yield* ClaudeProbeCache.ClaudeProbeCache; assert.equal(yield* cache.capabilities(input("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/homes/work")), undefined); @@ -97,12 +106,18 @@ it.effect("retries a failed probe after a short TTL and keeps a success longer", assert.equal(query.mock.calls.length, 1); yield* TestClock.adjust("30 seconds"); - assert.match((yield* cache.capabilities(input("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/homes/work")))?.email ?? "", /work$/); + const withoutUsage = yield* cache.capabilities(input("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/homes/work")); + assert.isDefined(withoutUsage); + assert.isUndefined(withoutUsage?.usage); assert.equal(query.mock.calls.length, 2); + yield* TestClock.adjust("30 seconds"); + assert.isDefined((yield* cache.capabilities(input("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/homes/work")))?.usage); + assert.equal(query.mock.calls.length, 3); + yield* TestClock.adjust("1 minute"); yield* cache.capabilities(input("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/homes/work")); - assert.equal(query.mock.calls.length, 2); + assert.equal(query.mock.calls.length, 3); }).pipe(Effect.scoped, Effect.provide(testLayer)), ); diff --git a/apps/server/src/provider/Drivers/ClaudeProbeCache.ts b/apps/server/src/provider/Drivers/ClaudeProbeCache.ts index d8e111b94b9e..01f289955066 100644 --- a/apps/server/src/provider/Drivers/ClaudeProbeCache.ts +++ b/apps/server/src/provider/Drivers/ClaudeProbeCache.ts @@ -30,7 +30,8 @@ export type ClaudeProbeInput = Pick & }; const PROBE_TTL = Duration.minutes(5); -// A failed probe marks every instance on that home as unverified, so retry soon. +// A failed probe or usage read leaves every instance with that input unverified +// or without limits, so retry soon. const FAILED_PROBE_TTL = Duration.seconds(30); const MAX_CONCURRENT_PROBES = 3; // Keep this far above any real instance count. The cache evicts the least @@ -68,7 +69,7 @@ export const make = Effect.gen(function* () { { capacity: MAX_CACHED_PROBES, timeToLive: (exit) => - Exit.isSuccess(exit) && exit.value !== undefined ? PROBE_TTL : FAILED_PROBE_TTL, + Exit.isSuccess(exit) && exit.value?.usage !== undefined ? PROBE_TTL : FAILED_PROBE_TTL, }, ); return { From 9161e7dbcb8b652944a66ea69e25e043757a0166 Mon Sep 17 00:00:00 2001 From: Theo Browne Date: Fri, 25 Sep 2026 19:35:27 -0700 Subject: [PATCH 08/16] test(server): an empty Claude home and an explicit ~/.claude run separate probes Both resolve to ~/.claude, but an explicit CLAUDE_CONFIG_DIR is a separate login to the CLI. Pin this so a later key normalization cannot merge them. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../src/provider/Drivers/ClaudeProbeCache.test.ts | 15 +++++++++++++++ 1 file changed, 15 insertions(+) diff --git a/apps/server/src/provider/Drivers/ClaudeProbeCache.test.ts b/apps/server/src/provider/Drivers/ClaudeProbeCache.test.ts index 0365d6e04149..87b2897e230a 100644 --- a/apps/server/src/provider/Drivers/ClaudeProbeCache.test.ts +++ b/apps/server/src/provider/Drivers/ClaudeProbeCache.test.ts @@ -95,6 +95,21 @@ it.effect("different homes or instance env vars run separate probes", () => }).pipe(Effect.scoped, Effect.provide(testLayer)), ); +// Both point at ~/.claude, but an explicit CLAUDE_CONFIG_DIR is a separate +// login to the CLI (its own keychain entry and .claude.json). +it.effect("an empty home and an explicit ~/.claude run separate probes", () => + Effect.gen(function* () { + const query = yield* mockSdk(); + const cache = yield* ClaudeProbeCache.ClaudeProbeCache; + + yield* cache.capabilities(input("")); + const explicit = yield* cache.capabilities(input("~/.claude")); + + assert.equal(query.mock.calls.length, 2); + assert.match(explicit?.email ?? "", /\.claude$/); + }).pipe(Effect.scoped, Effect.provide(testLayer)), +); + it.effect("retries a failed probe or usage read after a short TTL and keeps a success longer", () => Effect.gen(function* () { const query = yield* mockSdk(); From ce9d3b301da9cea8637dbeb18cf5c4d634ce1174 Mon Sep 17 00:00:00 2001 From: Theo Browne Date: Fri, 25 Sep 2026 20:51:48 -0700 Subject: [PATCH 09/16] fix(server): shared Claude probes keep honest usage times and run without a gate Usage limits now carry the time the probe read them, not the time of the status check that reused a cached probe. A cached read that is older than the published limits (for example a turn's update) no longer replaces them. Remove the 3-probe gate. Sharing already runs one probe per distinct input, and with about 8 s per probe the gate made status for 9 homes take about 24 s instead of about 8 s. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../provider/Drivers/ClaudeProbeCache.test.ts | 39 +++---------- .../src/provider/Drivers/ClaudeProbeCache.ts | 21 +++---- .../Layers/ClaudeCapabilitiesProbe.test.ts | 1 + .../src/provider/Layers/ClaudeProvider.ts | 15 ++++- .../provider/Layers/ProviderRegistry.test.ts | 58 +++++++++++-------- .../src/provider/makeManagedServerProvider.ts | 3 + .../src/provider/providerUsageLimits.test.ts | 30 +++++++++- .../src/provider/providerUsageLimits.ts | 13 +++++ 8 files changed, 108 insertions(+), 72 deletions(-) diff --git a/apps/server/src/provider/Drivers/ClaudeProbeCache.test.ts b/apps/server/src/provider/Drivers/ClaudeProbeCache.test.ts index 87b2897e230a..d7ef1086bbcc 100644 --- a/apps/server/src/provider/Drivers/ClaudeProbeCache.test.ts +++ b/apps/server/src/provider/Drivers/ClaudeProbeCache.test.ts @@ -2,7 +2,6 @@ import * as ClaudeSdk from "@anthropic-ai/claude-agent-sdk"; import * as NodeServices from "@effect/platform-node/NodeServices"; import { assert, it } from "@effect/vitest"; import * as Effect from "effect/Effect"; -import * as Fiber from "effect/Fiber"; import * as Layer from "effect/Layer"; import * as TestClock from "effect/testing/TestClock"; import { vi } from "vite-plus/test"; @@ -37,19 +36,16 @@ const failedUsageQuery = () => }) as ReturnType; // Stands in for the SDK. Each probe reports the Claude home it was started -// with as the account email. Probes finish once `ready` resolves. -const mockSdk = (ready: Promise = Promise.resolve()) => +// with as the account email. +const mockSdk = () => Effect.gen(function* () { const query = vi.spyOn(ClaudeSdk, "query").mockImplementation( ({ options }) => ({ - initializationResult: async () => { - await ready; - return { - account: { email: options?.env?.CLAUDE_CONFIG_DIR ?? "" }, - commands: [{ name: "review", description: "Review changes", argumentHint: "" }], - }; - }, + initializationResult: async () => ({ + account: { email: options?.env?.CLAUDE_CONFIG_DIR ?? "" }, + commands: [{ name: "review", description: "Review changes", argumentHint: "" }], + }), usage_EXPERIMENTAL_MAY_CHANGE_DO_NOT_RELY_ON_THIS_API_YET: async () => ({ rate_limits_available: false, rate_limits: null, @@ -69,12 +65,15 @@ it.effect("instances with the same probe input share one probe", () => [cache.capabilities(input("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/homes/work")), cache.capabilities(input("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/homes/work"))], { concurrency: "unbounded" }, ); + yield* TestClock.adjust("2 minutes"); const later = yield* cache.capabilities(input("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/homes/work")); assert.equal(query.mock.calls.length, 1); assert.match(first?.email ?? "", /work$/); assert.deepEqual(second, first); + // A later reader sees the probe's own time, not its read time. assert.deepEqual(later, first); + assert.equal(later?.checkedAt, "1970-01-01T00:00:00.000Z"); }).pipe(Effect.scoped, Effect.provide(testLayer)), ); @@ -166,23 +165,3 @@ it.effect("keeps every input cached across refreshes when there are many instanc assert.equal(query.mock.calls.length, homes.length); }).pipe(Effect.scoped, Effect.provide(testLayer)), ); - -it.effect("runs at most 3 SDK probes at once", () => - Effect.gen(function* () { - const { promise: released, resolve: release } = Promise.withResolvers(); - const query = yield* mockSdk(released); - const cache = yield* ClaudeProbeCache.ClaudeProbeCache; - const homes = ["a", "b", "c", "d", "e"].map((name) => `/homes/${name}`); - - const probes = yield* Effect.forEach(homes, (home) => cache.capabilities(input(home)), { - concurrency: "unbounded", - }).pipe(Effect.forkChild); - yield* Effect.yieldNow; - assert.equal(query.mock.calls.length, 3); - - release(); - const results = yield* Fiber.join(probes); - assert.equal(query.mock.calls.length, 5); - assert.isTrue(results.every((result) => result !== undefined)); - }).pipe(Effect.scoped, Effect.provide(testLayer)), -); diff --git a/apps/server/src/provider/Drivers/ClaudeProbeCache.ts b/apps/server/src/provider/Drivers/ClaudeProbeCache.ts index 01f289955066..15f7cbbeb448 100644 --- a/apps/server/src/provider/Drivers/ClaudeProbeCache.ts +++ b/apps/server/src/provider/Drivers/ClaudeProbeCache.ts @@ -2,8 +2,10 @@ * One server-wide cache for the Claude capabilities probe. Claude instances * whose probe inputs match (same binary, home, cwd and instance env vars) * read the same account, so they share one cached result instead of each - * starting its own SDK session. A small gate limits how many SDK probes run - * at once across all instances. + * starting its own SDK session. Concurrent reads of one input join one probe, + * so a full refresh runs one probe per distinct input. There is no global + * limit on purpose: a probe takes seconds, so a queue would delay the status + * of every input past the limit. * * @module provider/Drivers/ClaudeProbeCache */ @@ -14,7 +16,6 @@ import * as Duration from "effect/Duration"; import * as Effect from "effect/Effect"; import * as Exit from "effect/Exit"; import * as Layer from "effect/Layer"; -import * as Semaphore from "effect/Semaphore"; import { type ClaudeCapabilitiesProbe, probeClaudeCapabilities } from "../Layers/ClaudeProvider.ts"; import { mergeProviderInstanceEnvironment } from "../ProviderInstanceEnvironment.ts"; @@ -33,7 +34,6 @@ const PROBE_TTL = Duration.minutes(5); // A failed probe or usage read leaves every instance with that input unverified // or without limits, so retry soon. const FAILED_PROBE_TTL = Duration.seconds(30); -const MAX_CONCURRENT_PROBES = 3; // Keep this far above any real instance count. The cache evicts the least // recently used key, and every refresh reads the keys in the same order, so a // cap below the live key count makes each refresh re-probe every instance. @@ -54,17 +54,12 @@ export class ClaudeProbeCache extends Context.Service< /** @public Service construction is part of the canonical Effect module API. */ export const make = Effect.gen(function* () { - const gate = yield* Semaphore.make(MAX_CONCURRENT_PROBES); - // The probe keeps its own timeout inside the gate, so time spent waiting - // for a permit cannot turn into a false failure. const cache = yield* Cache.makeWith( (input: ClaudeProbeInput) => - gate.withPermits(1)( - probeClaudeCapabilities( - input, - mergeProviderInstanceEnvironment(input.environment), - input.cwd, - ), + probeClaudeCapabilities( + input, + mergeProviderInstanceEnvironment(input.environment), + input.cwd, ), { capacity: MAX_CACHED_PROBES, diff --git a/apps/server/src/provider/Layers/ClaudeCapabilitiesProbe.test.ts b/apps/server/src/provider/Layers/ClaudeCapabilitiesProbe.test.ts index 232b8cc02d00..cf03625be768 100644 --- a/apps/server/src/provider/Layers/ClaudeCapabilitiesProbe.test.ts +++ b/apps/server/src/provider/Layers/ClaudeCapabilitiesProbe.test.ts @@ -161,6 +161,7 @@ it.layer(NodeServices.layer)("Claude capability probe SDK boundary", (it) => { rate_limits_available: true, rate_limits: { five_hour: { utilization: 12, resets_at: "2026-07-18T14:39:00Z" } }, }, + checkedAt: "1970-01-01T00:00:00.000Z", }); // @effect-diagnostics-next-line preferSchemaOverJson:off diff --git a/apps/server/src/provider/Layers/ClaudeProvider.ts b/apps/server/src/provider/Layers/ClaudeProvider.ts index b296a9fd1cda..3aa605bca94f 100644 --- a/apps/server/src/provider/Layers/ClaudeProvider.ts +++ b/apps/server/src/provider/Layers/ClaudeProvider.ts @@ -242,6 +242,11 @@ export type ClaudeCapabilitiesProbe = { * otherwise successful response mean the account has none (API key). */ readonly usage?: Pick; + /** + * When the probe read usage. Instances share cached probes, so usage + * limits carry this time, not the time of the status check that read it. + */ + readonly checkedAt: string; }; function parseClaudeInitializationCommands( @@ -373,6 +378,7 @@ const probeClaudeCapabilities = ( rate_limits: usageResult.success.rate_limits, } : undefined; + const checkedAt = DateTime.formatIso(yield* DateTime.now); const account = init.account as | { readonly email?: string; @@ -388,6 +394,7 @@ const probeClaudeCapabilities = ( apiProvider: account?.apiProvider, slashCommands: parseClaudeInitializationCommands(init.commands), ...(usage ? { usage } : {}), + checkedAt, } satisfies ClaudeCapabilitiesProbe; }), ), @@ -563,14 +570,16 @@ export const checkClaudeProviderStatus = Effect.fn("checkClaudeProviderStatus")( subscriptionType: capabilities.subscriptionType, authMethod: capabilities.tokenSource, }) ?? apiProviderAuthMetadata(capabilities.apiProvider); + const usageCheckedAt = capabilities.checkedAt; const usageLimits = !capabilities.usage - ? makeUnavailableUsageLimits({ checkedAt, reason: "probeFailed" }) + ? makeUnavailableUsageLimits({ checkedAt: usageCheckedAt, reason: "probeFailed" }) : scopedLimitNames ? yield* recordClaudeUsageResponse(scopedLimitNames, { response: capabilities.usage, - checkedAt, + checkedAt: usageCheckedAt, }) - : claudeUsageResponseToLimits({ response: capabilities.usage, checkedAt }).limits; + : claudeUsageResponseToLimits({ response: capabilities.usage, checkedAt: usageCheckedAt }) + .limits; const resetCredits = resolveResetCredits && capabilities.subscriptionType && diff --git a/apps/server/src/provider/Layers/ProviderRegistry.test.ts b/apps/server/src/provider/Layers/ProviderRegistry.test.ts index abc2cd42fdb7..f93458313a5a 100644 --- a/apps/server/src/provider/Layers/ProviderRegistry.test.ts +++ b/apps/server/src/provider/Layers/ProviderRegistry.test.ts @@ -23,7 +23,6 @@ import { ProviderInstanceId, ServerSettings, type ServerProvider, - type ServerProviderSlashCommand, type ServerSettings as ContractServerSettings, } from "@t3tools/contracts"; import * as PlatformError from "effect/PlatformError"; @@ -34,7 +33,7 @@ import { createModelCapabilities } from "@t3tools/shared/model"; import { applyServerSettingsPatch } from "@t3tools/shared/serverSettings"; import { checkCodexProviderStatus, type CodexAppServerProviderSnapshot } from "./CodexProvider.ts"; -import { checkClaudeProviderStatus } from "./ClaudeProvider.ts"; +import { type ClaudeCapabilitiesProbe, checkClaudeProviderStatus } from "./ClaudeProvider.ts"; import * as BackgroundPolicy from "../../background/BackgroundPolicy.ts"; import { AntigravityInstallation } from "../AntigravityInstallation.ts"; import * as ModelManifest from "../ModelManifest.ts"; @@ -146,15 +145,7 @@ function booleanDescriptor(id: string, label: string) { }; } -type TestClaudeCapabilities = { - readonly email: string | undefined; - readonly subscriptionType: string | undefined; - readonly tokenSource: string | undefined; - readonly apiProvider: string | undefined; - readonly slashCommands: ReadonlyArray; -}; - -function claudeCapabilities(overrides: Partial = {}) { +function claudeCapabilities(overrides: Partial = {}) { return () => Effect.succeed({ email: undefined, @@ -162,12 +153,13 @@ function claudeCapabilities(overrides: Partial = {}) { tokenSource: undefined, apiProvider: undefined, slashCommands: [], + checkedAt: "2026-09-25T12:00:00.000Z", ...overrides, }); } const noClaudeCapabilities = () => - Effect.sync(() => undefined as TestClaudeCapabilities | undefined); + Effect.sync(() => undefined as ClaudeCapabilitiesProbe | undefined); function mockHandle(result: { stdout: string; stderr: string; code: number }) { return ChildProcessSpawner.makeHandle({ @@ -2758,19 +2750,13 @@ it.layer(Layer.mergeAll(NodeServices.layer, ServerSettingsModule.layerTest(), Te it.effect("reads banked resets only for subscription logins", () => Effect.gen(function* () { - const check = (overrides: Partial) => + const check = (overrides: Partial) => checkClaudeProviderStatus( defaultClaudeSettings, - () => - Effect.succeed({ - email: undefined, - subscriptionType: undefined, - tokenSource: undefined, - apiProvider: undefined, - slashCommands: [], - usage: { rate_limits_available: true, rate_limits: {} }, - ...overrides, - }), + claudeCapabilities({ + usage: { rate_limits_available: true, rate_limits: {} }, + ...overrides, + }), undefined, undefined, undefined, @@ -2792,6 +2778,32 @@ it.layer(Layer.mergeAll(NodeServices.layer, ServerSettingsModule.layerTest(), Te ), ); + // Instances share cached probes, so the usage can be older than the check. + it.effect("dates usage limits by the probe, not the status check", () => + Effect.gen(function* () { + const probedAt = "2026-09-25T12:00:00.000Z"; + const check = (usage: ClaudeCapabilitiesProbe["usage"]) => + checkClaudeProviderStatus( + defaultClaudeSettings, + claudeCapabilities({ checkedAt: probedAt, ...(usage ? { usage } : {}) }), + ); + const read = yield* check({ rate_limits_available: true, rate_limits: {} }); + const failed = yield* check(undefined); + assert.notStrictEqual(read.checkedAt, probedAt); + assert.strictEqual(read.usageLimits?.checkedAt, probedAt); + assert.strictEqual(failed.usageLimits?.unavailable?.reason, "probeFailed"); + assert.strictEqual(failed.usageLimits?.checkedAt, probedAt); + }).pipe( + Effect.provide( + mockSpawnerLayer((args) => { + const joined = args.join(" "); + if (joined === "--version") return { stdout: "1.0.0\n", stderr: "", code: 0 }; + throw new Error(`Unexpected args: ${joined}`); + }), + ), + ), + ); + it.effect("does not duplicate Claude in full subscription labels", () => Effect.gen(function* () { const status = yield* checkClaudeProviderStatus( diff --git a/apps/server/src/provider/makeManagedServerProvider.ts b/apps/server/src/provider/makeManagedServerProvider.ts index 119c577474c8..1c3189d1f648 100644 --- a/apps/server/src/provider/makeManagedServerProvider.ts +++ b/apps/server/src/provider/makeManagedServerProvider.ts @@ -4,6 +4,7 @@ import { ServerSettingsError, } from "@t3tools/contracts"; import { resolveServerBackgroundActivitySettings } from "@t3tools/shared/backgroundActivitySettings"; +import * as Clock from "effect/Clock"; import * as Duration from "effect/Duration"; import * as Effect from "effect/Effect"; import * as Equal from "effect/Equal"; @@ -150,6 +151,7 @@ export const makeManagedServerProvider = Effect.fn("makeManagedServerProvider")( return state.snapshot; } + const checkStartedAt = yield* Clock.currentTimeMillis; const probedSnapshot = yield* input.checkProvider; const { snapshot: nextSnapshot, generation: nextGeneration } = yield* Ref.modify( snapshotStateRef, @@ -162,6 +164,7 @@ export const makeManagedServerProvider = Effect.fn("makeManagedServerProvider")( resolveUsageLimitsAfterProbe({ published: state.snapshot.usageLimits, probed: probedSnapshot.usageLimits, + checkStartedAt, }), ); return [ diff --git a/apps/server/src/provider/providerUsageLimits.test.ts b/apps/server/src/provider/providerUsageLimits.test.ts index 6e288ddd3a36..64a660e99ea0 100644 --- a/apps/server/src/provider/providerUsageLimits.test.ts +++ b/apps/server/src/provider/providerUsageLimits.test.ts @@ -1,3 +1,4 @@ +import type { ServerProviderUsageLimits } from "@t3tools/contracts"; import { describe, expect, it } from "vite-plus/test"; import { applyUsageLimitsUpdate, resolveUsageLimitsAfterProbe } from "./providerUsageLimits.ts"; @@ -79,11 +80,34 @@ describe("applyUsageLimitsUpdate", () => { }); describe("resolveUsageLimitsAfterProbe", () => { + const checkStartedAt = Date.parse(checkedAt); + it("keeps the last good windows through a failed probe but not an unsupported one", () => { const failed = { checkedAt, windows: [], unavailable: { reason: "probeFailed" as const } }; const unsupported = { checkedAt, windows: [], unavailable: { reason: "unsupported" as const } }; - expect(resolveUsageLimitsAfterProbe({ published, probed: failed })).toBe(published); - expect(resolveUsageLimitsAfterProbe({ published, probed: unsupported })).toBe(unsupported); - expect(resolveUsageLimitsAfterProbe({ published: undefined, probed: failed })).toBe(failed); + expect(resolveUsageLimitsAfterProbe({ published, probed: failed, checkStartedAt })).toBe( + published, + ); + expect(resolveUsageLimitsAfterProbe({ published, probed: unsupported, checkStartedAt })).toBe( + unsupported, + ); + expect( + resolveUsageLimitsAfterProbe({ published: undefined, probed: failed, checkStartedAt }), + ).toBe(failed); + }); + + it("keeps a turn update over an older cached read, not over a fresh probe", () => { + const startedAt = Date.parse("2026-09-03T12:05:00.000Z"); + const turnUpdate = { checkedAt: "2026-09-03T12:04:00.000Z", windows: [session] }; + // Claude instances share cached probes, so a read can predate the check. + const cached = { checkedAt: "2026-09-03T12:01:00.000Z", windows: [session, weekly] }; + // Codex dates its read at the check start, and a turn can update mid-probe. + const fresh = { checkedAt: "2026-09-03T12:05:00.000Z", windows: [session, weekly] }; + const midProbeUpdate = { checkedAt: "2026-09-03T12:05:02.000Z", windows: [session] }; + const resolve = (current: ServerProviderUsageLimits, probed: ServerProviderUsageLimits) => + resolveUsageLimitsAfterProbe({ published: current, probed, checkStartedAt: startedAt }); + expect(resolve(turnUpdate, cached)).toBe(turnUpdate); + expect(resolve(midProbeUpdate, fresh)).toBe(fresh); + expect(resolve(published, cached)).toBe(cached); }); }); diff --git a/apps/server/src/provider/providerUsageLimits.ts b/apps/server/src/provider/providerUsageLimits.ts index ea8d0d1d029f..d67286b9df25 100644 --- a/apps/server/src/provider/providerUsageLimits.ts +++ b/apps/server/src/provider/providerUsageLimits.ts @@ -121,14 +121,27 @@ function usageWindowEquals(a: ServerProviderUsageWindow, b: ServerProviderUsageW * per-window epoch bookkeeping needed to reconcile the two was more code * than the sub-second regression it prevented. The next runtime event * corrects it. + * + * Limits stamped before the check started came from a cache (Claude + * instances share one probe), so they can be minutes old. Those do not + * replace newer published limits, such as a turn's update. */ export function resolveUsageLimitsAfterProbe(input: { readonly published: ServerProviderUsageLimits | undefined; readonly probed: ServerProviderUsageLimits | undefined; + /** Epoch millis when the status check that produced `probed` started. */ + readonly checkStartedAt: number; }): ServerProviderUsageLimits | undefined { const { published, probed } = input; if (probed?.unavailable?.reason === "probeFailed" && published && !published.unavailable) { return published; } + if ( + published && + probed && + Date.parse(probed.checkedAt) < Math.min(input.checkStartedAt, Date.parse(published.checkedAt)) + ) { + return published; + } return probed; } From bbfa57571b38f16f24e8b285f88d2f460ae39202 Mon Sep 17 00:00:00 2001 From: Theo Browne Date: Fri, 25 Sep 2026 21:08:29 -0700 Subject: [PATCH 10/16] fix(server): cached Claude reads keep fresh reset credits, and a repeat probe failure backs off A check that reads a cached Claude probe kept the published limits whole, so it also kept the old reset credit list instead of the credits it just read. It now keeps the published windows and takes the fresh credits. A failed probe or usage read retried after 30 s every time, so a signed-out input probed on every 1 min check. Only the first failure in a row retries after 30 s now. A repeat failure waits 5 min, like a success. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../provider/Drivers/ClaudeProbeCache.test.ts | 32 ++++++++++++------- .../src/provider/Drivers/ClaudeProbeCache.ts | 23 ++++++++++--- .../src/provider/providerUsageLimits.test.ts | 9 +++++- .../src/provider/providerUsageLimits.ts | 6 ++-- 4 files changed, 50 insertions(+), 20 deletions(-) diff --git a/apps/server/src/provider/Drivers/ClaudeProbeCache.test.ts b/apps/server/src/provider/Drivers/ClaudeProbeCache.test.ts index d7ef1086bbcc..35063a5ae3e1 100644 --- a/apps/server/src/provider/Drivers/ClaudeProbeCache.test.ts +++ b/apps/server/src/provider/Drivers/ClaudeProbeCache.test.ts @@ -109,29 +109,37 @@ it.effect("an empty home and an explicit ~/.claude run separate probes", () => }).pipe(Effect.scoped, Effect.provide(testLayer)), ); -it.effect("retries a failed probe or usage read after a short TTL and keeps a success longer", () => +it.effect("retries a first failure after 30 seconds and a repeat failure after 5 minutes", () => Effect.gen(function* () { const query = yield* mockSdk(); query.mockImplementationOnce(failedQuery).mockImplementationOnce(failedUsageQuery); const cache = yield* ClaudeProbeCache.ClaudeProbeCache; + const read = cache.capabilities(input("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/homes/work")); - assert.equal(yield* cache.capabilities(input("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/homes/work")), undefined); - assert.equal(yield* cache.capabilities(input("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/homes/work")), undefined); + assert.equal(yield* read, undefined); + yield* TestClock.adjust("29 seconds"); + yield* read; assert.equal(query.mock.calls.length, 1); - yield* TestClock.adjust("30 seconds"); - const withoutUsage = yield* cache.capabilities(input("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/homes/work")); - assert.isDefined(withoutUsage); - assert.isUndefined(withoutUsage?.usage); + // A failed usage read is a failure too, and this one repeats. + yield* TestClock.adjust("1 second"); + assert.isUndefined((yield* read)?.usage); + assert.equal(query.mock.calls.length, 2); + yield* TestClock.adjust("4 minutes"); + yield* read; assert.equal(query.mock.calls.length, 2); - - yield* TestClock.adjust("30 seconds"); - assert.isDefined((yield* cache.capabilities(input("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/homes/work")))?.usage); - assert.equal(query.mock.calls.length, 3); yield* TestClock.adjust("1 minute"); - yield* cache.capabilities(input("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/homes/work")); + assert.isDefined((yield* read)?.usage); assert.equal(query.mock.calls.length, 3); + + // A success ends the streak, so the next failure retries soon again. + query.mockImplementationOnce(failedQuery); + yield* TestClock.adjust("5 minutes"); + assert.equal(yield* read, undefined); + yield* TestClock.adjust("30 seconds"); + assert.isDefined(yield* read); + assert.equal(query.mock.calls.length, 5); }).pipe(Effect.scoped, Effect.provide(testLayer)), ); diff --git a/apps/server/src/provider/Drivers/ClaudeProbeCache.ts b/apps/server/src/provider/Drivers/ClaudeProbeCache.ts index 15f7cbbeb448..cd19bbde2fbf 100644 --- a/apps/server/src/provider/Drivers/ClaudeProbeCache.ts +++ b/apps/server/src/provider/Drivers/ClaudeProbeCache.ts @@ -16,6 +16,7 @@ import * as Duration from "effect/Duration"; import * as Effect from "effect/Effect"; import * as Exit from "effect/Exit"; import * as Layer from "effect/Layer"; +import * as MutableHashSet from "effect/MutableHashSet"; import { type ClaudeCapabilitiesProbe, probeClaudeCapabilities } from "../Layers/ClaudeProvider.ts"; import { mergeProviderInstanceEnvironment } from "../ProviderInstanceEnvironment.ts"; @@ -32,8 +33,9 @@ export type ClaudeProbeInput = Pick & const PROBE_TTL = Duration.minutes(5); // A failed probe or usage read leaves every instance with that input unverified -// or without limits, so retry soon. -const FAILED_PROBE_TTL = Duration.seconds(30); +// or without limits, so the first failure retries soon. A repeat failure +// (signed out, no usage endpoint) waits the full TTL, like a success. +const FIRST_FAILURE_TTL = Duration.seconds(30); // Keep this far above any real instance count. The cache evicts the least // recently used key, and every refresh reads the keys in the same order, so a // cap below the live key count makes each refresh re-probe every instance. @@ -54,21 +56,32 @@ export class ClaudeProbeCache extends Context.Service< /** @public Service construction is part of the canonical Effect module API. */ export const make = Effect.gen(function* () { + // Inputs whose last probe failed. Compares keys like the cache does. + const failing = MutableHashSet.empty(); const cache = yield* Cache.makeWith( (input: ClaudeProbeInput) => probeClaudeCapabilities( input, mergeProviderInstanceEnvironment(input.environment), input.cwd, + ).pipe( + Effect.map((probe) => { + if (probe?.usage !== undefined) { + MutableHashSet.remove(failing, input); + return { probe, timeToLive: PROBE_TTL }; + } + const repeat = MutableHashSet.has(failing, input); + MutableHashSet.add(failing, input); + return { probe, timeToLive: repeat ? PROBE_TTL : FIRST_FAILURE_TTL }; + }), ), { capacity: MAX_CACHED_PROBES, - timeToLive: (exit) => - Exit.isSuccess(exit) && exit.value?.usage !== undefined ? PROBE_TTL : FAILED_PROBE_TTL, + timeToLive: (exit) => (Exit.isSuccess(exit) ? exit.value.timeToLive : FIRST_FAILURE_TTL), }, ); return { - capabilities: (input) => Cache.get(cache, input), + capabilities: (input) => Cache.get(cache, input).pipe(Effect.map(({ probe }) => probe)), invalidate: (input) => Cache.invalidate(cache, input), } satisfies ClaudeProbeCache["Service"]; }); diff --git a/apps/server/src/provider/providerUsageLimits.test.ts b/apps/server/src/provider/providerUsageLimits.test.ts index 64a660e99ea0..15a06e6e062f 100644 --- a/apps/server/src/provider/providerUsageLimits.test.ts +++ b/apps/server/src/provider/providerUsageLimits.test.ts @@ -106,7 +106,14 @@ describe("resolveUsageLimitsAfterProbe", () => { const midProbeUpdate = { checkedAt: "2026-09-03T12:05:02.000Z", windows: [session] }; const resolve = (current: ServerProviderUsageLimits, probed: ServerProviderUsageLimits) => resolveUsageLimitsAfterProbe({ published: current, probed, checkStartedAt: startedAt }); - expect(resolve(turnUpdate, cached)).toBe(turnUpdate); + expect(resolve(turnUpdate, cached)).toEqual(turnUpdate); + // The check read credits itself, so a spent credit does not survive. + const spent = { availableCount: 1, nextCreditId: "spent" }; + const current = { availableCount: 0 }; + expect( + resolve({ ...turnUpdate, resetCredits: spent }, { ...cached, resetCredits: current }), + ).toEqual({ ...turnUpdate, resetCredits: current }); + expect(resolve({ ...turnUpdate, resetCredits: spent }, cached)).toEqual(turnUpdate); expect(resolve(midProbeUpdate, fresh)).toBe(fresh); expect(resolve(published, cached)).toBe(cached); }); diff --git a/apps/server/src/provider/providerUsageLimits.ts b/apps/server/src/provider/providerUsageLimits.ts index d67286b9df25..617d791a47a1 100644 --- a/apps/server/src/provider/providerUsageLimits.ts +++ b/apps/server/src/provider/providerUsageLimits.ts @@ -124,7 +124,8 @@ function usageWindowEquals(a: ServerProviderUsageWindow, b: ServerProviderUsageW * * Limits stamped before the check started came from a cache (Claude * instances share one probe), so they can be minutes old. Those do not - * replace newer published limits, such as a turn's update. + * replace newer published windows, such as a turn's update. Reset credits + * still come from the probe, because each check reads them itself. */ export function resolveUsageLimitsAfterProbe(input: { readonly published: ServerProviderUsageLimits | undefined; @@ -141,7 +142,8 @@ export function resolveUsageLimitsAfterProbe(input: { probed && Date.parse(probed.checkedAt) < Math.min(input.checkStartedAt, Date.parse(published.checkedAt)) ) { - return published; + const { resetCredits: _published, ...windows } = published; + return probed.resetCredits ? { ...windows, resetCredits: probed.resetCredits } : windows; } return probed; } From acb67c2eb9208ef9c47548f1210ea0258f6bc1e5 Mon Sep 17 00:00:00 2001 From: Theo Browne Date: Fri, 25 Sep 2026 21:09:57 -0700 Subject: [PATCH 11/16] docs(server): say reset credits come from the probed limits Co-Authored-By: Claude Opus 5.5 (1M context) --- apps/server/src/provider/providerUsageLimits.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/apps/server/src/provider/providerUsageLimits.ts b/apps/server/src/provider/providerUsageLimits.ts index 617d791a47a1..a37b0a6839ed 100644 --- a/apps/server/src/provider/providerUsageLimits.ts +++ b/apps/server/src/provider/providerUsageLimits.ts @@ -125,7 +125,7 @@ function usageWindowEquals(a: ServerProviderUsageWindow, b: ServerProviderUsageW * Limits stamped before the check started came from a cache (Claude * instances share one probe), so they can be minutes old. Those do not * replace newer published windows, such as a turn's update. Reset credits - * still come from the probe, because each check reads them itself. + * still come from `probed`, because each check reads them itself. */ export function resolveUsageLimitsAfterProbe(input: { readonly published: ServerProviderUsageLimits | undefined; From f88b84be1d6fb6dc005c4cbf538734fd750a30ff Mon Sep 17 00:00:00 2001 From: Theo Browne Date: Fri, 25 Sep 2026 21:15:03 -0700 Subject: [PATCH 12/16] fix(server): bound the Claude probe cache's failure memory The set of inputs whose last probe failed only shrank on a later success, so failed inputs from old instance configs stayed forever. It now clears when it reaches the cache's cap. A clear costs each input one early retry. Co-Authored-By: Claude Opus 5.5 (1M context) --- apps/server/src/provider/Drivers/ClaudeProbeCache.ts | 2 ++ 1 file changed, 2 insertions(+) diff --git a/apps/server/src/provider/Drivers/ClaudeProbeCache.ts b/apps/server/src/provider/Drivers/ClaudeProbeCache.ts index cd19bbde2fbf..a7f9be868193 100644 --- a/apps/server/src/provider/Drivers/ClaudeProbeCache.ts +++ b/apps/server/src/provider/Drivers/ClaudeProbeCache.ts @@ -71,6 +71,8 @@ export const make = Effect.gen(function* () { return { probe, timeToLive: PROBE_TTL }; } const repeat = MutableHashSet.has(failing, input); + // Bounded like the cache. A clear costs each input one early retry. + if (MutableHashSet.size(failing) >= MAX_CACHED_PROBES) MutableHashSet.clear(failing); MutableHashSet.add(failing, input); return { probe, timeToLive: repeat ? PROBE_TTL : FIRST_FAILURE_TTL }; }), From df89a24bcc4c01520d3e9f0a1cea983e6a274736 Mon Sep 17 00:00:00 2001 From: Theo Browne Date: Fri, 25 Sep 2026 21:43:12 -0700 Subject: [PATCH 13/16] fix(server): new Claude instances probe fresh, and a replaced probe cannot start a failure streak A new or edited Claude instance drops a sibling's finished probe for its input, so it starts from its own probe as it did before the cache was shared. It still joins a probe in flight, so instances created together at boot share one. Each probe now writes only its own failure record, so a probe that invalidate replaced cannot mark its input as failing after a newer probe succeeded. Adds a registry test for the driver wiring: two instances with the same home and env run one SDK probe, one with different env runs its own, and an edited instance probes again. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../src/provider/Drivers/ClaudeDriver.ts | 3 + .../provider/Drivers/ClaudeProbeCache.test.ts | 64 ++++++++++++++++ .../src/provider/Drivers/ClaudeProbeCache.ts | 58 ++++++++------ .../ProviderInstanceRegistryLive.test.ts | 75 +++++++++++++++++++ 4 files changed, 179 insertions(+), 21 deletions(-) diff --git a/apps/server/src/provider/Drivers/ClaudeDriver.ts b/apps/server/src/provider/Drivers/ClaudeDriver.ts index 871f3d5bac1e..00a3075acdc1 100644 --- a/apps/server/src/provider/Drivers/ClaudeDriver.ts +++ b/apps/server/src/provider/Drivers/ClaudeDriver.ts @@ -178,6 +178,9 @@ export const ClaudeDriver: ProviderDriver = { cwd, environment, } satisfies ClaudeProbeCache.ClaudeProbeInput; + // A new or edited instance starts from its own probe, not a sibling's + // finished one. It still joins a probe that is in flight. + yield* probeCache.dropFinished(probeInput); // Start the TTL-gated refresh without delaying provider readiness. The // next check observes a remote manifest after the background fetch lands. diff --git a/apps/server/src/provider/Drivers/ClaudeProbeCache.test.ts b/apps/server/src/provider/Drivers/ClaudeProbeCache.test.ts index 35063a5ae3e1..d5cdc36217f2 100644 --- a/apps/server/src/provider/Drivers/ClaudeProbeCache.test.ts +++ b/apps/server/src/provider/Drivers/ClaudeProbeCache.test.ts @@ -2,6 +2,7 @@ import * as ClaudeSdk from "@anthropic-ai/claude-agent-sdk"; import * as NodeServices from "@effect/platform-node/NodeServices"; import { assert, it } from "@effect/vitest"; import * as Effect from "effect/Effect"; +import * as Fiber from "effect/Fiber"; import * as Layer from "effect/Layer"; import * as TestClock from "effect/testing/TestClock"; import { vi } from "vite-plus/test"; @@ -56,6 +57,20 @@ const mockSdk = () => return query; }); +// Keeps the next probe in flight until `release`, then fails it. +const holdNextProbe = (query: Effect.Success>) => { + const started = Promise.withResolvers(); + const release = Promise.withResolvers(); + query.mockImplementationOnce(() => { + started.resolve(); + return { + initializationResult: () => + release.promise.then(() => Promise.reject(new Error("probe timed out"))), + } as ReturnType; + }); + return { started: started.promise, release: () => release.resolve() }; +}; + it.effect("instances with the same probe input share one probe", () => Effect.gen(function* () { const query = yield* mockSdk(); @@ -143,6 +158,55 @@ it.effect("retries a first failure after 30 seconds and a repeat failure after 5 }).pipe(Effect.scoped, Effect.provide(testLayer)), ); +// `invalidate` can replace a probe that is still running. When that old probe +// fails after the new one succeeded, the next failure is still a first one. +it.effect("a replaced probe that fails late does not start a failure streak", () => + Effect.gen(function* () { + const query = yield* mockSdk(); + const held = holdNextProbe(query); + const cache = yield* ClaudeProbeCache.ClaudeProbeCache; + const read = cache.capabilities(input("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/homes/work")); + + const replaced = yield* Effect.forkChild(read); + yield* Effect.promise(() => held.started); + yield* cache.invalidate(input("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/homes/work")); + assert.isDefined((yield* read)?.usage); + held.release(); + assert.equal(yield* Fiber.join(replaced), undefined); + + query.mockImplementationOnce(failedQuery); + yield* TestClock.adjust("5 minutes"); + assert.equal(yield* read, undefined); + yield* TestClock.adjust("30 seconds"); + assert.isDefined(yield* read); + assert.equal(query.mock.calls.length, 4); + }).pipe(Effect.scoped, Effect.provide(testLayer)), +); + +it.effect("dropFinished re-probes a finished input but joins a probe in flight", () => + Effect.gen(function* () { + const query = yield* mockSdk(); + const cache = yield* ClaudeProbeCache.ClaudeProbeCache; + const read = cache.capabilities(input("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/homes/work")); + + yield* read; + yield* cache.dropFinished(input("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/homes/work")); + yield* read; + assert.equal(query.mock.calls.length, 2); + + const held = holdNextProbe(query); + yield* cache.invalidate(input("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/homes/work")); + const first = yield* Effect.forkChild(read); + yield* Effect.promise(() => held.started); + yield* cache.dropFinished(input("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/homes/work")); + const second = yield* Effect.forkChild(read); + held.release(); + yield* Fiber.join(first); + yield* Fiber.join(second); + assert.equal(query.mock.calls.length, 3); + }).pipe(Effect.scoped, Effect.provide(testLayer)), +); + it.effect("invalidate re-probes only that input", () => Effect.gen(function* () { const query = yield* mockSdk(); diff --git a/apps/server/src/provider/Drivers/ClaudeProbeCache.ts b/apps/server/src/provider/Drivers/ClaudeProbeCache.ts index a7f9be868193..1c9c024b1b99 100644 --- a/apps/server/src/provider/Drivers/ClaudeProbeCache.ts +++ b/apps/server/src/provider/Drivers/ClaudeProbeCache.ts @@ -16,7 +16,8 @@ import * as Duration from "effect/Duration"; import * as Effect from "effect/Effect"; import * as Exit from "effect/Exit"; import * as Layer from "effect/Layer"; -import * as MutableHashSet from "effect/MutableHashSet"; +import * as MutableHashMap from "effect/MutableHashMap"; +import * as Option from "effect/Option"; import { type ClaudeCapabilitiesProbe, probeClaudeCapabilities } from "../Layers/ClaudeProvider.ts"; import { mergeProviderInstanceEnvironment } from "../ProviderInstanceEnvironment.ts"; @@ -49,34 +50,43 @@ export class ClaudeProbeCache extends Context.Service< readonly capabilities: ( input: ClaudeProbeInput, ) => Effect.Effect; - /** Drop the result for `input`, so the next read probes again. */ + /** + * Drop the result for `input`, even one still in flight, so the next read + * starts a new probe. + */ readonly invalidate: (input: ClaudeProbeInput) => Effect.Effect; + /** + * Drop a finished result for `input`, but keep an in-flight probe to join. + * A new instance calls this, so creating or editing an instance still + * probes fresh, while instances created together at boot share one probe. + */ + readonly dropFinished: (input: ClaudeProbeInput) => Effect.Effect; } >()("t3/provider/Drivers/ClaudeProbeCache") {} /** @public Service construction is part of the canonical Effect module API. */ export const make = Effect.gen(function* () { - // Inputs whose last probe failed. Compares keys like the cache does. - const failing = MutableHashSet.empty(); + // The latest probe started for each input, and whether it failed. A probe + // that `invalidate` replaced can finish after a newer one, so each probe + // writes only its own record and cannot start a failure streak late. + const latestRuns = MutableHashMap.empty(); const cache = yield* Cache.makeWith( (input: ClaudeProbeInput) => - probeClaudeCapabilities( - input, - mergeProviderInstanceEnvironment(input.environment), - input.cwd, - ).pipe( - Effect.map((probe) => { - if (probe?.usage !== undefined) { - MutableHashSet.remove(failing, input); - return { probe, timeToLive: PROBE_TTL }; - } - const repeat = MutableHashSet.has(failing, input); - // Bounded like the cache. A clear costs each input one early retry. - if (MutableHashSet.size(failing) >= MAX_CACHED_PROBES) MutableHashSet.clear(failing); - MutableHashSet.add(failing, input); - return { probe, timeToLive: repeat ? PROBE_TTL : FIRST_FAILURE_TTL }; - }), - ), + Effect.gen(function* () { + const previous = MutableHashMap.get(latestRuns, input); + const repeat = Option.isSome(previous) && previous.value.failed; + // Bounded like the cache. A clear costs each input one early retry. + if (MutableHashMap.size(latestRuns) >= MAX_CACHED_PROBES) MutableHashMap.clear(latestRuns); + const run = { failed: false }; + MutableHashMap.set(latestRuns, input, run); + const probe = yield* probeClaudeCapabilities( + input, + mergeProviderInstanceEnvironment(input.environment), + input.cwd, + ); + run.failed = probe?.usage === undefined; + return { probe, timeToLive: run.failed && !repeat ? FIRST_FAILURE_TTL : PROBE_TTL }; + }), { capacity: MAX_CACHED_PROBES, timeToLive: (exit) => (Exit.isSuccess(exit) ? exit.value.timeToLive : FIRST_FAILURE_TTL), @@ -85,6 +95,12 @@ export const make = Effect.gen(function* () { return { capabilities: (input) => Cache.get(cache, input).pipe(Effect.map(({ probe }) => probe)), invalidate: (input) => Cache.invalidate(cache, input), + dropFinished: (input) => + Cache.getSuccess(cache, input).pipe( + Effect.flatMap((finished) => + Option.isSome(finished) ? Cache.invalidate(cache, input) : Effect.void, + ), + ), } satisfies ClaudeProbeCache["Service"]; }); diff --git a/apps/server/src/provider/Layers/ProviderInstanceRegistryLive.test.ts b/apps/server/src/provider/Layers/ProviderInstanceRegistryLive.test.ts index 0ed9ad315591..7554db4eaab3 100644 --- a/apps/server/src/provider/Layers/ProviderInstanceRegistryLive.test.ts +++ b/apps/server/src/provider/Layers/ProviderInstanceRegistryLive.test.ts @@ -22,6 +22,7 @@ * binaries. That keeps the assertions focused on registry routing * behaviour rather than the runtime details of each provider. */ +import * as ClaudeSdk from "@anthropic-ai/claude-agent-sdk"; import { describe, expect, it } from "@effect/vitest"; import * as NodeServices from "@effect/platform-node/NodeServices"; import { @@ -32,6 +33,7 @@ import { type OpenCodeSettings, ProviderDriverKind, type ProviderInstanceConfigMap, + type ProviderInstanceEnvironment, ProviderInstanceId, } from "@t3tools/contracts"; import { HostProcessPlatform, isHostWindows } from "@t3tools/shared/hostProcess"; @@ -42,6 +44,7 @@ import * as Layer from "effect/Layer"; import * as Path from "effect/Path"; import * as Stream from "effect/Stream"; import { HttpClient, HttpClientResponse } from "effect/unstable/http"; +import { vi } from "vite-plus/test"; import * as BackgroundPolicy from "../../background/BackgroundPolicy.ts"; import type { BuiltInDriversEnv } from "../builtInDrivers.ts"; @@ -61,6 +64,8 @@ import * as ResetCreditCoordinator from "./resetCreditCoordinator.ts"; import { NoOpProviderEventLoggers, ProviderEventLoggers } from "./ProviderEventLoggers.ts"; import { makeProviderInstanceRegistry } from "./ProviderInstanceRegistryLive.ts"; +vi.mock("@anthropic-ai/claude-agent-sdk", { spy: true }); + const TestHttpClientLive = Layer.succeed( HttpClient.HttpClient, HttpClient.make((request) => @@ -542,6 +547,76 @@ describe("ProviderInstanceRegistryLive — multi-instance codex slice", () => { }), ); + // Instances with the same binary, home, cwd and env vars read the same + // account, so they share one SDK probe. A new or edited instance probes again. + it.live("shares one Claude probe between instances with the same probe input", () => + Effect.gen(function* () { + if (yield* isHostWindows) return; + const fixtures = yield* makeTildeProviderFixtures(); + const probesMayFinish = Promise.withResolvers(); + const query = vi.spyOn(ClaudeSdk, "query").mockImplementation( + () => + ({ + initializationResult: async () => { + await probesMayFinish.promise; + return {}; + }, + usage_EXPERIMENTAL_MAY_CHANGE_DO_NOT_RELY_ON_THIS_API_YET: async () => ({ + rate_limits_available: false, + rate_limits: null, + }), + }) as ReturnType, + ); + // The module mock records every earlier test's real probes too. + query.mockClear(); + yield* Effect.addFinalizer(() => Effect.sync(() => query.mockRestore())); + const probedAccounts = () => + query.mock.calls + .map(([params]) => params.options?.env?.T3_TEST_ACCOUNT ?? "default") + .toSorted(); + const claude = (displayName: string, environment: ProviderInstanceEnvironment = []) => ({ + driver: ProviderDriverKind.make("claudeAgent"), + displayName, + enabled: true, + environment, + config: makeClaudeConfig({ + enabled: true, + binaryPath: fixtures.claudeBinaryPath, + homePath: fixtures.claudeHomePath, + }), + }); + const configMap: ProviderInstanceConfigMap = { + [ProviderInstanceId.make("claude_a")]: claude("A"), + [ProviderInstanceId.make("claude_b")]: claude("B"), + [ProviderInstanceId.make("claude_other")]: claude("Other", [ + { name: "T3_TEST_ACCOUNT", value: "other", sensitive: false }, + ]), + }; + const { registry, mutator } = yield* makeProviderInstanceRegistry({ + drivers: [ClaudeDriver], + configMap, + }); + const refreshAll = registry.listInstances.pipe( + Effect.flatMap((instances) => + Effect.forEach(instances, (instance) => instance.snapshot.refresh, { + concurrency: "unbounded", + }), + ), + ); + // Like boot: every instance exists while the first probes are in flight. + probesMayFinish.resolve(); + yield* refreshAll; + expect(probedAccounts()).toEqual(["default", "other"]); + + yield* mutator.reconcile({ + ...configMap, + [ProviderInstanceId.make("claude_a")]: claude("A renamed"), + }); + yield* refreshAll; + expect(probedAccounts()).toEqual(["default", "default", "other"]); + }).pipe(Effect.provide(testLayer)), + ); + it.live( "shadows instances whose driver is not registered in this build without failing boot", () => From 81204d068c7a5032c2b310a625c310f4428b0776 Mon Sep 17 00:00:00 2001 From: Theo Browne Date: Fri, 25 Sep 2026 22:00:06 -0700 Subject: [PATCH 14/16] refactor(server): shrink the shared Claude probe cache to one Effect Cache The shared cache had grown a failure-streak map, split failure TTLs, a dropFinished step for new instances, and a cached-read rule in resolveUsageLimitsAfterProbe that every provider went through. None of that is needed to stop instances on one home from each running their own probe. ClaudeProbeCache is now one server-wide Effect Cache keyed by the narrowed probe input. Each entry keeps 5 minutes, a failed probe included, like the old per-instance cache. Explicit refresh and the reset-credit re-probe still invalidate the key. Usage limits keep the check time, as on main. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../src/provider/Drivers/ClaudeDriver.ts | 10 +- .../provider/Drivers/ClaudeProbeCache.test.ts | 239 ------------------ .../src/provider/Drivers/ClaudeProbeCache.ts | 91 +------ .../Layers/ClaudeCapabilitiesProbe.test.ts | 1 - .../src/provider/Layers/ClaudeProvider.ts | 15 +- .../ProviderInstanceRegistryLive.test.ts | 10 +- .../provider/Layers/ProviderRegistry.test.ts | 58 ++--- .../src/provider/makeManagedServerProvider.ts | 3 - .../src/provider/providerUsageLimits.test.ts | 37 +-- .../src/provider/providerUsageLimits.ts | 15 -- 10 files changed, 51 insertions(+), 428 deletions(-) delete mode 100644 apps/server/src/provider/Drivers/ClaudeProbeCache.test.ts diff --git a/apps/server/src/provider/Drivers/ClaudeDriver.ts b/apps/server/src/provider/Drivers/ClaudeDriver.ts index 00a3075acdc1..cc6bae0a172b 100644 --- a/apps/server/src/provider/Drivers/ClaudeDriver.ts +++ b/apps/server/src/provider/Drivers/ClaudeDriver.ts @@ -14,6 +14,7 @@ * @module provider/Drivers/ClaudeDriver */ import { ClaudeSettings, ProviderDriverKind } from "@t3tools/contracts"; +import * as Cache from "effect/Cache"; import * as Crypto from "effect/Crypto"; import * as Effect from "effect/Effect"; import * as FileSystem from "effect/FileSystem"; @@ -178,9 +179,6 @@ export const ClaudeDriver: ProviderDriver = { cwd, environment, } satisfies ClaudeProbeCache.ClaudeProbeInput; - // A new or edited instance starts from its own probe, not a sibling's - // finished one. It still joins a probe that is in flight. - yield* probeCache.dropFinished(probeInput); // Start the TTL-gated refresh without delaying provider readiness. The // next check observes a remote manifest after the background fetch lands. @@ -190,7 +188,7 @@ export const ClaudeDriver: ProviderDriver = { Effect.flatMap((manifest) => checkClaudeProviderStatus( effectiveConfig, - () => probeCache.capabilities(probeInput), + () => Cache.get(probeCache, probeInput), processEnv, cwd, resolveClaudeModelCatalog(manifest), @@ -300,7 +298,7 @@ export const ClaudeDriver: ProviderDriver = { Effect.tap((outcome) => Effect.gen(function* () { const before = (yield* snapshot.getSnapshot).usageLimits?.checkedAt; - yield* probeCache.invalidate(probeInput); + yield* Cache.invalidate(probeCache, probeInput); const refreshed = yield* snapshot.refresh; const after = refreshed.usageLimits?.checkedAt; if ( @@ -331,7 +329,7 @@ export const ClaudeDriver: ProviderDriver = { accentColor, enabled, snapshot, - invalidateCaches: probeCache.invalidate(probeInput), + invalidateCaches: Cache.invalidate(probeCache, probeInput), snapshotForCwd, adapter, textGeneration, diff --git a/apps/server/src/provider/Drivers/ClaudeProbeCache.test.ts b/apps/server/src/provider/Drivers/ClaudeProbeCache.test.ts deleted file mode 100644 index d5cdc36217f2..000000000000 --- a/apps/server/src/provider/Drivers/ClaudeProbeCache.test.ts +++ /dev/null @@ -1,239 +0,0 @@ -import * as ClaudeSdk from "@anthropic-ai/claude-agent-sdk"; -import * as NodeServices from "@effect/platform-node/NodeServices"; -import { assert, it } from "@effect/vitest"; -import * as Effect from "effect/Effect"; -import * as Fiber from "effect/Fiber"; -import * as Layer from "effect/Layer"; -import * as TestClock from "effect/testing/TestClock"; -import { vi } from "vite-plus/test"; - -import * as ClaudeProbeCache from "./ClaudeProbeCache.ts"; - -vi.mock("@anthropic-ai/claude-agent-sdk", { spy: true }); - -const testLayer = ClaudeProbeCache.layer.pipe(Layer.provide(NodeServices.layer)); - -// A fresh object per call, so sharing depends on equal inputs, not identity. -const input = ( - homePath: string, - environment: ClaudeProbeCache.ClaudeProbeInput["environment"] = [], -): ClaudeProbeCache.ClaudeProbeInput => ({ - binaryPath: "claude", - homePath, - cwd: "/repo", - environment, -}); - -const failedQuery = () => - ({ - initializationResult: () => Promise.reject(new Error("not logged in")), - }) as ReturnType; - -const failedUsageQuery = () => - ({ - initializationResult: async () => ({}), - usage_EXPERIMENTAL_MAY_CHANGE_DO_NOT_RELY_ON_THIS_API_YET: () => - Promise.reject(new Error("usage unavailable")), - }) as ReturnType; - -// Stands in for the SDK. Each probe reports the Claude home it was started -// with as the account email. -const mockSdk = () => - Effect.gen(function* () { - const query = vi.spyOn(ClaudeSdk, "query").mockImplementation( - ({ options }) => - ({ - initializationResult: async () => ({ - account: { email: options?.env?.CLAUDE_CONFIG_DIR ?? "" }, - commands: [{ name: "review", description: "Review changes", argumentHint: "" }], - }), - usage_EXPERIMENTAL_MAY_CHANGE_DO_NOT_RELY_ON_THIS_API_YET: async () => ({ - rate_limits_available: false, - rate_limits: null, - }), - }) as ReturnType, - ); - yield* Effect.addFinalizer(() => Effect.sync(() => query.mockRestore())); - return query; - }); - -// Keeps the next probe in flight until `release`, then fails it. -const holdNextProbe = (query: Effect.Success>) => { - const started = Promise.withResolvers(); - const release = Promise.withResolvers(); - query.mockImplementationOnce(() => { - started.resolve(); - return { - initializationResult: () => - release.promise.then(() => Promise.reject(new Error("probe timed out"))), - } as ReturnType; - }); - return { started: started.promise, release: () => release.resolve() }; -}; - -it.effect("instances with the same probe input share one probe", () => - Effect.gen(function* () { - const query = yield* mockSdk(); - const cache = yield* ClaudeProbeCache.ClaudeProbeCache; - - const [first, second] = yield* Effect.all( - [cache.capabilities(input("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/homes/work")), cache.capabilities(input("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/homes/work"))], - { concurrency: "unbounded" }, - ); - yield* TestClock.adjust("2 minutes"); - const later = yield* cache.capabilities(input("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/homes/work")); - - assert.equal(query.mock.calls.length, 1); - assert.match(first?.email ?? "", /work$/); - assert.deepEqual(second, first); - // A later reader sees the probe's own time, not its read time. - assert.deepEqual(later, first); - assert.equal(later?.checkedAt, "1970-01-01T00:00:00.000Z"); - }).pipe(Effect.scoped, Effect.provide(testLayer)), -); - -it.effect("different homes or instance env vars run separate probes", () => - Effect.gen(function* () { - const query = yield* mockSdk(); - const cache = yield* ClaudeProbeCache.ClaudeProbeCache; - - const work = yield* cache.capabilities(input("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/homes/work")); - const personal = yield* cache.capabilities(input("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/homes/personal")); - yield* cache.capabilities( - input("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/homes/work", [{ name: "ANTHROPIC_API_KEY", value: "sk-test", sensitive: true }]), - ); - - assert.equal(query.mock.calls.length, 3); - assert.match(work?.email ?? "", /work$/); - assert.match(personal?.email ?? "", /personal$/); - }).pipe(Effect.scoped, Effect.provide(testLayer)), -); - -// Both point at ~/.claude, but an explicit CLAUDE_CONFIG_DIR is a separate -// login to the CLI (its own keychain entry and .claude.json). -it.effect("an empty home and an explicit ~/.claude run separate probes", () => - Effect.gen(function* () { - const query = yield* mockSdk(); - const cache = yield* ClaudeProbeCache.ClaudeProbeCache; - - yield* cache.capabilities(input("")); - const explicit = yield* cache.capabilities(input("~/.claude")); - - assert.equal(query.mock.calls.length, 2); - assert.match(explicit?.email ?? "", /\.claude$/); - }).pipe(Effect.scoped, Effect.provide(testLayer)), -); - -it.effect("retries a first failure after 30 seconds and a repeat failure after 5 minutes", () => - Effect.gen(function* () { - const query = yield* mockSdk(); - query.mockImplementationOnce(failedQuery).mockImplementationOnce(failedUsageQuery); - const cache = yield* ClaudeProbeCache.ClaudeProbeCache; - const read = cache.capabilities(input("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/homes/work")); - - assert.equal(yield* read, undefined); - yield* TestClock.adjust("29 seconds"); - yield* read; - assert.equal(query.mock.calls.length, 1); - - // A failed usage read is a failure too, and this one repeats. - yield* TestClock.adjust("1 second"); - assert.isUndefined((yield* read)?.usage); - assert.equal(query.mock.calls.length, 2); - yield* TestClock.adjust("4 minutes"); - yield* read; - assert.equal(query.mock.calls.length, 2); - - yield* TestClock.adjust("1 minute"); - assert.isDefined((yield* read)?.usage); - assert.equal(query.mock.calls.length, 3); - - // A success ends the streak, so the next failure retries soon again. - query.mockImplementationOnce(failedQuery); - yield* TestClock.adjust("5 minutes"); - assert.equal(yield* read, undefined); - yield* TestClock.adjust("30 seconds"); - assert.isDefined(yield* read); - assert.equal(query.mock.calls.length, 5); - }).pipe(Effect.scoped, Effect.provide(testLayer)), -); - -// `invalidate` can replace a probe that is still running. When that old probe -// fails after the new one succeeded, the next failure is still a first one. -it.effect("a replaced probe that fails late does not start a failure streak", () => - Effect.gen(function* () { - const query = yield* mockSdk(); - const held = holdNextProbe(query); - const cache = yield* ClaudeProbeCache.ClaudeProbeCache; - const read = cache.capabilities(input("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/homes/work")); - - const replaced = yield* Effect.forkChild(read); - yield* Effect.promise(() => held.started); - yield* cache.invalidate(input("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/homes/work")); - assert.isDefined((yield* read)?.usage); - held.release(); - assert.equal(yield* Fiber.join(replaced), undefined); - - query.mockImplementationOnce(failedQuery); - yield* TestClock.adjust("5 minutes"); - assert.equal(yield* read, undefined); - yield* TestClock.adjust("30 seconds"); - assert.isDefined(yield* read); - assert.equal(query.mock.calls.length, 4); - }).pipe(Effect.scoped, Effect.provide(testLayer)), -); - -it.effect("dropFinished re-probes a finished input but joins a probe in flight", () => - Effect.gen(function* () { - const query = yield* mockSdk(); - const cache = yield* ClaudeProbeCache.ClaudeProbeCache; - const read = cache.capabilities(input("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/homes/work")); - - yield* read; - yield* cache.dropFinished(input("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/homes/work")); - yield* read; - assert.equal(query.mock.calls.length, 2); - - const held = holdNextProbe(query); - yield* cache.invalidate(input("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/homes/work")); - const first = yield* Effect.forkChild(read); - yield* Effect.promise(() => held.started); - yield* cache.dropFinished(input("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/homes/work")); - const second = yield* Effect.forkChild(read); - held.release(); - yield* Fiber.join(first); - yield* Fiber.join(second); - assert.equal(query.mock.calls.length, 3); - }).pipe(Effect.scoped, Effect.provide(testLayer)), -); - -it.effect("invalidate re-probes only that input", () => - Effect.gen(function* () { - const query = yield* mockSdk(); - const cache = yield* ClaudeProbeCache.ClaudeProbeCache; - - yield* cache.capabilities(input("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/homes/work")); - yield* cache.capabilities(input("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/homes/personal")); - yield* cache.invalidate(input("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/homes/work")); - yield* cache.capabilities(input("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/homes/work")); - yield* cache.capabilities(input("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/homes/personal")); - - assert.equal(query.mock.calls.length, 3); - }).pipe(Effect.scoped, Effect.provide(testLayer)), -); - -it.effect("keeps every input cached across refreshes when there are many instances", () => - Effect.gen(function* () { - const query = yield* mockSdk(); - const cache = yield* ClaudeProbeCache.ClaudeProbeCache; - const homes = Array.from({ length: 100 }, (_, index) => `/homes/${index}`); - const refresh = Effect.forEach(homes, (home) => cache.capabilities(input(home)), { - concurrency: "unbounded", - }); - - yield* refresh; - yield* refresh; - - assert.equal(query.mock.calls.length, homes.length); - }).pipe(Effect.scoped, Effect.provide(testLayer)), -); diff --git a/apps/server/src/provider/Drivers/ClaudeProbeCache.ts b/apps/server/src/provider/Drivers/ClaudeProbeCache.ts index 1c9c024b1b99..c68e596c545e 100644 --- a/apps/server/src/provider/Drivers/ClaudeProbeCache.ts +++ b/apps/server/src/provider/Drivers/ClaudeProbeCache.ts @@ -1,11 +1,9 @@ /** * One server-wide cache for the Claude capabilities probe. Claude instances - * whose probe inputs match (same binary, home, cwd and instance env vars) - * read the same account, so they share one cached result instead of each - * starting its own SDK session. Concurrent reads of one input join one probe, - * so a full refresh runs one probe per distinct input. There is no global - * limit on purpose: a probe takes seconds, so a queue would delay the status - * of every input past the limit. + * with the same probe input read the same account, so they share one cached + * result, and concurrent reads of one input join one SDK probe. Each entry + * keeps 5 minutes, a failed probe (`undefined`) included, like the old + * per-instance cache. * * @module provider/Drivers/ClaudeProbeCache */ @@ -13,95 +11,34 @@ import type { ClaudeSettings, ProviderInstanceEnvironment } from "@t3tools/contr import * as Cache from "effect/Cache"; import * as Context from "effect/Context"; import * as Duration from "effect/Duration"; -import * as Effect from "effect/Effect"; -import * as Exit from "effect/Exit"; import * as Layer from "effect/Layer"; -import * as MutableHashMap from "effect/MutableHashMap"; -import * as Option from "effect/Option"; import { type ClaudeCapabilitiesProbe, probeClaudeCapabilities } from "../Layers/ClaudeProvider.ts"; import { mergeProviderInstanceEnvironment } from "../ProviderInstanceEnvironment.ts"; /** * Every instance input the probe reads, and also the cache key. The lookup - * gets only this value, so the probe cannot depend on an input the key leaves - * out. Keys compare structurally. + * gets only this value, so the probe cannot read an input the key leaves out. + * Keys compare structurally. */ export type ClaudeProbeInput = Pick & { readonly cwd: string; readonly environment: ProviderInstanceEnvironment; }; -const PROBE_TTL = Duration.minutes(5); -// A failed probe or usage read leaves every instance with that input unverified -// or without limits, so the first failure retries soon. A repeat failure -// (signed out, no usage endpoint) waits the full TTL, like a success. -const FIRST_FAILURE_TTL = Duration.seconds(30); -// Keep this far above any real instance count. The cache evicts the least -// recently used key, and every refresh reads the keys in the same order, so a -// cap below the live key count makes each refresh re-probe every instance. -// Stale keys only come from instance edits, and entries are small. -const MAX_CACHED_PROBES = 1024; - export class ClaudeProbeCache extends Context.Service< ClaudeProbeCache, - { - /** Cached probe result for `input`, or `undefined` when the probe failed. */ - readonly capabilities: ( - input: ClaudeProbeInput, - ) => Effect.Effect; - /** - * Drop the result for `input`, even one still in flight, so the next read - * starts a new probe. - */ - readonly invalidate: (input: ClaudeProbeInput) => Effect.Effect; - /** - * Drop a finished result for `input`, but keep an in-flight probe to join. - * A new instance calls this, so creating or editing an instance still - * probes fresh, while instances created together at boot share one probe. - */ - readonly dropFinished: (input: ClaudeProbeInput) => Effect.Effect; - } + Cache.Cache >()("t3/provider/Drivers/ClaudeProbeCache") {} /** @public Service construction is part of the canonical Effect module API. */ -export const make = Effect.gen(function* () { - // The latest probe started for each input, and whether it failed. A probe - // that `invalidate` replaced can finish after a newer one, so each probe - // writes only its own record and cannot start a failure streak late. - const latestRuns = MutableHashMap.empty(); - const cache = yield* Cache.makeWith( - (input: ClaudeProbeInput) => - Effect.gen(function* () { - const previous = MutableHashMap.get(latestRuns, input); - const repeat = Option.isSome(previous) && previous.value.failed; - // Bounded like the cache. A clear costs each input one early retry. - if (MutableHashMap.size(latestRuns) >= MAX_CACHED_PROBES) MutableHashMap.clear(latestRuns); - const run = { failed: false }; - MutableHashMap.set(latestRuns, input, run); - const probe = yield* probeClaudeCapabilities( - input, - mergeProviderInstanceEnvironment(input.environment), - input.cwd, - ); - run.failed = probe?.usage === undefined; - return { probe, timeToLive: run.failed && !repeat ? FIRST_FAILURE_TTL : PROBE_TTL }; - }), - { - capacity: MAX_CACHED_PROBES, - timeToLive: (exit) => (Exit.isSuccess(exit) ? exit.value.timeToLive : FIRST_FAILURE_TTL), - }, - ); - return { - capabilities: (input) => Cache.get(cache, input).pipe(Effect.map(({ probe }) => probe)), - invalidate: (input) => Cache.invalidate(cache, input), - dropFinished: (input) => - Cache.getSuccess(cache, input).pipe( - Effect.flatMap((finished) => - Option.isSome(finished) ? Cache.invalidate(cache, input) : Effect.void, - ), - ), - } satisfies ClaudeProbeCache["Service"]; +export const make = Cache.make({ + // Far above any real input count. Each refresh reads keys in the same + // order, so a cap below the live key count would re-probe every key. + capacity: 256, + timeToLive: Duration.minutes(5), + lookup: (input: ClaudeProbeInput) => + probeClaudeCapabilities(input, mergeProviderInstanceEnvironment(input.environment), input.cwd), }); export const layer = Layer.effect(ClaudeProbeCache, make); diff --git a/apps/server/src/provider/Layers/ClaudeCapabilitiesProbe.test.ts b/apps/server/src/provider/Layers/ClaudeCapabilitiesProbe.test.ts index cf03625be768..232b8cc02d00 100644 --- a/apps/server/src/provider/Layers/ClaudeCapabilitiesProbe.test.ts +++ b/apps/server/src/provider/Layers/ClaudeCapabilitiesProbe.test.ts @@ -161,7 +161,6 @@ it.layer(NodeServices.layer)("Claude capability probe SDK boundary", (it) => { rate_limits_available: true, rate_limits: { five_hour: { utilization: 12, resets_at: "2026-07-18T14:39:00Z" } }, }, - checkedAt: "1970-01-01T00:00:00.000Z", }); // @effect-diagnostics-next-line preferSchemaOverJson:off diff --git a/apps/server/src/provider/Layers/ClaudeProvider.ts b/apps/server/src/provider/Layers/ClaudeProvider.ts index 3aa605bca94f..b296a9fd1cda 100644 --- a/apps/server/src/provider/Layers/ClaudeProvider.ts +++ b/apps/server/src/provider/Layers/ClaudeProvider.ts @@ -242,11 +242,6 @@ export type ClaudeCapabilitiesProbe = { * otherwise successful response mean the account has none (API key). */ readonly usage?: Pick; - /** - * When the probe read usage. Instances share cached probes, so usage - * limits carry this time, not the time of the status check that read it. - */ - readonly checkedAt: string; }; function parseClaudeInitializationCommands( @@ -378,7 +373,6 @@ const probeClaudeCapabilities = ( rate_limits: usageResult.success.rate_limits, } : undefined; - const checkedAt = DateTime.formatIso(yield* DateTime.now); const account = init.account as | { readonly email?: string; @@ -394,7 +388,6 @@ const probeClaudeCapabilities = ( apiProvider: account?.apiProvider, slashCommands: parseClaudeInitializationCommands(init.commands), ...(usage ? { usage } : {}), - checkedAt, } satisfies ClaudeCapabilitiesProbe; }), ), @@ -570,16 +563,14 @@ export const checkClaudeProviderStatus = Effect.fn("checkClaudeProviderStatus")( subscriptionType: capabilities.subscriptionType, authMethod: capabilities.tokenSource, }) ?? apiProviderAuthMetadata(capabilities.apiProvider); - const usageCheckedAt = capabilities.checkedAt; const usageLimits = !capabilities.usage - ? makeUnavailableUsageLimits({ checkedAt: usageCheckedAt, reason: "probeFailed" }) + ? makeUnavailableUsageLimits({ checkedAt, reason: "probeFailed" }) : scopedLimitNames ? yield* recordClaudeUsageResponse(scopedLimitNames, { response: capabilities.usage, - checkedAt: usageCheckedAt, + checkedAt, }) - : claudeUsageResponseToLimits({ response: capabilities.usage, checkedAt: usageCheckedAt }) - .limits; + : claudeUsageResponseToLimits({ response: capabilities.usage, checkedAt }).limits; const resetCredits = resolveResetCredits && capabilities.subscriptionType && diff --git a/apps/server/src/provider/Layers/ProviderInstanceRegistryLive.test.ts b/apps/server/src/provider/Layers/ProviderInstanceRegistryLive.test.ts index 7554db4eaab3..9588853011f4 100644 --- a/apps/server/src/provider/Layers/ProviderInstanceRegistryLive.test.ts +++ b/apps/server/src/provider/Layers/ProviderInstanceRegistryLive.test.ts @@ -548,7 +548,7 @@ describe("ProviderInstanceRegistryLive — multi-instance codex slice", () => { ); // Instances with the same binary, home, cwd and env vars read the same - // account, so they share one SDK probe. A new or edited instance probes again. + // account, so they share one SDK probe. An explicit refresh re-probes it. it.live("shares one Claude probe between instances with the same probe input", () => Effect.gen(function* () { if (yield* isHostWindows) return; @@ -592,7 +592,7 @@ describe("ProviderInstanceRegistryLive — multi-instance codex slice", () => { { name: "T3_TEST_ACCOUNT", value: "other", sensitive: false }, ]), }; - const { registry, mutator } = yield* makeProviderInstanceRegistry({ + const { registry } = yield* makeProviderInstanceRegistry({ drivers: [ClaudeDriver], configMap, }); @@ -608,10 +608,8 @@ describe("ProviderInstanceRegistryLive — multi-instance codex slice", () => { yield* refreshAll; expect(probedAccounts()).toEqual(["default", "other"]); - yield* mutator.reconcile({ - ...configMap, - [ProviderInstanceId.make("claude_a")]: claude("A renamed"), - }); + const claudeA = yield* registry.getInstance(ProviderInstanceId.make("claude_a")); + yield* claudeA?.invalidateCaches ?? Effect.void; yield* refreshAll; expect(probedAccounts()).toEqual(["default", "default", "other"]); }).pipe(Effect.provide(testLayer)), diff --git a/apps/server/src/provider/Layers/ProviderRegistry.test.ts b/apps/server/src/provider/Layers/ProviderRegistry.test.ts index f93458313a5a..abc2cd42fdb7 100644 --- a/apps/server/src/provider/Layers/ProviderRegistry.test.ts +++ b/apps/server/src/provider/Layers/ProviderRegistry.test.ts @@ -23,6 +23,7 @@ import { ProviderInstanceId, ServerSettings, type ServerProvider, + type ServerProviderSlashCommand, type ServerSettings as ContractServerSettings, } from "@t3tools/contracts"; import * as PlatformError from "effect/PlatformError"; @@ -33,7 +34,7 @@ import { createModelCapabilities } from "@t3tools/shared/model"; import { applyServerSettingsPatch } from "@t3tools/shared/serverSettings"; import { checkCodexProviderStatus, type CodexAppServerProviderSnapshot } from "./CodexProvider.ts"; -import { type ClaudeCapabilitiesProbe, checkClaudeProviderStatus } from "./ClaudeProvider.ts"; +import { checkClaudeProviderStatus } from "./ClaudeProvider.ts"; import * as BackgroundPolicy from "../../background/BackgroundPolicy.ts"; import { AntigravityInstallation } from "../AntigravityInstallation.ts"; import * as ModelManifest from "../ModelManifest.ts"; @@ -145,7 +146,15 @@ function booleanDescriptor(id: string, label: string) { }; } -function claudeCapabilities(overrides: Partial = {}) { +type TestClaudeCapabilities = { + readonly email: string | undefined; + readonly subscriptionType: string | undefined; + readonly tokenSource: string | undefined; + readonly apiProvider: string | undefined; + readonly slashCommands: ReadonlyArray; +}; + +function claudeCapabilities(overrides: Partial = {}) { return () => Effect.succeed({ email: undefined, @@ -153,13 +162,12 @@ function claudeCapabilities(overrides: Partial = {}) { tokenSource: undefined, apiProvider: undefined, slashCommands: [], - checkedAt: "2026-09-25T12:00:00.000Z", ...overrides, }); } const noClaudeCapabilities = () => - Effect.sync(() => undefined as ClaudeCapabilitiesProbe | undefined); + Effect.sync(() => undefined as TestClaudeCapabilities | undefined); function mockHandle(result: { stdout: string; stderr: string; code: number }) { return ChildProcessSpawner.makeHandle({ @@ -2750,13 +2758,19 @@ it.layer(Layer.mergeAll(NodeServices.layer, ServerSettingsModule.layerTest(), Te it.effect("reads banked resets only for subscription logins", () => Effect.gen(function* () { - const check = (overrides: Partial) => + const check = (overrides: Partial) => checkClaudeProviderStatus( defaultClaudeSettings, - claudeCapabilities({ - usage: { rate_limits_available: true, rate_limits: {} }, - ...overrides, - }), + () => + Effect.succeed({ + email: undefined, + subscriptionType: undefined, + tokenSource: undefined, + apiProvider: undefined, + slashCommands: [], + usage: { rate_limits_available: true, rate_limits: {} }, + ...overrides, + }), undefined, undefined, undefined, @@ -2778,32 +2792,6 @@ it.layer(Layer.mergeAll(NodeServices.layer, ServerSettingsModule.layerTest(), Te ), ); - // Instances share cached probes, so the usage can be older than the check. - it.effect("dates usage limits by the probe, not the status check", () => - Effect.gen(function* () { - const probedAt = "2026-09-25T12:00:00.000Z"; - const check = (usage: ClaudeCapabilitiesProbe["usage"]) => - checkClaudeProviderStatus( - defaultClaudeSettings, - claudeCapabilities({ checkedAt: probedAt, ...(usage ? { usage } : {}) }), - ); - const read = yield* check({ rate_limits_available: true, rate_limits: {} }); - const failed = yield* check(undefined); - assert.notStrictEqual(read.checkedAt, probedAt); - assert.strictEqual(read.usageLimits?.checkedAt, probedAt); - assert.strictEqual(failed.usageLimits?.unavailable?.reason, "probeFailed"); - assert.strictEqual(failed.usageLimits?.checkedAt, probedAt); - }).pipe( - Effect.provide( - mockSpawnerLayer((args) => { - const joined = args.join(" "); - if (joined === "--version") return { stdout: "1.0.0\n", stderr: "", code: 0 }; - throw new Error(`Unexpected args: ${joined}`); - }), - ), - ), - ); - it.effect("does not duplicate Claude in full subscription labels", () => Effect.gen(function* () { const status = yield* checkClaudeProviderStatus( diff --git a/apps/server/src/provider/makeManagedServerProvider.ts b/apps/server/src/provider/makeManagedServerProvider.ts index 1c3189d1f648..119c577474c8 100644 --- a/apps/server/src/provider/makeManagedServerProvider.ts +++ b/apps/server/src/provider/makeManagedServerProvider.ts @@ -4,7 +4,6 @@ import { ServerSettingsError, } from "@t3tools/contracts"; import { resolveServerBackgroundActivitySettings } from "@t3tools/shared/backgroundActivitySettings"; -import * as Clock from "effect/Clock"; import * as Duration from "effect/Duration"; import * as Effect from "effect/Effect"; import * as Equal from "effect/Equal"; @@ -151,7 +150,6 @@ export const makeManagedServerProvider = Effect.fn("makeManagedServerProvider")( return state.snapshot; } - const checkStartedAt = yield* Clock.currentTimeMillis; const probedSnapshot = yield* input.checkProvider; const { snapshot: nextSnapshot, generation: nextGeneration } = yield* Ref.modify( snapshotStateRef, @@ -164,7 +162,6 @@ export const makeManagedServerProvider = Effect.fn("makeManagedServerProvider")( resolveUsageLimitsAfterProbe({ published: state.snapshot.usageLimits, probed: probedSnapshot.usageLimits, - checkStartedAt, }), ); return [ diff --git a/apps/server/src/provider/providerUsageLimits.test.ts b/apps/server/src/provider/providerUsageLimits.test.ts index 15a06e6e062f..6e288ddd3a36 100644 --- a/apps/server/src/provider/providerUsageLimits.test.ts +++ b/apps/server/src/provider/providerUsageLimits.test.ts @@ -1,4 +1,3 @@ -import type { ServerProviderUsageLimits } from "@t3tools/contracts"; import { describe, expect, it } from "vite-plus/test"; import { applyUsageLimitsUpdate, resolveUsageLimitsAfterProbe } from "./providerUsageLimits.ts"; @@ -80,41 +79,11 @@ describe("applyUsageLimitsUpdate", () => { }); describe("resolveUsageLimitsAfterProbe", () => { - const checkStartedAt = Date.parse(checkedAt); - it("keeps the last good windows through a failed probe but not an unsupported one", () => { const failed = { checkedAt, windows: [], unavailable: { reason: "probeFailed" as const } }; const unsupported = { checkedAt, windows: [], unavailable: { reason: "unsupported" as const } }; - expect(resolveUsageLimitsAfterProbe({ published, probed: failed, checkStartedAt })).toBe( - published, - ); - expect(resolveUsageLimitsAfterProbe({ published, probed: unsupported, checkStartedAt })).toBe( - unsupported, - ); - expect( - resolveUsageLimitsAfterProbe({ published: undefined, probed: failed, checkStartedAt }), - ).toBe(failed); - }); - - it("keeps a turn update over an older cached read, not over a fresh probe", () => { - const startedAt = Date.parse("2026-09-03T12:05:00.000Z"); - const turnUpdate = { checkedAt: "2026-09-03T12:04:00.000Z", windows: [session] }; - // Claude instances share cached probes, so a read can predate the check. - const cached = { checkedAt: "2026-09-03T12:01:00.000Z", windows: [session, weekly] }; - // Codex dates its read at the check start, and a turn can update mid-probe. - const fresh = { checkedAt: "2026-09-03T12:05:00.000Z", windows: [session, weekly] }; - const midProbeUpdate = { checkedAt: "2026-09-03T12:05:02.000Z", windows: [session] }; - const resolve = (current: ServerProviderUsageLimits, probed: ServerProviderUsageLimits) => - resolveUsageLimitsAfterProbe({ published: current, probed, checkStartedAt: startedAt }); - expect(resolve(turnUpdate, cached)).toEqual(turnUpdate); - // The check read credits itself, so a spent credit does not survive. - const spent = { availableCount: 1, nextCreditId: "spent" }; - const current = { availableCount: 0 }; - expect( - resolve({ ...turnUpdate, resetCredits: spent }, { ...cached, resetCredits: current }), - ).toEqual({ ...turnUpdate, resetCredits: current }); - expect(resolve({ ...turnUpdate, resetCredits: spent }, cached)).toEqual(turnUpdate); - expect(resolve(midProbeUpdate, fresh)).toBe(fresh); - expect(resolve(published, cached)).toBe(cached); + expect(resolveUsageLimitsAfterProbe({ published, probed: failed })).toBe(published); + expect(resolveUsageLimitsAfterProbe({ published, probed: unsupported })).toBe(unsupported); + expect(resolveUsageLimitsAfterProbe({ published: undefined, probed: failed })).toBe(failed); }); }); diff --git a/apps/server/src/provider/providerUsageLimits.ts b/apps/server/src/provider/providerUsageLimits.ts index a37b0a6839ed..ea8d0d1d029f 100644 --- a/apps/server/src/provider/providerUsageLimits.ts +++ b/apps/server/src/provider/providerUsageLimits.ts @@ -121,29 +121,14 @@ function usageWindowEquals(a: ServerProviderUsageWindow, b: ServerProviderUsageW * per-window epoch bookkeeping needed to reconcile the two was more code * than the sub-second regression it prevented. The next runtime event * corrects it. - * - * Limits stamped before the check started came from a cache (Claude - * instances share one probe), so they can be minutes old. Those do not - * replace newer published windows, such as a turn's update. Reset credits - * still come from `probed`, because each check reads them itself. */ export function resolveUsageLimitsAfterProbe(input: { readonly published: ServerProviderUsageLimits | undefined; readonly probed: ServerProviderUsageLimits | undefined; - /** Epoch millis when the status check that produced `probed` started. */ - readonly checkStartedAt: number; }): ServerProviderUsageLimits | undefined { const { published, probed } = input; if (probed?.unavailable?.reason === "probeFailed" && published && !published.unavailable) { return published; } - if ( - published && - probed && - Date.parse(probed.checkedAt) < Math.min(input.checkStartedAt, Date.parse(published.checkedAt)) - ) { - const { resetCredits: _published, ...windows } = published; - return probed.resetCredits ? { ...windows, resetCredits: probed.resetCredits } : windows; - } return probed; } From cc29c1fdcc3d0b754db0d34f1e0ae736541dfd16 Mon Sep 17 00:00:00 2001 From: Theo Browne Date: Fri, 25 Sep 2026 22:23:16 -0700 Subject: [PATCH 15/16] fix(server): a new or rebuilt Claude instance probes fresh Before the shared cache, a config edit rebuilt the instance with an empty cache, so it probed at once. Now create drops the shared entry for its input when that entry already finished. An in-flight probe is joined, so instances that start together at boot still run one probe per input. Tests: a rebuild after the probe finished runs a fresh probe, an instance added while a probe is in flight joins it, and a reset re-probe reaches siblings on the same input. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../src/provider/Drivers/ClaudeDriver.ts | 7 ++ .../ProviderInstanceRegistryLive.test.ts | 104 +++++++++++------- 2 files changed, 74 insertions(+), 37 deletions(-) diff --git a/apps/server/src/provider/Drivers/ClaudeDriver.ts b/apps/server/src/provider/Drivers/ClaudeDriver.ts index cc6bae0a172b..a5405c0396be 100644 --- a/apps/server/src/provider/Drivers/ClaudeDriver.ts +++ b/apps/server/src/provider/Drivers/ClaudeDriver.ts @@ -18,6 +18,7 @@ import * as Cache from "effect/Cache"; import * as Crypto from "effect/Crypto"; import * as Effect from "effect/Effect"; import * as FileSystem from "effect/FileSystem"; +import * as Option from "effect/Option"; import * as Path from "effect/Path"; import * as Schema from "effect/Schema"; import { HttpClient } from "effect/unstable/http"; @@ -179,6 +180,12 @@ export const ClaudeDriver: ProviderDriver = { cwd, environment, } satisfies ClaudeProbeCache.ClaudeProbeInput; + // A new or rebuilt instance probes fresh, so a config edit never shows + // an older result. An in-flight probe is joined instead, so instances + // that start together at boot still run one probe. + if (Option.isSome(yield* Cache.getSuccess(probeCache, probeInput))) { + yield* Cache.invalidate(probeCache, probeInput); + } // Start the TTL-gated refresh without delaying provider readiness. The // next check observes a remote manifest after the background fetch lands. diff --git a/apps/server/src/provider/Layers/ProviderInstanceRegistryLive.test.ts b/apps/server/src/provider/Layers/ProviderInstanceRegistryLive.test.ts index 9588853011f4..e284545253c5 100644 --- a/apps/server/src/provider/Layers/ProviderInstanceRegistryLive.test.ts +++ b/apps/server/src/provider/Layers/ProviderInstanceRegistryLive.test.ts @@ -493,33 +493,40 @@ describe("ProviderInstanceRegistryLive — multi-instance codex slice", () => { }), ); const instanceId = ProviderInstanceId.make("claude_reset"); + // A sibling with the same probe input shares the cached probe. + const siblingId = ProviderInstanceId.make("claude_reset_sibling"); + const entry = { + driver: ProviderDriverKind.make("claudeAgent"), + enabled: true, + environment: [ + { name: "T3_CLAUDE_RESET_MARKER", value: marker, sensitive: false }, + ...(claim.usageFailsAfterClaim + ? [{ name: "T3_CLAUDE_USAGE_FAILS_AFTER_CLAIM", value: "1", sensitive: false }] + : []), + ], + config: makeClaudeConfig({ + enabled: true, + binaryPath: fixtures.claudeBinaryPath, + homePath: fixtures.claudeHomePath, + }), + }; const { registry } = yield* makeProviderInstanceRegistry({ drivers: [ClaudeDriver], - configMap: { - [instanceId]: { - driver: ProviderDriverKind.make("claudeAgent"), - enabled: true, - environment: [ - { name: "T3_CLAUDE_RESET_MARKER", value: marker, sensitive: false }, - ...(claim.usageFailsAfterClaim - ? [{ name: "T3_CLAUDE_USAGE_FAILS_AFTER_CLAIM", value: "1", sensitive: false }] - : []), - ], - config: makeClaudeConfig({ - enabled: true, - binaryPath: fixtures.claudeBinaryPath, - homePath: fixtures.claudeHomePath, - }), - }, - }, + configMap: { [instanceId]: entry, [siblingId]: entry }, }).pipe(Effect.provideService(HttpClient.HttpClient, client)); const instance = yield* registry.getInstance(instanceId); + const sibling = yield* registry.getInstance(siblingId); expect(instance).toBeDefined(); + expect(sibling).toBeDefined(); const before = yield* instance!.snapshot.refresh; expect(before.usageLimits?.windows[0]?.usedPercent).toBe(100); expect(before.usageLimits?.resetCredits?.nextCreditId).toBe("grant_a"); const outcome = yield* instance!.consumeResetCredit!().pipe(Effect.result); - return { outcome, after: yield* instance!.snapshot.getSnapshot }; + return { + outcome, + after: yield* instance!.snapshot.getSnapshot, + siblingAfter: yield* sibling!.snapshot.refresh, + }; }).pipe( // macOS logins live in the Keychain, where resets are never read. Effect.provideService(HostProcessPlatform, "linux"), @@ -537,6 +544,16 @@ describe("ProviderInstanceRegistryLive — multi-instance codex slice", () => { }), ); + it.live("a Claude reset re-probes the usage that siblings share", () => + Effect.gen(function* () { + const { siblingAfter } = yield* redeemClaudeReset({ + result: "reset", + usageFailsAfterClaim: false, + }); + expect(siblingAfter.usageLimits?.windows[0]?.usedPercent).toBe(0); + }), + ); + it.live("reports Claude's answer when a claim changed nothing and the re-probe fails", () => Effect.gen(function* () { const { outcome } = yield* redeemClaudeReset({ @@ -548,25 +565,27 @@ describe("ProviderInstanceRegistryLive — multi-instance codex slice", () => { ); // Instances with the same binary, home, cwd and env vars read the same - // account, so they share one SDK probe. An explicit refresh re-probes it. + // account, so they share one SDK probe. An explicit refresh or a rebuild + // re-probes it. it.live("shares one Claude probe between instances with the same probe input", () => Effect.gen(function* () { if (yield* isHostWindows) return; const fixtures = yield* makeTildeProviderFixtures(); + const firstProbeStarted = Promise.withResolvers(); const probesMayFinish = Promise.withResolvers(); - const query = vi.spyOn(ClaudeSdk, "query").mockImplementation( - () => - ({ - initializationResult: async () => { - await probesMayFinish.promise; - return {}; - }, - usage_EXPERIMENTAL_MAY_CHANGE_DO_NOT_RELY_ON_THIS_API_YET: async () => ({ - rate_limits_available: false, - rate_limits: null, - }), - }) as ReturnType, - ); + const query = vi.spyOn(ClaudeSdk, "query").mockImplementation(() => { + firstProbeStarted.resolve(); + return { + initializationResult: async () => { + await probesMayFinish.promise; + return {}; + }, + usage_EXPERIMENTAL_MAY_CHANGE_DO_NOT_RELY_ON_THIS_API_YET: async () => ({ + rate_limits_available: false, + rate_limits: null, + }), + } as ReturnType; + }); // The module mock records every earlier test's real probes too. query.mockClear(); yield* Effect.addFinalizer(() => Effect.sync(() => query.mockRestore())); @@ -585,16 +604,18 @@ describe("ProviderInstanceRegistryLive — multi-instance codex slice", () => { homePath: fixtures.claudeHomePath, }), }); + const claudeAId = ProviderInstanceId.make("claude_a"); + const bootConfigMap: ProviderInstanceConfigMap = { [claudeAId]: claude("A") }; const configMap: ProviderInstanceConfigMap = { - [ProviderInstanceId.make("claude_a")]: claude("A"), + ...bootConfigMap, [ProviderInstanceId.make("claude_b")]: claude("B"), [ProviderInstanceId.make("claude_other")]: claude("Other", [ { name: "T3_TEST_ACCOUNT", value: "other", sensitive: false }, ]), }; - const { registry } = yield* makeProviderInstanceRegistry({ + const { registry, mutator } = yield* makeProviderInstanceRegistry({ drivers: [ClaudeDriver], - configMap, + configMap: bootConfigMap, }); const refreshAll = registry.listInstances.pipe( Effect.flatMap((instances) => @@ -603,15 +624,24 @@ describe("ProviderInstanceRegistryLive — multi-instance codex slice", () => { }), ), ); - // Like boot: every instance exists while the first probes are in flight. + // Like boot: the other instances start while claude_a's probe is in flight. + yield* Effect.promise(() => firstProbeStarted.promise); + yield* mutator.reconcile(configMap); probesMayFinish.resolve(); yield* refreshAll; expect(probedAccounts()).toEqual(["default", "other"]); - const claudeA = yield* registry.getInstance(ProviderInstanceId.make("claude_a")); + const claudeA = yield* registry.getInstance(claudeAId); yield* claudeA?.invalidateCaches ?? Effect.void; yield* refreshAll; expect(probedAccounts()).toEqual(["default", "default", "other"]); + + // A config edit rebuilds claude_a after its probe finished. + yield* mutator.reconcile({ ...configMap, [claudeAId]: claude("A renamed") }); + const rebuiltA = yield* registry.getInstance(claudeAId); + expect(rebuiltA).not.toBe(claudeA); + yield* rebuiltA?.snapshot.refresh ?? Effect.void; + expect(probedAccounts()).toEqual(["default", "default", "default", "other"]); }).pipe(Effect.provide(testLayer)), ); From 10c54c9487f6d15a8fb1030a5443cf06b26acdb4 Mon Sep 17 00:00:00 2001 From: Theo Browne Date: Sat, 26 Sep 2026 01:44:21 -0700 Subject: [PATCH 16/16] docs(server): drop a stale reference from the Claude probe cache doc Co-Authored-By: Claude Opus 5.5 (1M context) --- apps/server/src/provider/Drivers/ClaudeProbeCache.ts | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/apps/server/src/provider/Drivers/ClaudeProbeCache.ts b/apps/server/src/provider/Drivers/ClaudeProbeCache.ts index c68e596c545e..a2d3837e60bf 100644 --- a/apps/server/src/provider/Drivers/ClaudeProbeCache.ts +++ b/apps/server/src/provider/Drivers/ClaudeProbeCache.ts @@ -2,8 +2,7 @@ * One server-wide cache for the Claude capabilities probe. Claude instances * with the same probe input read the same account, so they share one cached * result, and concurrent reads of one input join one SDK probe. Each entry - * keeps 5 minutes, a failed probe (`undefined`) included, like the old - * per-instance cache. + * keeps 5 minutes, a failed probe (`undefined`) included. * * @module provider/Drivers/ClaudeProbeCache */