Edit on GitHub

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.

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
quotas
web_search_providers
web_fetch_providers
academic_search_providers
oa_resolver
llm_clients
store
search_orchestrator
paper_search_orchestrator
fetch_orchestrator
@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()