Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions apps/server/src/provider/Drivers/ClaudeDriver.ts
Original file line number Diff line number Diff line change
Expand Up @@ -261,6 +261,7 @@ export const ClaudeDriver: ProviderDriver<ClaudeSettings, ClaudeDriverEnv> = {
accentColor,
enabled,
snapshot,
invalidateCaches: Cache.invalidateAll(capabilitiesProbeCache),
snapshotForCwd,
adapter,
textGeneration,
Expand Down
5 changes: 3 additions & 2 deletions apps/server/src/provider/Drivers/CursorDriver.ts
Original file line number Diff line number Diff line change
Expand Up @@ -134,11 +134,11 @@ export const CursorDriver: ProviderDriver<CursorSettings, CursorDriverEnv> = {

const textGeneration = yield* makeCursorTextGeneration(effectiveConfig, processEnv);

const discoverModels = yield* makeCursorModelDiscovery(effectiveConfig, processEnv);
const modelDiscovery = yield* makeCursorModelDiscovery(effectiveConfig, processEnv);
const checkProvider = checkCursorProviderStatus(
effectiveConfig,
processEnv,
discoverModels,
modelDiscovery.discover,
).pipe(
Effect.flatMap((snapshot) =>
effectiveConfig.enabled && snapshot.installed && snapshot.auth.status === "authenticated"
Expand Down Expand Up @@ -217,6 +217,7 @@ export const CursorDriver: ProviderDriver<CursorSettings, CursorDriverEnv> = {
accentColor,
enabled,
snapshot,
invalidateCaches: modelDiscovery.invalidate,
snapshotForCwd: (cwd) =>
!effectiveConfig.enabled
? snapshot.getSnapshot
Expand Down
6 changes: 5 additions & 1 deletion apps/server/src/provider/Layers/CursorProvider.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -742,7 +742,7 @@ describe("discoverCursorModelsViaAcp", () => {
apiEndpoint: "",
customModels: [],
};
const discover = yield* makeCursorModelDiscovery(settings, {
const { discover, invalidate } = yield* makeCursorModelDiscovery(settings, {
...process.env,
T3_ACP_REQUEST_LOG_PATH: requestLogPath,
});
Expand All @@ -755,6 +755,10 @@ describe("discoverCursorModelsViaAcp", () => {
yield* fileSystem.writeFileString(requestLogPath, "");
expect(yield* discover(about)).toEqual(first);
expect(yield* fileSystem.readFileString(requestLogPath)).toBe("");
yield* invalidate;
expect(yield* discover(about)).toEqual(first);
expect(yield* fileSystem.readFileString(requestLogPath)).toContain("initialize");
yield* fileSystem.writeFileString(requestLogPath, "");
yield* discover({ ...about, version: "2026.08.12" });
expect(yield* fileSystem.readFileString(requestLogPath)).toContain("initialize");
yield* fileSystem.writeFileString(requestLogPath, "");
Expand Down
7 changes: 5 additions & 2 deletions apps/server/src/provider/Layers/CursorProvider.ts
Original file line number Diff line number Diff line change
Expand Up @@ -678,8 +678,11 @@ export const makeCursorModelDiscovery = Effect.fn("makeCursorModelDiscovery")(fu
Exit.isSuccess(exit) && exit.value.length > 0 ? Duration.minutes(30) : Duration.zero,
},
);
return (about: Pick<CursorAboutResult, "version" | "auth">) =>
Cache.get(cache, JSON.stringify([about.version, about.auth]));
return {
discover: (about: Pick<CursorAboutResult, "version" | "auth">) =>
Cache.get(cache, JSON.stringify([about.version, about.auth])),
invalidate: Cache.invalidateAll(cache),
};
});

function getCursorFallbackModels(
Expand Down
79 changes: 79 additions & 0 deletions apps/server/src/provider/ModelManifest.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -342,6 +342,84 @@ const serviceLayers = (input: {
);

describe("ModelManifest service", () => {
it.live("explicit refresh bypasses fresh memory and disk caches", () => {
let fetchCount = 0;
const updated: ModelManifestData = {
...REMOTE_MANIFEST,
currentModels: { codex: ["gpt-reloaded"] },
};
return Effect.gen(function* () {
const service = yield* make;
assert.deepStrictEqual(yield* service.refresh, REMOTE_MANIFEST);
assert.deepStrictEqual(yield* service.refresh, REMOTE_MANIFEST);
assert.strictEqual(fetchCount, 1);

const rebooted = yield* make;
assert.deepStrictEqual(yield* rebooted.refresh, REMOTE_MANIFEST);
assert.strictEqual(fetchCount, 1);
assert.deepStrictEqual(yield* rebooted.forceRefresh, updated);
assert.strictEqual(fetchCount, 2);
assert.deepStrictEqual(yield* rebooted.current, updated);
assert.deepStrictEqual(yield* (yield* make).current, updated);
}).pipe(
Effect.scoped,
Effect.provide(
serviceLayers({
prefix: "model-manifest-force-refresh-test",
response: () => Response.json(fetchCount++ === 0 ? REMOTE_MANIFEST : updated),
}),
),
);
});

it.live("explicit refresh retries immediately after failure and preserves last-good data", () => {
let fetchCount = 0;
return Effect.gen(function* () {
const service = yield* make;
assert.deepStrictEqual(yield* service.refresh, REMOTE_MANIFEST);
assert.deepStrictEqual(yield* service.forceRefresh, REMOTE_MANIFEST);
assert.deepStrictEqual(yield* service.current, REMOTE_MANIFEST);
assert.deepStrictEqual(yield* (yield* make).current, REMOTE_MANIFEST);
assert.strictEqual(fetchCount, 2);
assert.deepStrictEqual(yield* service.forceRefresh, REMOTE_MANIFEST);
assert.strictEqual(fetchCount, 3);
}).pipe(
Effect.scoped,
Effect.provide(
serviceLayers({
prefix: "model-manifest-force-retry-test",
response: () =>
fetchCount++ === 1
? new Response(null, { status: 503 })
: Response.json(REMOTE_MANIFEST),
}),
),
);
});

it.live("explicit refresh bypasses the retry delay after an initial failure", () => {
let fetchCount = 0;
return Effect.gen(function* () {
const service = yield* make;
assert.deepStrictEqual(yield* service.refresh, BUNDLED_MODEL_MANIFEST);
assert.deepStrictEqual(yield* service.refresh, BUNDLED_MODEL_MANIFEST);
assert.strictEqual(fetchCount, 1);
assert.deepStrictEqual(yield* service.forceRefresh, REMOTE_MANIFEST);
assert.strictEqual(fetchCount, 2);
}).pipe(
Effect.scoped,
Effect.provide(
serviceLayers({
prefix: "model-manifest-force-initial-retry-test",
response: () =>
fetchCount++ === 0
? new Response(null, { status: 503 })
: Response.json(REMOTE_MANIFEST),
}),
),
);
});

it.live("prefers a fetched manifest over the bundle and caches it to disk", () =>
Effect.gen(function* () {
const service = yield* make;
Expand Down Expand Up @@ -457,6 +535,7 @@ describe("ModelManifest service", () => {
),
);
assert.deepStrictEqual(yield* service.refresh, BUNDLED_MODEL_MANIFEST);
assert.deepStrictEqual(yield* service.forceRefresh, BUNDLED_MODEL_MANIFEST);
assert.strictEqual(fetchCount, 0);
}).pipe(
Effect.scoped,
Expand Down
10 changes: 7 additions & 3 deletions apps/server/src/provider/ModelManifest.ts
Original file line number Diff line number Diff line change
Expand Up @@ -315,6 +315,8 @@ export class ModelManifest extends Context.Service<
readonly current: Effect.Effect<ModelManifestData>;
/** Manifest after a TTL-gated remote refresh; never fails. */
readonly refresh: Effect.Effect<ModelManifestData>;
/** Explicit refresh bypasses freshness and retry timers, retaining last-good data. */
readonly forceRefresh: Effect.Effect<ModelManifestData>;
/** Forks `refresh` into the service's own scope. Drivers call this from
* provider checks: the fetch is process-shared state, so it must survive
* the teardown of whichever instance happened to trigger it. */
Expand All @@ -326,6 +328,7 @@ export class ModelManifest extends Context.Service<
const BundledOnlyModelManifest: ModelManifest["Service"] = {
current: Effect.succeed(BUNDLED_MODEL_MANIFEST),
refresh: Effect.succeed(BUNDLED_MODEL_MANIFEST),
forceRefresh: Effect.succeed(BUNDLED_MODEL_MANIFEST),
refreshInBackground: Effect.void,
};

Expand Down Expand Up @@ -369,16 +372,16 @@ export const make = Effect.gen(function* () {
}),
);

const refresh = Effect.fn("ModelManifest.refresh")(function* () {
const refresh = Effect.fn("ModelManifest.refresh")(function* (force = false) {
yield* ensureDiskCacheLoaded;
const now = yield* Clock.currentTimeMillis;
// A timestamp in the future means the wall clock moved backwards (the
// disk cache crosses restarts, so monotonic time cannot cover it). Treat
// it as expired: the refetch rewrites both timestamps and self-heals.
const isWithin = (sinceMs: number | null, windowMs: number) =>
sinceMs !== null && now >= sinceMs && now - sinceMs < windowMs;
if (isWithin(fetchedAtMs, MANIFEST_TTL_MS)) return manifest;
if (isWithin(lastAttemptMs, MANIFEST_RETRY_MS)) return manifest;
if (!force && isWithin(fetchedAtMs, MANIFEST_TTL_MS)) return manifest;
if (!force && isWithin(lastAttemptMs, MANIFEST_RETRY_MS)) return manifest;

// The same switch that gates provider CLI update checks. It stops network
// fetches only: a manifest already cached on disk from an earlier fetch
Expand Down Expand Up @@ -413,6 +416,7 @@ export const make = Effect.gen(function* () {
return ModelManifest.of({
current: ensureDiskCacheLoaded.pipe(Effect.map(() => manifest)),
refresh: guardedRefresh,
forceRefresh: refreshSemaphore.withPermits(1)(refresh(true)),
refreshInBackground: Effect.forkIn(guardedRefresh, serviceScope).pipe(Effect.asVoid),
});
});
Expand Down
2 changes: 2 additions & 0 deletions apps/server/src/provider/ProviderDriver.ts
Original file line number Diff line number Diff line change
Expand Up @@ -74,6 +74,8 @@ export interface ProviderInstance {
readonly snapshot: ServerProviderShape;
readonly snapshotForCwd?: (cwd: string) => Effect.Effect<ServerProvider, ProviderDriverError>;
readonly refreshModels?: () => Effect.Effect<void, ProviderDriverError>;
/** Invalidate T3-owned discovery caches before an explicit provider refresh. */
readonly invalidateCaches?: Effect.Effect<void>;
/**
* Redeem one banked rate-limit reset credit on the signed-in account, then
* re-probe so the snapshot reflects the cleared windows. Account-level,
Expand Down
104 changes: 103 additions & 1 deletion apps/server/src/server.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -132,6 +132,7 @@ import { OrchestrationEventStoreLive } from "./persistence/Layers/OrchestrationE
import { OrchestrationEventStore } from "./persistence/Services/OrchestrationEventStore.ts";
import { PersistenceSqlError } from "./persistence/Errors.ts";
import * as ProviderRegistry from "./provider/Services/ProviderRegistry.ts";
import * as ModelManifest from "./provider/ModelManifest.ts";
import * as ProviderService from "./provider/Services/ProviderService.ts";
import { ProviderAuthService } from "./provider/Services/ProviderAuthService.ts";
import { ProviderInstanceRegistry } from "./provider/Services/ProviderInstanceRegistry.ts";
Expand All @@ -142,7 +143,10 @@ import {
import type { ProviderInstance } from "./provider/ProviderDriver.ts";
import * as ProviderSessionDirectory from "./provider/Services/ProviderSessionDirectory.ts";
import { ProviderAdapterRequestError } from "./provider/Errors.ts";
import { makeManualOnlyProviderMaintenanceCapabilities } from "./provider/providerMaintenance.ts";
import {
makeManualOnlyProviderMaintenanceCapabilities,
ProviderVersionCache,
} from "./provider/providerMaintenance.ts";
import * as ServerLifecycleEvents from "./serverLifecycleEvents.ts";
import * as ServerRuntimeStartup from "./serverRuntimeStartup.ts";
import * as ServiceLauncherClient from "./cloud/serviceLauncherClient.ts";
Expand Down Expand Up @@ -516,6 +520,7 @@ const buildAppUnderTest = (options?: {
keybindings?: Partial<Keybindings.Keybindings["Service"]>;
environmentTheme?: Partial<EnvironmentTheme.EnvironmentThemeService["Service"]>;
providerRegistry?: Partial<ProviderRegistry.ProviderRegistry["Service"]>;
modelManifest?: Partial<ModelManifest.ModelManifest["Service"]>;
usageLimitSources?: Partial<UsageLimitSources.UsageLimitSources["Service"]>;
providerService?: Partial<ProviderService.ProviderService["Service"]>;
providerAuth?: Partial<ProviderAuthService["Service"]>;
Expand Down Expand Up @@ -790,6 +795,10 @@ const buildAppUnderTest = (options?: {
),
Layer.provide(
Layer.mergeAll(
Layer.mock(ModelManifest.ModelManifest)({
forceRefresh: Effect.succeed(ModelManifest.BUNDLED_MODEL_MANIFEST),
...options?.layers?.modelManifest,
}),
Layer.mock(ProviderRegistry.ProviderRegistry)({
getProviders: Effect.succeed([]),
refresh: () => Effect.succeed([]),
Expand Down Expand Up @@ -6337,6 +6346,99 @@ it.layer(NodeServices.layer)("server router seam", (it) => {
}).pipe(Effect.provide(NodeHttpServer.layerTest)),
);

for (const mode of ["all", "targeted", "background"] as const) {
it.effect(`provider refresh invalidates T3 caches before probing (${mode})`, () => {
const driver = ProviderDriverKind.make("codex");
const instanceIds = [ProviderInstanceId.make("codex"), ProviderInstanceId.make("codex_work")];
const packageNames = ["@example/personal", "@example/work"];
const versionCache = new Map(
packageNames.map((name) => [
name,
{
expiresAt: Number.MAX_SAFE_INTEGER,
version: "1.0.0",
},
]),
);
const invalidated: string[] = [];
const freshMaintenance: string[] = [];
let manifestRefreshed = false;
let probed = false;
const instances = instanceIds.map(
(instanceId, index) =>
({
instanceId,
driverKind: driver,
continuationIdentity: { driverKind: driver, continuationKey: instanceId },
displayName: undefined,
enabled: true,
invalidateCaches: Effect.sync(() => {
invalidated.push(instanceId);
}),
snapshot: {
resolveMaintenance: (options) =>
Effect.sync(() => {
assert.isTrue(options?.fresh);
freshMaintenance.push(instanceId);
return makeManualOnlyProviderMaintenanceCapabilities({
provider: driver,
packageName: packageNames[index]!,
});
}),
getSnapshot: Effect.never,
refresh: Effect.never,
streamChanges: Stream.empty,
applyUsageLimits: () => Effect.void,
},
adapter: {} as ProviderInstance["adapter"],
textGeneration: {} as ProviderInstance["textGeneration"],
}) satisfies ProviderInstance,
);
const expected =
mode === "background" ? [] : mode === "targeted" ? [instanceIds[1]!] : instanceIds;
const probe = Effect.sync(() => {
probed = true;
assert.equal(manifestRefreshed, mode !== "background");
assert.deepEqual(invalidated.toSorted(), expected.toSorted());
assert.deepEqual(freshMaintenance.toSorted(), expected.toSorted());
for (let index = 0; index < instanceIds.length; index++) {
assert.equal(
versionCache.has(packageNames[index]!),
!expected.includes(instanceIds[index]!),
);
}
return [];
});
return Effect.gen(function* () {
yield* buildAppUnderTest({
layers: {
modelManifest: {
forceRefresh: Effect.sync(() => {
manifestRefreshed = true;
return ModelManifest.BUNDLED_MODEL_MANIFEST;
}),
},
providerInstanceRegistry: { listInstances: Effect.succeed(instances) },
providerRegistry: { refresh: () => probe, refreshInstance: () => probe },
},
});
const wsUrl = yield* getWsServerUrl("/ws");
yield* Effect.scoped(
withWsRpcClient(wsUrl, (client) =>
client[WS_METHODS.serverRefreshProviders]({
...(mode === "targeted" ? { instanceId: instanceIds[1]! } : {}),
...(mode !== "background" ? { refreshModels: true } : {}),
}),
),
);
assert.isTrue(probed);
}).pipe(
Effect.provideService(ProviderVersionCache, versionCache),
Effect.provide(NodeHttpServer.layerTest),
);
});
}

it.effect("serves config on reconnect without starting provider probes", () =>
Effect.gen(function* () {
const refresh = vi.fn(() => Effect.never);
Expand Down
26 changes: 26 additions & 0 deletions apps/server/src/ws.ts
Original file line number Diff line number Diff line change
Expand Up @@ -109,6 +109,8 @@ import {
observeRpcStreamEffect as instrumentRpcStreamEffect,
} from "./observability/RpcInstrumentation.ts";
import * as ProviderRegistry from "./provider/Services/ProviderRegistry.ts";
import * as ModelManifest from "./provider/ModelManifest.ts";
import * as ProviderMaintenance from "./provider/providerMaintenance.ts";
import * as ProviderService from "./provider/Services/ProviderService.ts";
import * as ProviderSessionDirectory from "./provider/Services/ProviderSessionDirectory.ts";
import * as ProviderMaintenanceRunner from "./provider/providerMaintenanceRunner.ts";
Expand Down Expand Up @@ -561,6 +563,8 @@ const makeWsRpcLayer = (
yield* Effect.context<Effect.Services<ReturnType<typeof remoteSshDeviceHosts>>>();
const portDiscovery = yield* PortScanner.PortDiscovery;
const providerRegistry = yield* ProviderRegistry.ProviderRegistry;
const modelManifest = yield* ModelManifest.ModelManifest;
const providerVersionCache = yield* ProviderMaintenance.ProviderVersionCache;
const providerService = yield* ProviderService.ProviderService;
const providerSessionDirectory = yield* ProviderSessionDirectory.ProviderSessionDirectory;
const providerMaintenanceRunner = yield* ProviderMaintenanceRunner.ProviderMaintenanceRunner;
Expand Down Expand Up @@ -2361,6 +2365,28 @@ const makeWsRpcLayer = (
observeRpcEffect(
WS_METHODS.serverRefreshProviders,
Effect.gen(function* () {
// Only explicit catalog refreshes bypass T3's caches. Workspace
// discovery and background status checks retain their timers.
if (input.refreshModels) {
yield* modelManifest.forceRefresh;
const instances = yield* providerInstances.listInstances;
yield* Effect.forEach(
instances.filter(
(instance) =>
input.instanceId === undefined || input.instanceId === instance.instanceId,
),
(instance) =>
Effect.gen(function* () {
yield* instance.invalidateCaches ?? Effect.void;
const maintenance = yield* instance.snapshot.resolveMaintenance({
fresh: true,
});
if (maintenance.packageName)
providerVersionCache.delete(maintenance.packageName);
}),
{ concurrency: "unbounded", discard: true },
);
}
// An untargeted refresh is "re-read everything's status", which
// includes quota from configured usage-limit sources. Awaited,
// not forked: the RPC scope closes on return and would
Expand Down
Loading
Loading