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)
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
async def
fetch( self, url: agent_search_gateway.url_normalization.NormalizedURL) -> agent_search_gateway.providers.contracts.URLFetchCandidate:
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)