Skip to content
Merged
Show file tree
Hide file tree
Changes from 2 commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions sdk/servicebus/azure-servicebus/CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,9 @@

## 7.0.0b6 (Unreleased)

**Breaking Changes**

* `ServiceBusClient.close()` now closes spawned senders and receivers.

## 7.0.0b5 (2020-08-10)

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2,11 +2,11 @@
# Copyright (c) Microsoft Corporation. All rights reserved.
# Licensed under the MIT License. See License.txt in the project root for license information.
# --------------------------------------------------------------------------------------------
from typing import Any, TYPE_CHECKING
from typing import Any, List, TYPE_CHECKING

import uamqp

from ._base_handler import _parse_conn_str, ServiceBusSharedKeyCredential
from ._base_handler import _parse_conn_str, ServiceBusSharedKeyCredential, BaseHandler
from ._servicebus_sender import ServiceBusSender
from ._servicebus_receiver import ServiceBusReceiver
from ._servicebus_session_receiver import ServiceBusSessionReceiver
Expand Down Expand Up @@ -69,6 +69,7 @@ def __init__(
self._auth_uri = "{}/{}".format(self._auth_uri, self._entity_name)
# Internal flag for switching whether to apply connection sharing, pending fix in uamqp library
self._connection_sharing = False
self._handlers = [] # type: List[BaseHandler]

def __enter__(self):
if self._connection_sharing:
Expand All @@ -89,10 +90,15 @@ def _create_uamqp_connection(self):
def close(self):
# type: () -> None
"""
Close down the ServiceBus client and the underlying connection.
Close down the ServiceBus client.
All spawned senders, receivers and underlying connection will be shutdown.

:return: None
"""
for handler in self._handlers:
handler.close()
Comment thread
KieranBrantnerMagee marked this conversation as resolved.
Outdated
self._handlers.clear()

if self._connection_sharing and self._connection:
self._connection.destroy()

Expand Down Expand Up @@ -157,7 +163,7 @@ def get_queue_sender(self, queue_name, **kwargs):

"""
# pylint: disable=protected-access
return ServiceBusSender(
handler = ServiceBusSender(
fully_qualified_namespace=self.fully_qualified_namespace,
queue_name=queue_name,
credential=self._credential,
Expand All @@ -168,6 +174,8 @@ def get_queue_sender(self, queue_name, **kwargs):
user_agent=self._config.user_agent,
**kwargs
)
self._handlers.append(handler)
return handler

def get_queue_receiver(self, queue_name, **kwargs):
# type: (str, Any) -> ServiceBusReceiver
Expand Down Expand Up @@ -205,7 +213,7 @@ def get_queue_receiver(self, queue_name, **kwargs):

"""
# pylint: disable=protected-access
return ServiceBusReceiver(
handler = ServiceBusReceiver(
fully_qualified_namespace=self.fully_qualified_namespace,
queue_name=queue_name,
credential=self._credential,
Expand All @@ -216,6 +224,8 @@ def get_queue_receiver(self, queue_name, **kwargs):
user_agent=self._config.user_agent,
**kwargs
)
self._handlers.append(handler)
return handler

def get_queue_deadletter_receiver(self, queue_name, **kwargs):
# type: (str, Any) -> ServiceBusReceiver
Expand Down Expand Up @@ -265,7 +275,7 @@ def get_queue_deadletter_receiver(self, queue_name, **kwargs):
queue_name=queue_name,
transfer_deadletter=kwargs.get('transfer_deadletter', False)
)
return ServiceBusReceiver(
handler = ServiceBusReceiver(
fully_qualified_namespace=self.fully_qualified_namespace,
entity_name=entity_name,
credential=self._credential,
Expand All @@ -277,6 +287,8 @@ def get_queue_deadletter_receiver(self, queue_name, **kwargs):
user_agent=self._config.user_agent,
**kwargs
)
self._handlers.append(handler)
return handler

def get_topic_sender(self, topic_name, **kwargs):
# type: (str, Any) -> ServiceBusSender
Expand All @@ -300,7 +312,7 @@ def get_topic_sender(self, topic_name, **kwargs):
:caption: Create a new instance of the ServiceBusSender from ServiceBusClient.

"""
return ServiceBusSender(
handler = ServiceBusSender(
fully_qualified_namespace=self.fully_qualified_namespace,
topic_name=topic_name,
credential=self._credential,
Expand All @@ -311,6 +323,8 @@ def get_topic_sender(self, topic_name, **kwargs):
user_agent=self._config.user_agent,
**kwargs
)
self._handlers.append(handler)
return handler

def get_subscription_receiver(self, topic_name, subscription_name, **kwargs):
# type: (str, str, Any) -> ServiceBusReceiver
Expand Down Expand Up @@ -353,7 +367,7 @@ def get_subscription_receiver(self, topic_name, subscription_name, **kwargs):

"""
# pylint: disable=protected-access
return ServiceBusReceiver(
handler = ServiceBusReceiver(
fully_qualified_namespace=self.fully_qualified_namespace,
topic_name=topic_name,
subscription_name=subscription_name,
Expand All @@ -365,6 +379,8 @@ def get_subscription_receiver(self, topic_name, subscription_name, **kwargs):
user_agent=self._config.user_agent,
**kwargs
)
self._handlers.append(handler)
return handler

def get_subscription_deadletter_receiver(self, topic_name, subscription_name, **kwargs):
# type: (str, str, Any) -> ServiceBusReceiver
Expand Down Expand Up @@ -416,7 +432,7 @@ def get_subscription_deadletter_receiver(self, topic_name, subscription_name, **
subscription_name=subscription_name,
transfer_deadletter=kwargs.get('transfer_deadletter', False)
)
return ServiceBusReceiver(
handler = ServiceBusReceiver(
fully_qualified_namespace=self.fully_qualified_namespace,
entity_name=entity_name,
credential=self._credential,
Expand All @@ -428,6 +444,8 @@ def get_subscription_deadletter_receiver(self, topic_name, subscription_name, **
user_agent=self._config.user_agent,
**kwargs
)
self._handlers.append(handler)
return handler

def get_subscription_session_receiver(self, topic_name, subscription_name, session_id=None, **kwargs):
# type: (str, str, str, Any) -> ServiceBusReceiver
Expand Down Expand Up @@ -473,7 +491,7 @@ def get_subscription_session_receiver(self, topic_name, subscription_name, sessi

"""
# pylint: disable=protected-access
return ServiceBusSessionReceiver(
handler = ServiceBusSessionReceiver(
fully_qualified_namespace=self.fully_qualified_namespace,
topic_name=topic_name,
subscription_name=subscription_name,
Expand All @@ -486,6 +504,8 @@ def get_subscription_session_receiver(self, topic_name, subscription_name, sessi
user_agent=self._config.user_agent,
**kwargs
)
self._handlers.append(handler)
return handler

def get_queue_session_receiver(self, queue_name, session_id=None, **kwargs):
# type: (str, str, Any) -> ServiceBusSessionReceiver
Expand Down Expand Up @@ -526,7 +546,7 @@ def get_queue_session_receiver(self, queue_name, session_id=None, **kwargs):

"""
# pylint: disable=protected-access
return ServiceBusSessionReceiver(
handler = ServiceBusSessionReceiver(
fully_qualified_namespace=self.fully_qualified_namespace,
queue_name=queue_name,
credential=self._credential,
Expand All @@ -538,3 +558,5 @@ def get_queue_session_receiver(self, queue_name, session_id=None, **kwargs):
user_agent=self._config.user_agent,
**kwargs
)
self._handlers.append(handler)
return handler
Original file line number Diff line number Diff line change
Expand Up @@ -2,12 +2,12 @@
# Copyright (c) Microsoft Corporation. All rights reserved.
# Licensed under the MIT License. See License.txt in the project root for license information.
# --------------------------------------------------------------------------------------------
from typing import Any, TYPE_CHECKING
from typing import Any, List, TYPE_CHECKING

import uamqp

from .._base_handler import _parse_conn_str
from ._base_handler_async import ServiceBusSharedKeyCredential
from ._base_handler_async import ServiceBusSharedKeyCredential, BaseHandler
from ._servicebus_sender_async import ServiceBusSender
from ._servicebus_receiver_async import ServiceBusReceiver
from ._servicebus_session_receiver_async import ServiceBusSessionReceiver
Expand Down Expand Up @@ -71,6 +71,7 @@ def __init__(
self._auth_uri = "{}/{}".format(self._auth_uri, self._entity_name)
# Internal flag for switching whether to apply connection sharing, pending fix in uamqp library
self._connection_sharing = False
self._handlers = [] # type: List[BaseHandler]

async def __aenter__(self):
if self._connection_sharing:
Expand Down Expand Up @@ -133,9 +134,14 @@ async def close(self):
# type: () -> None
"""
Close down the ServiceBus client.
All spawned senders, receivers and underlying connection will be shutdown.

:return: None
"""
for handler in self._handlers:
await handler.close()
self._handlers.clear()

if self._connection_sharing and self._connection:
await self._connection.destroy_async()

Expand All @@ -159,7 +165,7 @@ def get_queue_sender(self, queue_name, **kwargs):

"""
# pylint: disable=protected-access
return ServiceBusSender(
handler = ServiceBusSender(
fully_qualified_namespace=self.fully_qualified_namespace,
queue_name=queue_name,
credential=self._credential,
Expand All @@ -170,6 +176,8 @@ def get_queue_sender(self, queue_name, **kwargs):
user_agent=self._config.user_agent,
**kwargs
)
self._handlers.append(handler)
return handler

def get_queue_receiver(self, queue_name, **kwargs):
# type: (str, Any) -> ServiceBusReceiver
Expand Down Expand Up @@ -206,7 +214,7 @@ def get_queue_receiver(self, queue_name, **kwargs):

"""
# pylint: disable=protected-access
return ServiceBusReceiver(
handler = ServiceBusReceiver(
fully_qualified_namespace=self.fully_qualified_namespace,
queue_name=queue_name,
credential=self._credential,
Expand All @@ -217,6 +225,8 @@ def get_queue_receiver(self, queue_name, **kwargs):
user_agent=self._config.user_agent,
**kwargs
)
self._handlers.append(handler)
return handler

def get_queue_deadletter_receiver(self, queue_name, **kwargs):
# type: (str, Any) -> ServiceBusReceiver
Expand Down Expand Up @@ -266,7 +276,7 @@ def get_queue_deadletter_receiver(self, queue_name, **kwargs):
queue_name=queue_name,
transfer_deadletter=kwargs.get('transfer_deadletter', False)
)
return ServiceBusReceiver(
handler = ServiceBusReceiver(
fully_qualified_namespace=self.fully_qualified_namespace,
entity_name=entity_name,
credential=self._credential,
Expand All @@ -278,6 +288,8 @@ def get_queue_deadletter_receiver(self, queue_name, **kwargs):
user_agent=self._config.user_agent,
**kwargs
)
self._handlers.append(handler)
return handler

def get_topic_sender(self, topic_name, **kwargs):
# type: (str, Any) -> ServiceBusSender
Expand All @@ -301,7 +313,7 @@ def get_topic_sender(self, topic_name, **kwargs):
:caption: Create a new instance of the ServiceBusSender from ServiceBusClient.

"""
return ServiceBusSender(
handler = ServiceBusSender(
fully_qualified_namespace=self.fully_qualified_namespace,
topic_name=topic_name,
credential=self._credential,
Expand All @@ -312,6 +324,8 @@ def get_topic_sender(self, topic_name, **kwargs):
user_agent=self._config.user_agent,
**kwargs
)
self._handlers.append(handler)
return handler

def get_subscription_receiver(self, topic_name, subscription_name, **kwargs):
# type: (str, str, Any) -> ServiceBusReceiver
Expand Down Expand Up @@ -354,7 +368,7 @@ def get_subscription_receiver(self, topic_name, subscription_name, **kwargs):

"""
# pylint: disable=protected-access
return ServiceBusReceiver(
handler = ServiceBusReceiver(
fully_qualified_namespace=self.fully_qualified_namespace,
topic_name=topic_name,
subscription_name=subscription_name,
Expand All @@ -366,6 +380,8 @@ def get_subscription_receiver(self, topic_name, subscription_name, **kwargs):
user_agent=self._config.user_agent,
**kwargs
)
self._handlers.append(handler)
return handler

def get_subscription_deadletter_receiver(self, topic_name, subscription_name, **kwargs):
# type: (str, str, Any) -> ServiceBusReceiver
Expand Down Expand Up @@ -417,7 +433,7 @@ def get_subscription_deadletter_receiver(self, topic_name, subscription_name, **
subscription_name=subscription_name,
transfer_deadletter=kwargs.get('transfer_deadletter', False)
)
return ServiceBusReceiver(
handler = ServiceBusReceiver(
fully_qualified_namespace=self.fully_qualified_namespace,
entity_name=entity_name,
credential=self._credential,
Expand All @@ -429,6 +445,8 @@ def get_subscription_deadletter_receiver(self, topic_name, subscription_name, **
user_agent=self._config.user_agent,
**kwargs
)
self._handlers.append(handler)
return handler

def get_subscription_session_receiver(self, topic_name, subscription_name, session_id=None, **kwargs):
# type: (str, str, str, Any) -> ServiceBusReceiver
Expand Down Expand Up @@ -474,7 +492,7 @@ def get_subscription_session_receiver(self, topic_name, subscription_name, sessi

"""
# pylint: disable=protected-access
return ServiceBusSessionReceiver(
handler = ServiceBusSessionReceiver(
fully_qualified_namespace=self.fully_qualified_namespace,
topic_name=topic_name,
subscription_name=subscription_name,
Expand All @@ -487,6 +505,8 @@ def get_subscription_session_receiver(self, topic_name, subscription_name, sessi
user_agent=self._config.user_agent,
**kwargs
)
self._handlers.append(handler)
return handler

def get_queue_session_receiver(self, queue_name, session_id=None, **kwargs):
# type: (str, str, Any) -> ServiceBusSessionReceiver
Expand Down Expand Up @@ -526,7 +546,7 @@ def get_queue_session_receiver(self, queue_name, session_id=None, **kwargs):

"""
# pylint: disable=protected-access
return ServiceBusSessionReceiver(
handler = ServiceBusSessionReceiver(
fully_qualified_namespace=self.fully_qualified_namespace,
queue_name=queue_name,
credential=self._credential,
Expand All @@ -538,3 +558,5 @@ def get_queue_session_receiver(self, queue_name, session_id=None, **kwargs):
user_agent=self._config.user_agent,
**kwargs
)
self._handlers.append(handler)
return handler
Loading