agent_search_gateway.runtime
Resolved runtime assembly for the foreground daemon.
1"""Resolved runtime assembly for the foreground daemon.""" 2 3import logging 4from collections.abc import Callable, Mapping 5from typing import cast 6 7import httpx 8 9from .academic.aggregator import PaperAggregator 10from .concurrency import ProviderQuotaManager 11from .config import ( 12 ResolvedAcademicProviderConfig, 13 ResolvedConfig, 14 ResolvedOAResolverConfig, 15 ResolvedWebProviderConfig, 16) 17from .errors import ConfigFailure, ErrorCode 18from .llm.stages import LLMStages 19from .observability import log_event 20from .orchestrators.fetch import FetchOrchestrator 21from .orchestrators.paper import PaperSearchOrchestrator 22from .orchestrators.search import SearchOrchestrator 23from .paths import RuntimePaths 24from .providers.academic.defaults import ( 25 build_default_academic_registry, 26 build_default_oa_resolver_registry, 27) 28from .providers.academic.registry import AcademicProviderRegistry, OAResolverRegistry 29from .providers.contracts import ( 30 AcademicSearchProvider, 31 KeywordSearchProvider, 32 LLMClient, 33 OAResolver, 34 URLFetchProvider, 35) 36from .providers.defaults import build_default_registry 37from .providers.http import HttpJsonExecutor 38from .providers.openai_chat import OpenAIChatCompletionsClient 39from .providers.registry import ProviderRegistry 40from .result_writer import ResultWriter 41from .scheduler.fetch import FetchScheduler 42from .url_store import URLStore 43 44HttpClientFactory = Callable[[], httpx.AsyncClient] 45_RESERVED_WEB_ADAPTER_KWARGS = frozenset({"name", "http_executor", "secret"}) 46_RESERVED_ACADEMIC_ADAPTER_KWARGS = frozenset({"name", "executor", "api_key", "contact_email"}) 47 48 49class Runtime: 50 """Owns one in-memory gateway runtime and its transport clients.""" 51 52 def __init__( 53 self, 54 *, 55 quotas: ProviderQuotaManager, 56 web_search_providers: tuple[KeywordSearchProvider, ...], 57 web_fetch_providers: tuple[URLFetchProvider, ...], 58 academic_search_providers: tuple[AcademicSearchProvider, ...], 59 oa_resolver: OAResolver | None, 60 llm_clients: Mapping[str, LLMClient], 61 store: URLStore, 62 search_orchestrator: SearchOrchestrator, 63 paper_search_orchestrator: PaperSearchOrchestrator, 64 fetch_orchestrator: FetchOrchestrator, 65 web_http_executors: tuple[HttpJsonExecutor, ...], 66 academic_http_executors: tuple[HttpJsonExecutor, ...], 67 ) -> None: 68 self.quotas = quotas 69 self.web_search_providers = web_search_providers 70 self.web_fetch_providers = web_fetch_providers 71 self.academic_search_providers = academic_search_providers 72 self.oa_resolver = oa_resolver 73 self.llm_clients = dict(llm_clients) 74 self.store = store 75 self.search_orchestrator = search_orchestrator 76 self.paper_search_orchestrator = paper_search_orchestrator 77 self.fetch_orchestrator = fetch_orchestrator 78 self._web_http_executors = web_http_executors 79 self._academic_http_executors = academic_http_executors 80 self._closed = False 81 82 @classmethod 83 def build( 84 cls, 85 config: ResolvedConfig, 86 paths: RuntimePaths, 87 *, 88 registry: ProviderRegistry | None = None, 89 academic_registry: AcademicProviderRegistry | None = None, 90 oa_resolver_registry: OAResolverRegistry | None = None, 91 http_client_factory: HttpClientFactory = httpx.AsyncClient, 92 ) -> "Runtime": 93 provider_registry = registry or build_default_registry() 94 academic_provider_registry = academic_registry or build_default_academic_registry() 95 resolver_registry = oa_resolver_registry or build_default_oa_resolver_registry() 96 enabled_web = tuple( 97 item for item in config.web.providers if item.enable_search or item.enable_fetch 98 ) 99 enabled_academic = tuple(item for item in config.academic.providers if item.enabled) 100 quotas = ProviderQuotaManager( 101 web_limits={item.name: item.max_concurrency for item in enabled_web}, 102 llm_limits={item.name: item.max_concurrency for item in config.llm.providers}, 103 academic_limits={item.name: item.max_concurrency for item in enabled_academic}, 104 ) 105 web_search, web_fetch, web_executors = cls._build_web_providers( 106 enabled_web, 107 provider_registry, 108 config, 109 http_client_factory, 110 ) 111 academic_search, academic_executors = cls._build_academic_providers( 112 enabled_academic, 113 academic_provider_registry, 114 config, 115 http_client_factory, 116 ) 117 oa_resolver, resolver_executors = cls._build_oa_resolver( 118 config.oa_resolver, 119 resolver_registry, 120 config, 121 http_client_factory, 122 ) 123 llm_clients = cls._build_llm_clients(config, quotas, http_client_factory) 124 stages = LLMStages( 125 llm_clients, 126 judge=config.llm.judge, 127 safety=config.llm.safety, 128 content_clean=config.llm.content_clean, 129 focus_summary=config.llm.focus_summary, 130 ) 131 store = URLStore() 132 result_writer = ResultWriter(paths.results_dir) 133 paper_aggregator = PaperAggregator( 134 ( 135 *(provider.name for provider in academic_search), 136 *(f"llm:{item.provider}" for item in config.llm.search_invocations), 137 ) 138 ) 139 search_orchestrator = SearchOrchestrator( 140 keyword_providers=web_search, 141 llm_invocations=config.llm.search_invocations, 142 quotas=quotas, 143 stages=stages, 144 store=store, 145 result_writer=result_writer, 146 paper_aggregator=paper_aggregator, 147 paper_resolver=oa_resolver, 148 ) 149 paper_search_orchestrator = PaperSearchOrchestrator( 150 providers=academic_search, 151 quotas=quotas, 152 aggregator=paper_aggregator, 153 resolver=oa_resolver, 154 store=store, 155 result_writer=result_writer, 156 ) 157 fetch_orchestrator = FetchOrchestrator( 158 store=store, 159 scheduler=FetchScheduler(web_fetch, quotas, stages), 160 stages=stages, 161 ) 162 runtime = cls( 163 quotas=quotas, 164 web_search_providers=web_search, 165 web_fetch_providers=web_fetch, 166 academic_search_providers=academic_search, 167 oa_resolver=oa_resolver, 168 llm_clients=llm_clients, 169 store=store, 170 search_orchestrator=search_orchestrator, 171 paper_search_orchestrator=paper_search_orchestrator, 172 fetch_orchestrator=fetch_orchestrator, 173 web_http_executors=web_executors, 174 academic_http_executors=(*academic_executors, *resolver_executors), 175 ) 176 log_event( 177 logging.getLogger(__name__), 178 logging.DEBUG, 179 "runtime_built", 180 web_providers=",".join(item.name for item in enabled_web) or "-", 181 web_provider_count=len(enabled_web), 182 llm_providers=",".join(item.name for item in config.llm.providers) or "-", 183 llm_provider_count=len(config.llm.providers), 184 academic_providers=",".join(item.name for item in enabled_academic) or "-", 185 academic_provider_count=len(enabled_academic), 186 oa_resolver=config.oa_resolver.name if config.oa_resolver is not None else "-", 187 web_limits=",".join(f"{item.name}:{item.max_concurrency}" for item in enabled_web) 188 or "-", 189 llm_limits=",".join( 190 f"{item.name}:{item.max_concurrency}" for item in config.llm.providers 191 ) 192 or "-", 193 academic_limits=",".join( 194 f"{item.name}:{item.max_concurrency}" for item in enabled_academic 195 ) 196 or "-", 197 ) 198 return runtime 199 200 @classmethod 201 def _build_web_providers( 202 cls, 203 providers: tuple[ResolvedWebProviderConfig, ...], 204 registry: ProviderRegistry, 205 config: ResolvedConfig, 206 http_client_factory: HttpClientFactory, 207 ) -> tuple[ 208 tuple[KeywordSearchProvider, ...], 209 tuple[URLFetchProvider, ...], 210 tuple[HttpJsonExecutor, ...], 211 ]: 212 search: list[KeywordSearchProvider] = [] 213 fetch: list[URLFetchProvider] = [] 214 executors: list[HttpJsonExecutor] = [] 215 for provider_config in providers: 216 registration = registry.get(provider_config.name) 217 if registration is None: 218 raise ConfigFailure( 219 ErrorCode.CONFIG_ERROR, 220 f"Invalid enabled web provider: {provider_config.name}", 221 ) 222 credential_kwargs: dict[str, object] = {} 223 if registration.requires_api_key: 224 private_value = provider_config.secret 225 if private_value is None: 226 raise ConfigFailure( 227 ErrorCode.CONFIG_ERROR, 228 f"Invalid enabled web provider: {provider_config.name}", 229 ) 230 credential_kwargs["secret"] = private_value 231 reserved = set(provider_config.options) & _RESERVED_WEB_ADAPTER_KWARGS 232 if reserved: 233 names = ", ".join(sorted(reserved)) 234 raise ConfigFailure( 235 ErrorCode.CONFIG_ERROR, 236 f"Reserved config key(s) for web provider {provider_config.name}: {names}", 237 ) 238 executor = HttpJsonExecutor( 239 http_client_factory(), 240 config.retry, 241 provider_name=provider_config.name, 242 ) 243 kwargs: dict[str, object] = { 244 "name": provider_config.name, 245 "http_executor": executor, 246 } 247 kwargs.update(credential_kwargs) 248 kwargs.update(provider_config.options) 249 try: 250 adapter = registration.factory(**kwargs) 251 except TypeError as exc: 252 raise ConfigFailure( 253 ErrorCode.CONFIG_ERROR, 254 f"Invalid configuration for web provider {provider_config.name}", 255 ) from exc 256 executors.append(executor) 257 if provider_config.enable_search: 258 search.append(cast(KeywordSearchProvider, adapter)) 259 if provider_config.enable_fetch: 260 fetch.append(cast(URLFetchProvider, adapter)) 261 return tuple(search), tuple(fetch), tuple(executors) 262 263 @classmethod 264 def _build_academic_providers( 265 cls, 266 providers: tuple[ResolvedAcademicProviderConfig, ...], 267 registry: AcademicProviderRegistry, 268 config: ResolvedConfig, 269 http_client_factory: HttpClientFactory, 270 ) -> tuple[tuple[AcademicSearchProvider, ...], tuple[HttpJsonExecutor, ...]]: 271 search: list[AcademicSearchProvider] = [] 272 executors: list[HttpJsonExecutor] = [] 273 for provider_config in providers: 274 registration = registry.get(provider_config.name) 275 if registration is None: 276 raise ConfigFailure( 277 ErrorCode.CONFIG_ERROR, 278 f"Invalid enabled academic provider: {provider_config.name}", 279 ) 280 reserved = set(provider_config.options) & _RESERVED_ACADEMIC_ADAPTER_KWARGS 281 if reserved: 282 names = ", ".join(sorted(reserved)) 283 raise ConfigFailure( 284 ErrorCode.CONFIG_ERROR, 285 f"Reserved config key(s) for academic provider {provider_config.name}: {names}", 286 ) 287 executor = HttpJsonExecutor( 288 http_client_factory(), 289 config.retry, 290 provider_name=provider_config.name, 291 ) 292 kwargs: dict[str, object] = {"executor": executor} 293 if provider_config.api_key is not None: 294 kwargs["api_key"] = provider_config.api_key 295 if provider_config.contact_email is not None: 296 kwargs["contact_email"] = provider_config.contact_email 297 kwargs.update(provider_config.options) 298 try: 299 adapter = registration.factory(**kwargs) 300 except TypeError as exc: 301 raise ConfigFailure( 302 ErrorCode.CONFIG_ERROR, 303 f"Invalid configuration for academic provider {provider_config.name}", 304 ) from exc 305 search.append(cast(AcademicSearchProvider, adapter)) 306 executors.append(executor) 307 return tuple(search), tuple(executors) 308 309 @classmethod 310 def _build_oa_resolver( 311 cls, 312 resolver_config: ResolvedOAResolverConfig | None, 313 registry: OAResolverRegistry, 314 config: ResolvedConfig, 315 http_client_factory: HttpClientFactory, 316 ) -> tuple[OAResolver | None, tuple[HttpJsonExecutor, ...]]: 317 if resolver_config is None: 318 return None, () 319 registration = registry.get(resolver_config.name) 320 if registration is None: 321 raise ConfigFailure( 322 ErrorCode.CONFIG_ERROR, 323 f"Invalid enabled OA resolver: {resolver_config.name}", 324 ) 325 reserved = set(resolver_config.options) & _RESERVED_ACADEMIC_ADAPTER_KWARGS 326 if reserved: 327 names = ", ".join(sorted(reserved)) 328 raise ConfigFailure( 329 ErrorCode.CONFIG_ERROR, 330 f"Reserved config key(s) for OA resolver {resolver_config.name}: {names}", 331 ) 332 executor = HttpJsonExecutor( 333 http_client_factory(), 334 config.retry, 335 provider_name=resolver_config.name, 336 ) 337 kwargs: dict[str, object] = {"executor": executor} 338 if resolver_config.api_key is not None: 339 kwargs["api_key"] = resolver_config.api_key 340 if resolver_config.contact_email is not None: 341 kwargs["contact_email"] = resolver_config.contact_email 342 kwargs.update(resolver_config.options) 343 try: 344 adapter = registration.factory(**kwargs) 345 except TypeError as exc: 346 raise ConfigFailure( 347 ErrorCode.CONFIG_ERROR, 348 f"Invalid configuration for OA resolver {resolver_config.name}", 349 ) from exc 350 return cast(OAResolver, adapter), (executor,) 351 352 @staticmethod 353 def _build_llm_clients( 354 config: ResolvedConfig, 355 quotas: ProviderQuotaManager, 356 http_client_factory: HttpClientFactory, 357 ) -> dict[str, LLMClient]: 358 clients: dict[str, LLMClient] = {} 359 for provider_config in config.llm.providers: 360 private_value = provider_config.secret 361 executor = HttpJsonExecutor( 362 http_client_factory(), 363 config.retry, 364 provider_name=provider_config.name, 365 ) 366 clients[provider_config.name] = OpenAIChatCompletionsClient( 367 name=provider_config.name, 368 api_url=provider_config.api_url, 369 secret=private_value, 370 executor=executor, 371 quota=quotas.get_llm(provider_config.name), 372 retry_policy=config.retry, 373 ) 374 return clients 375 376 async def aclose(self) -> None: 377 if self._closed: 378 return 379 self._closed = True 380 for executor in self._web_http_executors: 381 await executor.aclose() 382 for executor in self._academic_http_executors: 383 await executor.aclose() 384 for client in self.llm_clients.values(): 385 await client.aclose() 386 387 def __repr__(self) -> str: 388 resolver_name = self.oa_resolver.name if self.oa_resolver is not None else "-" 389 return ( 390 "Runtime(" 391 f"web_search={len(self.web_search_providers)}, " 392 f"web_fetch={len(self.web_fetch_providers)}, " 393 f"academic_search={len(self.academic_search_providers)}, " 394 f"oa_resolver={resolver_name}, " 395 f"llm_clients={len(self.llm_clients)}" 396 ")" 397 )
HttpClientFactory =
collections.abc.Callable[[], httpx.AsyncClient]
class
Runtime:
50class Runtime: 51 """Owns one in-memory gateway runtime and its transport clients.""" 52 53 def __init__( 54 self, 55 *, 56 quotas: ProviderQuotaManager, 57 web_search_providers: tuple[KeywordSearchProvider, ...], 58 web_fetch_providers: tuple[URLFetchProvider, ...], 59 academic_search_providers: tuple[AcademicSearchProvider, ...], 60 oa_resolver: OAResolver | None, 61 llm_clients: Mapping[str, LLMClient], 62 store: URLStore, 63 search_orchestrator: SearchOrchestrator, 64 paper_search_orchestrator: PaperSearchOrchestrator, 65 fetch_orchestrator: FetchOrchestrator, 66 web_http_executors: tuple[HttpJsonExecutor, ...], 67 academic_http_executors: tuple[HttpJsonExecutor, ...], 68 ) -> None: 69 self.quotas = quotas 70 self.web_search_providers = web_search_providers 71 self.web_fetch_providers = web_fetch_providers 72 self.academic_search_providers = academic_search_providers 73 self.oa_resolver = oa_resolver 74 self.llm_clients = dict(llm_clients) 75 self.store = store 76 self.search_orchestrator = search_orchestrator 77 self.paper_search_orchestrator = paper_search_orchestrator 78 self.fetch_orchestrator = fetch_orchestrator 79 self._web_http_executors = web_http_executors 80 self._academic_http_executors = academic_http_executors 81 self._closed = False 82 83 @classmethod 84 def build( 85 cls, 86 config: ResolvedConfig, 87 paths: RuntimePaths, 88 *, 89 registry: ProviderRegistry | None = None, 90 academic_registry: AcademicProviderRegistry | None = None, 91 oa_resolver_registry: OAResolverRegistry | None = None, 92 http_client_factory: HttpClientFactory = httpx.AsyncClient, 93 ) -> "Runtime": 94 provider_registry = registry or build_default_registry() 95 academic_provider_registry = academic_registry or build_default_academic_registry() 96 resolver_registry = oa_resolver_registry or build_default_oa_resolver_registry() 97 enabled_web = tuple( 98 item for item in config.web.providers if item.enable_search or item.enable_fetch 99 ) 100 enabled_academic = tuple(item for item in config.academic.providers if item.enabled) 101 quotas = ProviderQuotaManager( 102 web_limits={item.name: item.max_concurrency for item in enabled_web}, 103 llm_limits={item.name: item.max_concurrency for item in config.llm.providers}, 104 academic_limits={item.name: item.max_concurrency for item in enabled_academic}, 105 ) 106 web_search, web_fetch, web_executors = cls._build_web_providers( 107 enabled_web, 108 provider_registry, 109 config, 110 http_client_factory, 111 ) 112 academic_search, academic_executors = cls._build_academic_providers( 113 enabled_academic, 114 academic_provider_registry, 115 config, 116 http_client_factory, 117 ) 118 oa_resolver, resolver_executors = cls._build_oa_resolver( 119 config.oa_resolver, 120 resolver_registry, 121 config, 122 http_client_factory, 123 ) 124 llm_clients = cls._build_llm_clients(config, quotas, http_client_factory) 125 stages = LLMStages( 126 llm_clients, 127 judge=config.llm.judge, 128 safety=config.llm.safety, 129 content_clean=config.llm.content_clean, 130 focus_summary=config.llm.focus_summary, 131 ) 132 store = URLStore() 133 result_writer = ResultWriter(paths.results_dir) 134 paper_aggregator = PaperAggregator( 135 ( 136 *(provider.name for provider in academic_search), 137 *(f"llm:{item.provider}" for item in config.llm.search_invocations), 138 ) 139 ) 140 search_orchestrator = SearchOrchestrator( 141 keyword_providers=web_search, 142 llm_invocations=config.llm.search_invocations, 143 quotas=quotas, 144 stages=stages, 145 store=store, 146 result_writer=result_writer, 147 paper_aggregator=paper_aggregator, 148 paper_resolver=oa_resolver, 149 ) 150 paper_search_orchestrator = PaperSearchOrchestrator( 151 providers=academic_search, 152 quotas=quotas, 153 aggregator=paper_aggregator, 154 resolver=oa_resolver, 155 store=store, 156 result_writer=result_writer, 157 ) 158 fetch_orchestrator = FetchOrchestrator( 159 store=store, 160 scheduler=FetchScheduler(web_fetch, quotas, stages), 161 stages=stages, 162 ) 163 runtime = cls( 164 quotas=quotas, 165 web_search_providers=web_search, 166 web_fetch_providers=web_fetch, 167 academic_search_providers=academic_search, 168 oa_resolver=oa_resolver, 169 llm_clients=llm_clients, 170 store=store, 171 search_orchestrator=search_orchestrator, 172 paper_search_orchestrator=paper_search_orchestrator, 173 fetch_orchestrator=fetch_orchestrator, 174 web_http_executors=web_executors, 175 academic_http_executors=(*academic_executors, *resolver_executors), 176 ) 177 log_event( 178 logging.getLogger(__name__), 179 logging.DEBUG, 180 "runtime_built", 181 web_providers=",".join(item.name for item in enabled_web) or "-", 182 web_provider_count=len(enabled_web), 183 llm_providers=",".join(item.name for item in config.llm.providers) or "-", 184 llm_provider_count=len(config.llm.providers), 185 academic_providers=",".join(item.name for item in enabled_academic) or "-", 186 academic_provider_count=len(enabled_academic), 187 oa_resolver=config.oa_resolver.name if config.oa_resolver is not None else "-", 188 web_limits=",".join(f"{item.name}:{item.max_concurrency}" for item in enabled_web) 189 or "-", 190 llm_limits=",".join( 191 f"{item.name}:{item.max_concurrency}" for item in config.llm.providers 192 ) 193 or "-", 194 academic_limits=",".join( 195 f"{item.name}:{item.max_concurrency}" for item in enabled_academic 196 ) 197 or "-", 198 ) 199 return runtime 200 201 @classmethod 202 def _build_web_providers( 203 cls, 204 providers: tuple[ResolvedWebProviderConfig, ...], 205 registry: ProviderRegistry, 206 config: ResolvedConfig, 207 http_client_factory: HttpClientFactory, 208 ) -> tuple[ 209 tuple[KeywordSearchProvider, ...], 210 tuple[URLFetchProvider, ...], 211 tuple[HttpJsonExecutor, ...], 212 ]: 213 search: list[KeywordSearchProvider] = [] 214 fetch: list[URLFetchProvider] = [] 215 executors: list[HttpJsonExecutor] = [] 216 for provider_config in providers: 217 registration = registry.get(provider_config.name) 218 if registration is None: 219 raise ConfigFailure( 220 ErrorCode.CONFIG_ERROR, 221 f"Invalid enabled web provider: {provider_config.name}", 222 ) 223 credential_kwargs: dict[str, object] = {} 224 if registration.requires_api_key: 225 private_value = provider_config.secret 226 if private_value is None: 227 raise ConfigFailure( 228 ErrorCode.CONFIG_ERROR, 229 f"Invalid enabled web provider: {provider_config.name}", 230 ) 231 credential_kwargs["secret"] = private_value 232 reserved = set(provider_config.options) & _RESERVED_WEB_ADAPTER_KWARGS 233 if reserved: 234 names = ", ".join(sorted(reserved)) 235 raise ConfigFailure( 236 ErrorCode.CONFIG_ERROR, 237 f"Reserved config key(s) for web provider {provider_config.name}: {names}", 238 ) 239 executor = HttpJsonExecutor( 240 http_client_factory(), 241 config.retry, 242 provider_name=provider_config.name, 243 ) 244 kwargs: dict[str, object] = { 245 "name": provider_config.name, 246 "http_executor": executor, 247 } 248 kwargs.update(credential_kwargs) 249 kwargs.update(provider_config.options) 250 try: 251 adapter = registration.factory(**kwargs) 252 except TypeError as exc: 253 raise ConfigFailure( 254 ErrorCode.CONFIG_ERROR, 255 f"Invalid configuration for web provider {provider_config.name}", 256 ) from exc 257 executors.append(executor) 258 if provider_config.enable_search: 259 search.append(cast(KeywordSearchProvider, adapter)) 260 if provider_config.enable_fetch: 261 fetch.append(cast(URLFetchProvider, adapter)) 262 return tuple(search), tuple(fetch), tuple(executors) 263 264 @classmethod 265 def _build_academic_providers( 266 cls, 267 providers: tuple[ResolvedAcademicProviderConfig, ...], 268 registry: AcademicProviderRegistry, 269 config: ResolvedConfig, 270 http_client_factory: HttpClientFactory, 271 ) -> tuple[tuple[AcademicSearchProvider, ...], tuple[HttpJsonExecutor, ...]]: 272 search: list[AcademicSearchProvider] = [] 273 executors: list[HttpJsonExecutor] = [] 274 for provider_config in providers: 275 registration = registry.get(provider_config.name) 276 if registration is None: 277 raise ConfigFailure( 278 ErrorCode.CONFIG_ERROR, 279 f"Invalid enabled academic provider: {provider_config.name}", 280 ) 281 reserved = set(provider_config.options) & _RESERVED_ACADEMIC_ADAPTER_KWARGS 282 if reserved: 283 names = ", ".join(sorted(reserved)) 284 raise ConfigFailure( 285 ErrorCode.CONFIG_ERROR, 286 f"Reserved config key(s) for academic provider {provider_config.name}: {names}", 287 ) 288 executor = HttpJsonExecutor( 289 http_client_factory(), 290 config.retry, 291 provider_name=provider_config.name, 292 ) 293 kwargs: dict[str, object] = {"executor": executor} 294 if provider_config.api_key is not None: 295 kwargs["api_key"] = provider_config.api_key 296 if provider_config.contact_email is not None: 297 kwargs["contact_email"] = provider_config.contact_email 298 kwargs.update(provider_config.options) 299 try: 300 adapter = registration.factory(**kwargs) 301 except TypeError as exc: 302 raise ConfigFailure( 303 ErrorCode.CONFIG_ERROR, 304 f"Invalid configuration for academic provider {provider_config.name}", 305 ) from exc 306 search.append(cast(AcademicSearchProvider, adapter)) 307 executors.append(executor) 308 return tuple(search), tuple(executors) 309 310 @classmethod 311 def _build_oa_resolver( 312 cls, 313 resolver_config: ResolvedOAResolverConfig | None, 314 registry: OAResolverRegistry, 315 config: ResolvedConfig, 316 http_client_factory: HttpClientFactory, 317 ) -> tuple[OAResolver | None, tuple[HttpJsonExecutor, ...]]: 318 if resolver_config is None: 319 return None, () 320 registration = registry.get(resolver_config.name) 321 if registration is None: 322 raise ConfigFailure( 323 ErrorCode.CONFIG_ERROR, 324 f"Invalid enabled OA resolver: {resolver_config.name}", 325 ) 326 reserved = set(resolver_config.options) & _RESERVED_ACADEMIC_ADAPTER_KWARGS 327 if reserved: 328 names = ", ".join(sorted(reserved)) 329 raise ConfigFailure( 330 ErrorCode.CONFIG_ERROR, 331 f"Reserved config key(s) for OA resolver {resolver_config.name}: {names}", 332 ) 333 executor = HttpJsonExecutor( 334 http_client_factory(), 335 config.retry, 336 provider_name=resolver_config.name, 337 ) 338 kwargs: dict[str, object] = {"executor": executor} 339 if resolver_config.api_key is not None: 340 kwargs["api_key"] = resolver_config.api_key 341 if resolver_config.contact_email is not None: 342 kwargs["contact_email"] = resolver_config.contact_email 343 kwargs.update(resolver_config.options) 344 try: 345 adapter = registration.factory(**kwargs) 346 except TypeError as exc: 347 raise ConfigFailure( 348 ErrorCode.CONFIG_ERROR, 349 f"Invalid configuration for OA resolver {resolver_config.name}", 350 ) from exc 351 return cast(OAResolver, adapter), (executor,) 352 353 @staticmethod 354 def _build_llm_clients( 355 config: ResolvedConfig, 356 quotas: ProviderQuotaManager, 357 http_client_factory: HttpClientFactory, 358 ) -> dict[str, LLMClient]: 359 clients: dict[str, LLMClient] = {} 360 for provider_config in config.llm.providers: 361 private_value = provider_config.secret 362 executor = HttpJsonExecutor( 363 http_client_factory(), 364 config.retry, 365 provider_name=provider_config.name, 366 ) 367 clients[provider_config.name] = OpenAIChatCompletionsClient( 368 name=provider_config.name, 369 api_url=provider_config.api_url, 370 secret=private_value, 371 executor=executor, 372 quota=quotas.get_llm(provider_config.name), 373 retry_policy=config.retry, 374 ) 375 return clients 376 377 async def aclose(self) -> None: 378 if self._closed: 379 return 380 self._closed = True 381 for executor in self._web_http_executors: 382 await executor.aclose() 383 for executor in self._academic_http_executors: 384 await executor.aclose() 385 for client in self.llm_clients.values(): 386 await client.aclose() 387 388 def __repr__(self) -> str: 389 resolver_name = self.oa_resolver.name if self.oa_resolver is not None else "-" 390 return ( 391 "Runtime(" 392 f"web_search={len(self.web_search_providers)}, " 393 f"web_fetch={len(self.web_fetch_providers)}, " 394 f"academic_search={len(self.academic_search_providers)}, " 395 f"oa_resolver={resolver_name}, " 396 f"llm_clients={len(self.llm_clients)}" 397 ")" 398 )
Owns one in-memory gateway runtime and its transport clients.
Runtime( *, quotas: agent_search_gateway.concurrency.ProviderQuotaManager, web_search_providers: tuple[agent_search_gateway.providers.contracts.KeywordSearchProvider, ...], web_fetch_providers: tuple[agent_search_gateway.providers.contracts.URLFetchProvider, ...], academic_search_providers: tuple[agent_search_gateway.providers.contracts.AcademicSearchProvider, ...], oa_resolver: agent_search_gateway.providers.contracts.OAResolver | None, llm_clients: Mapping[str, agent_search_gateway.providers.contracts.LLMClient], store: agent_search_gateway.url_store.URLStore, search_orchestrator: agent_search_gateway.orchestrators.search.SearchOrchestrator, paper_search_orchestrator: agent_search_gateway.orchestrators.paper.PaperSearchOrchestrator, fetch_orchestrator: agent_search_gateway.orchestrators.fetch.FetchOrchestrator, web_http_executors: tuple[agent_search_gateway.providers.http.HttpJsonExecutor, ...], academic_http_executors: tuple[agent_search_gateway.providers.http.HttpJsonExecutor, ...])
53 def __init__( 54 self, 55 *, 56 quotas: ProviderQuotaManager, 57 web_search_providers: tuple[KeywordSearchProvider, ...], 58 web_fetch_providers: tuple[URLFetchProvider, ...], 59 academic_search_providers: tuple[AcademicSearchProvider, ...], 60 oa_resolver: OAResolver | None, 61 llm_clients: Mapping[str, LLMClient], 62 store: URLStore, 63 search_orchestrator: SearchOrchestrator, 64 paper_search_orchestrator: PaperSearchOrchestrator, 65 fetch_orchestrator: FetchOrchestrator, 66 web_http_executors: tuple[HttpJsonExecutor, ...], 67 academic_http_executors: tuple[HttpJsonExecutor, ...], 68 ) -> None: 69 self.quotas = quotas 70 self.web_search_providers = web_search_providers 71 self.web_fetch_providers = web_fetch_providers 72 self.academic_search_providers = academic_search_providers 73 self.oa_resolver = oa_resolver 74 self.llm_clients = dict(llm_clients) 75 self.store = store 76 self.search_orchestrator = search_orchestrator 77 self.paper_search_orchestrator = paper_search_orchestrator 78 self.fetch_orchestrator = fetch_orchestrator 79 self._web_http_executors = web_http_executors 80 self._academic_http_executors = academic_http_executors 81 self._closed = False
@classmethod
def
build( cls, config: agent_search_gateway.config.ResolvedConfig, paths: agent_search_gateway.paths.RuntimePaths, *, registry: agent_search_gateway.providers.registry.ProviderRegistry | None = None, academic_registry: agent_search_gateway.providers.academic.registry.AcademicProviderRegistry | None = None, oa_resolver_registry: agent_search_gateway.providers.academic.registry.OAResolverRegistry | None = None, http_client_factory: Callable[[], httpx.AsyncClient] = <class 'httpx.AsyncClient'>) -> Runtime:
83 @classmethod 84 def build( 85 cls, 86 config: ResolvedConfig, 87 paths: RuntimePaths, 88 *, 89 registry: ProviderRegistry | None = None, 90 academic_registry: AcademicProviderRegistry | None = None, 91 oa_resolver_registry: OAResolverRegistry | None = None, 92 http_client_factory: HttpClientFactory = httpx.AsyncClient, 93 ) -> "Runtime": 94 provider_registry = registry or build_default_registry() 95 academic_provider_registry = academic_registry or build_default_academic_registry() 96 resolver_registry = oa_resolver_registry or build_default_oa_resolver_registry() 97 enabled_web = tuple( 98 item for item in config.web.providers if item.enable_search or item.enable_fetch 99 ) 100 enabled_academic = tuple(item for item in config.academic.providers if item.enabled) 101 quotas = ProviderQuotaManager( 102 web_limits={item.name: item.max_concurrency for item in enabled_web}, 103 llm_limits={item.name: item.max_concurrency for item in config.llm.providers}, 104 academic_limits={item.name: item.max_concurrency for item in enabled_academic}, 105 ) 106 web_search, web_fetch, web_executors = cls._build_web_providers( 107 enabled_web, 108 provider_registry, 109 config, 110 http_client_factory, 111 ) 112 academic_search, academic_executors = cls._build_academic_providers( 113 enabled_academic, 114 academic_provider_registry, 115 config, 116 http_client_factory, 117 ) 118 oa_resolver, resolver_executors = cls._build_oa_resolver( 119 config.oa_resolver, 120 resolver_registry, 121 config, 122 http_client_factory, 123 ) 124 llm_clients = cls._build_llm_clients(config, quotas, http_client_factory) 125 stages = LLMStages( 126 llm_clients, 127 judge=config.llm.judge, 128 safety=config.llm.safety, 129 content_clean=config.llm.content_clean, 130 focus_summary=config.llm.focus_summary, 131 ) 132 store = URLStore() 133 result_writer = ResultWriter(paths.results_dir) 134 paper_aggregator = PaperAggregator( 135 ( 136 *(provider.name for provider in academic_search), 137 *(f"llm:{item.provider}" for item in config.llm.search_invocations), 138 ) 139 ) 140 search_orchestrator = SearchOrchestrator( 141 keyword_providers=web_search, 142 llm_invocations=config.llm.search_invocations, 143 quotas=quotas, 144 stages=stages, 145 store=store, 146 result_writer=result_writer, 147 paper_aggregator=paper_aggregator, 148 paper_resolver=oa_resolver, 149 ) 150 paper_search_orchestrator = PaperSearchOrchestrator( 151 providers=academic_search, 152 quotas=quotas, 153 aggregator=paper_aggregator, 154 resolver=oa_resolver, 155 store=store, 156 result_writer=result_writer, 157 ) 158 fetch_orchestrator = FetchOrchestrator( 159 store=store, 160 scheduler=FetchScheduler(web_fetch, quotas, stages), 161 stages=stages, 162 ) 163 runtime = cls( 164 quotas=quotas, 165 web_search_providers=web_search, 166 web_fetch_providers=web_fetch, 167 academic_search_providers=academic_search, 168 oa_resolver=oa_resolver, 169 llm_clients=llm_clients, 170 store=store, 171 search_orchestrator=search_orchestrator, 172 paper_search_orchestrator=paper_search_orchestrator, 173 fetch_orchestrator=fetch_orchestrator, 174 web_http_executors=web_executors, 175 academic_http_executors=(*academic_executors, *resolver_executors), 176 ) 177 log_event( 178 logging.getLogger(__name__), 179 logging.DEBUG, 180 "runtime_built", 181 web_providers=",".join(item.name for item in enabled_web) or "-", 182 web_provider_count=len(enabled_web), 183 llm_providers=",".join(item.name for item in config.llm.providers) or "-", 184 llm_provider_count=len(config.llm.providers), 185 academic_providers=",".join(item.name for item in enabled_academic) or "-", 186 academic_provider_count=len(enabled_academic), 187 oa_resolver=config.oa_resolver.name if config.oa_resolver is not None else "-", 188 web_limits=",".join(f"{item.name}:{item.max_concurrency}" for item in enabled_web) 189 or "-", 190 llm_limits=",".join( 191 f"{item.name}:{item.max_concurrency}" for item in config.llm.providers 192 ) 193 or "-", 194 academic_limits=",".join( 195 f"{item.name}:{item.max_concurrency}" for item in enabled_academic 196 ) 197 or "-", 198 ) 199 return runtime
async def
aclose(self) -> None:
377 async def aclose(self) -> None: 378 if self._closed: 379 return 380 self._closed = True 381 for executor in self._web_http_executors: 382 await executor.aclose() 383 for executor in self._academic_http_executors: 384 await executor.aclose() 385 for client in self.llm_clients.values(): 386 await client.aclose()