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
21 changes: 15 additions & 6 deletions open-sse/services/quotaFetchThrottle.ts
Original file line number Diff line number Diff line change
Expand Up @@ -65,17 +65,21 @@ export class MinIntervalThrottle {
});
try {
await prev;
// The chain is single-file here, so this read cannot race another acquire.
const now = this.clock.now();
if (this.lastStart !== 0) {
const jitter = this.jitterMs > 0 ? Math.floor(this.rand() * this.jitterMs) : 0;
const wait = this.lastStart + this.minIntervalMs + jitter - now;
if (wait > 0) await this.clock.sleep(wait);
}
this.lastStart = this.clock.now();
const startAt =
this.lastStart === 0 ? now : this.lastStart + this.minIntervalMs + this.nextJitter();
this.lastStart = startAt;
const wait = startAt - now;
if (wait > 0) await this.clock.sleep(wait);
} finally {
release();
}
}

private nextJitter(): number {
return this.jitterMs > 0 ? Math.floor(this.rand() * this.jitterMs) : 0;
}
}

/**
Expand Down Expand Up @@ -116,3 +120,8 @@ export function throttleQuotaFetch(): Promise<void> {
export function resetQuotaFetchThrottle(): void {
_sharedThrottle = null;
}

/** Test-only: install a throttle as the process-wide gate. */
export function setQuotaFetchThrottleForTests(throttle: MinIntervalThrottle | null): void {
_sharedThrottle = throttle;
}
6 changes: 6 additions & 0 deletions open-sse/services/usage.ts
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,7 @@ import {
import { getCursorUsage } from "./usage/cursor.ts";
import { getKimiUsage } from "./usage/kimi.ts";
import { getCodexUsage } from "./usage/codex.ts";
import { throttleQuotaFetch } from "./quotaFetchThrottle.ts";
import { getClaudeUsage, getClaudePlanLabel } from "./usage/claude.ts";
import { getKiroUsage, buildKiroUsageResult, discoverKiroProfileArn } from "./usage/kiro.ts";
// Re-exported para os testes kiro-* (importam de services/usage).
Expand Down Expand Up @@ -154,6 +155,11 @@ export async function getUsageForProvider(
case "claude":
return await getClaudeUsage(accessToken);
case "codex":
// /me/status and the quota-cache refresh reach this fetch with no
// caller-side pacing. Gate the start here so those paths, and any later
// one, cannot burst the usage endpoint. Callers that already acquired the
// shared gate wait at most one more interval.
await throttleQuotaFetch();
return await getCodexUsage(accessToken, providerSpecificData);
case "cursor":
return await getCursorUsage(accessToken || "", providerSpecificData);
Expand Down
10 changes: 9 additions & 1 deletion open-sse/services/usage/codex.ts
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,14 @@ const CODEX_CONFIG = {
usageUrl: "https://chatgpt.com/backend-api/wham/usage",
};

type CodexUsageFetch = (input: string, init: RequestInit) => Promise<Response>;
let codexUsageFetch: CodexUsageFetch = (input, init) => fetch(input, init);

/** Test-only: replace the usage request. Pass null to restore global fetch. */
export function setCodexUsageFetchForTests(fetchImpl: CodexUsageFetch | null): void {
codexUsageFetch = fetchImpl ?? ((input, init) => fetch(input, init));
}

/**
* Codex (OpenAI) Usage - Fetch from ChatGPT backend API
* IMPORTANT: Uses persisted workspaceId from OAuth to ensure correct workspace binding.
Expand Down Expand Up @@ -46,7 +54,7 @@ export async function getCodexUsage(
headers["chatgpt-account-id"] = accountId;
}

const response = await fetch(CODEX_CONFIG.usageUrl, {
const response = await codexUsageFetch(CODEX_CONFIG.usageUrl, {
method: "GET",
headers,
});
Expand Down
135 changes: 135 additions & 0 deletions tests/unit/codex-usage-status-throttle-15172.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,135 @@
/**
* #15172 — Codex usage reads that reach /wham/usage through getUsageForProvider
* (the /me/status fan-out and the quota-cache background refresh both do) must
* take the same min-interval gate as the dedicated quota fetchers. Auto-Ping
* already acquires that gate itself, so the extra wait on its path is at most
* one interval.
*/
import test from "node:test";
import assert from "node:assert/strict";

import { getUsageForProvider } from "../../open-sse/services/usage.ts";
import { setCodexUsageFetchForTests } from "../../open-sse/services/usage/codex.ts";
import {
MinIntervalThrottle,
resetQuotaFetchThrottle,
setQuotaFetchThrottleForTests,
type ThrottleClock,
} from "../../open-sse/services/quotaFetchThrottle.ts";
import {
createQuotaAutoPingState,
runQuotaAutoPingTick,
} from "../../src/lib/services/quotaAutoPing.ts";

const INTERVAL_MS = 250;

class FakeClock implements ThrottleClock {
t = 1_000;
now = () => this.t;
sleep = async (ms: number) => {
this.t += ms;
};
}

function installSharedThrottle(clock: FakeClock | null) {
const shared = new MinIntervalThrottle({
minIntervalMs: INTERVAL_MS,
jitterMs: 0,
...(clock ? { clock } : {}),
});
setQuotaFetchThrottleForTests(shared);
return () => resetQuotaFetchThrottle();
}

function installUsageFetch(clock: FakeClock | null) {
const starts: number[] = [];
setCodexUsageFetchForTests(async () => {
starts.push(clock ? clock.now() : Date.now());
return new Response(
JSON.stringify({
plan_type: "plus",
rate_limit: {
primary_window: { used_percent: 10 },
secondary_window: { used_percent: 20 },
},
}),
{ status: 200, headers: { "content-type": "application/json" } }
);
});
return {
starts,
restore: () => setCodexUsageFetchForTests(null),
};
}

test("#15172 concurrent Codex usage fetches start at least one interval apart", async () => {
const restoreThrottle = installSharedThrottle(null);
const { starts, restore } = installUsageFetch(null);

try {
const connections = ["a", "b", "c", "d"].map((id) => ({
id,
provider: "codex",
accessToken: `token-${id}`,
providerSpecificData: { workspaceId: `ws-${id}` },
}));
await Promise.all(connections.map((connection) => getUsageForProvider(connection)));

assert.equal(starts.length, 4);
for (let i = 1; i < starts.length; i++) {
const gap = starts[i] - starts[i - 1];
assert.ok(
gap >= INTERVAL_MS - 30,
`gap ${gap}ms between fetch ${i - 1} and ${i} is below one interval`
);
}
} finally {
restore();
restoreThrottle();
}
});

test("#15172 Auto-Ping's own gate adds at most one extra interval", async () => {
const clock = new FakeClock();
const restoreThrottle = installSharedThrottle(clock);
const { starts, restore } = installUsageFetch(clock);

const deps = {
getSettings: async () => ({
codexAutoPing: { connections: { "codex-1": true, "codex-2": true } },
}),
getProviderConnections: async () => [
{ id: "codex-1", provider: "codex", authType: "oauth", accessToken: "token-1" },
{ id: "codex-2", provider: "codex", authType: "oauth", accessToken: "token-2" },
],
updateProviderConnection: async () => null,
refreshAndUpdateCredentials: async (connection: unknown) => ({ connection }),
getCodexUsage: (accessToken?: string, providerSpecificData?: Record<string, unknown>) =>
getUsageForProvider({
provider: "codex",
accessToken,
providerSpecificData,
}),
throttleQuotaFetch: () =>
new MinIntervalThrottle({ minIntervalMs: INTERVAL_MS, jitterMs: 0, clock }).acquire(),
resolveProxyForConnection: async () => ({ proxy: null }),
runWithProxyContext: async (_proxy: unknown, callback: () => Promise<unknown>) => callback(),
getExecutor: () => ({
execute: async () => ({ response: { ok: true, text: async () => "" } }),
}),
canExecuteProvider: () => true,
isConnectionUnavailableToAuxiliaryActivity: async () => false,
resolvePingModel: async () => "gpt-5-codex",
};

try {
await runQuotaAutoPingTick(deps as never, createQuotaAutoPingState(), () => clock.now());
assert.equal(starts.length, 2);
const spread = starts[1] - starts[0];
assert.ok(spread >= INTERVAL_MS, `spread ${spread}ms is below one interval`);
assert.ok(spread <= INTERVAL_MS * 2, `spread ${spread}ms is more than one extra interval`);
} finally {
restore();
restoreThrottle();
}
});
Loading