diff options
Diffstat (limited to 'contrib/python')
| -rw-r--r-- | contrib/python/ydb/py3/.dist-info/METADATA | 2 | ||||
| -rw-r--r-- | contrib/python/ydb/py3/ya.make | 3 | ||||
| -rw-r--r-- | contrib/python/ydb/py3/ydb/_constants.py | 2 | ||||
| -rw-r--r-- | contrib/python/ydb/py3/ydb/_grpc/grpcwrapper/common_utils.py | 33 | ||||
| -rw-r--r-- | contrib/python/ydb/py3/ydb/_topic_reader/topic_reader_asyncio.py | 10 | ||||
| -rw-r--r-- | contrib/python/ydb/py3/ydb/_topic_writer/topic_writer_asyncio.py | 8 | ||||
| -rw-r--r-- | contrib/python/ydb/py3/ydb/aio/query/session.py | 5 | ||||
| -rw-r--r-- | contrib/python/ydb/py3/ydb/query/session.py | 14 | ||||
| -rw-r--r-- | contrib/python/ydb/py3/ydb/ydb_version.py | 2 |
9 files changed, 60 insertions, 19 deletions
diff --git a/contrib/python/ydb/py3/.dist-info/METADATA b/contrib/python/ydb/py3/.dist-info/METADATA index b93ac851639..09428b518f6 100644 --- a/contrib/python/ydb/py3/.dist-info/METADATA +++ b/contrib/python/ydb/py3/.dist-info/METADATA @@ -1,6 +1,6 @@ Metadata-Version: 2.1 Name: ydb -Version: 3.21.6 +Version: 3.21.7 Summary: YDB Python SDK Home-page: http://github.com/ydb-platform/ydb-python-sdk Author: Yandex LLC diff --git a/contrib/python/ydb/py3/ya.make b/contrib/python/ydb/py3/ya.make index 7e224e5c310..59a91bb5816 100644 --- a/contrib/python/ydb/py3/ya.make +++ b/contrib/python/ydb/py3/ya.make @@ -2,7 +2,7 @@ PY3_LIBRARY() -VERSION(3.21.6) +VERSION(3.21.7) LICENSE(Apache-2.0) @@ -24,6 +24,7 @@ PY_SRCS( TOP_LEVEL ydb/__init__.py ydb/_apis.py + ydb/_constants.py ydb/_errors.py ydb/_grpc/__init__.py ydb/_grpc/common/__init__.py diff --git a/contrib/python/ydb/py3/ydb/_constants.py b/contrib/python/ydb/py3/ydb/_constants.py new file mode 100644 index 00000000000..af52c4e3d42 --- /dev/null +++ b/contrib/python/ydb/py3/ydb/_constants.py @@ -0,0 +1,2 @@ +DEFAULT_INITIAL_RESPONSE_TIMEOUT = 600 +DEFAULT_LONG_STREAM_TIMEOUT = 31536000 # year diff --git a/contrib/python/ydb/py3/ydb/_grpc/grpcwrapper/common_utils.py b/contrib/python/ydb/py3/ydb/_grpc/grpcwrapper/common_utils.py index 95a5744313e..004faf1524f 100644 --- a/contrib/python/ydb/py3/ydb/_grpc/grpcwrapper/common_utils.py +++ b/contrib/python/ydb/py3/ydb/_grpc/grpcwrapper/common_utils.py @@ -33,6 +33,8 @@ except ImportError: from contrib.ydb.public.api.protos import ydb_topic_pb2, ydb_issue_message_pb2 from ... import issues, connection +from ...settings import BaseRequestSettings +from ..._constants import DEFAULT_LONG_STREAM_TIMEOUT class IFromProto(abc.ABC): @@ -130,7 +132,7 @@ class SyncToAsyncIterator: class IGrpcWrapperAsyncIO(abc.ABC): @abc.abstractmethod - async def receive(self) -> Any: + async def receive(self, timeout: Optional[int] = None) -> Any: ... @abc.abstractmethod @@ -160,6 +162,13 @@ class GrpcWrapperAsyncIO(IGrpcWrapperAsyncIO): self._stream_call = None self._wait_executor = None + self._stream_settings: BaseRequestSettings = ( + BaseRequestSettings() + .with_operation_timeout(DEFAULT_LONG_STREAM_TIMEOUT) + .with_cancel_after(DEFAULT_LONG_STREAM_TIMEOUT) + .with_timeout(DEFAULT_LONG_STREAM_TIMEOUT) + ) + def __del__(self): self._clean_executor(wait=False) @@ -187,6 +196,7 @@ class GrpcWrapperAsyncIO(IGrpcWrapperAsyncIO): requests_iterator, stub, method, + settings=self._stream_settings, ) self._stream_call = stream_call self.from_server_grpc = stream_call.__aiter__() @@ -195,14 +205,29 @@ class GrpcWrapperAsyncIO(IGrpcWrapperAsyncIO): requests_iterator = AsyncQueueToSyncIteratorAsyncIO(self.from_client_grpc) self._wait_executor = concurrent.futures.ThreadPoolExecutor(max_workers=1) - stream_call = await to_thread(driver, requests_iterator, stub, method, executor=self._wait_executor) + stream_call = await to_thread( + driver, + requests_iterator, + stub, + method, + executor=self._wait_executor, + settings=self._stream_settings, + ) self._stream_call = stream_call self.from_server_grpc = SyncToAsyncIterator(stream_call.__iter__(), self._wait_executor) - async def receive(self) -> Any: + async def receive(self, timeout: Optional[int] = None) -> Any: # todo handle grpc exceptions and convert it to internal exceptions try: - grpc_message = await self.from_server_grpc.__anext__() + if timeout is None: + grpc_message = await self.from_server_grpc.__anext__() + else: + + async def get_response(): + return await self.from_server_grpc.__anext__() + + grpc_message = await asyncio.wait_for(get_response(), timeout) + except (grpc.RpcError, grpc.aio.AioRpcError) as e: raise connection._rpc_error_handler(self._connection_state, e) diff --git a/contrib/python/ydb/py3/ydb/_topic_reader/topic_reader_asyncio.py b/contrib/python/ydb/py3/ydb/_topic_reader/topic_reader_asyncio.py index 7baadacb3e0..b855a80b99c 100644 --- a/contrib/python/ydb/py3/ydb/_topic_reader/topic_reader_asyncio.py +++ b/contrib/python/ydb/py3/ydb/_topic_reader/topic_reader_asyncio.py @@ -38,6 +38,8 @@ from ..query.base import TxEvent if typing.TYPE_CHECKING: from ..query.transaction import BaseQueryTxContext +from .._constants import DEFAULT_INITIAL_RESPONSE_TIMEOUT + logger = logging.getLogger(__name__) @@ -490,7 +492,13 @@ class ReaderStream: logger.debug("reader stream %s send init request", self._id) stream.write(StreamReadMessage.FromClient(client_message=init_message)) - init_response = await stream.receive() # type: StreamReadMessage.FromServer + try: + init_response = await stream.receive( + timeout=DEFAULT_INITIAL_RESPONSE_TIMEOUT + ) # type: StreamReadMessage.FromServer + except asyncio.TimeoutError: + raise TopicReaderError("Timeout waiting for init response") + if isinstance(init_response.server_message, StreamReadMessage.InitResponse): self._session_id = init_response.server_message.session_id logger.debug("reader stream %s initialized session=%s", self._id, self._session_id) diff --git a/contrib/python/ydb/py3/ydb/_topic_writer/topic_writer_asyncio.py b/contrib/python/ydb/py3/ydb/_topic_writer/topic_writer_asyncio.py index eeecbfd2e81..d39606d1797 100644 --- a/contrib/python/ydb/py3/ydb/_topic_writer/topic_writer_asyncio.py +++ b/contrib/python/ydb/py3/ydb/_topic_writer/topic_writer_asyncio.py @@ -49,6 +49,8 @@ from ..query.base import TxEvent if typing.TYPE_CHECKING: from ..query.transaction import BaseQueryTxContext +from .._constants import DEFAULT_INITIAL_RESPONSE_TIMEOUT + logger = logging.getLogger(__name__) @@ -799,7 +801,11 @@ class WriterAsyncIOStream: logger.debug("writer stream %s send init request", self._id) stream.write(StreamWriteMessage.FromClient(init_message)) - resp = await stream.receive() + try: + resp = await stream.receive(timeout=DEFAULT_INITIAL_RESPONSE_TIMEOUT) + except asyncio.TimeoutError: + raise TopicWriterError("Timeout waiting for init response") + self._ensure_ok(resp) if not isinstance(resp, StreamWriteMessage.InitResponse): raise TopicWriterError("Unexpected answer for init request: %s" % resp) diff --git a/contrib/python/ydb/py3/ydb/aio/query/session.py b/contrib/python/ydb/py3/ydb/aio/query/session.py index fe857878a54..7a7ba5baef2 100644 --- a/contrib/python/ydb/py3/ydb/aio/query/session.py +++ b/contrib/python/ydb/py3/ydb/aio/query/session.py @@ -15,10 +15,11 @@ from ..._grpc.grpcwrapper import ydb_query_public_types as _ydb_query_public from ...query import base from ...query.session import ( BaseQuerySession, - DEFAULT_ATTACH_FIRST_RESP_TIMEOUT, QuerySessionStateEnum, ) +from ..._constants import DEFAULT_INITIAL_RESPONSE_TIMEOUT + class QuerySession(BaseQuerySession): """Session object for Query Service. It is not recommended to control @@ -47,7 +48,7 @@ class QuerySession(BaseQuerySession): try: first_response = await _utilities.get_first_message_with_timeout( self._status_stream, - DEFAULT_ATTACH_FIRST_RESP_TIMEOUT, + DEFAULT_INITIAL_RESPONSE_TIMEOUT, ) if first_response.status != issues.StatusCode.SUCCESS: raise RuntimeError("Failed to attach session") diff --git a/contrib/python/ydb/py3/ydb/query/session.py b/contrib/python/ydb/py3/ydb/query/session.py index 3cc6c13d4e9..5cfbdc6ceba 100644 --- a/contrib/python/ydb/py3/ydb/query/session.py +++ b/contrib/python/ydb/py3/ydb/query/session.py @@ -18,12 +18,10 @@ from .._grpc.grpcwrapper import ydb_query_public_types as _ydb_query_public from .transaction import QueryTxContext - -logger = logging.getLogger(__name__) +from .._constants import DEFAULT_INITIAL_RESPONSE_TIMEOUT, DEFAULT_LONG_STREAM_TIMEOUT -DEFAULT_ATTACH_FIRST_RESP_TIMEOUT = 600 -DEFAULT_ATTACH_LONG_TIMEOUT = 31536000 # year +logger = logging.getLogger(__name__) class QuerySessionStateEnum(enum.Enum): @@ -142,9 +140,9 @@ class BaseQuerySession: self._state = QuerySessionState(settings) self._attach_settings: BaseRequestSettings = ( BaseRequestSettings() - .with_operation_timeout(DEFAULT_ATTACH_LONG_TIMEOUT) - .with_cancel_after(DEFAULT_ATTACH_LONG_TIMEOUT) - .with_timeout(DEFAULT_ATTACH_LONG_TIMEOUT) + .with_operation_timeout(DEFAULT_LONG_STREAM_TIMEOUT) + .with_cancel_after(DEFAULT_LONG_STREAM_TIMEOUT) + .with_timeout(DEFAULT_LONG_STREAM_TIMEOUT) ) self._last_query_stats = None @@ -233,7 +231,7 @@ class QuerySession(BaseQuerySession): _stream = None - def _attach(self, first_resp_timeout: int = DEFAULT_ATTACH_FIRST_RESP_TIMEOUT) -> None: + def _attach(self, first_resp_timeout: int = DEFAULT_INITIAL_RESPONSE_TIMEOUT) -> None: self._stream = self._attach_call() status_stream = _utilities.SyncResponseIterator( self._stream, diff --git a/contrib/python/ydb/py3/ydb/ydb_version.py b/contrib/python/ydb/py3/ydb/ydb_version.py index 9062621ad98..d953c2ebcdd 100644 --- a/contrib/python/ydb/py3/ydb/ydb_version.py +++ b/contrib/python/ydb/py3/ydb/ydb_version.py @@ -1 +1 @@ -VERSION = "3.21.6" +VERSION = "3.21.7" |
