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
11 changes: 10 additions & 1 deletion server/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,16 @@ export function createApp(): express.Express {
app.use(helmet());
app.use(
cors({
exposedHeaders: ['X-Log-Count', 'Mcp-Session-Id', 'mcp-session-id'],
exposedHeaders: [
'X-Log-Count',
'Mcp-Session-Id',
'mcp-session-id',
'Retry-After',
'retry-after-ms',
'X-RateLimit-Remaining',
'X-RateLimit-Reset',
'x-ratelimit-remaining',
],
allowedHeaders: [
'Content-Type',
'Authorization',
Expand Down
44 changes: 43 additions & 1 deletion server/src/routes/proxy.ts
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,43 @@ function getProxyConfig(): ProxyConfig | null {

const router = Router();

// Helper: build API URL handling baseUrl already ending in version prefix
// 透传上游与限流相关的响应头,供前端(含直连模式)统一识别 Retry-After 与剩余配额。
const UPSTREAM_RATE_LIMIT_HEADERS = new Set([
'retry-after',
'retry-after-ms',
'x-ratelimit-remaining-requests',
'x-ratelimit-remaining-tokens',
'x-ratelimit-reset-requests',
'x-ratelimit-reset-tokens',
'x-ratelimit-limit-requests',
'x-ratelimit-limit-tokens',
'x-ratelimit-remaining',
'x-ratelimit-reset',
'x-ratelimit-limit',
'ratelimit-remaining',
'ratelimit-reset',
'ratelimit-limit',
'anthropic-ratelimit-requests-remaining',
'anthropic-ratelimit-requests-reset',
'anthropic-ratelimit-input-tokens-remaining',
'anthropic-ratelimit-input-tokens-reset',
'anthropic-ratelimit-output-tokens-remaining',
'anthropic-ratelimit-output-tokens-reset',
'x-should-retry',
]);

function relayRateLimitHeaders(res: import('express').Response, headers: Record<string, string> | undefined): void {
if (!headers) return;
for (const [key, value] of Object.entries(headers)) {
if (UPSTREAM_RATE_LIMIT_HEADERS.has(key.toLowerCase()) && typeof value === 'string' && value) {
try {
res.setHeader(key, value);
} catch { /* ignore invalid header name */ }
}
}
}

// Helper: build API URL handling base already ending in version prefix
function buildApiUrl(baseUrl: string, pathWithVersion: string): string {
const baseUrlWithSlash = baseUrl.endsWith('/') ? baseUrl : `${baseUrl}/`;
const versionPrefix = pathWithVersion.split('/')[0] || '';
Expand Down Expand Up @@ -104,6 +140,7 @@ router.post('/api/proxy/github/*', async (req, res) => {

const proxyConfig = getProxyConfig();
const result = await proxyRequest({ url: targetUrl, method, headers, body: body.body, proxyConfig });
relayRateLimitHeaders(res, result.headers);
res.status(result.status).json(result.data);
} catch (err) {
logger.errorFromError('proxy.github', 'GitHub proxy error', err);
Expand Down Expand Up @@ -175,6 +212,7 @@ router.post('/api/proxy/github-raw', async (req, res) => {

// Raw content is text/plain, forward as-is (not JSON-wrapped)
const contentType = String(result.headers['content-type'] || 'text/plain');
relayRateLimitHeaders(res, result.headers);
res.status(result.status).type(contentType).send(typeof result.data === 'string' ? result.data : JSON.stringify(result.data));
} catch (err) {
logger.errorFromError('proxy.github-raw', 'GitHub raw proxy error', err);
Expand Down Expand Up @@ -302,6 +340,7 @@ router.post('/api/proxy/ai', async (req, res) => {
allowPrivate,
});

relayRateLimitHeaders(res, result.headers);
res.status(result.status).json(result.data);
} catch (err) {
logger.errorFromError('proxy.ai', 'AI proxy error', err);
Expand Down Expand Up @@ -366,6 +405,7 @@ router.post('/api/proxy/webdav', async (req, res) => {
allowPrivate: true,
});

relayRateLimitHeaders(res, result.headers);
res.status(result.status).json(result.data);
} catch (err) {
logger.errorFromError('proxy.webdav', 'WebDAV proxy error', err);
Expand Down Expand Up @@ -406,6 +446,7 @@ router.post('/api/proxy/github/search/repositories', async (req, res) => {

const proxyConfig = getProxyConfig();
const result = await proxyRequest({ url: targetUrl, method: 'GET', headers, proxyConfig });
relayRateLimitHeaders(res, result.headers);
res.status(result.status).json(result.data);
} catch (err) {
logger.errorFromError('proxy.github.search', 'GitHub search repositories proxy error', err);
Expand Down Expand Up @@ -446,6 +487,7 @@ router.post('/api/proxy/github/search/users', async (req, res) => {

const proxyConfig = getProxyConfig();
const result = await proxyRequest({ url: targetUrl, method: 'GET', headers, proxyConfig });
relayRateLimitHeaders(res, result.headers);
res.status(result.status).json(result.data);
} catch (err) {
logger.errorFromError('proxy.github.search', 'GitHub search users proxy error', err);
Expand Down
30 changes: 30 additions & 0 deletions server/tests/routes/githubProxyRoute.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -139,4 +139,34 @@ describe('GitHub proxy routes', () => {
expect(options.headers['Content-Length']).toBeUndefined();
expect(options.headers['X-Custom']).toBe('kept');
});

it('relays upstream rate-limit headers to the client and skips unrelated ones', async () => {
const app = createTestApp();
proxyRequestMock.mockResolvedValueOnce({
status: 429,
data: { error: { message: 'rate limited' } },
headers: {
'content-type': 'application/json',
'retry-after-ms': '2500',
'retry-after': '5',
'x-ratelimit-remaining-requests': '12',
'x-ratelimit-remaining-tokens': '900',
'x-debug-id': 'secret-upstream-id',
'set-cookie': 'sid=abc',
},
});

const res = await request(app)
.post('/api/proxy/github/gists/abc123')
.send({ method: 'GET' })
.expect(429);

expect(res.headers['retry-after-ms']).toBe('2500');
expect(res.headers['retry-after']).toBe('5');
expect(res.headers['x-ratelimit-remaining-requests']).toBe('12');
expect(res.headers['x-ratelimit-remaining-tokens']).toBe('900');
// 非限流相关头不应透传
expect(res.headers['x-debug-id']).toBeUndefined();
expect(res.headers['set-cookie']).toBeUndefined();
});
});
4 changes: 4 additions & 0 deletions src/components/DiscoveryView.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -793,6 +793,10 @@ export const DiscoveryView: React.FC = React.memo(() => {
const aiService = new AIService(activeConfig, language);
const optimizer = new AIAnalysisOptimizer({
initialConcurrency: activeConfig.concurrency || 3,
rateLimiter: {
maxConcurrency: 0,
requestsPerMinute: activeConfig.requestsPerMinute || 0,
},
});
setAnalysisOptimizer(optimizer);

Expand Down
8 changes: 8 additions & 0 deletions src/components/RepositoryList.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -322,6 +322,10 @@ export const RepositoryList: React.FC<RepositoryListProps> = ({
maxRetries: 3,
retryDelayBaseMs: 1000,
enableAdaptiveConcurrency: true,
rateLimiter: {
maxConcurrency: 0,
requestsPerMinute: activeConfig.requestsPerMinute || 0,
},
});

try {
Expand Down Expand Up @@ -675,6 +679,10 @@ export const RepositoryList: React.FC<RepositoryListProps> = ({
maxRetries: 3,
retryDelayBaseMs: 1000,
enableAdaptiveConcurrency: true,
rateLimiter: {
maxConcurrency: 0,
requestsPerMinute: activeConfig.requestsPerMinute || 0,
},
});

try {
Expand Down
188 changes: 188 additions & 0 deletions src/services/aiAnalysisOptimizer.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,188 @@
import { describe, it, expect, vi } from 'vitest';
import { AIAnalysisOptimizer, AnalysisTask } from './aiAnalysisOptimizer';
import type { AIService } from './aiService';
import type { GitHubApiService } from './githubApi';
import { AIRequestError } from './aiService';
import type { Repository } from '../types';

function makeRepo(id: number): Repository {
return {
id,
name: `repo-${id}`,
full_name: `acme/repo-${id}`,
description: null,
html_url: `https://github.com/acme/repo-${id}`,
stargazers_count: 0,
forks_count: 0,
forks: 0,
language: 'TypeScript',
created_at: '2024-01-01T00:00:00Z',
updated_at: '2024-06-01T00:00:00Z',
pushed_at: '2024-06-01T00:00:00Z',
owner: { login: 'owner', avatar_url: '' },
topics: [],
};
}

function makeTask(id: number): AnalysisTask {
return { repo: makeRepo(id), readmeContent: '# readme', retries: 0 };
}

describe('AIAnalysisOptimizer 与共享限流器集成', () => {
it('遭遇 429 时记录到限流器、等待后重试成功,并清零计数', async () => {
const optimizer = new AIAnalysisOptimizer({
maxRetries: 2,
retryDelayBaseMs: 1,
rateLimiter: {
maxConcurrency: 0,
requestsPerMinute: 0,
cooldownThreshold: 3,
backoffBaseMs: 10,
backoffCapMs: 120,
maxRetryAfterMs: 50,
},
});

const analyze = vi.fn()
.mockRejectedValueOnce(new AIRequestError('rate limited', 429, 200000))
.mockResolvedValueOnce({ summary: 'ok', tags: [], platforms: [] });
const fakeAi = { analyzeRepository: analyze } as unknown as AIService;

const result = await optimizer.analyzeWithRetry(makeTask(1), fakeAi, []);

expect(result.success).toBe(true);
expect(analyze).toHaveBeenCalledTimes(2);
// 成功后将连续 429 计数复位
expect(optimizer.limiter.getStatus().consecutiveRateLimits).toBe(0);
});

it('连续 429 超过阈值触发熔断,重试耗尽后返回失败', async () => {
const optimizer = new AIAnalysisOptimizer({
maxRetries: 2,
retryDelayBaseMs: 1,
rateLimiter: {
maxConcurrency: 0,
requestsPerMinute: 0,
cooldownThreshold: 2,
backoffBaseMs: 5,
backoffCapMs: 50,
maxRetryAfterMs: 30,
},
});

const analyze = vi.fn().mockRejectedValue(new AIRequestError('too many requests', 429));
const fakeAi = { analyzeRepository: analyze } as unknown as AIService;

const result = await optimizer.analyzeWithRetry(makeTask(2), fakeAi, []);

expect(result.success).toBe(false);
expect(result.error).toBeInstanceOf(AIRequestError);
// 3 次调用 >= 阈值 2 => 熔断打开
expect(optimizer.limiter.getStatus().circuitOpen).toBe(true);
expect(optimizer.limiter.getStatus().consecutiveRateLimits).toBe(3);
});

it('非限流错误不受 RateLimiter 冷却影响', async () => {
const optimizer = new AIAnalysisOptimizer({
maxRetries: 1,
rateLimiter: { maxConcurrency: 0, requestsPerMinute: 0, cooldownThreshold: 2 },
});

const analyze = vi.fn()
.mockRejectedValueOnce(new Error('some server error'))
.mockResolvedValueOnce({ summary: 'ok', tags: [], platforms: [] });
const fakeAi = { analyzeRepository: analyze } as unknown as AIService;

const result = await optimizer.analyzeWithRetry(makeTask(3), fakeAi, []);

expect(result.success).toBe(true);
expect(analyze).toHaveBeenCalledTimes(2);
expect(optimizer.limiter.getStatus().consecutiveRateLimits).toBe(0);
});

it('abort() 立即停止阻塞在限流等待与重试延迟中的 worker', async () => {
const optimizer = new AIAnalysisOptimizer({
maxRetries: 3,
retryDelayBaseMs: 10000,
rateLimiter: {
maxConcurrency: 0,
requestsPerMinute: 0,
cooldownThreshold: 1,
backoffBaseMs: 10000,
backoffCapMs: 10000,
maxRetryAfterMs: 10000,
},
});

const analyze = vi.fn().mockRejectedValue(new AIRequestError('too many requests', 429));
const fakeAi = { analyzeRepository: analyze } as unknown as AIService;

// 第一个任务触发 429 后陷入长重试延迟;第二个任务阻塞在限流器 acquire
const first = optimizer.analyzeWithRetry(makeTask(4), fakeAi, []);
const second = optimizer.analyzeWithRetry(makeTask(5), fakeAi, []);

// 给两个任务时间进入阻塞状态,然后中止整个批次
await new Promise(r => setTimeout(r, 100));
const start = Date.now();
optimizer.abort();

const [r1, r2] = await Promise.all([first, second]);
const elapsed = Date.now() - start;

// abort 应立刻穿过冷却等待 / 重试延迟(总时长被设为 10s)
expect(elapsed).toBeLessThan(500);
expect(r1.success).toBe(false);
expect(r1.error?.message).toBe('Analysis aborted');
expect(r2.success).toBe(false);
expect(r2.error?.message).toBe('Analysis aborted');
});

it('abort() 立即结束流水线,即使 README 请求仍在飞行', async () => {
const optimizer = new AIAnalysisOptimizer({
maxRetries: 0,
rateLimiter: { maxConcurrency: 0, requestsPerMinute: 0 },
});

const analyze = vi.fn().mockResolvedValue({ summary: 'ok', tags: [], platforms: [] });
const fakeAi = { analyzeRepository: analyze } as unknown as AIService;

// README 拉取在收到 batch 信号中止前一直挂起:验证信号确实被传到了请求 API
let resolveReadmeStarted = () => {};
const readmeStarted = new Promise<void>(resolve => {
resolveReadmeStarted = resolve;
});
const readmeAborted = vi.fn();
const githubApi = {
getRepositoryReadme: vi.fn((_owner: string, _repo: string, signal?: AbortSignal) => {
resolveReadmeStarted();
return new Promise<string>((_resolve, reject) => {
const onAbort = () => {
readmeAborted();
reject(new Error('Aborted'));
};
if (!signal) {
reject(new Error('Missing AbortSignal'));
} else if (signal.aborted) {
onAbort();
} else {
signal.addEventListener('abort', onAbort, { once: true });
}
});
}),
} as unknown as GitHubApiService;

const pending = optimizer.analyzeRepositoriesPipelined(
[makeRepo(6), makeRepo(7)], githubApi, fakeAi, []);

// 等 README 请求真正启动后再中止整个批次,避免固定延时与请求启动竞速
await readmeStarted;
const start = Date.now();
optimizer.abort();
await pending;

expect(Date.now() - start).toBeLessThan(500);
expect(githubApi.getRepositoryReadme).toHaveBeenCalled();
// 断言 README 请求确实收到了中止信号,而非流水线提前返回放任其挂起
expect(readmeAborted).toHaveBeenCalled();
});
});
Loading