From e115bbd51eb3e861baab3bcb07eb2538331881bd Mon Sep 17 00:00:00 2001 From: Yunhao Ling Date: Tue, 18 Feb 2020 19:00:03 -0800 Subject: [PATCH 01/24] api proposal for t2 sb --- sdk/servicebus/azure-servicebus/api_review.py | 576 ++++++++++++++++++ 1 file changed, 576 insertions(+) create mode 100644 sdk/servicebus/azure-servicebus/api_review.py diff --git a/sdk/servicebus/azure-servicebus/api_review.py b/sdk/servicebus/azure-servicebus/api_review.py new file mode 100644 index 000000000000..fa37eadd8148 --- /dev/null +++ b/sdk/servicebus/azure-servicebus/api_review.py @@ -0,0 +1,576 @@ +# Approach 1 + +## APIs: +class SeviceBusSenderClient: + def __init__( + self, + fully_qualified_namespace : str, + entity_name : str, + credential : TokenCredential, + logging_enable: bool = False, + http_proxy: dict = None, + transport_type: TransportType = None, + retry_total : int = 3, + ) -> None: + @classmethod + def from_queue( + cls, + fully_qualified_namepsace : str, + queue_name : str, + credential : TokenCredential, + **kwargs + ) -> SeviceBusSenderClient: + @classmethod + def from_topic( + cls, + fully_qualified_namespace : str, + topic_name : str, + credential : TokenCredential, + **kwargs + ) -> SeviceBusSenderClient: + @classmethod + def from_connection_string( + cls, + conn_str : str, + entity_name : str = None, + **kwargs + ) -> SeviceBusSenderClient: + + def __enter__(self): + def __exit__(self): + + def close(self) -> None: + + def get_properties(self) -> Dict[str, Any]: + + def list_session_ids(self) -> List[str]: + + def send( + self, + messages, + session_id : str = None, + message_timeout : float = None + ) -> None: + def schedule( + self, + messages, + schedule_time : datetime, + session_id : str = None + ) -> List[int]: + def cancel_scheduled_messages(self, sequence_numbers : List[int]) -> None: + + +class ServiceBusReceiverClient: + def __init__( + self, + fully_qualified_namespace : str, + entity_name : str, + subscription_name : str = None, + credential : TokenCredential, + logging_enable: bool = False, + http_proxy: dict = None, + transport_type: TransportType = None, + retry_total : int = 3, + ) -> None: + @classmethod + def from_queue( + cls, + fully_qualified_namespace : str, + queue_name : str, + credential : TokenCredential, + **kwargs + ) -> ServiceBusReceiverClient: + @classmethod + def from_topic_subscription( + cls, + fully_qualified_namespace : str, + topic_name : str, + subscription_name : str, + credential: TokenCredential, + **kwargs + ) -> ServiceBusReceiverClient: + @classmethod + def from_connection_string( + cls, + conn_str : str, + entity_name : str = None, + subscription_name : str = None, + **kwargs + )-> ServiceBusReceiverClient: + + def __enter__(self): + def __exit__(self): + + def close(self) -> None: + + def get_properties(self) -> Dict[str, any]: + + def __iter__(self): + def __next__(self): + def next(self): + + def peek( + self, + session : Union[str, Session] = None, + count : int = 1, + start_from : int = None + ) -> List[PeekMessage]: + def receive_deferred_messages( + self, + session : Union[str, Session] = None, + sequence_numbers : List[int], + mode : ReceiveSettleMode =ReceiveSettleMode.PeekLock + ) -> List[DeferredMessage]: + def receive( + self, + session : Union[str, Session] = None, + max_batch_size : int = None, + timeout : float = None + ) -> List[Message]: # Pull mode receive + def settle_deferred_messages( + self, + settlement : str, # TODO: In T1 settlement is just string like 'completed', 'suspended', 'abandoned', can we improve this parameter to be Enum or remove this method? + messages : List[DeferredMessage], + **kwargs + ) -> None: # Batch settle deferred messages + + def list_session_ids(self) -> List[str]: + def get_session(self, session_id : str = "NEXT_AVAILABLE") -> Session: # raise Error when called on non-session entity + + # Rule APIs, raise Error when called on Queue + def add_rule(self, rule_name : str, filter : str) -> None: # TODO: figure out what filter is + def remove_rule(self, rule_name : str) -> None: + def get_rules(self) -> List[str] -> Dict[str, Any]: + + +class Session: + def get_session_id(self) -> str: + def get_session_state(self) -> str: # This can be made to property if property is more pythonic way + def set_session_state(self, state : str) -> None: # This can be made to property if property is more pythonic way + def renew_lock(self) -> None: + @property + def expired(self) -> bool: + + +class Transaction: + def __init__(self) -> None: + def __enter__(self): + def __exit__(self): + def begin(self) -> None: + def commit(self) -> None: + def abort(self) -> None: + + +class Message: + def __init__(self, body : str, encoding : str = 'UTF-8', **kwargs) -> None: + def __str__(self): + + # @properties + def settled(self) -> bool: # read-only + def enqueued_time(self) -> datetime: # read-only + def scheduled_enqueue_time(self) -> datetime: # read-only + def sequence_number(self) -> int: # read-only + def partition_id(self) -> str: # read-only + def locked_until(self) -> datetime: # read-only + def expired(self) -> bool: # read-only + def lock_token(self) -> str: # read-only + def body(self) -> Union[bytes, Generator[bytes]]: # read-only + def partition_key(self, value : str): + def partition_key(self) -> str: + def via_partition_key(self, value: str): + def via_partition_key(self) -> str: + def sessio_id(self, value : str): + def sessio_id(self) -> str: + def time_to_live(self, value : Union[float, timedelta]): + def time_to_live(self) -> datetime.timedelta: + def annotations(self, value : dict[str, Any]): + def annotations(self) -> dict[str, Any]: + def user_properties(self, value : dict[str, Any]): + def user_properties(self) -> dict[str, Any]: + def enqueue_sequence_number(self, value : int): + def enqueue_sequence_number(self) -> int: + + # Methods + def schedule(self, schedule_time : datetime) -> None: + def renew_lock(self) -> None: + def complete(self) -> None: + def dead_letter(self, description : str = None) -> None: + def abandon(self) -> None: + def defer(self) -> None: + + +class BatchMessage(Message): +# inherited from Message + + +class PeekMessage(Message): +# inherited from Message, and cannot settle/lock the message + +class DeferredMessage(Message): + # interited from Message + # its own @properties + def settled(self) -> bool: # read-only + +class ReceiveSettleMode: + PeekLock + ReceiveAndDelete + + +class TransportType(Enum): + Amqp + AmqpOverWebsocket + + +class ServiceBusSharedKeyCredential: + def __init__(self, policy, key): + def get_token(self, *scopes, **kwargs): + + +## Sampe Code: + + +### creation of differen entity clients + +# Queue Sender +queue_sender = SeviceBusSenderClient.from_queue( + fully_qualified_namepsace="fully_qualified_namepsace", + queue_name="queue_name", + credential=ServiceBusSharedKeyCredential("policy", "key") +) + +queue_sender = SeviceBusSenderClient.from_connection_string( + conn_str="conn_str", + entity_name="queue_name" +) + +queue_sender = SeviceBusSenderClient( + fully_qualified_namepsace="fully_qualified_namepsace", + entity_name="queue_name", + ServiceBusSharedKeyCredential("policy", "key") +) + +# Queue Receiver +queue_receiver = SeviceBusReceiverClient.from_queue( + fully_qualified_namepsace="fully_qualified_namepsace", + queue_name="queue_name", + credential=ServiceBusSharedKeyCredential("policy", "key") +) + +queue_receiver = SeviceBusReceiverClient.from_connection_string( + conn_str="conn_str", + entity_name="queue_name" +) + +queue_receiver = SeviceBusReceiverClient( + fully_qualified_namepsace="fully_qualified_namepsace", + entity_name="queue_name", + credential=ServiceBusSharedKeyCredential("policy", "key") +) + +# Topic Sender +topic_sender = SeviceBusSenderClient.from_topic( + fully_qualified_namepsace="fully_qualified_namepsace", + topic_name="topic_name", + credential=ServiceBusSharedKeyCredential("policy", "key") +) + +topic_sender = SeviceBusSenderClient.from_connection_string( + conn_str="conn_str", + entity_name="topic_name" +) + +topic_sender = SeviceBusSenderClient( + fully_qualified_namepsace="fully_qualified_namepsace", + entity_name="topic_name", + credential=ServiceBusSharedKeyCredential("policy", "key") +) + +# Subscription Receiver +subscription_receiver = SeviceBusReceiverClient.from_topic_subscription( + fully_qualified_namepsace="fully_qualified_namespace", + topic_name="topic_name", + subscription_name="subscription_name", + credential=ServiceBusSharedKeyCredential("policy", "key") +) + +subscription_receiver = SeviceBusReceiverClient.from_connection_string( + conn_str="conn_str", + entity_name="topic_name", + subscription_name="subscription_name" +) + +subscription_receiver = SeviceBusReceiverClient( + fully_qualified_namepsace="fully_qualified_namespace", + entity_name="topic_name", + subscription_name="subcription_name", + credential=ServiceBusSharedKeyCredential("policy", "key") +) + + +### Send and receive +# 1 Send + +queue_sender = SeviceBusSenderClient.from_queue( + fully_qualified_namepsace="fully_qualified_namepsace", + queue_name="queue_name", + credential=ServiceBusSharedKeyCredential("policy", "key") +) + +''' +topic_sender = SeviceBusSenderClient.from_queue( + fully_qualified_namepsace="fully_qualified_namepsace", + topic_name="topic_name", + credential=ServiceBusSharedKeyCredential("policy", "key") +) +''' + +with queue_sender: + # 1.1 send single + message = Message("Test message") + queue_sender.send(message) + + # 1.2 send list + messages = [Message("Test message {}".format(i)) for i in range(10)] + queue_sender.send(messages) + + # 1.3 schedule + schedule_message = Message("scheduled_message") + enqueue_time = datetime.utcnow() + timedelta(minutes=10) + sequence_number = queue_sender.schedule(schedule_message, enqueue_time) + + # 1.4 cancel schedule + queue_sender.cancel_scheduled_messages(*sequence_number) + +# 2 Receive + +queue_receiver = ServiceBusReceiverClient.from_queue( + fully_qualified_namespace="fully_qualified_namespace", + queue_name="queue_name", + credential=ServiceBusSharedKeyCredential("policy", "key") +) + +''' +subscription_receiver = ServiceBusReceiverClient.from_queue( + fully_qualified_namespace="fully_qualified_namespace", + topic_name="topic_name", + subscription_name="subscription_name" + credential=ServiceBusSharedKeyCredential("policy", "key") +) +''' + +# 2.1 peek +with queue_receiver: + msgs = queue_receiver.peek(count=10) + +# 2.2 iterator receive +with queue_receiver: + for msg in queue_receiver: + print(message) + msg.complete() + # msg.renew_lock() + # msg.abondon() + # msg.defer() + # msg.dead_letter() + +# 2.3 pull mode receive +with queue_receiver: + msgs = queue_receiver.receive(max_batch_size=10) + for msg in msgs: + msg.complete() + +# 2.4 receive deferred letter +with queue_receiver: + defered_sequence_numbers = [xxxx] + deferred = queue_receiver.receive_deferred_messages(defered_sequence_numbers) + queue_receiver.settle_deferred_messages('completed', deferred) + + +### Session operation + + +# 1 Send +topic_sender = SeviceBusSenderClient.from_topic +( + fully_qualified_namespace="fully_qualified_namepsace", + topic_name="topic_name", + credential=ServiceBusSharedKeyCredential("policy", "key") +) + +with topic_sender: + # 1.1 send single + message = Message("Test message") + # message.session_id = "test_session" + topic_sender.send(message, session="test_session") + + # 1.2 send list + messages = [Message("Test message {}".format(i)) for i in range(10)] + topic_sender.send(messages, session="test_session") + + # 1.3 schedule + schedule_message = Message("scheduled_message") + enqueue_time = datetime.utcnow() + timedelta(minutes=10) + sequence_number = topic_sender.schedule(schedule_message, enqueue_time, session="test_session") + + # 1.4 cancel schedule + topic_sender.cancel_scheduled_messages(*sequence_number) + +# 2 Receive +subscription_receiver = ServiceBusReceiverClient.from_topic_subscription +( + fully_qualified_namespace="fully_qualified_namespace", + topic_name="topic_name", + subscription_name="subscription_name", + credential=ServiceBusSharedKeyCredential("policy", "key") +) + +# 2.1 peek +with subscription_receiver: + msgs = subscription_receiver.peek(count=10, session="test_session") + +# 2.2 iterator receive +with subscription_receiver: + session = subscription_receiver.get_session("test_session") + # Option 2: + # session_ids = subscription_receiver.list_session() + # session = subscription_receiver.get_session(session_ids[0]) + session.set_session_state("START") + for msg in subscription_receiver(session=session): + print(message) + msg.complete() + session.renew_lock() + if str(msg) == "shutdown": + session.set_session_state("STOP") + break + +# 2.3 pull mode receive +with subscription_receiver: + session = subscription_receiver.get_session("test_session") + msgs = subscription_receiver.receive(max_batch_size=10, session=session) + session.set_session_state("BEGIN") + for msg in msgs: + msg.complete() + session.renew_lock() + session.set_session_state("END") + +# 2.4 receive deferred letter +with subscription_receiver: + session = subscription_receiver.get_session("test_session") + if session.get_session_state() == "DONE": + defered_sequence_numbers = [xxxx] + deferred = subscription_receiver.receive_deferred_messages(defered_sequence_numbers, session=session) + subscription_receiver.settle_deferred_messages('completed', deferred) + +### Rule operations + +subscription_receiver = ServiceBusReceiverClient.from_topic_subscription +( + fully_qualified_namespace="fully_qualified_namespace", + topic_name="topic_name", + subscription_name="subscription_name", + credential=ServiceBusSharedKeyCredential("policy", "key")) + +with subscription_receiver: + subscription_receiver.get_rules() + subscription_receiver.add_rule("test_rule", "1=1") + subscription_receiver.remove_rule("test_rule") + +### Transaction operations + +subscription_receiver = ServiceBusReceiverClient.from_topic_subscription +( + fully_qualified_namespace="fully_qualified_namespace", + topic_name="topic_name", + subscription_name="subscription_name", + credential=ServiceBusSharedKeyCredential("policy", "key") +) + +with subscription_receiver: + transaction = subscription_receiver.get_transaction() # TODO, no idea on how to get a transaction for now + with transaction: + try: + subscription_receiver.add_rule("test_rule", "1=1") + msgs = subscription_receiver.receive(max_batch_size=1000) + for msg in msgs: + msg.abandon() + subscription_receiver.remove_rule("test_rule") + session_msgs = subscription_receiver.receive(max_batch_size=1000, session="test_session") + for msg in session_msgs: + msg.complete() + except: + transaction.abort() + + +# Approach for Session APIs - Option 1 +## APIs +### Put session methods on ServiceBusReceiverClient +class ServiceBusReceiverClient: + def list_sessions(self): + def set_session_state(self, session_id, session_state): + def get_session_state(self, session_id): + def renew_lock(self, session_id): + @property + def expired(self, session_id): + + +## Sample Code: +# iterator receive +with subscription_receiver: + session_id = "test_session" + subscription_receiver.set_session_state(session=session_id, session_state="START") + for msg in subscription_receiver(session=session_id): + print(message) + msg.complete() + subscription_receiver.renew_lock(session=session_id) + if str(msg) == "shutdown": + subscription_receiver.set_session_state(session=session_id, session_state="STOP") + break + +# pull mode receive +with subscription_receiver: + session_id = "test_session" + msgs = subscription_receiver.receive(max_batch_size=10, session=session_id) + subscription_receiver.set_session_state(session=session_id, session_state="BEGIN") + for msg in msgs: + msg.complete() + subscription_receiver.renew_lock(session=session_id) + subscription_receiver.set_session_state(session=session_id, session_state="END") + +# receive deferred letter +with subscription_receiver: + if subscription_receiver.get_session_state(session="test_session") == "DONE": + defered_sequence_numbers = [xxxx] + deferred = subscription_receiver.receive_deferred_messages(defered_sequence_numbers, session=session) + subscription_receiver.settle_deferred_messages('completed', deferred) + + +# Approach for Rule APIs - Option 1 +## APIs +### 3-Clients approach + +class ServiceBusSenderClient: +class QueueReceiverClient: +class SubscriptionReceiverClient: + def add_rule(self, rule_name, filter): + def remove_rule(self, rule_name): + def get_rules(self): + + + +# Approach for Rule APIs - Option 2 +## APIs +class ServiceBusReceiverClient: + def get_subscription_rule_manager(self): # raise Error when called on Queue + +class SubscriptionRuleManager: + def add_rule(self, rule_name, filter): + def remove_rule(self, rule_name): + def get_rules(self): + + +# Approach for Rule APIs - Option 3 +## APIs +class SubscriptionRuleManager: + def from_receiver_client_for_subscription(receiver_client): + def add_rule(self, rule_name, filter): + def remove_rule(self, rule_name): + def get_rules(self): From 67d0b0450110a9b4c871ecb49bb0440928e18451 Mon Sep 17 00:00:00 2001 From: Yunhao Ling Date: Tue, 18 Feb 2020 19:05:28 -0800 Subject: [PATCH 02/24] minor update --- sdk/servicebus/azure-servicebus/api_review.py | 9 ++++++--- 1 file changed, 6 insertions(+), 3 deletions(-) diff --git a/sdk/servicebus/azure-servicebus/api_review.py b/sdk/servicebus/azure-servicebus/api_review.py index fa37eadd8148..3de068e781b9 100644 --- a/sdk/servicebus/azure-servicebus/api_review.py +++ b/sdk/servicebus/azure-servicebus/api_review.py @@ -1,4 +1,4 @@ -# Approach 1 +# Two-Clients Approach ## APIs: class SeviceBusSenderClient: @@ -206,11 +206,13 @@ class BatchMessage(Message): class PeekMessage(Message): # inherited from Message, and cannot settle/lock the message + class DeferredMessage(Message): # interited from Message # its own @properties def settled(self) -> bool: # read-only + class ReceiveSettleMode: PeekLock ReceiveAndDelete @@ -226,8 +228,9 @@ def __init__(self, policy, key): def get_token(self, *scopes, **kwargs): -## Sampe Code: +################################################### Split Line ################################################### +## Sampe Code: ### creation of differen entity clients @@ -499,6 +502,7 @@ def get_token(self, *scopes, **kwargs): except: transaction.abort() +################################################### Split Line ################################################### # Approach for Session APIs - Option 1 ## APIs @@ -555,7 +559,6 @@ def remove_rule(self, rule_name): def get_rules(self): - # Approach for Rule APIs - Option 2 ## APIs class ServiceBusReceiverClient: From b598e9d9a30488348e9bead373127c939b0f847d Mon Sep 17 00:00:00 2001 From: Yunhao Ling Date: Tue, 18 Feb 2020 19:06:17 -0800 Subject: [PATCH 03/24] minor update --- sdk/servicebus/azure-servicebus/api_review.py | 6 ++++++ 1 file changed, 6 insertions(+) diff --git a/sdk/servicebus/azure-servicebus/api_review.py b/sdk/servicebus/azure-servicebus/api_review.py index 3de068e781b9..a9d5ea63122a 100644 --- a/sdk/servicebus/azure-servicebus/api_review.py +++ b/sdk/servicebus/azure-servicebus/api_review.py @@ -547,6 +547,8 @@ def expired(self, session_id): subscription_receiver.settle_deferred_messages('completed', deferred) +################################################### Split Line ################################################### + # Approach for Rule APIs - Option 1 ## APIs ### 3-Clients approach @@ -559,6 +561,8 @@ def remove_rule(self, rule_name): def get_rules(self): +################################################### Split Line ################################################### + # Approach for Rule APIs - Option 2 ## APIs class ServiceBusReceiverClient: @@ -570,6 +574,8 @@ def remove_rule(self, rule_name): def get_rules(self): +################################################### Split Line ################################################### + # Approach for Rule APIs - Option 3 ## APIs class SubscriptionRuleManager: From 979f2ec9b371d7b1c437b09a88340d5c67626019 Mon Sep 17 00:00:00 2001 From: Yunhao Ling Date: Tue, 18 Feb 2020 19:23:14 -0800 Subject: [PATCH 04/24] add mode to receive method --- sdk/servicebus/azure-servicebus/api_review.py | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/sdk/servicebus/azure-servicebus/api_review.py b/sdk/servicebus/azure-servicebus/api_review.py index a9d5ea63122a..bf54f9e26b0c 100644 --- a/sdk/servicebus/azure-servicebus/api_review.py +++ b/sdk/servicebus/azure-servicebus/api_review.py @@ -125,7 +125,8 @@ def receive( self, session : Union[str, Session] = None, max_batch_size : int = None, - timeout : float = None + timeout : float = None, + mode : ReceiveSettleMode =ReceiveSettleMode.PeekLock ) -> List[Message]: # Pull mode receive def settle_deferred_messages( self, From d853ca49e00ab280252e44f278a4c1a35c19f5c7 Mon Sep 17 00:00:00 2001 From: Yunhao Ling Date: Wed, 19 Feb 2020 11:23:24 -0800 Subject: [PATCH 05/24] put session into constructor --- sdk/servicebus/azure-servicebus/api_review.py | 72 +++++++++---------- 1 file changed, 36 insertions(+), 36 deletions(-) diff --git a/sdk/servicebus/azure-servicebus/api_review.py b/sdk/servicebus/azure-servicebus/api_review.py index bf54f9e26b0c..76fa50d98357 100644 --- a/sdk/servicebus/azure-servicebus/api_review.py +++ b/sdk/servicebus/azure-servicebus/api_review.py @@ -48,14 +48,14 @@ def list_session_ids(self) -> List[str]: def send( self, messages, - session_id : str = None, + session : str = None, message_timeout : float = None ) -> None: def schedule( self, messages, schedule_time : datetime, - session_id : str = None + session : str = None ) -> List[int]: def cancel_scheduled_messages(self, sequence_numbers : List[int]) -> None: @@ -66,6 +66,7 @@ def __init__( fully_qualified_namespace : str, entity_name : str, subscription_name : str = None, + session : str = None, credential : TokenCredential, logging_enable: bool = False, http_proxy: dict = None, @@ -78,6 +79,7 @@ def from_queue( fully_qualified_namespace : str, queue_name : str, credential : TokenCredential, + session: str = None, **kwargs ) -> ServiceBusReceiverClient: @classmethod @@ -87,6 +89,7 @@ def from_topic_subscription( topic_name : str, subscription_name : str, credential: TokenCredential, + session: str = None, **kwargs ) -> ServiceBusReceiverClient: @classmethod @@ -95,6 +98,7 @@ def from_connection_string( conn_str : str, entity_name : str = None, subscription_name : str = None, + session: str = None, **kwargs )-> ServiceBusReceiverClient: @@ -111,19 +115,16 @@ def next(self): def peek( self, - session : Union[str, Session] = None, count : int = 1, start_from : int = None ) -> List[PeekMessage]: def receive_deferred_messages( self, - session : Union[str, Session] = None, sequence_numbers : List[int], mode : ReceiveSettleMode =ReceiveSettleMode.PeekLock ) -> List[DeferredMessage]: def receive( self, - session : Union[str, Session] = None, max_batch_size : int = None, timeout : float = None, mode : ReceiveSettleMode =ReceiveSettleMode.PeekLock @@ -135,8 +136,8 @@ def settle_deferred_messages( **kwargs ) -> None: # Batch settle deferred messages - def list_session_ids(self) -> List[str]: - def get_session(self, session_id : str = "NEXT_AVAILABLE") -> Session: # raise Error when called on non-session entity + def list_sessions(self) -> List[str]: + def get_session(self) -> Session: # raise Error when called on non-session entity # Rule APIs, raise Error when called on Queue def add_rule(self, rule_name : str, filter : str) -> None: # TODO: figure out what filter is @@ -268,7 +269,8 @@ def get_token(self, *scopes, **kwargs): queue_receiver = SeviceBusReceiverClient( fully_qualified_namepsace="fully_qualified_namepsace", entity_name="queue_name", - credential=ServiceBusSharedKeyCredential("policy", "key") + credential=ServiceBusSharedKeyCredential("policy", "key"), + session="test_session" ) # Topic Sender @@ -307,7 +309,8 @@ def get_token(self, *scopes, **kwargs): fully_qualified_namepsace="fully_qualified_namespace", entity_name="topic_name", subscription_name="subcription_name", - credential=ServiceBusSharedKeyCredential("policy", "key") + credential=ServiceBusSharedKeyCredential("policy", "key"), + session="test_session" ) @@ -424,21 +427,19 @@ def get_token(self, *scopes, **kwargs): fully_qualified_namespace="fully_qualified_namespace", topic_name="topic_name", subscription_name="subscription_name", - credential=ServiceBusSharedKeyCredential("policy", "key") + credential=ServiceBusSharedKeyCredential("policy", "key"), + session="test_session" ) # 2.1 peek with subscription_receiver: - msgs = subscription_receiver.peek(count=10, session="test_session") + msgs = subscription_receiver.peek(count=10) # 2.2 iterator receive with subscription_receiver: - session = subscription_receiver.get_session("test_session") - # Option 2: - # session_ids = subscription_receiver.list_session() - # session = subscription_receiver.get_session(session_ids[0]) + session = subscription_receiver.get_session() session.set_session_state("START") - for msg in subscription_receiver(session=session): + for msg in subscription_receiver: print(message) msg.complete() session.renew_lock() @@ -448,8 +449,8 @@ def get_token(self, *scopes, **kwargs): # 2.3 pull mode receive with subscription_receiver: - session = subscription_receiver.get_session("test_session") - msgs = subscription_receiver.receive(max_batch_size=10, session=session) + session = subscription_receiver.get_session() + msgs = subscription_receiver.receive(max_batch_size=10) session.set_session_state("BEGIN") for msg in msgs: msg.complete() @@ -458,10 +459,10 @@ def get_token(self, *scopes, **kwargs): # 2.4 receive deferred letter with subscription_receiver: - session = subscription_receiver.get_session("test_session") + session = subscription_receiver.get_session() if session.get_session_state() == "DONE": defered_sequence_numbers = [xxxx] - deferred = subscription_receiver.receive_deferred_messages(defered_sequence_numbers, session=session) + deferred = subscription_receiver.receive_deferred_messages(defered_sequence_numbers) subscription_receiver.settle_deferred_messages('completed', deferred) ### Rule operations @@ -497,7 +498,7 @@ def get_token(self, *scopes, **kwargs): for msg in msgs: msg.abandon() subscription_receiver.remove_rule("test_rule") - session_msgs = subscription_receiver.receive(max_batch_size=1000, session="test_session") + session_msgs = subscription_receiver.receive(max_batch_size=1000) for msg in session_msgs: msg.complete() except: @@ -510,41 +511,40 @@ def get_token(self, *scopes, **kwargs): ### Put session methods on ServiceBusReceiverClient class ServiceBusReceiverClient: def list_sessions(self): - def set_session_state(self, session_id, session_state): - def get_session_state(self, session_id): - def renew_lock(self, session_id): + def set_session_state(self, session_state): + def get_session_state(self): + def renew_lock(self): @property - def expired(self, session_id): + def expired(self): ## Sample Code: # iterator receive with subscription_receiver: - session_id = "test_session" - subscription_receiver.set_session_state(session=session_id, session_state="START") - for msg in subscription_receiver(session=session_id): + subscription_receiver.set_session_state(session_state="START") + for msg in subscription_receiver: print(message) msg.complete() - subscription_receiver.renew_lock(session=session_id) + subscription_receiver.renew_lock() if str(msg) == "shutdown": - subscription_receiver.set_session_state(session=session_id, session_state="STOP") + subscription_receiver.set_session_state(session_state="STOP") break # pull mode receive with subscription_receiver: session_id = "test_session" - msgs = subscription_receiver.receive(max_batch_size=10, session=session_id) - subscription_receiver.set_session_state(session=session_id, session_state="BEGIN") + msgs = subscription_receiver.receive(max_batch_size=10) + subscription_receiver.set_session_state(session_state="BEGIN") for msg in msgs: msg.complete() - subscription_receiver.renew_lock(session=session_id) - subscription_receiver.set_session_state(session=session_id, session_state="END") + subscription_receiver.renew_lock() + subscription_receiver.set_session_state(session_state="END") # receive deferred letter with subscription_receiver: - if subscription_receiver.get_session_state(session="test_session") == "DONE": + if subscription_receiver.get_session_state() == "DONE": defered_sequence_numbers = [xxxx] - deferred = subscription_receiver.receive_deferred_messages(defered_sequence_numbers, session=session) + deferred = subscription_receiver.receive_deferred_messages(defered_sequence_numbers) subscription_receiver.settle_deferred_messages('completed', deferred) From 0d83238b4d7648e43e876bd55a02466ab1896d1e Mon Sep 17 00:00:00 2001 From: Yunhao Ling Date: Wed, 19 Feb 2020 15:28:29 -0800 Subject: [PATCH 06/24] move receive mode param into constructor --- sdk/servicebus/azure-servicebus/api_review.py | 8 +++++--- 1 file changed, 5 insertions(+), 3 deletions(-) diff --git a/sdk/servicebus/azure-servicebus/api_review.py b/sdk/servicebus/azure-servicebus/api_review.py index 76fa50d98357..5d57500ed214 100644 --- a/sdk/servicebus/azure-servicebus/api_review.py +++ b/sdk/servicebus/azure-servicebus/api_review.py @@ -65,9 +65,10 @@ def __init__( self, fully_qualified_namespace : str, entity_name : str, + credential: TokenCredential, subscription_name : str = None, session : str = None, - credential : TokenCredential, + mode : ReceiveSettleMode = ReceiveSettleMode.PeekLock logging_enable: bool = False, http_proxy: dict = None, transport_type: TransportType = None, @@ -80,6 +81,7 @@ def from_queue( queue_name : str, credential : TokenCredential, session: str = None, + mode : ReceiveSettleMode = ReceiveSettleMode.PeekLock **kwargs ) -> ServiceBusReceiverClient: @classmethod @@ -90,6 +92,7 @@ def from_topic_subscription( subscription_name : str, credential: TokenCredential, session: str = None, + mode : ReceiveSettleMode = ReceiveSettleMode.PeekLock **kwargs ) -> ServiceBusReceiverClient: @classmethod @@ -99,6 +102,7 @@ def from_connection_string( entity_name : str = None, subscription_name : str = None, session: str = None, + mode : ReceiveSettleMode = ReceiveSettleMode.PeekLock **kwargs )-> ServiceBusReceiverClient: @@ -121,13 +125,11 @@ def peek( def receive_deferred_messages( self, sequence_numbers : List[int], - mode : ReceiveSettleMode =ReceiveSettleMode.PeekLock ) -> List[DeferredMessage]: def receive( self, max_batch_size : int = None, timeout : float = None, - mode : ReceiveSettleMode =ReceiveSettleMode.PeekLock ) -> List[Message]: # Pull mode receive def settle_deferred_messages( self, From 18aa2353d7fe1c3673274c0bbd8c23d26ed129b4 Mon Sep 17 00:00:00 2001 From: Yunhao Ling Date: Mon, 24 Feb 2020 23:04:27 -0800 Subject: [PATCH 07/24] fix typo --- sdk/servicebus/azure-servicebus/api_review.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/sdk/servicebus/azure-servicebus/api_review.py b/sdk/servicebus/azure-servicebus/api_review.py index 5d57500ed214..3d099529f574 100644 --- a/sdk/servicebus/azure-servicebus/api_review.py +++ b/sdk/servicebus/azure-servicebus/api_review.py @@ -183,8 +183,8 @@ def partition_key(self, value : str): def partition_key(self) -> str: def via_partition_key(self, value: str): def via_partition_key(self) -> str: - def sessio_id(self, value : str): - def sessio_id(self) -> str: + def session_id(self, value : str): + def session_id(self) -> str: def time_to_live(self, value : Union[float, timedelta]): def time_to_live(self) -> datetime.timedelta: def annotations(self, value : dict[str, Any]): From 3b914d5717eb1e3a1420b3b878f90f683ff85bdd Mon Sep 17 00:00:00 2001 From: Yunhao Ling Date: Tue, 25 Feb 2020 09:59:59 -0800 Subject: [PATCH 08/24] minor update on send --- sdk/servicebus/azure-servicebus/api_review.py | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/sdk/servicebus/azure-servicebus/api_review.py b/sdk/servicebus/azure-servicebus/api_review.py index 3d099529f574..8c336d0367cf 100644 --- a/sdk/servicebus/azure-servicebus/api_review.py +++ b/sdk/servicebus/azure-servicebus/api_review.py @@ -47,17 +47,17 @@ def list_session_ids(self) -> List[str]: def send( self, - messages, + message : Union[Message,BatchMessage], session : str = None, message_timeout : float = None ) -> None: def schedule( self, - messages, + message : Union[Message,BatchMessage], schedule_time : datetime, session : str = None ) -> List[int]: - def cancel_scheduled_messages(self, sequence_numbers : List[int]) -> None: + def cancel_scheduled_messages(self, sequence_number : Union[int, List[int]]) -> None: class ServiceBusReceiverClient: @@ -138,7 +138,7 @@ def settle_deferred_messages( **kwargs ) -> None: # Batch settle deferred messages - def list_sessions(self) -> List[str]: + def list_session_ids(self) -> List[str]: def get_session(self) -> Session: # raise Error when called on non-session entity # Rule APIs, raise Error when called on Queue From a724188236fcd3a560d500e44d38a7dd71afb328 Mon Sep 17 00:00:00 2001 From: Yunhao Ling Date: Tue, 25 Feb 2020 23:52:21 -0800 Subject: [PATCH 09/24] add batch message, retry option, sb connection --- sdk/servicebus/azure-servicebus/api_review.py | 323 +++--------------- .../azure-servicebus/api_samples.py | 306 +++++++++++++++++ 2 files changed, 344 insertions(+), 285 deletions(-) create mode 100644 sdk/servicebus/azure-servicebus/api_samples.py diff --git a/sdk/servicebus/azure-servicebus/api_review.py b/sdk/servicebus/azure-servicebus/api_review.py index 8c336d0367cf..aa3e034adda4 100644 --- a/sdk/servicebus/azure-servicebus/api_review.py +++ b/sdk/servicebus/azure-servicebus/api_review.py @@ -11,6 +11,9 @@ def __init__( http_proxy: dict = None, transport_type: TransportType = None, retry_total : int = 3, + retry_backoff_factor : float = 30, + retry_backoff_maximum : int = 120 # cur_retry_backoff = retry_backoff_factor * (2^cur_retry_time) + connection : ServiceBusConnection = None ) -> None: @classmethod def from_queue( @@ -45,18 +48,25 @@ def get_properties(self) -> Dict[str, Any]: def list_session_ids(self) -> List[str]: + def create_batch( + self, + session_id : str = None, + max_size_in_bytes : int = None + ) -> BatchMessage: + def send( self, - message : Union[Message,BatchMessage], - session : str = None, + message : Union[Message, BatchMessage], + session_id : str = None, message_timeout : float = None ) -> None: def schedule( self, - message : Union[Message,BatchMessage], + message : Union[Message, BatchMessage], schedule_time : datetime, - session : str = None + session_id : str = None ) -> List[int]: + def cancel_scheduled_messages(self, sequence_number : Union[int, List[int]]) -> None: @@ -67,12 +77,15 @@ def __init__( entity_name : str, credential: TokenCredential, subscription_name : str = None, - session : str = None, + session_id : str = None, mode : ReceiveSettleMode = ReceiveSettleMode.PeekLock logging_enable: bool = False, http_proxy: dict = None, transport_type: TransportType = None, retry_total : int = 3, + retry_backoff_factor: float = 30, + retry_backoff_maximum: int = 120 # cur_retry_backoff = retry_backoff_factor * (2^cur_retry_time) + connection: Connection = None ) -> None: @classmethod def from_queue( @@ -80,7 +93,7 @@ def from_queue( fully_qualified_namespace : str, queue_name : str, credential : TokenCredential, - session: str = None, + session_id: str = None, mode : ReceiveSettleMode = ReceiveSettleMode.PeekLock **kwargs ) -> ServiceBusReceiverClient: @@ -91,7 +104,7 @@ def from_topic_subscription( topic_name : str, subscription_name : str, credential: TokenCredential, - session: str = None, + session_id: str = None, mode : ReceiveSettleMode = ReceiveSettleMode.PeekLock **kwargs ) -> ServiceBusReceiverClient: @@ -101,7 +114,7 @@ def from_connection_string( conn_str : str, entity_name : str = None, subscription_name : str = None, - session: str = None, + session_id: str = None, mode : ReceiveSettleMode = ReceiveSettleMode.PeekLock **kwargs )-> ServiceBusReceiverClient: @@ -130,6 +143,7 @@ def receive( self, max_batch_size : int = None, timeout : float = None, + auto_complete : bool = False ) -> List[Message]: # Pull mode receive def settle_deferred_messages( self, @@ -139,7 +153,7 @@ def settle_deferred_messages( ) -> None: # Batch settle deferred messages def list_session_ids(self) -> List[str]: - def get_session(self) -> Session: # raise Error when called on non-session entity + def get_session(self) -> ServiceBusSession: # raise Error when called on non-session entity # Rule APIs, raise Error when called on Queue def add_rule(self, rule_name : str, filter : str) -> None: # TODO: figure out what filter is @@ -147,7 +161,7 @@ def remove_rule(self, rule_name : str) -> None: def get_rules(self) -> List[str] -> Dict[str, Any]: -class Session: +class ServiceBusSession: def get_session_id(self) -> str: def get_session_state(self) -> str: # This can be made to property if property is more pythonic way def set_session_state(self, state : str) -> None: # This can be made to property if property is more pythonic way @@ -205,6 +219,7 @@ def defer(self) -> None: class BatchMessage(Message): # inherited from Message + def try_add(self, message: Message) -> None: class PeekMessage(Message): @@ -227,285 +242,23 @@ class TransportType(Enum): AmqpOverWebsocket +class ServiceBusConnection: + def __int__( + self, + fully_qualified_namespace, + entity_name, + token_credential + ) -> None: + def from_connection_string( + self, + conn_str, + entity_name=None + ) -> ServiceBusConnection: + class ServiceBusSharedKeyCredential: def __init__(self, policy, key): def get_token(self, *scopes, **kwargs): - -################################################### Split Line ################################################### - -## Sampe Code: - -### creation of differen entity clients - -# Queue Sender -queue_sender = SeviceBusSenderClient.from_queue( - fully_qualified_namepsace="fully_qualified_namepsace", - queue_name="queue_name", - credential=ServiceBusSharedKeyCredential("policy", "key") -) - -queue_sender = SeviceBusSenderClient.from_connection_string( - conn_str="conn_str", - entity_name="queue_name" -) - -queue_sender = SeviceBusSenderClient( - fully_qualified_namepsace="fully_qualified_namepsace", - entity_name="queue_name", - ServiceBusSharedKeyCredential("policy", "key") -) - -# Queue Receiver -queue_receiver = SeviceBusReceiverClient.from_queue( - fully_qualified_namepsace="fully_qualified_namepsace", - queue_name="queue_name", - credential=ServiceBusSharedKeyCredential("policy", "key") -) - -queue_receiver = SeviceBusReceiverClient.from_connection_string( - conn_str="conn_str", - entity_name="queue_name" -) - -queue_receiver = SeviceBusReceiverClient( - fully_qualified_namepsace="fully_qualified_namepsace", - entity_name="queue_name", - credential=ServiceBusSharedKeyCredential("policy", "key"), - session="test_session" -) - -# Topic Sender -topic_sender = SeviceBusSenderClient.from_topic( - fully_qualified_namepsace="fully_qualified_namepsace", - topic_name="topic_name", - credential=ServiceBusSharedKeyCredential("policy", "key") -) - -topic_sender = SeviceBusSenderClient.from_connection_string( - conn_str="conn_str", - entity_name="topic_name" -) - -topic_sender = SeviceBusSenderClient( - fully_qualified_namepsace="fully_qualified_namepsace", - entity_name="topic_name", - credential=ServiceBusSharedKeyCredential("policy", "key") -) - -# Subscription Receiver -subscription_receiver = SeviceBusReceiverClient.from_topic_subscription( - fully_qualified_namepsace="fully_qualified_namespace", - topic_name="topic_name", - subscription_name="subscription_name", - credential=ServiceBusSharedKeyCredential("policy", "key") -) - -subscription_receiver = SeviceBusReceiverClient.from_connection_string( - conn_str="conn_str", - entity_name="topic_name", - subscription_name="subscription_name" -) - -subscription_receiver = SeviceBusReceiverClient( - fully_qualified_namepsace="fully_qualified_namespace", - entity_name="topic_name", - subscription_name="subcription_name", - credential=ServiceBusSharedKeyCredential("policy", "key"), - session="test_session" -) - - -### Send and receive -# 1 Send - -queue_sender = SeviceBusSenderClient.from_queue( - fully_qualified_namepsace="fully_qualified_namepsace", - queue_name="queue_name", - credential=ServiceBusSharedKeyCredential("policy", "key") -) - -''' -topic_sender = SeviceBusSenderClient.from_queue( - fully_qualified_namepsace="fully_qualified_namepsace", - topic_name="topic_name", - credential=ServiceBusSharedKeyCredential("policy", "key") -) -''' - -with queue_sender: - # 1.1 send single - message = Message("Test message") - queue_sender.send(message) - - # 1.2 send list - messages = [Message("Test message {}".format(i)) for i in range(10)] - queue_sender.send(messages) - - # 1.3 schedule - schedule_message = Message("scheduled_message") - enqueue_time = datetime.utcnow() + timedelta(minutes=10) - sequence_number = queue_sender.schedule(schedule_message, enqueue_time) - - # 1.4 cancel schedule - queue_sender.cancel_scheduled_messages(*sequence_number) - -# 2 Receive - -queue_receiver = ServiceBusReceiverClient.from_queue( - fully_qualified_namespace="fully_qualified_namespace", - queue_name="queue_name", - credential=ServiceBusSharedKeyCredential("policy", "key") -) - -''' -subscription_receiver = ServiceBusReceiverClient.from_queue( - fully_qualified_namespace="fully_qualified_namespace", - topic_name="topic_name", - subscription_name="subscription_name" - credential=ServiceBusSharedKeyCredential("policy", "key") -) -''' - -# 2.1 peek -with queue_receiver: - msgs = queue_receiver.peek(count=10) - -# 2.2 iterator receive -with queue_receiver: - for msg in queue_receiver: - print(message) - msg.complete() - # msg.renew_lock() - # msg.abondon() - # msg.defer() - # msg.dead_letter() - -# 2.3 pull mode receive -with queue_receiver: - msgs = queue_receiver.receive(max_batch_size=10) - for msg in msgs: - msg.complete() - -# 2.4 receive deferred letter -with queue_receiver: - defered_sequence_numbers = [xxxx] - deferred = queue_receiver.receive_deferred_messages(defered_sequence_numbers) - queue_receiver.settle_deferred_messages('completed', deferred) - - -### Session operation - - -# 1 Send -topic_sender = SeviceBusSenderClient.from_topic -( - fully_qualified_namespace="fully_qualified_namepsace", - topic_name="topic_name", - credential=ServiceBusSharedKeyCredential("policy", "key") -) - -with topic_sender: - # 1.1 send single - message = Message("Test message") - # message.session_id = "test_session" - topic_sender.send(message, session="test_session") - - # 1.2 send list - messages = [Message("Test message {}".format(i)) for i in range(10)] - topic_sender.send(messages, session="test_session") - - # 1.3 schedule - schedule_message = Message("scheduled_message") - enqueue_time = datetime.utcnow() + timedelta(minutes=10) - sequence_number = topic_sender.schedule(schedule_message, enqueue_time, session="test_session") - - # 1.4 cancel schedule - topic_sender.cancel_scheduled_messages(*sequence_number) - -# 2 Receive -subscription_receiver = ServiceBusReceiverClient.from_topic_subscription -( - fully_qualified_namespace="fully_qualified_namespace", - topic_name="topic_name", - subscription_name="subscription_name", - credential=ServiceBusSharedKeyCredential("policy", "key"), - session="test_session" -) - -# 2.1 peek -with subscription_receiver: - msgs = subscription_receiver.peek(count=10) - -# 2.2 iterator receive -with subscription_receiver: - session = subscription_receiver.get_session() - session.set_session_state("START") - for msg in subscription_receiver: - print(message) - msg.complete() - session.renew_lock() - if str(msg) == "shutdown": - session.set_session_state("STOP") - break - -# 2.3 pull mode receive -with subscription_receiver: - session = subscription_receiver.get_session() - msgs = subscription_receiver.receive(max_batch_size=10) - session.set_session_state("BEGIN") - for msg in msgs: - msg.complete() - session.renew_lock() - session.set_session_state("END") - -# 2.4 receive deferred letter -with subscription_receiver: - session = subscription_receiver.get_session() - if session.get_session_state() == "DONE": - defered_sequence_numbers = [xxxx] - deferred = subscription_receiver.receive_deferred_messages(defered_sequence_numbers) - subscription_receiver.settle_deferred_messages('completed', deferred) - -### Rule operations - -subscription_receiver = ServiceBusReceiverClient.from_topic_subscription -( - fully_qualified_namespace="fully_qualified_namespace", - topic_name="topic_name", - subscription_name="subscription_name", - credential=ServiceBusSharedKeyCredential("policy", "key")) - -with subscription_receiver: - subscription_receiver.get_rules() - subscription_receiver.add_rule("test_rule", "1=1") - subscription_receiver.remove_rule("test_rule") - -### Transaction operations - -subscription_receiver = ServiceBusReceiverClient.from_topic_subscription -( - fully_qualified_namespace="fully_qualified_namespace", - topic_name="topic_name", - subscription_name="subscription_name", - credential=ServiceBusSharedKeyCredential("policy", "key") -) - -with subscription_receiver: - transaction = subscription_receiver.get_transaction() # TODO, no idea on how to get a transaction for now - with transaction: - try: - subscription_receiver.add_rule("test_rule", "1=1") - msgs = subscription_receiver.receive(max_batch_size=1000) - for msg in msgs: - msg.abandon() - subscription_receiver.remove_rule("test_rule") - session_msgs = subscription_receiver.receive(max_batch_size=1000) - for msg in session_msgs: - msg.complete() - except: - transaction.abort() - ################################################### Split Line ################################################### # Approach for Session APIs - Option 1 diff --git a/sdk/servicebus/azure-servicebus/api_samples.py b/sdk/servicebus/azure-servicebus/api_samples.py new file mode 100644 index 000000000000..6f98d443253e --- /dev/null +++ b/sdk/servicebus/azure-servicebus/api_samples.py @@ -0,0 +1,306 @@ + +## Sampe Code: + +### creation of differen entity clients + +# Queue Sender +queue_sender = SeviceBusSenderClient.from_queue( + fully_qualified_namepsace="fully_qualified_namepsace", + queue_name="queue_name", + credential=ServiceBusSharedKeyCredential("policy", "key") +) + +queue_sender = SeviceBusSenderClient.from_connection_string( + conn_str="conn_str", + entity_name="queue_name" +) + +queue_sender = SeviceBusSenderClient( + fully_qualified_namepsace="fully_qualified_namepsace", + entity_name="queue_name", + ServiceBusSharedKeyCredential("policy", "key") +) + +# Queue Receiver +queue_receiver = SeviceBusReceiverClient.from_queue( + fully_qualified_namepsace="fully_qualified_namepsace", + queue_name="queue_name", + credential=ServiceBusSharedKeyCredential("policy", "key") +) + +queue_receiver = SeviceBusReceiverClient.from_connection_string( + conn_str="conn_str", + entity_name="queue_name" +) + +queue_receiver = SeviceBusReceiverClient( + fully_qualified_namepsace="fully_qualified_namepsace", + entity_name="queue_name", + credential=ServiceBusSharedKeyCredential("policy", "key"), + session_id="test_session" +) + +# Topic Sender +topic_sender = SeviceBusSenderClient.from_topic( + fully_qualified_namepsace="fully_qualified_namepsace", + topic_name="topic_name", + credential=ServiceBusSharedKeyCredential("policy", "key") +) + +topic_sender = SeviceBusSenderClient.from_connection_string( + conn_str="conn_str", + entity_name="topic_name" +) + +topic_sender = SeviceBusSenderClient( + fully_qualified_namepsace="fully_qualified_namepsace", + entity_name="topic_name", + credential=ServiceBusSharedKeyCredential("policy", "key") +) + +# Subscription Receiver +subscription_receiver = SeviceBusReceiverClient.from_topic_subscription( + fully_qualified_namepsace="fully_qualified_namespace", + topic_name="topic_name", + subscription_name="subscription_name", + credential=ServiceBusSharedKeyCredential("policy", "key") +) + +subscription_receiver = SeviceBusReceiverClient.from_connection_string( + conn_str="conn_str", + entity_name="topic_name", + subscription_name="subscription_name" +) + +subscription_receiver = SeviceBusReceiverClient( + fully_qualified_namepsace="fully_qualified_namespace", + entity_name="topic_name", + subscription_name="subcription_name", + credential=ServiceBusSharedKeyCredential("policy", "key"), + session_id="test_session" +) + + +### Send and receive +# 1 Send + +queue_sender = SeviceBusSenderClient.from_queue( + fully_qualified_namepsace="fully_qualified_namepsace", + queue_name="queue_name", + credential=ServiceBusSharedKeyCredential("policy", "key") +) + +''' +topic_sender = SeviceBusSenderClient.from_queue( + fully_qualified_namepsace="fully_qualified_namepsace", + topic_name="topic_name", + credential=ServiceBusSharedKeyCredential("policy", "key") +) +''' + +with queue_sender: + # 1.1 send single + message = Message("Test message") + queue_sender.send(message) + + # 1.2 send batch + batch_message = topic_sender.create_batch(max_size_in_bytes=256 * 1204) + while True: + try: + batch_message.try_add(Message("Test message")) + except ValueError: + break + + topic_sender.send(batch_message) + + # 1.3 schedule + schedule_message = Message("scheduled_message") + + enqueue_time = datetime.utcnow() + timedelta(minutes=10) + sequence_number = topic_sender.schedule(schedule_message, enqueue_time) + + batch_schedule_message = topic_sender.create_batch(max_size_in_bytes=256 * 1204) + while True: + try: + batch_message.try_add(Message("Test message")) + except ValueError: + break + + batch_sequence_numbers = topic_sender.schedule(batch_schedule_message, enqueue_time) + + # 1.4 cancel schedule + topic_sender.cancel_scheduled_messages(sequence_number) + topic_sender.cancel_scheduled_messages(batch_sequence_numbers) + +# 2 Receive + +queue_receiver = ServiceBusReceiverClient.from_queue( + fully_qualified_namespace="fully_qualified_namespace", + queue_name="queue_name", + credential=ServiceBusSharedKeyCredential("policy", "key") +) + +''' +subscription_receiver = ServiceBusReceiverClient.from_queue( + fully_qualified_namespace="fully_qualified_namespace", + topic_name="topic_name", + subscription_name="subscription_name" + credential=ServiceBusSharedKeyCredential("policy", "key") +) +''' + +# 2.1 peek +with queue_receiver: + msgs = queue_receiver.peek(count=10) + +# 2.2 iterator receive +with queue_receiver: + for msg in queue_receiver: + print(message) + msg.complete() + # msg.renew_lock() + # msg.abondon() + # msg.defer() + # msg.dead_letter() + +# 2.3 pull mode receive +with queue_receiver: + msgs = queue_receiver.receive(max_batch_size=10) + for msg in msgs: + msg.complete() + +# 2.4 receive deferred letter +with queue_receiver: + defered_sequence_numbers = [xxxx] + deferred = queue_receiver.receive_deferred_messages(defered_sequence_numbers) + queue_receiver.settle_deferred_messages('completed', deferred) + + +### Session operation + + +# 1 Send +topic_sender = SeviceBusSenderClient.from_topic +( + fully_qualified_namespace="fully_qualified_namepsace", + topic_name="topic_name", + credential=ServiceBusSharedKeyCredential("policy", "key") +) + +with topic_sender: + # 1.1 send single + message = Message("Test message") + # message.session_id = "test_session" + topic_sender.send(message, session_id="test_session") + + # 1.2 send batch + batch_message = topic_sender.create_batch(max_size_in_bytes=256 * 1204, seesion_id="test_session") + while True: + try: + batch_message.try_add(Message("Test message")) + except ValueError: + break + + topic_sender.send(batch_message) + + # 1.3 schedule + schedule_message = Message("scheduled_message") + + enqueue_time = datetime.utcnow() + timedelta(minutes=10) + sequence_number = topic_sender.schedule(schedule_message, enqueue_time, session_id="test_session") + + batch_schedule_message = topic_sender.create_batch(max_size_in_bytes=256 * 1204) + while True: + try: + batch_message.try_add(Message("Test message")) + except ValueError: + break + + batch_sequence_numbers = topic_sender.schedule(batch_schedule_message, enqueue_time) + + # 1.4 cancel schedule + topic_sender.cancel_scheduled_messages(sequence_number) + topic_sender.cancel_scheduled_messages(batch_sequence_numbers) + +# 2 Receive +subscription_receiver = ServiceBusReceiverClient.from_topic_subscription +( + fully_qualified_namespace="fully_qualified_namespace", + topic_name="topic_name", + subscription_name="subscription_name", + credential=ServiceBusSharedKeyCredential("policy", "key"), + session_id="test_session" +) + +# 2.1 peek +with subscription_receiver: + msgs = subscription_receiver.peek(count=10) + +# 2.2 iterator receive +with subscription_receiver: + session = subscription_receiver.get_session() + session.set_session_state("START") + for msg in subscription_receiver: + print(message) + msg.complete() + session.renew_lock() + if str(msg) == "shutdown": + session.set_session_state("STOP") + break + +# 2.3 pull mode receive +with subscription_receiver: + session = subscription_receiver.get_session() + msgs = subscription_receiver.receive(max_batch_size=10) + session.set_session_state("BEGIN") + for msg in msgs: + msg.complete() + session.renew_lock() + session.set_session_state("END") + +# 2.4 receive deferred letter +with subscription_receiver: + session = subscription_receiver.get_session() + if session.get_session_state() == "DONE": + defered_sequence_numbers = [1,2,3,4,5] + deferred = subscription_receiver.receive_deferred_messages(defered_sequence_numbers) + subscription_receiver.settle_deferred_messages('completed', deferred) + +### Rule operations + +subscription_receiver = ServiceBusReceiverClient.from_topic_subscription +( + fully_qualified_namespace="fully_qualified_namespace", + topic_name="topic_name", + subscription_name="subscription_name", + credential=ServiceBusSharedKeyCredential("policy", "key")) + +with subscription_receiver: + subscription_receiver.get_rules() + subscription_receiver.add_rule("test_rule", "1=1") + subscription_receiver.remove_rule("test_rule") + +### Transaction operations + +subscription_receiver = ServiceBusReceiverClient.from_topic_subscription +( + fully_qualified_namespace="fully_qualified_namespace", + topic_name="topic_name", + subscription_name="subscription_name", + credential=ServiceBusSharedKeyCredential("policy", "key") +) + +with subscription_receiver: + transaction = subscription_receiver.get_transaction() # TODO, no idea on how to get a transaction for now + with transaction: + try: + subscription_receiver.add_rule("test_rule", "1=1") + msgs = subscription_receiver.receive(max_batch_size=1000) + for msg in msgs: + msg.abandon() + subscription_receiver.remove_rule("test_rule") + session_msgs = subscription_receiver.receive(max_batch_size=1000) + for msg in session_msgs: + msg.complete() + except: + transaction.abort() \ No newline at end of file From bf5ece0112cb0f0fc27c936dbe5cb7889a67172f Mon Sep 17 00:00:00 2001 From: Yunhao Ling Date: Wed, 26 Feb 2020 23:36:01 -0800 Subject: [PATCH 10/24] update as per discussion --- sdk/servicebus/azure-servicebus/api_review.py | 144 ++++++++---------- .../azure-servicebus/api_samples.py | 61 +++++--- 2 files changed, 103 insertions(+), 102 deletions(-) diff --git a/sdk/servicebus/azure-servicebus/api_review.py b/sdk/servicebus/azure-servicebus/api_review.py index aa3e034adda4..adcccd6f4e36 100644 --- a/sdk/servicebus/azure-servicebus/api_review.py +++ b/sdk/servicebus/azure-servicebus/api_review.py @@ -46,8 +46,6 @@ def close(self) -> None: def get_properties(self) -> Dict[str, Any]: - def list_session_ids(self) -> List[str]: - def create_batch( self, session_id : str = None, @@ -63,7 +61,7 @@ def send( def schedule( self, message : Union[Message, BatchMessage], - schedule_time : datetime, + schedule_time_utc : datetime, session_id : str = None ) -> List[int]: @@ -143,17 +141,12 @@ def receive( self, max_batch_size : int = None, timeout : float = None, - auto_complete : bool = False - ) -> List[Message]: # Pull mode receive - def settle_deferred_messages( - self, - settlement : str, # TODO: In T1 settlement is just string like 'completed', 'suspended', 'abandoned', can we improve this parameter to be Enum or remove this method? - messages : List[DeferredMessage], - **kwargs - ) -> None: # Batch settle deferred messages + ) -> List[ReceivedMessage]: # Pull mode receive + + + @property + def session(self) -> ServiceBusSession: - def list_session_ids(self) -> List[str]: - def get_session(self) -> ServiceBusSession: # raise Error when called on non-session entity # Rule APIs, raise Error when called on Queue def add_rule(self, rule_name : str, filter : str) -> None: # TODO: figure out what filter is @@ -162,11 +155,12 @@ def get_rules(self) -> List[str] -> Dict[str, Any]: class ServiceBusSession: - def get_session_id(self) -> str: + @property def get_session_state(self) -> str: # This can be made to property if property is more pythonic way def set_session_state(self, state : str) -> None: # This can be made to property if property is more pythonic way def renew_lock(self) -> None: @property + def session_id(self) -> str: def expired(self) -> bool: @@ -184,14 +178,6 @@ def __init__(self, body : str, encoding : str = 'UTF-8', **kwargs) -> None: def __str__(self): # @properties - def settled(self) -> bool: # read-only - def enqueued_time(self) -> datetime: # read-only - def scheduled_enqueue_time(self) -> datetime: # read-only - def sequence_number(self) -> int: # read-only - def partition_id(self) -> str: # read-only - def locked_until(self) -> datetime: # read-only - def expired(self) -> bool: # read-only - def lock_token(self) -> str: # read-only def body(self) -> Union[bytes, Generator[bytes]]: # read-only def partition_key(self, value : str): def partition_key(self) -> str: @@ -209,7 +195,28 @@ def enqueue_sequence_number(self, value : int): def enqueue_sequence_number(self) -> int: # Methods - def schedule(self, schedule_time : datetime) -> None: + def schedule(self, schedule_time_utc : datetime) -> None: + + + +class BatchMessage(Message): +# inherited from Message + def add(self, message: Message) -> None: + + +class ReceivedMessage(Message): + + # @properties + def settled(self) -> bool: # read-only + def enqueued_time(self) -> datetime: # read-only + def scheduled_enqueue_time(self) -> datetime: # read-only + def sequence_number(self) -> int: # read-only + def partition_id(self) -> str: # read-only + def locked_until(self) -> datetime: # read-only + def expired(self) -> bool: # read-only + def lock_token(self) -> str: # read-only + + # methods def renew_lock(self) -> None: def complete(self) -> None: def dead_letter(self, description : str = None) -> None: @@ -217,19 +224,12 @@ def abandon(self) -> None: def defer(self) -> None: -class BatchMessage(Message): -# inherited from Message - def try_add(self, message: Message) -> None: - +class PeekMessage(ReceivedMessage): +# inherited from ReceivedMessage, and cannot settle/lock the message -class PeekMessage(Message): -# inherited from Message, and cannot settle/lock the message - -class DeferredMessage(Message): - # interited from Message - # its own @properties - def settled(self) -> bool: # read-only +class DeferredMessage(ReceivedMessage): +# interited from ReceivedMessage class ReceiveSettleMode: @@ -259,49 +259,37 @@ class ServiceBusSharedKeyCredential: def __init__(self, policy, key): def get_token(self, *scopes, **kwargs): -################################################### Split Line ################################################### +class ServiceBusError(Exception): -# Approach for Session APIs - Option 1 -## APIs -### Put session methods on ServiceBusReceiverClient -class ServiceBusReceiverClient: - def list_sessions(self): - def set_session_state(self, session_state): - def get_session_state(self): - def renew_lock(self): - @property - def expired(self): - - -## Sample Code: -# iterator receive -with subscription_receiver: - subscription_receiver.set_session_state(session_state="START") - for msg in subscription_receiver: - print(message) - msg.complete() - subscription_receiver.renew_lock() - if str(msg) == "shutdown": - subscription_receiver.set_session_state(session_state="STOP") - break - -# pull mode receive -with subscription_receiver: - session_id = "test_session" - msgs = subscription_receiver.receive(max_batch_size=10) - subscription_receiver.set_session_state(session_state="BEGIN") - for msg in msgs: - msg.complete() - subscription_receiver.renew_lock() - subscription_receiver.set_session_state(session_state="END") - -# receive deferred letter -with subscription_receiver: - if subscription_receiver.get_session_state() == "DONE": - defered_sequence_numbers = [xxxx] - deferred = subscription_receiver.receive_deferred_messages(defered_sequence_numbers) - subscription_receiver.settle_deferred_messages('completed', deferred) +class ServiceBusClient: # MGMT APIs + def __init__( + self, + fully_qualified_namespace : str, + credential : TokenCredential, + **kwargs + ): + @classmethod + def from_connection_string(cls, conn_str, **kwargs): + + def create_queue(self, queue_name, **kwargs): + def delete_queue(self, queue_name, fail_not_exist=False): + def list_queues(self) -> List[str]: + def get_queue_sender(self, queue_name, **kwargs) -> ServiceBusSenderClient: + def get_queue_receiver(self, queue_name, **kwargs) -> ServiceBusReceiverClient: + + def create_topic(self, topic_name, **kwargs): + def delete_topic(self, topic_name, fail_not_exist=False): + def list_topics(self): + def get_topic_sender(self, topic_name) -> ServiceBusSenderClient: + + def create_subscription(self, topic_name, subscription_name): + def delete_subscription(self, topic_name, subscription_name, fail_not_exist=False): + def list_subscriptions(self, topic_name): + def get_subscription_receiver(self, topic_name, subscription_name) -> ServiceBusReceiverClient: + + def list_session_ids_of_queue(self, queue_name): + def list_session_ids_of_subscription(self, topic_name, subscription_name): ################################################### Split Line ################################################### @@ -322,7 +310,8 @@ def get_rules(self): # Approach for Rule APIs - Option 2 ## APIs class ServiceBusReceiverClient: - def get_subscription_rule_manager(self): # raise Error when called on Queue + @property + def subscription_rule_manager(self): # raise Error when called on Queue class SubscriptionRuleManager: def add_rule(self, rule_name, filter): @@ -335,7 +324,8 @@ def get_rules(self): # Approach for Rule APIs - Option 3 ## APIs class SubscriptionRuleManager: - def from_receiver_client_for_subscription(receiver_client): + @classmethod + def from_receiver_client_for_subscription(cls, receiver_client): def add_rule(self, rule_name, filter): def remove_rule(self, rule_name): def get_rules(self): diff --git a/sdk/servicebus/azure-servicebus/api_samples.py b/sdk/servicebus/azure-servicebus/api_samples.py index 6f98d443253e..06435632e52a 100644 --- a/sdk/servicebus/azure-servicebus/api_samples.py +++ b/sdk/servicebus/azure-servicebus/api_samples.py @@ -107,7 +107,7 @@ batch_message = topic_sender.create_batch(max_size_in_bytes=256 * 1204) while True: try: - batch_message.try_add(Message("Test message")) + batch_message.add(Message("Test message")) except ValueError: break @@ -116,17 +116,20 @@ # 1.3 schedule schedule_message = Message("scheduled_message") - enqueue_time = datetime.utcnow() + timedelta(minutes=10) - sequence_number = topic_sender.schedule(schedule_message, enqueue_time) + enqueue_time_utc = datetime.utcnow() + timedelta(minutes=10) + sequence_number = topic_sender.schedule(schedule_message, enqueue_time_utc) + + running = True batch_schedule_message = topic_sender.create_batch(max_size_in_bytes=256 * 1204) - while True: + + while running: try: - batch_message.try_add(Message("Test message")) + batch_message.add(Message("Test message")) except ValueError: - break + batch_sequence_numbers = topic_sender.schedule(batch_schedule_message, enqueue_time_utc) + batch_schedule_message = topic_sender.create_batch(max_size_in_bytes=256 * 1204) - batch_sequence_numbers = topic_sender.schedule(batch_schedule_message, enqueue_time) # 1.4 cancel schedule topic_sender.cancel_scheduled_messages(sequence_number) @@ -152,6 +155,8 @@ # 2.1 peek with queue_receiver: msgs = queue_receiver.peek(count=10) + for msg in msgs: + print(msg) # 2.2 iterator receive with queue_receiver: @@ -171,9 +176,10 @@ # 2.4 receive deferred letter with queue_receiver: - defered_sequence_numbers = [xxxx] - deferred = queue_receiver.receive_deferred_messages(defered_sequence_numbers) - queue_receiver.settle_deferred_messages('completed', deferred) + defered_sequence_numbers = [1,2,3,4,5,6] + deferred_messages = queue_receiver.receive_deferred_messages(defered_sequence_numbers) + for i in deferred_messages: + deferred_messages.abandon() ### Session operation @@ -194,10 +200,10 @@ topic_sender.send(message, session_id="test_session") # 1.2 send batch - batch_message = topic_sender.create_batch(max_size_in_bytes=256 * 1204, seesion_id="test_session") + batch_message = topic_sender.create_batch(max_size_in_bytes=256 * 1204, session_id="test_session") while True: try: - batch_message.try_add(Message("Test message")) + batch_message.add(Message("Test message")) except ValueError: break @@ -206,17 +212,19 @@ # 1.3 schedule schedule_message = Message("scheduled_message") - enqueue_time = datetime.utcnow() + timedelta(minutes=10) - sequence_number = topic_sender.schedule(schedule_message, enqueue_time, session_id="test_session") + enqueue_time_utc = datetime.utcnow() + timedelta(minutes=10) + sequence_number = topic_sender.schedule(schedule_message, enqueue_time_utc) - batch_schedule_message = topic_sender.create_batch(max_size_in_bytes=256 * 1204) - while True: + running = True + + batch_schedule_message = topic_sender.create_batch(max_size_in_bytes=256 * 1204, session_id="test_session") + + while running: try: - batch_message.try_add(Message("Test message")) + batch_message.add(Message("Test message")) except ValueError: - break - - batch_sequence_numbers = topic_sender.schedule(batch_schedule_message, enqueue_time) + batch_sequence_numbers = topic_sender.schedule(batch_schedule_message, enqueue_time_utc) + batch_schedule_message = topic_sender.create_batch(max_size_in_bytes=256 * 1204) # 1.4 cancel schedule topic_sender.cancel_scheduled_messages(sequence_number) @@ -238,7 +246,7 @@ # 2.2 iterator receive with subscription_receiver: - session = subscription_receiver.get_session() + session = subscription_receiver.session session.set_session_state("START") for msg in subscription_receiver: print(message) @@ -250,21 +258,24 @@ # 2.3 pull mode receive with subscription_receiver: - session = subscription_receiver.get_session() + session = subscription_receiver.session msgs = subscription_receiver.receive(max_batch_size=10) + session.set_session_state("BEGIN") for msg in msgs: msg.complete() session.renew_lock() session.set_session_state("END") + # 2.4 receive deferred letter with subscription_receiver: - session = subscription_receiver.get_session() + session = subscription_receiver.session if session.get_session_state() == "DONE": defered_sequence_numbers = [1,2,3,4,5] - deferred = subscription_receiver.receive_deferred_messages(defered_sequence_numbers) - subscription_receiver.settle_deferred_messages('completed', deferred) + deferred_messages = subscription_receiver.receive_deferred_messages(defered_sequence_numbers) + for msg in deferred_messages: + deferred_messages.abandon() ### Rule operations From 175f7eb38083fd8e5e477899411f097978347f9f Mon Sep 17 00:00:00 2001 From: Yunhao Ling Date: Thu, 27 Feb 2020 18:23:59 -0800 Subject: [PATCH 11/24] fix typo and remove entity name from connection constructor --- sdk/servicebus/azure-servicebus/api_review.py | 18 +++++++++++------- 1 file changed, 11 insertions(+), 7 deletions(-) diff --git a/sdk/servicebus/azure-servicebus/api_review.py b/sdk/servicebus/azure-servicebus/api_review.py index adcccd6f4e36..c2ec5a4e8c19 100644 --- a/sdk/servicebus/azure-servicebus/api_review.py +++ b/sdk/servicebus/azure-servicebus/api_review.py @@ -1,7 +1,7 @@ # Two-Clients Approach ## APIs: -class SeviceBusSenderClient: +class ServiceBusSenderClient: def __init__( self, fully_qualified_namespace : str, @@ -22,7 +22,7 @@ def from_queue( queue_name : str, credential : TokenCredential, **kwargs - ) -> SeviceBusSenderClient: + ) -> ServiceBusSenderClient: @classmethod def from_topic( cls, @@ -30,14 +30,14 @@ def from_topic( topic_name : str, credential : TokenCredential, **kwargs - ) -> SeviceBusSenderClient: + ) -> ServiceBusSenderClient: @classmethod def from_connection_string( cls, conn_str : str, entity_name : str = None, **kwargs - ) -> SeviceBusSenderClient: + ) -> ServiceBusSenderClient: def __enter__(self): def __exit__(self): @@ -246,15 +246,19 @@ class ServiceBusConnection: def __int__( self, fully_qualified_namespace, - entity_name, token_credential ) -> None: def from_connection_string( self, - conn_str, - entity_name=None + conn_str ) -> ServiceBusConnection: + def __enter__(self): + def __exit__(self): + + def open(self): + def close(self): + class ServiceBusSharedKeyCredential: def __init__(self, policy, key): def get_token(self, *scopes, **kwargs): From 1fa61280227a6bf22eb51940fdcfd5a26257bd9f Mon Sep 17 00:00:00 2001 From: Yunhao Ling Date: Thu, 27 Feb 2020 22:40:29 -0800 Subject: [PATCH 12/24] fix typo in samples --- sdk/servicebus/azure-servicebus/api_samples.py | 18 +++++++++--------- 1 file changed, 9 insertions(+), 9 deletions(-) diff --git a/sdk/servicebus/azure-servicebus/api_samples.py b/sdk/servicebus/azure-servicebus/api_samples.py index 06435632e52a..cdd85d39ce09 100644 --- a/sdk/servicebus/azure-servicebus/api_samples.py +++ b/sdk/servicebus/azure-servicebus/api_samples.py @@ -4,18 +4,18 @@ ### creation of differen entity clients # Queue Sender -queue_sender = SeviceBusSenderClient.from_queue( +queue_sender = ServiceBusSenderClient.from_queue( fully_qualified_namepsace="fully_qualified_namepsace", queue_name="queue_name", credential=ServiceBusSharedKeyCredential("policy", "key") ) -queue_sender = SeviceBusSenderClient.from_connection_string( +queue_sender = ServiceBusSenderClient.from_connection_string( conn_str="conn_str", entity_name="queue_name" ) -queue_sender = SeviceBusSenderClient( +queue_sender = ServiceBusSenderClient( fully_qualified_namepsace="fully_qualified_namepsace", entity_name="queue_name", ServiceBusSharedKeyCredential("policy", "key") @@ -41,18 +41,18 @@ ) # Topic Sender -topic_sender = SeviceBusSenderClient.from_topic( +topic_sender = ServiceBusSenderClient.from_topic( fully_qualified_namepsace="fully_qualified_namepsace", topic_name="topic_name", credential=ServiceBusSharedKeyCredential("policy", "key") ) -topic_sender = SeviceBusSenderClient.from_connection_string( +topic_sender = ServiceBusSenderClient.from_connection_string( conn_str="conn_str", entity_name="topic_name" ) -topic_sender = SeviceBusSenderClient( +topic_sender = ServiceBusSenderClient( fully_qualified_namepsace="fully_qualified_namepsace", entity_name="topic_name", credential=ServiceBusSharedKeyCredential("policy", "key") @@ -84,14 +84,14 @@ ### Send and receive # 1 Send -queue_sender = SeviceBusSenderClient.from_queue( +queue_sender = ServiceBusSenderClient.from_queue( fully_qualified_namepsace="fully_qualified_namepsace", queue_name="queue_name", credential=ServiceBusSharedKeyCredential("policy", "key") ) ''' -topic_sender = SeviceBusSenderClient.from_queue( +topic_sender = ServiceBusSenderClient.from_queue( fully_qualified_namepsace="fully_qualified_namepsace", topic_name="topic_name", credential=ServiceBusSharedKeyCredential("policy", "key") @@ -186,7 +186,7 @@ # 1 Send -topic_sender = SeviceBusSenderClient.from_topic +topic_sender = ServiceBusSenderClient.from_topic ( fully_qualified_namespace="fully_qualified_namepsace", topic_name="topic_name", From 0558c32520a83714ae5cc96e303d778788b321a3 Mon Sep 17 00:00:00 2001 From: Yunhao Ling Date: Mon, 2 Mar 2020 11:41:50 -0800 Subject: [PATCH 13/24] remove factory methods and update sample accordingly --- sdk/servicebus/azure-servicebus/api_review.py | 52 +++------------- .../azure-servicebus/api_samples.py | 61 +++++-------------- 2 files changed, 25 insertions(+), 88 deletions(-) diff --git a/sdk/servicebus/azure-servicebus/api_review.py b/sdk/servicebus/azure-servicebus/api_review.py index c2ec5a4e8c19..bbb6fa3202c8 100644 --- a/sdk/servicebus/azure-servicebus/api_review.py +++ b/sdk/servicebus/azure-servicebus/api_review.py @@ -5,8 +5,9 @@ class ServiceBusSenderClient: def __init__( self, fully_qualified_namespace : str, - entity_name : str, - credential : TokenCredential, + credential: TokenCredential, + queue_name: str = None, + topic_name: str = None, logging_enable: bool = False, http_proxy: dict = None, transport_type: TransportType = None, @@ -15,27 +16,11 @@ def __init__( retry_backoff_maximum : int = 120 # cur_retry_backoff = retry_backoff_factor * (2^cur_retry_time) connection : ServiceBusConnection = None ) -> None: - @classmethod - def from_queue( - cls, - fully_qualified_namepsace : str, - queue_name : str, - credential : TokenCredential, - **kwargs - ) -> ServiceBusSenderClient: - @classmethod - def from_topic( - cls, - fully_qualified_namespace : str, - topic_name : str, - credential : TokenCredential, - **kwargs - ) -> ServiceBusSenderClient: - @classmethod def from_connection_string( cls, conn_str : str, - entity_name : str = None, + queue_name: str = None, + topic_name: str = None, **kwargs ) -> ServiceBusSenderClient: @@ -72,8 +57,9 @@ class ServiceBusReceiverClient: def __init__( self, fully_qualified_namespace : str, - entity_name : str, credential: TokenCredential, + queue_name: str = None, + topic_name: str = None, subscription_name : str = None, session_id : str = None, mode : ReceiveSettleMode = ReceiveSettleMode.PeekLock @@ -86,31 +72,11 @@ def __init__( connection: Connection = None ) -> None: @classmethod - def from_queue( - cls, - fully_qualified_namespace : str, - queue_name : str, - credential : TokenCredential, - session_id: str = None, - mode : ReceiveSettleMode = ReceiveSettleMode.PeekLock - **kwargs - ) -> ServiceBusReceiverClient: - @classmethod - def from_topic_subscription( - cls, - fully_qualified_namespace : str, - topic_name : str, - subscription_name : str, - credential: TokenCredential, - session_id: str = None, - mode : ReceiveSettleMode = ReceiveSettleMode.PeekLock - **kwargs - ) -> ServiceBusReceiverClient: - @classmethod def from_connection_string( cls, conn_str : str, - entity_name : str = None, + queue_name: str = None, + topic_name: str = None, subscription_name : str = None, session_id: str = None, mode : ReceiveSettleMode = ReceiveSettleMode.PeekLock diff --git a/sdk/servicebus/azure-servicebus/api_samples.py b/sdk/servicebus/azure-servicebus/api_samples.py index cdd85d39ce09..e6a2a5e2c5fd 100644 --- a/sdk/servicebus/azure-servicebus/api_samples.py +++ b/sdk/servicebus/azure-servicebus/api_samples.py @@ -4,77 +4,52 @@ ### creation of differen entity clients # Queue Sender -queue_sender = ServiceBusSenderClient.from_queue( - fully_qualified_namepsace="fully_qualified_namepsace", - queue_name="queue_name", - credential=ServiceBusSharedKeyCredential("policy", "key") -) - queue_sender = ServiceBusSenderClient.from_connection_string( conn_str="conn_str", - entity_name="queue_name" + queue_name="queue_name" ) queue_sender = ServiceBusSenderClient( fully_qualified_namepsace="fully_qualified_namepsace", - entity_name="queue_name", + queue_name="queue_name", ServiceBusSharedKeyCredential("policy", "key") ) # Queue Receiver -queue_receiver = SeviceBusReceiverClient.from_queue( - fully_qualified_namepsace="fully_qualified_namepsace", - queue_name="queue_name", - credential=ServiceBusSharedKeyCredential("policy", "key") -) - queue_receiver = SeviceBusReceiverClient.from_connection_string( conn_str="conn_str", - entity_name="queue_name" + queue_name="queue_name" ) queue_receiver = SeviceBusReceiverClient( fully_qualified_namepsace="fully_qualified_namepsace", - entity_name="queue_name", + queue_name="queue_name", credential=ServiceBusSharedKeyCredential("policy", "key"), session_id="test_session" ) # Topic Sender -topic_sender = ServiceBusSenderClient.from_topic( - fully_qualified_namepsace="fully_qualified_namepsace", - topic_name="topic_name", - credential=ServiceBusSharedKeyCredential("policy", "key") -) - topic_sender = ServiceBusSenderClient.from_connection_string( conn_str="conn_str", - entity_name="topic_name" + topic_name="topic_name" ) topic_sender = ServiceBusSenderClient( fully_qualified_namepsace="fully_qualified_namepsace", - entity_name="topic_name", - credential=ServiceBusSharedKeyCredential("policy", "key") -) - -# Subscription Receiver -subscription_receiver = SeviceBusReceiverClient.from_topic_subscription( - fully_qualified_namepsace="fully_qualified_namespace", topic_name="topic_name", - subscription_name="subscription_name", credential=ServiceBusSharedKeyCredential("policy", "key") ) +# Subscription Receiver subscription_receiver = SeviceBusReceiverClient.from_connection_string( conn_str="conn_str", - entity_name="topic_name", + topic_name="topic_name", subscription_name="subscription_name" ) subscription_receiver = SeviceBusReceiverClient( fully_qualified_namepsace="fully_qualified_namespace", - entity_name="topic_name", + topic_name="topic_name", subscription_name="subcription_name", credential=ServiceBusSharedKeyCredential("policy", "key"), session_id="test_session" @@ -84,14 +59,14 @@ ### Send and receive # 1 Send -queue_sender = ServiceBusSenderClient.from_queue( +queue_sender = ServiceBusSenderClient( fully_qualified_namepsace="fully_qualified_namepsace", queue_name="queue_name", credential=ServiceBusSharedKeyCredential("policy", "key") ) ''' -topic_sender = ServiceBusSenderClient.from_queue( +topic_sender = ServiceBusSenderClient( fully_qualified_namepsace="fully_qualified_namepsace", topic_name="topic_name", credential=ServiceBusSharedKeyCredential("policy", "key") @@ -137,14 +112,14 @@ # 2 Receive -queue_receiver = ServiceBusReceiverClient.from_queue( +queue_receiver = ServiceBusReceiverClient( fully_qualified_namespace="fully_qualified_namespace", queue_name="queue_name", credential=ServiceBusSharedKeyCredential("policy", "key") ) ''' -subscription_receiver = ServiceBusReceiverClient.from_queue( +subscription_receiver = ServiceBusReceiverClient( fully_qualified_namespace="fully_qualified_namespace", topic_name="topic_name", subscription_name="subscription_name" @@ -186,8 +161,7 @@ # 1 Send -topic_sender = ServiceBusSenderClient.from_topic -( +topic_sender = ServiceBusSenderClient( fully_qualified_namespace="fully_qualified_namepsace", topic_name="topic_name", credential=ServiceBusSharedKeyCredential("policy", "key") @@ -231,8 +205,7 @@ topic_sender.cancel_scheduled_messages(batch_sequence_numbers) # 2 Receive -subscription_receiver = ServiceBusReceiverClient.from_topic_subscription -( +subscription_receiver = ServiceBusReceiverClient( fully_qualified_namespace="fully_qualified_namespace", topic_name="topic_name", subscription_name="subscription_name", @@ -279,8 +252,7 @@ ### Rule operations -subscription_receiver = ServiceBusReceiverClient.from_topic_subscription -( +subscription_receiver = ServiceBusReceiverClient( fully_qualified_namespace="fully_qualified_namespace", topic_name="topic_name", subscription_name="subscription_name", @@ -293,8 +265,7 @@ ### Transaction operations -subscription_receiver = ServiceBusReceiverClient.from_topic_subscription -( +subscription_receiver = ServiceBusReceiverClient( fully_qualified_namespace="fully_qualified_namespace", topic_name="topic_name", subscription_name="subscription_name", From c21dd5e7637948a54d7ea36354268a9efd807c40 Mon Sep 17 00:00:00 2001 From: Yunhao Ling Date: Wed, 4 Mar 2020 12:04:11 -0800 Subject: [PATCH 14/24] introduce top level sb client --- sdk/servicebus/azure-servicebus/api_review.py | 81 +++++++++--------- .../azure-servicebus/api_samples.py | 85 +++++++++++++------ 2 files changed, 98 insertions(+), 68 deletions(-) diff --git a/sdk/servicebus/azure-servicebus/api_review.py b/sdk/servicebus/azure-servicebus/api_review.py index bbb6fa3202c8..143881646464 100644 --- a/sdk/servicebus/azure-servicebus/api_review.py +++ b/sdk/servicebus/azure-servicebus/api_review.py @@ -1,7 +1,38 @@ -# Two-Clients Approach - ## APIs: -class ServiceBusSenderClient: +class ServiceBusClient: # MGMT APIs + def __init__( + self, + fully_qualified_namespace : str, + credential : TokenCredential, + **kwargs + ): + + ########## Preivew 1 scope ########## + @classmethod + def from_connection_string(cls, conn_str, **kwargs): + def get_queue_sender(self, queue_name, **kwargs) -> ServiceBusSender: + def get_queue_receiver(self, queue_name, **kwargs) -> ServiceBusReceiver: + ########## Preivew 1 scope ########## + + def get_topic_sender(self, topic_name) -> ServiceBusSender: + def get_subscription_receiver(self, topic_name, subscription_name) -> ServiceBusReceiver: + + def create_queue(self, queue_name, **kwargs): + def delete_queue(self, queue_name, fail_not_exist=False): + def list_queues(self) -> List[str]: + + def create_topic(self, topic_name, **kwargs): + def delete_topic(self, topic_name, fail_not_exist=False): + def list_topics(self): + + def create_subscription(self, topic_name, subscription_name): + def delete_subscription(self, topic_name, subscription_name, fail_not_exist=False): + def list_subscriptions(self, topic_name): + + def list_session_ids_of_queue(self, queue_name): + def list_session_ids_of_subscription(self, topic_name, subscription_name): + +class ServiceBusSender: def __init__( self, fully_qualified_namespace : str, @@ -22,7 +53,7 @@ def from_connection_string( queue_name: str = None, topic_name: str = None, **kwargs - ) -> ServiceBusSenderClient: + ) -> ServiceBusSender: def __enter__(self): def __exit__(self): @@ -53,7 +84,7 @@ def schedule( def cancel_scheduled_messages(self, sequence_number : Union[int, List[int]]) -> None: -class ServiceBusReceiverClient: +class ServiceBusReceiver: def __init__( self, fully_qualified_namespace : str, @@ -81,7 +112,7 @@ def from_connection_string( session_id: str = None, mode : ReceiveSettleMode = ReceiveSettleMode.PeekLock **kwargs - )-> ServiceBusReceiverClient: + )-> ServiceBusReceiver: def __enter__(self): def __exit__(self): @@ -231,45 +262,15 @@ def get_token(self, *scopes, **kwargs): class ServiceBusError(Exception): -class ServiceBusClient: # MGMT APIs - def __init__( - self, - fully_qualified_namespace : str, - credential : TokenCredential, - **kwargs - ): - - @classmethod - def from_connection_string(cls, conn_str, **kwargs): - - def create_queue(self, queue_name, **kwargs): - def delete_queue(self, queue_name, fail_not_exist=False): - def list_queues(self) -> List[str]: - def get_queue_sender(self, queue_name, **kwargs) -> ServiceBusSenderClient: - def get_queue_receiver(self, queue_name, **kwargs) -> ServiceBusReceiverClient: - - def create_topic(self, topic_name, **kwargs): - def delete_topic(self, topic_name, fail_not_exist=False): - def list_topics(self): - def get_topic_sender(self, topic_name) -> ServiceBusSenderClient: - - def create_subscription(self, topic_name, subscription_name): - def delete_subscription(self, topic_name, subscription_name, fail_not_exist=False): - def list_subscriptions(self, topic_name): - def get_subscription_receiver(self, topic_name, subscription_name) -> ServiceBusReceiverClient: - - def list_session_ids_of_queue(self, queue_name): - def list_session_ids_of_subscription(self, topic_name, subscription_name): - ################################################### Split Line ################################################### # Approach for Rule APIs - Option 1 ## APIs ### 3-Clients approach -class ServiceBusSenderClient: -class QueueReceiverClient: -class SubscriptionReceiverClient: +class ServiceBusSender: +class QueueReceiver: +class SubscriptionReceiver: def add_rule(self, rule_name, filter): def remove_rule(self, rule_name): def get_rules(self): @@ -279,7 +280,7 @@ def get_rules(self): # Approach for Rule APIs - Option 2 ## APIs -class ServiceBusReceiverClient: +class ServiceBusReceiver: @property def subscription_rule_manager(self): # raise Error when called on Queue diff --git a/sdk/servicebus/azure-servicebus/api_samples.py b/sdk/servicebus/azure-servicebus/api_samples.py index e6a2a5e2c5fd..3a9227dd7c27 100644 --- a/sdk/servicebus/azure-servicebus/api_samples.py +++ b/sdk/servicebus/azure-servicebus/api_samples.py @@ -3,13 +3,29 @@ ### creation of differen entity clients +# Top Level ServiceBusClient + +## from connection string + +sb_client = ServiceBusClient.from_connection_string(conn_str="conn_str") +sb_client = ServiceBusClient( + fully_qualified_namepsace="fully_qualified_namepsace", + crendetial=ServiceBusSharedKeyCredential("policy", "key") +) + +queue_sender = sb_client.get_queue_sender(queue_name="queue_name") +queue_receiver = sb_client.get_queue_receiver(queue_name="queue_name") +topic_sender = sb_client.get_topic_sender(topic_name="topic_name") +subscription_receiver = sb_client.get_subscription_sender(topic_name="topic_name", subscription_name="subscription_name") + # Queue Sender -queue_sender = ServiceBusSenderClient.from_connection_string( + +queue_sender = ServiceBusSender.from_connection_string( conn_str="conn_str", queue_name="queue_name" ) -queue_sender = ServiceBusSenderClient( +queue_sender = ServiceBusSender( fully_qualified_namepsace="fully_qualified_namepsace", queue_name="queue_name", ServiceBusSharedKeyCredential("policy", "key") @@ -29,12 +45,12 @@ ) # Topic Sender -topic_sender = ServiceBusSenderClient.from_connection_string( +topic_sender = ServiceBusSender.from_connection_string( conn_str="conn_str", topic_name="topic_name" ) -topic_sender = ServiceBusSenderClient( +topic_sender = ServiceBusSender( fully_qualified_namepsace="fully_qualified_namepsace", topic_name="topic_name", credential=ServiceBusSharedKeyCredential("policy", "key") @@ -58,15 +74,17 @@ ### Send and receive # 1 Send +sb_client = ServiceBusClient.from_connection_string(conn_str="conn_str") +queue_sender = sb_client.get_queue_sender(queue_name="queue_name") -queue_sender = ServiceBusSenderClient( - fully_qualified_namepsace="fully_qualified_namepsace", - queue_name="queue_name", - credential=ServiceBusSharedKeyCredential("policy", "key") -) +# queue_sender = ServiceBusSender( +# fully_qualified_namepsace="fully_qualified_namepsace", +# queue_name="queue_name", +# credential=ServiceBusSharedKeyCredential("policy", "key") +# ) ''' -topic_sender = ServiceBusSenderClient( +topic_sender = ServiceBusSender( fully_qualified_namepsace="fully_qualified_namepsace", topic_name="topic_name", credential=ServiceBusSharedKeyCredential("policy", "key") @@ -111,15 +129,17 @@ topic_sender.cancel_scheduled_messages(batch_sequence_numbers) # 2 Receive +sb_client = ServiceBusClient.from_connection_string(conn_str="conn_str") +queue_receiver = sb_client.get_queue_receiver(queue_name="queue_name") -queue_receiver = ServiceBusReceiverClient( - fully_qualified_namespace="fully_qualified_namespace", - queue_name="queue_name", - credential=ServiceBusSharedKeyCredential("policy", "key") -) +# queue_receiver = ServiceBusReceiver( +# fully_qualified_namespace="fully_qualified_namespace", +# queue_name="queue_name", +# credential=ServiceBusSharedKeyCredential("policy", "key") +# ) ''' -subscription_receiver = ServiceBusReceiverClient( +subscription_receiver = ServiceBusReceiver( fully_qualified_namespace="fully_qualified_namespace", topic_name="topic_name", subscription_name="subscription_name" @@ -161,11 +181,14 @@ # 1 Send -topic_sender = ServiceBusSenderClient( - fully_qualified_namespace="fully_qualified_namepsace", - topic_name="topic_name", - credential=ServiceBusSharedKeyCredential("policy", "key") -) +sb_client = ServiceBusClient.from_connection_string(conn_str="conn_str") +topic_sender = sb_client.get_topic_sender(topic_name="topic_name") + +# topic_sender = ServiceBusSender( +# fully_qualified_namespace="fully_qualified_namepsace", +# topic_name="topic_name", +# credential=ServiceBusSharedKeyCredential("policy", "key") +# ) with topic_sender: # 1.1 send single @@ -205,14 +228,20 @@ topic_sender.cancel_scheduled_messages(batch_sequence_numbers) # 2 Receive -subscription_receiver = ServiceBusReceiverClient( - fully_qualified_namespace="fully_qualified_namespace", +sb_client = ServiceBusClient.from_connection_string(conn_str="conn_str") +subscription_receiver = sb_client.get_subscription_receiver( topic_name="topic_name", - subscription_name="subscription_name", - credential=ServiceBusSharedKeyCredential("policy", "key"), - session_id="test_session" + subscription_name="subscription_name" ) +# subscription_receiver = ServiceBusReceiver( +# fully_qualified_namespace="fully_qualified_namespace", +# topic_name="topic_name", +# subscription_name="subscription_name", +# credential=ServiceBusSharedKeyCredential("policy", "key"), +# session_id="test_session" +# ) + # 2.1 peek with subscription_receiver: msgs = subscription_receiver.peek(count=10) @@ -252,7 +281,7 @@ ### Rule operations -subscription_receiver = ServiceBusReceiverClient( +subscription_receiver = ServiceBusReceiver( fully_qualified_namespace="fully_qualified_namespace", topic_name="topic_name", subscription_name="subscription_name", @@ -265,7 +294,7 @@ ### Transaction operations -subscription_receiver = ServiceBusReceiverClient( +subscription_receiver = ServiceBusReceiver( fully_qualified_namespace="fully_qualified_namespace", topic_name="topic_name", subscription_name="subscription_name", From c1f212f1b76cb14951ddbe098757833acb621a7a Mon Sep 17 00:00:00 2001 From: Yunhao Ling Date: Wed, 4 Mar 2020 14:50:27 -0800 Subject: [PATCH 15/24] update api and samples as per discussion --- sdk/servicebus/azure-servicebus/api_review.py | 94 +---- .../azure-servicebus/api_samples.py | 330 ++++++------------ 2 files changed, 128 insertions(+), 296 deletions(-) diff --git a/sdk/servicebus/azure-servicebus/api_review.py b/sdk/servicebus/azure-servicebus/api_review.py index 143881646464..23a417e6e741 100644 --- a/sdk/servicebus/azure-servicebus/api_review.py +++ b/sdk/servicebus/azure-servicebus/api_review.py @@ -4,33 +4,22 @@ def __init__( self, fully_qualified_namespace : str, credential : TokenCredential, - **kwargs + logging_enable: bool = False, + http_proxy: dict = None, + transport_type: TransportType = None, ): - ########## Preivew 1 scope ########## + def __enter__(self): + def __exit__(self, *kwargs): + def close(self): + @classmethod def from_connection_string(cls, conn_str, **kwargs): + def get_queue_sender(self, queue_name, **kwargs) -> ServiceBusSender: def get_queue_receiver(self, queue_name, **kwargs) -> ServiceBusReceiver: - ########## Preivew 1 scope ########## - - def get_topic_sender(self, topic_name) -> ServiceBusSender: - def get_subscription_receiver(self, topic_name, subscription_name) -> ServiceBusReceiver: - - def create_queue(self, queue_name, **kwargs): - def delete_queue(self, queue_name, fail_not_exist=False): - def list_queues(self) -> List[str]: - - def create_topic(self, topic_name, **kwargs): - def delete_topic(self, topic_name, fail_not_exist=False): - def list_topics(self): - - def create_subscription(self, topic_name, subscription_name): - def delete_subscription(self, topic_name, subscription_name, fail_not_exist=False): - def list_subscriptions(self, topic_name): - - def list_session_ids_of_queue(self, queue_name): - def list_session_ids_of_subscription(self, topic_name, subscription_name): + def get_topic_sender(self, topic_name, **kwargs) -> ServiceBusSender: + def get_subscription_receiver(self, topic_name, subscription_name, **kwargs) -> ServiceBusReceiver: class ServiceBusSender: def __init__( @@ -45,8 +34,8 @@ def __init__( retry_total : int = 3, retry_backoff_factor : float = 30, retry_backoff_maximum : int = 120 # cur_retry_backoff = retry_backoff_factor * (2^cur_retry_time) - connection : ServiceBusConnection = None ) -> None: + @classmethod def from_connection_string( cls, conn_str : str, @@ -64,7 +53,6 @@ def get_properties(self) -> Dict[str, Any]: def create_batch( self, - session_id : str = None, max_size_in_bytes : int = None ) -> BatchMessage: @@ -100,7 +88,6 @@ def __init__( retry_total : int = 3, retry_backoff_factor: float = 30, retry_backoff_maximum: int = 120 # cur_retry_backoff = retry_backoff_factor * (2^cur_retry_time) - connection: Connection = None ) -> None: @classmethod def from_connection_string( @@ -180,8 +167,6 @@ def partition_key(self, value : str): def partition_key(self) -> str: def via_partition_key(self, value: str): def via_partition_key(self) -> str: - def session_id(self, value : str): - def session_id(self) -> str: def time_to_live(self, value : Union[float, timedelta]): def time_to_live(self) -> datetime.timedelta: def annotations(self, value : dict[str, Any]): @@ -212,6 +197,7 @@ def partition_id(self) -> str: # read-only def locked_until(self) -> datetime: # read-only def expired(self) -> bool: # read-only def lock_token(self) -> str: # read-only + def session_id(self) -> str: # read-only # methods def renew_lock(self) -> None: @@ -239,64 +225,8 @@ class TransportType(Enum): AmqpOverWebsocket -class ServiceBusConnection: - def __int__( - self, - fully_qualified_namespace, - token_credential - ) -> None: - def from_connection_string( - self, - conn_str - ) -> ServiceBusConnection: - - def __enter__(self): - def __exit__(self): - - def open(self): - def close(self): - class ServiceBusSharedKeyCredential: def __init__(self, policy, key): def get_token(self, *scopes, **kwargs): class ServiceBusError(Exception): - -################################################### Split Line ################################################### - -# Approach for Rule APIs - Option 1 -## APIs -### 3-Clients approach - -class ServiceBusSender: -class QueueReceiver: -class SubscriptionReceiver: - def add_rule(self, rule_name, filter): - def remove_rule(self, rule_name): - def get_rules(self): - - -################################################### Split Line ################################################### - -# Approach for Rule APIs - Option 2 -## APIs -class ServiceBusReceiver: - @property - def subscription_rule_manager(self): # raise Error when called on Queue - -class SubscriptionRuleManager: - def add_rule(self, rule_name, filter): - def remove_rule(self, rule_name): - def get_rules(self): - - -################################################### Split Line ################################################### - -# Approach for Rule APIs - Option 3 -## APIs -class SubscriptionRuleManager: - @classmethod - def from_receiver_client_for_subscription(cls, receiver_client): - def add_rule(self, rule_name, filter): - def remove_rule(self, rule_name): - def get_rules(self): diff --git a/sdk/servicebus/azure-servicebus/api_samples.py b/sdk/servicebus/azure-servicebus/api_samples.py index 3a9227dd7c27..734b494c6175 100644 --- a/sdk/servicebus/azure-servicebus/api_samples.py +++ b/sdk/servicebus/azure-servicebus/api_samples.py @@ -1,44 +1,42 @@ +## Sample Code: -## Sampe Code: - -### creation of differen entity clients +### creation of different entity clients # Top Level ServiceBusClient - -## from connection string - sb_client = ServiceBusClient.from_connection_string(conn_str="conn_str") sb_client = ServiceBusClient( - fully_qualified_namepsace="fully_qualified_namepsace", + fully_qualified_namespace="fully_qualified_namespace", crendetial=ServiceBusSharedKeyCredential("policy", "key") ) queue_sender = sb_client.get_queue_sender(queue_name="queue_name") queue_receiver = sb_client.get_queue_receiver(queue_name="queue_name") topic_sender = sb_client.get_topic_sender(topic_name="topic_name") -subscription_receiver = sb_client.get_subscription_sender(topic_name="topic_name", subscription_name="subscription_name") +subscription_receiver = sb_client.get_subscription_receiver( + topic_name="topic_name", + subscription_name="subscription_name" +) # Queue Sender - queue_sender = ServiceBusSender.from_connection_string( conn_str="conn_str", queue_name="queue_name" ) queue_sender = ServiceBusSender( - fully_qualified_namepsace="fully_qualified_namepsace", + fully_qualified_namespace="fully_qualified_namespace", queue_name="queue_name", ServiceBusSharedKeyCredential("policy", "key") ) # Queue Receiver -queue_receiver = SeviceBusReceiverClient.from_connection_string( +queue_receiver = ServiceBusReceiver.from_connection_string( conn_str="conn_str", queue_name="queue_name" ) -queue_receiver = SeviceBusReceiverClient( - fully_qualified_namepsace="fully_qualified_namepsace", +queue_receiver = ServiceBusReceiver( + fully_qualified_namespace="fully_qualified_namespace", queue_name="queue_name", credential=ServiceBusSharedKeyCredential("policy", "key"), session_id="test_session" @@ -51,20 +49,20 @@ ) topic_sender = ServiceBusSender( - fully_qualified_namepsace="fully_qualified_namepsace", + fully_qualified_namespace="fully_qualified_namespace", topic_name="topic_name", credential=ServiceBusSharedKeyCredential("policy", "key") ) # Subscription Receiver -subscription_receiver = SeviceBusReceiverClient.from_connection_string( +subscription_receiver = ServiceBusReceiver.from_connection_string( conn_str="conn_str", topic_name="topic_name", subscription_name="subscription_name" ) -subscription_receiver = SeviceBusReceiverClient( - fully_qualified_namepsace="fully_qualified_namespace", +subscription_receiver = ServiceBusReceiver( + fully_qualified_namespace="fully_qualified_namespace", topic_name="topic_name", subscription_name="subcription_name", credential=ServiceBusSharedKeyCredential("policy", "key"), @@ -77,241 +75,145 @@ sb_client = ServiceBusClient.from_connection_string(conn_str="conn_str") queue_sender = sb_client.get_queue_sender(queue_name="queue_name") -# queue_sender = ServiceBusSender( -# fully_qualified_namepsace="fully_qualified_namepsace", -# queue_name="queue_name", -# credential=ServiceBusSharedKeyCredential("policy", "key") -# ) - -''' -topic_sender = ServiceBusSender( - fully_qualified_namepsace="fully_qualified_namepsace", - topic_name="topic_name", - credential=ServiceBusSharedKeyCredential("policy", "key") -) -''' - -with queue_sender: - # 1.1 send single - message = Message("Test message") - queue_sender.send(message) +with sb_client: + with queue_sender: + # 1.1 send single + message = Message("Test message") + queue_sender.send(message) - # 1.2 send batch - batch_message = topic_sender.create_batch(max_size_in_bytes=256 * 1204) - while True: - try: - batch_message.add(Message("Test message")) - except ValueError: - break + # 1.2 send batch + batch_message = topic_sender.create_batch(max_size_in_bytes=256 * 1204) + while True: + try: + batch_message.add(Message("Test message")) + except ValueError: + break - topic_sender.send(batch_message) + topic_sender.send(batch_message) - # 1.3 schedule - schedule_message = Message("scheduled_message") + # 1.3 schedule + schedule_message = Message("scheduled_message") - enqueue_time_utc = datetime.utcnow() + timedelta(minutes=10) - sequence_number = topic_sender.schedule(schedule_message, enqueue_time_utc) + schedule_time_utc = datetime.utcnow() + timedelta(minutes=10) + sequence_number = topic_sender.schedule(schedule_message, schedule_time_utc) - running = True + batch_schedule_message = topic_sender.create_batch(max_size_in_bytes=256 * 1204) - batch_schedule_message = topic_sender.create_batch(max_size_in_bytes=256 * 1204) + while True: + try: + batch_message.add(Message("Test message")) + except ValueError: + break - while running: - try: - batch_message.add(Message("Test message")) - except ValueError: - batch_sequence_numbers = topic_sender.schedule(batch_schedule_message, enqueue_time_utc) - batch_schedule_message = topic_sender.create_batch(max_size_in_bytes=256 * 1204) + batch_sequence_numbers = topic_sender.schedule(batch_schedule_message, schedule_time_utc) - - # 1.4 cancel schedule - topic_sender.cancel_scheduled_messages(sequence_number) - topic_sender.cancel_scheduled_messages(batch_sequence_numbers) + # 1.4 cancel schedule + topic_sender.cancel_scheduled_messages(sequence_number) + topic_sender.cancel_scheduled_messages(batch_sequence_numbers) # 2 Receive sb_client = ServiceBusClient.from_connection_string(conn_str="conn_str") queue_receiver = sb_client.get_queue_receiver(queue_name="queue_name") -# queue_receiver = ServiceBusReceiver( -# fully_qualified_namespace="fully_qualified_namespace", -# queue_name="queue_name", -# credential=ServiceBusSharedKeyCredential("policy", "key") -# ) - -''' -subscription_receiver = ServiceBusReceiver( - fully_qualified_namespace="fully_qualified_namespace", - topic_name="topic_name", - subscription_name="subscription_name" - credential=ServiceBusSharedKeyCredential("policy", "key") -) -''' - # 2.1 peek -with queue_receiver: - msgs = queue_receiver.peek(count=10) - for msg in msgs: - print(msg) +with sb_client: + with queue_receiver: + msgs = queue_receiver.peek(count=10) + for msg in msgs: + print(msg) # 2.2 iterator receive -with queue_receiver: - for msg in queue_receiver: - print(message) - msg.complete() - # msg.renew_lock() - # msg.abondon() - # msg.defer() - # msg.dead_letter() +with sb_client: + with queue_receiver: + for msg in queue_receiver: + print(message) + msg.complete() + # msg.renew_lock() + # msg.abondon() + # msg.defer() + # msg.dead_letter() # 2.3 pull mode receive -with queue_receiver: - msgs = queue_receiver.receive(max_batch_size=10) - for msg in msgs: - msg.complete() +with sb_client: + with queue_receiver: + msgs = queue_receiver.receive(max_batch_size=10) + for msg in msgs: + msg.complete() # 2.4 receive deferred letter -with queue_receiver: - defered_sequence_numbers = [1,2,3,4,5,6] - deferred_messages = queue_receiver.receive_deferred_messages(defered_sequence_numbers) - for i in deferred_messages: - deferred_messages.abandon() +with sb_client: + with queue_receiver: + defered_sequence_numbers = [1,2,3,4,5,6] + deferred_messages = queue_receiver.receive_deferred_messages(defered_sequence_numbers) + for i in deferred_messages: + deferred_messages.abandon() ### Session operation - - # 1 Send sb_client = ServiceBusClient.from_connection_string(conn_str="conn_str") topic_sender = sb_client.get_topic_sender(topic_name="topic_name") -# topic_sender = ServiceBusSender( -# fully_qualified_namespace="fully_qualified_namepsace", -# topic_name="topic_name", -# credential=ServiceBusSharedKeyCredential("policy", "key") -# ) +with sb_client: + with topic_sender: + # 1.1 send single + message = Message("Test message") + topic_sender.send(message, session_id="test_session") -with topic_sender: - # 1.1 send single - message = Message("Test message") - # message.session_id = "test_session" - topic_sender.send(message, session_id="test_session") + # 1.2 send batch + batch_message = topic_sender.create_batch() + while True: + try: + batch_message.add(Message("Test message")) + except ValueError: + break - # 1.2 send batch - batch_message = topic_sender.create_batch(max_size_in_bytes=256 * 1204, session_id="test_session") - while True: - try: - batch_message.add(Message("Test message")) - except ValueError: - break + topic_sender.send(batch_message, session_id="test_session") - topic_sender.send(batch_message) + # 1.3 schedule + schedule_message = Message("Scheduled message") + schedule_time_utc = datetime.utcnow() + timedelta(minutes=10) + sequence_number = topic_sender.schedule(schedule_message, schedule_time_utc, session_id="test_session") - # 1.3 schedule - schedule_message = Message("scheduled_message") + batch_schedule_message = topic_sender.create_batch() - enqueue_time_utc = datetime.utcnow() + timedelta(minutes=10) - sequence_number = topic_sender.schedule(schedule_message, enqueue_time_utc) + batch_sequence_numbers = topic_sender.schedule( + message=batch_schedule_message, + schedule_time_utc=schedule_time_utc, + session_id="test_session" + ) - running = True - - batch_schedule_message = topic_sender.create_batch(max_size_in_bytes=256 * 1204, session_id="test_session") - - while running: - try: - batch_message.add(Message("Test message")) - except ValueError: - batch_sequence_numbers = topic_sender.schedule(batch_schedule_message, enqueue_time_utc) - batch_schedule_message = topic_sender.create_batch(max_size_in_bytes=256 * 1204) - - # 1.4 cancel schedule - topic_sender.cancel_scheduled_messages(sequence_number) - topic_sender.cancel_scheduled_messages(batch_sequence_numbers) + # 1.4 cancel schedule + topic_sender.cancel_scheduled_messages(sequence_number) + topic_sender.cancel_scheduled_messages(batch_sequence_numbers) # 2 Receive sb_client = ServiceBusClient.from_connection_string(conn_str="conn_str") subscription_receiver = sb_client.get_subscription_receiver( - topic_name="topic_name", - subscription_name="subscription_name" -) - -# subscription_receiver = ServiceBusReceiver( -# fully_qualified_namespace="fully_qualified_namespace", -# topic_name="topic_name", -# subscription_name="subscription_name", -# credential=ServiceBusSharedKeyCredential("policy", "key"), -# session_id="test_session" -# ) - -# 2.1 peek -with subscription_receiver: - msgs = subscription_receiver.peek(count=10) - -# 2.2 iterator receive -with subscription_receiver: - session = subscription_receiver.session - session.set_session_state("START") - for msg in subscription_receiver: - print(message) - msg.complete() - session.renew_lock() - if str(msg) == "shutdown": - session.set_session_state("STOP") - break - -# 2.3 pull mode receive -with subscription_receiver: - session = subscription_receiver.session - msgs = subscription_receiver.receive(max_batch_size=10) - - session.set_session_state("BEGIN") - for msg in msgs: - msg.complete() - session.renew_lock() - session.set_session_state("END") - - -# 2.4 receive deferred letter -with subscription_receiver: - session = subscription_receiver.session - if session.get_session_state() == "DONE": - defered_sequence_numbers = [1,2,3,4,5] - deferred_messages = subscription_receiver.receive_deferred_messages(defered_sequence_numbers) - for msg in deferred_messages: - deferred_messages.abandon() - -### Rule operations - -subscription_receiver = ServiceBusReceiver( - fully_qualified_namespace="fully_qualified_namespace", topic_name="topic_name", subscription_name="subscription_name", - credential=ServiceBusSharedKeyCredential("policy", "key")) - -with subscription_receiver: - subscription_receiver.get_rules() - subscription_receiver.add_rule("test_rule", "1=1") - subscription_receiver.remove_rule("test_rule") - -### Transaction operations - -subscription_receiver = ServiceBusReceiver( - fully_qualified_namespace="fully_qualified_namespace", - topic_name="topic_name", - subscription_name="subscription_name", - credential=ServiceBusSharedKeyCredential("policy", "key") + session_id="test_session" ) -with subscription_receiver: - transaction = subscription_receiver.get_transaction() # TODO, no idea on how to get a transaction for now - with transaction: - try: - subscription_receiver.add_rule("test_rule", "1=1") - msgs = subscription_receiver.receive(max_batch_size=1000) - for msg in msgs: - msg.abandon() - subscription_receiver.remove_rule("test_rule") - session_msgs = subscription_receiver.receive(max_batch_size=1000) - for msg in session_msgs: - msg.complete() - except: - transaction.abort() \ No newline at end of file +# 2.1 iterator receive +with sb_client: + with subscription_receiver: + session = subscription_receiver.session + session.set_session_state("START") + for msg in subscription_receiver: + if session.expired: + break + print(message) + msg.complete() + session.renew_lock() + +# 2.2 pull mode receive +with sb_client: + with subscription_receiver: + session = subscription_receiver.session + msgs = subscription_receiver.receive(max_batch_size=10) + session.set_session_state("BEGIN") + for msg in msgs: + msg.complete() + session.renew_lock() + session.set_session_state("END") From 16eb476c69f10e73c51795bd9ebef9813aac363b Mon Sep 17 00:00:00 2001 From: Yunhao Ling Date: Thu, 5 Mar 2020 18:22:26 -0800 Subject: [PATCH 16/24] update message inheritance --- sdk/servicebus/azure-servicebus/api_review.py | 35 +++++++++++++------ 1 file changed, 25 insertions(+), 10 deletions(-) diff --git a/sdk/servicebus/azure-servicebus/api_review.py b/sdk/servicebus/azure-servicebus/api_review.py index 23a417e6e741..f969a326ac72 100644 --- a/sdk/servicebus/azure-servicebus/api_review.py +++ b/sdk/servicebus/azure-servicebus/api_review.py @@ -1,5 +1,5 @@ ## APIs: -class ServiceBusClient: # MGMT APIs +class ServiceBusClient: def __init__( self, fully_qualified_namespace : str, @@ -127,13 +127,11 @@ def receive( timeout : float = None, ) -> List[ReceivedMessage]: # Pull mode receive - @property def session(self) -> ServiceBusSession: - # Rule APIs, raise Error when called on Queue - def add_rule(self, rule_name : str, filter : str) -> None: # TODO: figure out what filter is + def add_rule(self, rule_name : str, filter : str) -> None: def remove_rule(self, rule_name : str) -> None: def get_rules(self) -> List[str] -> Dict[str, Any]: @@ -180,12 +178,18 @@ def enqueue_sequence_number(self) -> int: def schedule(self, schedule_time_utc : datetime) -> None: - class BatchMessage(Message): # inherited from Message def add(self, message: Message) -> None: +class PeekMessage(Message): + def enqueued_time(self) -> datetime: # read-only + def scheduled_enqueue_time(self) -> datetime: # read-only + def sequence_number(self) -> int: # read-only + def partition_id(self) -> str: # read-only + def session_id(self) -> str: # read-only + class ReceivedMessage(Message): # @properties @@ -207,12 +211,23 @@ def abandon(self) -> None: def defer(self) -> None: -class PeekMessage(ReceivedMessage): -# inherited from ReceivedMessage, and cannot settle/lock the message - - class DeferredMessage(ReceivedMessage): -# interited from ReceivedMessage + # @properties + def settled(self) -> bool: # read-only + def enqueued_time(self) -> datetime: # read-only + def scheduled_enqueue_time(self) -> datetime: # read-only + def sequence_number(self) -> int: # read-only + def partition_id(self) -> str: # read-only + def locked_until(self) -> datetime: # read-only + def expired(self) -> bool: # read-only + def lock_token(self) -> str: # read-only + def session_id(self) -> str: # read-only + + # methods + def renew_lock(self) -> None: + def complete(self) -> None: + def dead_letter(self, description: str = None) -> None: + def abandon(self) -> None: class ReceiveSettleMode: From 187700f7a13b17be2dfc99b9884abcc7ae81f13a Mon Sep 17 00:00:00 2001 From: Yunhao Ling Date: Thu, 5 Mar 2020 23:31:13 -0800 Subject: [PATCH 17/24] fix transport type being None --- sdk/servicebus/azure-servicebus/api_review.py | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/sdk/servicebus/azure-servicebus/api_review.py b/sdk/servicebus/azure-servicebus/api_review.py index f969a326ac72..f81e7ce93713 100644 --- a/sdk/servicebus/azure-servicebus/api_review.py +++ b/sdk/servicebus/azure-servicebus/api_review.py @@ -6,7 +6,7 @@ def __init__( credential : TokenCredential, logging_enable: bool = False, http_proxy: dict = None, - transport_type: TransportType = None, + transport_type: TransportType = TransportType.Amqp, ): def __enter__(self): @@ -30,7 +30,7 @@ def __init__( topic_name: str = None, logging_enable: bool = False, http_proxy: dict = None, - transport_type: TransportType = None, + transport_type: TransportType = TransportType.Amqp, retry_total : int = 3, retry_backoff_factor : float = 30, retry_backoff_maximum : int = 120 # cur_retry_backoff = retry_backoff_factor * (2^cur_retry_time) @@ -84,7 +84,7 @@ def __init__( mode : ReceiveSettleMode = ReceiveSettleMode.PeekLock logging_enable: bool = False, http_proxy: dict = None, - transport_type: TransportType = None, + transport_type: TransportType = TransportType.Amqp, retry_total : int = 3, retry_backoff_factor: float = 30, retry_backoff_maximum: int = 120 # cur_retry_backoff = retry_backoff_factor * (2^cur_retry_time) @@ -244,4 +244,5 @@ class ServiceBusSharedKeyCredential: def __init__(self, policy, key): def get_token(self, *scopes, **kwargs): + class ServiceBusError(Exception): From 190a9128a3b3e0ee2a38cd00e28211b6af73a5fe Mon Sep 17 00:00:00 2001 From: "Adam Ling (MSFT)" <47871814+yunhaoling@users.noreply.github.com> Date: Fri, 6 Mar 2020 11:39:39 -0800 Subject: [PATCH 18/24] Update sdk/servicebus/azure-servicebus/api_samples.py Co-Authored-By: Richard Park <51494936+richardpark-msft@users.noreply.github.com> --- sdk/servicebus/azure-servicebus/api_samples.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/sdk/servicebus/azure-servicebus/api_samples.py b/sdk/servicebus/azure-servicebus/api_samples.py index 734b494c6175..51010422e127 100644 --- a/sdk/servicebus/azure-servicebus/api_samples.py +++ b/sdk/servicebus/azure-servicebus/api_samples.py @@ -129,7 +129,7 @@ print(message) msg.complete() # msg.renew_lock() - # msg.abondon() + # msg.abandon() # msg.defer() # msg.dead_letter() From 68b0e4a4d22d5171677649c63aa56b9e02c92efe Mon Sep 17 00:00:00 2001 From: Yunhao Ling Date: Fri, 6 Mar 2020 13:14:56 -0800 Subject: [PATCH 19/24] remove rule and transaction --- sdk/servicebus/azure-servicebus/api_review.py | 17 ----------------- 1 file changed, 17 deletions(-) diff --git a/sdk/servicebus/azure-servicebus/api_review.py b/sdk/servicebus/azure-servicebus/api_review.py index f81e7ce93713..3be5150e8d3c 100644 --- a/sdk/servicebus/azure-servicebus/api_review.py +++ b/sdk/servicebus/azure-servicebus/api_review.py @@ -130,11 +130,6 @@ def receive( @property def session(self) -> ServiceBusSession: - # Rule APIs, raise Error when called on Queue - def add_rule(self, rule_name : str, filter : str) -> None: - def remove_rule(self, rule_name : str) -> None: - def get_rules(self) -> List[str] -> Dict[str, Any]: - class ServiceBusSession: @property @@ -146,15 +141,6 @@ def session_id(self) -> str: def expired(self) -> bool: -class Transaction: - def __init__(self) -> None: - def __enter__(self): - def __exit__(self): - def begin(self) -> None: - def commit(self) -> None: - def abort(self) -> None: - - class Message: def __init__(self, body : str, encoding : str = 'UTF-8', **kwargs) -> None: def __str__(self): @@ -243,6 +229,3 @@ class TransportType(Enum): class ServiceBusSharedKeyCredential: def __init__(self, policy, key): def get_token(self, *scopes, **kwargs): - - -class ServiceBusError(Exception): From 86e19ef868534e2014dcae73520238eb13b69ca7 Mon Sep 17 00:00:00 2001 From: Yunhao Ling Date: Thu, 12 Mar 2020 16:57:54 -0700 Subject: [PATCH 20/24] update as per discussion with Rayma --- sdk/servicebus/azure-servicebus/api_review.py | 25 ------------------- 1 file changed, 25 deletions(-) diff --git a/sdk/servicebus/azure-servicebus/api_review.py b/sdk/servicebus/azure-servicebus/api_review.py index 3be5150e8d3c..090f5704bfb2 100644 --- a/sdk/servicebus/azure-servicebus/api_review.py +++ b/sdk/servicebus/azure-servicebus/api_review.py @@ -49,8 +49,6 @@ def __exit__(self): def close(self) -> None: - def get_properties(self) -> Dict[str, Any]: - def create_batch( self, max_size_in_bytes : int = None @@ -59,14 +57,12 @@ def create_batch( def send( self, message : Union[Message, BatchMessage], - session_id : str = None, message_timeout : float = None ) -> None: def schedule( self, message : Union[Message, BatchMessage], schedule_time_utc : datetime, - session_id : str = None ) -> List[int]: def cancel_scheduled_messages(self, sequence_number : Union[int, List[int]]) -> None: @@ -106,8 +102,6 @@ def __exit__(self): def close(self) -> None: - def get_properties(self) -> Dict[str, any]: - def __iter__(self): def __next__(self): def next(self): @@ -197,25 +191,6 @@ def abandon(self) -> None: def defer(self) -> None: -class DeferredMessage(ReceivedMessage): - # @properties - def settled(self) -> bool: # read-only - def enqueued_time(self) -> datetime: # read-only - def scheduled_enqueue_time(self) -> datetime: # read-only - def sequence_number(self) -> int: # read-only - def partition_id(self) -> str: # read-only - def locked_until(self) -> datetime: # read-only - def expired(self) -> bool: # read-only - def lock_token(self) -> str: # read-only - def session_id(self) -> str: # read-only - - # methods - def renew_lock(self) -> None: - def complete(self) -> None: - def dead_letter(self, description: str = None) -> None: - def abandon(self) -> None: - - class ReceiveSettleMode: PeekLock ReceiveAndDelete From 785bea98475eea4a081780fd57f411467f7db285 Mon Sep 17 00:00:00 2001 From: Yunhao Ling Date: Fri, 13 Mar 2020 00:41:09 -0700 Subject: [PATCH 21/24] minor update --- sdk/servicebus/azure-servicebus/api_review.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/sdk/servicebus/azure-servicebus/api_review.py b/sdk/servicebus/azure-servicebus/api_review.py index 090f5704bfb2..13001217d814 100644 --- a/sdk/servicebus/azure-servicebus/api_review.py +++ b/sdk/servicebus/azure-servicebus/api_review.py @@ -108,8 +108,8 @@ def next(self): def peek( self, - count : int = 1, - start_from : int = None + message_count : int = 1, + sequence_number : int = None ) -> List[PeekMessage]: def receive_deferred_messages( self, From 30c37b53b551331ae57169ed3b6135a590e3717e Mon Sep 17 00:00:00 2001 From: Kieran Brantner-Magee Date: Wed, 18 Mar 2020 10:30:22 -0700 Subject: [PATCH 22/24] add session_id onto message, and autorenewer class. --- sdk/servicebus/azure-servicebus/api_review.py | 11 ++++++++++- 1 file changed, 10 insertions(+), 1 deletion(-) diff --git a/sdk/servicebus/azure-servicebus/api_review.py b/sdk/servicebus/azure-servicebus/api_review.py index 13001217d814..1855859b247b 100644 --- a/sdk/servicebus/azure-servicebus/api_review.py +++ b/sdk/servicebus/azure-servicebus/api_review.py @@ -136,13 +136,15 @@ def expired(self) -> bool: class Message: - def __init__(self, body : str, encoding : str = 'UTF-8', **kwargs) -> None: + def __init__(self, body : str, encoding : str = 'UTF-8', session_id : str = None, **kwargs) -> None: def __str__(self): # @properties def body(self) -> Union[bytes, Generator[bytes]]: # read-only def partition_key(self, value : str): def partition_key(self) -> str: + def session_id(self, value : str): + def session_id(self) -> str: def via_partition_key(self, value: str): def via_partition_key(self) -> str: def time_to_live(self, value : Union[float, timedelta]): @@ -204,3 +206,10 @@ class TransportType(Enum): class ServiceBusSharedKeyCredential: def __init__(self, policy, key): def get_token(self, *scopes, **kwargs): + + +class AutoLockRenew: + def __init__(self, executor=None, max_workers=None): + def __enter__(self): + def __exit__(self, *args): + def register(self, renewable, timeout=300): \ No newline at end of file From 17ae9c9e57613dfe2f3d01045dfc13067ffbd63f Mon Sep 17 00:00:00 2001 From: Yunhao Ling Date: Thu, 16 Apr 2020 23:41:14 -0700 Subject: [PATCH 23/24] add rule manager api --- sdk/servicebus/azure-servicebus/api_review.py | 62 +++++++++++++++++++ 1 file changed, 62 insertions(+) diff --git a/sdk/servicebus/azure-servicebus/api_review.py b/sdk/servicebus/azure-servicebus/api_review.py index 13001217d814..6bd14c40bc77 100644 --- a/sdk/servicebus/azure-servicebus/api_review.py +++ b/sdk/servicebus/azure-servicebus/api_review.py @@ -21,6 +21,11 @@ def get_queue_receiver(self, queue_name, **kwargs) -> ServiceBusReceiver: def get_topic_sender(self, topic_name, **kwargs) -> ServiceBusSender: def get_subscription_receiver(self, topic_name, subscription_name, **kwargs) -> ServiceBusReceiver: + def get_queue_session_receiver(self, queue_name, session_id, **kwargs) -> ServiceBusSessionReceiver: + def get_subscription_session_receiver(self, topic_name, subscription_name, session_id, **kwargs) -> ServiceBusSessionReceiver: + + def get_subscription_rule_manager(self, topic_name, subscription_name) -> SubscriptionRuleManager: + class ServiceBusSender: def __init__( self, @@ -121,6 +126,8 @@ def receive( timeout : float = None, ) -> List[ReceivedMessage]: # Pull mode receive + +class ServiceBusSessionReceiver(ServiceBusReceiver): @property def session(self) -> ServiceBusSession: @@ -204,3 +211,58 @@ class TransportType(Enum): class ServiceBusSharedKeyCredential: def __init__(self, policy, key): def get_token(self, *scopes, **kwargs): + + +class SubscriptionRuleManager: + def __init__( + self, + fully_qualified_namespace: str, + topic_name: str, + subscription_name: str, + credential: TokenCredential, + logging_enable: bool = False, + http_proxy: dict = None, + transport_type: TransportType = TransportType.Amqp + ) -> SubscriptionRuleManager + + @classmethod + def from_connection_string(cls, conn_str, subscription_name, topic_name=None, **kwargs) -> SubscriptionRuleManager: + + def get_rules(self, timeout: float = None) -> list[RuleDescription]: + def add_rule(self, rule_name: str, filter: Union(bool, str, CorrelationFilter), sql_rule_action_expression: str, timeout: float = None): + def remove_rule(self, rule_name: str, timeout: float = None): + + +class RuleDescription: + #TODO: is it necessary to make it separate class or just a list[dict]? + def __init__(self): + @property + def filter(self) -> Union(str, CorrelationFilter): # TODO: return multiple types OK? + def action(self) -> str: + def name(self) -> str: + + +class CorrelationFilter: + # TODO: is it necessary to make it separate class, depending on the fact that there're so many fields, + # and how people would use this, I personally prefer it being a class + def __init__(self): + + @property + def correlation_id(self, val : str): + def correlation_id(self) -> str: + def message_id(self, val : str): + def message_id(self) -> str: + def to(self, val: str): + def to(self) -> str: + def reply_to(self, val : str): + def reply_to(self) -> str: + def label(self, val: str): + def label(self) -> str: + def session_id(self, val: str): + def session_id(self) -> str: + def reply_to_session_id(self, val: str): + def reply_to_session_id(self) -> str: + def content_type(self, val: str): + def content_type(self) -> str: + def user_properties(self, val: dict): + def user_properties(self) -> dict: From 5c58b7ccd46766d1c96f58e2b1dd70ed3af48772 Mon Sep 17 00:00:00 2001 From: Yunhao Ling Date: Fri, 17 Apr 2020 00:03:10 -0700 Subject: [PATCH 24/24] add api and sample for rule management --- sdk/servicebus/azure-servicebus/api_review.py | 3 ++ .../azure-servicebus/api_samples.py | 35 +++++++++++++++++++ 2 files changed, 38 insertions(+) diff --git a/sdk/servicebus/azure-servicebus/api_review.py b/sdk/servicebus/azure-servicebus/api_review.py index 6bd14c40bc77..456f692edfde 100644 --- a/sdk/servicebus/azure-servicebus/api_review.py +++ b/sdk/servicebus/azure-servicebus/api_review.py @@ -214,6 +214,9 @@ def get_token(self, *scopes, **kwargs): class SubscriptionRuleManager: + # sql and correlation filter apply to user properties and system properties + # https://docs.microsoft.com/en-us/azure/service-bus-messaging/topic-filters + # rule action: https://docs.microsoft.com/en-us/azure/service-bus-messaging/service-bus-messaging-sql-rule-action def __init__( self, fully_qualified_namespace: str, diff --git a/sdk/servicebus/azure-servicebus/api_samples.py b/sdk/servicebus/azure-servicebus/api_samples.py index 51010422e127..146893210fb5 100644 --- a/sdk/servicebus/azure-servicebus/api_samples.py +++ b/sdk/servicebus/azure-servicebus/api_samples.py @@ -217,3 +217,38 @@ msg.complete() session.renew_lock() session.set_session_state("END") + +### Rule operation + +# 1. get rules +subscription_rule_manager = sb_client.get_subscription_rule_manager("topic_name", "subscription_name") +rules = subscription_rule_manager.get_rules() +for rule in rules: + print(rule.name) + print(rule.filter) # If it's type of correlation, print the detail info + print(rule.action) + +# 2. add rule + +# 2.1 add boolean filter rule + +subscription_rule_manager.add_rule("rule_name", filter=True, sql_rule_action_expression=None) + +# 2.2 add sql filter rule + +message = Message("msg") +message.correlation_id = 2 +message.user_properties = {"priority" : 2} +subscription_rule_manager.add_rule("rule_name", filter="priority >= 1 AND sys.correlation-id = 2", sql_rule_action_expression=None) + +# 2.3 add correlation filter rule + +message = Message("msg") +message.correlation_id = 2 +message.user_properties = {"key" : "value"} + +correlation_filter = CorrelationFilter() +correlation_filter.user_properties = {"key": "value"} +correlation_filter.correlation_id = 2 + +subscription_rule_manager.add_rule("rule_name", filter=CorrelationFilter, sql_rule_action_expression="SET key = 'new value'")