Coverage for /usr/local/lib/python3.10/site-packages/opal_common-0.0.0-py3.10.egg/opal_common/fetcher/providers/fastapi_rpc_fetch_provider.py: 0%

30 statements  

« prev     ^ index     » next       coverage.py v7.15.2, created at 2026-10-10 11:54 +0000

1"""Simple HTTP get data fetcher using requests supports.""" 

2 

3from typing import ClassVar, Set 

4 

5from fastapi_websocket_rpc.rpc_methods import RpcMethodsBase 

6from fastapi_websocket_rpc.websocket_rpc_client import WebSocketRpcClient 

7from opal_common.fetcher.events import FetcherConfig, FetchEvent 

8from opal_common.fetcher.fetch_provider import BaseFetchProvider 

9from opal_common.fetcher.logger import get_logger 

10from opal_common.http_utils import redact_url 

11 

12logger = get_logger("rpc_fetch_provider") 

13 

14 

15class FastApiRpcFetchConfig(FetcherConfig): 

16 """Config for FastApiRpcFetchConfig's Adding HTTP headers.""" 

17 

18 # ``rpc_arguments`` may carry credentials - mask it in repr/str. 

19 _redacted_repr_fields: ClassVar[Set[str]] = {"rpc_arguments"} 

20 

21 rpc_method_name: str 

22 rpc_arguments: dict 

23 

24 

25class FastApiRpcFetchEvent(FetchEvent): 

26 fetcher: str = "FastApiRpcFetchProvider" 

27 config: FastApiRpcFetchConfig 

28 

29 

30class FastApiRpcFetchProvider(BaseFetchProvider): 

31 def __init__(self, event: FastApiRpcFetchEvent) -> None: 

32 self._event: FastApiRpcFetchEvent 

33 super().__init__(event) 

34 

35 def parse_event(self, event: FetchEvent) -> FastApiRpcFetchEvent: 

36 return FastApiRpcFetchEvent( 

37 **event.dict(exclude={"config"}), config=event.config 

38 ) 

39 

40 async def _fetch_(self): 

41 assert ( 

42 self._event is not None 

43 ), "FastApiRpcFetchEvent not provided for FastApiRpcFetchProvider" 

44 args = self._event.config.rpc_arguments 

45 method = self._event.config.rpc_method_name 

46 result = None 

47 # Note: ``args`` (rpc_arguments) may carry credentials - never log it. 

48 logger.info( 

49 f"{self.__class__.__name__} fetching from {redact_url(self._url)} with RPC call {method}" 

50 ) 

51 async with WebSocketRpcClient( 

52 self._url, 

53 # we don't expose anything to the server 

54 RpcMethodsBase(), 

55 default_response_timeout=4, 

56 ) as client: 

57 result = await client.call(method, args) 

58 return result