refactor(ai): optimized credential lookup and stream chunk processing
- Memoize credential lookup and filtered usage report output in RemoteAuthCredentialStore for O(1) steady-state performance. - Optimize Devin stream chunk processing to avoid redundant buffer copies during steady-state reads. - Add comprehensive integration test covering usage report filter memoization and cache invalidation on snapshot changes.
This commit is contained in:
@@ -253,6 +253,10 @@ export class RemoteAuthCredentialStore implements AuthCredentialStore {
|
||||
#usageInflight?: Promise<UsageReport[] | null>;
|
||||
#credentialBlockReconcileAfter: Map<string, number> = new Map();
|
||||
#usageCacheEpoch = 0;
|
||||
/** Per-snapshot lookup of oauth credentials by provider; rebuilt when `#snapshot` is replaced. */
|
||||
#usageFilterLookup?: { snapshot: SnapshotResponse; byProvider: Map<Provider, OAuthCredential[]> };
|
||||
/** Memoized `#filterUsageReports` output, keyed on (input identity, lookup identity). */
|
||||
#usageFilterResult?: { input: UsageReport[]; byProvider: Map<Provider, OAuthCredential[]>; 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<Provider, OAuthCredential[]>();
|
||||
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 {
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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();
|
||||
|
||||
Reference in New Issue
Block a user