feat(ai): prevented usage cache from replaying stale values after invalidation
- Replace legacy credential-listing invalidation logic with prefix-based and exhaustive usage cache clearing. - Update auth broker usage stale route to await storage cache invalidation. - Add tests verifying explicit cache invalidation behavior prevents replaying last-good data on failure.
This commit is contained in:
@@ -741,7 +741,7 @@ export function startAuthBroker(opts: AuthBrokerServerOptions): AuthBrokerServer
|
||||
}
|
||||
if (req.method === "POST" && pathname === "/v1/usage/stale") {
|
||||
try {
|
||||
opts.storage.invalidateUsageCache?.();
|
||||
await opts.storage.invalidateUsageCache?.();
|
||||
logger.info("auth-broker usage cache invalidated", { peer });
|
||||
return json(200, { ok: true });
|
||||
} catch (error) {
|
||||
|
||||
@@ -782,6 +782,7 @@ interface UsageCache {
|
||||
get<T>(key: string): UsageCacheEntry<T> | undefined;
|
||||
getStale<T>(key: string): UsageCacheEntry<T> | undefined;
|
||||
set<T>(key: string, entry: UsageCacheEntry<T>): void;
|
||||
deletePrefix(prefix: string): boolean;
|
||||
cleanup?(): void;
|
||||
}
|
||||
|
||||
@@ -1185,6 +1186,12 @@ class AuthStorageUsageCache implements UsageCache {
|
||||
this.store.setCache(`${USAGE_CACHE_PREFIX}${key}`, payload, Math.floor(durableExpiresAt / 1000));
|
||||
}
|
||||
|
||||
deletePrefix(prefix: string): boolean {
|
||||
if (!this.store.deleteCachePrefix) return false;
|
||||
this.store.deleteCachePrefix(`${USAGE_CACHE_PREFIX}${prefix}`);
|
||||
return true;
|
||||
}
|
||||
|
||||
cleanup(): void {
|
||||
this.store.cleanExpiredCache();
|
||||
}
|
||||
@@ -2955,20 +2962,22 @@ export class AuthStorage {
|
||||
return baseUrl?.trim().replace(/\/+$/, "") ?? "";
|
||||
}
|
||||
|
||||
#usageCacheProviderKey(provider: Provider): string {
|
||||
const versionOverride = USAGE_REPORT_CACHE_KEY_VERSION_OVERRIDES[provider];
|
||||
return versionOverride === undefined ? provider : `${versionOverride}:${provider}`;
|
||||
}
|
||||
|
||||
#buildUsageReportCacheKey(request: UsageRequestDescriptor): string {
|
||||
const baseUrl = this.#normalizeUsageBaseUrl(request.baseUrl) || "default";
|
||||
const identity = this.#buildUsageCacheIdentity(request.credential);
|
||||
const versionOverride = USAGE_REPORT_CACHE_KEY_VERSION_OVERRIDES[request.provider];
|
||||
const providerKey = versionOverride === undefined ? request.provider : `${versionOverride}:${request.provider}`;
|
||||
const providerKey = this.#usageCacheProviderKey(request.provider);
|
||||
return `report:${providerKey}:${baseUrl}:${identity}`;
|
||||
}
|
||||
|
||||
#buildUsageReportsCacheKey(requests: ReadonlyArray<UsageRequestDescriptor>): string {
|
||||
const snapshot = requests
|
||||
.map(request => {
|
||||
const versionOverride = USAGE_REPORT_CACHE_KEY_VERSION_OVERRIDES[request.provider];
|
||||
const providerKey =
|
||||
versionOverride === undefined ? request.provider : `${versionOverride}:${request.provider}`;
|
||||
const providerKey = this.#usageCacheProviderKey(request.provider);
|
||||
return `${providerKey}:${this.#normalizeUsageBaseUrl(request.baseUrl) || "default"}:${this.#buildUsageCacheIdentity(request.credential)}`;
|
||||
})
|
||||
.sort()
|
||||
@@ -5729,31 +5738,30 @@ export class AuthStorage {
|
||||
}
|
||||
|
||||
/**
|
||||
* Force-invalidate cached usage reports so the next fetch retrieves fresh
|
||||
* values from upstream providers. If `provider` is specified, only that
|
||||
* provider's credentials are invalidated; otherwise, all credentials in the
|
||||
* store are invalidated.
|
||||
* Drop report snapshots for a user-requested refresh so a failed probe
|
||||
* cannot replay the pre-invalidation last-good value.
|
||||
*/
|
||||
async #clearUsageReportCache(provider?: string): Promise<void> {
|
||||
this.#usageCacheEpoch += 1;
|
||||
const prefix = provider ? `report:${this.#usageCacheProviderKey(provider)}:` : "report:";
|
||||
if (this.#usageCache.deletePrefix(prefix)) return;
|
||||
|
||||
// Third-party stores may not support prefix deletion. Clear every active
|
||||
// request key instead, including API-key and environment credentials.
|
||||
const requests = await this.#collectUsageRequests();
|
||||
for (const request of requests) {
|
||||
if (provider && request.provider !== provider) continue;
|
||||
this.#usageCache.set(this.#buildUsageReportCacheKey(request), { value: null, expiresAt: 0 });
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Discard cached usage reports before a user-requested refresh. The next
|
||||
* read probes upstream; a failure reports no fresh usage instead of replaying
|
||||
* an invalidated last-good snapshot.
|
||||
*/
|
||||
async invalidateUsageCache(provider?: string, signal?: AbortSignal): Promise<void> {
|
||||
if (provider) {
|
||||
this.#invalidateUsageReportCache(provider);
|
||||
} else {
|
||||
this.#usageCacheEpoch += 1;
|
||||
const expired = Date.now() - 1;
|
||||
try {
|
||||
const credentials = this.#store.listAuthCredentials();
|
||||
for (const entry of credentials) {
|
||||
if (entry.credential.type !== "oauth") continue;
|
||||
const cacheKey = this.#buildUsageReportCacheKey(
|
||||
this.#buildUsageRequestForOauth(entry.provider, entry.credential),
|
||||
);
|
||||
const existing = this.#usageCache.getStale<UsageReport | null>(cacheKey);
|
||||
this.#usageCache.set(cacheKey, { value: existing?.value ?? null, expiresAt: expired });
|
||||
}
|
||||
} catch (err) {
|
||||
logger.debug("Failed to list auth credentials for complete usage cache invalidation", { err });
|
||||
}
|
||||
}
|
||||
await this.#clearUsageReportCache(provider);
|
||||
|
||||
if (this.#store.invalidateUsageCache) {
|
||||
await this.#store.invalidateUsageCache(signal).catch(err => {
|
||||
|
||||
@@ -18,7 +18,7 @@ import {
|
||||
AuthStorage,
|
||||
type StoredAuthCredential,
|
||||
} from "@oh-my-pi/pi-ai/auth-storage";
|
||||
import type { UsageLimit, UsageReport } from "@oh-my-pi/pi-ai/usage";
|
||||
import type { UsageLimit, UsageProvider, UsageReport } from "@oh-my-pi/pi-ai/usage";
|
||||
import { alibabaTokenPlanUsageProvider } from "@oh-my-pi/pi-ai/usage/alibaba-token-plan";
|
||||
import * as claudeUsage from "@oh-my-pi/pi-ai/usage/claude";
|
||||
import { serializeAlibabaTokenPlanCredential } from "@oh-my-pi/pi-catalog/wire/alibaba-token-plan";
|
||||
@@ -326,6 +326,71 @@ describe("AuthStorage usage cache: last-good failure fallback", () => {
|
||||
expect(third).toHaveLength(1);
|
||||
expect(calls).toBe(3);
|
||||
});
|
||||
|
||||
it("does not replay last-good quota after an explicit invalidation", async () => {
|
||||
let calls = 0;
|
||||
vi.spyOn(claudeUsage.claudeUsageProvider, "fetchUsage").mockImplementation(async () => {
|
||||
calls += 1;
|
||||
return calls === 1 ? makeReport("a@example.com") : null;
|
||||
});
|
||||
|
||||
expect(anthropicReports(await storage.fetchUsageReports())).toHaveLength(1);
|
||||
await storage.invalidateUsageCache();
|
||||
|
||||
expect(anthropicReports(await storage.fetchUsageReports())).toHaveLength(0);
|
||||
expect(calls).toBe(2);
|
||||
});
|
||||
});
|
||||
|
||||
|
||||
describe("AuthStorage usage cache: explicit invalidation", () => {
|
||||
it("clears cached API-key reports before the next usage read", async () => {
|
||||
const store = makeStore([
|
||||
{
|
||||
id: 1,
|
||||
provider: "zai",
|
||||
credential: { type: "api_key", key: "zai-key" },
|
||||
disabledCause: null,
|
||||
},
|
||||
]);
|
||||
let calls = 0;
|
||||
const usageProvider: UsageProvider = {
|
||||
id: "zai",
|
||||
supports: params => params.provider === "zai" && params.credential.type === "api_key",
|
||||
async fetchUsage() {
|
||||
calls += 1;
|
||||
return {
|
||||
provider: "zai",
|
||||
fetchedAt: Date.now(),
|
||||
limits: [
|
||||
{
|
||||
id: "zai:requests:5h",
|
||||
label: "Z.AI Request Quota",
|
||||
scope: { provider: "zai", windowId: "5h" },
|
||||
amount: { used: calls === 1 ? 80 : 20, limit: 100, unit: "requests" },
|
||||
status: "ok",
|
||||
},
|
||||
],
|
||||
};
|
||||
},
|
||||
};
|
||||
const storage = new AuthStorage(store, {
|
||||
usageProviderResolver: provider => (provider === "zai" ? usageProvider : undefined),
|
||||
});
|
||||
await storage.reload();
|
||||
try {
|
||||
const initial = await storage.fetchUsageReports();
|
||||
expect(initial?.[0]?.limits[0]?.amount.used).toBe(80);
|
||||
|
||||
await storage.invalidateUsageCache();
|
||||
|
||||
const refreshed = await storage.fetchUsageReports();
|
||||
expect(refreshed?.[0]?.limits[0]?.amount.used).toBe(20);
|
||||
expect(calls).toBe(2);
|
||||
} finally {
|
||||
storage.close();
|
||||
}
|
||||
});
|
||||
});
|
||||
describe("AuthStorage usage cache: provider failure policy", () => {
|
||||
it("drops stale QwenCloud quota after the optional console session expires", async () => {
|
||||
|
||||
Reference in New Issue
Block a user