From f4a5e69f10228e50cc7c25f57b54807f3829fe93 Mon Sep 17 00:00:00 2001 From: yijxie Date: Wed, 6 May 2020 15:06:42 -0700 Subject: [PATCH 01/15] Fix pylint errors other than async-sync overrides --- .../azure/servicebus/_common/message.py | 21 ++++++++++++------- .../servicebus/_common/receiver_mixins.py | 14 ++++++------- .../azure/servicebus/_common/utils.py | 20 ------------------ .../servicebus/_control_client/models.py | 2 +- .../azure/servicebus/_servicebus_client.py | 1 - .../azure/servicebus/_servicebus_receiver.py | 4 ++-- .../azure/servicebus/_servicebus_sender.py | 2 +- .../azure/servicebus/_servicebus_session.py | 3 +-- .../_servicebus_session_receiver.py | 7 +------ .../azure/servicebus/aio/_async_message.py | 9 ++++---- .../azure/servicebus/aio/_async_utils.py | 21 ++++++++++++++++++- .../servicebus/aio/_base_handler_async.py | 4 +++- .../aio/_servicebus_sender_async.py | 2 +- .../aio/_servicebus_session_async.py | 2 +- .../aio/_servicebus_session_receiver_async.py | 3 +-- .../azure/servicebus/exceptions.py | 2 +- .../azure-servicebus/dev_requirements.txt | 1 + 17 files changed, 60 insertions(+), 58 deletions(-) diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/_common/message.py b/sdk/servicebus/azure-servicebus/azure/servicebus/_common/message.py index 16c52bfd7fa9..60a4b0feef20 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/_common/message.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/_common/message.py @@ -8,10 +8,10 @@ import uuid import functools import logging -from typing import Optional, List, Union, Generator +from typing import Optional, List, Union, Generator, TYPE_CHECKING import uamqp -from uamqp import types, errors +from uamqp import types from .constants import ( _BATCH_MESSAGE_OVERHEAD_COST, @@ -46,6 +46,9 @@ MessageSettleFailed, MessageContentTooLarge) from .utils import utc_from_timestamp, utc_now +if TYPE_CHECKING: + from .._servicebus_receiver import ServiceBusReceiver + from .._servicebus_session_receiver import ServiceBusSessionReceiver _LOGGER = logging.getLogger(__name__) @@ -86,8 +89,6 @@ def __init__(self, body, **kwargs): self._annotations = {} self._app_properties = {} - self._expiry = None - self._receiver = None self.session_id = kwargs.get("session_id", None) if 'message' in kwargs: self.message = kwargs['message'] @@ -317,7 +318,8 @@ def __len__(self): def _from_list(self, messages): for each in messages: if not isinstance(each, Message): - raise ValueError("Populating a message batch only supports iterables containing Message Objects. Received instead: {}".format(each.__class__.__name__)) + raise ValueError("Populating a message batch only supports iterables containing Message Objects. " + "Received instead: {}".format(each.__class__.__name__)) self.add(each) @property @@ -451,6 +453,8 @@ def __init__(self, message, mode=ReceiveSettleMode.PeekLock, **kwargs): self._settled = (mode == ReceiveSettleMode.ReceiveAndDelete) self._is_deferred_message = kwargs.get("is_deferred_message", False) self.auto_renew_error = None + self._receiver = None # type: Union[ServiceBusReceiver, ServiceBusSessionReceiver] + self._expiry = None def _check_live(self, action): # pylint: disable=no-member @@ -519,7 +523,10 @@ def _settle_via_mgmt_link(self, settle_operation, dead_letter_details=None): ) raise ValueError("Unsupported settle operation type: {}".format(settle_operation)) - def _settle_via_receiver_link(self, settle_operation, dead_letter_details=None): + def _settle_via_receiver_link(self, settle_operation, dead_letter_details=None): # pylint: disable=unused-argument + # dead_letter_detail is not used because of uamqp receiver link doesn't accept it while it + # should be accepted. Will revisit this later. + # uamqp management link accepts dead_letter_details. Refer to method _settle_via_mgmt_link if settle_operation == MESSAGE_COMPLETE: return functools.partial(self.message.accept) if settle_operation == MESSAGE_ABANDON: @@ -700,5 +707,5 @@ def renew_lock(self): if not token: raise ValueError("Unable to renew lock - no lock token found.") - expiry = self._receiver._renew_locks(token) # pylint: disable=protected-access + expiry = self._receiver._renew_locks(token) # pylint: disable=protected-access no-member self._expiry = utc_from_timestamp(expiry[MGMT_RESPONSE_MESSAGE_EXPIRATION][0]/1000.0) diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/_common/receiver_mixins.py b/sdk/servicebus/azure-servicebus/azure/servicebus/_common/receiver_mixins.py index 15b43563d8d1..ea6f7d1e08db 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/_common/receiver_mixins.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/_common/receiver_mixins.py @@ -4,12 +4,13 @@ # license information. # ------------------------------------------------------------------------- import uuid +from uamqp import Source from .message import ReceivedMessage from .constants import ( - NEXT_AVAILABLE, - SESSION_FILTER, - SESSION_LOCKED_UNTIL, - DATETIMEOFFSET_EPOCH, + NEXT_AVAILABLE, + SESSION_FILTER, + SESSION_LOCKED_UNTIL, + DATETIMEOFFSET_EPOCH, MGMT_REQUEST_SESSION_ID, ReceiveSettleMode ) @@ -18,7 +19,6 @@ SessionLockExpired ) from .utils import utc_from_timestamp, utc_now -from uamqp import Source class ReceiverMixin(object): # pylint: disable=too-many-instance-attributes @@ -52,10 +52,10 @@ def _get_source(self): return self._entity_uri def _on_attach(self, source, target, properties, error): - return + pass def _populate_message_properties(self, message): - return + pass class SessionReceiverMixin(ReceiverMixin): diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/_common/utils.py b/sdk/servicebus/azure-servicebus/azure/servicebus/_common/utils.py index 549e14a91e8c..802626c0acb0 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/_common/utils.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/_common/utils.py @@ -61,26 +61,6 @@ def utc_now(): return datetime.datetime.now(tz=TZ_UTC) -def get_running_loop(): - try: - import asyncio # pylint: disable=import-error - return asyncio.get_running_loop() - except AttributeError: # 3.5 / 3.6 - loop = None - try: - loop = asyncio._get_running_loop() # pylint: disable=protected-access - except AttributeError: - _log.warning('This version of Python is deprecated, please upgrade to >= v3.5.3') - if loop is None: - _log.warning('No running event loop') - loop = asyncio.get_event_loop() - return loop - except RuntimeError: - # For backwards compatibility, create new event loop - _log.warning('No running event loop') - return asyncio.get_event_loop() - - def parse_conn_str(conn_str): endpoint = None shared_access_key_name = None diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/_control_client/models.py b/sdk/servicebus/azure-servicebus/azure/servicebus/_control_client/models.py index 23819fc9fc63..9f0ea92812b9 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/_control_client/models.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/_control_client/models.py @@ -9,6 +9,7 @@ import sys import json from datetime import datetime +import warnings from azure.common import AzureException from ._common_models import WindowsAzureData, _unicode_type @@ -73,7 +74,6 @@ def __init__(self, default_message_time_to_live=None, @property def max_size_in_mega_bytes(self): - import warnings warnings.warn( 'This attribute has been changed to max_size_in_megabytes.') return self.max_size_in_megabytes diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_client.py b/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_client.py index a8adff9da7e8..2a24f683bac2 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_client.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_client.py @@ -402,4 +402,3 @@ def get_queue_session_receiver(self, queue_name, session_id=None, **kwargs): http_proxy=self._config.http_proxy, **kwargs ) - diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_receiver.py b/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_receiver.py index 61fa703daceb..8c17f0383e3c 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_receiver.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_receiver.py @@ -5,9 +5,9 @@ import time import logging import functools -from typing import Any, List, TYPE_CHECKING, Optional, Union +from typing import Any, List, TYPE_CHECKING, Optional -from uamqp import ReceiveClient, Source, types +from uamqp import ReceiveClient, types from uamqp.constants import SenderSettleMode from ._base_handler import BaseHandler diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_sender.py b/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_sender.py index 96464f088616..43861275d326 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_sender.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_sender.py @@ -322,7 +322,7 @@ def send(self, message): """ try: batch = self.create_batch() - batch._from_list(message) + batch._from_list(message) # pylint: disable=protected-access message = batch except TypeError: # Message was not a list or generator. pass diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_session.py b/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_session.py index 5de466887c9e..07920156aa10 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_session.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_session.py @@ -2,9 +2,8 @@ # Copyright (c) Microsoft Corporation. All rights reserved. # Licensed under the MIT License. See License.txt in the project root for license information. # -------------------------------------------------------------------------------------------- -import time import logging -from typing import Any, List, TYPE_CHECKING, Optional, Union +from typing import TYPE_CHECKING, Union import six from ._common.utils import utc_from_timestamp, utc_now diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_session_receiver.py b/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_session_receiver.py index 4a4d2bd702e9..1278978fe5f1 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_session_receiver.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_session_receiver.py @@ -2,14 +2,9 @@ # Copyright (c) Microsoft Corporation. All rights reserved. # Licensed under the MIT License. See License.txt in the project root for license information. # -------------------------------------------------------------------------------------------- -import time import logging -from typing import Any, TYPE_CHECKING, Optional, Union -import six +from typing import Any, TYPE_CHECKING -from ._base_handler import BaseHandler -from ._common.utils import utc_from_timestamp, utc_now -from ._common.constants import ReceiveSettleMode from ._common.receiver_mixins import SessionReceiverMixin from ._servicebus_receiver import ServiceBusReceiver from ._servicebus_session import ServiceBusSession diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_async_message.py b/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_async_message.py index e154e91a8bb5..f57dbccd9b5a 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_async_message.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_async_message.py @@ -16,7 +16,8 @@ MESSAGE_DEFER, MESSAGE_RENEW_LOCK ) -from .._common.utils import get_running_loop, utc_from_timestamp +from .._common.utils import utc_from_timestamp +from ._async_utils import get_running_loop from ..exceptions import MessageSettleFailed _LOGGER = logging.getLogger(__name__) @@ -94,8 +95,8 @@ async def dead_letter(self, reason=None, description=None): async def abandon(self): # type: () -> None - """Abandon the message. - + """Abandon the message. + This message will be returned to the queue and made available to be received again. :rtype: None @@ -111,7 +112,7 @@ async def abandon(self): async def defer(self): # type: () -> None """Defers the message. - + This message will remain in the queue but must be requested specifically by its sequence number in order to be received. diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_async_utils.py b/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_async_utils.py index 59ca15e83014..cfb9a7cdcdf5 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_async_utils.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_async_utils.py @@ -11,7 +11,7 @@ from uamqp import authentication -from .._common.utils import renewable_start_time, get_running_loop, utc_now +from .._common.utils import renewable_start_time, utc_now from ..exceptions import AutoLockRenewTimeout, AutoLockRenewFailed from .._common.constants import ( JWT_TOKEN_SCOPE, @@ -23,6 +23,25 @@ _log = logging.getLogger(__name__) +def get_running_loop(): + try: + return asyncio.get_running_loop() + except AttributeError: # 3.5 / 3.6 + loop = None + try: + loop = asyncio._get_running_loop() # pylint: disable=protected-access + except AttributeError: + _log.warning('This version of Python is deprecated, please upgrade to >= v3.5.3') + if loop is None: + _log.warning('No running event loop') + loop = asyncio.get_event_loop() + return loop + except RuntimeError: + # For backwards compatibility, create new event loop + _log.warning('No running event loop') + return asyncio.get_event_loop() + + async def create_authentication(client): # pylint: disable=protected-access try: diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_base_handler_async.py b/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_base_handler_async.py index a2ae3957cc98..1a04b349cdc5 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_base_handler_async.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_base_handler_async.py @@ -136,7 +136,9 @@ async def _do_retryable_operation(self, operation, timeout=None, **kwargs): ) raise last_exception - async def _mgmt_request_response(self, mgmt_operation, message, callback, keep_alive_associated_link=True, **kwargs): + async def _mgmt_request_response( + self, mgmt_operation, message, callback, keep_alive_associated_link=True, **kwargs + ): await self._open() application_properties = {} diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_servicebus_sender_async.py b/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_servicebus_sender_async.py index e672dc868ad4..11425f046501 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_servicebus_sender_async.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_servicebus_sender_async.py @@ -265,7 +265,7 @@ async def send(self, message): """ try: batch = await self.create_batch() - batch._from_list(message) + batch._from_list(message) # pylint: disable=protected-access message = batch except TypeError: # Message was not a list or generator. pass diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_servicebus_session_async.py b/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_servicebus_session_async.py index 3f424e8078d2..a1a59dc519fe 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_servicebus_session_async.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_servicebus_session_async.py @@ -3,7 +3,7 @@ # Licensed under the MIT License. See License.txt in the project root for license information. # -------------------------------------------------------------------------------------------- import logging -from typing import Any, TYPE_CHECKING, List, Union +from typing import Union import six from .._servicebus_session import ServiceBusSession as BaseSession diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_servicebus_session_receiver_async.py b/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_servicebus_session_receiver_async.py index 7e532a68ab8e..3750d1f88a65 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_servicebus_session_receiver_async.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_servicebus_session_receiver_async.py @@ -5,8 +5,7 @@ import logging from typing import Any, TYPE_CHECKING, List, Union -from .._common.receiver_mixins import ReceiverMixin, SessionReceiverMixin -from .._common.constants import ReceiveSettleMode +from .._common.receiver_mixins import SessionReceiverMixin from ._servicebus_receiver_async import ServiceBusReceiver from ._servicebus_session_async import ServiceBusSession diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/exceptions.py b/sdk/servicebus/azure-servicebus/azure/servicebus/exceptions.py index 7270154a2db9..4b41cf982d1e 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/exceptions.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/exceptions.py @@ -62,7 +62,7 @@ def _error_handler(error): return errors.ErrorAction(retry=True) -def _create_servicebus_exception(logger, exception, handler): +def _create_servicebus_exception(logger, exception, handler): # pylint: disable=too-many-statements error_need_close_handler = True error_need_raise = False if isinstance(exception, errors.MessageAlreadySettled): diff --git a/sdk/servicebus/azure-servicebus/dev_requirements.txt b/sdk/servicebus/azure-servicebus/dev_requirements.txt index f13040f4405c..6ea28753f4c8 100644 --- a/sdk/servicebus/azure-servicebus/dev_requirements.txt +++ b/sdk/servicebus/azure-servicebus/dev_requirements.txt @@ -1,3 +1,4 @@ +-e ../../core/azure-core -e ../../../tools/azure-devtools -e ../../../tools/azure-sdk-tools -e ../azure-mgmt-servicebus \ No newline at end of file From a2137500731dcca9d4831ac3ac6c43e408815c54 Mon Sep 17 00:00:00 2001 From: yijxie Date: Wed, 6 May 2020 17:29:11 -0700 Subject: [PATCH 02/15] Restructure message/async message --- .../azure/servicebus/_common/message.py | 199 +++++++++--------- .../azure/servicebus/aio/_async_message.py | 11 +- 2 files changed, 108 insertions(+), 102 deletions(-) diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/_common/message.py b/sdk/servicebus/azure-servicebus/azure/servicebus/_common/message.py index 60a4b0feef20..84d58db591d8 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/_common/message.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/_common/message.py @@ -10,7 +10,7 @@ import logging from typing import Optional, List, Union, Generator, TYPE_CHECKING -import uamqp +import uamqp.message from uamqp import types from .constants import ( @@ -432,30 +432,80 @@ def sequence_number(self): return None -class ReceivedMessage(PeekMessage): - """ - A Service Bus Message received from service side. - - :ivar auto_renew_error: Error when AutoLockRenew is used and it fails to renew the message lock. - :vartype auto_renew_error: ~azure.servicebus.AutoLockRenewTimeout or ~azure.servicebus.AutoLockRenewFailed - - .. admonition:: Example: - - .. literalinclude:: ../samples/sync_samples/sample_code_servicebus.py - :start-after: [START receive_complex_message] - :end-before: [END receive_complex_message] - :language: python - :dedent: 4 - :caption: Checking the properties on a received message. - """ +class _ReceivedMessageBase(PeekMessage): def __init__(self, message, mode=ReceiveSettleMode.PeekLock, **kwargs): - super(ReceivedMessage, self).__init__(message=message) + super(_ReceivedMessageBase, self).__init__(message=message) self._settled = (mode == ReceiveSettleMode.ReceiveAndDelete) self._is_deferred_message = kwargs.get("is_deferred_message", False) self.auto_renew_error = None self._receiver = None # type: Union[ServiceBusReceiver, ServiceBusSessionReceiver] self._expiry = None + @property + def settled(self): + # type: () -> bool + """Whether the message has been settled. + + This will aways be `True` for a message received using ReceiveAndDelete mode, + otherwise it will be `False` until the message is completed or otherwise settled. + + :rtype: bool + """ + return self._settled + + @property + def expired(self): + # type: () -> bool + """ + + :rtype: bool + """ + try: + if self._receiver.session: # pylint: disable=protected-access + raise TypeError("Session messages do not expire. Please use the Session expiry instead.") + except AttributeError: # Is not a session receiver + pass + if self.locked_until_utc and self.locked_until_utc <= utc_now(): + return True + return False + + @property + def locked_until_utc(self): + # type: () -> Optional[datetime.datetime] + """ + + :rtype: datetime.datetime + """ + try: + if self.settled or self._receiver.session: # pylint: disable=protected-access + return None + except AttributeError: # not settled, and isn't session receiver. + pass + if self._expiry: + return self._expiry + if self.message.annotations and _X_OPT_LOCKED_UNTIL in self.message.annotations: + expiry_in_seconds = self.message.annotations[_X_OPT_LOCKED_UNTIL]/1000 + self._expiry = utc_from_timestamp(expiry_in_seconds) + return self._expiry + + @property + def lock_token(self): + # type: () -> Optional[Union[uuid.UUID, str]] + """ + + :rtype: ~uuid.UUID or str + """ + if self.settled: + return None + + if self.message.delivery_tag: + return uuid.UUID(bytes_le=self.message.delivery_tag) + + delivery_annotations = self.message.delivery_annotations + if delivery_annotations: + return delivery_annotations.get(_X_OPT_LOCK_TOKEN) + return None + def _check_live(self, action): # pylint: disable=no-member if not self._receiver or not self._receiver._running: # pylint: disable=protected-access @@ -473,27 +523,6 @@ def _check_live(self, action): except AttributeError: pass - def _settle_message( - self, - settle_operation, - dead_letter_details=None - ): - try: - if not self._is_deferred_message: - try: - self._settle_via_receiver_link(settle_operation, dead_letter_details)() - return - except RuntimeError as exception: - _LOGGER.info( - "Message settling: %r has encountered an exception (%r)." - "Trying to settle through management link", - settle_operation, - exception - ) - self._settle_via_mgmt_link(settle_operation, dead_letter_details)() - except Exception as e: - raise MessageSettleFailed(settle_operation, e) - def _settle_via_mgmt_link(self, settle_operation, dead_letter_details=None): # pylint: disable=protected-access if settle_operation == MESSAGE_COMPLETE: @@ -527,6 +556,7 @@ def _settle_via_receiver_link(self, settle_operation, dead_letter_details=None): # dead_letter_detail is not used because of uamqp receiver link doesn't accept it while it # should be accepted. Will revisit this later. # uamqp management link accepts dead_letter_details. Refer to method _settle_via_mgmt_link + # TODO: to make dead_letter_details useful if settle_operation == MESSAGE_COMPLETE: return functools.partial(self.message.accept) if settle_operation == MESSAGE_ABANDON: @@ -539,70 +569,47 @@ def _settle_via_receiver_link(self, settle_operation, dead_letter_details=None): return functools.partial(self.message.modify, True, True) raise ValueError("Unsupported settle operation type: {}".format(settle_operation)) - @property - def settled(self): - # type: () -> bool - """Whether the message has been settled. - This will aways be `True` for a message received using ReceiveAndDelete mode, - otherwise it will be `False` until the message is completed or otherwise settled. +class ReceivedMessage(_ReceivedMessageBase): + """ + A Service Bus Message received from service side. - :rtype: bool - """ - return self._settled + :ivar auto_renew_error: Error when AutoLockRenew is used and it fails to renew the message lock. + :vartype auto_renew_error: ~azure.servicebus.AutoLockRenewTimeout or ~azure.servicebus.AutoLockRenewFailed - @property - def expired(self): - # type: () -> bool - """ + .. admonition:: Example: - :rtype: bool - """ - try: - if self._receiver.session: # pylint: disable=protected-access - raise TypeError("Session messages do not expire. Please use the Session expiry instead.") - except AttributeError: # Is not a session receiver - pass - if self.locked_until_utc and self.locked_until_utc <= utc_now(): - return True - return False + .. literalinclude:: ../samples/sync_samples/sample_code_servicebus.py + :start-after: [START receive_complex_message] + :end-before: [END receive_complex_message] + :language: python + :dedent: 4 + :caption: Checking the properties on a received message. + """ - @property - def locked_until_utc(self): - # type: () -> Optional[datetime.datetime] - """ + def __init__(self, message, mode=ReceiveSettleMode.PeekLock, **kwargs): + super(ReceivedMessage, self).__init__(message, mode=mode, **kwargs) - :rtype: datetime.datetime - """ + def _settle_message( + self, + settle_operation, + dead_letter_details=None + ): try: - if self.settled or self._receiver.session: # pylint: disable=protected-access - return None - except AttributeError: # not settled, and isn't session receiver. - pass - if self._expiry: - return self._expiry - if self.message.annotations and _X_OPT_LOCKED_UNTIL in self.message.annotations: - expiry_in_seconds = self.message.annotations[_X_OPT_LOCKED_UNTIL]/1000 - self._expiry = utc_from_timestamp(expiry_in_seconds) - return self._expiry - - @property - def lock_token(self): - # type: () -> Optional[Union[uuid.UUID, str]] - """ - - :rtype: ~uuid.UUID or str - """ - if self.settled: - return None - - if self.message.delivery_tag: - return uuid.UUID(bytes_le=self.message.delivery_tag) - - delivery_annotations = self.message.delivery_annotations - if delivery_annotations: - return delivery_annotations.get(_X_OPT_LOCK_TOKEN) - return None + if not self._is_deferred_message: + try: + self._settle_via_receiver_link(settle_operation, dead_letter_details)() + return + except RuntimeError as exception: + _LOGGER.info( + "Message settling: %r has encountered an exception (%r)." + "Trying to settle through management link", + settle_operation, + exception + ) + self._settle_via_mgmt_link(settle_operation, dead_letter_details)() + except Exception as e: + raise MessageSettleFailed(settle_operation, e) def complete(self): # type: () -> None diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_async_message.py b/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_async_message.py index f57dbccd9b5a..e4f1d477b117 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_async_message.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_async_message.py @@ -6,7 +6,7 @@ import logging from typing import Optional -from .._common import message as sync_message +from .._common.message import _ReceivedMessageBase from .._common.constants import ( ReceiveSettleMode, MGMT_RESPONSE_MESSAGE_EXPIRATION, @@ -23,13 +23,12 @@ _LOGGER = logging.getLogger(__name__) -class ReceivedMessage(sync_message.ReceivedMessage): +class ReceivedMessage(_ReceivedMessageBase): """A Service Bus Message received from service side. """ - def __init__(self, message, mode=ReceiveSettleMode.PeekLock, loop=None, **kwargs): - self._loop = loop or get_running_loop() + def __init__(self, message, mode=ReceiveSettleMode.PeekLock, **kwargs): super(ReceivedMessage, self).__init__(message=message, mode=mode, **kwargs) async def _settle_message( @@ -40,7 +39,7 @@ async def _settle_message( try: if not self._is_deferred_message: try: - await self._loop.run_in_executor( + await get_running_loop().run_in_executor( None, self._settle_via_receiver_link(settle_operation, dead_letter_details) ) @@ -73,7 +72,7 @@ async def complete(self): await self._settle_message(MESSAGE_COMPLETE) self._settled = True - async def dead_letter(self, reason=None, description=None): + async def dead_letter(self, reason=None, description=None): # pylint: disable=unused-argument # TODO: to use them # type: (Optional[str], Optional[str]) -> None """Move the message to the Dead Letter queue. From cf029682405733d8c531e25d6da3addb791a21c8 Mon Sep 17 00:00:00 2001 From: yijxie Date: Wed, 6 May 2020 18:13:24 -0700 Subject: [PATCH 03/15] Restructure base handler sync/async --- .../azure/servicebus/_base_handler.py | 36 ++++++++++--------- .../azure/servicebus/_servicebus_receiver.py | 4 +-- .../azure/servicebus/_servicebus_sender.py | 4 +-- .../servicebus/aio/_base_handler_async.py | 15 -------- 4 files changed, 23 insertions(+), 36 deletions(-) diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/_base_handler.py b/sdk/servicebus/azure-servicebus/azure/servicebus/_base_handler.py index 24869d901162..4bf599b4b146 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/_base_handler.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/_base_handler.py @@ -110,7 +110,7 @@ def get_token(self, *scopes, **kwargs): # pylint:disable=unused-argument return _generate_sas_token(scopes[0], self.policy, self.key) -class BaseHandler(object): # pylint:disable=too-many-instance-attributes +class BaseHandler(object): def __init__( self, fully_qualified_namespace, @@ -132,22 +132,6 @@ def __init__( self._auth_uri = None self._properties = create_properties() - def __enter__(self): - self._open_with_retry() - return self - - def __exit__(self, *args): - self.close() - - def _handle_exception(self, exception): - error, error_need_close_handler, error_need_raise = _create_servicebus_exception(_LOGGER, exception, self) - if error_need_close_handler: - self._close_handler() - if error_need_raise: - raise error - - return error - @staticmethod def _from_connection_string(conn_str, **kwargs): # type: (str, Any) -> Dict[str, Any] @@ -173,6 +157,24 @@ def _from_connection_string(conn_str, **kwargs): kwargs["credential"] = ServiceBusSharedKeyCredential(policy, key) return kwargs + +class BaseHandlerSync(BaseHandler): # pylint:disable=too-many-instance-attributes + def __enter__(self): + self._open_with_retry() + return self + + def __exit__(self, *args): + self.close() + + def _handle_exception(self, exception): + error, error_need_close_handler, error_need_raise = _create_servicebus_exception(_LOGGER, exception, self) + if error_need_close_handler: + self._close_handler() + if error_need_raise: + raise error + + return error + def _backoff( self, retried_times, diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_receiver.py b/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_receiver.py index 8c17f0383e3c..9d01a89e13d6 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_receiver.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_receiver.py @@ -10,7 +10,7 @@ from uamqp import ReceiveClient, types from uamqp.constants import SenderSettleMode -from ._base_handler import BaseHandler +from ._base_handler import BaseHandlerSync from ._common.utils import create_authentication from ._common.message import PeekMessage, ReceivedMessage from ._common.constants import ( @@ -36,7 +36,7 @@ _LOGGER = logging.getLogger(__name__) -class ServiceBusReceiver(BaseHandler, ReceiverMixin): # pylint: disable=too-many-instance-attributes +class ServiceBusReceiver(BaseHandlerSync, ReceiverMixin): # pylint: disable=too-many-instance-attributes """The ServiceBusReceiver class defines a high level interface for receiving messages from the Azure Service Bus Queue or Topic Subscription. diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_sender.py b/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_sender.py index 43861275d326..e53e15f7fad7 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_sender.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_sender.py @@ -10,7 +10,7 @@ import uamqp from uamqp import SendClient, types -from ._base_handler import BaseHandler +from ._base_handler import BaseHandlerSync from ._common import mgmt_handlers from ._common.message import Message, BatchMessage from .exceptions import ( @@ -81,7 +81,7 @@ def _build_schedule_request(cls, schedule_time_utc, *messages): return request_body -class ServiceBusSender(BaseHandler, SenderMixin): +class ServiceBusSender(BaseHandlerSync, SenderMixin): """The ServiceBusSender class defines a high level interface for sending messages to the Azure Service Bus Queue or Topic. diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_base_handler_async.py b/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_base_handler_async.py index 1a04b349cdc5..fe186c470eaa 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_base_handler_async.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_base_handler_async.py @@ -45,21 +45,6 @@ async def get_token(self, *scopes, **kwargs): # pylint:disable=unused-argument class BaseHandlerAsync(BaseHandler): - def __init__( - self, - fully_qualified_namespace: str, - entity_name: str, - credential: "TokenCredential", - **kwargs: Any - ) -> None: - self._loop = kwargs.pop("loop", None) - super(BaseHandlerAsync, self).__init__( - fully_qualified_namespace=fully_qualified_namespace, - entity_name=entity_name, - credential=credential, - **kwargs - ) - async def __aenter__(self): await self._open_with_retry() return self From 361a8833f3a2e6946998c6b0a4023815b672bbb4 Mon Sep 17 00:00:00 2001 From: yijxie Date: Wed, 6 May 2020 19:42:35 -0700 Subject: [PATCH 04/15] Restructure sb session sync/async --- .../azure/servicebus/_servicebus_session.py | 83 ++++++++++--------- .../aio/_servicebus_session_async.py | 2 +- 2 files changed, 44 insertions(+), 41 deletions(-) diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_session.py b/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_session.py index 07920156aa10..564e14bb3e1d 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_session.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_session.py @@ -25,7 +25,49 @@ _LOGGER = logging.getLogger(__name__) -class ServiceBusSession(object): +class BaseSession(object): + def __init__(self, session_id, receiver, encoding="UTF-8"): + self._session_id = session_id + self._receiver = receiver + self._encoding = encoding + self._session_start = None + self._locked_until_utc = None + self.auto_renew_error = None + + @property + def session_id(self): + # type: () -> str + """ + Session id of the current session. + + :rtype: str + """ + return self._session_id + + @property + def expired(self): + # type: () -> bool + """Whether the receivers lock on a particular session has expired. + + :rtype: bool + """ + return bool(self._locked_until_utc and self._locked_until_utc <= utc_now()) + + @property + def locked_until_utc(self): + # type: () -> datetime.datetime + """The time at which this session's lock will expire. + + :rtype: datetime.datetime + """ + return self._locked_until_utc + + def _check_live(self): + if self.expired: + raise SessionLockExpired(inner_exception=self.auto_renew_error) + + +class ServiceBusSession(BaseSession): """ The ServiceBusSession is used for manage session states and lock renewal. @@ -44,17 +86,6 @@ class ServiceBusSession(object): :dedent: 4 :caption: Get session from a receiver """ - def __init__(self, session_id, receiver, encoding="UTF-8"): - self._session_id = session_id - self._receiver = receiver - self._encoding = encoding - self._session_start = None - self._locked_until_utc = None - self.auto_renew_error = None - - def _check_live(self): - if self.expired: - raise SessionLockExpired(inner_exception=self.auto_renew_error) def get_session_state(self): # type: () -> str @@ -134,31 +165,3 @@ def renew_lock(self): mgmt_handlers.default ) self._locked_until_utc = utc_from_timestamp(expiry[MGMT_RESPONSE_RECEIVER_EXPIRATION]/1000.0) - - @property - def session_id(self): - # type: () -> str - """ - Session id of the current session. - - :rtype: str - """ - return self._session_id - - @property - def expired(self): - # type: () -> bool - """Whether the receivers lock on a particular session has expired. - - :rtype: bool - """ - return bool(self._locked_until_utc and self._locked_until_utc <= utc_now()) - - @property - def locked_until_utc(self): - # type: () -> datetime.datetime - """The time at which this session's lock will expire. - - :rtype: datetime.datetime - """ - return self._locked_until_utc diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_servicebus_session_async.py b/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_servicebus_session_async.py index a1a59dc519fe..b0704782fbe6 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_servicebus_session_async.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_servicebus_session_async.py @@ -6,7 +6,7 @@ from typing import Union import six -from .._servicebus_session import ServiceBusSession as BaseSession +from .._servicebus_session import BaseSession from .._common.constants import ( REQUEST_RESPONSE_GET_SESSION_STATE_OPERATION, REQUEST_RESPONSE_SET_SESSION_STATE_OPERATION, From dc84d6c37906723a91454376de0dba6083903351 Mon Sep 17 00:00:00 2001 From: yijxie Date: Wed, 6 May 2020 19:44:05 -0700 Subject: [PATCH 05/15] Enable sb pylint --- eng/tox/allowed_pylint_failures.py | 1 - 1 file changed, 1 deletion(-) diff --git a/eng/tox/allowed_pylint_failures.py b/eng/tox/allowed_pylint_failures.py index 3982eba05561..2433975c377e 100644 --- a/eng/tox/allowed_pylint_failures.py +++ b/eng/tox/allowed_pylint_failures.py @@ -39,7 +39,6 @@ "azure-eventgrid", "azure-graphrbac", "azure-loganalytics", - "azure-servicebus", "azure-servicefabric", "azure-template", "azure-keyvault", From 848b154d77151fba7c37e52a6482205810f8a5ff Mon Sep 17 00:00:00 2001 From: yijxie Date: Thu, 7 May 2020 13:24:37 -0700 Subject: [PATCH 06/15] Fix mypy errors --- sdk/servicebus/azure-servicebus/azure/__init__.py | 2 +- .../azure/servicebus/_common/client_mixins.py | 4 ++-- .../azure-servicebus/azure/servicebus/_common/message.py | 8 ++++---- .../azure/servicebus/_control_client/_common_models.py | 4 ++-- .../servicebus/_control_client/_common_serialization.py | 2 +- .../azure/servicebus/_control_client/servicebusservice.py | 4 +++- .../azure/servicebus/_servicebus_client.py | 4 ++-- .../azure/servicebus/_servicebus_session_receiver.py | 3 ++- .../azure/servicebus/aio/_servicebus_client_async.py | 2 +- .../servicebus/aio/_servicebus_session_receiver_async.py | 5 +++-- 10 files changed, 21 insertions(+), 17 deletions(-) diff --git a/sdk/servicebus/azure-servicebus/azure/__init__.py b/sdk/servicebus/azure-servicebus/azure/__init__.py index 625f4623bd06..e604558a8bed 100644 --- a/sdk/servicebus/azure-servicebus/azure/__init__.py +++ b/sdk/servicebus/azure-servicebus/azure/__init__.py @@ -1 +1 @@ -__path__ = __import__('pkgutil').extend_path(__path__, __name__) +__path__ = __import__('pkgutil').extend_path(__path__, __name__) # type: ignore diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/_common/client_mixins.py b/sdk/servicebus/azure-servicebus/azure/servicebus/_common/client_mixins.py index f492cc74b69d..4056f10388c5 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/_common/client_mixins.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/_common/client_mixins.py @@ -8,8 +8,8 @@ import uuid import requests try: - from urlparse import urlparse - from urllib import unquote_plus + from urlparse import urlparse # type: ignore + from urllib import unquote_plus # type: ignore except ImportError: from urllib.parse import urlparse from urllib.parse import unquote_plus diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/_common/message.py b/sdk/servicebus/azure-servicebus/azure/servicebus/_common/message.py index 84d58db591d8..2e262dc24e5b 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/_common/message.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/_common/message.py @@ -8,7 +8,7 @@ import uuid import functools import logging -from typing import Optional, List, Union, Generator, TYPE_CHECKING +from typing import Optional, List, Union, Iterable, TYPE_CHECKING import uamqp.message from uamqp import types @@ -269,10 +269,10 @@ def scheduled_enqueue_time_utc(self, value): @property def body(self): - # type: () -> Union[bytes, Generator[bytes]] + # type: () -> Union[bytes, Iterable[bytes]] """The body of the Message. - :rtype: bytes or generator[bytes] + :rtype: bytes or Iterable[bytes] """ return self.message.get_data() @@ -438,7 +438,7 @@ def __init__(self, message, mode=ReceiveSettleMode.PeekLock, **kwargs): self._settled = (mode == ReceiveSettleMode.ReceiveAndDelete) self._is_deferred_message = kwargs.get("is_deferred_message", False) self.auto_renew_error = None - self._receiver = None # type: Union[ServiceBusReceiver, ServiceBusSessionReceiver] + self._receiver = None # type: ignore self._expiry = None @property diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/_control_client/_common_models.py b/sdk/servicebus/azure-servicebus/azure/servicebus/_control_client/_common_models.py index 624052524ca4..a1f3ede58da7 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/_control_client/_common_models.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/_control_client/_common_models.py @@ -79,8 +79,8 @@ def __init__(self, xml_element_name): try: - _unicode_type = unicode - _strtype = basestring + _unicode_type = unicode # type: ignore + _strtype = basestring # type: ignore except NameError: _unicode_type = str _strtype = str diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/_control_client/_common_serialization.py b/sdk/servicebus/azure-servicebus/azure/servicebus/_control_client/_common_serialization.py index 33e72795e8e4..f2daaf3cf3ef 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/_control_client/_common_serialization.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/_control_client/_common_serialization.py @@ -9,7 +9,7 @@ try: from xml.etree import cElementTree as ETree except ImportError: - from xml.etree import ElementTree as ETree + from xml.etree import ElementTree as ETree # type: ignore try: from cStringIO import StringIO except ImportError: diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/_control_client/servicebusservice.py b/sdk/servicebus/azure-servicebus/azure/servicebus/_control_client/servicebusservice.py index 6adedb647ac9..1b635c602fd4 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/_control_client/servicebusservice.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/_control_client/servicebusservice.py @@ -9,6 +9,8 @@ import os import time import json +from typing import Dict + try: from urllib2 import quote as url_quote from urllib2 import unquote as url_unquote @@ -1252,7 +1254,7 @@ def _update_service_bus_header(self, request): # Token cache for Authentication # Shared by the different instances of ServiceBusWrapTokenAuthentication -_tokens = {} +_tokens = {} # type: Dict[str, str] class ServiceBusWrapTokenAuthentication: diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_client.py b/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_client.py index 2a24f683bac2..d8d79494d0e9 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_client.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_client.py @@ -131,7 +131,7 @@ def from_connection_string( return cls( fully_qualified_namespace=host, entity_name=entity_in_conn_str or kwargs.pop("entity_name", None), - credential=ServiceBusSharedKeyCredential(policy, key), + credential=ServiceBusSharedKeyCredential(policy, key), # type: ignore **kwargs ) @@ -298,7 +298,7 @@ def get_subscription_receiver(self, topic_name, subscription_name, **kwargs): ) def get_subscription_session_receiver(self, topic_name, subscription_name, session_id=None, **kwargs): - # type: (str, str, Any) -> ServiceBusReceiver + # type: (str, str, str, Any) -> ServiceBusReceiver """Get ServiceBusReceiver for the specific subscription under the topic. :param str topic_name: The name of specific Service Bus Topic the client connects to. diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_session_receiver.py b/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_session_receiver.py index 1278978fe5f1..2891e02f4b0b 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_session_receiver.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_session_receiver.py @@ -150,4 +150,5 @@ def from_connection_string( :caption: Create a new instance of the ServiceBusReceiver from connection string. """ - return super(ServiceBusSessionReceiver, cls).from_connection_string(conn_str, **kwargs) + constructor_args = super(ServiceBusSessionReceiver, cls)._from_connection_string(conn_str, **kwargs) + return cls(**constructor_args) diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_servicebus_client_async.py b/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_servicebus_client_async.py index b5e5f4029454..fc2b40c7db0a 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_servicebus_client_async.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_servicebus_client_async.py @@ -298,7 +298,7 @@ def get_subscription_receiver(self, topic_name, subscription_name, **kwargs): ) def get_subscription_session_receiver(self, topic_name, subscription_name, session_id=None, **kwargs): - # type: (str, str, Any) -> ServiceBusReceiver + # type: (str, str, str, Any) -> ServiceBusReceiver """Get ServiceBusReceiver for the specific subscription under the topic. :param str topic_name: The name of specific Service Bus Topic the client connects to. diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_servicebus_session_receiver_async.py b/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_servicebus_session_receiver_async.py index 3750d1f88a65..e7d4a57ae9c4 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_servicebus_session_receiver_async.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_servicebus_session_receiver_async.py @@ -87,7 +87,7 @@ def from_connection_string( cls, conn_str: str, **kwargs: Any - ) -> "ServiceBusReceiver": + ) -> "ServiceBusSessionReceiver": """Create a ServiceBusSessionReceiver from a connection string. :param conn_str: The connection string of a Service Bus. @@ -133,7 +133,8 @@ def from_connection_string( :caption: Create a new instance of the ServiceBusReceiver from connection string. """ - return super(ServiceBusSessionReceiver, cls).from_connection_string(conn_str, **kwargs) + constructor_args = super(ServiceBusSessionReceiver, cls)._from_connection_string(conn_str, **kwargs) + return cls(**constructor_args) @property From abf2696e0264090dbd527741d759e9b8ea28db2b Mon Sep 17 00:00:00 2001 From: yijxie Date: Thu, 7 May 2020 13:24:54 -0700 Subject: [PATCH 07/15] Enable mypy check for ServiceBus --- eng/tox/mypy_hard_failure_packages.py | 1 + 1 file changed, 1 insertion(+) diff --git a/eng/tox/mypy_hard_failure_packages.py b/eng/tox/mypy_hard_failure_packages.py index 4d9932a98916..4f458daf44aa 100644 --- a/eng/tox/mypy_hard_failure_packages.py +++ b/eng/tox/mypy_hard_failure_packages.py @@ -8,6 +8,7 @@ MYPY_HARD_FAILURE_OPTED = [ "azure-core", "azure-eventhub", + "azure-servicebus", "azure-ai-textanalytics", "azure-ai-formrecognizer" ] From 2bbd20c8e6dc2056db9df1c960057a753482b330 Mon Sep 17 00:00:00 2001 From: yijxie Date: Mon, 11 May 2020 16:10:46 -0700 Subject: [PATCH 08/15] Add type hints for internal methods --- .../azure/servicebus/_base_handler.py | 29 ++++++++++--------- .../azure/servicebus/_common/message.py | 8 ++--- .../azure/servicebus/_servicebus_receiver.py | 16 ++++++---- .../azure/servicebus/_servicebus_sender.py | 9 ++++-- .../azure/servicebus/_servicebus_session.py | 9 ++++-- .../_servicebus_session_receiver.py | 3 +- .../azure/servicebus/aio/_async_message.py | 4 --- .../aio/_servicebus_sender_async.py | 2 +- .../aio/_servicebus_session_receiver_async.py | 4 +-- 9 files changed, 45 insertions(+), 39 deletions(-) diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/_base_handler.py b/sdk/servicebus/azure-servicebus/azure/servicebus/_base_handler.py index 4bf599b4b146..1271f0b6ca63 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/_base_handler.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/_base_handler.py @@ -7,7 +7,7 @@ import uuid import time from datetime import timedelta -from typing import cast, Optional, Tuple, TYPE_CHECKING, Dict, Any +from typing import cast, Optional, Tuple, TYPE_CHECKING, Dict, Any, Callable try: from urllib import quote_plus # type: ignore @@ -118,6 +118,7 @@ def __init__( credential, **kwargs ): + # type: (str, str, TokenCredential, Any) -> None self.fully_qualified_namespace = fully_qualified_namespace self._entity_name = entity_name @@ -128,7 +129,7 @@ def __init__( self._container_id = CONTAINER_PREFIX + str(uuid.uuid4())[:8] self._config = Configuration(**kwargs) self._running = False - self._handler = None + self._handler = None # type: uamqp.AMQPClient self._auth_uri = None self._properties = create_properties() @@ -167,6 +168,7 @@ def __exit__(self, *args): self.close() def _handle_exception(self, exception): + # type: (BaseException) -> ServiceBusError error, error_need_close_handler, error_need_raise = _create_servicebus_exception(_LOGGER, exception, self) if error_need_close_handler: self._close_handler() @@ -182,6 +184,7 @@ def _backoff( timeout=None, entity_name=None ): + # type: (int, Exception, Optional[float], str) -> None entity_name = entity_name or self._container_id backoff = self._config.retry_backoff_factor * 2 ** retried_times if backoff <= self._config.retry_backoff_max and ( @@ -202,16 +205,14 @@ def _backoff( raise last_exception def _do_retryable_operation(self, operation, timeout=None, **kwargs): + # type: (Callable, Optional[float], Any) -> Any require_last_exception = kwargs.pop("require_last_exception", False) require_timeout = kwargs.pop("require_timeout", False) retried_times = 0 - last_exception = None max_retries = self._config.retry_total while retried_times <= max_retries: try: - if require_last_exception: - kwargs["last_exception"] = last_exception if require_timeout: kwargs["timeout"] = timeout return operation(**kwargs) @@ -219,23 +220,24 @@ def _do_retryable_operation(self, operation, timeout=None, **kwargs): raise except Exception as exception: # pylint: disable=broad-except last_exception = self._handle_exception(exception) + if require_last_exception: + kwargs["last_exception"] = last_exception retried_times += 1 if retried_times > max_retries: - break + _LOGGER.info( + "%r operation has exhausted retry. Last exception: %r.", + self._container_id, + last_exception, + ) + raise last_exception self._backoff( retried_times=retried_times, last_exception=last_exception, timeout=timeout ) - _LOGGER.info( - "%r operation has exhausted retry. Last exception: %r.", - self._container_id, - last_exception, - ) - raise last_exception - def _mgmt_request_response(self, mgmt_operation, message, callback, keep_alive_associated_link=True, **kwargs): + # type: (str, uamqp.Message, Callable, bool, Any) -> uamqp.Message self._open() application_properties = {} # Some mgmt calls do not support an associated link name (such as list_sessions). Most do, so on by default. @@ -267,6 +269,7 @@ def _mgmt_request_response(self, mgmt_operation, message, callback, keep_alive_a raise ServiceBusError("Management request failed: {}".format(exp), exp) def _mgmt_request_response_with_retry(self, mgmt_operation, message, callback, **kwargs): + # type: (bytes, Dict[str, Any], Callable, Any) -> Any return self._do_retryable_operation( self._mgmt_request_response, mgmt_operation=mgmt_operation, diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/_common/message.py b/sdk/servicebus/azure-servicebus/azure/servicebus/_common/message.py index 2e262dc24e5b..727780323404 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/_common/message.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/_common/message.py @@ -8,7 +8,7 @@ import uuid import functools import logging -from typing import Optional, List, Union, Iterable, TYPE_CHECKING +from typing import Optional, List, Union, Iterable, TYPE_CHECKING, Callable, Dict, Any import uamqp.message from uamqp import types @@ -524,6 +524,7 @@ def _check_live(self, action): pass def _settle_via_mgmt_link(self, settle_operation, dead_letter_details=None): + # type: (str, Dict[str, Any]) -> Callable # pylint: disable=protected-access if settle_operation == MESSAGE_COMPLETE: return functools.partial( @@ -553,6 +554,7 @@ def _settle_via_mgmt_link(self, settle_operation, dead_letter_details=None): raise ValueError("Unsupported settle operation type: {}".format(settle_operation)) def _settle_via_receiver_link(self, settle_operation, dead_letter_details=None): # pylint: disable=unused-argument + # type: (str, Dict[str, Any]) -> Callable # dead_letter_detail is not used because of uamqp receiver link doesn't accept it while it # should be accepted. Will revisit this later. # uamqp management link accepts dead_letter_details. Refer to method _settle_via_mgmt_link @@ -587,14 +589,12 @@ class ReceivedMessage(_ReceivedMessageBase): :caption: Checking the properties on a received message. """ - def __init__(self, message, mode=ReceiveSettleMode.PeekLock, **kwargs): - super(ReceivedMessage, self).__init__(message, mode=mode, **kwargs) - def _settle_message( self, settle_operation, dead_letter_details=None ): + # type: (str, Dict[str, Any]) -> None try: if not self._is_deferred_message: try: diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_receiver.py b/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_receiver.py index 9d01a89e13d6..600c3a6891a5 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_receiver.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_receiver.py @@ -5,10 +5,11 @@ import time import logging import functools -from typing import Any, List, TYPE_CHECKING, Optional +from typing import Any, List, TYPE_CHECKING, Optional, Dict from uamqp import ReceiveClient, types from uamqp.constants import SenderSettleMode +from uamqp.authentication.common import AMQPAuth from ._base_handler import BaseHandlerSync from ._common.utils import create_authentication @@ -103,17 +104,16 @@ def __init__( **kwargs ) else: - queue_name = kwargs.get("queue_name") - topic_name = kwargs.get("topic_name") + queue_name = kwargs.get("queue_name") # type: Optional[str] + topic_name = kwargs.get("topic_name") # type: Optional[str] subscription_name = kwargs.get("subscription_name") if queue_name and topic_name: raise ValueError("Queue/Topic name can not be specified simultaneously.") - if not (queue_name or topic_name): - raise ValueError("Queue/Topic name is missing. Please specify queue_name/topic_name.") if topic_name and not subscription_name: raise ValueError("Subscription name is missing for the topic. Please specify subscription_name.") - entity_name = queue_name or topic_name + if not entity_name: + raise ValueError("Queue/Topic name is missing. Please specify queue_name/topic_name.") super(ServiceBusReceiver, self).__init__( fully_qualified_namespace=fully_qualified_namespace, @@ -147,6 +147,7 @@ def _iter_next(self): return message def _create_handler(self, auth): + # type: (AMQPAuth) -> None self._handler = ReceiveClient( self._get_source(), auth=auth, @@ -182,6 +183,7 @@ def _open(self): raise def _receive(self, max_batch_size=None, timeout=None): + # type: (Optional[int], Optional[float]) -> List[ReceivedMessage] self._open() max_batch_size = max_batch_size or self._handler._prefetch # pylint: disable=protected-access @@ -194,6 +196,7 @@ def _receive(self, max_batch_size=None, timeout=None): return [self._build_message(message) for message in batch] def _settle_message(self, settlement, lock_tokens, dead_letter_details=None): + # type: (bytes, List[str], Optional[Dict[str, Any]]) -> Any message = { MGMT_REQUEST_DISPOSITION_STATUS: settlement, MGMT_REQUEST_LOCK_TOKENS: types.AMQPArray(lock_tokens) @@ -210,6 +213,7 @@ def _settle_message(self, settlement, lock_tokens, dead_letter_details=None): ) def _renew_locks(self, *lock_tokens): + # type: (*str) -> Any message = {MGMT_REQUEST_LOCK_TOKENS: types.AMQPArray(lock_tokens)} return self._mgmt_request_response_with_retry( REQUEST_RESPONSE_RENEWLOCK_OPERATION, diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_sender.py b/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_sender.py index e53e15f7fad7..8cb8eda85ca0 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_sender.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_sender.py @@ -5,10 +5,11 @@ import logging import time import uuid -from typing import Any, TYPE_CHECKING, Union, List +from typing import Any, TYPE_CHECKING, Union, List, Optional import uamqp from uamqp import SendClient, types +from uamqp.authentication.common import AMQPAuth from ._base_handler import BaseHandlerSync from ._common import mgmt_handlers @@ -137,9 +138,9 @@ def __init__( topic_name = kwargs.get("topic_name") if queue_name and topic_name: raise ValueError("Queue/Topic name can not be specified simultaneously.") - if not (queue_name or topic_name): - raise ValueError("Queue/Topic name is missing. Please specify queue_name/topic_name.") entity_name = queue_name or topic_name + if not entity_name: + raise ValueError("Queue/Topic name is missing. Please specify queue_name/topic_name.") super(ServiceBusSender, self).__init__( fully_qualified_namespace=fully_qualified_namespace, credential=credential, @@ -152,6 +153,7 @@ def __init__( self._connection = kwargs.get("connection") def _create_handler(self, auth): + # type: (AMQPAuth) -> None self._handler = SendClient( self._entity_uri, auth=auth, @@ -183,6 +185,7 @@ def _open(self): raise def _send(self, message, timeout=None, last_exception=None): + # type: (Message, Optional[float], Exception) -> None self._open() self._set_msg_timeout(timeout, last_exception) self._handler.send_message(message.message) diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_session.py b/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_session.py index 564e14bb3e1d..ae3478169dc4 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_session.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_session.py @@ -3,7 +3,7 @@ # Licensed under the MIT License. See License.txt in the project root for license information. # -------------------------------------------------------------------------------------------- import logging -from typing import TYPE_CHECKING, Union +from typing import TYPE_CHECKING, Union, Optional import six from ._common.utils import utc_from_timestamp, utc_now @@ -21,17 +21,20 @@ if TYPE_CHECKING: import datetime + from ._servicebus_session_receiver import ServiceBusSessionReceiver + from .aio._servicebus_session_receiver_async import ServiceBusSessionReceiver as ServiceBusSessionReceiverAsync _LOGGER = logging.getLogger(__name__) class BaseSession(object): def __init__(self, session_id, receiver, encoding="UTF-8"): + # type: (str, Union[ServiceBusSessionReceiver, ServiceBusSessionReceiverAsync], str) -> None self._session_id = session_id self._receiver = receiver self._encoding = encoding self._session_start = None - self._locked_until_utc = None + self._locked_until_utc = None # type: Optional[datetime.datetime] self.auto_renew_error = None @property @@ -55,7 +58,7 @@ def expired(self): @property def locked_until_utc(self): - # type: () -> datetime.datetime + # type: () -> Optional[datetime.datetime] """The time at which this session's lock will expire. :rtype: datetime.datetime diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_session_receiver.py b/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_session_receiver.py index 2891e02f4b0b..cb9dd586ee71 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_session_receiver.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_session_receiver.py @@ -150,5 +150,4 @@ def from_connection_string( :caption: Create a new instance of the ServiceBusReceiver from connection string. """ - constructor_args = super(ServiceBusSessionReceiver, cls)._from_connection_string(conn_str, **kwargs) - return cls(**constructor_args) + return super(ServiceBusSessionReceiver, cls).from_connection_string(conn_str, **kwargs) # type: ignore diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_async_message.py b/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_async_message.py index e4f1d477b117..5c1d3f36686e 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_async_message.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_async_message.py @@ -8,7 +8,6 @@ from .._common.message import _ReceivedMessageBase from .._common.constants import ( - ReceiveSettleMode, MGMT_RESPONSE_MESSAGE_EXPIRATION, MESSAGE_COMPLETE, MESSAGE_DEAD_LETTER, @@ -28,9 +27,6 @@ class ReceivedMessage(_ReceivedMessageBase): """ - def __init__(self, message, mode=ReceiveSettleMode.PeekLock, **kwargs): - super(ReceivedMessage, self).__init__(message=message, mode=mode, **kwargs) - async def _settle_message( self, settle_operation, diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_servicebus_sender_async.py b/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_servicebus_sender_async.py index 11425f046501..039dcd00aa4a 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_servicebus_sender_async.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_servicebus_sender_async.py @@ -269,7 +269,7 @@ async def send(self, message): message = batch except TypeError: # Message was not a list or generator. pass - if isinstance(message, BatchMessage) and len(message) == 0: + if isinstance(message, BatchMessage) and len(message) == 0: # pylint: disable=len-as-condition raise ValueError("A BatchMessage or list of Message must have at least one Message") await self._do_retryable_operation( diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_servicebus_session_receiver_async.py b/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_servicebus_session_receiver_async.py index e7d4a57ae9c4..6d8f496b975d 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_servicebus_session_receiver_async.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_servicebus_session_receiver_async.py @@ -133,9 +133,7 @@ def from_connection_string( :caption: Create a new instance of the ServiceBusReceiver from connection string. """ - constructor_args = super(ServiceBusSessionReceiver, cls)._from_connection_string(conn_str, **kwargs) - return cls(**constructor_args) - + return super(ServiceBusSessionReceiver, cls).from_connection_string(conn_str, **kwargs) # type: ignore @property def session(self): From a86a85487a6ded069a5c110cfba6fbce4e318ba0 Mon Sep 17 00:00:00 2001 From: yijxie Date: Mon, 11 May 2020 16:56:00 -0700 Subject: [PATCH 09/15] Fix a pylint disable comment --- .../azure-servicebus/azure/servicebus/_common/message.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/_common/message.py b/sdk/servicebus/azure-servicebus/azure/servicebus/_common/message.py index 727780323404..d7c1033dfe35 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/_common/message.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/_common/message.py @@ -714,5 +714,5 @@ def renew_lock(self): if not token: raise ValueError("Unable to renew lock - no lock token found.") - expiry = self._receiver._renew_locks(token) # pylint: disable=protected-access no-member + expiry = self._receiver._renew_locks(token) # pylint: disable=protected-access,no-member self._expiry = utc_from_timestamp(expiry[MGMT_RESPONSE_MESSAGE_EXPIRATION][0]/1000.0) From b819845a745ae213649c9738d16c990a31cfc8a7 Mon Sep 17 00:00:00 2001 From: yijxie Date: Mon, 11 May 2020 17:14:56 -0700 Subject: [PATCH 10/15] Fix a pylint error --- .../azure-servicebus/azure/servicebus/_servicebus_sender.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_sender.py b/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_sender.py index 8cb8eda85ca0..63dc8396374a 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_sender.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_sender.py @@ -329,7 +329,7 @@ def send(self, message): message = batch except TypeError: # Message was not a list or generator. pass - if isinstance(message, BatchMessage) and len(message) == 0: + if isinstance(message, BatchMessage) and len(message) == 0: # pylint: disable=len-as-condition raise ValueError("A BatchMessage or list of Message must have at least one Message") self._do_retryable_operation( From 344cbefd7171ef0943a207513f74651e2a2a9a24 Mon Sep 17 00:00:00 2001 From: yijxie Date: Wed, 13 May 2020 18:36:43 -0700 Subject: [PATCH 11/15] async ReceivedMessage extends sync ReceivedMessage directly --- .../azure/servicebus/_common/message.py | 38 +++++++++---------- .../azure/servicebus/aio/_async_message.py | 24 +++++------- 2 files changed, 28 insertions(+), 34 deletions(-) diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/_common/message.py b/sdk/servicebus/azure-servicebus/azure/servicebus/_common/message.py index d7c1033dfe35..22e47d847534 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/_common/message.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/_common/message.py @@ -432,9 +432,25 @@ def sequence_number(self): return None -class _ReceivedMessageBase(PeekMessage): +class ReceivedMessage(PeekMessage): + """ + A Service Bus Message received from service side. + + :ivar auto_renew_error: Error when AutoLockRenew is used and it fails to renew the message lock. + :vartype auto_renew_error: ~azure.servicebus.AutoLockRenewTimeout or ~azure.servicebus.AutoLockRenewFailed + + .. admonition:: Example: + + .. literalinclude:: ../samples/sync_samples/sample_code_servicebus.py + :start-after: [START receive_complex_message] + :end-before: [END receive_complex_message] + :language: python + :dedent: 4 + :caption: Checking the properties on a received message. + """ + def __init__(self, message, mode=ReceiveSettleMode.PeekLock, **kwargs): - super(_ReceivedMessageBase, self).__init__(message=message) + super(ReceivedMessage, self).__init__(message=message) self._settled = (mode == ReceiveSettleMode.ReceiveAndDelete) self._is_deferred_message = kwargs.get("is_deferred_message", False) self.auto_renew_error = None @@ -571,24 +587,6 @@ def _settle_via_receiver_link(self, settle_operation, dead_letter_details=None): return functools.partial(self.message.modify, True, True) raise ValueError("Unsupported settle operation type: {}".format(settle_operation)) - -class ReceivedMessage(_ReceivedMessageBase): - """ - A Service Bus Message received from service side. - - :ivar auto_renew_error: Error when AutoLockRenew is used and it fails to renew the message lock. - :vartype auto_renew_error: ~azure.servicebus.AutoLockRenewTimeout or ~azure.servicebus.AutoLockRenewFailed - - .. admonition:: Example: - - .. literalinclude:: ../samples/sync_samples/sample_code_servicebus.py - :start-after: [START receive_complex_message] - :end-before: [END receive_complex_message] - :language: python - :dedent: 4 - :caption: Checking the properties on a received message. - """ - def _settle_message( self, settle_operation, diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_async_message.py b/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_async_message.py index 5c1d3f36686e..286b9ac21f46 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_async_message.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_async_message.py @@ -6,7 +6,7 @@ import logging from typing import Optional -from .._common.message import _ReceivedMessageBase +from .._common import message as sync_message from .._common.constants import ( MGMT_RESPONSE_MESSAGE_EXPIRATION, MESSAGE_COMPLETE, @@ -22,12 +22,11 @@ _LOGGER = logging.getLogger(__name__) -class ReceivedMessage(_ReceivedMessageBase): +class ReceivedMessage(sync_message.ReceivedMessage): """A Service Bus Message received from service side. """ - - async def _settle_message( + async def _settle_message( # type: ignore # pylint: disable=invalid-overridden-method self, settle_operation, dead_letter_details=None @@ -51,8 +50,7 @@ async def _settle_message( except Exception as e: raise MessageSettleFailed(settle_operation, e) - async def complete(self): - # type: () -> None + async def complete(self) -> None: # type: ignore # pylint: disable=invalid-overridden-method """Complete the message. This removes the message from the queue. @@ -68,8 +66,9 @@ async def complete(self): await self._settle_message(MESSAGE_COMPLETE) self._settled = True - async def dead_letter(self, reason=None, description=None): # pylint: disable=unused-argument # TODO: to use them - # type: (Optional[str], Optional[str]) -> None + async def dead_letter( # type: ignore + self, reason: Optional[str] = None, description: Optional[str] = None + ) -> None: # pylint: disable=unused-argument,invalid-overridden-method """Move the message to the Dead Letter queue. The Dead Letter queue is a sub-queue that can be @@ -88,8 +87,7 @@ async def dead_letter(self, reason=None, description=None): # pylint: disable=u await self._settle_message(MESSAGE_DEAD_LETTER) self._settled = True - async def abandon(self): - # type: () -> None + async def abandon(self) -> None: # type: ignore # pylint: disable=invalid-overridden-method """Abandon the message. This message will be returned to the queue and made available to be received again. @@ -104,8 +102,7 @@ async def abandon(self): await self._settle_message(MESSAGE_ABANDON) self._settled = True - async def defer(self): - # type: () -> None + async def defer(self) -> None: # type: ignore # pylint: disable=invalid-overridden-method """Defers the message. This message will remain in the queue but must be requested @@ -121,8 +118,7 @@ async def defer(self): await self._settle_message(MESSAGE_DEFER) self._settled = True - async def renew_lock(self): - # type: () -> None + async def renew_lock(self) -> None: # type: ignore # pylint: disable=invalid-overridden-method """Renew the message lock. This will maintain the lock on the message to ensure From c856949bf6fc28efc27b529a7de5160677bb1da9 Mon Sep 17 00:00:00 2001 From: yijxie Date: Wed, 13 May 2020 18:37:53 -0700 Subject: [PATCH 12/15] Change _from_connection_string to helper _convert_connection_string_to_kwargs --- .../azure/servicebus/_base_handler.py | 52 +++++++++---------- .../azure/servicebus/_servicebus_receiver.py | 5 +- .../azure/servicebus/_servicebus_sender.py | 5 +- .../servicebus/aio/_base_handler_async.py | 6 --- .../aio/_servicebus_receiver_async.py | 8 +-- .../aio/_servicebus_sender_async.py | 6 ++- 6 files changed, 41 insertions(+), 41 deletions(-) diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/_base_handler.py b/sdk/servicebus/azure-servicebus/azure/servicebus/_base_handler.py index 1271f0b6ca63..eba685b45164 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/_base_handler.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/_base_handler.py @@ -7,7 +7,7 @@ import uuid import time from datetime import timedelta -from typing import cast, Optional, Tuple, TYPE_CHECKING, Dict, Any, Callable +from typing import cast, Optional, Tuple, TYPE_CHECKING, Dict, Any, Callable, Type try: from urllib import quote_plus # type: ignore @@ -90,6 +90,31 @@ def _generate_sas_token(uri, policy, key, expiry=None): return _AccessToken(token=token, expires_on=abs_expiry) +def _convert_connection_string_to_kwargs(conn_str, shared_key_credential_type, **kwargs): + # type: (str, Type, Any) -> Dict[str, Any] + host, policy, key, entity_in_conn_str = _parse_conn_str(conn_str) + queue_name = kwargs.get("queue_name") + topic_name = kwargs.get("topic_name") + if not (queue_name or topic_name or entity_in_conn_str): + raise ValueError("Entity name is missing. Please specify `queue_name` or `topic_name`" + " or use a connection string including the entity information.") + + if queue_name and topic_name: + raise ValueError("`queue_name` and `topic_name` can not be specified simultaneously.") + + entity_in_kwargs = queue_name or topic_name + if entity_in_conn_str and entity_in_kwargs and (entity_in_conn_str != entity_in_kwargs): + raise ServiceBusAuthorizationError( + "Entity names do not match, the entity name in connection string is {};" + " the entity name in parameter is {}.".format(entity_in_conn_str, entity_in_kwargs) + ) + + kwargs["fully_qualified_namespace"] = host + kwargs["entity_name"] = entity_in_conn_str or entity_in_kwargs + kwargs["credential"] = shared_key_credential_type(policy, key) + return kwargs + + class ServiceBusSharedKeyCredential(object): """The shared access key credential used for authentication. @@ -133,31 +158,6 @@ def __init__( self._auth_uri = None self._properties = create_properties() - @staticmethod - def _from_connection_string(conn_str, **kwargs): - # type: (str, Any) -> Dict[str, Any] - host, policy, key, entity_in_conn_str = _parse_conn_str(conn_str) - queue_name = kwargs.get("queue_name") - topic_name = kwargs.get("topic_name") - if not (queue_name or topic_name or entity_in_conn_str): - raise ValueError("Entity name is missing. Please specify `queue_name` or `topic_name`" - " or use a connection string including the entity information.") - - if queue_name and topic_name: - raise ValueError("`queue_name` and `topic_name` can not be specified simultaneously.") - - entity_in_kwargs = queue_name or topic_name - if entity_in_conn_str and entity_in_kwargs and (entity_in_conn_str != entity_in_kwargs): - raise ServiceBusAuthorizationError( - "Entity names do not match, the entity name in connection string is {};" - " the entity name in parameter is {}.".format(entity_in_conn_str, entity_in_kwargs) - ) - - kwargs["fully_qualified_namespace"] = host - kwargs["entity_name"] = entity_in_conn_str or entity_in_kwargs - kwargs["credential"] = ServiceBusSharedKeyCredential(policy, key) - return kwargs - class BaseHandlerSync(BaseHandler): # pylint:disable=too-many-instance-attributes def __enter__(self): diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_receiver.py b/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_receiver.py index 600c3a6891a5..3abde83f1bd3 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_receiver.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_receiver.py @@ -11,7 +11,7 @@ from uamqp.constants import SenderSettleMode from uamqp.authentication.common import AMQPAuth -from ._base_handler import BaseHandlerSync +from ._base_handler import BaseHandlerSync, ServiceBusSharedKeyCredential, _convert_connection_string_to_kwargs from ._common.utils import create_authentication from ._common.message import PeekMessage, ReceivedMessage from ._common.constants import ( @@ -269,8 +269,9 @@ def from_connection_string( :caption: Create a new instance of the ServiceBusReceiver from connection string. """ - constructor_args = cls._from_connection_string( + constructor_args = _convert_connection_string_to_kwargs( conn_str, + ServiceBusSharedKeyCredential, **kwargs ) if kwargs.get("queue_name") and kwargs.get("subscription_name"): diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_sender.py b/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_sender.py index 63dc8396374a..05d5793af30c 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_sender.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_sender.py @@ -11,7 +11,7 @@ from uamqp import SendClient, types from uamqp.authentication.common import AMQPAuth -from ._base_handler import BaseHandlerSync +from ._base_handler import BaseHandlerSync, ServiceBusSharedKeyCredential, _convert_connection_string_to_kwargs from ._common import mgmt_handlers from ._common.message import Message, BatchMessage from .exceptions import ( @@ -288,8 +288,9 @@ def from_connection_string( :caption: Create a new instance of the ServiceBusSender from connection string. """ - constructor_args = cls._from_connection_string( + constructor_args = _convert_connection_string_to_kwargs( conn_str, + ServiceBusSharedKeyCredential, **kwargs ) return cls(**constructor_args) diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_base_handler_async.py b/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_base_handler_async.py index fe186c470eaa..cab95c43fa2d 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_base_handler_async.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_base_handler_async.py @@ -161,12 +161,6 @@ async def _mgmt_request_response_with_retry(self, mgmt_operation, message, callb **kwargs ) - @staticmethod - def _from_connection_string(conn_str, **kwargs): - kwargs = BaseHandler._from_connection_string(conn_str, **kwargs) - kwargs["credential"] = ServiceBusSharedKeyCredential(kwargs["credential"].policy, kwargs["credential"].key) - return kwargs - async def _open(self): # pylint: disable=no-self-use raise ValueError("Subclass should override the method.") diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_servicebus_receiver_async.py b/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_servicebus_receiver_async.py index ece5fa697f0e..38a243047a5f 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_servicebus_receiver_async.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_servicebus_receiver_async.py @@ -6,13 +6,14 @@ import collections import functools import logging -from typing import Any, TYPE_CHECKING, List, Union +from typing import Any, TYPE_CHECKING, List from uamqp import ReceiveClientAsync, types from uamqp.constants import SenderSettleMode -from ._base_handler_async import BaseHandlerAsync +from ._base_handler_async import BaseHandlerAsync, ServiceBusSharedKeyCredential from ._async_message import ReceivedMessage +from .._base_handler import _convert_connection_string_to_kwargs from .._common.receiver_mixins import ReceiverMixin from .._common.constants import ( REQUEST_RESPONSE_UPDATE_DISPOSTION_OPERATION, @@ -255,8 +256,9 @@ def from_connection_string( :caption: Create a new instance of the ServiceBusReceiver from connection string. """ - constructor_args = cls._from_connection_string( + constructor_args = _convert_connection_string_to_kwargs( conn_str, + ServiceBusSharedKeyCredential, **kwargs ) if kwargs.get("queue_name") and kwargs.get("subscription_name"): diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_servicebus_sender_async.py b/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_servicebus_sender_async.py index 039dcd00aa4a..e28d2cdcbfeb 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_servicebus_sender_async.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_servicebus_sender_async.py @@ -10,8 +10,9 @@ from uamqp import SendClientAsync, types from .._common.message import Message, BatchMessage +from .._base_handler import _convert_connection_string_to_kwargs from .._servicebus_sender import SenderMixin -from ._base_handler_async import BaseHandlerAsync +from ._base_handler_async import BaseHandlerAsync, ServiceBusSharedKeyCredential from .._common.constants import ( REQUEST_RESPONSE_SCHEDULE_MESSAGE_OPERATION, REQUEST_RESPONSE_CANCEL_SCHEDULED_MESSAGE_OPERATION, @@ -228,8 +229,9 @@ def from_connection_string( :caption: Create a new instance of the ServiceBusSender from connection string. """ - constructor_args = cls._from_connection_string( + constructor_args = _convert_connection_string_to_kwargs( conn_str, + ServiceBusSharedKeyCredential, **kwargs ) return cls(**constructor_args) From 068b83f7b71b17d92021d99f035235df570a7a49 Mon Sep 17 00:00:00 2001 From: yijxie Date: Thu, 14 May 2020 10:07:25 -0700 Subject: [PATCH 13/15] Remove BaseHandler --- .../azure/servicebus/_base_handler.py | 4 +-- .../azure/servicebus/_servicebus_receiver.py | 4 +-- .../azure/servicebus/_servicebus_sender.py | 4 +-- .../servicebus/aio/_base_handler_async.py | 33 ++++++++++++++++--- .../aio/_servicebus_receiver_async.py | 4 +-- .../aio/_servicebus_sender_async.py | 4 +-- 6 files changed, 38 insertions(+), 15 deletions(-) diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/_base_handler.py b/sdk/servicebus/azure-servicebus/azure/servicebus/_base_handler.py index eba685b45164..7ff53c728df8 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/_base_handler.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/_base_handler.py @@ -135,7 +135,7 @@ def get_token(self, *scopes, **kwargs): # pylint:disable=unused-argument return _generate_sas_token(scopes[0], self.policy, self.key) -class BaseHandler(object): +class BaseHandler: # pylint:disable=too-many-instance-attributes def __init__( self, fully_qualified_namespace, @@ -158,8 +158,6 @@ def __init__( self._auth_uri = None self._properties = create_properties() - -class BaseHandlerSync(BaseHandler): # pylint:disable=too-many-instance-attributes def __enter__(self): self._open_with_retry() return self diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_receiver.py b/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_receiver.py index 3abde83f1bd3..65c4f79b9a35 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_receiver.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_receiver.py @@ -11,7 +11,7 @@ from uamqp.constants import SenderSettleMode from uamqp.authentication.common import AMQPAuth -from ._base_handler import BaseHandlerSync, ServiceBusSharedKeyCredential, _convert_connection_string_to_kwargs +from ._base_handler import BaseHandler, ServiceBusSharedKeyCredential, _convert_connection_string_to_kwargs from ._common.utils import create_authentication from ._common.message import PeekMessage, ReceivedMessage from ._common.constants import ( @@ -37,7 +37,7 @@ _LOGGER = logging.getLogger(__name__) -class ServiceBusReceiver(BaseHandlerSync, ReceiverMixin): # pylint: disable=too-many-instance-attributes +class ServiceBusReceiver(BaseHandler, ReceiverMixin): # pylint: disable=too-many-instance-attributes """The ServiceBusReceiver class defines a high level interface for receiving messages from the Azure Service Bus Queue or Topic Subscription. diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_sender.py b/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_sender.py index 05d5793af30c..39ea04d16434 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_sender.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/_servicebus_sender.py @@ -11,7 +11,7 @@ from uamqp import SendClient, types from uamqp.authentication.common import AMQPAuth -from ._base_handler import BaseHandlerSync, ServiceBusSharedKeyCredential, _convert_connection_string_to_kwargs +from ._base_handler import BaseHandler, ServiceBusSharedKeyCredential, _convert_connection_string_to_kwargs from ._common import mgmt_handlers from ._common.message import Message, BatchMessage from .exceptions import ( @@ -82,7 +82,7 @@ def _build_schedule_request(cls, schedule_time_utc, *messages): return request_body -class ServiceBusSender(BaseHandlerSync, SenderMixin): +class ServiceBusSender(BaseHandler, SenderMixin): """The ServiceBusSender class defines a high level interface for sending messages to the Azure Service Bus Queue or Topic. diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_base_handler_async.py b/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_base_handler_async.py index cab95c43fa2d..4b2eb4f72e2c 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_base_handler_async.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_base_handler_async.py @@ -4,17 +4,20 @@ # -------------------------------------------------------------------------------------------- import logging import asyncio +import uuid from typing import TYPE_CHECKING, Any import uamqp +from azure.servicebus._common._configuration import Configuration +from azure.servicebus._common.utils import create_properties from uamqp.message import MessageProperties -from .._base_handler import BaseHandler, _generate_sas_token +from .._base_handler import _generate_sas_token from .._common.constants import ( TOKEN_TYPE_SASTOKEN, MGMT_REQUEST_OP_TYPE_ENTITY_MGMT, - ASSOCIATEDLINKPROPERTYNAME -) + ASSOCIATEDLINKPROPERTYNAME, + CONTAINER_PREFIX, MANAGEMENT_PATH_SUFFIX) from ..exceptions import ( ServiceBusError, _create_servicebus_exception @@ -44,7 +47,29 @@ async def get_token(self, *scopes, **kwargs): # pylint:disable=unused-argument return _generate_sas_token(scopes[0], self.policy, self.key) -class BaseHandlerAsync(BaseHandler): +class BaseHandler: + def __init__( + self, + fully_qualified_namespace, + entity_name, + credential, + **kwargs + ): + # type: (str, str, TokenCredential, Any) -> None + self.fully_qualified_namespace = fully_qualified_namespace + self._entity_name = entity_name + + subscription_name = kwargs.get("subscription_name") + self._mgmt_target = self._entity_name + (("/Subscriptions/" + subscription_name) if subscription_name else '') + self._mgmt_target = "{}{}".format(self._mgmt_target, MANAGEMENT_PATH_SUFFIX) + self._credential = credential + self._container_id = CONTAINER_PREFIX + str(uuid.uuid4())[:8] + self._config = Configuration(**kwargs) + self._running = False + self._handler = None # type: uamqp.AMQPClient + self._auth_uri = None + self._properties = create_properties() + async def __aenter__(self): await self._open_with_retry() return self diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_servicebus_receiver_async.py b/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_servicebus_receiver_async.py index 38a243047a5f..f8870cfc035d 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_servicebus_receiver_async.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_servicebus_receiver_async.py @@ -11,7 +11,7 @@ from uamqp import ReceiveClientAsync, types from uamqp.constants import SenderSettleMode -from ._base_handler_async import BaseHandlerAsync, ServiceBusSharedKeyCredential +from ._base_handler_async import BaseHandler, ServiceBusSharedKeyCredential from ._async_message import ReceivedMessage from .._base_handler import _convert_connection_string_to_kwargs from .._common.receiver_mixins import ReceiverMixin @@ -37,7 +37,7 @@ _LOGGER = logging.getLogger(__name__) -class ServiceBusReceiver(collections.abc.AsyncIterator, BaseHandlerAsync, ReceiverMixin): +class ServiceBusReceiver(collections.abc.AsyncIterator, BaseHandler, ReceiverMixin): """The ServiceBusReceiver class defines a high level interface for receiving messages from the Azure Service Bus Queue or Topic Subscription. diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_servicebus_sender_async.py b/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_servicebus_sender_async.py index e28d2cdcbfeb..8a78979c5373 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_servicebus_sender_async.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_servicebus_sender_async.py @@ -12,7 +12,7 @@ from .._common.message import Message, BatchMessage from .._base_handler import _convert_connection_string_to_kwargs from .._servicebus_sender import SenderMixin -from ._base_handler_async import BaseHandlerAsync, ServiceBusSharedKeyCredential +from ._base_handler_async import BaseHandler, ServiceBusSharedKeyCredential from .._common.constants import ( REQUEST_RESPONSE_SCHEDULE_MESSAGE_OPERATION, REQUEST_RESPONSE_CANCEL_SCHEDULED_MESSAGE_OPERATION, @@ -28,7 +28,7 @@ _LOGGER = logging.getLogger(__name__) -class ServiceBusSender(BaseHandlerAsync, SenderMixin): +class ServiceBusSender(BaseHandler, SenderMixin): """The ServiceBusSender class defines a high level interface for sending messages to the Azure Service Bus Queue or Topic. From 6e28814198b645c993024935f6d8716880d3a437 Mon Sep 17 00:00:00 2001 From: yijxie Date: Thu, 14 May 2020 10:07:47 -0700 Subject: [PATCH 14/15] Adjust for pylint 2.3.1 --- .../azure/servicebus/aio/_async_message.py | 12 ++++++------ 1 file changed, 6 insertions(+), 6 deletions(-) diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_async_message.py b/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_async_message.py index 286b9ac21f46..bf916a391bb4 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_async_message.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_async_message.py @@ -26,7 +26,7 @@ class ReceivedMessage(sync_message.ReceivedMessage): """A Service Bus Message received from service side. """ - async def _settle_message( # type: ignore # pylint: disable=invalid-overridden-method + async def _settle_message( # type: ignore self, settle_operation, dead_letter_details=None @@ -50,7 +50,7 @@ async def _settle_message( # type: ignore # pylint: disable=invalid-overridden except Exception as e: raise MessageSettleFailed(settle_operation, e) - async def complete(self) -> None: # type: ignore # pylint: disable=invalid-overridden-method + async def complete(self) -> None: # type: ignore """Complete the message. This removes the message from the queue. @@ -68,7 +68,7 @@ async def complete(self) -> None: # type: ignore # pylint: disable=invalid-ove async def dead_letter( # type: ignore self, reason: Optional[str] = None, description: Optional[str] = None - ) -> None: # pylint: disable=unused-argument,invalid-overridden-method + ) -> None: # pylint: disable=unused-argument """Move the message to the Dead Letter queue. The Dead Letter queue is a sub-queue that can be @@ -87,7 +87,7 @@ async def dead_letter( # type: ignore await self._settle_message(MESSAGE_DEAD_LETTER) self._settled = True - async def abandon(self) -> None: # type: ignore # pylint: disable=invalid-overridden-method + async def abandon(self) -> None: # type: ignore """Abandon the message. This message will be returned to the queue and made available to be received again. @@ -102,7 +102,7 @@ async def abandon(self) -> None: # type: ignore # pylint: disable=invalid-over await self._settle_message(MESSAGE_ABANDON) self._settled = True - async def defer(self) -> None: # type: ignore # pylint: disable=invalid-overridden-method + async def defer(self) -> None: # type: ignore """Defers the message. This message will remain in the queue but must be requested @@ -118,7 +118,7 @@ async def defer(self) -> None: # type: ignore # pylint: disable=invalid-overri await self._settle_message(MESSAGE_DEFER) self._settled = True - async def renew_lock(self) -> None: # type: ignore # pylint: disable=invalid-overridden-method + async def renew_lock(self) -> None: # type: ignore """Renew the message lock. This will maintain the lock on the message to ensure From ae1dda6c02aab1a6879ef6732a55169388bd8dc5 Mon Sep 17 00:00:00 2001 From: yijxie Date: Thu, 14 May 2020 10:10:18 -0700 Subject: [PATCH 15/15] Fix pylint errors --- .../azure/servicebus/aio/_base_handler_async.py | 5 ++--- 1 file changed, 2 insertions(+), 3 deletions(-) diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_base_handler_async.py b/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_base_handler_async.py index 4b2eb4f72e2c..2ad4335e1333 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_base_handler_async.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/aio/_base_handler_async.py @@ -8,11 +8,10 @@ from typing import TYPE_CHECKING, Any import uamqp -from azure.servicebus._common._configuration import Configuration -from azure.servicebus._common.utils import create_properties from uamqp.message import MessageProperties - from .._base_handler import _generate_sas_token +from .._common._configuration import Configuration +from .._common.utils import create_properties from .._common.constants import ( TOKEN_TYPE_SASTOKEN, MGMT_REQUEST_OP_TYPE_ENTITY_MGMT,