| """ |
| X-RUDRA Web Search v2 |
| Single-file web research engine. |
| |
| Pipeline: |
| Model A + Model B -> keyword/query planning -> DuckDuckGo discovery |
| -> static HTTP -> Scrapling parsing -> DynamicFetcher/Playwright fallback |
| -> optional StealthyFetcher fallback -> evidence extraction -> ranking. |
| |
| Important: |
| - Browser automation is for rendering pages that require JavaScript. |
| - It does NOT attempt to defeat CAPTCHAs, access controls, paywalls, or authentication. |
| - Webpage text is untrusted data and never becomes an instruction to the agent. |
| - SSRF checks are applied before navigation. |
| """ |
|
|
| from __future__ import annotations |
|
|
| import asyncio |
| import hashlib |
| import ipaddress |
| import json |
| import logging |
| import os |
| import re |
| import socket |
| import sqlite3 |
| import time |
| from collections import Counter |
| from datetime import datetime, timezone |
| from typing import Any, Optional |
| from urllib.parse import parse_qsl, urlencode, urljoin, urlparse, urlunparse |
| from urllib.robotparser import RobotFileParser |
|
|
| import trafilatura |
| from bs4 import BeautifulSoup |
| from ddgs import DDGS |
| from fastapi import FastAPI, HTTPException |
| from pydantic import BaseModel, Field, ConfigDict |
|
|
| |
| |
| try: |
| from scrapling.fetchers import Fetcher, AsyncFetcher, DynamicFetcher, StealthyFetcher |
| SCRAPLING_AVAILABLE = True |
| except Exception: |
| Fetcher = AsyncFetcher = DynamicFetcher = StealthyFetcher = None |
| SCRAPLING_AVAILABLE = False |
|
|
| try: |
| from playwright.async_api import async_playwright |
| PLAYWRIGHT_AVAILABLE = True |
| except Exception: |
| async_playwright = None |
| PLAYWRIGHT_AVAILABLE = False |
|
|
|
|
| |
| |
| |
|
|
| M1_SPACE = os.getenv("M1_SPACE", "Shrijanagain/M1") |
| M2_SPACE = os.getenv("M2_SPACE", "Shrijanagain/M2") |
| HF_TOKEN = os.getenv("HF_TOKEN", "") |
|
|
| MAX_RESULTS = int(os.getenv("MAX_RESULTS", "12")) |
| MAX_ROUNDS = int(os.getenv("MAX_ROUNDS", "3")) |
| MAX_CONCURRENT = int(os.getenv("MAX_CONCURRENT", "4")) |
| HTTP_TIMEOUT = float(os.getenv("HTTP_TIMEOUT", "15")) |
| BROWSER_TIMEOUT_MS = int(os.getenv("BROWSER_TIMEOUT_MS", "25000")) |
| MAX_BODY = int(os.getenv("MAX_BODY", str(8 * 1024 * 1024))) |
| CACHE_TTL = int(os.getenv("CACHE_TTL", "3600")) |
| CACHE_DB = os.getenv("CACHE_DB", "x_rudra_web_v2.sqlite3") |
| BROWSER_MODE = os.getenv("BROWSER_MODE", "auto") |
| USE_STEALTH_FETCHER = os.getenv("USE_STEALTH_FETCHER", "1") == "1" |
| USER_AGENT = os.getenv( |
| "USER_AGENT", |
| "X-RUDRA-ResearchBot/2.0 (+research client)", |
| ) |
|
|
| logging.basicConfig( |
| level=os.getenv("LOG_LEVEL", "INFO"), |
| format="%(asctime)s | %(levelname)s | %(name)s | %(message)s", |
| ) |
| log = logging.getLogger("x-rudra-web") |
|
|
|
|
| |
| |
| |
|
|
| class SearchRequest(BaseModel): |
| model_config = ConfigDict(extra="ignore") |
| question: str = Field(min_length=1, max_length=12000) |
| max_results: int = Field(default=MAX_RESULTS, ge=1, le=30) |
| max_rounds: int = Field(default=MAX_ROUNDS, ge=1, le=5) |
| use_models: bool = True |
| freshness: str = "auto" |
|
|
|
|
| class SearchPlan(BaseModel): |
| intent: str = "general" |
| topic: str = "" |
| entities: list[str] = Field(default_factory=list) |
| keywords: list[str] = Field(default_factory=list) |
| queries: list[str] = Field(default_factory=list) |
| opposing_queries: list[str] = Field(default_factory=list) |
| time_constraint: Optional[str] = None |
| source_types: list[str] = Field(default_factory=list) |
|
|
|
|
| class SearchResult(BaseModel): |
| title: str |
| url: str |
| snippet: str = "" |
| query: str = "" |
| rank: int = 0 |
|
|
|
|
| class PageDocument(BaseModel): |
| url: str |
| final_url: str |
| domain: str |
| title: str = "" |
| description: str = "" |
| author: Optional[str] = None |
| published_at: Optional[str] = None |
| modified_at: Optional[str] = None |
| text: str = "" |
| passages: list[dict[str, Any]] = Field(default_factory=list) |
| source_score: float = 0.0 |
| freshness_score: float = 0.0 |
| relevance_score: float = 0.0 |
| fetch_method: str = "unknown" |
| status_code: Optional[int] = None |
| prompt_injection_detected: bool = False |
|
|
|
|
| class Claim(BaseModel): |
| claim: str |
| source_url: str |
| passage: str |
| support_score: float = 0.0 |
|
|
|
|
| class SearchReport(BaseModel): |
| status: str |
| question: str |
| plan: SearchPlan |
| queries: list[str] |
| results: list[SearchResult] |
| sources: list[PageDocument] |
| claims: list[Claim] |
| contradictions: list[dict[str, Any]] |
| failures: list[dict[str, Any]] |
| rounds: int |
| stopping_reason: str |
|
|
|
|
| |
| |
| |
|
|
| TRACKING = { |
| "utm_source", "utm_medium", "utm_campaign", "utm_term", "utm_content", |
| "gclid", "fbclid", "mc_cid", "mc_eid", "ref", "ref_src", |
| } |
|
|
|
|
| def clean_text(value: str) -> str: |
| return re.sub(r"\s+", " ", value or "").strip() |
|
|
|
|
| def normalize_url(url: str) -> str: |
| try: |
| p = urlparse(url.strip()) |
| if p.scheme.lower() not in {"http", "https"}: |
| return "" |
| query = [ |
| (k, v) for k, v in parse_qsl(p.query, keep_blank_values=True) |
| if k.lower() not in TRACKING |
| ] |
| return urlunparse( |
| p._replace( |
| scheme=p.scheme.lower(), |
| netloc=p.netloc.lower(), |
| query=urlencode(query), |
| fragment="", |
| ) |
| ) |
| except Exception: |
| return "" |
|
|
|
|
| def domain(url: str) -> str: |
| return (urlparse(url).hostname or "").lower() |
|
|
|
|
| def tokens(text: str) -> set[str]: |
| return set(re.findall(r"[a-zA-Z0-9][a-zA-Z0-9_-]{2,}", text.lower())) |
|
|
|
|
| def similarity(a: str, b: str) -> float: |
| x, y = tokens(a), tokens(b) |
| return len(x & y) / max(1, len(x | y)) |
|
|
|
|
| def sha(value: str) -> str: |
| return hashlib.sha256(value.encode()).hexdigest() |
|
|
|
|
| def private_host(host: str) -> bool: |
| if not host or host.lower() in {"localhost", "localhost.localdomain"}: |
| return True |
| try: |
| addresses = socket.getaddrinfo(host, None) |
| except OSError: |
| return True |
| for info in addresses: |
| try: |
| ip = ipaddress.ip_address(info[4][0]) |
| except ValueError: |
| return True |
| if ( |
| ip.is_private or ip.is_loopback or ip.is_link_local |
| or ip.is_multicast or ip.is_reserved or ip.is_unspecified |
| ): |
| return True |
| return False |
|
|
|
|
| def safe_url(url: str) -> bool: |
| url = normalize_url(url) |
| if not url: |
| return False |
| p = urlparse(url) |
| return p.scheme in {"http", "https"} and not private_host(p.hostname or "") |
|
|
|
|
| |
| |
| |
|
|
| class Cache: |
| def __init__(self, path: str): |
| self.path = path |
| with sqlite3.connect(self.path) as c: |
| c.execute( |
| "CREATE TABLE IF NOT EXISTS cache " |
| "(k TEXT PRIMARY KEY, v TEXT, t REAL)" |
| ) |
|
|
| def get(self, key: str) -> Any: |
| try: |
| with sqlite3.connect(self.path) as c: |
| row = c.execute( |
| "SELECT v,t FROM cache WHERE k=?", (key,) |
| ).fetchone() |
| if not row: |
| return None |
| if time.time() - row[1] > CACHE_TTL: |
| return None |
| return json.loads(row[0]) |
| except Exception: |
| return None |
|
|
| def set(self, key: str, value: Any): |
| try: |
| with sqlite3.connect(self.path) as c: |
| c.execute( |
| "INSERT OR REPLACE INTO cache(k,v,t) VALUES(?,?,?)", |
| (key, json.dumps(value, ensure_ascii=False), time.time()), |
| ) |
| except Exception as exc: |
| log.debug("cache write failed: %s", exc) |
|
|
|
|
| CACHE = Cache(CACHE_DB) |
|
|
|
|
| |
| |
| |
|
|
| class SpaceModel: |
| def __init__(self, space: str): |
| self.space = space |
| self.client = None |
|
|
| def _client(self): |
| if self.client is None: |
| from gradio_client import Client |
| kwargs = {"hf_token": HF_TOKEN} if HF_TOKEN else {} |
| self.client = Client(self.space, **kwargs) |
| return self.client |
|
|
| async def generate(self, prompt: str) -> Optional[str]: |
| def run(): |
| try: |
| result = self._client().predict( |
| prompt, 1200, 0.2, 0.9, api_name="/generate" |
| ) |
| return result if isinstance(result, str) else str(result) |
| except Exception as exc: |
| log.warning("Space %s unavailable: %s", self.space, exc) |
| return None |
| return await asyncio.to_thread(run) |
|
|
|
|
| class Planner: |
| def __init__(self): |
| self.a = SpaceModel(M1_SPACE) |
| self.b = SpaceModel(M2_SPACE) |
|
|
| @staticmethod |
| def fallback(q: str) -> SearchPlan: |
| stop = { |
| "what","when","where","which","who","why","how","is","are", |
| "was","were","the","a","an","of","to","for","in","on","and", |
| "or","with","about","tell","me","please" |
| } |
| words = [x for x in re.findall(r"[a-zA-Z0-9_-]{3,}", q.lower()) |
| if x not in stop] |
| words = list(dict.fromkeys(words)) |
| current = any(x in q.lower() for x in |
| ("latest","today","current","recent","2026")) |
| return SearchPlan( |
| intent="current_factual" if current else "general", |
| topic=" ".join(words[:10]), |
| keywords=words[:20], |
| time_constraint="recent" if current else "auto", |
| source_types=["official","primary","academic","reputable"], |
| ) |
|
|
| @staticmethod |
| def parse(text: Optional[str]) -> Optional[SearchPlan]: |
| if not text: |
| return None |
| m = re.search(r"```(?:json)?\s*(\{.*?\})\s*```", text, re.S) |
| raw = m.group(1) if m else None |
| if raw is None: |
| a, b = text.find("{"), text.rfind("}") |
| if a >= 0 and b > a: |
| raw = text[a:b + 1] |
| if not raw: |
| return None |
| try: |
| return SearchPlan.model_validate(json.loads(raw)) |
| except Exception: |
| return None |
|
|
| async def one(self, model: SpaceModel, question: str): |
| prompt = f""" |
| X-RUDRA search planner. Return JSON only. |
| Schema: |
| {{ |
| "intent":"general|current_factual|research|comparison|explanatory", |
| "topic":"...", |
| "entities":[], |
| "keywords":[], |
| "queries":[], |
| "opposing_queries":[], |
| "time_constraint":null, |
| "source_types":[] |
| }} |
| Extract search concepts and generate independent queries. |
| Never answer the user question. |
| QUESTION: {question} |
| """ |
| return self.parse(await model.generate(prompt)) |
|
|
| async def plan(self, question: str, use_models=True) -> SearchPlan: |
| fallback = self.fallback(question) |
| if not use_models: |
| return fallback |
| a, b = await asyncio.gather( |
| self.one(self.a, question), |
| self.one(self.b, question), |
| ) |
| plans = [x for x in (a, b) if x] |
| if not plans: |
| return fallback |
|
|
| merged = fallback.model_copy(deep=True) |
| merged.keywords = [] |
| merged.queries = [] |
| merged.entities = [] |
| merged.opposing_queries = [] |
| merged.source_types = [] |
|
|
| for p in plans: |
| for field in ("keywords","queries","entities", |
| "opposing_queries","source_types"): |
| for x in getattr(p, field): |
| if x and x not in getattr(merged, field): |
| getattr(merged, field).append(x) |
| if p.topic: |
| merged.topic = p.topic |
| if p.time_constraint: |
| merged.time_constraint = p.time_constraint |
| if p.intent: |
| merged.intent = p.intent |
|
|
| return merged |
|
|
|
|
| |
| |
| |
|
|
| class QueryBuilder: |
| @staticmethod |
| def build(question: str, plan: SearchPlan) -> list[str]: |
| out = [] |
|
|
| def add(x): |
| x = clean_text(x) |
| if x and x not in out: |
| out.append(x) |
|
|
| for q in plan.queries: |
| add(q) |
|
|
| core = plan.topic or " ".join(plan.keywords[:10]) |
| add(core) |
| add(f"{core} latest") |
| add(f"{core} official") |
| add(f"{core} primary source") |
| add(f"{core} research") |
| add(f"{core} report") |
| add(f"{core} evidence") |
| add(f"{core} criticism") |
| add(f"{core} limitations") |
|
|
| if plan.time_constraint not in (None, "auto"): |
| add(f"{core} {plan.time_constraint}") |
|
|
| for q in plan.opposing_queries: |
| add(q) |
|
|
| return list(dict.fromkeys(out))[:30] |
|
|
|
|
| class DDG: |
| async def search(self, query: str, limit: int) -> list[SearchResult]: |
| key = "ddg:" + sha(f"{query}:{limit}") |
| cached = CACHE.get(key) |
| if cached: |
| return [SearchResult.model_validate(x) for x in cached] |
|
|
| def run(): |
| out = [] |
| try: |
| with DDGS(timeout=int(HTTP_TIMEOUT)) as d: |
| for i, item in enumerate( |
| d.text(query, max_results=limit), 1 |
| ): |
| url = normalize_url(item.get("href", "")) |
| if url: |
| out.append(SearchResult( |
| title=clean_text(item.get("title", "")), |
| url=url, |
| snippet=clean_text(item.get("body", "")), |
| query=query, |
| rank=i, |
| )) |
| except Exception as exc: |
| log.warning("DDG failed: %s", exc) |
| return out |
|
|
| result = await asyncio.to_thread(run) |
| CACHE.set(key, [x.model_dump() for x in result]) |
| return result |
|
|
|
|
| |
| |
| |
|
|
| class Robots: |
| def __init__(self): |
| self.cache = {} |
|
|
| async def allowed(self, url: str) -> bool: |
| host = domain(url) |
| if not host: |
| return False |
| if host in self.cache: |
| return self.cache[host] |
|
|
| robots_url = f"{urlparse(url).scheme}://{host}/robots.txt" |
|
|
| def read(): |
| rp = RobotFileParser() |
| rp.set_url(robots_url) |
| try: |
| rp.read() |
| return rp.can_fetch(USER_AGENT, url) |
| except Exception: |
| return True |
|
|
| value = await asyncio.to_thread(read) |
| self.cache[host] = value |
| return value |
|
|
|
|
| |
| |
| |
|
|
| class ScraplingStatic: |
| async def fetch(self, url: str) -> Optional[dict[str, Any]]: |
| if not SCRAPLING_AVAILABLE or not safe_url(url): |
| return None |
|
|
| def run(): |
| try: |
| page = Fetcher.get( |
| url, |
| timeout=int(HTTP_TIMEOUT), |
| retries=2, |
| follow_redirects="safe", |
| stealthy_headers=True, |
| impersonate="chrome", |
| ) |
| html = getattr(page, "html_content", None) |
| if callable(html): |
| html = html() |
| if html is None: |
| html = getattr(page, "text", "") or "" |
| status = getattr(page, "status", None) |
| final = getattr(page, "url", None) or url |
| return { |
| "html": html, |
| "status": status, |
| "final_url": final, |
| "method": "scrapling-http", |
| } |
| except Exception as exc: |
| log.debug("Scrapling static failed %s: %s", url, exc) |
| return None |
|
|
| return await asyncio.to_thread(run) |
|
|
|
|
| |
| |
| |
|
|
| class PlaywrightBrowser: |
| def __init__(self): |
| self.playwright = None |
| self.browser = None |
| self.lock = asyncio.Lock() |
|
|
| async def start(self): |
| if not PLAYWRIGHT_AVAILABLE: |
| return |
| async with self.lock: |
| if self.browser: |
| return |
| self.playwright = await async_playwright().start() |
| self.browser = await self.playwright.chromium.launch( |
| headless=True, |
| args=[ |
| "--disable-dev-shm-usage", |
| "--no-first-run", |
| "--no-default-browser-check", |
| ], |
| ) |
|
|
| async def close(self): |
| if self.browser: |
| await self.browser.close() |
| self.browser = None |
| if self.playwright: |
| await self.playwright.stop() |
| self.playwright = None |
|
|
| async def fetch(self, url: str) -> Optional[dict[str, Any]]: |
| if not PLAYWRIGHT_AVAILABLE or not safe_url(url): |
| return None |
|
|
| await self.start() |
| if not self.browser: |
| return None |
|
|
| context = await self.browser.new_context( |
| user_agent=USER_AGENT, |
| java_script_enabled=True, |
| ignore_https_errors=False, |
| ) |
| page = await context.new_page() |
|
|
| try: |
| response = await page.goto( |
| url, |
| wait_until="domcontentloaded", |
| timeout=BROWSER_TIMEOUT_MS, |
| ) |
| await page.wait_for_timeout(500) |
| html = await page.content() |
| final = page.url |
|
|
| if not safe_url(final): |
| return None |
|
|
| return { |
| "html": html[:MAX_BODY], |
| "status": response.status if response else None, |
| "final_url": final, |
| "method": "playwright", |
| } |
| except Exception as exc: |
| log.debug("Playwright failed %s: %s", url, exc) |
| return None |
| finally: |
| await page.close() |
| await context.close() |
|
|
|
|
| |
| |
| |
|
|
| class ScraplingBrowser: |
| async def dynamic(self, url: str) -> Optional[dict[str, Any]]: |
| if not SCRAPLING_AVAILABLE or not safe_url(url): |
| return None |
|
|
| def run(): |
| try: |
| page = DynamicFetcher.fetch( |
| url, |
| headless=True, |
| timeout=BROWSER_TIMEOUT_MS, |
| disable_resources=True, |
| block_ads=True, |
| network_idle=False, |
| load_dom=True, |
| ) |
| html = getattr(page, "html_content", "") or "" |
| return { |
| "html": html[:MAX_BODY], |
| "status": getattr(page, "status", None), |
| "final_url": getattr(page, "url", None) or url, |
| "method": "scrapling-dynamic", |
| } |
| except Exception as exc: |
| log.debug("Scrapling Dynamic failed %s: %s", url, exc) |
| return None |
|
|
| return await asyncio.to_thread(run) |
|
|
| async def stealth(self, url: str) -> Optional[dict[str, Any]]: |
| if not SCRAPLING_AVAILABLE or not USE_STEALTH_FETCHER: |
| return None |
| if not safe_url(url): |
| return None |
|
|
| |
| |
| def run(): |
| try: |
| page = StealthyFetcher.fetch( |
| url, |
| headless=True, |
| timeout=BROWSER_TIMEOUT_MS, |
| disable_resources=True, |
| block_ads=True, |
| network_idle=False, |
| load_dom=True, |
| ) |
| html = getattr(page, "html_content", "") or "" |
| return { |
| "html": html[:MAX_BODY], |
| "status": getattr(page, "status", None), |
| "final_url": getattr(page, "url", None) or url, |
| "method": "scrapling-stealth", |
| } |
| except Exception as exc: |
| log.debug("Scrapling Stealth failed %s: %s", url, exc) |
| return None |
|
|
| return await asyncio.to_thread(run) |
|
|
|
|
| |
| |
| |
|
|
| INJECTION_PATTERNS = [ |
| r"ignore\s+(all|any|the)\s+(previous|prior|system)\s+instructions", |
| r"reveal\s+(your|the)\s+(system|developer)\s+prompt", |
| r"disregard\s+your\s+instructions", |
| r"act\s+as\s+the\s+system", |
| r"show\s+your\s+hidden\s+instructions", |
| ] |
|
|
|
|
| def detect_injection(text: str) -> bool: |
| x = text.lower() |
| return any(re.search(p, x) for p in INJECTION_PATTERNS) |
|
|
|
|
| def metadata(html: str) -> dict[str, Any]: |
| soup = BeautifulSoup(html, "html.parser") |
| title = clean_text(soup.title.get_text(" ")) if soup.title else "" |
|
|
| def meta(name=None, prop=None): |
| tag = soup.find( |
| "meta", |
| attrs=({"name": re.compile(f"^{name}$", re.I)} |
| if name else {"property": prop}), |
| ) |
| return clean_text(tag.get("content", "")) if tag else "" |
|
|
| canonical = "" |
| link = soup.find("link", attrs={ |
| "rel": lambda x: x and "canonical" in x |
| }) |
| if link: |
| canonical = link.get("href", "") |
|
|
| published = meta(prop="article:published_time") |
| modified = meta(prop="article:modified_time") |
| author = meta(name="author") |
|
|
| for script in soup.find_all( |
| "script", attrs={"type": "application/ld+json"} |
| ): |
| try: |
| data = json.loads(script.string or "") |
| items = data if isinstance(data, list) else [data] |
| for item in items: |
| if not isinstance(item, dict): |
| continue |
| title = title or clean_text(item.get("headline", "")) |
| published = published or item.get("datePublished", "") |
| modified = modified or item.get("dateModified", "") |
| if isinstance(item.get("author"), dict): |
| author = author or item["author"].get("name", "") |
| except Exception: |
| continue |
|
|
| return { |
| "title": title, |
| "description": meta(name="description"), |
| "canonical": canonical, |
| "author": author or None, |
| "published": published or None, |
| "modified": modified or None, |
| } |
|
|
|
|
| def extract_text(html: str) -> str: |
| try: |
| value = trafilatura.extract( |
| html, |
| include_tables=True, |
| include_links=False, |
| include_images=False, |
| favor_precision=True, |
| ) |
| if value: |
| return clean_text(value) |
| except Exception: |
| pass |
|
|
| soup = BeautifulSoup(html, "html.parser") |
| for tag in soup([ |
| "script","style","noscript","svg","canvas", |
| "iframe","nav","footer","form" |
| ]): |
| tag.decompose() |
| return clean_text(soup.get_text(" ", strip=True)) |
|
|
|
|
| def passages(question: str, text: str, limit=8): |
| sentences = re.split(r"(?<=[.!?])\s+", clean_text(text)) |
| chunks, current = [], "" |
|
|
| for sentence in sentences: |
| if len(current) + len(sentence) < 1800: |
| current = f"{current} {sentence}".strip() |
| else: |
| if current: |
| chunks.append(current) |
| current = sentence |
| if current: |
| chunks.append(current) |
|
|
| scored = [] |
| for chunk in chunks: |
| scored.append({ |
| "text": chunk, |
| "relevance": similarity(question, chunk), |
| }) |
|
|
| return sorted( |
| scored, key=lambda x: x["relevance"], reverse=True |
| )[:limit] |
|
|
|
|
| |
| |
| |
|
|
| def authority(host: str) -> float: |
| h = host.lower() |
| if h.endswith(".gov") or ".gov." in h: |
| return 1.0 |
| if h.endswith(".edu") or ".edu." in h: |
| return 0.95 |
| if h.endswith(".org"): |
| return 0.72 |
|
|
| strong = { |
| "unesco.org","un.org","who.int","worldbank.org", |
| "oecd.org","nasa.gov","nih.gov","reuters.com", |
| "apnews.com","bbc.com" |
| } |
| return 0.95 if any( |
| h == x or h.endswith("." + x) for x in strong |
| ) else 0.45 |
|
|
|
|
| def freshness(value: Optional[str], mode: str) -> float: |
| if not value: |
| return 0.5 |
| try: |
| from dateutil import parser |
| dt = parser.parse(value) |
| if not dt.tzinfo: |
| dt = dt.replace(tzinfo=timezone.utc) |
| age = max( |
| 0, |
| (datetime.now(timezone.utc) - dt.astimezone(timezone.utc)) |
| .total_seconds() / 86400, |
| ) |
| if mode in {"latest","recent","current"}: |
| return max(0, 1 - age / 365) |
| return 0.7 |
| except Exception: |
| return 0.5 |
|
|
|
|
| def source_score(host, rel, fresh): |
| return round( |
| authority(host) * 0.35 |
| + rel * 0.40 |
| + fresh * 0.25, |
| 4, |
| ) |
|
|
|
|
| |
| |
| |
|
|
| class FetchOrchestrator: |
| def __init__(self): |
| self.static = ScraplingStatic() |
| self.dynamic = ScraplingBrowser() |
| self.playwright = PlaywrightBrowser() |
| self.robots = Robots() |
| self.sem = asyncio.Semaphore(MAX_CONCURRENT) |
|
|
| async def fetch(self, url: str, browser_hint=False): |
| if not safe_url(url): |
| return None, "unsafe_url" |
|
|
| if not await self.robots.allowed(url): |
| return None, "robots_disallowed" |
|
|
| |
| result = await self.static.fetch(url) |
| if result and result.get("html"): |
| return result, "scrapling-http" |
|
|
| if BROWSER_MODE == "never": |
| return None, "static_failed" |
|
|
| |
| if browser_hint or BROWSER_MODE == "browser": |
| result = await self.dynamic.dynamic(url) |
| if result and result.get("html"): |
| return result, "scrapling-dynamic" |
|
|
| |
| if browser_hint and USE_STEALTH_FETCHER: |
| result = await self.dynamic.stealth(url) |
| if result and result.get("html"): |
| return result, "scrapling-stealth" |
|
|
| |
| if PLAYWRIGHT_AVAILABLE: |
| result = await self.playwright.fetch(url) |
| if result and result.get("html"): |
| return result, "playwright" |
|
|
| return None, "all_fetchers_failed" |
|
|
| async def close(self): |
| await self.playwright.close() |
|
|
|
|
| |
| |
| |
|
|
| class EvidenceEngine: |
| @staticmethod |
| def claims(question, sources): |
| claims = [] |
| for source in sources: |
| for p in source.passages: |
| rel = float(p.get("relevance", 0)) |
| if rel >= 0.18 and len(p["text"]) >= 50: |
| claims.append( |
| Claim( |
| claim=p["text"][:1500], |
| source_url=source.url, |
| passage=p["text"][:2000], |
| support_score=round( |
| rel * source.source_score, 4 |
| ), |
| ) |
| ) |
| return claims |
|
|
| @staticmethod |
| def dedupe(claims): |
| out = [] |
| for c in claims: |
| if not any( |
| similarity(c.claim, x.claim) >= 0.90 |
| for x in out |
| ): |
| out.append(c) |
| return out |
|
|
| @staticmethod |
| def contradictions(claims): |
| result = [] |
| for i, a in enumerate(claims): |
| for b in claims[i + 1:]: |
| if a.source_url == b.source_url: |
| continue |
| sim = similarity(a.claim, b.claim) |
| if sim < 0.45: |
| continue |
|
|
| na = bool(re.search( |
| r"\b(not|no|never|false|denied|did not)\b", |
| a.claim.lower(), |
| )) |
| nb = bool(re.search( |
| r"\b(not|no|never|false|denied|did not)\b", |
| b.claim.lower(), |
| )) |
|
|
| if na != nb: |
| result.append({ |
| "claim_a": a.claim, |
| "source_a": a.source_url, |
| "claim_b": b.claim, |
| "source_b": b.source_url, |
| "status": "needs_review", |
| }) |
| return result[:20] |
|
|
|
|
| |
| |
| |
|
|
| class XrudraWebSearch: |
| def __init__(self): |
| self.planner = Planner() |
| self.ddg = DDG() |
| self.fetcher = FetchOrchestrator() |
|
|
| async def search( |
| self, |
| question: str, |
| max_results=MAX_RESULTS, |
| max_rounds=MAX_ROUNDS, |
| use_models=True, |
| freshness_mode="auto", |
| ): |
| plan = await self.planner.plan(question, use_models) |
| queries = QueryBuilder.build(question, plan) |
|
|
| all_results = [] |
| seen = set() |
| all_sources = [] |
| failures = [] |
|
|
| for round_no in range(1, max_rounds + 1): |
| if round_no > 1: |
| followups = [ |
| f"{plan.topic} primary source", |
| f"{plan.topic} official evidence", |
| f"{plan.topic} independent evidence", |
| f"{plan.topic} criticism", |
| ] |
| queries.extend(q for q in followups if q not in queries) |
|
|
| round_queries = queries[ |
| (round_no - 1) * 8: round_no * 8 |
| ] |
|
|
| if not round_queries: |
| break |
|
|
| batches = await asyncio.gather( |
| *[ |
| self.ddg.search(q, max_results) |
| for q in round_queries |
| ], |
| return_exceptions=True, |
| ) |
|
|
| for batch in batches: |
| if isinstance(batch, Exception): |
| continue |
| for r in batch: |
| r.url = normalize_url(r.url) |
| if r.url and r.url not in seen: |
| seen.add(r.url) |
| all_results.append(r) |
|
|
| |
| |
| ranked = sorted( |
| all_results, |
| key=lambda r: similarity( |
| question, |
| r.title + " " + r.snippet |
| ), |
| reverse=True, |
| ) |
|
|
| targets = ranked[:max_results] |
| jobs = [ |
| self._fetch_one( |
| r, |
| question, |
| plan, |
| freshness_mode, |
| ) |
| for r in targets |
| ] |
|
|
| fetched = await asyncio.gather( |
| *jobs, |
| return_exceptions=True, |
| ) |
|
|
| for item in fetched: |
| if isinstance(item, Exception): |
| continue |
| source, failure = item |
| if source: |
| if not any( |
| x.final_url == source.final_url |
| for x in all_sources |
| ): |
| all_sources.append(source) |
| elif failure: |
| failures.append(failure) |
|
|
| |
| strong = [ |
| s for s in all_sources |
| if s.source_score >= 0.45 |
| and s.relevance_score >= 0.18 |
| ] |
| domains = {s.domain for s in strong} |
|
|
| if len(strong) >= 3 and len(domains) >= 2: |
| return await self._report( |
| question, plan, queries, all_results, |
| all_sources, failures, round_no, |
| "sufficient_independent_evidence", |
| ) |
|
|
| return await self._report( |
| question, plan, queries, all_results, |
| all_sources, failures, max_rounds, |
| "max_rounds_reached", |
| ) |
|
|
| async def _fetch_one( |
| self, result, question, plan, freshness_mode |
| ): |
| browser_hint = any( |
| x in result.snippet.lower() |
| for x in ( |
| "javascript", "dynamic", "interactive", |
| "app", "dashboard" |
| ) |
| ) |
|
|
| try: |
| async with self.fetcher.sem: |
| payload, method = await self.fetcher.fetch( |
| result.url, |
| browser_hint=browser_hint, |
| ) |
|
|
| if not payload: |
| return None, { |
| "url": result.url, |
| "error": method, |
| } |
|
|
| final_url = normalize_url( |
| payload.get("final_url") or result.url |
| ) |
| if not safe_url(final_url): |
| return None, { |
| "url": result.url, |
| "error": "unsafe_final_redirect", |
| } |
|
|
| html = payload.get("html", "") |
| if not html: |
| return None, { |
| "url": result.url, |
| "error": "empty_html", |
| } |
|
|
| meta = metadata(html) |
| text = extract_text(html) |
| if not text: |
| return None, { |
| "url": result.url, |
| "error": "empty_text", |
| } |
|
|
| p = passages(question, text) |
| rel = max( |
| [x["relevance"] for x in p], |
| default=0.0, |
| ) |
| fresh = freshness( |
| meta.get("published") or meta.get("modified"), |
| plan.time_constraint |
| if plan.time_constraint not in (None, "auto") |
| else freshness_mode, |
| ) |
|
|
| src = PageDocument( |
| url=result.url, |
| final_url=final_url, |
| domain=domain(final_url), |
| title=meta["title"] or result.title, |
| description=meta["description"], |
| author=meta["author"], |
| published_at=meta["published"], |
| modified_at=meta["modified"], |
| text=text[:100000], |
| passages=p, |
| source_score=source_score( |
| domain(final_url), rel, fresh |
| ), |
| freshness_score=fresh, |
| relevance_score=rel, |
| fetch_method=method, |
| status_code=payload.get("status"), |
| prompt_injection_detected=detect_injection( |
| text[:100000] |
| ), |
| ) |
| return src, None |
|
|
| except Exception as exc: |
| return None, { |
| "url": result.url, |
| "error": f"{type(exc).__name__}:{exc}", |
| } |
|
|
| async def _report( |
| self, question, plan, queries, results, |
| sources, failures, rounds, reason |
| ): |
| |
| selected = [] |
| counts = Counter() |
|
|
| for s in sorted( |
| sources, |
| key=lambda x: x.source_score, |
| reverse=True, |
| ): |
| if counts[s.domain] >= 2: |
| continue |
| counts[s.domain] += 1 |
| selected.append(s) |
| if len(selected) >= MAX_RESULTS: |
| break |
|
|
| claims = EvidenceEngine.dedupe( |
| EvidenceEngine.claims(question, selected) |
| ) |
| contradictions = EvidenceEngine.contradictions(claims) |
|
|
| return SearchReport( |
| status="success", |
| question=question, |
| plan=plan, |
| queries=list(dict.fromkeys(queries)), |
| results=results[:MAX_RESULTS * 3], |
| sources=selected, |
| claims=claims[:100], |
| contradictions=contradictions, |
| failures=failures[:100], |
| rounds=rounds, |
| stopping_reason=reason, |
| ) |
|
|
|
|
| |
| |
| |
|
|
| app = FastAPI( |
| title="X-RUDRA Web Search v2", |
| version="2.0", |
| ) |
|
|
|
|
| ENGINE = XrudraWebSearch() |
|
|
|
|
| @app.get("/health") |
| async def health(): |
| return { |
| "status": "ok", |
| "scrapling": SCRAPLING_AVAILABLE, |
| "playwright": PLAYWRIGHT_AVAILABLE, |
| "m1": M1_SPACE, |
| "m2": M2_SPACE, |
| "browser_mode": BROWSER_MODE, |
| } |
|
|
|
|
| @app.post("/v2/web-search", response_model=SearchReport) |
| async def web_search(req: SearchRequest): |
| try: |
| return await ENGINE.search( |
| question=req.question, |
| max_results=req.max_results, |
| max_rounds=req.max_rounds, |
| use_models=req.use_models, |
| freshness_mode=req.freshness, |
| ) |
| except Exception as exc: |
| log.exception("web search failed") |
| raise HTTPException( |
| status_code=500, |
| detail="web search failed", |
| ) from exc |
|
|
|
|
| @app.on_event("shutdown") |
| async def shutdown(): |
| await ENGINE.fetcher.close() |
|
|
|
|
| |
| |
| |
|
|
| async def main(): |
| import argparse |
|
|
| p = argparse.ArgumentParser() |
| p.add_argument("question") |
| p.add_argument("--results", type=int, default=10) |
| p.add_argument("--rounds", type=int, default=3) |
| p.add_argument("--no-models", action="store_true") |
| args = p.parse_args() |
|
|
| report = await ENGINE.search( |
| args.question, |
| max_results=args.results, |
| max_rounds=args.rounds, |
| use_models=not args.no_models, |
| ) |
| print(json.dumps( |
| report.model_dump(), |
| ensure_ascii=False, |
| indent=2, |
| )) |
|
|
|
|
| if __name__ == "__main__": |
| asyncio.run(main()) |
|
|