investigate_ddg_blocks.txt
investigate_ddg_blocks.txt
// ==================== FILE: src\swarm\worker.ts ====================
// src/swarm/worker.ts
import { LlmCallManager } from "../utils/llm";
import { logLlmDiagnostics } from "../utils/tokens";
import { extractYouTubeTranscript } from "../net/youtube-extractor";
import {
searchSerper,
searchBraveApi,
searchOpenAlex,
searchCrossref,
searchArxiv,
searchGdelt,
searchReferenceSites,
} from "../net/api-engines";
import {
searchDDG,
searchDDGPaginated,
DdgRateLimiter,
sharedDdgLimiter,
} from "../net/ddg";
import { multiEngineSearch, SearchEngine } from "../net/search-engines";
import { SearchHealthTracker } from "./health";
import { fetchPage, sleep } from "../net/http";
import {
extractPage,
contentFingerprint,
computeRelevance,
} from "../net/extractor";
import {
isPdfUrl,
isPdfContentType,
extractPdf,
} from "../net/pdf-extractor";
import {
scoreCandidate,
rankCandidates,
scoreOutlinks,
} from "../scoring/authority";
import {
SwarmTask,
WorkerResult,
CrawledSource,
ScoredCandidate,
SourceTier,
StatusFn,
WarnFn,
SearchHit,
ExtractedPage,
} from "../types";
import { harvestLocalSources } from "../local/search";
import {
BATCH_INTER_FETCH_DELAY_MS,
MIN_USEFUL_WORD_COUNT,
} from "../constants";
import { normalizeUrl } from "./visited-cache";
export interface CrawlMetrics {
ddgQueries: number;
ddgHits: number;
mutatedQueriesTried: number;
mutationAccepted: number;
mutationHits: number;
extraEngineQueries: number;
extraEngineHits: number;
rawHits: number;
dedupedHits: number;
rankedCandidates: number;
fetchCandidates: number;
fetchAttempts: number;
fetchFailures: number;
acceptedSources: number;
skippedLowWordCount: number;
skippedOffTopic: number;
skippedVeryOffTopic: number;
skippedDuplicateContent: number;
skippedVisited: number;
skippedDomainCap: number;
skippedAvoided: number;
skippedBlacklisted: number;
cacheChecks: number;
cacheHits: number;
cacheAccepted: number;
cacheRejectedDuplicate: number;
cacheRejectedOffTopic: number;
cacheRejectedLowWordCount: number;
cacheWrites: number;
followedLinks: number;
crossWorkerDiscoveriesUsed: number;
localSourcesAccepted: number;
}
export interface SharedCrawlState {
readonly visitedUrls: ReadonlySet<string>;
readonly contentHashes: ReadonlySet<string>;
readonly domainCounts: ReadonlyMap<string, number>;
readonly domainFailures: ReadonlyMap<string, number>;
addVisited(url: string): void;
addHash(hash: string): void;
incrementDomain(url: string): void;
domainCount(url: string): number;
noteFailure(url: string, reason: string): void;
noteDomainFailure(url: string): void;
isDomainBlacklisted(url: string): boolean;
shouldAvoidUrl(url: string): boolean;
pushDiscovery(url: string, title: string, fromWorker: string): void;
drainDiscoveries(limit: number): ReadonlyArray<{ url: string; title: string }>;
isRecentlyVisited(url: string): boolean;
getCachedSource(url: string): CrawledSource | null;
markVisitedPersistent(source: CrawledSource): void;
addWebCacheDocument(source: CrawledSource): void;
getMetricsSnapshot(): Readonly<CrawlMetrics>;
mergeMetrics(delta: Partial<CrawlMetrics>): void;
}
const MUTATION_STRATEGIES: ReadonlyArray<(q: string) => string> = [
(q) => `"${q}"`,
(q) => `${q} explained`,
(q) => `${q} guide overview`,
(q) => `${q} research`,
(q) => q.split(" ").slice(0, 10).join(" "),
(q) => q.replace(/\b(how|what|why|when)\b/gi, " ").trim(),
];
type SearchHitLike = {
url: string;
title: string;
snippet: string;
query: string;
discoveredBy?: string;
requestedRoute?: string; // <--- ADD THIS
actualBackend?: string; // <--- ADD THIS
resultDomain?: string; // <--- ADD THIS
};
function isRateLimitError(msg: string): boolean {
const lower = msg.toLowerCase();
return (
lower.includes("429") ||
lower.includes("too many requests") ||
lower.includes("rate limit")
);
}
function zeroMetrics(): CrawlMetrics {
return {
ddgQueries: 0, ddgHits: 0, mutatedQueriesTried: 0, mutationAccepted: 0,
mutationHits: 0, extraEngineQueries: 0, extraEngineHits: 0, rawHits: 0,
dedupedHits: 0, rankedCandidates: 0, fetchCandidates: 0, fetchAttempts: 0,
fetchFailures: 0, acceptedSources: 0, skippedLowWordCount: 0, skippedOffTopic: 0,
skippedVeryOffTopic: 0, skippedDuplicateContent: 0, skippedVisited: 0, skippedDomainCap: 0,
skippedAvoided: 0, skippedBlacklisted: 0, cacheChecks: 0, cacheHits: 0, cacheAccepted: 0,
cacheRejectedDuplicate: 0, cacheRejectedOffTopic: 0, cacheRejectedLowWordCount: 0,
cacheWrites: 0, followedLinks: 0, crossWorkerDiscoveriesUsed: 0, localSourcesAccepted: 0,
};
}
export async function runWorker(
task: SwarmTask,
state: SharedCrawlState,
signal: AbortSignal,
status: StatusFn,
warn: WarnFn,
topicKws: ReadonlyArray<string> = [],
limiter: DdgRateLimiter = sharedDdgLimiter,
health: SearchHealthTracker,
llmManager: LlmCallManager
): Promise<WorkerResult> {
const sources: CrawledSource[] = [];
const errors: string[] = [];
const queriesExecuted: string[] = [];
const roleTag = `[${task.label}]`;
const metrics = zeroMetrics();
status(`${roleTag} Starting - ${task.queries.length} queries, budget=${task.pageBudget}, links=${task.followLinks ? "on" : "off"}`);
if (task.enableLocalSources) {
const localBudget = Math.max(2, Math.ceil(task.pageBudget * 0.3));
const localSources = harvestLocalSources(
task.queries, task.role, task.label, localBudget, task.contentLimit, task.localLibraryIds, task.roleLibraryMap,
);
if (localSources.length > 0) {
for (const src of localSources) {
state.addVisited(src.url);
state.addHash(contentFingerprint(src.text));
sources.push(src);
}
metrics.localSourcesAccepted += localSources.length;
status(`${roleTag} Local sources: ${localSources.length} chunks`);
}
}
const allHits: SearchHitLike[] = [];
let ddgBlocked = false;
for (const query of task.queries) {
if (signal.aborted) break;
// 9. QUERY PLANNING RULES - Reject bad queries
// Relaxed Query Quality Check: Reject only empty, corrupt, or absurdly long query strings
const qTrimmed = query.trim();
if (!qTrimmed || qTrimmed.length < 3 || qTrimmed.length > 250) {
warn(`[Query Quality] Rejected out-of-bounds query length (${qTrimmed.length} chars): "${qTrimmed.slice(0, 50)}..."`);
continue;
}
let ddgHits: ReadonlyArray<SearchHit> = [];
let effectiveHitCount = 0;
// 4. DDG FAILURE DEFINITION & HEALTH TRACKING
if (task.extraEngines.includes("ddg") && health.isDdgAvailable()) {
metrics.ddgQueries++;
try {
if (task.searchPages > 1) {
ddgHits = await searchDDGPaginated(query, task.searchResultsPerQuery, task.searchPages, task.safeSearch, signal, limiter, task.timeRange ?? "all");
} else {
ddgHits = await searchDDG(query, task.searchResultsPerQuery, task.safeSearch, signal, limiter, task.timeRange ?? "all");
}
// 5. LOWERED DDG TRIPWIRE: Accept even 1 result to support tight encyclopedia operators
if (ddgHits.length < 1) {
throw new Error("INVALID_RESULT_SET (< 1 result)");
}
health.recordDdgSuccess(ddgHits.length);
metrics.ddgHits += ddgHits.length;
effectiveHitCount = ddgHits.length;
// Track Attribution
for (const h of ddgHits) {
let domain = "";
try { domain = new URL(h.url).hostname.replace(/^www\./, ""); } catch {}
allHits.push({
...h,
query,
requestedRoute: "DDG",
actualBackend: "DDG",
resultDomain: domain
});
}
queriesExecuted.push(query);
status(`${roleTag} DDG: "${query}" -> ${ddgHits.length} results`);
} catch (err: unknown) {
if (isAbortError(err)) break;
const msg = errorMessage(err);
health.recordDdgFailure(msg);
warn(`${roleTag} DDG failed (${health.ddg.consecutiveFailures}/5): ${msg}`);
}
}
// Query Mutation (LLM)
const isHighlyRestrictive = query.includes("site:");
if (
ddgHits.length < task.queryMutationThreshold &&
!signal.aborted &&
!ddgBlocked &&
!isHighlyRestrictive &&
llmManager.canCall(task.id)
) {
llmManager.recordCall(task.id);
const contextSnippets = allHits.filter((h) => h.query === query).slice(0, 3).map((h) => h.snippet);
const llmMutated = await getLLMQueryMutation(query, contextSnippets, signal);
const mutationsToTry = [llmMutated, ...MUTATION_STRATEGIES.map((s) => s(query))].filter((m): m is string => !!m && m !== query);
let bestMutation: { query: string; hits: ReadonlyArray<SearchHit> } | null = null;
const baseline = mutationQuality(ddgHits);
for (const mutated of mutationsToTry) {
if (signal.aborted) break;
metrics.mutatedQueriesTried++;
metrics.ddgQueries++;
try {
const mutHits = await searchDDG(mutated, task.searchResultsPerQuery, task.safeSearch, signal, limiter, task.timeRange ?? "all");
metrics.ddgHits += mutHits.length;
if (mutationQuality(mutHits) > baseline) {
if (!bestMutation || mutationQuality(mutHits) > mutationQuality(bestMutation.hits)) {
bestMutation = { query: mutated, hits: mutHits };
}
}
} catch (err) {
const msg = errorMessage(err);
if (isRateLimitError(msg)) {
ddgBlocked = true;
warn(`${roleTag} Mutation RATE LIMIT (429)! Falling back...`);
break;
}
}
}
if (bestMutation) {
metrics.mutationAccepted++;
metrics.mutationHits += bestMutation.hits.length;
for (const h of bestMutation.hits) allHits.push({ ...h, query: bestMutation.query });
queriesExecuted.push(bestMutation.query);
effectiveHitCount = Math.max(effectiveHitCount, bestMutation.hits.length);
}
} else if (llmManager.canCall(task.id) === false && !isHighlyRestrictive && !ddgBlocked) {
warn(`${roleTag} LLM call budget exhausted, skipping mutation.`);
}
// Fallback to SearxNG or API if DDG is dead or hit 0
const shouldUseExtraEngines = task.extraEngines.length > 0 && !signal.aborted && (ddgBlocked || effectiveHitCount < Math.ceil(task.searchResultsPerQuery * 0.6) || task.role === "academic" || task.role === "primary" || task.role === "regulatory");
if (shouldUseExtraEngines) {
metrics.extraEngineQueries++;
try {
const extraHits = await multiEngineSearch(
query, Math.min(task.searchResultsPerQuery, 8),
task.extraEngines as ReadonlyArray<SearchEngine>,
signal, () => limiter, task.timeRange ?? "all",
{ serperApiKey: (task as any).serperApiKey, braveApiKey: (task as any).braveApiKey }
);
metrics.extraEngineHits += extraHits.length;
for (const h of extraHits) allHits.push({ ...h, query });
if (extraHits.length > 0) status(`${roleTag} -> ${extraHits.length} extra results`);
} catch {}
}
}
metrics.rawHits = allHits.length;
if (signal.aborted || allHits.length === 0) {
metrics.acceptedSources = sources.length;
state.mergeMetrics(metrics);
return { taskId: task.id, role: task.role, label: task.label, sources, queries: queriesExecuted, errors };
}
const deduped = deduplicateByUrl(allHits);
metrics.dedupedHits = deduped.length;
const scored = deduped.map((h) => {
const sc = scoreCandidate(h, h.query);
return { ...sc, discoveredBy: h.discoveredBy };
});
const filtered = task.preferredTiers ? scored.filter((c) => task.preferredTiers!.includes(c.tier)) : scored;
const poolSize = task.pageBudget * task.candidatePoolMultiplier;
const ranked = rankCandidates(filtered.length > 0 ? filtered : scored, poolSize);
const candidates = capCandidatesPerHost(ranked, 2);
metrics.rankedCandidates = candidates.length;
metrics.fetchCandidates = candidates.length;
await fetchBatch(candidates, task, state, signal, status, warn, sources, errors, roleTag, topicKws, metrics, llmManager);
if (task.followLinks && sources.length > 0 && sources.length < task.pageBudget) {
for (let depth = 1; depth <= task.linkCrawlDepth; depth++) {
if (sources.length >= task.pageBudget || signal.aborted) break;
const budget = Math.min(task.pageBudget - sources.length, task.maxLinksToFollow);
if (budget <= 0) break;
const sourcesForLinks = depth === 1 ? sources : sources.slice(-budget * 2);
const newCount = await followLinks(sourcesForLinks, task, state, signal, status, warn, sources, errors, roleTag, budget, topicKws, depth, metrics, llmManager);
metrics.followedLinks += newCount;
if (newCount === 0) break;
}
}
for (const src of sources) {
for (const link of src.outlinks.slice(0, 5)) {
state.pushDiscovery(link.href, link.text, task.label);
}
}
metrics.acceptedSources = sources.length;
state.mergeMetrics(metrics);
status(`✅ ${roleTag} Mission Complete - ${sources.length} high-quality sources collected!`);
return { taskId: task.id, role: task.role, label: task.label, sources, queries: queriesExecuted, errors };
}
async function fetchBatch(
candidates: ReadonlyArray<ScoredCandidate>,
task: SwarmTask,
state: SharedCrawlState,
signal: AbortSignal,
status: StatusFn,
warn: WarnFn,
results: CrawledSource[],
errors: string[],
tag: string,
topicKws: ReadonlyArray<string>,
metrics: CrawlMetrics,
llmManager: LlmCallManager
): Promise<void> {
let idx = 0;
const concurrency = task.workerConcurrency;
const domainCap = task.maxPagesPerDomain;
const minRelevance = task.minRelevanceScore;
while (results.length < task.pageBudget && idx < candidates.length && !signal.aborted) {
const slice = candidates.slice(idx, idx + concurrency);
idx += concurrency;
const batch: ScoredCandidate[] = [];
for (const c of slice) {
const normalized = normalizeUrl(c.url);
if (state.visitedUrls.has(normalized)) { metrics.skippedVisited++; continue; }
if (state.domainCount(c.url) >= domainCap) { metrics.skippedDomainCap++; continue; }
if (state.shouldAvoidUrl(c.url)) { metrics.skippedAvoided++; continue; }
if (state.isDomainBlacklisted(c.url)) { metrics.skippedBlacklisted++; continue; }
batch.push(c);
}
if (batch.length === 0) continue;
for (const c of batch) state.addVisited(c.url);
const settled = await Promise.allSettled(
batch.map((c) => resolveCandidate(c, task, state, topicKws, signal, metrics)),
);
for (let i = 0; i < settled.length; i++) {
const candidate = batch[i];
const settledResult = settled[i];
if (signal.aborted) return;
if (settledResult.status === "rejected") {
if (!isAbortError(settledResult.reason)) {
const msg = errorMessage(settledResult.reason);
metrics.fetchFailures++;
state.noteFailure(candidate.url, msg);
state.noteDomainFailure(candidate.url);
warn(`${tag} Failed: ${truncUrl(candidate.url)} - ${msg}`);
errors.push(`fetch:${candidate.url}: ${msg}`);
}
continue;
}
const { page, fromCache } = settledResult.value;
if (page.wordCount < MIN_USEFUL_WORD_COUNT) { metrics.skippedLowWordCount++; if (fromCache) metrics.cacheRejectedLowWordCount++; continue; }
if (page.relevanceScore < minRelevance * 0.5) { metrics.skippedVeryOffTopic++; if (fromCache) metrics.cacheRejectedOffTopic++; else state.noteDomainFailure(candidate.url); continue; }
if (page.relevanceScore < minRelevance) { metrics.skippedOffTopic++; if (fromCache) metrics.cacheRejectedOffTopic++; continue; }
const fp = contentFingerprint(page.text);
if (state.contentHashes.has(fp)) { metrics.skippedDuplicateContent++; if (fromCache) metrics.cacheRejectedDuplicate++; continue; }
state.addHash(fp);
state.incrementDomain(candidate.url);
if (fromCache) {
metrics.cacheAccepted++;
} else {
state.markVisitedPersistent(page);
metrics.cacheWrites++;
state.addWebCacheDocument(page);
}
results.push(page);
llmManager.recordProgress(); // Tell the watchdog we found something!
status(`${tag} [${results.length}/${task.pageBudget}] ${fromCache ? "[cache] " : ""}(rel=${page.relevanceScore.toFixed(2)}) ${page.title.slice(0, 60)}`);
if (results.length >= task.pageBudget) return;
}
if (idx < candidates.length && results.length < task.pageBudget) {
await sleep(BATCH_INTER_FETCH_DELAY_MS);
}
}
}
async function resolveCandidate(
candidate: ScoredCandidate,
task: SwarmTask,
state: SharedCrawlState,
topicKws: ReadonlyArray<string>,
signal: AbortSignal,
metrics: CrawlMetrics,
): Promise<{ page: CrawledSource; fromCache: boolean }> {
metrics.cacheChecks++;
const cached = state.getCachedSource(candidate.url);
if (cached) {
metrics.cacheHits++;
return {
page: adaptCachedSourceForTask(cached, candidate.query, candidate.snippet, task, topicKws),
fromCache: true,
};
}
metrics.fetchAttempts++;
return {
page: await fetchAndExtract(candidate.url, candidate.query, candidate.snippet, task, topicKws, signal, candidate.discoveredBy),
fromCache: false,
};
}
function adaptCachedSourceForTask(
cached: CrawledSource,
query: string,
snippet: string,
task: SwarmTask,
topicKws: ReadonlyArray<string>,
): CrawledSource {
const { domainScore, freshnessScore, tier } = scoreCandidate({ url: cached.url, title: cached.title, snippet: cached.description }, query);
const relevanceScore = computeRelevance(cached.text, cached.title, snippet, topicKws);
return {
...cached,
url: normalizeUrl(cached.url),
finalUrl: normalizeUrl(cached.finalUrl ?? cached.url),
sourceQuery: query,
workerRole: task.role,
workerLabel: task.label,
domainScore,
freshnessScore,
tier: tier as SourceTier,
relevanceScore,
};
}
async function followLinks(
existingSources: ReadonlyArray<CrawledSource>,
task: SwarmTask,
state: SharedCrawlState,
signal: AbortSignal,
status: StatusFn,
warn: WarnFn,
results: CrawledSource[],
errors: string[],
tag: string,
budget: number,
topicKws: ReadonlyArray<string>,
depth: number,
metrics: CrawlMetrics,
llmManager: LlmCallManager
): Promise<number> {
const allLinks = existingSources.flatMap((s) => s.outlinks);
const linkKws = task.queries.join(" ").toLowerCase().split(/\s+/).filter((w) => w.length > 3).slice(0, 12);
const scored = scoreOutlinks(allLinks, linkKws, state.visitedUrls, task.maxLinksToEvaluate);
const toFollow = scored.slice(0, task.maxLinksToFollow);
if (toFollow.length === 0) return 0;
status(`${tag} Following ${toFollow.length} link(s) (depth ${depth})…`);
const before = results.length;
const linkCandidates = toFollow.map((l) => scoreCandidate({ url: l.href, title: "", snippet: "" }, task.queries[0] ?? ""));
await fetchBatch(linkCandidates, { ...task, pageBudget: results.length + budget }, state, signal, status, warn, results, errors, tag, topicKws, metrics, llmManager);
return results.length - before;
}
async function fetchAndExtract(
url: string,
query: string,
snippet: string,
task: SwarmTask,
topicKws: ReadonlyArray<string>,
signal: AbortSignal,
discoveredBy?: string,
): Promise<CrawledSource> {
let page: ExtractedPage;
if (/youtube\.com|youtu\.be/i.test(url)) {
const ytData = await extractYouTubeTranscript(url, signal, task.contentLimit);
if (!ytData) throw new Error("YouTube transcript extraction failed");
page = {
url: url,
finalUrl: url,
title: ytData.title,
description: ytData.description,
published: null,
text: ytData.text,
wordCount: ytData.text.split(/\s+/).filter(Boolean).length,
outlinks: [],
page: 1,
totalPages: 1,
};
} else {
const fetchResult = await fetchWithWaybackFallback(url, signal, task);
const { finalUrl } = fetchResult;
const isPdf = (fetchResult.rawBuffer && isPdfContentType(fetchResult.contentType)) || (!fetchResult.rawBuffer && isPdfUrl(url));
if (isPdf && fetchResult.rawBuffer) {
page = await extractPdf(fetchResult.rawBuffer, url, finalUrl, task.contentLimit, false);
} else if (isPdf && fetchResult.html && fetchResult.html.startsWith("%PDF")) {
const buf = Buffer.from(fetchResult.html, "binary");
page = await extractPdf(buf, url, finalUrl, task.contentLimit, false);
} else {
page = extractPage(fetchResult.html, url, finalUrl, task.contentLimit, task.maxOutlinksPerPage);
}
}
const { domainScore, freshnessScore, tier } = scoreCandidate({ url, title: page.title, snippet: page.description }, query);
const relevanceScore = computeRelevance(page.text, page.title, snippet, topicKws);
return {
url: normalizeUrl(page.url),
finalUrl: normalizeUrl(page.finalUrl),
title: page.title,
description: page.description,
published: page.published,
text: page.text,
wordCount: page.wordCount,
outlinks: page.outlinks,
sourceQuery: query,
workerRole: task.role,
workerLabel: task.label,
domainScore,
freshnessScore,
tier: tier as SourceTier,
relevanceScore,
origin: "web" as const,
discoveredBy: discoveredBy ?? "unknown",
page: page.page,
totalPages: page.totalPages,
};
}
async function getLLMQueryMutation(
query: string,
topSnippets: string[],
signal: AbortSignal,
): Promise<string | null> {
try {
const endpoint = "http://localhost:1234/v1/chat/completions";
const prompt = `You are a deep research assistant. The initial query "${query}" yielded these snippets:\n${topSnippets.slice(0, 3).join("\n")}\n\nGenerate ONE highly specific, alternative search query to find missing technical details or counter-arguments. Return ONLY the raw query string, no quotes or explanations.`;
logLlmDiagnostics("getLLMQueryMutation", prompt);
const res = await fetch(endpoint, {
method: "POST",
headers: { "Content-Type": "application/json" },
signal,
body: JSON.stringify({
model: "local-model",
messages: [{ role: "user", content: prompt }],
temperature: 0.7,
max_tokens: 40,
}),
});
if (!res.ok) return null;
const data = await res.json();
const mutated = data.choices?.[0]?.message?.content?.trim();
return mutated ? mutated.replace(/^[\"']|[\"']$/g, "") : null;
} catch {
return null;
}
}
async function fetchWithWaybackFallback(
url: string,
signal: AbortSignal,
task: SwarmTask // <-- We add task here to grab the config!
): Promise<Awaited<ReturnType<typeof fetchPage>>> {
try {
// TIER A: Standard Fetch (Tries normally first)
return await fetchPage(url, signal);
} catch (err: any) {
if (isAbortError(err)) throw err;
// TIER B: FlareSolverr (Optional Power-User Bypass)
const fsUrl = (task as any).flaresolverrUrl as string | undefined;
if (fsUrl && fsUrl.trim() !== "") {
try {
const fsRes = await fetch(fsUrl, {
method: "POST",
headers: { "Content-Type": "application/json" },
body: JSON.stringify({ cmd: "request.get", url: url, maxTimeout: 15000 }),
signal
});
const fsData = await fsRes.json();
if (fsData?.solution?.response) {
return {
html: fsData.solution.response,
finalUrl: url,
contentType: "text/html",
rawBuffer: undefined // FlareSolverr only returns HTML, not raw PDFs
};
}
} catch (fsErr) {
// Silently fail and drop down to Wayback Machine
}
}
// TIER C: Wayback Machine Archive
try {
const api = `https://archive.org/wayback/available?url=${encodeURIComponent(url)}`;
const apiRes = await fetch(api, { signal });
if (!apiRes.ok) throw new Error("Wayback API failed");
const data = await apiRes.json();
const wbUrl = data?.archived_snapshots?.closest?.url;
if (wbUrl) {
return await fetchPage(wbUrl, signal);
}
} catch {
// ignore
}
throw err;
}
}
function uniqueHostCount(hits: ReadonlyArray<{ url: string }>): number {
const hosts = new Set<string>();
for (const h of hits) {
try {
hosts.add(new URL(h.url).hostname.replace(/^www\./, ""));
} catch {
// ignore
}
}
return hosts.size;
}
function mutationQuality(hits: ReadonlyArray<{ url: string }>): number {
return uniqueHostCount(hits) * 10 + Math.min(hits.length, 10);
}
function deduplicateByUrl<T extends { url: string }>(items: ReadonlyArray<T>): T[] {
const seen = new Set<string>();
return items.filter((item) => {
const key = normalizeUrl(item.url);
if (seen.has(key)) return false;
seen.add(key);
return true;
});
}
function capCandidatesPerHost(
candidates: ReadonlyArray<ScoredCandidate>,
maxPerHost: number,
): ScoredCandidate[] {
const perHost = new Map<string, number>();
const kept: ScoredCandidate[] = [];
for (const c of candidates) {
const host = safeHostname(c.url);
const count = perHost.get(host) ?? 0;
if (count >= maxPerHost) continue;
perHost.set(host, count + 1);
kept.push(c);
}
return kept;
}
function safeHostname(url: string): string {
try {
return new URL(url).hostname.replace(/^www\./, "");
} catch {
return "";
}
}
function truncUrl(url: string, max = 70): string {
return url.length > max ? url.slice(0, max) + "…" : url;
}
function isAbortError(err: unknown): boolean {
return err instanceof DOMException && err.name === "AbortError";
}
function errorMessage(err: unknown): string {
return err instanceof Error ? err.message : String(err ?? "unknown");
}
// ==================== FILE: src\net\ddg.ts ====================
import { fetchPage } from "./http";
export class DdgRateLimiter {
private lastRequest = 0;
private minDelay: number;
constructor(minDelayMs: number = 2500) {
this.minDelay = minDelayMs;
}
async acquire(): Promise<void> {
const now = Date.now();
const elapsed = now - this.lastRequest;
if (elapsed < this.minDelay) {
await new Promise((resolve) => setTimeout(resolve, this.minDelay - elapsed));
}
this.lastRequest = Date.now();
}
}
export const sharedDdgLimiter = new DdgRateLimiter(2500);
export class DdgLimiterPool {
private limiters: DdgRateLimiter[] = [];
private currentIndex = 0;
constructor(numLanes: number, minDelayMs: number = 2500) {
for (let i = 0; i < numLanes; i++) {
this.limiters.push(new DdgRateLimiter(minDelayMs));
}
}
next(): DdgRateLimiter {
const limiter = this.limiters[this.currentIndex];
this.currentIndex = (this.currentIndex + 1) % this.limiters.length;
return limiter;
}
}
export function resetThrottle(): void {}
/**
* Bulletproof multi-tier DDG Search Waterfall:
* Tier A: DDG Lite POST endpoint (fastest, mimics ddgr)
* Tier B: DDG HTML fallback endpoint (handles strict blocks)
* Tier C: Graceful recovery (returns empty array instead of throwing crash errors)
*/
export async function searchDDG(
query: string,
maxResults: number,
safeSearch: "strict" | "moderate" | "off" = "moderate",
signal?: AbortSignal,
limiter?: DdgRateLimiter,
timeRange: string = "all"
): Promise<ReadonlyArray<{ url: string; title: string; snippet: string }>> {
if (limiter) await limiter.acquire();
if (signal?.aborted) throw new DOMException("Aborted", "AbortError");
const safeParam = safeSearch === "strict" ? "1" : safeSearch === "off" ? "-1" : "0";
const dfParam = timeRange === "all" ? "" : `&df=${timeRange}`;
const formData = `q=${encodeURIComponent(query)}&kp=${safeParam}${dfParam}`;
// TIER A: DDG Lite Endpoint
try {
const res = await fetchPage("https://lite.duckduckgo.com/lite/", signal!, {
method: "POST",
body: formData,
headers: {
"Content-Type": "application/x-www-form-urlencoded",
"Referer": "https://lite.duckduckgo.com/",
"User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/124.0.0.0 Safari/537.36",
"Accept": "text/html,application/xhtml+xml,application/xml;q=0.9,image/avif,image/webp,image/apng,*/*;q=0.8",
"Accept-Language": "en-US,en;q=0.9",
"Sec-Ch-Ua": '"Chromium";v="124", "Google Chrome";v="124", "Not-A.Brand";v="99"',
"Sec-Ch-Ua-Mobile": "?0",
"Sec-Ch-Ua-Platform": '"Windows"',
"Sec-Fetch-Dest": "document",
"Sec-Fetch-Mode": "navigate",
"Sec-Fetch-Site": "same-origin",
"Sec-Fetch-User": "?1",
"Upgrade-Insecure-Requests": "1"
}
});
const hits = parseDDGResults(res.html, maxResults);
if (hits.length > 0) return hits;
} catch (err) {
if (signal?.aborted) throw err;
console.warn(`[DDG Tier A] Failed for query "${query}": ${err instanceof Error ? err.message : String(err)}. Trying Tier B...`);
}
// TIER B: DDG HTML Endpoint Fallback
try {
if (limiter) await limiter.acquire();
const htmlFallbackUrl = `https://html.duckduckgo.com/html/?q=${encodeURIComponent(query)}`;
const resB = await fetchPage(htmlFallbackUrl, signal!, {
method: "GET",
headers: {
"Content-Type": "application/x-www-form-urlencoded",
"Referer": "https://lite.duckduckgo.com/",
"User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/124.0.0.0 Safari/537.36",
"Accept": "text/html,application/xhtml+xml,application/xml;q=0.9,image/avif,image/webp,image/apng,*/*;q=0.8",
"Accept-Language": "en-US,en;q=0.9",
"Sec-Ch-Ua": '"Chromium";v="124", "Google Chrome";v="124", "Not-A.Brand";v="99"',
"Sec-Ch-Ua-Mobile": "?0",
"Sec-Ch-Ua-Platform": '"Windows"',
"Sec-Fetch-Dest": "document",
"Sec-Fetch-Mode": "navigate",
"Sec-Fetch-Site": "same-origin",
"Sec-Fetch-User": "?1",
"Upgrade-Insecure-Requests": "1"
}
});
const hitsB = parseHTMLDDGResults(resB.html, maxResults);
if (hitsB.length > 0) return hitsB;
} catch (errB) {
if (signal?.aborted) throw errB;
console.warn(`[DDG Tier B] Fallback also failed for query "${query}": ${errB instanceof Error ? errB.message : String(errB)}`);
}
// TIER C: Graceful exit (Returns empty array so the worker's outer health system handles it cleanly without crashing)
return [];
}
export async function searchDDGPaginated(
query: string,
maxResultsPerPage: number,
pages: number,
safeSearch: "strict" | "moderate" | "off" = "moderate",
signal?: AbortSignal,
limiter?: DdgRateLimiter,
timeRange: string = "all"
): Promise<ReadonlyArray<{ url: string; title: string; snippet: string }>> {
const allHits: { url: string; title: string; snippet: string }[] = [];
for (let p = 1; p <= pages; p++) {
if (signal?.aborted) break;
const safeParam = safeSearch === "strict" ? "1" : safeSearch === "off" ? "-1" : "0";
const dfParam = timeRange === "all" ? "" : `&df=${timeRange}`;
const formData = `q=${encodeURIComponent(query)}&kp=${safeParam}${dfParam}&s=${(p - 1) * maxResultsPerPage}`;
try {
if (limiter) await limiter.acquire();
const res = await fetchPage("https://lite.duckduckgo.com/lite/", signal!, {
method: "POST",
body: formData,
headers: {
"Content-Type": "application/x-www-form-urlencoded",
"Referer": "https://lite.duckduckgo.com/"
}
});
const hits = parseDDGResults(res.html, maxResultsPerPage);
allHits.push(...hits);
if (hits.length < maxResultsPerPage) break;
} catch {
break;
}
}
return allHits;
}
function parseDDGResults(html: string, maxResults: number): { url: string; title: string; snippet: string }[] {
const hits: { url: string; title: string; snippet: string }[] = [];
const seen = new Set<string>();
// Resilient regex pattern matching DDG Lite result rows
const resultRe = /<a[^>]*rel="nofollow"[^>]*href="([^"]+)"[^>]*>([\s\S]*?)<\/a>[\s\S]*?<td class="result-snippet">([\s\S]*?)<\/td>/gi;
let match: RegExpExecArray | null;
while (hits.length < maxResults && (match = resultRe.exec(html)) !== null) {
let url = match[1];
const title = match[2].replace(/<[^>]+>/g, "").trim();
const snippet = match[3].replace(/<[^>]+>/g, "").trim();
// Clean up redirect wrappers if present
if (url.includes("uddg=")) {
const matchUrl = url.match(/uddg=([^&]+)/);
if (matchUrl) url = decodeURIComponent(matchUrl[1]);
}
if (url.includes("duckduckgo.com")) continue;
if (!url.startsWith("http")) continue;
if (seen.has(url)) continue;
seen.add(url);
hits.push({ url, title, snippet });
}
return hits;
}
function parseHTMLDDGResults(html: string, maxResults: number): { url: string; title: string; snippet: string }[] {
const hits: { url: string; title: string; snippet: string }[] = [];
const seen = new Set<string>();
// Secondary parser for standard DDG HTML results page layout
const resultRe = /<a class="result__url" href="([^"]+)"[^>]*>([\s\S]*?)<\/a>[\s\S]*?<a class="result__snippet"[^>]*>([\s\S]*?)<\/a>/gi;
let match: RegExpExecArray | null;
while (hits.length < maxResults && (match = resultRe.exec(html)) !== null) {
let url = match[1];
const title = match[2].replace(/<[^>]+>/g, "").trim();
const snippet = match[3].replace(/<[^>]+>/g, "").trim();
if (url.includes("duckduckgo.com")) continue;
if (!url.startsWith("http") && url.startsWith("//")) url = "https:" + url;
if (!url.startsWith("http")) continue;
if (seen.has(url)) continue;
seen.add(url);
hits.push({ url, title, snippet });
}
return hits;
}
// ==================== FILE: src\net\http.ts ====================
// src/net/http.ts
import * as https from "node:https";
import * as http from "node:http";
import { setServers } from "node:dns";
const DNS_RESOLVERS = ["1.1.1.1", "1.0.0.1", "8.8.8.8", "8.8.4.4", "9.9.9.9"];
setServers(DNS_RESOLVERS);
const UA_POOL: ReadonlyArray<string> = [
"Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/124.0.0.0 Safari/537.36",
"Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/123.0.0.0 Safari/537.36",
"Mozilla/5.0 (Windows NT 10.0; Win64; x64; rv:125.0) Gecko/20100101 Firefox/125.0",
"Mozilla/5.0 (X11; Linux x86_64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/124.0.0.0 Safari/537.36",
];
function randomUA(): string {
return UA_POOL[Math.floor(Math.random() * UA_POOL.length)];
}
// SSRF Guard: Block internal/local IPs
function isPrivateUrl(url: string): boolean {
try {
const u = new URL(url);
const host = u.hostname;
if (host === "localhost" || host === "127.0.0.1" || host === "0.0.0.0") return true;
if (host.startsWith("10.") || host.startsWith("192.168.")) return true;
if (host.startsWith("172.")) {
const parts = host.split(".");
const second = parseInt(parts[1], 10);
if (second >= 16 && second <= 31) return true;
}
return false;
} catch {
return true;
}
}
export function buildBrowserHeaders(url?: string): Record<string, string> {
const ua = randomUA();
const headers: Record<string, string> = {
"User-Agent": ua,
Accept: "text/html,application/xhtml+xml,application/xml;q=0.9,image/avif,image/webp,*/*;q=0.8",
"Accept-Language": "en-US,en;q=0.9",
"Accept-Encoding": "gzip, deflate, br",
DNT: "1",
Connection: "keep-alive",
"Upgrade-Insecure-Requests": "1",
"Sec-Fetch-Dest": "document",
"Sec-Fetch-Mode": "navigate",
"Sec-Fetch-Site": "none",
"Sec-Fetch-User": "?1",
};
if (url) {
try {
const parsed = new URL(url);
headers["Referer"] = `${parsed.protocol}//${parsed.hostname}/`;
} catch {}
}
return headers;
}
export function buildDDGHeaders(): Record<string, string> {
return buildBrowserHeaders();
}
export interface FetchResult {
readonly html: string;
readonly finalUrl: string;
readonly contentType?: string;
readonly rawBuffer?: Buffer;
}
const FETCH_TIMEOUT_MS = 12000;
const FETCH_MAX_RETRIES = 3;
// Proxy support (optional)
const PROXY_URL = process.env.HTTP_PROXY || process.env.HTTPS_PROXY;
const proxyAgent = PROXY_URL ? new (require("https-proxy-agent").HttpsProxyAgent)(PROXY_URL) : undefined;
export interface FetchOptions {
method?: 'GET' | 'POST';
body?: string;
headers?: Record<string, string>;
}
export async function fetchPage(
url: string,
signal: AbortSignal,
options: FetchOptions = {}
): Promise<FetchResult> {
if (isPrivateUrl(url)) throw new Error("SSRF Guard: Blocked internal URL");
let lastError: unknown;
const { method = 'GET', body, headers: customHeaders = {} } = options;
for (let attempt = 0; attempt < FETCH_MAX_RETRIES; attempt++) {
if (signal.aborted) throw new DOMException("Aborted", "AbortError");
// Delay to prevent burst blocking
await new Promise((resolve) => setTimeout(resolve, 2500));
try {
const headers = { ...buildBrowserHeaders(url), ...customHeaders };
const controller = new AbortController();
const timeout = setTimeout(() => controller.abort(), FETCH_TIMEOUT_MS);
const res = await fetch(url, {
method,
body,
signal: controller.signal,
headers,
redirect: "follow"
});
clearTimeout(timeout);
// If blocked by Cloudflare/WAF, fall back to Wayback Machine immediately
if ([403, 429, 503].includes(res.status)) {
throw new Error(`HTTP ${res.status}`);
}
const contentType = res.headers.get("content-type") || "";
const finalUrl = res.url || url;
if (contentType.includes("application/pdf") || contentType.includes("application/octet-stream")) {
const arrayBuf = await res.arrayBuffer();
return { html: "", finalUrl, contentType, rawBuffer: Buffer.from(arrayBuf) };
}
const html = await res.text();
// Only check for very specific captcha markers, avoid generic words like "blocked"
if (html.includes("cf-challenge") || html.includes("hcaptcha") || html.includes("recaptcha")) {
throw new Error("Captcha page detected");
}
return { html, finalUrl, contentType };
} catch (err: unknown) {
if (err instanceof DOMException && err.name === "AbortError") throw err;
lastError = err;
// Exponential backoff
const backoffTime = 5000 * Math.pow(2, attempt);
await new Promise((resolve) => setTimeout(resolve, backoffTime));
}
}
// If direct fetch fails, try 10 Web Archives
return fetchFromArchives(url, signal);
}
async function fetchFromArchives(url: string, signal: AbortSignal): Promise<FetchResult> {
const encoded = encodeURIComponent(url);
const archives = [
`https://web.archive.org/web/2/${url}`,
`https://webcache.googleusercontent.com/search?q=cache:${encoded}&strip=1`,
`https://archive.today/newest/${url}`,
`https://cc.bingj.com/cache.aspx?q=${url}`,
`https://webcache.googleusercontent.com/search?q=cache:${url}`,
`http://web.archive.org/web/2024/${url}`,
`https://freezedry.com/${url}`,
`https://corsproxy.io/?${encoded}`,
`https://api.allorigins.win/raw?url=${encoded}`,
`https://cors-anywhere.herokuapp.com/${url}`
];
for (const cacheUrl of archives) {
if (signal.aborted) throw new DOMException("Aborted", "AbortError");
try {
const controller = new AbortController();
const timeout = setTimeout(() => controller.abort(), 8000);
const res = await fetch(cacheUrl, {
signal: controller.signal,
headers: buildBrowserHeaders(cacheUrl)
});
clearTimeout(timeout);
if (res.ok) {
const html = await res.text();
if (html.length > 500) {
return { html, finalUrl: url, contentType: res.headers.get("content-type") || "" };
}
}
} catch {}
}
throw new Error(`Failed to fetch ${url}: Blocked and all archives failed`);
}
export function safeHostname(url: string): string {
try {
return new URL(url).hostname;
} catch {
return "";
}
}
export function sleep(ms: number): Promise<void> {
return new Promise((resolve) => setTimeout(resolve, ms));
}
// ==================== FILE: src\planning\planner.ts ====================
import { logLlmDiagnostics } from "../utils/tokens";
import { LMStudioClient } from "@lmstudio/sdk";
import {
QueryPlan,
WorkerRole,
DynamicWorkerSpec,
AdaptiveGapPlan,
CrawledSource,
AgentMessage,
StatusFn,
} from "../types";
import { DIMENSIONS, detectGaps, gapFillQueries } from "./dimensions";
import {
DepthProfile,
AI_PLANNING_MAX_TOKENS,
AI_PLANNING_TEMPERATURE,
AI_PLANNING_TIMEOUT_MS,
AI_MIN_ACCEPTABLE_QUERIES,
AI_DECOMPOSITION_MAX_TOKENS,
AI_DECOMPOSITION_TEMPERATURE,
AI_DECOMPOSITION_TIMEOUT_MS,
AI_FINDINGS_SUMMARY_MAX_TOKENS,
AI_FINDINGS_SUMMARY_TEMPERATURE,
FINDINGS_SUMMARY_SOURCE_CHARS,
DECOMPOSITION_MIN_WORKERS,
QUERY_LINE_MIN_LEN,
QUERY_LINE_MAX_LEN,
SYSTEM_INSTRUCTIONS,
} from "../constants";
async function callLoadedModel(
prompt: string,
maxTokens: number = AI_PLANNING_MAX_TOKENS,
temperature: number = AI_PLANNING_TEMPERATURE,
timeoutMs: number = AI_PLANNING_TIMEOUT_MS,
): Promise<string | null> {
try {
const client = new LMStudioClient();
const models = await Promise.race<Awaited<ReturnType<typeof client.llm.listLoaded>>>([
client.llm.listLoaded(),
new Promise<never>((_, reject) =>
setTimeout(() => reject(new Error("timeout")), timeoutMs),
),
]);
if (!Array.isArray(models) || models.length === 0) return null;
const model = await client.llm.model(models[0].identifier);
const stream = model.respond(
[
{ role: "system", content: SYSTEM_INSTRUCTIONS },
{ role: "user", content: prompt },
],
{
maxTokens,
temperature,
},
);
let result = "";
for await (const chunk of stream) result += chunk.content ?? "";
return result.trim() || null;
} catch {
return null;
}
}
function parseLines(raw: string, maxLines: number = 6): ReadonlyArray<string> {
return raw
.split(/\n/)
.map((line) => line.replace(/^\d+[.)]\s*|^[-*•]\s*/, "").trim())
.filter(
(line) =>
line.length > QUERY_LINE_MIN_LEN && line.length < QUERY_LINE_MAX_LEN,
)
.filter((line, idx, arr) => arr.indexOf(line) === idx)
.slice(0, maxLines);
}
const VALID_ROLES: ReadonlyArray<WorkerRole> = [
"breadth",
"depth",
"recency",
"academic",
"critical",
"statistical",
"regulatory",
"technical",
"primary",
"comparative",
];
function makeDecompositionPrompt(
topic: string,
focusAreas: ReadonlyArray<string>,
profile: DepthProfile,
): string {
const focus = focusAreas.length
? `\nFocus areas: ${focusAreas.join(", ")}`
: "";
return `You are a research decomposition system. Given a research topic, output a strict JSON array of specialized worker agents.
Topic: "${topic}"${focus}
Each worker object MUST have:
"role": one of "breadth", "depth", "recency", "academic", "critical", "statistical", "regulatory", "technical", "primary", "comparative"
"label": descriptive name (e.g., "Clinical Evidence Researcher", "Policy Critic")
"queries": array of ${Math.min(profile.maxQueriesPerWorker, 6)}-${profile.maxQueriesPerWorker} specific, natural search queries
"budgetWeight": number 0.1-0.4
"followLinks": boolean
"preferredTiers": array of strings
Rules:
1. Output ${DECOMPOSITION_MIN_WORKERS} to ${profile.maxDecompositionWorkers} workers.
2. Queries must be human-like search strings (under 8 words). DO NOT just mash the topic and role together.
3. Output ONLY valid JSON. No conversational text.
JSON:`;
}
async function aiDecompose(
topic: string,
focusAreas: ReadonlyArray<string>,
status: StatusFn,
profile: DepthProfile,
): Promise<ReadonlyArray<DynamicWorkerSpec> | null> {
const raw = await callLoadedModel(
makeDecompositionPrompt(topic, focusAreas, profile),
AI_DECOMPOSITION_MAX_TOKENS,
AI_DECOMPOSITION_TEMPERATURE,
AI_DECOMPOSITION_TIMEOUT_MS,
);
if (!raw) return null;
try {
const jsonStr = raw.replace(/```json\s*|```\s*/g, "").trim();
const parsed = JSON.parse(jsonStr);
if (!Array.isArray(parsed) || parsed.length < DECOMPOSITION_MIN_WORKERS)
return null;
const specs: DynamicWorkerSpec[] = [];
for (const item of parsed.slice(0, profile.maxDecompositionWorkers)) {
const role = VALID_ROLES.includes(item.role) ? item.role : "breadth";
const queries = Array.isArray(item.queries)
? item.queries
.filter((q: unknown) => typeof q === "string" && q.length > 3)
.slice(0, profile.maxQueriesPerWorker)
: [];
if (queries.length < 2) continue;
specs.push({
role: role as WorkerRole,
label:
typeof item.label === "string"
? item.label.slice(0, 60)
: `${role} worker`,
queries,
budgetWeight:
typeof item.budgetWeight === "number"
? Math.max(0.05, Math.min(0.5, item.budgetWeight))
: 0.2,
followLinks: item.followLinks === true,
preferredTiers: Array.isArray(item.preferredTiers)
? item.preferredTiers
: undefined,
});
}
if (specs.length < DECOMPOSITION_MIN_WORKERS) return null;
const totalWeight = specs.reduce((sum, s) => sum + s.budgetWeight, 0);
const normalised = specs.map((s) => ({
...s,
budgetWeight: s.budgetWeight / totalWeight,
}));
status(`AI decomposed topic into ${normalised.length} specialised workers`);
return normalised;
} catch {
return null;
}
}
function makeRolePlanPrompt(
role: WorkerRole,
topic: string,
focusAreas: ReadonlyArray<string>,
profile: DepthProfile,
): string {
const roleDescriptions: Readonly<Record<WorkerRole, string>> = {
breadth: "broad coverage - primary facts, key individuals, organizations, and canonical entities",
depth: "deep dive - experimental methodologies, field investigations, and primary case data",
recency: "recent developments - publications, peer reviews, and verifiable records from 2024-2026",
academic: "academic papers - peer-reviewed studies, university archives, and institutional research",
critical: "skeptical analysis - methodological critiques, counter-arguments, fraud investigations, and alternative explanations",
statistical: "statistics and quantitative findings - sample sizes, replication rates, and surveys",
regulatory: "governing standards, institutional positions, and formal guidelines",
technical: "technical mechanisms and specific evaluation frameworks",
primary: "primary documentation - field notes, first-hand interviews, and official recordings",
comparative: "comparative models - alternative scientific, physiological, or psychological hypotheses",
};
const focus = focusAreas.length ? `\nFocus areas: ${focusAreas.join(", ")}` : "";
return `You are an expert investigative search planner.
Target Topic: "${topic}"${focus}
Assigned Role: ${roleDescriptions[role]}
Generate exactly ${profile.maxQueriesPerWorker} precision search queries.
STRICT QUERY RULES:
1. Target canonical entities, notable researchers, research institutes, and key published works directly.
2. If the topic is empirical/scientific, append negative terms when appropriate (e.g., -livestream -forum -llm) to filter junk.
3. If this role is "critical", explicitly hunt for critiques, failed replications, and skeptics.
4. Output ONLY the queries, one per line. No numbering, quotes, or Markdown formatting.
Queries:`;
}
const ROLE_DIMENSIONS: Readonly<Record<WorkerRole, ReadonlyArray<string>>> = {
breadth: ["overview", "applications", "history", "economics"],
depth: ["mechanism", "evidence", "expert"],
recency: ["current", "future"],
academic: ["evidence", "expert", "mechanism"],
critical: ["challenges", "controversy", "comparison"],
statistical: ["evidence", "economics", "overview"],
regulatory: ["challenges", "controversy", "current"],
technical: ["mechanism", "applications", "evidence"],
primary: ["evidence", "expert", "history"],
comparative: ["comparison", "challenges", "applications"],
};
function dimensionFallbackQueries(
role: WorkerRole,
topic: string,
focusAreas: ReadonlyArray<string>,
maxQueries: number,
): ReadonlyArray<string> {
const dimIds = ROLE_DIMENSIONS[role];
const dims = DIMENSIONS.filter((d) => dimIds.includes(d.id));
const queries: string[] = [];
// BUGFIX: Actually use the shortenTopic function to strip the massive prompt
const shortTopic = shortenTopic(topic);
for (const dim of dims) {
for (const q of dim.queries(shortTopic)) {
if (!queries.includes(q)) queries.push(q);
}
}
for (const area of focusAreas) {
const q = `${shortTopic} ${area}`;
if (!queries.includes(q)) queries.push(q);
}
return queries.slice(0, maxQueries);
}
export async function buildQueryPlan(
topic: string,
focusAreas: ReadonlyArray<string>,
useAI: boolean,
status: StatusFn,
profile: DepthProfile,
): Promise<QueryPlan> {
const CORE_ROLES: ReadonlyArray<WorkerRole> = [
"breadth",
"depth",
"recency",
"academic",
"critical",
];
const EXTENDED_ROLES: ReadonlyArray<WorkerRole> = [
"statistical",
"regulatory",
"technical",
"primary",
"comparative",
];
let roles: ReadonlyArray<WorkerRole>;
if (profile.depthRounds >= 10) {
roles = [...CORE_ROLES, ...EXTENDED_ROLES];
} else if (profile.depthRounds >= 5) {
roles = [...CORE_ROLES, "technical", "comparative", "statistical"];
} else {
roles = CORE_ROLES;
}
const queriesByRole: Partial<Record<WorkerRole, ReadonlyArray<string>>> = {};
let usedAI = false;
let dynamicSpecs: ReadonlyArray<DynamicWorkerSpec> | undefined;
if (useAI) {
status("AI task decomposition - analysing topic for specialised workers...");
const specs = await aiDecompose(topic, focusAreas, status, profile);
if (specs && specs.length >= DECOMPOSITION_MIN_WORKERS) {
dynamicSpecs = specs;
usedAI = true;
for (const spec of specs) {
queriesByRole[spec.role] = spec.queries;
}
} else {
status("AI planning queries for each swarm worker...");
}
const uncoveredRoles = roles.filter((r) => !queriesByRole[r]?.length);
if (uncoveredRoles.length > 0) {
const results = await Promise.allSettled(
uncoveredRoles.map(async (role) => ({
role,
queries: await callLoadedModel(
makeRolePlanPrompt(role, topic, focusAreas, profile),
),
})),
);
for (const result of results) {
if (result.status !== "fulfilled") continue;
const { role, queries: raw } = result.value;
if (!raw) continue;
const parsed = parseLines(raw, profile.maxQueriesPerWorker);
if (parsed.length >= AI_MIN_ACCEPTABLE_QUERIES) {
queriesByRole[role] = parsed;
usedAI = true;
}
}
}
if (usedAI) {
status(
`AI generated queries for ${Object.keys(queriesByRole).length} worker role(s)`,
);
} else {
status("AI unavailable, using dimension-based query planning");
}
}
for (const role of roles) {
if (!queriesByRole[role] || queriesByRole[role]!.length === 0) {
queriesByRole[role] = dimensionFallbackQueries(
role,
topic,
focusAreas,
profile.maxQueriesPerWorker,
);
}
}
const topicKeywords = extractKeywords(topic);
// PRE-SEARCH PLANNING (If AI is enabled)
if (useAI) {
const planPrompt = `Topic: ${topic}\nFocus: ${focusAreas.join(", ")}\n
Generate a JSON object with:
1. "subQuestions": 3-5 specific questions to answer.
2. "likelySourceTypes": Array of strings (e.g., "academic", "news", "reference").
3. "stopConditions": Array of strings (e.g., "Found 3 sources confirming Q3 revenue").`;
logLlmDiagnostics("buildQueryPlan-PreSearch", planPrompt);
}
return {
queriesByRole: queriesByRole as Record<WorkerRole, ReadonlyArray<string>>,
usedAI,
topicKeywords,
dynamicSpecs,
};
}
export async function summariseFindings(
sources: ReadonlyArray<CrawledSource>,
topic: string,
useAI: boolean,
status: StatusFn,
): Promise<ReadonlyArray<AgentMessage>> {
if (!useAI || sources.length === 0) return [];
const sourceSummaries = sources
.slice(0, 20)
.map(
(s, i) =>
`[${i + 1}] ${s.workerLabel}: ${s.title} - ${s.text.slice(0, FINDINGS_SUMMARY_SOURCE_CHARS)}`,
)
.join("\n\n");
const prompt = `You are a research coordinator. A team of research agents collected these sources on "${topic}":
${sourceSummaries}
Summarise:
The 3-5 most important findings discovered so far (one line each)
3-5 specific questions or angles that were NOT covered and need follow-up
Output format:
FINDINGS:
finding 1
finding 2
...
FOLLOW_UP:
question 1
question 2
...`;
const raw = await callLoadedModel(
prompt,
AI_FINDINGS_SUMMARY_MAX_TOKENS,
AI_FINDINGS_SUMMARY_TEMPERATURE,
);
if (!raw) return [];
const findings: string[] = [];
const followUps: string[] = [];
let section: "findings" | "followup" | null = null;
for (const line of raw.split("\n")) {
const trimmed = line.trim();
if (/^FINDINGS:/i.test(trimmed)) {
section = "findings";
continue;
}
if (/^FOLLOW.?UP:/i.test(trimmed)) {
section = "followup";
continue;
}
const item = trimmed.replace(/^[-*•]\s*/, "").trim();
if (item.length < 5) continue;
if (section === "findings") findings.push(item);
if (section === "followup") followUps.push(item);
}
if (findings.length === 0 && followUps.length === 0) return [];
status(
`AI summarised ${findings.length} key findings, ${followUps.length} follow-up suggestions`,
);
return [
{
fromWorker: "round-coordinator",
keyFindings: findings,
suggestedFollowUps: followUps,
},
];
}
const GAP_ROLE_MAP: Readonly<
Record<
string,
{
role: WorkerRole;
followLinks: boolean;
tiers?: ReadonlyArray<import("../types").SourceTier>;
}
>
> = {
overview: { role: "breadth", followLinks: false },
mechanism: { role: "technical", followLinks: true },
history: { role: "breadth", followLinks: false },
current: { role: "recency", followLinks: false },
applications: { role: "breadth", followLinks: false },
challenges: { role: "critical", followLinks: false },
comparison: { role: "comparative", followLinks: false },
evidence: {
role: "academic",
followLinks: true,
tiers: ["academic", "government", "reference"],
},
expert: { role: "primary", followLinks: true, tiers: ["academic", "news"] },
future: { role: "recency", followLinks: false },
controversy: { role: "critical", followLinks: false },
economics: { role: "statistical", followLinks: false },
};
export async function buildAdaptiveGapFill(
topic: string,
coveredIds: ReadonlyArray<string>,
priorMessages: ReadonlyArray<AgentMessage>,
useAI: boolean,
status: StatusFn,
profile: DepthProfile,
): Promise<ReadonlyArray<AdaptiveGapPlan>> {
const gaps = detectGaps(coveredIds);
if (gaps.length === 0) {
status("All research dimensions covered - no gap queries needed");
return [];
}
status(`Gaps: ${gaps.map((g) => g.label).join(", ")}`);
const byRole = new Map<
WorkerRole,
{
dimIds: string[];
dimLabels: string[];
queries: string[];
followLinks: boolean;
tiers?: ReadonlyArray<import("../types").SourceTier>;
}
>();
// BUGFIX: Shorten the topic here too so the fallback queries are clean
const shortTopic = shortenTopic(topic);
for (const gap of gaps) {
const mapping = GAP_ROLE_MAP[gap.id] ?? {
role: "breadth" as WorkerRole,
followLinks: false,
};
const existing = byRole.get(mapping.role) ?? {
dimIds: [],
dimLabels: [],
queries: [],
followLinks: mapping.followLinks,
tiers: mapping.tiers,
};
existing.dimIds.push(gap.id);
existing.dimLabels.push(gap.label);
existing.queries.push(...gap.queries(shortTopic));
byRole.set(mapping.role, existing);
}
if (useAI) {
const followUpContext = priorMessages
.flatMap((m) => m.suggestedFollowUps)
.slice(0, 6);
const entries = Array.from(byRole.entries());
const aiResults = await Promise.allSettled(
entries.map(async ([role, group]) => {
const queryCount = Math.min(
group.dimLabels.length * 3,
profile.maxGapFillQueries,
);
const prompt = `You are an expert Google searcher. A research session on "${topic}" is missing these angles:
${group.dimLabels.join(", ")}
${followUpContext.length > 0 ? `Previous round suggested exploring:\n${followUpContext.join("\n")}\n` : ""}
Generate ${queryCount} specific, natural search engine queries to fill these gaps.
Rules:
1. DO NOT just append words to the topic string. Write human-like search queries.
2. Keep queries short and concise.
3. Return ONLY the queries, one per line. No numbering, no prefixes.
Queries:`;
return { role, raw: await callLoadedModel(prompt) };
}),
);
for (const result of aiResults) {
if (result.status !== "fulfilled" || !result.value.raw) continue;
const { role, raw } = result.value;
const parsed = parseLines(raw, profile.maxGapFillQueries);
if (parsed.length >= 2) {
const group = byRole.get(role);
if (group) group.queries = [...parsed];
}
}
}
const plans: AdaptiveGapPlan[] = [];
for (const [role, group] of byRole) {
plans.push({
role,
label: `Gap-fill: ${group.dimLabels.slice(0, 3).join(", ")}`,
queries: group.queries.slice(0, profile.maxGapFillQueries),
followLinks: group.followLinks,
preferredTiers: group.tiers,
});
}
status(`${plans.length} adaptive gap-fill worker(s) planned`);
return plans;
}
const STOP_WORDS = new Set([
"the", "a", "an", "is", "in", "of", "and", "or", "for", "to", "how", "what",
"why", "when", "does", "with", "from", "that", "could", "which", "about",
"their", "this", "these", "those", "would", "should", "current",
"hypothetical", "scenarios", "lead",
]);
function extractKeywords(topic: string): ReadonlyArray<string> {
return topic
.toLowerCase()
.replace(/[^a-z0-9\s]/g, "")
.split(/\s+/)
.filter((w) => w.length > 2 && !STOP_WORDS.has(w))
.slice(0, 8);
}
function shortenTopic(topic: string): string {
const colonIdx = topic.indexOf(":");
const dashIdx = topic.indexOf(" - ");
const sepIdx = colonIdx > 3 ? colonIdx : dashIdx > 3 ? dashIdx : -1;
let core: string;
if (sepIdx > 3 && sepIdx < topic.length * 0.6) {
core = topic.slice(0, sepIdx).trim();
} else {
core = topic;
}
const words = core
.replace(/[,;()]/g, " ")
.split(/\s+/)
.filter((w) => w.length > 1 && !STOP_WORDS.has(w.toLowerCase()));
let result = "";
let count = 0;
for (const w of words) {
if (count >= 6 || result.length + w.length > 58) break;
result += (result ? " " : "") + w;
count++;
}
if (result.length < 5) {
result = topic.split(/\s+/).slice(0, 5).join(" ");
}
return result;
}
// ==================== FILE: src\swarm\worker.ts ====================
// src/swarm/worker.ts
import { LlmCallManager } from "../utils/llm";
import { logLlmDiagnostics } from "../utils/tokens";
import { extractYouTubeTranscript } from "../net/youtube-extractor";
import {
searchSerper,
searchBraveApi,
searchOpenAlex,
searchCrossref,
searchArxiv,
searchGdelt,
searchReferenceSites,
} from "../net/api-engines";
import {
searchDDG,
searchDDGPaginated,
DdgRateLimiter,
sharedDdgLimiter,
} from "../net/ddg";
import { multiEngineSearch, SearchEngine } from "../net/search-engines";
import { SearchHealthTracker } from "./health";
import { fetchPage, sleep } from "../net/http";
import {
extractPage,
contentFingerprint,
computeRelevance,
} from "../net/extractor";
import {
isPdfUrl,
isPdfContentType,
extractPdf,
} from "../net/pdf-extractor";
import {
scoreCandidate,
rankCandidates,
scoreOutlinks,
} from "../scoring/authority";
import {
SwarmTask,
WorkerResult,
CrawledSource,
ScoredCandidate,
SourceTier,
StatusFn,
WarnFn,
SearchHit,
ExtractedPage,
} from "../types";
import { harvestLocalSources } from "../local/search";
import {
BATCH_INTER_FETCH_DELAY_MS,
MIN_USEFUL_WORD_COUNT,
} from "../constants";
import { normalizeUrl } from "./visited-cache";
export interface CrawlMetrics {
ddgQueries: number;
ddgHits: number;
mutatedQueriesTried: number;
mutationAccepted: number;
mutationHits: number;
extraEngineQueries: number;
extraEngineHits: number;
rawHits: number;
dedupedHits: number;
rankedCandidates: number;
fetchCandidates: number;
fetchAttempts: number;
fetchFailures: number;
acceptedSources: number;
skippedLowWordCount: number;
skippedOffTopic: number;
skippedVeryOffTopic: number;
skippedDuplicateContent: number;
skippedVisited: number;
skippedDomainCap: number;
skippedAvoided: number;
skippedBlacklisted: number;
cacheChecks: number;
cacheHits: number;
cacheAccepted: number;
cacheRejectedDuplicate: number;
cacheRejectedOffTopic: number;
cacheRejectedLowWordCount: number;
cacheWrites: number;
followedLinks: number;
crossWorkerDiscoveriesUsed: number;
localSourcesAccepted: number;
}
export interface SharedCrawlState {
readonly visitedUrls: ReadonlySet<string>;
readonly contentHashes: ReadonlySet<string>;
readonly domainCounts: ReadonlyMap<string, number>;
readonly domainFailures: ReadonlyMap<string, number>;
addVisited(url: string): void;
addHash(hash: string): void;
incrementDomain(url: string): void;
domainCount(url: string): number;
noteFailure(url: string, reason: string): void;
noteDomainFailure(url: string): void;
isDomainBlacklisted(url: string): boolean;
shouldAvoidUrl(url: string): boolean;
pushDiscovery(url: string, title: string, fromWorker: string): void;
drainDiscoveries(limit: number): ReadonlyArray<{ url: string; title: string }>;
isRecentlyVisited(url: string): boolean;
getCachedSource(url: string): CrawledSource | null;
markVisitedPersistent(source: CrawledSource): void;
addWebCacheDocument(source: CrawledSource): void;
getMetricsSnapshot(): Readonly<CrawlMetrics>;
mergeMetrics(delta: Partial<CrawlMetrics>): void;
}
const MUTATION_STRATEGIES: ReadonlyArray<(q: string) => string> = [
(q) => `"${q}"`,
(q) => `${q} explained`,
(q) => `${q} guide overview`,
(q) => `${q} research`,
(q) => q.split(" ").slice(0, 10).join(" "),
(q) => q.replace(/\b(how|what|why|when)\b/gi, " ").trim(),
];
type SearchHitLike = {
url: string;
title: string;
snippet: string;
query: string;
discoveredBy?: string;
requestedRoute?: string; // <--- ADD THIS
actualBackend?: string; // <--- ADD THIS
resultDomain?: string; // <--- ADD THIS
};
function isRateLimitError(msg: string): boolean {
const lower = msg.toLowerCase();
return (
lower.includes("429") ||
lower.includes("too many requests") ||
lower.includes("rate limit")
);
}
function zeroMetrics(): CrawlMetrics {
return {
ddgQueries: 0, ddgHits: 0, mutatedQueriesTried: 0, mutationAccepted: 0,
mutationHits: 0, extraEngineQueries: 0, extraEngineHits: 0, rawHits: 0,
dedupedHits: 0, rankedCandidates: 0, fetchCandidates: 0, fetchAttempts: 0,
fetchFailures: 0, acceptedSources: 0, skippedLowWordCount: 0, skippedOffTopic: 0,
skippedVeryOffTopic: 0, skippedDuplicateContent: 0, skippedVisited: 0, skippedDomainCap: 0,
skippedAvoided: 0, skippedBlacklisted: 0, cacheChecks: 0, cacheHits: 0, cacheAccepted: 0,
cacheRejectedDuplicate: 0, cacheRejectedOffTopic: 0, cacheRejectedLowWordCount: 0,
cacheWrites: 0, followedLinks: 0, crossWorkerDiscoveriesUsed: 0, localSourcesAccepted: 0,
};
}
export async function runWorker(
task: SwarmTask,
state: SharedCrawlState,
signal: AbortSignal,
status: StatusFn,
warn: WarnFn,
topicKws: ReadonlyArray<string> = [],
limiter: DdgRateLimiter = sharedDdgLimiter,
health: SearchHealthTracker,
llmManager: LlmCallManager
): Promise<WorkerResult> {
const sources: CrawledSource[] = [];
const errors: string[] = [];
const queriesExecuted: string[] = [];
const roleTag = `[${task.label}]`;
const metrics = zeroMetrics();
status(`${roleTag} Starting - ${task.queries.length} queries, budget=${task.pageBudget}, links=${task.followLinks ? "on" : "off"}`);
if (task.enableLocalSources) {
const localBudget = Math.max(2, Math.ceil(task.pageBudget * 0.3));
const localSources = harvestLocalSources(
task.queries, task.role, task.label, localBudget, task.contentLimit, task.localLibraryIds, task.roleLibraryMap,
);
if (localSources.length > 0) {
for (const src of localSources) {
state.addVisited(src.url);
state.addHash(contentFingerprint(src.text));
sources.push(src);
}
metrics.localSourcesAccepted += localSources.length;
status(`${roleTag} Local sources: ${localSources.length} chunks`);
}
}
const allHits: SearchHitLike[] = [];
let ddgBlocked = false;
for (const query of task.queries) {
if (signal.aborted) break;
// 9. QUERY PLANNING RULES - Reject bad queries
// Relaxed Query Quality Check: Reject only empty, corrupt, or absurdly long query strings
const qTrimmed = query.trim();
if (!qTrimmed || qTrimmed.length < 3 || qTrimmed.length > 250) {
warn(`[Query Quality] Rejected out-of-bounds query length (${qTrimmed.length} chars): "${qTrimmed.slice(0, 50)}..."`);
continue;
}
let ddgHits: ReadonlyArray<SearchHit> = [];
let effectiveHitCount = 0;
// 4. DDG FAILURE DEFINITION & HEALTH TRACKING
if (task.extraEngines.includes("ddg") && health.isDdgAvailable()) {
metrics.ddgQueries++;
try {
if (task.searchPages > 1) {
ddgHits = await searchDDGPaginated(query, task.searchResultsPerQuery, task.searchPages, task.safeSearch, signal, limiter, task.timeRange ?? "all");
} else {
ddgHits = await searchDDG(query, task.searchResultsPerQuery, task.safeSearch, signal, limiter, task.timeRange ?? "all");
}
// 5. LOWERED DDG TRIPWIRE: Accept even 1 result to support tight encyclopedia operators
if (ddgHits.length < 1) {
throw new Error("INVALID_RESULT_SET (< 1 result)");
}
health.recordDdgSuccess(ddgHits.length);
metrics.ddgHits += ddgHits.length;
effectiveHitCount = ddgHits.length;
// Track Attribution
for (const h of ddgHits) {
let domain = "";
try { domain = new URL(h.url).hostname.replace(/^www\./, ""); } catch {}
allHits.push({
...h,
query,
requestedRoute: "DDG",
actualBackend: "DDG",
resultDomain: domain
});
}
queriesExecuted.push(query);
status(`${roleTag} DDG: "${query}" -> ${ddgHits.length} results`);
} catch (err: unknown) {
if (isAbortError(err)) break;
const msg = errorMessage(err);
health.recordDdgFailure(msg);
warn(`${roleTag} DDG failed (${health.ddg.consecutiveFailures}/5): ${msg}`);
}
}
// Query Mutation (LLM)
const isHighlyRestrictive = query.includes("site:");
if (
ddgHits.length < task.queryMutationThreshold &&
!signal.aborted &&
!ddgBlocked &&
!isHighlyRestrictive &&
llmManager.canCall(task.id)
) {
llmManager.recordCall(task.id);
const contextSnippets = allHits.filter((h) => h.query === query).slice(0, 3).map((h) => h.snippet);
const llmMutated = await getLLMQueryMutation(query, contextSnippets, signal);
const mutationsToTry = [llmMutated, ...MUTATION_STRATEGIES.map((s) => s(query))].filter((m): m is string => !!m && m !== query);
let bestMutation: { query: string; hits: ReadonlyArray<SearchHit> } | null = null;
const baseline = mutationQuality(ddgHits);
for (const mutated of mutationsToTry) {
if (signal.aborted) break;
metrics.mutatedQueriesTried++;
metrics.ddgQueries++;
try {
const mutHits = await searchDDG(mutated, task.searchResultsPerQuery, task.safeSearch, signal, limiter, task.timeRange ?? "all");
metrics.ddgHits += mutHits.length;
if (mutationQuality(mutHits) > baseline) {
if (!bestMutation || mutationQuality(mutHits) > mutationQuality(bestMutation.hits)) {
bestMutation = { query: mutated, hits: mutHits };
}
}
} catch (err) {
const msg = errorMessage(err);
if (isRateLimitError(msg)) {
ddgBlocked = true;
warn(`${roleTag} Mutation RATE LIMIT (429)! Falling back...`);
break;
}
}
}
if (bestMutation) {
metrics.mutationAccepted++;
metrics.mutationHits += bestMutation.hits.length;
for (const h of bestMutation.hits) allHits.push({ ...h, query: bestMutation.query });
queriesExecuted.push(bestMutation.query);
effectiveHitCount = Math.max(effectiveHitCount, bestMutation.hits.length);
}
} else if (llmManager.canCall(task.id) === false && !isHighlyRestrictive && !ddgBlocked) {
warn(`${roleTag} LLM call budget exhausted, skipping mutation.`);
}
// Fallback to SearxNG or API if DDG is dead or hit 0
const shouldUseExtraEngines = task.extraEngines.length > 0 && !signal.aborted && (ddgBlocked || effectiveHitCount < Math.ceil(task.searchResultsPerQuery * 0.6) || task.role === "academic" || task.role === "primary" || task.role === "regulatory");
if (shouldUseExtraEngines) {
metrics.extraEngineQueries++;
try {
const extraHits = await multiEngineSearch(
query, Math.min(task.searchResultsPerQuery, 8),
task.extraEngines as ReadonlyArray<SearchEngine>,
signal, () => limiter, task.timeRange ?? "all",
{ serperApiKey: (task as any).serperApiKey, braveApiKey: (task as any).braveApiKey }
);
metrics.extraEngineHits += extraHits.length;
for (const h of extraHits) allHits.push({ ...h, query });
if (extraHits.length > 0) status(`${roleTag} -> ${extraHits.length} extra results`);
} catch {}
}
}
metrics.rawHits = allHits.length;
if (signal.aborted || allHits.length === 0) {
metrics.acceptedSources = sources.length;
state.mergeMetrics(metrics);
return { taskId: task.id, role: task.role, label: task.label, sources, queries: queriesExecuted, errors };
}
const deduped = deduplicateByUrl(allHits);
metrics.dedupedHits = deduped.length;
const scored = deduped.map((h) => {
const sc = scoreCandidate(h, h.query);
return { ...sc, discoveredBy: h.discoveredBy };
});
const filtered = task.preferredTiers ? scored.filter((c) => task.preferredTiers!.includes(c.tier)) : scored;
const poolSize = task.pageBudget * task.candidatePoolMultiplier;
const ranked = rankCandidates(filtered.length > 0 ? filtered : scored, poolSize);
const candidates = capCandidatesPerHost(ranked, 2);
metrics.rankedCandidates = candidates.length;
metrics.fetchCandidates = candidates.length;
await fetchBatch(candidates, task, state, signal, status, warn, sources, errors, roleTag, topicKws, metrics, llmManager);
if (task.followLinks && sources.length > 0 && sources.length < task.pageBudget) {
for (let depth = 1; depth <= task.linkCrawlDepth; depth++) {
if (sources.length >= task.pageBudget || signal.aborted) break;
const budget = Math.min(task.pageBudget - sources.length, task.maxLinksToFollow);
if (budget <= 0) break;
const sourcesForLinks = depth === 1 ? sources : sources.slice(-budget * 2);
const newCount = await followLinks(sourcesForLinks, task, state, signal, status, warn, sources, errors, roleTag, budget, topicKws, depth, metrics, llmManager);
metrics.followedLinks += newCount;
if (newCount === 0) break;
}
}
for (const src of sources) {
for (const link of src.outlinks.slice(0, 5)) {
state.pushDiscovery(link.href, link.text, task.label);
}
}
metrics.acceptedSources = sources.length;
state.mergeMetrics(metrics);
status(`✅ ${roleTag} Mission Complete - ${sources.length} high-quality sources collected!`);
return { taskId: task.id, role: task.role, label: task.label, sources, queries: queriesExecuted, errors };
}
async function fetchBatch(
candidates: ReadonlyArray<ScoredCandidate>,
task: SwarmTask,
state: SharedCrawlState,
signal: AbortSignal,
status: StatusFn,
warn: WarnFn,
results: CrawledSource[],
errors: string[],
tag: string,
topicKws: ReadonlyArray<string>,
metrics: CrawlMetrics,
llmManager: LlmCallManager
): Promise<void> {
let idx = 0;
const concurrency = task.workerConcurrency;
const domainCap = task.maxPagesPerDomain;
const minRelevance = task.minRelevanceScore;
while (results.length < task.pageBudget && idx < candidates.length && !signal.aborted) {
const slice = candidates.slice(idx, idx + concurrency);
idx += concurrency;
const batch: ScoredCandidate[] = [];
for (const c of slice) {
const normalized = normalizeUrl(c.url);
if (state.visitedUrls.has(normalized)) { metrics.skippedVisited++; continue; }
if (state.domainCount(c.url) >= domainCap) { metrics.skippedDomainCap++; continue; }
if (state.shouldAvoidUrl(c.url)) { metrics.skippedAvoided++; continue; }
if (state.isDomainBlacklisted(c.url)) { metrics.skippedBlacklisted++; continue; }
batch.push(c);
}
if (batch.length === 0) continue;
for (const c of batch) state.addVisited(c.url);
const settled = await Promise.allSettled(
batch.map((c) => resolveCandidate(c, task, state, topicKws, signal, metrics)),
);
for (let i = 0; i < settled.length; i++) {
const candidate = batch[i];
const settledResult = settled[i];
if (signal.aborted) return;
if (settledResult.status === "rejected") {
if (!isAbortError(settledResult.reason)) {
const msg = errorMessage(settledResult.reason);
metrics.fetchFailures++;
state.noteFailure(candidate.url, msg);
state.noteDomainFailure(candidate.url);
warn(`${tag} Failed: ${truncUrl(candidate.url)} - ${msg}`);
errors.push(`fetch:${candidate.url}: ${msg}`);
}
continue;
}
const { page, fromCache } = settledResult.value;
if (page.wordCount < MIN_USEFUL_WORD_COUNT) { metrics.skippedLowWordCount++; if (fromCache) metrics.cacheRejectedLowWordCount++; continue; }
if (page.relevanceScore < minRelevance * 0.5) { metrics.skippedVeryOffTopic++; if (fromCache) metrics.cacheRejectedOffTopic++; else state.noteDomainFailure(candidate.url); continue; }
if (page.relevanceScore < minRelevance) { metrics.skippedOffTopic++; if (fromCache) metrics.cacheRejectedOffTopic++; continue; }
const fp = contentFingerprint(page.text);
if (state.contentHashes.has(fp)) { metrics.skippedDuplicateContent++; if (fromCache) metrics.cacheRejectedDuplicate++; continue; }
state.addHash(fp);
state.incrementDomain(candidate.url);
if (fromCache) {
metrics.cacheAccepted++;
} else {
state.markVisitedPersistent(page);
metrics.cacheWrites++;
state.addWebCacheDocument(page);
}
results.push(page);
llmManager.recordProgress(); // Tell the watchdog we found something!
status(`${tag} [${results.length}/${task.pageBudget}] ${fromCache ? "[cache] " : ""}(rel=${page.relevanceScore.toFixed(2)}) ${page.title.slice(0, 60)}`);
if (results.length >= task.pageBudget) return;
}
if (idx < candidates.length && results.length < task.pageBudget) {
await sleep(BATCH_INTER_FETCH_DELAY_MS);
}
}
}
async function resolveCandidate(
candidate: ScoredCandidate,
task: SwarmTask,
state: SharedCrawlState,
topicKws: ReadonlyArray<string>,
signal: AbortSignal,
metrics: CrawlMetrics,
): Promise<{ page: CrawledSource; fromCache: boolean }> {
metrics.cacheChecks++;
const cached = state.getCachedSource(candidate.url);
if (cached) {
metrics.cacheHits++;
return {
page: adaptCachedSourceForTask(cached, candidate.query, candidate.snippet, task, topicKws),
fromCache: true,
};
}
metrics.fetchAttempts++;
return {
page: await fetchAndExtract(candidate.url, candidate.query, candidate.snippet, task, topicKws, signal, candidate.discoveredBy),
fromCache: false,
};
}
function adaptCachedSourceForTask(
cached: CrawledSource,
query: string,
snippet: string,
task: SwarmTask,
topicKws: ReadonlyArray<string>,
): CrawledSource {
const { domainScore, freshnessScore, tier } = scoreCandidate({ url: cached.url, title: cached.title, snippet: cached.description }, query);
const relevanceScore = computeRelevance(cached.text, cached.title, snippet, topicKws);
return {
...cached,
url: normalizeUrl(cached.url),
finalUrl: normalizeUrl(cached.finalUrl ?? cached.url),
sourceQuery: query,
workerRole: task.role,
workerLabel: task.label,
domainScore,
freshnessScore,
tier: tier as SourceTier,
relevanceScore,
};
}
async function followLinks(
existingSources: ReadonlyArray<CrawledSource>,
task: SwarmTask,
state: SharedCrawlState,
signal: AbortSignal,
status: StatusFn,
warn: WarnFn,
results: CrawledSource[],
errors: string[],
tag: string,
budget: number,
topicKws: ReadonlyArray<string>,
depth: number,
metrics: CrawlMetrics,
llmManager: LlmCallManager
): Promise<number> {
const allLinks = existingSources.flatMap((s) => s.outlinks);
const linkKws = task.queries.join(" ").toLowerCase().split(/\s+/).filter((w) => w.length > 3).slice(0, 12);
const scored = scoreOutlinks(allLinks, linkKws, state.visitedUrls, task.maxLinksToEvaluate);
const toFollow = scored.slice(0, task.maxLinksToFollow);
if (toFollow.length === 0) return 0;
status(`${tag} Following ${toFollow.length} link(s) (depth ${depth})…`);
const before = results.length;
const linkCandidates = toFollow.map((l) => scoreCandidate({ url: l.href, title: "", snippet: "" }, task.queries[0] ?? ""));
await fetchBatch(linkCandidates, { ...task, pageBudget: results.length + budget }, state, signal, status, warn, results, errors, tag, topicKws, metrics, llmManager);
return results.length - before;
}
async function fetchAndExtract(
url: string,
query: string,
snippet: string,
task: SwarmTask,
topicKws: ReadonlyArray<string>,
signal: AbortSignal,
discoveredBy?: string,
): Promise<CrawledSource> {
let page: ExtractedPage;
if (/youtube\.com|youtu\.be/i.test(url)) {
const ytData = await extractYouTubeTranscript(url, signal, task.contentLimit);
if (!ytData) throw new Error("YouTube transcript extraction failed");
page = {
url: url,
finalUrl: url,
title: ytData.title,
description: ytData.description,
published: null,
text: ytData.text,
wordCount: ytData.text.split(/\s+/).filter(Boolean).length,
outlinks: [],
page: 1,
totalPages: 1,
};
} else {
const fetchResult = await fetchWithWaybackFallback(url, signal, task);
const { finalUrl } = fetchResult;
const isPdf = (fetchResult.rawBuffer && isPdfContentType(fetchResult.contentType)) || (!fetchResult.rawBuffer && isPdfUrl(url));
if (isPdf && fetchResult.rawBuffer) {
page = await extractPdf(fetchResult.rawBuffer, url, finalUrl, task.contentLimit, false);
} else if (isPdf && fetchResult.html && fetchResult.html.startsWith("%PDF")) {
const buf = Buffer.from(fetchResult.html, "binary");
page = await extractPdf(buf, url, finalUrl, task.contentLimit, false);
} else {
page = extractPage(fetchResult.html, url, finalUrl, task.contentLimit, task.maxOutlinksPerPage);
}
}
const { domainScore, freshnessScore, tier } = scoreCandidate({ url, title: page.title, snippet: page.description }, query);
const relevanceScore = computeRelevance(page.text, page.title, snippet, topicKws);
return {
url: normalizeUrl(page.url),
finalUrl: normalizeUrl(page.finalUrl),
title: page.title,
description: page.description,
published: page.published,
text: page.text,
wordCount: page.wordCount,
outlinks: page.outlinks,
sourceQuery: query,
workerRole: task.role,
workerLabel: task.label,
domainScore,
freshnessScore,
tier: tier as SourceTier,
relevanceScore,
origin: "web" as const,
discoveredBy: discoveredBy ?? "unknown",
page: page.page,
totalPages: page.totalPages,
};
}
async function getLLMQueryMutation(
query: string,
topSnippets: string[],
signal: AbortSignal,
): Promise<string | null> {
try {
const endpoint = "http://localhost:1234/v1/chat/completions";
const prompt = `You are a deep research assistant. The initial query "${query}" yielded these snippets:\n${topSnippets.slice(0, 3).join("\n")}\n\nGenerate ONE highly specific, alternative search query to find missing technical details or counter-arguments. Return ONLY the raw query string, no quotes or explanations.`;
logLlmDiagnostics("getLLMQueryMutation", prompt);
const res = await fetch(endpoint, {
method: "POST",
headers: { "Content-Type": "application/json" },
signal,
body: JSON.stringify({
model: "local-model",
messages: [{ role: "user", content: prompt }],
temperature: 0.7,
max_tokens: 40,
}),
});
if (!res.ok) return null;
const data = await res.json();
const mutated = data.choices?.[0]?.message?.content?.trim();
return mutated ? mutated.replace(/^[\"']|[\"']$/g, "") : null;
} catch {
return null;
}
}
async function fetchWithWaybackFallback(
url: string,
signal: AbortSignal,
task: SwarmTask // <-- We add task here to grab the config!
): Promise<Awaited<ReturnType<typeof fetchPage>>> {
try {
// TIER A: Standard Fetch (Tries normally first)
return await fetchPage(url, signal);
} catch (err: any) {
if (isAbortError(err)) throw err;
// TIER B: FlareSolverr (Optional Power-User Bypass)
const fsUrl = (task as any).flaresolverrUrl as string | undefined;
if (fsUrl && fsUrl.trim() !== "") {
try {
const fsRes = await fetch(fsUrl, {
method: "POST",
headers: { "Content-Type": "application/json" },
body: JSON.stringify({ cmd: "request.get", url: url, maxTimeout: 15000 }),
signal
});
const fsData = await fsRes.json();
if (fsData?.solution?.response) {
return {
html: fsData.solution.response,
finalUrl: url,
contentType: "text/html",
rawBuffer: undefined // FlareSolverr only returns HTML, not raw PDFs
};
}
} catch (fsErr) {
// Silently fail and drop down to Wayback Machine
}
}
// TIER C: Wayback Machine Archive
try {
const api = `https://archive.org/wayback/available?url=${encodeURIComponent(url)}`;
const apiRes = await fetch(api, { signal });
if (!apiRes.ok) throw new Error("Wayback API failed");
const data = await apiRes.json();
const wbUrl = data?.archived_snapshots?.closest?.url;
if (wbUrl) {
return await fetchPage(wbUrl, signal);
}
} catch {
// ignore
}
throw err;
}
}
function uniqueHostCount(hits: ReadonlyArray<{ url: string }>): number {
const hosts = new Set<string>();
for (const h of hits) {
try {
hosts.add(new URL(h.url).hostname.replace(/^www\./, ""));
} catch {
// ignore
}
}
return hosts.size;
}
function mutationQuality(hits: ReadonlyArray<{ url: string }>): number {
return uniqueHostCount(hits) * 10 + Math.min(hits.length, 10);
}
function deduplicateByUrl<T extends { url: string }>(items: ReadonlyArray<T>): T[] {
const seen = new Set<string>();
return items.filter((item) => {
const key = normalizeUrl(item.url);
if (seen.has(key)) return false;
seen.add(key);
return true;
});
}
function capCandidatesPerHost(
candidates: ReadonlyArray<ScoredCandidate>,
maxPerHost: number,
): ScoredCandidate[] {
const perHost = new Map<string, number>();
const kept: ScoredCandidate[] = [];
for (const c of candidates) {
const host = safeHostname(c.url);
const count = perHost.get(host) ?? 0;
if (count >= maxPerHost) continue;
perHost.set(host, count + 1);
kept.push(c);
}
return kept;
}
function safeHostname(url: string): string {
try {
return new URL(url).hostname.replace(/^www\./, "");
} catch {
return "";
}
}
function truncUrl(url: string, max = 70): string {
return url.length > max ? url.slice(0, max) + "…" : url;
}
function isAbortError(err: unknown): boolean {
return err instanceof DOMException && err.name === "AbortError";
}
function errorMessage(err: unknown): string {
return err instanceof Error ? err.message : String(err ?? "unknown");
}
// ==================== FILE: src\net\ddg.ts ====================
import { fetchPage } from "./http";
export class DdgRateLimiter {
private lastRequest = 0;
private minDelay: number;
constructor(minDelayMs: number = 2500) {
this.minDelay = minDelayMs;
}
async acquire(): Promise<void> {
const now = Date.now();
const elapsed = now - this.lastRequest;
if (elapsed < this.minDelay) {
await new Promise((resolve) => setTimeout(resolve, this.minDelay - elapsed));
}
this.lastRequest = Date.now();
}
}
export const sharedDdgLimiter = new DdgRateLimiter(2500);
export class DdgLimiterPool {
private limiters: DdgRateLimiter[] = [];
private currentIndex = 0;
constructor(numLanes: number, minDelayMs: number = 2500) {
for (let i = 0; i < numLanes; i++) {
this.limiters.push(new DdgRateLimiter(minDelayMs));
}
}
next(): DdgRateLimiter {
const limiter = this.limiters[this.currentIndex];
this.currentIndex = (this.currentIndex + 1) % this.limiters.length;
return limiter;
}
}
export function resetThrottle(): void {}
/**
* Bulletproof multi-tier DDG Search Waterfall:
* Tier A: DDG Lite POST endpoint (fastest, mimics ddgr)
* Tier B: DDG HTML fallback endpoint (handles strict blocks)
* Tier C: Graceful recovery (returns empty array instead of throwing crash errors)
*/
export async function searchDDG(
query: string,
maxResults: number,
safeSearch: "strict" | "moderate" | "off" = "moderate",
signal?: AbortSignal,
limiter?: DdgRateLimiter,
timeRange: string = "all"
): Promise<ReadonlyArray<{ url: string; title: string; snippet: string }>> {
if (limiter) await limiter.acquire();
if (signal?.aborted) throw new DOMException("Aborted", "AbortError");
const safeParam = safeSearch === "strict" ? "1" : safeSearch === "off" ? "-1" : "0";
const dfParam = timeRange === "all" ? "" : `&df=${timeRange}`;
const formData = `q=${encodeURIComponent(query)}&kp=${safeParam}${dfParam}`;
// TIER A: DDG Lite Endpoint
try {
const res = await fetchPage("https://lite.duckduckgo.com/lite/", signal!, {
method: "POST",
body: formData,
headers: {
"Content-Type": "application/x-www-form-urlencoded",
"Referer": "https://lite.duckduckgo.com/",
"User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/124.0.0.0 Safari/537.36",
"Accept": "text/html,application/xhtml+xml,application/xml;q=0.9,image/avif,image/webp,image/apng,*/*;q=0.8",
"Accept-Language": "en-US,en;q=0.9",
"Sec-Ch-Ua": '"Chromium";v="124", "Google Chrome";v="124", "Not-A.Brand";v="99"',
"Sec-Ch-Ua-Mobile": "?0",
"Sec-Ch-Ua-Platform": '"Windows"',
"Sec-Fetch-Dest": "document",
"Sec-Fetch-Mode": "navigate",
"Sec-Fetch-Site": "same-origin",
"Sec-Fetch-User": "?1",
"Upgrade-Insecure-Requests": "1"
}
});
const hits = parseDDGResults(res.html, maxResults);
if (hits.length > 0) return hits;
} catch (err) {
if (signal?.aborted) throw err;
console.warn(`[DDG Tier A] Failed for query "${query}": ${err instanceof Error ? err.message : String(err)}. Trying Tier B...`);
}
// TIER B: DDG HTML Endpoint Fallback
try {
if (limiter) await limiter.acquire();
const htmlFallbackUrl = `https://html.duckduckgo.com/html/?q=${encodeURIComponent(query)}`;
const resB = await fetchPage(htmlFallbackUrl, signal!, {
method: "GET",
headers: {
"Content-Type": "application/x-www-form-urlencoded",
"Referer": "https://lite.duckduckgo.com/",
"User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/124.0.0.0 Safari/537.36",
"Accept": "text/html,application/xhtml+xml,application/xml;q=0.9,image/avif,image/webp,image/apng,*/*;q=0.8",
"Accept-Language": "en-US,en;q=0.9",
"Sec-Ch-Ua": '"Chromium";v="124", "Google Chrome";v="124", "Not-A.Brand";v="99"',
"Sec-Ch-Ua-Mobile": "?0",
"Sec-Ch-Ua-Platform": '"Windows"',
"Sec-Fetch-Dest": "document",
"Sec-Fetch-Mode": "navigate",
"Sec-Fetch-Site": "same-origin",
"Sec-Fetch-User": "?1",
"Upgrade-Insecure-Requests": "1"
}
});
const hitsB = parseHTMLDDGResults(resB.html, maxResults);
if (hitsB.length > 0) return hitsB;
} catch (errB) {
if (signal?.aborted) throw errB;
console.warn(`[DDG Tier B] Fallback also failed for query "${query}": ${errB instanceof Error ? errB.message : String(errB)}`);
}
// TIER C: Graceful exit (Returns empty array so the worker's outer health system handles it cleanly without crashing)
return [];
}
export async function searchDDGPaginated(
query: string,
maxResultsPerPage: number,
pages: number,
safeSearch: "strict" | "moderate" | "off" = "moderate",
signal?: AbortSignal,
limiter?: DdgRateLimiter,
timeRange: string = "all"
): Promise<ReadonlyArray<{ url: string; title: string; snippet: string }>> {
const allHits: { url: string; title: string; snippet: string }[] = [];
for (let p = 1; p <= pages; p++) {
if (signal?.aborted) break;
const safeParam = safeSearch === "strict" ? "1" : safeSearch === "off" ? "-1" : "0";
const dfParam = timeRange === "all" ? "" : `&df=${timeRange}`;
const formData = `q=${encodeURIComponent(query)}&kp=${safeParam}${dfParam}&s=${(p - 1) * maxResultsPerPage}`;
try {
if (limiter) await limiter.acquire();
const res = await fetchPage("https://lite.duckduckgo.com/lite/", signal!, {
method: "POST",
body: formData,
headers: {
"Content-Type": "application/x-www-form-urlencoded",
"Referer": "https://lite.duckduckgo.com/"
}
});
const hits = parseDDGResults(res.html, maxResultsPerPage);
allHits.push(...hits);
if (hits.length < maxResultsPerPage) break;
} catch {
break;
}
}
return allHits;
}
function parseDDGResults(html: string, maxResults: number): { url: string; title: string; snippet: string }[] {
const hits: { url: string; title: string; snippet: string }[] = [];
const seen = new Set<string>();
// Resilient regex pattern matching DDG Lite result rows
const resultRe = /<a[^>]*rel="nofollow"[^>]*href="([^"]+)"[^>]*>([\s\S]*?)<\/a>[\s\S]*?<td class="result-snippet">([\s\S]*?)<\/td>/gi;
let match: RegExpExecArray | null;
while (hits.length < maxResults && (match = resultRe.exec(html)) !== null) {
let url = match[1];
const title = match[2].replace(/<[^>]+>/g, "").trim();
const snippet = match[3].replace(/<[^>]+>/g, "").trim();
// Clean up redirect wrappers if present
if (url.includes("uddg=")) {
const matchUrl = url.match(/uddg=([^&]+)/);
if (matchUrl) url = decodeURIComponent(matchUrl[1]);
}
if (url.includes("duckduckgo.com")) continue;
if (!url.startsWith("http")) continue;
if (seen.has(url)) continue;
seen.add(url);
hits.push({ url, title, snippet });
}
return hits;
}
function parseHTMLDDGResults(html: string, maxResults: number): { url: string; title: string; snippet: string }[] {
const hits: { url: string; title: string; snippet: string }[] = [];
const seen = new Set<string>();
// Secondary parser for standard DDG HTML results page layout
const resultRe = /<a class="result__url" href="([^"]+)"[^>]*>([\s\S]*?)<\/a>[\s\S]*?<a class="result__snippet"[^>]*>([\s\S]*?)<\/a>/gi;
let match: RegExpExecArray | null;
while (hits.length < maxResults && (match = resultRe.exec(html)) !== null) {
let url = match[1];
const title = match[2].replace(/<[^>]+>/g, "").trim();
const snippet = match[3].replace(/<[^>]+>/g, "").trim();
if (url.includes("duckduckgo.com")) continue;
if (!url.startsWith("http") && url.startsWith("//")) url = "https:" + url;
if (!url.startsWith("http")) continue;
if (seen.has(url)) continue;
seen.add(url);
hits.push({ url, title, snippet });
}
return hits;
}
// ==================== FILE: src\net\http.ts ====================
// src/net/http.ts
import * as https from "node:https";
import * as http from "node:http";
import { setServers } from "node:dns";
const DNS_RESOLVERS = ["1.1.1.1", "1.0.0.1", "8.8.8.8", "8.8.4.4", "9.9.9.9"];
setServers(DNS_RESOLVERS);
const UA_POOL: ReadonlyArray<string> = [
"Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/124.0.0.0 Safari/537.36",
"Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/123.0.0.0 Safari/537.36",
"Mozilla/5.0 (Windows NT 10.0; Win64; x64; rv:125.0) Gecko/20100101 Firefox/125.0",
"Mozilla/5.0 (X11; Linux x86_64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/124.0.0.0 Safari/537.36",
];
function randomUA(): string {
return UA_POOL[Math.floor(Math.random() * UA_POOL.length)];
}
// SSRF Guard: Block internal/local IPs
function isPrivateUrl(url: string): boolean {
try {
const u = new URL(url);
const host = u.hostname;
if (host === "localhost" || host === "127.0.0.1" || host === "0.0.0.0") return true;
if (host.startsWith("10.") || host.startsWith("192.168.")) return true;
if (host.startsWith("172.")) {
const parts = host.split(".");
const second = parseInt(parts[1], 10);
if (second >= 16 && second <= 31) return true;
}
return false;
} catch {
return true;
}
}
export function buildBrowserHeaders(url?: string): Record<string, string> {
const ua = randomUA();
const headers: Record<string, string> = {
"User-Agent": ua,
Accept: "text/html,application/xhtml+xml,application/xml;q=0.9,image/avif,image/webp,*/*;q=0.8",
"Accept-Language": "en-US,en;q=0.9",
"Accept-Encoding": "gzip, deflate, br",
DNT: "1",
Connection: "keep-alive",
"Upgrade-Insecure-Requests": "1",
"Sec-Fetch-Dest": "document",
"Sec-Fetch-Mode": "navigate",
"Sec-Fetch-Site": "none",
"Sec-Fetch-User": "?1",
};
if (url) {
try {
const parsed = new URL(url);
headers["Referer"] = `${parsed.protocol}//${parsed.hostname}/`;
} catch {}
}
return headers;
}
export function buildDDGHeaders(): Record<string, string> {
return buildBrowserHeaders();
}
export interface FetchResult {
readonly html: string;
readonly finalUrl: string;
readonly contentType?: string;
readonly rawBuffer?: Buffer;
}
const FETCH_TIMEOUT_MS = 12000;
const FETCH_MAX_RETRIES = 3;
// Proxy support (optional)
const PROXY_URL = process.env.HTTP_PROXY || process.env.HTTPS_PROXY;
const proxyAgent = PROXY_URL ? new (require("https-proxy-agent").HttpsProxyAgent)(PROXY_URL) : undefined;
export interface FetchOptions {
method?: 'GET' | 'POST';
body?: string;
headers?: Record<string, string>;
}
export async function fetchPage(
url: string,
signal: AbortSignal,
options: FetchOptions = {}
): Promise<FetchResult> {
if (isPrivateUrl(url)) throw new Error("SSRF Guard: Blocked internal URL");
let lastError: unknown;
const { method = 'GET', body, headers: customHeaders = {} } = options;
for (let attempt = 0; attempt < FETCH_MAX_RETRIES; attempt++) {
if (signal.aborted) throw new DOMException("Aborted", "AbortError");
// Delay to prevent burst blocking
await new Promise((resolve) => setTimeout(resolve, 2500));
try {
const headers = { ...buildBrowserHeaders(url), ...customHeaders };
const controller = new AbortController();
const timeout = setTimeout(() => controller.abort(), FETCH_TIMEOUT_MS);
const res = await fetch(url, {
method,
body,
signal: controller.signal,
headers,
redirect: "follow"
});
clearTimeout(timeout);
// If blocked by Cloudflare/WAF, fall back to Wayback Machine immediately
if ([403, 429, 503].includes(res.status)) {
throw new Error(`HTTP ${res.status}`);
}
const contentType = res.headers.get("content-type") || "";
const finalUrl = res.url || url;
if (contentType.includes("application/pdf") || contentType.includes("application/octet-stream")) {
const arrayBuf = await res.arrayBuffer();
return { html: "", finalUrl, contentType, rawBuffer: Buffer.from(arrayBuf) };
}
const html = await res.text();
// Only check for very specific captcha markers, avoid generic words like "blocked"
if (html.includes("cf-challenge") || html.includes("hcaptcha") || html.includes("recaptcha")) {
throw new Error("Captcha page detected");
}
return { html, finalUrl, contentType };
} catch (err: unknown) {
if (err instanceof DOMException && err.name === "AbortError") throw err;
lastError = err;
// Exponential backoff
const backoffTime = 5000 * Math.pow(2, attempt);
await new Promise((resolve) => setTimeout(resolve, backoffTime));
}
}
// If direct fetch fails, try 10 Web Archives
return fetchFromArchives(url, signal);
}
async function fetchFromArchives(url: string, signal: AbortSignal): Promise<FetchResult> {
const encoded = encodeURIComponent(url);
const archives = [
`https://web.archive.org/web/2/${url}`,
`https://webcache.googleusercontent.com/search?q=cache:${encoded}&strip=1`,
`https://archive.today/newest/${url}`,
`https://cc.bingj.com/cache.aspx?q=${url}`,
`https://webcache.googleusercontent.com/search?q=cache:${url}`,
`http://web.archive.org/web/2024/${url}`,
`https://freezedry.com/${url}`,
`https://corsproxy.io/?${encoded}`,
`https://api.allorigins.win/raw?url=${encoded}`,
`https://cors-anywhere.herokuapp.com/${url}`
];
for (const cacheUrl of archives) {
if (signal.aborted) throw new DOMException("Aborted", "AbortError");
try {
const controller = new AbortController();
const timeout = setTimeout(() => controller.abort(), 8000);
const res = await fetch(cacheUrl, {
signal: controller.signal,
headers: buildBrowserHeaders(cacheUrl)
});
clearTimeout(timeout);
if (res.ok) {
const html = await res.text();
if (html.length > 500) {
return { html, finalUrl: url, contentType: res.headers.get("content-type") || "" };
}
}
} catch {}
}
throw new Error(`Failed to fetch ${url}: Blocked and all archives failed`);
}
export function safeHostname(url: string): string {
try {
return new URL(url).hostname;
} catch {
return "";
}
}
export function sleep(ms: number): Promise<void> {
return new Promise((resolve) => setTimeout(resolve, ms));
}
// ==================== FILE: src\planning\planner.ts ====================
import { logLlmDiagnostics } from "../utils/tokens";
import { LMStudioClient } from "@lmstudio/sdk";
import {
QueryPlan,
WorkerRole,
DynamicWorkerSpec,
AdaptiveGapPlan,
CrawledSource,
AgentMessage,
StatusFn,
} from "../types";
import { DIMENSIONS, detectGaps, gapFillQueries } from "./dimensions";
import {
DepthProfile,
AI_PLANNING_MAX_TOKENS,
AI_PLANNING_TEMPERATURE,
AI_PLANNING_TIMEOUT_MS,
AI_MIN_ACCEPTABLE_QUERIES,
AI_DECOMPOSITION_MAX_TOKENS,
AI_DECOMPOSITION_TEMPERATURE,
AI_DECOMPOSITION_TIMEOUT_MS,
AI_FINDINGS_SUMMARY_MAX_TOKENS,
AI_FINDINGS_SUMMARY_TEMPERATURE,
FINDINGS_SUMMARY_SOURCE_CHARS,
DECOMPOSITION_MIN_WORKERS,
QUERY_LINE_MIN_LEN,
QUERY_LINE_MAX_LEN,
SYSTEM_INSTRUCTIONS,
} from "../constants";
async function callLoadedModel(
prompt: string,
maxTokens: number = AI_PLANNING_MAX_TOKENS,
temperature: number = AI_PLANNING_TEMPERATURE,
timeoutMs: number = AI_PLANNING_TIMEOUT_MS,
): Promise<string | null> {
try {
const client = new LMStudioClient();
const models = await Promise.race<Awaited<ReturnType<typeof client.llm.listLoaded>>>([
client.llm.listLoaded(),
new Promise<never>((_, reject) =>
setTimeout(() => reject(new Error("timeout")), timeoutMs),
),
]);
if (!Array.isArray(models) || models.length === 0) return null;
const model = await client.llm.model(models[0].identifier);
const stream = model.respond(
[
{ role: "system", content: SYSTEM_INSTRUCTIONS },
{ role: "user", content: prompt },
],
{
maxTokens,
temperature,
},
);
let result = "";
for await (const chunk of stream) result += chunk.content ?? "";
return result.trim() || null;
} catch {
return null;
}
}
function parseLines(raw: string, maxLines: number = 6): ReadonlyArray<string> {
return raw
.split(/\n/)
.map((line) => line.replace(/^\d+[.)]\s*|^[-*•]\s*/, "").trim())
.filter(
(line) =>
line.length > QUERY_LINE_MIN_LEN && line.length < QUERY_LINE_MAX_LEN,
)
.filter((line, idx, arr) => arr.indexOf(line) === idx)
.slice(0, maxLines);
}
const VALID_ROLES: ReadonlyArray<WorkerRole> = [
"breadth",
"depth",
"recency",
"academic",
"critical",
"statistical",
"regulatory",
"technical",
"primary",
"comparative",
];
function makeDecompositionPrompt(
topic: string,
focusAreas: ReadonlyArray<string>,
profile: DepthProfile,
): string {
const focus = focusAreas.length
? `\nFocus areas: ${focusAreas.join(", ")}`
: "";
return `You are a research decomposition system. Given a research topic, output a strict JSON array of specialized worker agents.
Topic: "${topic}"${focus}
Each worker object MUST have:
"role": one of "breadth", "depth", "recency", "academic", "critical", "statistical", "regulatory", "technical", "primary", "comparative"
"label": descriptive name (e.g., "Clinical Evidence Researcher", "Policy Critic")
"queries": array of ${Math.min(profile.maxQueriesPerWorker, 6)}-${profile.maxQueriesPerWorker} specific, natural search queries
"budgetWeight": number 0.1-0.4
"followLinks": boolean
"preferredTiers": array of strings
Rules:
1. Output ${DECOMPOSITION_MIN_WORKERS} to ${profile.maxDecompositionWorkers} workers.
2. Queries must be human-like search strings (under 8 words). DO NOT just mash the topic and role together.
3. Output ONLY valid JSON. No conversational text.
JSON:`;
}
async function aiDecompose(
topic: string,
focusAreas: ReadonlyArray<string>,
status: StatusFn,
profile: DepthProfile,
): Promise<ReadonlyArray<DynamicWorkerSpec> | null> {
const raw = await callLoadedModel(
makeDecompositionPrompt(topic, focusAreas, profile),
AI_DECOMPOSITION_MAX_TOKENS,
AI_DECOMPOSITION_TEMPERATURE,
AI_DECOMPOSITION_TIMEOUT_MS,
);
if (!raw) return null;
try {
const jsonStr = raw.replace(/```json\s*|```\s*/g, "").trim();
const parsed = JSON.parse(jsonStr);
if (!Array.isArray(parsed) || parsed.length < DECOMPOSITION_MIN_WORKERS)
return null;
const specs: DynamicWorkerSpec[] = [];
for (const item of parsed.slice(0, profile.maxDecompositionWorkers)) {
const role = VALID_ROLES.includes(item.role) ? item.role : "breadth";
const queries = Array.isArray(item.queries)
? item.queries
.filter((q: unknown) => typeof q === "string" && q.length > 3)
.slice(0, profile.maxQueriesPerWorker)
: [];
if (queries.length < 2) continue;
specs.push({
role: role as WorkerRole,
label:
typeof item.label === "string"
? item.label.slice(0, 60)
: `${role} worker`,
queries,
budgetWeight:
typeof item.budgetWeight === "number"
? Math.max(0.05, Math.min(0.5, item.budgetWeight))
: 0.2,
followLinks: item.followLinks === true,
preferredTiers: Array.isArray(item.preferredTiers)
? item.preferredTiers
: undefined,
});
}
if (specs.length < DECOMPOSITION_MIN_WORKERS) return null;
const totalWeight = specs.reduce((sum, s) => sum + s.budgetWeight, 0);
const normalised = specs.map((s) => ({
...s,
budgetWeight: s.budgetWeight / totalWeight,
}));
status(`AI decomposed topic into ${normalised.length} specialised workers`);
return normalised;
} catch {
return null;
}
}
function makeRolePlanPrompt(
role: WorkerRole,
topic: string,
focusAreas: ReadonlyArray<string>,
profile: DepthProfile,
): string {
const roleDescriptions: Readonly<Record<WorkerRole, string>> = {
breadth: "broad coverage - primary facts, key individuals, organizations, and canonical entities",
depth: "deep dive - experimental methodologies, field investigations, and primary case data",
recency: "recent developments - publications, peer reviews, and verifiable records from 2024-2026",
academic: "academic papers - peer-reviewed studies, university archives, and institutional research",
critical: "skeptical analysis - methodological critiques, counter-arguments, fraud investigations, and alternative explanations",
statistical: "statistics and quantitative findings - sample sizes, replication rates, and surveys",
regulatory: "governing standards, institutional positions, and formal guidelines",
technical: "technical mechanisms and specific evaluation frameworks",
primary: "primary documentation - field notes, first-hand interviews, and official recordings",
comparative: "comparative models - alternative scientific, physiological, or psychological hypotheses",
};
const focus = focusAreas.length ? `\nFocus areas: ${focusAreas.join(", ")}` : "";
return `You are an expert investigative search planner.
Target Topic: "${topic}"${focus}
Assigned Role: ${roleDescriptions[role]}
Generate exactly ${profile.maxQueriesPerWorker} precision search queries.
STRICT QUERY RULES:
1. Target canonical entities, notable researchers, research institutes, and key published works directly.
2. If the topic is empirical/scientific, append negative terms when appropriate (e.g., -livestream -forum -llm) to filter junk.
3. If this role is "critical", explicitly hunt for critiques, failed replications, and skeptics.
4. Output ONLY the queries, one per line. No numbering, quotes, or Markdown formatting.
Queries:`;
}
const ROLE_DIMENSIONS: Readonly<Record<WorkerRole, ReadonlyArray<string>>> = {
breadth: ["overview", "applications", "history", "economics"],
depth: ["mechanism", "evidence", "expert"],
recency: ["current", "future"],
academic: ["evidence", "expert", "mechanism"],
critical: ["challenges", "controversy", "comparison"],
statistical: ["evidence", "economics", "overview"],
regulatory: ["challenges", "controversy", "current"],
technical: ["mechanism", "applications", "evidence"],
primary: ["evidence", "expert", "history"],
comparative: ["comparison", "challenges", "applications"],
};
function dimensionFallbackQueries(
role: WorkerRole,
topic: string,
focusAreas: ReadonlyArray<string>,
maxQueries: number,
): ReadonlyArray<string> {
const dimIds = ROLE_DIMENSIONS[role];
const dims = DIMENSIONS.filter((d) => dimIds.includes(d.id));
const queries: string[] = [];
// BUGFIX: Actually use the shortenTopic function to strip the massive prompt
const shortTopic = shortenTopic(topic);
for (const dim of dims) {
for (const q of dim.queries(shortTopic)) {
if (!queries.includes(q)) queries.push(q);
}
}
for (const area of focusAreas) {
const q = `${shortTopic} ${area}`;
if (!queries.includes(q)) queries.push(q);
}
return queries.slice(0, maxQueries);
}
export async function buildQueryPlan(
topic: string,
focusAreas: ReadonlyArray<string>,
useAI: boolean,
status: StatusFn,
profile: DepthProfile,
): Promise<QueryPlan> {
const CORE_ROLES: ReadonlyArray<WorkerRole> = [
"breadth",
"depth",
"recency",
"academic",
"critical",
];
const EXTENDED_ROLES: ReadonlyArray<WorkerRole> = [
"statistical",
"regulatory",
"technical",
"primary",
"comparative",
];
let roles: ReadonlyArray<WorkerRole>;
if (profile.depthRounds >= 10) {
roles = [...CORE_ROLES, ...EXTENDED_ROLES];
} else if (profile.depthRounds >= 5) {
roles = [...CORE_ROLES, "technical", "comparative", "statistical"];
} else {
roles = CORE_ROLES;
}
const queriesByRole: Partial<Record<WorkerRole, ReadonlyArray<string>>> = {};
let usedAI = false;
let dynamicSpecs: ReadonlyArray<DynamicWorkerSpec> | undefined;
if (useAI) {
status("AI task decomposition - analysing topic for specialised workers...");
const specs = await aiDecompose(topic, focusAreas, status, profile);
if (specs && specs.length >= DECOMPOSITION_MIN_WORKERS) {
dynamicSpecs = specs;
usedAI = true;
for (const spec of specs) {
queriesByRole[spec.role] = spec.queries;
}
} else {
status("AI planning queries for each swarm worker...");
}
const uncoveredRoles = roles.filter((r) => !queriesByRole[r]?.length);
if (uncoveredRoles.length > 0) {
const results = await Promise.allSettled(
uncoveredRoles.map(async (role) => ({
role,
queries: await callLoadedModel(
makeRolePlanPrompt(role, topic, focusAreas, profile),
),
})),
);
for (const result of results) {
if (result.status !== "fulfilled") continue;
const { role, queries: raw } = result.value;
if (!raw) continue;
const parsed = parseLines(raw, profile.maxQueriesPerWorker);
if (parsed.length >= AI_MIN_ACCEPTABLE_QUERIES) {
queriesByRole[role] = parsed;
usedAI = true;
}
}
}
if (usedAI) {
status(
`AI generated queries for ${Object.keys(queriesByRole).length} worker role(s)`,
);
} else {
status("AI unavailable, using dimension-based query planning");
}
}
for (const role of roles) {
if (!queriesByRole[role] || queriesByRole[role]!.length === 0) {
queriesByRole[role] = dimensionFallbackQueries(
role,
topic,
focusAreas,
profile.maxQueriesPerWorker,
);
}
}
const topicKeywords = extractKeywords(topic);
// PRE-SEARCH PLANNING (If AI is enabled)
if (useAI) {
const planPrompt = `Topic: ${topic}\nFocus: ${focusAreas.join(", ")}\n
Generate a JSON object with:
1. "subQuestions": 3-5 specific questions to answer.
2. "likelySourceTypes": Array of strings (e.g., "academic", "news", "reference").
3. "stopConditions": Array of strings (e.g., "Found 3 sources confirming Q3 revenue").`;
logLlmDiagnostics("buildQueryPlan-PreSearch", planPrompt);
}
return {
queriesByRole: queriesByRole as Record<WorkerRole, ReadonlyArray<string>>,
usedAI,
topicKeywords,
dynamicSpecs,
};
}
export async function summariseFindings(
sources: ReadonlyArray<CrawledSource>,
topic: string,
useAI: boolean,
status: StatusFn,
): Promise<ReadonlyArray<AgentMessage>> {
if (!useAI || sources.length === 0) return [];
const sourceSummaries = sources
.slice(0, 20)
.map(
(s, i) =>
`[${i + 1}] ${s.workerLabel}: ${s.title} - ${s.text.slice(0, FINDINGS_SUMMARY_SOURCE_CHARS)}`,
)
.join("\n\n");
const prompt = `You are a research coordinator. A team of research agents collected these sources on "${topic}":
${sourceSummaries}
Summarise:
The 3-5 most important findings discovered so far (one line each)
3-5 specific questions or angles that were NOT covered and need follow-up
Output format:
FINDINGS:
finding 1
finding 2
...
FOLLOW_UP:
question 1
question 2
...`;
const raw = await callLoadedModel(
prompt,
AI_FINDINGS_SUMMARY_MAX_TOKENS,
AI_FINDINGS_SUMMARY_TEMPERATURE,
);
if (!raw) return [];
const findings: string[] = [];
const followUps: string[] = [];
let section: "findings" | "followup" | null = null;
for (const line of raw.split("\n")) {
const trimmed = line.trim();
if (/^FINDINGS:/i.test(trimmed)) {
section = "findings";
continue;
}
if (/^FOLLOW.?UP:/i.test(trimmed)) {
section = "followup";
continue;
}
const item = trimmed.replace(/^[-*•]\s*/, "").trim();
if (item.length < 5) continue;
if (section === "findings") findings.push(item);
if (section === "followup") followUps.push(item);
}
if (findings.length === 0 && followUps.length === 0) return [];
status(
`AI summarised ${findings.length} key findings, ${followUps.length} follow-up suggestions`,
);
return [
{
fromWorker: "round-coordinator",
keyFindings: findings,
suggestedFollowUps: followUps,
},
];
}
const GAP_ROLE_MAP: Readonly<
Record<
string,
{
role: WorkerRole;
followLinks: boolean;
tiers?: ReadonlyArray<import("../types").SourceTier>;
}
>
> = {
overview: { role: "breadth", followLinks: false },
mechanism: { role: "technical", followLinks: true },
history: { role: "breadth", followLinks: false },
current: { role: "recency", followLinks: false },
applications: { role: "breadth", followLinks: false },
challenges: { role: "critical", followLinks: false },
comparison: { role: "comparative", followLinks: false },
evidence: {
role: "academic",
followLinks: true,
tiers: ["academic", "government", "reference"],
},
expert: { role: "primary", followLinks: true, tiers: ["academic", "news"] },
future: { role: "recency", followLinks: false },
controversy: { role: "critical", followLinks: false },
economics: { role: "statistical", followLinks: false },
};
export async function buildAdaptiveGapFill(
topic: string,
coveredIds: ReadonlyArray<string>,
priorMessages: ReadonlyArray<AgentMessage>,
useAI: boolean,
status: StatusFn,
profile: DepthProfile,
): Promise<ReadonlyArray<AdaptiveGapPlan>> {
const gaps = detectGaps(coveredIds);
if (gaps.length === 0) {
status("All research dimensions covered - no gap queries needed");
return [];
}
status(`Gaps: ${gaps.map((g) => g.label).join(", ")}`);
const byRole = new Map<
WorkerRole,
{
dimIds: string[];
dimLabels: string[];
queries: string[];
followLinks: boolean;
tiers?: ReadonlyArray<import("../types").SourceTier>;
}
>();
// BUGFIX: Shorten the topic here too so the fallback queries are clean
const shortTopic = shortenTopic(topic);
for (const gap of gaps) {
const mapping = GAP_ROLE_MAP[gap.id] ?? {
role: "breadth" as WorkerRole,
followLinks: false,
};
const existing = byRole.get(mapping.role) ?? {
dimIds: [],
dimLabels: [],
queries: [],
followLinks: mapping.followLinks,
tiers: mapping.tiers,
};
existing.dimIds.push(gap.id);
existing.dimLabels.push(gap.label);
existing.queries.push(...gap.queries(shortTopic));
byRole.set(mapping.role, existing);
}
if (useAI) {
const followUpContext = priorMessages
.flatMap((m) => m.suggestedFollowUps)
.slice(0, 6);
const entries = Array.from(byRole.entries());
const aiResults = await Promise.allSettled(
entries.map(async ([role, group]) => {
const queryCount = Math.min(
group.dimLabels.length * 3,
profile.maxGapFillQueries,
);
const prompt = `You are an expert Google searcher. A research session on "${topic}" is missing these angles:
${group.dimLabels.join(", ")}
${followUpContext.length > 0 ? `Previous round suggested exploring:\n${followUpContext.join("\n")}\n` : ""}
Generate ${queryCount} specific, natural search engine queries to fill these gaps.
Rules:
1. DO NOT just append words to the topic string. Write human-like search queries.
2. Keep queries short and concise.
3. Return ONLY the queries, one per line. No numbering, no prefixes.
Queries:`;
return { role, raw: await callLoadedModel(prompt) };
}),
);
for (const result of aiResults) {
if (result.status !== "fulfilled" || !result.value.raw) continue;
const { role, raw } = result.value;
const parsed = parseLines(raw, profile.maxGapFillQueries);
if (parsed.length >= 2) {
const group = byRole.get(role);
if (group) group.queries = [...parsed];
}
}
}
const plans: AdaptiveGapPlan[] = [];
for (const [role, group] of byRole) {
plans.push({
role,
label: `Gap-fill: ${group.dimLabels.slice(0, 3).join(", ")}`,
queries: group.queries.slice(0, profile.maxGapFillQueries),
followLinks: group.followLinks,
preferredTiers: group.tiers,
});
}
status(`${plans.length} adaptive gap-fill worker(s) planned`);
return plans;
}
const STOP_WORDS = new Set([
"the", "a", "an", "is", "in", "of", "and", "or", "for", "to", "how", "what",
"why", "when", "does", "with", "from", "that", "could", "which", "about",
"their", "this", "these", "those", "would", "should", "current",
"hypothetical", "scenarios", "lead",
]);
function extractKeywords(topic: string): ReadonlyArray<string> {
return topic
.toLowerCase()
.replace(/[^a-z0-9\s]/g, "")
.split(/\s+/)
.filter((w) => w.length > 2 && !STOP_WORDS.has(w))
.slice(0, 8);
}
function shortenTopic(topic: string): string {
const colonIdx = topic.indexOf(":");
const dashIdx = topic.indexOf(" - ");
const sepIdx = colonIdx > 3 ? colonIdx : dashIdx > 3 ? dashIdx : -1;
let core: string;
if (sepIdx > 3 && sepIdx < topic.length * 0.6) {
core = topic.slice(0, sepIdx).trim();
} else {
core = topic;
}
const words = core
.replace(/[,;()]/g, " ")
.split(/\s+/)
.filter((w) => w.length > 1 && !STOP_WORDS.has(w.toLowerCase()));
let result = "";
let count = 0;
for (const w of words) {
if (count >= 6 || result.length + w.length > 58) break;
result += (result ? " " : "") + w;
count++;
}
if (result.length < 5) {
result = topic.split(/\s+/).slice(0, 5).join(" ");
}
return result;
}