hermes-agent/plugins/web/ddgs/provider.py
kshitij 646e71a9be fix: sanitize subprocess env for DDGS worker
os.environ.copy() passes all Hermes secrets (gateway tokens, API keys,
dashboard session tokens) into the DDGS child process. Use
_sanitize_subprocess_env() to strip Hermes-managed secrets before
spawning the worker.
2026-07-21 11:51:01 +05:30

360 lines
13 KiB
Python

"""DuckDuckGo search — plugin form (via the ``ddgs`` package).
Subclasses the plugin-facing :class:`agent.web_search_provider.WebSearchProvider`.
The legacy in-tree module ``tools.web_providers.ddgs`` was removed in the
same commit that moved this code under ``plugins/``; this file is now the
canonical implementation.
The ``ddgs`` package is an optional dependency. ``is_available()`` reflects
whether the package is importable; the plugin still registers either way so
``hermes tools`` can prompt the user to install it.
Isolation note (#68096): ``ddgs``/``primp`` can block inside native code while
holding the Python GIL. A ``ThreadPoolExecutor`` + ``future.result(timeout=…)``
cap (see #52118) cannot fire in that state — the waiter never reacquires the
GIL — so the whole Hermes process freezes through Ctrl+C/SIGTERM. Each search
therefore runs in a disposable child process the parent can terminate/kill.
"""
from __future__ import annotations
import concurrent.futures as cf
import json
import logging
import os
import subprocess
import sys
import time
from typing import Any, Dict, Optional
from agent.web_search_provider import WebSearchProvider
logger = logging.getLogger(__name__)
# Overall wall-clock cap for a single ddgs search. The DDGS constructor's
# ``timeout`` only bounds individual HTTP requests; ddgs's multi-engine retry
# loop has no overall cap, so a slow/rate-limited DuckDuckGo response can hang
# the (single, shared) agent loop indefinitely (#36776). Enforce a hard cap
# here by killing a disposable worker process (#68096).
_SEARCH_TIMEOUT_SECS = 30
# How often the parent polls stdout / interrupt flag while waiting.
_POLL_INTERVAL_SECS = 0.1
# After terminate(), wait this long before escalating to kill().
_TERMINATE_GRACE_SECS = 1.0
class _SearchInterrupted(Exception):
"""Raised when tools.interrupt.is_interrupted() trips during a search wait."""
def _run_ddgs_search(query: str, safe_limit: int) -> list[dict[str, Any]]:
"""Run the blocking ddgs query and return normalized hits.
Module-level (not a closure) so the child worker can import it and so
tests can patch it for in-process unit tests. ``DDGS(timeout=…)`` bounds
each individual HTTP request; the overall wall-clock cap is enforced by
the parent via process timeout (#68096).
"""
from ddgs import DDGS # type: ignore
results: list[dict[str, Any]] = []
with DDGS(timeout=10) as client:
for i, hit in enumerate(client.text(query, max_results=safe_limit)):
if i >= safe_limit:
break
url = str(hit.get("href") or hit.get("url") or "")
results.append(
{
"title": str(hit.get("title", "")),
"url": url,
"description": str(hit.get("body", "")),
"position": i + 1,
}
)
return results
# Optional test-only hook name forwarded to the child (see _search_worker.py).
# Production search() never sets this.
_test_hook: Optional[str] = None
# Last worker Popen started by ``_run_ddgs_search_bounded`` (test reap checks).
_last_worker_proc: Optional[subprocess.Popen] = None
def _plugins_path_entry() -> str:
"""Return the ``sys.path`` entry that makes ``import plugins`` work.
Prefer the live ``plugins`` package location over counting ``dirname``s from
this file — that stays correct for source checkouts and site-packages.
"""
try:
import plugins as plugins_pkg
pkg_file = getattr(plugins_pkg, "__file__", None)
if pkg_file:
return os.path.dirname(os.path.dirname(os.path.abspath(pkg_file)))
except Exception: # noqa: BLE001 — fall through to path-walk fallback
pass
return os.path.dirname(
os.path.dirname(
os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
)
)
def _terminate_and_reap(
proc: Optional[subprocess.Popen],
*,
grace: float = _TERMINATE_GRACE_SECS,
) -> None:
"""Terminate a worker, escalate to kill, and wait so no orphan remains.
Does not close the parent's pipe ends — the caller must finish any
``communicate()``/reader first. Closing stdout while another thread is
blocked in ``read()`` deadlocks on some platforms.
"""
if proc is None:
return
def _wait_until_dead(seconds: float) -> bool:
deadline = time.monotonic() + seconds
while time.monotonic() < deadline:
if proc.poll() is not None:
return True
time.sleep(0.05)
return proc.poll() is not None
try:
if proc.poll() is None:
proc.terminate()
_wait_until_dead(grace)
if proc.poll() is None:
proc.kill()
if not _wait_until_dead(grace):
logger.warning("DDGS worker pid=%s did not exit after kill", proc.pid)
except Exception as exc: # noqa: BLE001 — best-effort cleanup
logger.debug("DDGS worker reap error: %s", exc)
def _run_ddgs_search_bounded(query: str, safe_limit: int) -> list[dict[str, Any]]:
"""Run ``_run_ddgs_search`` in a disposable process with a hard deadline.
The parent never joins the child while it may be inside native code holding
*its* GIL — it only polls a communicator thread and, on timeout/interrupt,
terminates the child OS process. Raises ``TimeoutError``,
``_SearchInterrupted``, or ``RuntimeError``.
"""
# Imported lazily so plugin import stays light for ``hermes tools`` probes.
from tools.interrupt import is_interrupted
global _last_worker_proc
request: dict[str, Any] = {"query": query, "safe_limit": safe_limit}
if _test_hook:
request["test_hook"] = _test_hook
from tools.environments.local import _sanitize_subprocess_env
env = _sanitize_subprocess_env(dict(os.environ))
if _test_hook:
env["HERMES_DDGS_ALLOW_TEST_HOOKS"] = "1"
# Running the worker as a script puts ``plugins/web/ddgs/`` on ``sys.path[0]``,
# which breaks ``import plugins...``. Prepend the path entry that makes the
# live ``plugins`` package importable (source tree or site-packages).
child_pythonpath = env.get("PYTHONPATH", "")
path_entry = _plugins_path_entry()
if path_entry and path_entry not in child_pythonpath.split(os.pathsep):
env["PYTHONPATH"] = (
path_entry + os.pathsep + child_pythonpath if child_pythonpath else path_entry
)
worker_path = os.path.join(os.path.dirname(os.path.abspath(__file__)), "_search_worker.py")
# Platform-only spawn knobs — stdin/stdout/stderr must stay as explicit
# keyword args on the Popen call so scripts/check_subprocess_stdin.py can
# see them (TUI gateway inherits stdin; #14036).
extra_kwargs: dict[str, Any] = {}
if sys.platform == "win32":
# New process group so terminate/kill reach the worker cleanly on Windows.
extra_kwargs["creationflags"] = subprocess.CREATE_NEW_PROCESS_GROUP
else:
# Own session so a hung primp/libcurl grandchild can be reaped with the worker.
extra_kwargs["start_new_session"] = True
proc = subprocess.Popen(
[sys.executable, worker_path],
stdin=subprocess.PIPE,
stdout=subprocess.PIPE,
# DEVNULL avoids the classic deadlock where a chatty child fills the
# stderr pipe buffer while the parent only drains stdout.
stderr=subprocess.DEVNULL,
env=env,
text=True,
**extra_kwargs,
)
_last_worker_proc = proc
# ``communicate`` runs in a side thread so the parent can poll interrupt /
# deadline without blocking. Killing the child unblocks communicate.
pool = cf.ThreadPoolExecutor(max_workers=1)
fut = pool.submit(proc.communicate, json.dumps(request))
timed_out = False
interrupted = False
raw = ""
try:
deadline = time.monotonic() + _SEARCH_TIMEOUT_SECS
while True:
if is_interrupted():
interrupted = True
break
remaining = deadline - time.monotonic()
if remaining <= 0:
timed_out = True
break
try:
out, _err = fut.result(timeout=min(_POLL_INTERVAL_SECS, remaining))
raw = out or ""
break
except cf.TimeoutError:
continue
finally:
_terminate_and_reap(proc)
# After kill, communicate should return promptly; don't block forever.
if not fut.done():
try:
out, _err = fut.result(timeout=_TERMINATE_GRACE_SECS)
if not raw:
raw = out or ""
except Exception: # noqa: BLE001
pass
pool.shutdown(wait=False, cancel_futures=True)
if interrupted:
raise _SearchInterrupted("DuckDuckGo search interrupted")
if timed_out:
raise TimeoutError(
f"DuckDuckGo search timed out after {_SEARCH_TIMEOUT_SECS}s"
)
raw = raw.strip()
if not raw:
raise RuntimeError(
f"DDGS worker exited without a result (code={proc.poll()})"
)
try:
envelope = json.loads(raw)
except json.JSONDecodeError as exc:
raise RuntimeError(
f"DDGS worker returned invalid JSON: {raw[:200]!r}"
) from exc
if not isinstance(envelope, dict):
raise RuntimeError(f"DDGS worker returned an invalid envelope: {envelope!r}")
if envelope.get("ok"):
results = envelope.get("results") or []
if not isinstance(results, list):
raise RuntimeError("DDGS worker returned non-list results")
return results
raise RuntimeError(str(envelope.get("error") or "DDGS worker failed"))
class DDGSWebSearchProvider(WebSearchProvider):
"""DuckDuckGo HTML-scrape search provider.
No API key needed. Rate limits are enforced server-side by DuckDuckGo;
the provider surfaces ``DuckDuckGoSearchException`` and other ddgs errors
as ``{"success": False, "error": ...}`` rather than raising.
"""
@property
def name(self) -> str:
return "ddgs"
@property
def display_name(self) -> str:
return "DuckDuckGo (ddgs)"
def is_available(self) -> bool:
"""Return True when the ``ddgs`` package is importable.
Probes the import once; cheap because Python caches the import. Must
NOT perform network I/O — runs at tool-registration time and on every
``hermes tools`` paint.
"""
try:
import ddgs # noqa: F401
return True
except ImportError:
return False
def supports_search(self) -> bool:
return True
def supports_extract(self) -> bool:
return False
def search(self, query: str, limit: int = 5) -> Dict[str, Any]:
"""Execute a DuckDuckGo search and return normalized results.
The synchronous ``ddgs`` call runs in a disposable child process with
a hard wall-clock timeout (``_SEARCH_TIMEOUT_SECS``) so a hung native
``primp`` call cannot freeze the Hermes process (#36776, #68096).
"""
try:
import ddgs # type: ignore # noqa: F401 — availability probe
except ImportError:
return {
"success": False,
"error": "ddgs package is not installed — run `pip install ddgs`",
}
# DDGS().text yields at most `max_results` items; we cap defensively
# in case the package ignores the hint.
safe_limit = max(1, int(limit))
try:
web_results = _run_ddgs_search_bounded(query, safe_limit)
except TimeoutError:
logger.warning(
"DDGS search timed out after %ds for query: %r",
_SEARCH_TIMEOUT_SECS,
query,
)
return {
"success": False,
"error": (
f"DuckDuckGo search timed out after {_SEARCH_TIMEOUT_SECS}s — "
"DuckDuckGo may be rate-limiting or slow. Try again later "
"or switch to a different search provider."
),
}
except _SearchInterrupted:
logger.info("DDGS search interrupted for query: %r", query)
return {
"success": False,
"error": "DuckDuckGo search interrupted",
}
except Exception as exc: # noqa: BLE001 — ddgs raises its own exceptions
logger.warning("DDGS search error: %s", exc)
return {"success": False, "error": f"DuckDuckGo search failed: {exc}"}
logger.info(
"DDGS search '%s': %d results (limit %d)", query, len(web_results), limit
)
return {"success": True, "data": {"web": web_results}}
def get_setup_schema(self) -> Dict[str, Any]:
return {
"name": "DuckDuckGo (ddgs)",
"badge": "free · no key · search only",
"tag": "Search via the ddgs Python package — no API key (pair with any extract provider)",
"env_vars": [],
# Trigger `_run_post_setup("ddgs")` after the user picks this row
# so the ddgs Python package gets pip-installed on first selection.
"post_setup": "ddgs",
}