Edit on GitHub

agent_search_gateway.providers.web.decodo

Decodo Google Search template and Markdown scrape adapter.

 1"""Decodo Google Search template and Markdown scrape adapter."""
 2
 3from ...errors import ExecutionFailure
 4from ...observability import SecretValue
 5from ...providers.contracts import KeywordSearchHit, URLFetchCandidate
 6from ...url_normalization import NormalizedURL
 7from .common import (
 8    HttpRequester,
 9    configured_string,
10    endpoint,
11    failure,
12    non_empty_string,
13    optional_string,
14    require_list,
15    require_object,
16)
17
18
19class DecodoAdapter:
20    def __init__(
21        self,
22        *,
23        name: str,
24        api_url: str,
25        secret: SecretValue,
26        http_executor: HttpRequester,
27    ) -> None:
28        self.name = name
29        self._api_url = configured_string(api_url, "api_url").rstrip("/")
30        self._secret = secret
31        self._http = http_executor
32
33    @property
34    def _headers(self) -> dict[str, str]:
35        return {"Authorization": f"Basic {self._secret.reveal()}"}
36
37    async def search(self, query: str) -> list[KeywordSearchHit]:
38        payload = await self._http.request_json(
39            "POST",
40            endpoint(self._api_url, "/v2/scrape"),
41            stage="search",
42            headers=self._headers,
43            json_body={"target": "google_search", "query": query, "parse": True},
44        )
45        results = self._organic_results(payload)
46        hits: list[KeywordSearchHit] = []
47        for item in results:
48            try:
49                result = require_object(item, self.name, "search", "result")
50                hits.append(
51                    KeywordSearchHit(
52                        url=non_empty_string(result.get("url"), self.name, "search", "result.url"),
53                        title=optional_string(
54                            result.get("title"), self.name, "search", "result.title"
55                        ),
56                        snippet=optional_string(
57                            result.get("desc"), self.name, "search", "result.desc"
58                        ),
59                    )
60                )
61            except ExecutionFailure:
62                continue
63        return hits
64
65    async def fetch(self, url: NormalizedURL) -> URLFetchCandidate:
66        text = await self._http.request_text(
67            "POST",
68            endpoint(self._api_url, "/v2/scrape"),
69            stage="fetch",
70            headers=self._headers,
71            json_body={"url": str(url), "markdown": True},
72        )
73        if not text.strip():
74            raise failure(self.name, "fetch", "page body is empty")
75        return URLFetchCandidate(text, text)
76
77    def _organic_results(self, payload: object) -> list[object]:
78        root = require_object(payload, self.name, "search", "response")
79        pages = require_list(root.get("results"), self.name, "search", "results")
80        if not pages:
81            raise failure(self.name, "search", "results must contain a parsed page")
82        page = require_object(pages[0], self.name, "search", "result page")
83        if page.get("status_code") != 200:
84            raise failure(self.name, "search", "provider reported failure")
85        content = require_object(page.get("content"), self.name, "search", "content")
86        errors = require_list(content.get("errors", []), self.name, "search", "content.errors")
87        if errors:
88            raise failure(self.name, "search", "provider reported failure")
89        parsed_wrapper = require_object(
90            content.get("results"), self.name, "search", "content.results"
91        )
92        parsed = require_object(
93            parsed_wrapper.get("results"), self.name, "search", "content.results.results"
94        )
95        return require_list(parsed.get("organic"), self.name, "search", "organic")
class DecodoAdapter:
20class DecodoAdapter:
21    def __init__(
22        self,
23        *,
24        name: str,
25        api_url: str,
26        secret: SecretValue,
27        http_executor: HttpRequester,
28    ) -> None:
29        self.name = name
30        self._api_url = configured_string(api_url, "api_url").rstrip("/")
31        self._secret = secret
32        self._http = http_executor
33
34    @property
35    def _headers(self) -> dict[str, str]:
36        return {"Authorization": f"Basic {self._secret.reveal()}"}
37
38    async def search(self, query: str) -> list[KeywordSearchHit]:
39        payload = await self._http.request_json(
40            "POST",
41            endpoint(self._api_url, "/v2/scrape"),
42            stage="search",
43            headers=self._headers,
44            json_body={"target": "google_search", "query": query, "parse": True},
45        )
46        results = self._organic_results(payload)
47        hits: list[KeywordSearchHit] = []
48        for item in results:
49            try:
50                result = require_object(item, self.name, "search", "result")
51                hits.append(
52                    KeywordSearchHit(
53                        url=non_empty_string(result.get("url"), self.name, "search", "result.url"),
54                        title=optional_string(
55                            result.get("title"), self.name, "search", "result.title"
56                        ),
57                        snippet=optional_string(
58                            result.get("desc"), self.name, "search", "result.desc"
59                        ),
60                    )
61                )
62            except ExecutionFailure:
63                continue
64        return hits
65
66    async def fetch(self, url: NormalizedURL) -> URLFetchCandidate:
67        text = await self._http.request_text(
68            "POST",
69            endpoint(self._api_url, "/v2/scrape"),
70            stage="fetch",
71            headers=self._headers,
72            json_body={"url": str(url), "markdown": True},
73        )
74        if not text.strip():
75            raise failure(self.name, "fetch", "page body is empty")
76        return URLFetchCandidate(text, text)
77
78    def _organic_results(self, payload: object) -> list[object]:
79        root = require_object(payload, self.name, "search", "response")
80        pages = require_list(root.get("results"), self.name, "search", "results")
81        if not pages:
82            raise failure(self.name, "search", "results must contain a parsed page")
83        page = require_object(pages[0], self.name, "search", "result page")
84        if page.get("status_code") != 200:
85            raise failure(self.name, "search", "provider reported failure")
86        content = require_object(page.get("content"), self.name, "search", "content")
87        errors = require_list(content.get("errors", []), self.name, "search", "content.errors")
88        if errors:
89            raise failure(self.name, "search", "provider reported failure")
90        parsed_wrapper = require_object(
91            content.get("results"), self.name, "search", "content.results"
92        )
93        parsed = require_object(
94            parsed_wrapper.get("results"), self.name, "search", "content.results.results"
95        )
96        return require_list(parsed.get("organic"), self.name, "search", "organic")
DecodoAdapter( *, name: str, api_url: str, secret: agent_search_gateway.observability.SecretValue, http_executor: agent_search_gateway.providers.web.common.HttpRequester)
21    def __init__(
22        self,
23        *,
24        name: str,
25        api_url: str,
26        secret: SecretValue,
27        http_executor: HttpRequester,
28    ) -> None:
29        self.name = name
30        self._api_url = configured_string(api_url, "api_url").rstrip("/")
31        self._secret = secret
32        self._http = http_executor
name
async def search( self, query: str) -> list[agent_search_gateway.providers.contracts.KeywordSearchHit]:
38    async def search(self, query: str) -> list[KeywordSearchHit]:
39        payload = await self._http.request_json(
40            "POST",
41            endpoint(self._api_url, "/v2/scrape"),
42            stage="search",
43            headers=self._headers,
44            json_body={"target": "google_search", "query": query, "parse": True},
45        )
46        results = self._organic_results(payload)
47        hits: list[KeywordSearchHit] = []
48        for item in results:
49            try:
50                result = require_object(item, self.name, "search", "result")
51                hits.append(
52                    KeywordSearchHit(
53                        url=non_empty_string(result.get("url"), self.name, "search", "result.url"),
54                        title=optional_string(
55                            result.get("title"), self.name, "search", "result.title"
56                        ),
57                        snippet=optional_string(
58                            result.get("desc"), self.name, "search", "result.desc"
59                        ),
60                    )
61                )
62            except ExecutionFailure:
63                continue
64        return hits
66    async def fetch(self, url: NormalizedURL) -> URLFetchCandidate:
67        text = await self._http.request_text(
68            "POST",
69            endpoint(self._api_url, "/v2/scrape"),
70            stage="fetch",
71            headers=self._headers,
72            json_body={"url": str(url), "markdown": True},
73        )
74        if not text.strip():
75            raise failure(self.name, "fetch", "page body is empty")
76        return URLFetchCandidate(text, text)