Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
100 changes: 100 additions & 0 deletions examples/deepeyes_v2_agentic/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -105,6 +105,106 @@ and propagates it to every Ray worker via `--runtime-env-json`. Set
`APPTAINER_IMAGE_PATH` explicitly to override (e.g. shared NFS path for
multi-node).

## Text search

Text search defaults to a deterministic, offline mock. No search service, network
or API key is needed; model inference and the Python sandbox retain their normal
requirements. Select a real backend with one YAML file:

```bash
export DEEPEYES_V2_SEARCH_CONFIG=/path/to/search.yaml
```

For a Search-R1-compatible service:

```yaml
backend: retriever
retriever:
url: http://your-retriever:8000/retrieve
```

The client sends `POST {"queries": [query], "topk": k, "return_scores": true}`
and accepts both raw corpus documents and `document`/`score` wrappers inside
`result[0]`. It preserves `contents` as the snippet. Missing titles use the first
line of the passage; missing links remain empty rather than inventing URLs.

For [Tavily](https://docs.tavily.com/documentation/api-reference/endpoint/search),
set `DEEPEYES_V2_SEARCH_API_KEY` in your environment and use:

```yaml
backend: external
external:
endpoint: https://api.tavily.com/search
method: POST
auth:
header: Authorization
prefix: "Bearer "
request_map:
query: query
size: max_results
response_map:
results: results
title: title
link: url
snippet: content
date: published_date
```

`external` supports GET query parameters or POST JSON, configurable header
authentication, and dot-separated response paths. Omit `auth` for an unauthenticated
endpoint and omit the date mapping when the API does not provide one. Request
mapping names must be distinct; no provider SDK is required.

Common settings (all optional):

| Setting | Default | Meaning |
| ------------- | ------- | --------------------------------------------------------------- |
| `backend` | `mock` | `mock`, `retriever`, or `external` |
| `top_k` | `5` | Maximum results; an explicit `search(query, size)` overrides it |
| `timeout_s` | `10` | HTTP connect/read timeout, not a total wall-clock deadline |
| `max_retries` | `2` | Retries after the first request; zero means one attempt |
| `backoff_s` | `0.5` | Initial retry delay; doubles up to 2 seconds |
| `trust_env` | `false` | Enable environment proxy/netrc settings when needed |

Successful calls return `elapsed_time` in seconds and `data`, a list of
`{title, link, snippet, date}` records. Date is a string or null. Empty results
are successful; malformed responses are errors. Connection failures, timeouts,
HTTP 408/429 and 5xx statuses are retried. Other HTTP failures (including
redirects), malformed JSON and configuration errors fail immediately. Failure
returns the existing `"Error"` sentinel, which becomes a non-terminal
`search_failed` observation. A broken real backend never falls back to mock.
Diagnostic logs identify the error without printing credentials or response bodies.

Both training launchers forward the configuration path and API key to Ray workers.
The YAML must be readable at that path on every node; forwarding a path does not
upload its file. Credentials are excluded from shell tracing, but Ray runtime
configuration and submission process arguments may be visible to cluster
administrators. Use worker-provisioned credentials if that exposure is unsuitable.

CPU tests use a local HTTP server for the Search-R1 and external protocols, plus a
scripted model to exercise the real Agent loop. They need no search service, GPU
or Apptainer instance:

```bash
python -m pytest tests/examples/deepeyes_v2_agentic -q
```

To check your configured service separately, from the repository root:

```bash
PYTHONPATH=examples/deepeyes_v2_agentic:. python - <<'PYTHON'
from app.search_utils import search

result = search("What is reinforcement learning?", size=3)
assert result != "Error", "See the search diagnostic above"
print(result)
PYTHON
```

The protocol tests do not validate a live E5/FAISS deployment or API credentials.
The repository's native Search-R1 server uses CUDA; the search client itself can
run on a CPU host, including macOS.

## Image-search cache (optional, only for the `search` split)

The `<tool_call>image_search</tool_call>` branch hits a precomputed
Expand Down
7 changes: 3 additions & 4 deletions examples/deepeyes_v2_agentic/app/env_deepeyes_v2.py
Original file line number Diff line number Diff line change
Expand Up @@ -444,7 +444,7 @@ def _dispatch_search(self, tool_name: str, tool_args: Any) -> dict:
}

# tool_name == "search"
query = tool_args["query"] if isinstance(tool_args, dict) and "query" in tool_args else str(tool_args)
query = tool_args.get("query") if isinstance(tool_args, dict) else tool_args
result = search(query)
if result == "Error":
return {"status": "error", "result": "Error", "images": []}
Expand All @@ -458,9 +458,8 @@ def _dispatch_search(self, tool_name: str, tool_args: Any) -> dict:
if page.get("snippet") is not None:
snippet = "\n" + page["snippet"]
snippets.append(f"{idx + 1}. [{page['title']}]({page['link']}){date_published}{snippet}")
content = (
f"A Google search for '{query}' found {len(snippets)} results:"
f"\n\n## Web Results\n" + "\n\n".join(snippets)
content = f"A Web search for '{query}' found {len(snippets)} results:\n\n## Web Results\n" + "\n\n".join(
snippets
)
except (KeyError, TypeError) as exc:
return {
Expand Down
239 changes: 239 additions & 0 deletions examples/deepeyes_v2_agentic/app/search_backends.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,239 @@
# Copyright (c) 2026 Relax Authors. All Rights Reserved.

"""Text search adapters configured by DEEPEYES_V2_SEARCH_CONFIG."""

from __future__ import annotations

import math
import os
import time
from pathlib import Path
from typing import Any
from urllib.parse import urlsplit


DEFAULTS = {
"backend": "mock",
"top_k": 5,
"timeout_s": 10.0,
"max_retries": 2,
"backoff_s": 0.5,
"trust_env": False,
}


class SearchError(ValueError):
"""Search failure whose message is safe to log."""


def _mapping(value: Any, name: str, allowed: set[str]) -> dict:
if not isinstance(value, dict):
raise SearchError(f"{name} must be a mapping")
if value.keys() - allowed:
raise SearchError(f"{name} contains unknown fields")
return value


def _text(value: Any, name: str, *, empty: bool = False) -> str:
if not isinstance(value, str) or (not empty and not value.strip()):
raise SearchError(f"{name} must be a {'non-empty ' if not empty else ''}string")
return value


def _integer(value: Any, name: str, minimum: int) -> int:
if type(value) is not int or value < minimum:
raise SearchError(f"{name} must be an integer >= {minimum}")
return value


def _load_config() -> dict:
config = dict(DEFAULTS)
path = os.environ.get("DEEPEYES_V2_SEARCH_CONFIG", "").strip()
if path:
import yaml

try:
source = yaml.safe_load(Path(path).expanduser().read_text(encoding="utf-8"))
except (OSError, UnicodeError, yaml.YAMLError) as exc:
raise SearchError(f"cannot read DEEPEYES_V2_SEARCH_CONFIG ({type(exc).__name__})") from None
config.update(_mapping(source, "search config", set(DEFAULTS) | {"retriever", "external"}))
if config["backend"] not in ("mock", "retriever", "external"):
raise SearchError("backend must be mock, retriever or external")
_integer(config["top_k"], "top_k", 1)
_integer(config["max_retries"], "max_retries", 0)
for name in ("timeout_s", "backoff_s"):
value = config[name]
if type(value) not in (int, float) or not math.isfinite(value) or value < 0:
raise SearchError(f"{name} must be a finite non-negative number")
if config["timeout_s"] == 0:
raise SearchError("timeout_s must be positive")
if not isinstance(config["trust_env"], bool):
raise SearchError("trust_env must be boolean")
return config


def _endpoint(value: Any, name: str) -> str:
value = _text(value, name).strip()
try:
parsed = urlsplit(value)
valid = parsed.scheme in ("http", "https") and parsed.hostname and not parsed.username and not parsed.password
parsed.port
except ValueError:
valid = False
if not valid:
raise SearchError(f"{name} must be an HTTP(S) URL without credentials")
return value


def _request(config: dict, url: str, method: str, payload: dict, headers: dict) -> Any:
import requests

delay = min(config["backoff_s"], 2.0)
attempts = config["max_retries"] + 1
with requests.Session() as session:
session.trust_env = config["trust_env"]
for attempt in range(attempts):
try:
with session.request(
method,
url,
headers=headers,
**{"params" if method == "GET" else "json": payload},
timeout=config["timeout_s"],
allow_redirects=False,
) as response:
status = response.status_code
if 200 <= status < 300:
try:
return response.json()
except ValueError:
raise SearchError("search response is not valid JSON") from None
if status not in (408, 429) and not 500 <= status < 600:
raise SearchError(f"search HTTP status {status}")
reason = f"HTTP {status}"
except (requests.Timeout, requests.ConnectionError) as exc:
reason = type(exc).__name__
except requests.RequestException as exc:
raise SearchError(f"search request failed ({type(exc).__name__})") from None
if attempt + 1 < attempts:
time.sleep(delay)
delay = min(delay * 2, 2.0)
raise SearchError(f"search failed after {attempts} attempts ({reason})")


def _date(value: Any) -> str | None:
return None if value is None else _text(value, "result.date", empty=True)


def _retriever_rows(body: Any, size: int) -> list[dict]:
batches = body.get("result") if isinstance(body, dict) else None
if not isinstance(batches, list) or len(batches) != 1 or not isinstance(batches[0], list):
raise SearchError("retriever response must contain one result batch")
rows = []
for item in batches[0][:size]:
doc = item.get("document", item) if isinstance(item, dict) else None
if not isinstance(doc, dict):
raise SearchError("retriever document must be a mapping")
content = _text(doc.get("contents"), "retriever document.contents")
title = doc.get("title")
if title is not None:
title = _text(title, "result.title", empty=True)
title = title or content.strip().splitlines()[0].strip().strip('"')
link = doc.get("url")
if link is None:
link = doc.get("link", "")
rows.append(
{
"title": title,
"link": _text("" if link is None else link, "result.link", empty=True),
"snippet": content,
"date": _date(doc.get("date")),
}
)
return rows


def _field(row: Any, path: str) -> Any:
for key in path.split("."):
if not isinstance(row, dict):
return None
row = row.get(key)
return row


def _external(config: dict, query: str, size: int) -> list[dict]:
external = _mapping(
config.get("external"), "external", {"endpoint", "method", "auth", "request_map", "response_map"}
)
url = _endpoint(external.get("endpoint"), "external.endpoint")
method = _text(external.get("method", "POST"), "external.method").upper()
if method not in ("GET", "POST"):
raise SearchError("external.method must be GET or POST")
request_map = _mapping(external.get("request_map"), "external.request_map", {"query", "size"})
response_map = _mapping(
external.get("response_map"), "external.response_map", {"results", "title", "link", "snippet", "date"}
)
for name in ("query", "size"):
_text(request_map.get(name), f"external.request_map.{name}")
if request_map["query"] == request_map["size"]:
raise SearchError("external.request_map fields must be distinct")
for name in ("results", "title", "link", "snippet"):
_text(response_map.get(name), f"external.response_map.{name}")
if "date" in response_map:
_text(response_map["date"], "external.response_map.date")
headers = {"Accept": "application/json"}
if "auth" in external:
auth = _mapping(external["auth"], "external.auth", {"header", "prefix"})
header = _text(auth.get("header"), "external.auth.header")
prefix = _text(auth.get("prefix", ""), "external.auth.prefix", empty=True)
secret = os.environ.get("DEEPEYES_V2_SEARCH_API_KEY", "")
if not secret.strip() or any(c in header + prefix + secret for c in "\r\n"):
raise SearchError("external authentication is missing or invalid")
headers[header] = prefix + secret
body = _request(config, url, method, {request_map["query"]: query, request_map["size"]: size}, headers)
items = _field(body, response_map["results"])
if not isinstance(items, list):
raise SearchError("external results must be a list")
rows = []
for item in items[:size]:
row = {
name: _text(_field(item, response_map[name]), f"result.{name}", empty=True)
for name in ("title", "link", "snippet")
}
row["date"] = _date(_field(item, response_map["date"])) if "date" in response_map else None
rows.append(row)
return rows


def run_search(query: str, size: int | None = None) -> dict:
"""Return normalized search results, or raise SearchError on failure."""
query = _text(query, "query").strip()
config = _load_config()
size = config["top_k"] if size is None else _integer(size, "size", 1)
if config["backend"] == "mock":
return {
"elapsed_time": 0.0,
"data": [
{
"title": f"Offline mock result {i + 1}",
"link": f"https://example.invalid/{i + 1}",
"snippet": f"Deterministic offline result for: {query}",
"date": None,
}
for i in range(size)
],
}
started = time.monotonic()
if config["backend"] == "retriever":
retriever = _mapping(config.get("retriever"), "retriever", {"url"})
body = _request(
config,
_endpoint(retriever.get("url"), "retriever.url"),
"POST",
{"queries": [query], "topk": size, "return_scores": True},
{"Accept": "application/json"},
)
rows = _retriever_rows(body, size)
else:
rows = _external(config, query, size)
return {"elapsed_time": time.monotonic() - started, "data": rows}
Loading
Loading