mirror of
https://github.com/tiennm99/ccs.git
synced 2026-10-11 03:13:12 +00:00
feat(websearch): add managed mcp runtime and provider cooldowns
- switch ccs-websearch MCP to stdio-first framing and keep legacy compatibility - add cooldown persistence and bounded retries for provider failures - surface cooldown status in readiness output and regression tests
This commit is contained in:
1 parent
cae5b710b8
commit
7f83e041b7
10 files changed
+1115
-63
No files matched your search
@@ -30,6 +30,13 @@ const DDG_URL = 'https://html.duckduckgo.com/html/';
|
||||
const BRAVE_URL = 'https://api.search.brave.com/res/v1/web/search';
|
||||
const USER_AGENT =
|
||||
'Mozilla/5.0 (Macintosh; Intel Mac OS X 14_7_2) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/120.0.0.0 Safari/537.36';
|
||||
const PROVIDER_STATE_FILE = 'websearch-provider-state.json';
|
||||
const SHORT_RETRY_AFTER_MAX_SEC = 3;
|
||||
const TRANSIENT_RETRY_DELAY_MS = 750;
|
||||
const TRANSIENT_RETRY_ATTEMPTS = 1;
|
||||
const DEFAULT_RATE_LIMIT_COOLDOWN_SEC = 120;
|
||||
const DEFAULT_QUOTA_COOLDOWN_SEC = 900;
|
||||
const MAX_PROVIDER_COOLDOWN_SEC = 60 * 60;
|
||||
|
||||
const SHARED_INSTRUCTIONS = `Instructions:
|
||||
1. Search the web for current, up-to-date information
|
||||
@@ -101,6 +108,111 @@ function getSafeTracePrefixes() {
|
||||
];
|
||||
}
|
||||
|
||||
function getProviderStatePath() {
|
||||
return path.join(getCcsDirPath(), 'cache', PROVIDER_STATE_FILE);
|
||||
}
|
||||
|
||||
function readProviderState() {
|
||||
try {
|
||||
const statePath = getProviderStatePath();
|
||||
if (!fs.existsSync(statePath)) {
|
||||
return { cooldowns: {} };
|
||||
}
|
||||
|
||||
const parsed = JSON.parse(fs.readFileSync(statePath, 'utf8'));
|
||||
const cooldowns =
|
||||
parsed && typeof parsed === 'object' && parsed.cooldowns && typeof parsed.cooldowns === 'object'
|
||||
? parsed.cooldowns
|
||||
: {};
|
||||
return { cooldowns };
|
||||
} catch {
|
||||
return { cooldowns: {} };
|
||||
}
|
||||
}
|
||||
|
||||
function writeProviderState(state) {
|
||||
try {
|
||||
const statePath = getProviderStatePath();
|
||||
fs.mkdirSync(path.dirname(statePath), { recursive: true });
|
||||
const tempPath = `${statePath}.${process.pid}.${Date.now()}.tmp`;
|
||||
fs.writeFileSync(tempPath, JSON.stringify(state, null, 2) + '\n', 'utf8');
|
||||
fs.renameSync(tempPath, statePath);
|
||||
} catch {
|
||||
// Best-effort only.
|
||||
}
|
||||
}
|
||||
|
||||
function sanitizeProviderState(state) {
|
||||
const now = Date.now();
|
||||
const nextCooldowns = {};
|
||||
let changed = false;
|
||||
|
||||
for (const [providerId, entry] of Object.entries(state.cooldowns || {})) {
|
||||
if (!entry || typeof entry !== 'object') {
|
||||
changed = true;
|
||||
continue;
|
||||
}
|
||||
|
||||
const until = Number.parseInt(String(entry.until || ''), 10);
|
||||
if (!Number.isFinite(until) || until <= now) {
|
||||
changed = true;
|
||||
continue;
|
||||
}
|
||||
|
||||
nextCooldowns[providerId] = {
|
||||
until,
|
||||
reason: typeof entry.reason === 'string' ? entry.reason : 'rate_limited',
|
||||
updatedAt: Number.parseInt(String(entry.updatedAt || ''), 10) || now,
|
||||
sourceError: typeof entry.sourceError === 'string' ? entry.sourceError : '',
|
||||
};
|
||||
}
|
||||
|
||||
return {
|
||||
state: { cooldowns: nextCooldowns },
|
||||
changed,
|
||||
};
|
||||
}
|
||||
|
||||
function getProviderCooldown(providerId) {
|
||||
const { state, changed } = sanitizeProviderState(readProviderState());
|
||||
if (changed) {
|
||||
writeProviderState(state);
|
||||
}
|
||||
|
||||
return state.cooldowns[providerId] || null;
|
||||
}
|
||||
|
||||
function clearProviderCooldown(providerId) {
|
||||
const { state } = sanitizeProviderState(readProviderState());
|
||||
if (!(providerId in state.cooldowns)) {
|
||||
return;
|
||||
}
|
||||
|
||||
delete state.cooldowns[providerId];
|
||||
writeProviderState(state);
|
||||
}
|
||||
|
||||
function applyProviderCooldown(providerId, cooldownSec, reason, sourceError) {
|
||||
const clampedCooldownSec = Math.max(
|
||||
1,
|
||||
Math.min(MAX_PROVIDER_COOLDOWN_SEC, Math.floor(cooldownSec))
|
||||
);
|
||||
const { state } = sanitizeProviderState(readProviderState());
|
||||
const until = Date.now() + clampedCooldownSec * 1000;
|
||||
state.cooldowns[providerId] = {
|
||||
until,
|
||||
reason,
|
||||
updatedAt: Date.now(),
|
||||
sourceError: sourceError || '',
|
||||
};
|
||||
writeProviderState(state);
|
||||
return until;
|
||||
}
|
||||
|
||||
function sleep(ms) {
|
||||
return new Promise((resolve) => setTimeout(resolve, ms));
|
||||
}
|
||||
|
||||
function getAllowedTraceFileOverride() {
|
||||
const configured = (process.env.CCS_WEBSEARCH_TRACE_FILE || '').trim();
|
||||
if (!configured) {
|
||||
@@ -146,6 +258,42 @@ function traceWebSearchEvent(event, payload = {}) {
|
||||
}
|
||||
}
|
||||
|
||||
function readHeaderValue(headers, headerName) {
|
||||
if (!headers) {
|
||||
return '';
|
||||
}
|
||||
|
||||
if (typeof headers.get === 'function') {
|
||||
return headers.get(headerName) || '';
|
||||
}
|
||||
|
||||
const direct = headers[headerName] ?? headers[String(headerName).toLowerCase()];
|
||||
if (Array.isArray(direct)) {
|
||||
return direct[0] || '';
|
||||
}
|
||||
return typeof direct === 'string' ? direct : '';
|
||||
}
|
||||
|
||||
function parseRetryAfterSeconds(rawValue) {
|
||||
const value = String(rawValue || '').trim();
|
||||
if (!value) {
|
||||
return null;
|
||||
}
|
||||
|
||||
const asSeconds = Number.parseInt(value, 10);
|
||||
if (Number.isFinite(asSeconds) && asSeconds > 0) {
|
||||
return asSeconds;
|
||||
}
|
||||
|
||||
const asDate = Date.parse(value);
|
||||
if (Number.isFinite(asDate)) {
|
||||
const deltaSec = Math.ceil((asDate - Date.now()) / 1000);
|
||||
return deltaSec > 0 ? deltaSec : null;
|
||||
}
|
||||
|
||||
return null;
|
||||
}
|
||||
|
||||
function getQueryFingerprint(query) {
|
||||
const normalizedQuery = typeof query === 'string' ? query.trim() : '';
|
||||
return {
|
||||
@@ -371,6 +519,8 @@ async function tryBraveSearch(query, timeoutSec = DEFAULT_TIMEOUT_SEC) {
|
||||
return {
|
||||
success: false,
|
||||
error: `Brave Search returned ${response.status}: ${body.slice(0, 160)}`,
|
||||
statusCode: response.status,
|
||||
retryAfterSec: parseRetryAfterSeconds(readHeaderValue(response.headers, 'retry-after')),
|
||||
};
|
||||
}
|
||||
|
||||
@@ -422,7 +572,12 @@ async function tryExaSearch(query, timeoutSec = DEFAULT_TIMEOUT_SEC) {
|
||||
|
||||
if (!response.ok) {
|
||||
const body = await response.text();
|
||||
return { success: false, error: `Exa returned ${response.status}: ${body.slice(0, 160)}` };
|
||||
return {
|
||||
success: false,
|
||||
error: `Exa returned ${response.status}: ${body.slice(0, 160)}`,
|
||||
statusCode: response.status,
|
||||
retryAfterSec: parseRetryAfterSeconds(readHeaderValue(response.headers, 'retry-after')),
|
||||
};
|
||||
}
|
||||
|
||||
const body = await response.json();
|
||||
@@ -474,7 +629,12 @@ async function tryTavilySearch(query, timeoutSec = DEFAULT_TIMEOUT_SEC) {
|
||||
|
||||
if (!response.ok) {
|
||||
const body = await response.text();
|
||||
return { success: false, error: `Tavily returned ${response.status}: ${body.slice(0, 160)}` };
|
||||
return {
|
||||
success: false,
|
||||
error: `Tavily returned ${response.status}: ${body.slice(0, 160)}`,
|
||||
statusCode: response.status,
|
||||
retryAfterSec: parseRetryAfterSeconds(readHeaderValue(response.headers, 'retry-after')),
|
||||
};
|
||||
}
|
||||
|
||||
const body = await response.json();
|
||||
@@ -511,7 +671,12 @@ async function tryDuckDuckGoSearch(query, timeoutSec = DEFAULT_TIMEOUT_SEC) {
|
||||
);
|
||||
|
||||
if (!response.ok) {
|
||||
return { success: false, error: `DuckDuckGo returned ${response.status}` };
|
||||
return {
|
||||
success: false,
|
||||
error: `DuckDuckGo returned ${response.status}`,
|
||||
statusCode: response.status,
|
||||
retryAfterSec: parseRetryAfterSeconds(readHeaderValue(response.headers, 'retry-after')),
|
||||
};
|
||||
}
|
||||
|
||||
const html = await response.text();
|
||||
@@ -720,8 +885,174 @@ function getConfiguredProviders() {
|
||||
];
|
||||
}
|
||||
|
||||
function looksLikeQuotaExhaustion(errorMessage) {
|
||||
const lower = String(errorMessage || '').toLowerCase();
|
||||
return (
|
||||
(lower.includes('quota') &&
|
||||
(lower.includes('exceed') ||
|
||||
lower.includes('exhaust') ||
|
||||
lower.includes('deplet') ||
|
||||
lower.includes('limit') ||
|
||||
lower.includes('used up'))) ||
|
||||
lower.includes('insufficient credits') ||
|
||||
lower.includes('credit balance') ||
|
||||
lower.includes('out of credits') ||
|
||||
lower.includes('billing hard limit') ||
|
||||
lower.includes('monthly usage cap')
|
||||
);
|
||||
}
|
||||
|
||||
function looksLikeTransientFailure(errorMessage) {
|
||||
const lower = String(errorMessage || '').toLowerCase();
|
||||
return (
|
||||
lower.includes('timed out') ||
|
||||
lower.includes('timeout') ||
|
||||
lower.includes('temporarily unavailable') ||
|
||||
lower.includes('service unavailable') ||
|
||||
lower.includes('bad gateway') ||
|
||||
lower.includes('gateway timeout') ||
|
||||
lower.includes('socket hang up') ||
|
||||
lower.includes('econnreset') ||
|
||||
lower.includes('fetch failed') ||
|
||||
lower.includes('network')
|
||||
);
|
||||
}
|
||||
|
||||
function classifyProviderFailure(result) {
|
||||
const errorMessage = String(result.error || '');
|
||||
const statusCode =
|
||||
Number.isFinite(result.statusCode) && result.statusCode > 0 ? result.statusCode : null;
|
||||
const retryAfterSec = Number.isFinite(result.retryAfterSec) ? result.retryAfterSec : null;
|
||||
|
||||
if (looksLikeQuotaExhaustion(errorMessage)) {
|
||||
return {
|
||||
kind: 'cooldown',
|
||||
reason: 'quota_exhausted',
|
||||
cooldownSec: retryAfterSec || DEFAULT_QUOTA_COOLDOWN_SEC,
|
||||
retryAfterSec,
|
||||
};
|
||||
}
|
||||
|
||||
if (statusCode === 429 || /too many requests|rate limit/i.test(errorMessage)) {
|
||||
if (retryAfterSec && retryAfterSec <= SHORT_RETRY_AFTER_MAX_SEC) {
|
||||
return {
|
||||
kind: 'retry',
|
||||
delayMs: retryAfterSec * 1000,
|
||||
reason: 'rate_limited_short_backoff',
|
||||
retryAfterSec,
|
||||
};
|
||||
}
|
||||
|
||||
return {
|
||||
kind: 'cooldown',
|
||||
reason: 'rate_limited',
|
||||
cooldownSec: retryAfterSec || DEFAULT_RATE_LIMIT_COOLDOWN_SEC,
|
||||
retryAfterSec,
|
||||
};
|
||||
}
|
||||
|
||||
if (
|
||||
(statusCode && [502, 503, 504].includes(statusCode)) ||
|
||||
looksLikeTransientFailure(errorMessage)
|
||||
) {
|
||||
return {
|
||||
kind: 'retry',
|
||||
delayMs: TRANSIENT_RETRY_DELAY_MS,
|
||||
reason: 'transient_failure',
|
||||
retryAfterSec,
|
||||
};
|
||||
}
|
||||
|
||||
return {
|
||||
kind: 'fail',
|
||||
reason: 'non_retryable',
|
||||
retryAfterSec,
|
||||
};
|
||||
}
|
||||
|
||||
async function runProviderWithPolicy(provider, query, timeoutSec, fingerprint) {
|
||||
for (let attempt = 0; attempt <= TRANSIENT_RETRY_ATTEMPTS; attempt += 1) {
|
||||
traceWebSearchEvent('websearch_provider_attempt', {
|
||||
source: 'provider',
|
||||
providerId: provider.id,
|
||||
providerName: provider.name,
|
||||
attempt: attempt + 1,
|
||||
...fingerprint,
|
||||
});
|
||||
|
||||
const result = await provider.fn(query, timeoutSec);
|
||||
if (result.success) {
|
||||
clearProviderCooldown(provider.id);
|
||||
return result;
|
||||
}
|
||||
|
||||
const policy = classifyProviderFailure(result);
|
||||
if (policy.kind === 'retry' && attempt < TRANSIENT_RETRY_ATTEMPTS) {
|
||||
traceWebSearchEvent('websearch_provider_retry_scheduled', {
|
||||
source: 'provider',
|
||||
providerId: provider.id,
|
||||
providerName: provider.name,
|
||||
attempt: attempt + 1,
|
||||
delayMs: policy.delayMs,
|
||||
reason: policy.reason,
|
||||
retryAfterSec: policy.retryAfterSec,
|
||||
...fingerprint,
|
||||
});
|
||||
await sleep(policy.delayMs);
|
||||
continue;
|
||||
}
|
||||
|
||||
if (policy.kind === 'retry' && policy.reason === 'rate_limited_short_backoff') {
|
||||
const cooldownSec = policy.retryAfterSec || DEFAULT_RATE_LIMIT_COOLDOWN_SEC;
|
||||
const until = applyProviderCooldown(provider.id, cooldownSec, 'rate_limited', result.error);
|
||||
traceWebSearchEvent('websearch_provider_cooldown_applied', {
|
||||
source: 'provider',
|
||||
providerId: provider.id,
|
||||
providerName: provider.name,
|
||||
cooldownUntil: until,
|
||||
cooldownSec,
|
||||
reason: 'rate_limited',
|
||||
retryAfterSec: policy.retryAfterSec,
|
||||
afterRetryExhausted: true,
|
||||
...fingerprint,
|
||||
});
|
||||
return {
|
||||
...result,
|
||||
error: `${result.error} (cooldown ${cooldownSec}s)`,
|
||||
};
|
||||
}
|
||||
|
||||
if (policy.kind === 'cooldown') {
|
||||
const until = applyProviderCooldown(
|
||||
provider.id,
|
||||
policy.cooldownSec,
|
||||
policy.reason,
|
||||
result.error
|
||||
);
|
||||
traceWebSearchEvent('websearch_provider_cooldown_applied', {
|
||||
source: 'provider',
|
||||
providerId: provider.id,
|
||||
providerName: provider.name,
|
||||
cooldownUntil: until,
|
||||
cooldownSec: policy.cooldownSec,
|
||||
reason: policy.reason,
|
||||
retryAfterSec: policy.retryAfterSec,
|
||||
...fingerprint,
|
||||
});
|
||||
return {
|
||||
...result,
|
||||
error: `${result.error} (cooldown ${policy.cooldownSec}s)`,
|
||||
};
|
||||
}
|
||||
|
||||
return result;
|
||||
}
|
||||
|
||||
return { success: false, error: 'Provider retry policy exhausted' };
|
||||
}
|
||||
|
||||
function getActiveProviders() {
|
||||
return getConfiguredProviders().filter((provider) => provider.available());
|
||||
return getConfiguredProviders().filter((provider) => !getProviderCooldown(provider.id) && provider.available());
|
||||
}
|
||||
|
||||
function getActiveProviderIds() {
|
||||
@@ -733,8 +1064,29 @@ function hasAnyActiveProviders() {
|
||||
}
|
||||
|
||||
async function runLocalWebSearch(query, timeoutSec = DEFAULT_TIMEOUT_SEC) {
|
||||
const activeProviders = getActiveProviders();
|
||||
const fingerprint = getQueryFingerprint(query);
|
||||
const configuredProviders = getConfiguredProviders();
|
||||
const activeProviders = [];
|
||||
|
||||
for (const provider of configuredProviders) {
|
||||
const cooldown = getProviderCooldown(provider.id);
|
||||
if (cooldown) {
|
||||
traceWebSearchEvent('websearch_provider_cooldown_skip', {
|
||||
source: 'provider',
|
||||
providerId: provider.id,
|
||||
providerName: provider.name,
|
||||
cooldownUntil: cooldown.until,
|
||||
cooldownReason: cooldown.reason,
|
||||
remainingMs: Math.max(0, cooldown.until - Date.now()),
|
||||
...fingerprint,
|
||||
});
|
||||
continue;
|
||||
}
|
||||
|
||||
if (provider.available()) {
|
||||
activeProviders.push(provider);
|
||||
}
|
||||
}
|
||||
|
||||
debug(
|
||||
`Enabled providers: ${activeProviders.map((provider) => provider.name).join(', ') || 'none'}`
|
||||
@@ -757,13 +1109,7 @@ async function runLocalWebSearch(query, timeoutSec = DEFAULT_TIMEOUT_SEC) {
|
||||
const errors = [];
|
||||
for (const provider of activeProviders) {
|
||||
debug(`Trying ${provider.name}`);
|
||||
traceWebSearchEvent('websearch_provider_attempt', {
|
||||
source: 'provider',
|
||||
providerId: provider.id,
|
||||
providerName: provider.name,
|
||||
...fingerprint,
|
||||
});
|
||||
const result = await provider.fn(query, timeoutSec);
|
||||
const result = await runProviderWithPolicy(provider, query, timeoutSec, fingerprint);
|
||||
if (result.success) {
|
||||
traceWebSearchEvent('websearch_provider_success', {
|
||||
source: 'provider',
|
||||
@@ -890,8 +1236,10 @@ module.exports = {
|
||||
runLocalWebSearch,
|
||||
shouldSkipHook,
|
||||
getActiveProviderIds,
|
||||
classifyProviderFailure,
|
||||
getQueryFingerprint,
|
||||
getSkipReason,
|
||||
parseRetryAfterSeconds,
|
||||
traceWebSearchEvent,
|
||||
tryExaSearch,
|
||||
tryTavilySearch,
|
||||
|
||||
@@ -61,9 +61,7 @@ function getTools() {
|
||||
}
|
||||
|
||||
function writeMessage(message) {
|
||||
const body = Buffer.from(JSON.stringify(message), 'utf8');
|
||||
process.stdout.write(`Content-Length: ${body.length}\r\n\r\n`);
|
||||
process.stdout.write(body);
|
||||
process.stdout.write(`${JSON.stringify(message)}\n`);
|
||||
}
|
||||
|
||||
function writeResponse(id, result) {
|
||||
@@ -262,26 +260,46 @@ function writeSessionSummary(exitCodeOrSignal) {
|
||||
|
||||
function parseMessages() {
|
||||
while (true) {
|
||||
const headerEnd = inputBuffer.indexOf('\r\n\r\n');
|
||||
if (headerEnd === -1) {
|
||||
return;
|
||||
}
|
||||
let body;
|
||||
const startsWithLegacyHeaders = inputBuffer
|
||||
.slice(0, Math.min(inputBuffer.length, 32))
|
||||
.toString('utf8')
|
||||
.toLowerCase()
|
||||
.startsWith('content-length:');
|
||||
|
||||
const headerText = inputBuffer.slice(0, headerEnd).toString('utf8');
|
||||
const contentLengthMatch = headerText.match(/content-length:\s*(\d+)/i);
|
||||
if (!contentLengthMatch) {
|
||||
inputBuffer = Buffer.alloc(0);
|
||||
return;
|
||||
}
|
||||
if (startsWithLegacyHeaders) {
|
||||
const headerEnd = inputBuffer.indexOf('\r\n\r\n');
|
||||
if (headerEnd === -1) {
|
||||
return;
|
||||
}
|
||||
|
||||
const contentLength = Number.parseInt(contentLengthMatch[1], 10);
|
||||
const messageEnd = headerEnd + 4 + contentLength;
|
||||
if (inputBuffer.length < messageEnd) {
|
||||
return;
|
||||
}
|
||||
const headerText = inputBuffer.slice(0, headerEnd).toString('utf8');
|
||||
const contentLengthMatch = headerText.match(/content-length:\s*(\d+)/i);
|
||||
if (!contentLengthMatch) {
|
||||
inputBuffer = Buffer.alloc(0);
|
||||
return;
|
||||
}
|
||||
|
||||
const body = inputBuffer.slice(headerEnd + 4, messageEnd).toString('utf8');
|
||||
inputBuffer = inputBuffer.slice(messageEnd);
|
||||
const contentLength = Number.parseInt(contentLengthMatch[1], 10);
|
||||
const messageEnd = headerEnd + 4 + contentLength;
|
||||
if (inputBuffer.length < messageEnd) {
|
||||
return;
|
||||
}
|
||||
|
||||
body = inputBuffer.slice(headerEnd + 4, messageEnd).toString('utf8');
|
||||
inputBuffer = inputBuffer.slice(messageEnd);
|
||||
} else {
|
||||
const newlineIndex = inputBuffer.indexOf('\n');
|
||||
if (newlineIndex === -1) {
|
||||
return;
|
||||
}
|
||||
|
||||
body = inputBuffer.slice(0, newlineIndex).toString('utf8').replace(/\r$/, '').trim();
|
||||
inputBuffer = inputBuffer.slice(newlineIndex + 1);
|
||||
if (!body) {
|
||||
continue;
|
||||
}
|
||||
}
|
||||
|
||||
let message;
|
||||
try {
|
||||
|
||||
Reference in new issue
Block a user