diff --git a/plugins/web/ddgs/_search_worker.py b/plugins/web/ddgs/_search_worker.py new file mode 100644 index 000000000000..0521284c53c1 --- /dev/null +++ b/plugins/web/ddgs/_search_worker.py @@ -0,0 +1,113 @@ +"""DDGS search child-process entrypoint (#68096). + +Invoked as ``python plugins/web/ddgs/_search_worker.py`` (script path from the +parent provider). Reads one JSON request from stdin, writes one JSON envelope +to stdout, then exits. + +Request:: + {"query": str, "safe_limit": int} + +Envelope:: + {"ok": true, "results": [...]} + {"ok": false, "error": str} + +Optional test hooks (only when ``HERMES_DDGS_ALLOW_TEST_HOOKS=1``):: + {"query": ..., "safe_limit": ..., "test_hook": "sleep"|"gil"|"success"|"error"|"empty"} +""" + +from __future__ import annotations + +import json +import os +import sys +import time + + +def _hold_gil(secs: int) -> None: + """Block in a foreign call that keeps the GIL (ctypes.PyDLL). + + Mirrors native ``primp`` holding the interpreter lock. ``PyDLL`` (unlike + ``CDLL``/``WinDLL``) does not release the GIL around the call. + """ + import ctypes + + if sys.platform == "win32": + lib = ctypes.PyDLL("kernel32") + lib.Sleep.argtypes = [ctypes.c_uint] + lib.Sleep(int(secs * 1000)) + return + + lib = ctypes.PyDLL(None) + try: + sleep = lib.sleep + except AttributeError: # pragma: no cover — macOS libSystem fallback + sleep = ctypes.PyDLL("/usr/lib/libSystem.B.dylib").sleep + sleep.argtypes = [ctypes.c_uint] + sleep(int(secs)) + + +def _run_test_hook(hook: str) -> dict: + if hook == "sleep": + time.sleep(30) + return {"ok": False, "error": "sleep hook returned unexpectedly"} + if hook == "gil": + _hold_gil(30) + return {"ok": False, "error": "gil hook returned unexpectedly"} + if hook == "success": + return { + "ok": True, + "results": [ + { + "title": "Hit", + "url": "https://example.com", + "description": "body", + "position": 1, + } + ], + } + if hook == "empty": + return {"ok": True, "results": []} + if hook == "error": + return {"ok": False, "error": "RuntimeError: boom"} + return {"ok": False, "error": f"unknown test_hook: {hook!r}"} + + +def _write_envelope(envelope: dict) -> None: + json.dump(envelope, sys.stdout) + sys.stdout.flush() + + +def main() -> int: + try: + request = json.load(sys.stdin) + except Exception as exc: # noqa: BLE001 + _write_envelope({"ok": False, "error": f"invalid request: {exc}"}) + return 2 + + hook = request.get("test_hook") + if hook: + if os.environ.get("HERMES_DDGS_ALLOW_TEST_HOOKS") != "1": + _write_envelope( + {"ok": False, "error": "test_hook refused (hooks not enabled)"} + ) + return 3 + envelope = _run_test_hook(str(hook)) + _write_envelope(envelope) + return 0 if envelope.get("ok") else 1 + + query = str(request.get("query") or "") + safe_limit = max(1, int(request.get("safe_limit") or 1)) + try: + # Import inside main so script startup stays light / patchable. + from plugins.web.ddgs.provider import _run_ddgs_search + + results = _run_ddgs_search(query, safe_limit) + _write_envelope({"ok": True, "results": results}) + return 0 + except Exception as exc: # noqa: BLE001 + _write_envelope({"ok": False, "error": f"{type(exc).__name__}: {exc}"}) + return 1 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/plugins/web/ddgs/provider.py b/plugins/web/ddgs/provider.py index dcdeb0897ac9..28e4762bec7e 100644 --- a/plugins/web/ddgs/provider.py +++ b/plugins/web/ddgs/provider.py @@ -8,13 +8,24 @@ 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 concurrent.futures as cf +import json import logging -from typing import Any, Dict +import os +import subprocess +import sys +import time +from typing import Any, Dict, Optional from agent.web_search_provider import WebSearchProvider @@ -23,18 +34,28 @@ 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 and block every platform -# (#36776). Enforce a hard cap here via a worker thread. +# 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 tests can patch it directly without - spawning a real multi-second worker thread. ``DDGS(timeout=...)`` bounds + 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 caller via a future timeout. + the parent via process timeout (#68096). """ from ddgs import DDGS # type: ignore @@ -55,6 +76,190 @@ def _run_ddgs_search(query: str, safe_limit: int) -> list[dict[str, Any]]: 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 + + env = os.environ.copy() + 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. @@ -94,9 +299,9 @@ class DDGSWebSearchProvider(WebSearchProvider): def search(self, query: str, limit: int = 5) -> Dict[str, Any]: """Execute a DuckDuckGo search and return normalized results. - The synchronous ``ddgs`` call is run in a worker thread with a hard - wall-clock timeout (``_SEARCH_TIMEOUT_SECS``) so a hung search cannot - block the shared agent loop indefinitely (#36776). + 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 @@ -110,40 +315,35 @@ class DDGSWebSearchProvider(WebSearchProvider): # in case the package ignores the hint. safe_limit = max(1, int(limit)) - # A fresh single-worker pool per call (rather than a module-level one) - # is intentional: on timeout the blocking ddgs call cannot be cancelled - # and keeps running, so a shared pool would serialise every later search - # behind that hung worker. A per-call pool isolates each search from a - # previously-hung one. - pool = _cf.ThreadPoolExecutor(max_workers=1) try: - future = pool.submit(_run_ddgs_search, query, safe_limit) - try: - web_results = future.result(timeout=_SEARCH_TIMEOUT_SECS) - except _cf.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." - ), - } + 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}"} - finally: - # Return immediately without joining the worker. On timeout the - # already-running ddgs call can't be cancelled (cancel_futures only - # affects not-yet-started work), so the worker runs to completion - # on its own; it writes nothing shared, so leaking it is safe. - pool.shutdown(wait=False, cancel_futures=True) - logger.info("DDGS search '%s': %d results (limit %d)", query, len(web_results), limit) + 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]: