diff --git a/packages/ai/src/auth-broker/remote-store.ts b/packages/ai/src/auth-broker/remote-store.ts index a527fae82..631da491f 100644 --- a/packages/ai/src/auth-broker/remote-store.ts +++ b/packages/ai/src/auth-broker/remote-store.ts @@ -253,6 +253,10 @@ export class RemoteAuthCredentialStore implements AuthCredentialStore { #usageInflight?: Promise; #credentialBlockReconcileAfter: Map = new Map(); #usageCacheEpoch = 0; + /** Per-snapshot lookup of oauth credentials by provider; rebuilt when `#snapshot` is replaced. */ + #usageFilterLookup?: { snapshot: SnapshotResponse; byProvider: Map }; + /** Memoized `#filterUsageReports` output, keyed on (input identity, lookup identity). */ + #usageFilterResult?: { input: UsageReport[]; byProvider: Map; output: UsageReport[] }; #closed = false; /** * `true` once the SSE consumer received its first frame and hasn't dropped @@ -962,18 +966,39 @@ export class RemoteAuthCredentialStore implements AuthCredentialStore { return overlay ?? matched; } + /** + * Hot path — called per `getUsageReport()`/`fetchUsageReports()` (status-line + * refresh cadence). The oauth-credential lookup is memoized on `#snapshot` + * identity (every update site replaces the reference), and the filtered + * output on (reports identity, lookup identity) — `#loadUsageReports` + * serves the same array for 15s, so steady-state calls are O(1). + */ #filterUsageReports(reports: UsageReport[]): UsageReport[] { const accountPool = this.#accountPool; if (!accountPool) return reports; - return reports.filter(report => { + let lookup = this.#usageFilterLookup; + if (!lookup || lookup.snapshot !== this.#snapshot) { + const byProvider = new Map(); + for (const entry of this.#snapshot.credentials) { + if (entry.credential.type !== "oauth") continue; + const list = byProvider.get(entry.provider); + if (list) list.push(entry.credential); + else byProvider.set(entry.provider, [entry.credential]); + } + lookup = { snapshot: this.#snapshot, byProvider }; + this.#usageFilterLookup = lookup; + } + const memo = this.#usageFilterResult; + if (memo && memo.input === reports && memo.byProvider === lookup.byProvider) return memo.output; + const byProvider = lookup.byProvider; + const output = reports.filter(report => { if (!accountPool.has(report.provider)) return true; - return this.#snapshot.credentials.some( - entry => - entry.provider === report.provider && - entry.credential.type === "oauth" && - usageReportMatchesCredential(report, entry.credential), - ); + const credentials = byProvider.get(report.provider); + if (!credentials) return false; + return credentials.some(credential => usageReportMatchesCredential(report, credential)); }); + this.#usageFilterResult = { input: reports, byProvider, output }; + return output; } ingestUsageReport(provider: Provider, credential: OAuthCredential, report: UsageReport): boolean { diff --git a/packages/ai/src/providers/devin.ts b/packages/ai/src/providers/devin.ts index 7b5113b81..0863a2af9 100644 --- a/packages/ai/src/providers/devin.ts +++ b/packages/ai/src/providers/devin.ts @@ -217,7 +217,12 @@ export const streamDevin: StreamFunction<"devin-agent"> = ( for (;;) { const { done, value } = await reader.read(); if (value && value.length > 0) { - pending = Buffer.concat([pending, value]); + // Steady state drains fully per chunk; view the fresh reader chunk + // instead of copying it through Buffer.concat (see aws-eventstream.ts). + pending = + pending.length === 0 + ? Buffer.from(value.buffer, value.byteOffset, value.byteLength) + : Buffer.concat([pending, value]); } while (pending.length >= 5) { diff --git a/packages/ai/test/remote-auth-store.test.ts b/packages/ai/test/remote-auth-store.test.ts index f7013a34d..6e47604ae 100644 --- a/packages/ai/test/remote-auth-store.test.ts +++ b/packages/ai/test/remote-auth-store.test.ts @@ -1226,6 +1226,79 @@ describe("RemoteAuthCredentialStore + AuthStorage integration", () => { } }); + test("usage report filter memoizes per (reports, snapshot) and invalidates when the snapshot changes", async () => { + const brokerClient = new AuthBrokerClient({ url: "http://127.0.0.1:9", token: "unused" }); + const now = Date.now(); + const oauthCredential = { + type: "oauth" as const, + access: "oauth-access", + refresh: REMOTE_REFRESH_SENTINEL, + expires: now + 120_000, + accountId: "oauth-account", + email: "oauth@example.com", + }; + const matchingReport: UsageReport = { + provider: "anthropic", + fetchedAt: now, + limits: [], + metadata: { accountId: "oauth-account", email: "oauth@example.com" }, + }; + const strangerReport: UsageReport = { + provider: "anthropic", + fetchedAt: now, + limits: [], + metadata: { accountId: "stranger-account", email: "stranger@example.com" }, + }; + const nonPooledReport: UsageReport = { + provider: "openai-codex", + fetchedAt: now, + limits: [], + metadata: { accountId: "codex-account" }, + }; + const reports: UsageReport[] = [matchingReport, strangerReport, nonPooledReport]; + const fetchSpy = vi.spyOn(brokerClient, "fetchUsage").mockResolvedValue({ generatedAt: now, reports }); + vi.spyOn(brokerClient, "disableCredential").mockResolvedValue({ ok: true }); + const oauthIdentity = "email:oauth@example.com"; + const remoteStore = new RemoteAuthCredentialStore({ + client: brokerClient, + streamSnapshots: false, + accountPool: new Map([["anthropic", new Set([oauthIdentity])]]), + initialSnapshot: { + generation: 1, + generatedAt: now, + serverNowMs: now, + refresher: { enabled: false, intervalMs: 0, skewMs: 0, nextSweepInMs: Number.MAX_SAFE_INTEGER }, + credentials: [ + { + id: 1, + provider: "anthropic", + credential: oauthCredential, + identityKey: oauthIdentity, + rotatesInMs: null, + }, + ], + }, + }); + try { + // Semantics: matching oauth report kept, pooled report without a + // matching credential dropped, non-pooled provider passed through. + const first = await remoteStore.fetchUsageReports(); + expect(first).toEqual([matchingReport, nonPooledReport]); + // Same cached reports array + same snapshot → memoized output, same identity. + const second = await remoteStore.fetchUsageReports(); + expect(second).toBe(first!); + expect(fetchSpy).toHaveBeenCalledTimes(1); + // Snapshot replacement (credential removal) invalidates the memo: the + // previously matching report is no longer attributable and disappears. + remoteStore.deleteAuthCredential(1, "test"); + const third = await remoteStore.fetchUsageReports(); + expect(third).not.toBe(first!); + expect(third).toEqual([nonPooledReport]); + } finally { + remoteStore.close(); + } + }); + test("rejects a refreshed credential whose identity leaves the account pool", async () => { const brokerClient = new AuthBrokerClient({ url: "http://127.0.0.1:9", token: "unused" }); const now = Date.now();