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
57 changes: 56 additions & 1 deletion src/adapters/rateLimitError.test.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
import { describe, it, expect } from 'vitest';
import { classifyLimitResponse, detectRateLimit, rateLimitFromCodexHeaders, rateLimitFromHttpResponse, matchesRateLimitMessage, RateLimitError } from './rateLimitError.js';
import { classifyLimitResponse, detectRateLimit, parseRetryAfterSeconds, rateLimitFromCodexHeaders, rateLimitFromHttpResponse, matchesRateLimitMessage, RateLimitError } from './rateLimitError.js';
import { resolveLimitResponse, throttleWaitMs } from './throttleRetry.js';
import { isInfraError } from './errorClassification.js';
import { runAgenticLoop } from './agenticLoop.js';
Expand All @@ -21,6 +21,9 @@ describe('per-provider usage-limit recognition (INT-2520 audit)', () => {
['OpenRouter 402 relaying upstream BYOK balance',
'{"error":{"message":"Provider returned error","code":402,"metadata":{"raw":"{\\"code\\":402,\\"msg\\":\\"insufficient balance\\"}","provider_name":"AtlasCloud","is_byok":true}}}'],
['HTTP 429 too many requests (local)', 'Local API error (429): Too Many Requests'],
['OpenRouter 429 prose', 'Rate limit exceeded: 1000 requests per 1 day'],
['local 429 overloaded', '{"error":"Too Many Requests: server is overloaded"}'],
['local overloaded body', '{"error":"server is overloaded"}'],
];
for (const [name, output] of REAL_LIMIT_OUTPUTS) {
it(`detects: ${name}`, () => {
Expand All @@ -42,6 +45,9 @@ describe('per-provider usage-limit recognition (INT-2520 audit)', () => {
// anchored on the provider's JSON key rather than added as a bare substring.
"throw new Error('Insufficient balance') // wallet guard",
'if (res.status === 402) throw new Error("Insufficient balance for this transfer");',
// A 529 capacity blip is a single word, and it must stay an infra backoff:
// matching the bare word here would make every overload pause the scheduler.
'Anthropic API error: overloaded',
];
for (const b of benign) {
expect(matchesRateLimitMessage(b)).toBe(false);
Expand Down Expand Up @@ -379,6 +385,55 @@ describe('resolveLimitResponse gating (INT-2907)', () => {
});
});

describe('Retry-After parsing (AGT-3442)', () => {
it('parses delta-seconds and HTTP-date Retry-After values', () => {
expect(parseRetryAfterSeconds('120')).toBe(120);
expect(parseRetryAfterSeconds(' 45 ')).toBe(45);
// Prefix digits must not silently win over a malformed token.
expect(parseRetryAfterSeconds('60xyz')).toBeUndefined();
expect(parseRetryAfterSeconds('not-a-date')).toBeUndefined();

const future = new Date(Date.now() + 180_000);
const before = Math.floor(Date.now() / 1000);
const seconds = parseRetryAfterSeconds(future.toUTCString());
const after = Math.floor(Date.now() / 1000);
expect(seconds).toBeDefined();
const expected = Math.floor(future.getTime() / 1000);
expect(seconds!).toBeGreaterThanOrEqual(expected - after);
expect(seconds!).toBeLessThanOrEqual(expected - before);
});

it('surfaces an HTTP-date Retry-After as seconds-from-now via classifyLimitResponse', () => {
// Relative future date so the assertion does not rot as wall-clock moves.
const future = new Date(Date.now() + 120_000);
const headers = new Headers({ 'retry-after': future.toUTCString() });
const before = Math.floor(Date.now() / 1000);
const result = classifyLimitResponse(headers, '{}');
const after = Math.floor(Date.now() / 1000);
expect(result.quota).toBe(false); // no quota-exhausted body signature
const expected = Math.floor(future.getTime() / 1000);
expect(result.retryAfterSeconds!).toBeGreaterThanOrEqual(expected - after);
expect(result.retryAfterSeconds!).toBeLessThanOrEqual(expected - before);
});

it('rejects an unusable Retry-After instead of inventing a wait', () => {
// A date-shaped-but-invalid token must not become a NaN/0-second pause.
expect(classifyLimitResponse(new Headers({ 'retry-after': 'Fri, 99 Foo 9999' }), '').retryAfterSeconds)
.toBeUndefined();
});

it('sets RateLimitError.resetsAt from an HTTP-date Retry-After on a 429', () => {
const future = new Date(Date.now() + 300_000);
const headers = new Headers({ 'retry-after': future.toUTCString() });
const err = rateLimitFromHttpResponse(429, headers, '{"error":"rate limit"}');
expect(err).toBeInstanceOf(RateLimitError);
const expected = Math.floor(future.getTime() / 1000);
// ±1s: parseRetryAfterSeconds and extractResetsAt each sample Date.now().
expect(err!.resetsAt).toBeGreaterThanOrEqual(expected - 1);
expect(err!.resetsAt).toBeLessThanOrEqual(expected + 1);
});
});

describe('throttle backoff + downstream classification (INT-2907)', () => {
it('honors Retry-After, caps it, and otherwise escalates the backoff', () => {
// Backoff carries up to 1s of jitter so concurrent subagents don't retry in lockstep.
Expand Down
45 changes: 41 additions & 4 deletions src/adapters/rateLimitError.ts
Original file line number Diff line number Diff line change
Expand Up @@ -60,6 +60,12 @@ const RATE_LIMIT_SUBSTRINGS: readonly string[] = [
'purchase more credits', // codex CLI stdout error event
'exceeded your current quota',// OpenAI insufficient_quota human message
'too many requests', // HTTP 429 standard reason (local/lmstudio/others)
'rate limit exceeded', // OpenRouter/OpenAI 429 prose ("Rate limit exceeded: 1000 requests per 1 day")
// local/lmstudio 429 body ("server is overloaded"). Deliberately NOT the bare
// word "overloaded": Anthropic/OpenRouter report a 529 capacity blip with that
// exact single word, and that is an infra backoff, not a scheduler pause —
// matching it here would re-bucket every 529 ahead of isInfraError.
'server is overloaded',
];

// Regex signatures that need structure (co-occurrence / numeric context) to stay
Expand Down Expand Up @@ -110,6 +116,23 @@ export function parseResetsAtFromBody(text: string): number | undefined {
return m ? parseInt(m[1], 10) : undefined;
}

/**
* RFC 7231 §7.1.3: Retry-After is either 1*DIGIT delta-seconds or an HTTP-date.
* Returns seconds-from-now when parseable; undefined when the value is unusable.
*/
export function parseRetryAfterSeconds(value: string): number | undefined {
const trimmed = value.trim();
// Delta-seconds must be the entire token — parseInt("Fri, …") is NaN, but
// parseInt("60xyz") would silently accept a prefix, so require /^\d+$/.
if (/^\d+$/.test(trimmed)) {
const delta = parseInt(trimmed, 10);
return Number.isFinite(delta) ? delta : undefined;
}
const dateMs = Date.parse(trimmed);
if (!Number.isFinite(dateMs)) return undefined;
return Math.max(0, Math.floor(dateMs / 1000) - Math.floor(Date.now() / 1000));
}

/** Pull a unix reset timestamp (seconds) out of headers or a JSON body, if present. */
function extractResetsAt(headers: Headers | undefined, body: string): number | undefined {
const fromHeader = (k: string): number | undefined => {
Expand All @@ -119,16 +142,21 @@ function extractResetsAt(headers: Headers | undefined, body: string): number | u
};
// Only headers/fields that are genuinely UNIX-epoch seconds or seconds-from-now:
// - x-codex-primary-reset-at: epoch seconds
// - Retry-After: seconds-from-now (→ convert to epoch)
// - Retry-After: seconds-from-now (→ convert to epoch) OR an HTTP-date
// - body "resets_at": epoch seconds
// Deliberately NOT x-ratelimit-reset-requests/-tokens: OpenAI returns those as
// DURATION strings ("1s", "6ms", "2m59s"), not epoch — parseInt would yield a
// 1970 timestamp and defeat the pause. Omitting them falls back to the safe
// 60s default, which is correct rather than wrong. (INT-2520 review)
const codexReset = fromHeader('x-codex-primary-reset-at');
if (codexReset != null) return codexReset;
const retryAfter = fromHeader('retry-after');
if (retryAfter != null) return Math.floor(Date.now() / 1000) + retryAfter;
const retryAfter = headers?.get('retry-after');
if (retryAfter != null) {
// RFC 7231 §7.1.3 via parseRetryAfterSeconds (delta-seconds or HTTP-date).
// Convert seconds-from-now → absolute epoch for RateLimitError.resetsAt.
const seconds = parseRetryAfterSeconds(retryAfter);
if (seconds != null) return Math.floor(Date.now() / 1000) + seconds;
}
return parseResetsAtFromBody(body);
}

Expand Down Expand Up @@ -201,7 +229,16 @@ export function classifyLimitResponse(headers: Headers | undefined, body: string
return Number.isFinite(n) ? n : undefined;
};
const usedPercent = num('x-codex-primary-used-percent');
const retryAfterSeconds = num('retry-after');
// Retry-After is RFC 7231 delta-seconds OR an HTTP-date; parseInt on a date
// yields NaN and silently drops the server's wait. Fall back to the codex
// absolute reset epoch, converted to seconds-from-now. (AGT-3442)
let retryAfterSeconds: number | undefined;
const retryAfterHeader = headers?.get('retry-after');
if (retryAfterHeader != null) retryAfterSeconds = parseRetryAfterSeconds(retryAfterHeader);
if (retryAfterSeconds == null) {
const resetAt = num('x-codex-primary-reset-at');
if (resetAt != null) retryAfterSeconds = Math.max(0, resetAt - Math.floor(Date.now() / 1000));
}
const lower = body.toLowerCase();
const quota =
QUOTA_EXHAUSTED_SUBSTRINGS.some((s) => lower.includes(s)) ||
Expand Down
63 changes: 63 additions & 0 deletions src/adapters/webTools.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -58,6 +58,69 @@ describe('redirect method rewriting', () => {
expect(second.method).toBe('POST');
expect(second.body).toBeTruthy();
});

it('cancels the redirect response body before following the next hop', async () => {
const cancel = vi.fn(async () => {});
let calls = 0;
const f = vi.fn(async () => {
calls += 1;
if (calls === 1) {
return {
status: 302,
ok: false,
statusText: 'Found',
headers: new Headers({ location: 'https://example.com/next' }),
body: { cancel },
arrayBuffer: async () => new ArrayBuffer(0),
} as unknown as Response;
}
return new Response('ok', { status: 200, headers: { 'content-type': 'text/plain' } });
});
vi.stubGlobal('fetch', f);
await webFetch('https://example.com/start');
expect(cancel).toHaveBeenCalledOnce();
expect(f).toHaveBeenCalledTimes(2);
});
});

describe('redirect destination validation', () => {
it('refuses to forward credentials or a request body across origins', async () => {
const f = vi.fn(async () =>
new Response(null, { status: 302, headers: { location: 'https://evil.example/collect' } }),
);
vi.stubGlobal('fetch', f);
vi.stubEnv('TAVILY_KEY', 'secret-key');
const out = await webSearch('q', 1);
expect(out).toContain('Search failed');
expect(out).toMatch(/credentials or request body across origins/i);
// Only the first hop — never followed the cross-origin Location.
expect(f).toHaveBeenCalledTimes(1);
});

it('refuses a non-http(s) redirect Location before following', async () => {
const f = vi.fn(async () =>
new Response(null, { status: 302, headers: { location: 'file:///etc/passwd' } }),
);
vi.stubGlobal('fetch', f);
const out = await webFetch('https://example.com/start');
expect(out).toMatch(/Refusing redirect to non-http/i);
expect(f).toHaveBeenCalledTimes(1);
});

it('allows a same-origin redirect that carries a request body', async () => {
let calls = 0;
const f = vi.fn(async () => {
calls += 1;
return calls === 1
? new Response(null, { status: 307, headers: { location: 'https://api.tavily.com/next' } })
: new Response(JSON.stringify({ results: [] }), { status: 200 });
});
vi.stubGlobal('fetch', f);
vi.stubEnv('TAVILY_KEY', 'k');
const out = await webSearch('q', 1);
expect(out).toContain('No results');
expect(f).toHaveBeenCalledTimes(2);
});
});

describe('webFetch', () => {
Expand Down
25 changes: 23 additions & 2 deletions src/adapters/webTools.ts
Original file line number Diff line number Diff line change
Expand Up @@ -128,8 +128,29 @@ async function fetchWithTimeout(
if ([301, 302, 303, 307, 308].includes(response.status)) {
const location = response.headers.get('location');
if (!location) throw new Error('Redirect response has no location');
await response.body?.cancel();
const next = new URL(location, current);
// Release the redirect body before the next hop so undici does not keep
// the prior socket/buffer alive across a chain of Location responses.
if (response.body) {
try {
await response.body.cancel();
} catch {
// A body that cannot be cancelled (already locked/consumed) must not
// abort the hop; draining it releases the socket the same way.
await response.arrayBuffer().catch(() => undefined);
}
}
let next: URL;
try {
next = new URL(location, current);
} catch {
throw new Error(`Invalid redirect location: ${location}`);
}
// Validate the hop before following: only http(s) destinations are
// eligible (publicFetch would also reject, but fail closed here so a
// credentialed request never even attempts a file:/javascript: target).
if (next.protocol !== 'http:' && next.protocol !== 'https:') {
throw new Error(`Refusing redirect to non-http(s) URL (${next.protocol})`);
}
if (carriesSensitiveRequestData && next.origin !== initialOrigin) {
throw new Error('Refusing to forward credentials or request body across origins');
}
Expand Down
1 change: 1 addition & 0 deletions src/auth/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ export {
runOAuthPkceFlow,
loginAndSaveProfile,
DEFAULT_OPENAI_CLIENT_ID,
isLoopbackRemote,
type OAuthFlowResult,
type OAuthFlowOptions,
} from './oauthPkce.js';
Expand Down
27 changes: 27 additions & 0 deletions src/auth/oauthPkce.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,27 @@
import { describe, expect, it } from 'vitest';
import { isLoopbackRemote } from './oauthPkce.js';

/**
* The OpenAI redirect_uri is always the advertised localhost form. Callbacks may
* still arrive over IPv4 127.0.0.1 or IPv6 ::1, and Node reports an IPv4 client
* on a dual-stack socket in the IPv4-mapped form.
*/
describe('isLoopbackRemote (OpenAI PKCE callback)', () => {
it('accepts the IPv4 and IPv6 loopback remotes Node reports', () => {
expect(isLoopbackRemote('127.0.0.1')).toBe(true);
expect(isLoopbackRemote('::1')).toBe(true);
expect(isLoopbackRemote('::ffff:127.0.0.1')).toBe(true);
});

it('rejects non-loopback remotes that must never drive the callback', () => {
expect(isLoopbackRemote(undefined)).toBe(false);
expect(isLoopbackRemote('')).toBe(false);
expect(isLoopbackRemote('10.0.0.1')).toBe(false);
expect(isLoopbackRemote('192.168.1.1')).toBe(false);
expect(isLoopbackRemote('8.8.8.8')).toBe(false);
expect(isLoopbackRemote('fe80::1')).toBe(false);
expect(isLoopbackRemote('2001:db8::1')).toBe(false);
// A mapped non-loopback address is still a remote host.
expect(isLoopbackRemote('::ffff:10.0.0.1')).toBe(false);
});
});
31 changes: 31 additions & 0 deletions src/auth/oauthPkce.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@
// ============================================

import { createServer, type IncomingMessage, type ServerResponse } from 'node:http';
import { isIP } from 'node:net';
import { randomBytes, createHash } from 'node:crypto';
import { AuthProfileStore, type AuthProfile } from './oauthStore.js';
import { openBrowser } from './openBrowser.js';
Expand All @@ -26,6 +27,26 @@ const PROFILE_KEY = 'openai-gpt:default';
export const DEFAULT_OPENAI_CLIENT_ID = 'app_EMoamEEZ73f0CkXaXp7hrann';
const OAUTH_ORIGINATOR = 'openswarm';

/**
* Accept only the advertised loopback callback forms: IPv4 127.0.0.1 and IPv6
* ::1 (including the IPv4-mapped ::ffff:127.0.0.1 form Node reports when a
* dual-stack socket receives an IPv4 connection). Any other remote address is
* rejected so the callback server cannot be driven from off-host.
*
* The OpenAI redirect_uri stays `http://localhost:<port>/auth/callback` (the
* registered client value); browsers may still connect via ::1 or 127.0.0.1.
*/
export function isLoopbackRemote(address: string | undefined): boolean {
if (!address) return false;
const family = isIP(address);
if (family === 4) return address === '127.0.0.1';
if (family === 6) {
const lower = address.toLowerCase();
return lower === '::1' || lower === '::ffff:127.0.0.1';
}
return false;
}

// PKCE helpers

function generateCodeVerifier(): string {
Expand Down Expand Up @@ -113,6 +134,12 @@ export async function runOAuthPkceFlow(options: OAuthFlowOptions = {}): Promise<
return;
}

if (!isLoopbackRemote(req.socket.remoteAddress)) {
res.writeHead(403);
res.end('Forbidden');
return;
}

const url = new URL(req.url ?? '/', `http://127.0.0.1:${port}`);

if (url.pathname !== '/auth/callback') {
Expand Down Expand Up @@ -236,6 +263,10 @@ export async function runOAuthPkceFlow(options: OAuthFlowOptions = {}): Promise<
}
});

// Bound to the IPv4 loopback, not 'localhost': listen('localhost') resolves
// the name ONCE and binds whichever family wins (measured on macOS: ::1
// only, so 127.0.0.1 clients get ECONNREFUSED). isLoopbackRemote still
// accepts the IPv6 and IPv4-mapped forms so a dual-stack bind stays correct.
server.listen(port, '127.0.0.1', () => {
console.log(`[Auth] Callback server listening on http://127.0.0.1:${port}`);
console.log(`[Auth] 브라우저에서 OpenAI 로그인 페이지를 엽니다...`);
Expand Down
25 changes: 25 additions & 0 deletions src/cli/mcpCommand.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,31 @@ describe('parseServerSpec (INT-1953)', () => {
it('throws when nothing usable is given', () => {
expect(() => parseServerSpec(undefined, [])).toThrow();
});
it('rejects an unknown preset name before persistence (AGT-3442)', () => {
expect(() => parseServerSpec(undefined, [], 'not-a-real-preset')).toThrow(/unknown preset/);
});
it('rejects a malformed http(s) URL before persistence (AGT-3442)', () => {
expect(() => parseServerSpec('https://exa mple.com/mcp', [])).toThrow(/invalid server URL/);
});
it('rejects an empty server name before persistence (AGT-3442)', () => {
expect(() => addServer({ mcpServers: {} }, ' ', { command: 'x' })).toThrow(/non-empty server name/);
});
it('does not write the registry when add validation fails (AGT-3442)', () => {
const dir = mkdtempSync(join(tmpdir(), 'mcpcmd-'));
const path = join(dir, 'mcp.json');
try {
expect(() =>
runMcpCommand('add', 'bad', ['https://exa mple.com/mcp'], { path }),
).toThrow(/invalid server URL/);
expect(existsSync(path)).toBe(false);
expect(() =>
runMcpCommand('add', 'bad', [], { preset: 'not-a-real-preset', path }),
).toThrow(/unknown preset/);
expect(existsSync(path)).toBe(false);
} finally {
rmSync(dir, { recursive: true, force: true });
}
});
});

describe('addServer / removeServer / formatServerList', () => {
Expand Down
Loading
Loading