diff --git a/sdk/servicebus/azure-servicebus/api_review.py b/sdk/servicebus/azure-servicebus/api_review.py new file mode 100644 index 000000000000..6711642502c7 --- /dev/null +++ b/sdk/servicebus/azure-servicebus/api_review.py @@ -0,0 +1,285 @@ +## APIs: +class ServiceBusClient: + def __init__( + self, + fully_qualified_namespace : str, + credential : TokenCredential, + logging_enable: bool = False, + http_proxy: dict = None, + transport_type: TransportType = TransportType.Amqp, + ): + + 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: + 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, + fully_qualified_namespace : str, + credential: TokenCredential, + queue_name: str = None, + topic_name: str = None, + logging_enable: bool = False, + http_proxy: dict = 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) + ) -> None: + @classmethod + def from_connection_string( + cls, + conn_str : str, + queue_name: str = None, + topic_name: str = None, + **kwargs + ) -> ServiceBusSender: + + def __enter__(self): + def __exit__(self): + + def close(self) -> None: + + def create_batch( + self, + max_size_in_bytes : int = None + ) -> BatchMessage: + + def send( + self, + message : Union[Message, BatchMessage], + message_timeout : float = None + ) -> None: + def schedule( + self, + message : Union[Message, BatchMessage], + schedule_time_utc : datetime, + ) -> List[int]: + + def cancel_scheduled_messages(self, sequence_number : Union[int, List[int]]) -> None: + + +class ServiceBusReceiver: + def __init__( + self, + fully_qualified_namespace : str, + credential: TokenCredential, + queue_name: str = None, + topic_name: str = None, + subscription_name : str = None, + session_id : str = None, + mode : ReceiveSettleMode = ReceiveSettleMode.PeekLock + logging_enable: bool = False, + http_proxy: dict = 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) + ) -> None: + @classmethod + def from_connection_string( + cls, + conn_str : str, + queue_name: str = None, + topic_name: str = None, + subscription_name : str = None, + session_id: str = None, + mode : ReceiveSettleMode = ReceiveSettleMode.PeekLock + **kwargs + )-> ServiceBusReceiver: + + def __enter__(self): + def __exit__(self): + + def close(self) -> None: + + def __iter__(self): + def __next__(self): + def next(self): + + def peek( + self, + message_count : int = 1, + sequence_number : int = None + ) -> List[PeekMessage]: + def receive_deferred_messages( + self, + sequence_numbers : List[int], + ) -> List[DeferredMessage]: + def receive( + self, + max_batch_size : int = None, + timeout : float = None, + ) -> List[ReceivedMessage]: # Pull mode receive + + +class ServiceBusSessionReceiver(ServiceBusReceiver): + @property + def session(self) -> ServiceBusSession: + + +class ServiceBusSession: + @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: + + +class Message: + 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]): + 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_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 + 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: + def defer(self) -> None: + + +class ReceiveSettleMode: + PeekLock + ReceiveAndDelete + + +class TransportType(Enum): + Amqp + AmqpOverWebsocket + + +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): + + +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, + 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: + + + diff --git a/sdk/servicebus/azure-servicebus/api_samples.py b/sdk/servicebus/azure-servicebus/api_samples.py new file mode 100644 index 000000000000..146893210fb5 --- /dev/null +++ b/sdk/servicebus/azure-servicebus/api_samples.py @@ -0,0 +1,254 @@ +## Sample Code: + +### creation of different entity clients + +# Top Level ServiceBusClient +sb_client = ServiceBusClient.from_connection_string(conn_str="conn_str") +sb_client = ServiceBusClient( + 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_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_namespace="fully_qualified_namespace", + queue_name="queue_name", + ServiceBusSharedKeyCredential("policy", "key") +) + +# Queue Receiver +queue_receiver = ServiceBusReceiver.from_connection_string( + conn_str="conn_str", + queue_name="queue_name" +) + +queue_receiver = ServiceBusReceiver( + fully_qualified_namespace="fully_qualified_namespace", + queue_name="queue_name", + credential=ServiceBusSharedKeyCredential("policy", "key"), + session_id="test_session" +) + +# Topic Sender +topic_sender = ServiceBusSender.from_connection_string( + conn_str="conn_str", + topic_name="topic_name" +) + +topic_sender = ServiceBusSender( + fully_qualified_namespace="fully_qualified_namespace", + topic_name="topic_name", + credential=ServiceBusSharedKeyCredential("policy", "key") +) + +# Subscription Receiver +subscription_receiver = ServiceBusReceiver.from_connection_string( + conn_str="conn_str", + topic_name="topic_name", + subscription_name="subscription_name" +) + +subscription_receiver = ServiceBusReceiver( + fully_qualified_namespace="fully_qualified_namespace", + topic_name="topic_name", + subscription_name="subcription_name", + credential=ServiceBusSharedKeyCredential("policy", "key"), + session_id="test_session" +) + + +### 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") + +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 + + 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) + + 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 + + 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) + +# 2 Receive +sb_client = ServiceBusClient.from_connection_string(conn_str="conn_str") +queue_receiver = sb_client.get_queue_receiver(queue_name="queue_name") + +# 2.1 peek +with sb_client: + with queue_receiver: + msgs = queue_receiver.peek(count=10) + for msg in msgs: + print(msg) + +# 2.2 iterator receive +with sb_client: + with queue_receiver: + for msg in queue_receiver: + print(message) + msg.complete() + # msg.renew_lock() + # msg.abandon() + # msg.defer() + # msg.dead_letter() + +# 2.3 pull mode receive +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 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") + +with sb_client: + with topic_sender: + # 1.1 send single + message = Message("Test message") + 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 + + topic_sender.send(batch_message, session_id="test_session") + + # 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") + + batch_schedule_message = topic_sender.create_batch() + + batch_sequence_numbers = topic_sender.schedule( + message=batch_schedule_message, + schedule_time_utc=schedule_time_utc, + session_id="test_session" + ) + + # 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", + session_id="test_session" +) + +# 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") + +### 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'")