summaryrefslogtreecommitdiffstats
path: root/contrib/python
diff options
context:
space:
mode:
Diffstat (limited to 'contrib/python')
-rw-r--r--contrib/python/ydb/py3/.dist-info/METADATA2
-rw-r--r--contrib/python/ydb/py3/ya.make3
-rw-r--r--contrib/python/ydb/py3/ydb/_constants.py2
-rw-r--r--contrib/python/ydb/py3/ydb/_grpc/grpcwrapper/common_utils.py33
-rw-r--r--contrib/python/ydb/py3/ydb/_topic_reader/topic_reader_asyncio.py10
-rw-r--r--contrib/python/ydb/py3/ydb/_topic_writer/topic_writer_asyncio.py8
-rw-r--r--contrib/python/ydb/py3/ydb/aio/query/session.py5
-rw-r--r--contrib/python/ydb/py3/ydb/query/session.py14
-rw-r--r--contrib/python/ydb/py3/ydb/ydb_version.py2
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"