mirror of
https://github.com/tiennm99/ccs.git
synced 2026-10-11 03:13:12 +00:00
fix(proxy): keep undici timeouts above the upstream request timeout (#1524)
Sets undici headersTimeout/bodyTimeout to request_timeout+30s so the AbortController is the single authority on upstream request lifetime, preventing premature socket closes on slow self-hosted upstreams. Verified undici v5 ProxyAgent object-signature compat.
This commit is contained in:
1 parent
0d856fc0ab
commit
8f9795bce2
5 files changed
+116
-13
No files matched your search
@@ -327,6 +327,14 @@ That flag is respected by both:
|
||||
- CCS now preserves upstream rate-limit errors and retry headers
|
||||
- Empty or malformed provider JSON is returned as Anthropic-style `api_error`
|
||||
|
||||
### Slow upstreams: `socket connection was closed unexpectedly`
|
||||
|
||||
- Long-running upstreams (self-hosted LLMs with queue/prefill phases) can stay
|
||||
silent for minutes before the first or next token. The proxy already allows up
|
||||
to 10 minutes per request; set `CCS_OPENAI_PROXY_REQUEST_TIMEOUT_MS` in the
|
||||
profile settings to raise or lower that ceiling.
|
||||
- Restart the proxy after changing the setting.
|
||||
|
||||
### Requests route to the wrong model/profile
|
||||
|
||||
- Use an explicit selector such as `profile:model`
|
||||
|
||||
@@ -10,9 +10,16 @@ import {
|
||||
import { ProxySseStreamTransformer } from '../transformers/sse-stream-transformer';
|
||||
import { isAnthropicPassthroughProfile, resolveOpenAIChatCompletionsUrl } from '../upstream-url';
|
||||
import { createLogger } from '../../services/logging';
|
||||
import {
|
||||
createGlobalFetchProxyDispatcher,
|
||||
type UpstreamAgentTimeoutOptions,
|
||||
} from '../../utils/fetch-proxy-setup';
|
||||
import { pipeWebResponseToNode, readRawBody, writeJson } from './http-helpers';
|
||||
|
||||
const REQUEST_TIMEOUT_MS = 600_000;
|
||||
// Keep undici's per-phase timeouts above the explicit request timeout so the
|
||||
// AbortController is the single authority on when an upstream request dies.
|
||||
const UPSTREAM_TIMEOUT_GRACE_MS = 30_000;
|
||||
const DIRECT_OPENAI_REASONING_CHAT_MODEL = /^(?:gpt-5|o[134])(?:[-.]|$)/;
|
||||
const logger = createLogger('proxy:openai-compat:messages');
|
||||
|
||||
@@ -375,7 +382,7 @@ function buildFetchInit(
|
||||
signal: AbortSignal,
|
||||
incomingHeaders: http.IncomingHttpHeaders,
|
||||
preserveUserAgent: boolean,
|
||||
insecureDispatcher?: Dispatcher
|
||||
dispatcher?: Dispatcher
|
||||
): RequestInit {
|
||||
const init: RequestInit = {
|
||||
method: 'POST',
|
||||
@@ -384,8 +391,8 @@ function buildFetchInit(
|
||||
signal,
|
||||
};
|
||||
|
||||
if (insecureDispatcher) {
|
||||
(init as Record<string, unknown>).dispatcher = insecureDispatcher;
|
||||
if (dispatcher) {
|
||||
(init as Record<string, unknown>).dispatcher = dispatcher;
|
||||
}
|
||||
|
||||
return init;
|
||||
@@ -401,6 +408,30 @@ function getRequestTimeoutMs(): number {
|
||||
return Number.isFinite(parsed) && parsed > 0 ? parsed : REQUEST_TIMEOUT_MS;
|
||||
}
|
||||
|
||||
/**
|
||||
* undici defaults `headersTimeout`/`bodyTimeout` to 300s, which silently
|
||||
* undercuts {@link REQUEST_TIMEOUT_MS} (600s) and the
|
||||
* `CCS_OPENAI_PROXY_REQUEST_TIMEOUT_MS` override: slow upstreams (self-hosted
|
||||
* LLMs with long queue + prefill phases) get their socket closed at 300s with
|
||||
* a generic connection error instead of the proxy's timeout response. Every
|
||||
* dispatcher used for upstream fetches must carry these options.
|
||||
*/
|
||||
export function buildUpstreamAgentTimeouts(): UpstreamAgentTimeoutOptions {
|
||||
const ceiling = getRequestTimeoutMs() + UPSTREAM_TIMEOUT_GRACE_MS;
|
||||
return { headersTimeout: ceiling, bodyTimeout: ceiling };
|
||||
}
|
||||
|
||||
let defaultUpstreamDispatcher: Dispatcher | null = null;
|
||||
|
||||
function getDefaultUpstreamDispatcher(): Dispatcher {
|
||||
if (!defaultUpstreamDispatcher) {
|
||||
const timeouts = buildUpstreamAgentTimeouts();
|
||||
// Honor HTTP(S)_PROXY routing when configured; otherwise a plain Agent.
|
||||
defaultUpstreamDispatcher = createGlobalFetchProxyDispatcher(timeouts) ?? new Agent(timeouts);
|
||||
}
|
||||
return defaultUpstreamDispatcher;
|
||||
}
|
||||
|
||||
function formatTimeoutDuration(timeoutMs: number): string {
|
||||
return timeoutMs % 1000 === 0 ? `${timeoutMs / 1000} seconds` : `${timeoutMs}ms`;
|
||||
}
|
||||
@@ -538,19 +569,20 @@ export async function handleProxyMessagesRequest(
|
||||
const useProfileInsecureTls = upstream.route.profile.insecure === true;
|
||||
const ephemeralInsecureDispatcher =
|
||||
useProfileInsecureTls && !useSharedInsecureDispatcher
|
||||
? new Agent({ connect: { rejectUnauthorized: false } })
|
||||
? new Agent({ connect: { rejectUnauthorized: false }, ...buildUpstreamAgentTimeouts() })
|
||||
: undefined;
|
||||
const insecureTls = useSharedInsecureDispatcher || useProfileInsecureTls;
|
||||
const dispatcher = useSharedInsecureDispatcher
|
||||
? insecureDispatcher
|
||||
: useProfileInsecureTls
|
||||
? ephemeralInsecureDispatcher
|
||||
: undefined;
|
||||
: getDefaultUpstreamDispatcher();
|
||||
|
||||
try {
|
||||
logger.stage('dispatch', 'upstream.dispatch', 'Dispatching upstream fetch', {
|
||||
profileName: profile.profileName,
|
||||
routedProfileName: upstream.route.profile.profileName,
|
||||
insecureTls: dispatcher !== undefined,
|
||||
insecureTls,
|
||||
passthrough,
|
||||
});
|
||||
const upstreamUrl = resolveOpenAIChatCompletionsUrl(upstream.route.profile.baseUrl, {
|
||||
|
||||
@@ -5,6 +5,7 @@ import type { OpenAICompatProfileConfig } from '../profile-router';
|
||||
import { OPENAI_COMPAT_PROXY_SERVICE_NAME } from '../proxy-daemon-paths';
|
||||
import { createLogger, withRequestContext } from '../../services/logging';
|
||||
import {
|
||||
buildUpstreamAgentTimeouts,
|
||||
handleProxyMessagesRequest,
|
||||
handleProxyModelsRequest,
|
||||
validateIncomingProxyAuth,
|
||||
@@ -39,7 +40,7 @@ export function startOpenAICompatProxyServer(options: OpenAICompatProxyServerOpt
|
||||
port: options.port,
|
||||
});
|
||||
const insecureDispatcher = options.insecure
|
||||
? new Agent({ connect: { rejectUnauthorized: false } })
|
||||
? new Agent({ connect: { rejectUnauthorized: false }, ...buildUpstreamAgentTimeouts() })
|
||||
: undefined;
|
||||
const server = http.createServer((req, res) => {
|
||||
const requestId = resolveInboundRequestId(req.headers);
|
||||
|
||||
@@ -12,6 +12,8 @@ const FETCH_PROXY_PROTOCOLS = ['http:', 'https:'];
|
||||
type RoutingDispatchOptions = Parameters<Dispatcher['dispatch']>[0];
|
||||
type RoutingDispatchHandler = Parameters<Dispatcher['dispatch']>[1];
|
||||
|
||||
export type UpstreamAgentTimeoutOptions = Pick<Agent.Options, 'headersTimeout' | 'bodyTimeout'>;
|
||||
|
||||
type GlobalFetchProxyConfig = {
|
||||
httpProxyUrl?: string;
|
||||
httpsProxyUrl?: string;
|
||||
@@ -19,14 +21,23 @@ type GlobalFetchProxyConfig = {
|
||||
};
|
||||
|
||||
class RoutingProxyDispatcher extends Dispatcher {
|
||||
private readonly directDispatcher = new Agent();
|
||||
private readonly directDispatcher: Agent;
|
||||
private readonly httpProxyDispatcher: ProxyAgent | null;
|
||||
private readonly httpsProxyDispatcher: ProxyAgent | null;
|
||||
|
||||
constructor(httpProxyUrl: string | undefined, httpsProxyUrl: string | undefined) {
|
||||
constructor(
|
||||
httpProxyUrl: string | undefined,
|
||||
httpsProxyUrl: string | undefined,
|
||||
agentOptions: UpstreamAgentTimeoutOptions = {}
|
||||
) {
|
||||
super();
|
||||
this.httpProxyDispatcher = httpProxyUrl ? new ProxyAgent(httpProxyUrl) : null;
|
||||
this.httpsProxyDispatcher = httpsProxyUrl ? new ProxyAgent(httpsProxyUrl) : null;
|
||||
this.directDispatcher = new Agent(agentOptions);
|
||||
this.httpProxyDispatcher = httpProxyUrl
|
||||
? new ProxyAgent({ uri: httpProxyUrl, ...agentOptions })
|
||||
: null;
|
||||
this.httpsProxyDispatcher = httpsProxyUrl
|
||||
? new ProxyAgent({ uri: httpsProxyUrl, ...agentOptions })
|
||||
: null;
|
||||
}
|
||||
|
||||
dispatch(options: RoutingDispatchOptions, handler: RoutingDispatchHandler): boolean {
|
||||
@@ -114,14 +125,16 @@ class RoutingProxyDispatcher extends Dispatcher {
|
||||
}
|
||||
}
|
||||
|
||||
export function createGlobalFetchProxyDispatcher(): Dispatcher | null {
|
||||
export function createGlobalFetchProxyDispatcher(
|
||||
agentOptions: UpstreamAgentTimeoutOptions = {}
|
||||
): Dispatcher | null {
|
||||
const { httpProxyUrl, httpsProxyUrl } = resolveGlobalFetchProxyConfig();
|
||||
|
||||
if (!httpProxyUrl && !httpsProxyUrl) {
|
||||
return null;
|
||||
}
|
||||
|
||||
return new RoutingProxyDispatcher(httpProxyUrl, httpsProxyUrl);
|
||||
return new RoutingProxyDispatcher(httpProxyUrl, httpsProxyUrl, agentOptions);
|
||||
}
|
||||
|
||||
export function applyGlobalFetchProxy(): { enabled: boolean; error?: string } {
|
||||
|
||||
@@ -8,6 +8,7 @@ import { Agent } from 'undici';
|
||||
import { resolveOpenAICompatProfileConfig } from '../../../src/proxy/profile-router';
|
||||
import {
|
||||
attachDisconnectAbortHandlers,
|
||||
buildUpstreamAgentTimeouts,
|
||||
handleProxyMessagesRequest,
|
||||
} from '../../../src/proxy/server/messages-route';
|
||||
import { loadSettings } from '../../../src/utils/config-manager';
|
||||
@@ -171,7 +172,55 @@ describe('attachDisconnectAbortHandlers', () => {
|
||||
});
|
||||
});
|
||||
|
||||
describe('buildUpstreamAgentTimeouts', () => {
|
||||
afterEach(() => {
|
||||
delete process.env.CCS_OPENAI_PROXY_REQUEST_TIMEOUT_MS;
|
||||
});
|
||||
|
||||
it('keeps undici per-phase timeouts above the default request timeout', () => {
|
||||
const timeouts = buildUpstreamAgentTimeouts();
|
||||
expect(timeouts.headersTimeout).toBeGreaterThan(600_000);
|
||||
expect(timeouts.bodyTimeout).toBeGreaterThan(600_000);
|
||||
});
|
||||
|
||||
it('tracks CCS_OPENAI_PROXY_REQUEST_TIMEOUT_MS overrides', () => {
|
||||
process.env.CCS_OPENAI_PROXY_REQUEST_TIMEOUT_MS = '120000';
|
||||
const timeouts = buildUpstreamAgentTimeouts();
|
||||
expect(timeouts.headersTimeout).toBeGreaterThan(120_000);
|
||||
expect(timeouts.bodyTimeout).toBeGreaterThan(120_000);
|
||||
});
|
||||
});
|
||||
|
||||
describe('handleProxyMessagesRequest', () => {
|
||||
it('dispatches secure upstream fetches with an explicit dispatcher (not undici defaults)', async () => {
|
||||
const activeProfile = buildProfile('hf');
|
||||
let capturedDispatcher: unknown;
|
||||
|
||||
globalThis.fetch = (async (_input: RequestInfo | URL, init?: RequestInit) => {
|
||||
capturedDispatcher = (init as RequestInit & { dispatcher?: unknown })?.dispatcher;
|
||||
throw new Error('upstream exploded');
|
||||
}) as typeof globalThis.fetch;
|
||||
|
||||
const req = new FakeRequest({
|
||||
'x-api-key': 'local-token',
|
||||
});
|
||||
const res = new FakeResponse();
|
||||
const pending = handleProxyMessagesRequest(req as never, res as never, activeProfile, 'local-token');
|
||||
req.end(
|
||||
JSON.stringify({
|
||||
model: 'hf-default',
|
||||
stream: true,
|
||||
messages: [{ role: 'user', content: 'stay secure' }],
|
||||
})
|
||||
);
|
||||
await pending;
|
||||
|
||||
// A bare global-dispatcher fallback would reapply undici's 300s
|
||||
// headersTimeout/bodyTimeout and undercut the proxy's request timeout.
|
||||
expect(capturedDispatcher).toBeDefined();
|
||||
expect(res.statusCode).toBe(502);
|
||||
});
|
||||
|
||||
it('auto-passes through Kimi requests and preserves the coding-agent User-Agent', async () => {
|
||||
const activeProfile = buildProfile('kimic');
|
||||
let capturedInput: RequestInfo | URL | undefined;
|
||||
|
||||
Reference in new issue
Block a user