Files
Genarrative/apps/ai-game-creator-shell/scripts/llm-transient-fault-proxy.mjs
AIGameCreator App 58dad6e96d 补齐Swarm最终回复重试真实门禁
新增Supervisor最终回复瞬时故障、持久退避与Runner强杀真实验收入口
强化故障代理异步选择器、停止竞态和请求隐私回归
修复并行Agent文件写改删的短等待写锁与绝对路径脱敏
补充终态sidecar等待、并发写锁测试和六轮真实Provider证据
同步Runtime方案、实施计划与项目共享记忆文档
2026-07-20 05:50:24 +08:00

622 lines
18 KiB
JavaScript

import http from 'node:http';
import https from 'node:https';
const LOOPBACK_HOST = '127.0.0.1';
const LOOPBACK_NO_PROXY_ENTRIES = Object.freeze([
LOOPBACK_HOST,
'localhost',
'::1',
]);
const DEFAULT_FALLBACK_PORTS = Object.freeze(
Array.from({ length: 128 }, (_, index) => 62_000 + index),
);
export function withLoopbackNoProxy(environment) {
const entries = [environment.NO_PROXY, environment.no_proxy]
.flatMap((value) => (typeof value === 'string' ? value.split(',') : []))
.map((value) => value.trim())
.filter(Boolean);
const noProxy = [...new Set([...entries, ...LOOPBACK_NO_PROXY_ENTRIES])].join(
',',
);
return { ...environment, NO_PROXY: noProxy, no_proxy: noProxy };
}
function proxyError(message) {
const error = new Error(message);
error.name = 'LlmTransientFaultProxyError';
return error;
}
function parseUpstreamBaseUrl(value) {
let parsed;
try {
parsed = value instanceof URL ? new URL(value.href) : new URL(value);
} catch {
throw proxyError('upstreamBaseUrl must be a valid HTTP(S) URL');
}
if (
!['http:', 'https:'].includes(parsed.protocol) ||
parsed.username ||
parsed.password ||
parsed.hash
) {
throw proxyError(
'upstreamBaseUrl must use HTTP(S) without userinfo or a fragment',
);
}
return Object.freeze({
protocol: parsed.protocol,
hostname: parsed.hostname.replace(/^\[|\]$/gu, ''),
port: parsed.port || undefined,
hostHeader: parsed.host,
basePathname: parsed.pathname.replace(/\/+$/gu, ''),
});
}
function normalizeOptions(upstreamBaseUrlOrOptions, faultCount, extraOptions) {
const options =
typeof upstreamBaseUrlOrOptions === 'string' ||
upstreamBaseUrlOrOptions instanceof URL
? {
...extraOptions,
upstreamBaseUrl: upstreamBaseUrlOrOptions,
faultCount,
}
: upstreamBaseUrlOrOptions;
if (!options || typeof options !== 'object') {
throw proxyError('proxy options are required');
}
const normalizedFaultCount = options.faultCount ?? 1;
if (!Number.isSafeInteger(normalizedFaultCount) || normalizedFaultCount < 0) {
throw proxyError('faultCount must be a non-negative safe integer');
}
if (
options.holdAfterFault !== undefined &&
typeof options.holdAfterFault !== 'boolean'
) {
throw proxyError('holdAfterFault must be a boolean');
}
if (options.listen !== undefined && typeof options.listen !== 'function') {
throw proxyError('listen must be a function');
}
if (
options.shouldInjectFault !== undefined &&
typeof options.shouldInjectFault !== 'function'
) {
throw proxyError('shouldInjectFault must be a function');
}
const fallbackPorts = options.fallbackPorts ?? DEFAULT_FALLBACK_PORTS;
if (
!Array.isArray(fallbackPorts) ||
fallbackPorts.length === 0 ||
fallbackPorts.length > 512 ||
fallbackPorts.some(
(port) => !Number.isSafeInteger(port) || port < 49_152 || port > 65_535,
)
) {
throw proxyError('fallbackPorts must contain valid high ports');
}
return {
upstream: parseUpstreamBaseUrl(options.upstreamBaseUrl),
faultCount: normalizedFaultCount,
holdAfterFault: options.holdAfterFault ?? false,
fallbackPorts: [...new Set(fallbackPorts)],
listen: options.listen ?? listenOnLoopback,
shouldInjectFault: options.shouldInjectFault ?? null,
};
}
function listenOnLoopback(server, { host, port }) {
return new Promise((resolve, reject) => {
const cleanup = () => {
server.off('error', onError);
server.off('listening', onListening);
};
const onError = (error) => {
cleanup();
reject(error);
};
const onListening = () => {
cleanup();
resolve();
};
server.once('error', onError);
server.once('listening', onListening);
try {
server.listen({ host, port, exclusive: true });
} catch (error) {
cleanup();
reject(error);
}
});
}
function errorCode(error) {
return error && typeof error === 'object' && 'code' in error
? error.code
: undefined;
}
async function bindServer(server, listen, fallbackPorts) {
try {
await listen(server, { host: LOOPBACK_HOST, port: 0 });
return;
} catch (error) {
if (errorCode(error) !== 'EADDRINUSE') {
throw proxyError('unable to bind transient fault proxy on loopback');
}
}
for (const port of fallbackPorts) {
try {
await listen(server, { host: LOOPBACK_HOST, port });
return;
} catch (error) {
if (errorCode(error) !== 'EADDRINUSE') {
throw proxyError('unable to bind transient fault proxy on loopback');
}
}
}
throw proxyError('transient fault proxy fallback port pool is exhausted');
}
function isOriginFormPath(value) {
return (
typeof value === 'string' &&
value.startsWith('/') &&
!value.startsWith('//') &&
!hasAsciiControlCharacter(value)
);
}
function hasAsciiControlCharacter(value) {
for (const character of value) {
const codePoint = character.codePointAt(0);
if (codePoint <= 0x1f || codePoint === 0x7f) return true;
}
return false;
}
function isWithinBasePath(value, basePathname) {
let pathname;
try {
pathname = new URL(value, 'http://proxy.invalid').pathname;
} catch {
return false;
}
return (
basePathname === '' ||
pathname === basePathname ||
pathname.startsWith(`${basePathname}/`)
);
}
function resetSocket(socket) {
if (!socket || socket.destroyed) return;
try {
if (typeof socket.resetAndDestroy === 'function') {
socket.resetAndDestroy();
} else {
socket.destroy();
}
} catch {
socket.destroy();
}
}
function forwardHeaders(request, hostHeader) {
const headers = Object.create(null);
for (let index = 0; index < request.rawHeaders.length; index += 2) {
const name = request.rawHeaders[index];
const value = request.rawHeaders[index + 1];
if (!name || value === undefined) continue;
const lowerName = name.toLowerCase();
if (lowerName === 'host' || lowerName === 'proxy-connection') continue;
const existing = headers[lowerName];
if (existing === undefined) {
headers[lowerName] = value;
} else if (Array.isArray(existing)) {
existing.push(value);
} else {
headers[lowerName] = [existing, value];
}
}
headers.host = hostHeader;
return headers;
}
function failClosed(request, response, statusCode) {
const body = statusCode === 405 ? 'method not allowed' : 'request rejected';
request.resume();
response.shouldKeepAlive = false;
response.writeHead(statusCode, {
connection: 'close',
'content-length': Buffer.byteLength(body),
'content-type': 'text/plain; charset=utf-8',
});
response.end(body);
}
function closeServer(server) {
if (!server.listening) return Promise.resolve();
return new Promise((resolve) => {
server.close(() => resolve());
});
}
/**
* Starts a loopback-only proxy that resets the first configured POST requests.
* The object form also accepts holdAfterFault, fallbackPorts, and a test listen strategy.
*/
export async function startLlmTransientFaultProxy(
upstreamBaseUrlOrOptions,
faultCount = 1,
extraOptions = {},
) {
const options = normalizeOptions(
upstreamBaseUrlOrOptions,
faultCount,
extraOptions,
);
const downstreamSockets = new Set();
const upstreamRequests = new Set();
const upstreamSockets = new Set();
const heldForwardingWaiters = new Set();
const counterWaiters = new Set();
const requestLog = [];
let requestCount = 0;
let faultInjectedCount = 0;
let heldRequestCount = 0;
let forwardedRequestCount = 0;
let forwardingReleased = !options.holdAfterFault || options.faultCount === 0;
let stopping = false;
let stopped = false;
let stopPromise;
const stats = () =>
Object.freeze({
requestCount,
faultInjectedCount,
heldRequestCount,
forwardedRequestCount,
forwardingReleased,
stopped,
});
const requestLogSnapshot = () =>
Object.freeze(requestLog.map((entry) => Object.freeze({ ...entry })));
const notifyCounterWaiters = (kind) => {
for (const waiter of [...counterWaiters]) {
if (waiter.kind !== kind) continue;
clearTimeout(waiter.timer);
counterWaiters.delete(waiter);
waiter.resolve(stats());
}
};
const waitForCounter = (kind, timeoutMs) => {
if (!Number.isSafeInteger(timeoutMs) || timeoutMs <= 0) {
return Promise.reject(proxyError('timeoutMs must be a positive integer'));
}
const current = kind === 'fault' ? faultInjectedCount : heldRequestCount;
if (current > 0) return Promise.resolve(stats());
if (stopping || stopped) {
return Promise.reject(
proxyError('proxy stopped before the wait completed'),
);
}
return new Promise((resolve, reject) => {
const waiter = { kind, resolve, reject, timer: undefined };
waiter.timer = setTimeout(() => {
counterWaiters.delete(waiter);
reject(proxyError(`timed out waiting for proxy ${kind}`));
}, timeoutMs);
waiter.timer.unref?.();
counterWaiters.add(waiter);
});
};
const releaseHeldForwarding = (shouldForward) => {
for (const waiter of [...heldForwardingWaiters]) {
heldForwardingWaiters.delete(waiter);
waiter.complete(shouldForward);
}
};
const waitForForwardingRelease = (request, response) => {
if (forwardingReleased) return Promise.resolve(true);
if (stopping) return Promise.resolve(false);
return new Promise((resolve) => {
const onAborted = () => waiter.complete(false);
const onClosed = () => waiter.complete(false);
const waiter = {
complete: (shouldForward) => {
heldForwardingWaiters.delete(waiter);
request.off('aborted', onAborted);
response.off('close', onClosed);
resolve(shouldForward);
},
};
request.once('aborted', onAborted);
response.once('close', onClosed);
heldForwardingWaiters.add(waiter);
if (forwardingReleased) waiter.complete(true);
if (stopping) waiter.complete(false);
});
};
const sendBadGateway = (request, response) => {
if (stopping || response.destroyed) return;
if (response.headersSent) {
response.destroy();
return;
}
failClosed(request, response, 502);
};
const forwardRequest = (request, response, requestMetadata) => {
requestMetadata.forwardingStartedAtMs = Date.now();
const transport = options.upstream.protocol === 'https:' ? https : http;
let upstreamRequest;
try {
upstreamRequest = transport.request({
protocol: options.upstream.protocol,
hostname: options.upstream.hostname,
port: options.upstream.port,
method: request.method,
path: request.url,
headers: forwardHeaders(request, options.upstream.hostHeader),
agent: false,
setHost: false,
});
} catch {
sendBadGateway(request, response);
return;
}
forwardedRequestCount += 1;
upstreamRequests.add(upstreamRequest);
upstreamRequest.once('close', () =>
upstreamRequests.delete(upstreamRequest),
);
upstreamRequest.on('socket', (socket) => {
upstreamSockets.add(socket);
socket.once('close', () => upstreamSockets.delete(socket));
socket.on('error', () => {});
});
let upstreamResponse;
const stopUpstream = () => {
upstreamResponse?.destroy();
upstreamRequest.destroy();
};
request.once('aborted', stopUpstream);
request.once('error', stopUpstream);
response.once('error', stopUpstream);
response.once('close', () => {
if (!response.writableEnded) stopUpstream();
});
upstreamRequest.once('response', (receivedResponse) => {
upstreamResponse = receivedResponse;
receivedResponse.once('error', () => {
if (!response.destroyed) response.destroy();
});
receivedResponse.once('aborted', () => {
if (!response.destroyed) response.destroy();
});
if (stopping || response.destroyed) {
receivedResponse.destroy();
return;
}
try {
if (receivedResponse.statusMessage) {
response.writeHead(
receivedResponse.statusCode ?? 502,
receivedResponse.statusMessage,
receivedResponse.rawHeaders,
);
} else {
response.writeHead(
receivedResponse.statusCode ?? 502,
receivedResponse.rawHeaders,
);
}
} catch {
receivedResponse.destroy();
sendBadGateway(request, response);
return;
}
receivedResponse.pipe(response);
});
upstreamRequest.once('error', () => {
sendBadGateway(request, response);
});
request.pipe(upstreamRequest);
};
const handleRequest = async (request, response, expectsContinue) => {
requestCount += 1;
const requestMetadata = {
sequence: requestCount,
acceptedAtMs: Date.now(),
faultInjectedAtMs: null,
heldAtMs: null,
forwardingStartedAtMs: null,
};
requestLog.push(requestMetadata);
request.on('error', () => {});
response.on('error', () => {});
if (request.method !== 'POST') {
failClosed(request, response, 405);
return;
}
if (
!isOriginFormPath(request.url) ||
!isWithinBasePath(request.url, options.upstream.basePathname)
) {
failClosed(request, response, 400);
return;
}
let shouldInjectFault = faultInjectedCount < options.faultCount;
if (shouldInjectFault && options.shouldInjectFault) {
const decision = await options.shouldInjectFault(
Object.freeze({
sequence: requestMetadata.sequence,
acceptedAtMs: requestMetadata.acceptedAtMs,
}),
);
if (typeof decision !== 'boolean') {
throw proxyError('shouldInjectFault must resolve to a boolean');
}
shouldInjectFault = decision;
}
if (stopping || request.destroyed || response.destroyed) {
resetSocket(request.socket);
return;
}
if (
options.holdAfterFault &&
faultInjectedCount > 0 &&
!forwardingReleased
) {
heldRequestCount += 1;
requestMetadata.heldAtMs = Date.now();
notifyCounterWaiters('held');
const shouldForward = await waitForForwardingRelease(request, response);
if (!shouldForward) {
resetSocket(request.socket);
return;
}
}
if (shouldInjectFault && faultInjectedCount < options.faultCount) {
faultInjectedCount += 1;
requestMetadata.faultInjectedAtMs = Date.now();
notifyCounterWaiters('fault');
resetSocket(request.socket);
return;
}
if (stopping || request.destroyed || response.destroyed) {
resetSocket(request.socket);
return;
}
if (expectsContinue) response.writeContinue();
forwardRequest(request, response, requestMetadata);
};
const dispatchRequest = (request, response, expectsContinue = false) => {
void handleRequest(request, response, expectsContinue).catch(() => {
sendBadGateway(request, response);
});
};
const server = http.createServer();
server.on('request', (request, response) => {
dispatchRequest(request, response);
});
server.on('checkContinue', (request, response) => {
dispatchRequest(request, response, true);
});
server.on('checkExpectation', (request, response) => {
failClosed(request, response, 400);
});
server.on('connect', (_request, socket) => resetSocket(socket));
server.on('upgrade', (_request, socket) => resetSocket(socket));
server.on('clientError', (_error, socket) => resetSocket(socket));
server.on('connection', (socket) => {
downstreamSockets.add(socket);
socket.once('close', () => downstreamSockets.delete(socket));
socket.on('error', () => {});
});
try {
await bindServer(server, options.listen, options.fallbackPorts);
const address = server.address();
if (
!address ||
typeof address === 'string' ||
address.address !== LOOPBACK_HOST
) {
throw proxyError('transient fault proxy did not bind to loopback');
}
} catch (error) {
for (const socket of downstreamSockets) socket.destroy();
await closeServer(server);
if (error?.name === 'LlmTransientFaultProxyError') throw error;
throw proxyError('unable to start transient fault proxy');
}
server.on('error', () => {});
const address = server.address();
const port = address.port;
const url = `http://${LOOPBACK_HOST}:${port}`;
const baseUrl = `${url}${options.upstream.basePathname}`;
const releaseForwarding = () => {
if (stopping || forwardingReleased) return stats();
forwardingReleased = true;
releaseHeldForwarding(true);
return stats();
};
const stop = () => {
if (stopPromise) return stopPromise;
stopPromise = (async () => {
stopping = true;
releaseHeldForwarding(false);
for (const waiter of [...counterWaiters]) {
clearTimeout(waiter.timer);
counterWaiters.delete(waiter);
waiter.reject(proxyError('proxy stopped before the wait completed'));
}
const closePromise = closeServer(server);
for (const request of [...upstreamRequests]) request.destroy();
for (const socket of [...upstreamSockets]) socket.destroy();
for (const socket of [...downstreamSockets]) socket.destroy();
server.closeAllConnections?.();
await closePromise;
stopped = true;
})();
return stopPromise;
};
return Object.freeze({
url,
baseUrl,
port,
get stats() {
return stats();
},
getStats: stats,
getRequestLog: requestLogSnapshot,
waitForFault(timeoutMs = 5_000) {
return waitForCounter('fault', timeoutMs);
},
waitForHeldRequest(timeoutMs = 5_000) {
return waitForCounter('held', timeoutMs);
},
releaseForwarding,
stop,
});
}