From 0618690aaaad725158754c24227855cf6e2581f2 Mon Sep 17 00:00:00 2001 From: swathipil Date: Thu, 24 Oct 2024 15:44:12 -0700 Subject: [PATCH 1/7] [ServiceBus] add service specific message annotations to receiver logs --- sdk/servicebus/azure-servicebus/CHANGELOG.md | 2 ++ .../azure/servicebus/_pyamqp/aio/_receiver_async.py | 6 ++++++ .../azure-servicebus/azure/servicebus/_pyamqp/receiver.py | 6 ++++++ 3 files changed, 14 insertions(+) diff --git a/sdk/servicebus/azure-servicebus/CHANGELOG.md b/sdk/servicebus/azure-servicebus/CHANGELOG.md index b31709522598..27dd3a3865f7 100644 --- a/sdk/servicebus/azure-servicebus/CHANGELOG.md +++ b/sdk/servicebus/azure-servicebus/CHANGELOG.md @@ -12,6 +12,8 @@ ### Other Changes +- Added logging to track received messages. + ## 7.12.3 (2024-09-19) ### Bugs Fixed diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/_pyamqp/aio/_receiver_async.py b/sdk/servicebus/azure-servicebus/azure/servicebus/_pyamqp/aio/_receiver_async.py index aeb587aea2dd..fbc9f4c9a1d5 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/_pyamqp/aio/_receiver_async.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/_pyamqp/aio/_receiver_async.py @@ -78,6 +78,12 @@ async def _incoming_transfer(self, frame): self._received_payload = bytearray() else: message = decode_payload(frame[11]) + _LOGGER.debug( + "Received message: annotations: %r, header: %r", + message.message_annotations, + message.header, + extra=self.network_trace_params + ) delivery_state = await self._process_incoming_message(self._first_frame, message) if not frame[4] and delivery_state: # settled await self._outgoing_disposition( diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/_pyamqp/receiver.py b/sdk/servicebus/azure-servicebus/azure/servicebus/_pyamqp/receiver.py index e3cb8c77f54a..da1b81ba4db3 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/_pyamqp/receiver.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/_pyamqp/receiver.py @@ -76,6 +76,12 @@ def _incoming_transfer(self, frame): self._received_payload = bytearray() else: message = decode_payload(frame[11]) + _LOGGER.debug( + "Received message: annotations: %r, header: %r", + message.message_annotations, + message.header, + extra=self.network_trace_params + ) delivery_state = self._process_incoming_message(self._first_frame, message) if not frame[4] and delivery_state: # settled From 488d3acff2fc91788ff6c4ad9e0546727e8ca817 Mon Sep 17 00:00:00 2001 From: swathipil Date: Thu, 24 Oct 2024 15:50:51 -0700 Subject: [PATCH 2/7] add to eventhub + update readme logging to include thread formatting --- sdk/eventhub/azure-eventhub/CHANGELOG.md | 2 ++ sdk/eventhub/azure-eventhub/README.md | 2 ++ .../azure/eventhub/_pyamqp/aio/_receiver_async.py | 6 ++++++ .../azure-eventhub/azure/eventhub/_pyamqp/receiver.py | 6 ++++++ sdk/servicebus/azure-servicebus/README.md | 2 ++ 5 files changed, 18 insertions(+) diff --git a/sdk/eventhub/azure-eventhub/CHANGELOG.md b/sdk/eventhub/azure-eventhub/CHANGELOG.md index e5229ee5b8f9..97570fce0077 100644 --- a/sdk/eventhub/azure-eventhub/CHANGELOG.md +++ b/sdk/eventhub/azure-eventhub/CHANGELOG.md @@ -10,6 +10,8 @@ ### Other Changes +- Added logging to track received messages. + ## 5.12.2 (2024-10-02) ### Bugs Fixed diff --git a/sdk/eventhub/azure-eventhub/README.md b/sdk/eventhub/azure-eventhub/README.md index 73df7c38e5a8..30297dc498ae 100644 --- a/sdk/eventhub/azure-eventhub/README.md +++ b/sdk/eventhub/azure-eventhub/README.md @@ -453,6 +453,8 @@ import logging import sys handler = logging.StreamHandler(stream=sys.stdout) +log_fmt = logging.Formatter(fmt="%(asctime)s | %(threadName)s | %(levelname)s | %(name)s | %(message)s") +handler.setFormatter(log_fmt) logger = logging.getLogger('azure.eventhub') logger.setLevel(logging.DEBUG) logger.addHandler(handler) diff --git a/sdk/eventhub/azure-eventhub/azure/eventhub/_pyamqp/aio/_receiver_async.py b/sdk/eventhub/azure-eventhub/azure/eventhub/_pyamqp/aio/_receiver_async.py index aeb587aea2dd..fbc9f4c9a1d5 100644 --- a/sdk/eventhub/azure-eventhub/azure/eventhub/_pyamqp/aio/_receiver_async.py +++ b/sdk/eventhub/azure-eventhub/azure/eventhub/_pyamqp/aio/_receiver_async.py @@ -78,6 +78,12 @@ async def _incoming_transfer(self, frame): self._received_payload = bytearray() else: message = decode_payload(frame[11]) + _LOGGER.debug( + "Received message: annotations: %r, header: %r", + message.message_annotations, + message.header, + extra=self.network_trace_params + ) delivery_state = await self._process_incoming_message(self._first_frame, message) if not frame[4] and delivery_state: # settled await self._outgoing_disposition( diff --git a/sdk/eventhub/azure-eventhub/azure/eventhub/_pyamqp/receiver.py b/sdk/eventhub/azure-eventhub/azure/eventhub/_pyamqp/receiver.py index e3cb8c77f54a..da1b81ba4db3 100644 --- a/sdk/eventhub/azure-eventhub/azure/eventhub/_pyamqp/receiver.py +++ b/sdk/eventhub/azure-eventhub/azure/eventhub/_pyamqp/receiver.py @@ -76,6 +76,12 @@ def _incoming_transfer(self, frame): self._received_payload = bytearray() else: message = decode_payload(frame[11]) + _LOGGER.debug( + "Received message: annotations: %r, header: %r", + message.message_annotations, + message.header, + extra=self.network_trace_params + ) delivery_state = self._process_incoming_message(self._first_frame, message) if not frame[4] and delivery_state: # settled diff --git a/sdk/servicebus/azure-servicebus/README.md b/sdk/servicebus/azure-servicebus/README.md index c6e5732cca6f..cd03c01e8f01 100644 --- a/sdk/servicebus/azure-servicebus/README.md +++ b/sdk/servicebus/azure-servicebus/README.md @@ -424,6 +424,8 @@ import logging import sys handler = logging.StreamHandler(stream=sys.stdout) +log_fmt = logging.Formatter(fmt="%(asctime)s | %(threadName)s | %(levelname)s | %(name)s | %(message)s") +handler.setFormatter(log_fmt) logger = logging.getLogger('azure.servicebus') logger.setLevel(logging.DEBUG) logger.addHandler(handler) From aefa1a4e94257ea0520440b20fad75771a114eaf Mon Sep 17 00:00:00 2001 From: swathipil Date: Thu, 31 Oct 2024 12:32:28 -0700 Subject: [PATCH 3/7] only log service annotations header props --- .../servicebus/_pyamqp/aio/_receiver_async.py | 37 ++++++++++++++++--- .../azure/servicebus/_pyamqp/receiver.py | 37 ++++++++++++++++--- 2 files changed, 62 insertions(+), 12 deletions(-) diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/_pyamqp/aio/_receiver_async.py b/sdk/servicebus/azure-servicebus/azure/servicebus/_pyamqp/aio/_receiver_async.py index fbc9f4c9a1d5..8a1991222937 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/_pyamqp/aio/_receiver_async.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/_pyamqp/aio/_receiver_async.py @@ -7,6 +7,7 @@ import uuid import logging from typing import Optional, Union +from datetime import datetime, timezone from .._decode import decode_payload from ._link_async import Link @@ -78,12 +79,36 @@ async def _incoming_transfer(self, frame): self._received_payload = bytearray() else: message = decode_payload(frame[11]) - _LOGGER.debug( - "Received message: annotations: %r, header: %r", - message.message_annotations, - message.header, - extra=self.network_trace_params - ) + if self.network_trace and _LOGGER.isEnabledFor(logging.DEBUG): + seq_num = None + enqueued_time = None + locked_until = None + ttl = None + delivery_count = None + if message.message_annotations: + seq_num = message.message_annotations.get(b"x-opt-sequence-number") + try: + enqueued_time = message.message_annotations.get(b"x-opt-enqueued-time") + enqueued_time = datetime.fromtimestamp(enqueued_time/1000, tz=timezone.utc) + except TypeError: + enqueued_time = None + try: + locked_until = message.message_annotations.get(b"x-opt-locked-until") + enqueued_time = datetime.fromtimestamp(locked_until/1000, tz=timezone.utc) + except TypeError: + locked_until = None + if message.header: + ttl = message.header.ttl + delivery_count = message.header.delivery_count + _LOGGER.debug( + "Received message: seq-num: %r, enqd-utc: %r, lockd-til-utc: %r, ttl: %r, dlvry-cnt: %r", + seq_num, + enqueued_time, + locked_until, + ttl, + delivery_count, + extra=self.network_trace_params + ) delivery_state = await self._process_incoming_message(self._first_frame, message) if not frame[4] and delivery_state: # settled await self._outgoing_disposition( diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/_pyamqp/receiver.py b/sdk/servicebus/azure-servicebus/azure/servicebus/_pyamqp/receiver.py index da1b81ba4db3..1ccde912a29c 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/_pyamqp/receiver.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/_pyamqp/receiver.py @@ -7,6 +7,7 @@ import uuid import logging from typing import Optional, Union +from datetime import datetime, timezone from ._decode import decode_payload from .link import Link @@ -76,12 +77,36 @@ def _incoming_transfer(self, frame): self._received_payload = bytearray() else: message = decode_payload(frame[11]) - _LOGGER.debug( - "Received message: annotations: %r, header: %r", - message.message_annotations, - message.header, - extra=self.network_trace_params - ) + if self.network_trace and _LOGGER.isEnabledFor(logging.DEBUG): + seq_num = None + enqueued_time = None + locked_until = None + ttl = None + delivery_count = None + if message.message_annotations: + seq_num = message.message_annotations.get(b"x-opt-sequence-number") + try: + enqueued_time = message.message_annotations.get(b"x-opt-enqueued-time") + enqueued_time = datetime.fromtimestamp(enqueued_time/1000, tz=timezone.utc) + except TypeError: + enqueued_time = None + try: + locked_until = message.message_annotations.get(b"x-opt-locked-until") + enqueued_time = datetime.fromtimestamp(locked_until/1000, tz=timezone.utc) + except TypeError: + locked_until = None + if message.header: + ttl = message.header.ttl + delivery_count = message.header.delivery_count + _LOGGER.debug( + "Received message: seq-num: %r, enqd-utc: %r, lockd-til-utc: %r, ttl: %r, dlvry-cnt: %r", + seq_num, + enqueued_time, + locked_until, + ttl, + delivery_count, + extra=self.network_trace_params + ) delivery_state = self._process_incoming_message(self._first_frame, message) if not frame[4] and delivery_state: # settled From caf86a37929942d361d2d5f6bcde0c006082d20c Mon Sep 17 00:00:00 2001 From: swathipil Date: Thu, 31 Oct 2024 17:11:39 -0700 Subject: [PATCH 4/7] move logging to sdk layer --- .../azure/servicebus/_pyamqp/receiver.py | 30 ------------------- .../_transport/_pyamqp_transport.py | 11 +++++++ 2 files changed, 11 insertions(+), 30 deletions(-) diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/_pyamqp/receiver.py b/sdk/servicebus/azure-servicebus/azure/servicebus/_pyamqp/receiver.py index 1ccde912a29c..454d0bc6e6b7 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/_pyamqp/receiver.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/_pyamqp/receiver.py @@ -77,36 +77,6 @@ def _incoming_transfer(self, frame): self._received_payload = bytearray() else: message = decode_payload(frame[11]) - if self.network_trace and _LOGGER.isEnabledFor(logging.DEBUG): - seq_num = None - enqueued_time = None - locked_until = None - ttl = None - delivery_count = None - if message.message_annotations: - seq_num = message.message_annotations.get(b"x-opt-sequence-number") - try: - enqueued_time = message.message_annotations.get(b"x-opt-enqueued-time") - enqueued_time = datetime.fromtimestamp(enqueued_time/1000, tz=timezone.utc) - except TypeError: - enqueued_time = None - try: - locked_until = message.message_annotations.get(b"x-opt-locked-until") - enqueued_time = datetime.fromtimestamp(locked_until/1000, tz=timezone.utc) - except TypeError: - locked_until = None - if message.header: - ttl = message.header.ttl - delivery_count = message.header.delivery_count - _LOGGER.debug( - "Received message: seq-num: %r, enqd-utc: %r, lockd-til-utc: %r, ttl: %r, dlvry-cnt: %r", - seq_num, - enqueued_time, - locked_until, - ttl, - delivery_count, - extra=self.network_trace_params - ) delivery_state = self._process_incoming_message(self._first_frame, message) if not frame[4] and delivery_state: # settled diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/_transport/_pyamqp_transport.py b/sdk/servicebus/azure-servicebus/azure/servicebus/_transport/_pyamqp_transport.py index 0bdc78b254cd..7818703db7be 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/_transport/_pyamqp_transport.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/_transport/_pyamqp_transport.py @@ -10,6 +10,7 @@ import datetime from datetime import timezone from typing import Optional, Tuple, cast, List, TYPE_CHECKING, Any, Callable, Dict, Union, Iterator, Type +import logging from .._pyamqp import ( utils, @@ -102,6 +103,7 @@ from .._pyamqp.client import AMQPClient from .._pyamqp.message import MessageDict +_LOGGER = logging.getLogger(__name__) class _ServiceBusErrorPolicy(RetryPolicy): @@ -743,6 +745,15 @@ def build_received_message( amqp_transport=receiver._amqp_transport ) receiver._last_received_sequenced_number = message.sequence_number + if receiver._config.logging_enable and _LOGGER.isEnabledFor(logging.DEBUG): + _LOGGER.debug( + "Received message: seq-num: %r, enqd-utc: %r, lockd-til-utc: %r, ttl: %r, dlvry-cnt: %r", + message.sequence_number, + message.enqueued_time_utc, + message.locked_until_utc, + message.time_to_live, + message.delivery_count, + ) return message @staticmethod From e16056d7fcc63df973a519dc35e4c4e0b1d42fa9 Mon Sep 17 00:00:00 2001 From: swathipil Date: Thu, 31 Oct 2024 17:37:33 -0700 Subject: [PATCH 5/7] move msg logging into eh consumer --- sdk/eventhub/azure-eventhub/azure/eventhub/_consumer.py | 7 +++++++ .../azure/eventhub/_pyamqp/aio/_receiver_async.py | 6 ------ .../azure-eventhub/azure/eventhub/_pyamqp/receiver.py | 6 ------ .../azure-eventhub/azure/eventhub/aio/_consumer_async.py | 7 +++++++ 4 files changed, 14 insertions(+), 12 deletions(-) diff --git a/sdk/eventhub/azure-eventhub/azure/eventhub/_consumer.py b/sdk/eventhub/azure-eventhub/azure/eventhub/_consumer.py index 49a992caab56..df4a89b91756 100644 --- a/sdk/eventhub/azure-eventhub/azure/eventhub/_consumer.py +++ b/sdk/eventhub/azure-eventhub/azure/eventhub/_consumer.py @@ -185,6 +185,13 @@ def _next_message_in_buffer(self): # pylint:disable=protected-access message = self._message_buffer.popleft() event_data = EventData._from_message(message) + if self._client._config.network_tracing and _LOGGER.isEnabledFor(logging.DEBUG): + _LOGGER.debug( + "Received message: seq-num: %r, offset: %r, partition-key: %r", + event_data.sequence_number, + event_data.offset, + event_data.partition_key, + ) if self._amqp_transport.KIND == "uamqp": event_data._uamqp_message = message self._last_received_event = event_data diff --git a/sdk/eventhub/azure-eventhub/azure/eventhub/_pyamqp/aio/_receiver_async.py b/sdk/eventhub/azure-eventhub/azure/eventhub/_pyamqp/aio/_receiver_async.py index fbc9f4c9a1d5..aeb587aea2dd 100644 --- a/sdk/eventhub/azure-eventhub/azure/eventhub/_pyamqp/aio/_receiver_async.py +++ b/sdk/eventhub/azure-eventhub/azure/eventhub/_pyamqp/aio/_receiver_async.py @@ -78,12 +78,6 @@ async def _incoming_transfer(self, frame): self._received_payload = bytearray() else: message = decode_payload(frame[11]) - _LOGGER.debug( - "Received message: annotations: %r, header: %r", - message.message_annotations, - message.header, - extra=self.network_trace_params - ) delivery_state = await self._process_incoming_message(self._first_frame, message) if not frame[4] and delivery_state: # settled await self._outgoing_disposition( diff --git a/sdk/eventhub/azure-eventhub/azure/eventhub/_pyamqp/receiver.py b/sdk/eventhub/azure-eventhub/azure/eventhub/_pyamqp/receiver.py index da1b81ba4db3..e3cb8c77f54a 100644 --- a/sdk/eventhub/azure-eventhub/azure/eventhub/_pyamqp/receiver.py +++ b/sdk/eventhub/azure-eventhub/azure/eventhub/_pyamqp/receiver.py @@ -76,12 +76,6 @@ def _incoming_transfer(self, frame): self._received_payload = bytearray() else: message = decode_payload(frame[11]) - _LOGGER.debug( - "Received message: annotations: %r, header: %r", - message.message_annotations, - message.header, - extra=self.network_trace_params - ) delivery_state = self._process_incoming_message(self._first_frame, message) if not frame[4] and delivery_state: # settled diff --git a/sdk/eventhub/azure-eventhub/azure/eventhub/aio/_consumer_async.py b/sdk/eventhub/azure-eventhub/azure/eventhub/aio/_consumer_async.py index b49e19f7fc4e..b776aacef9ac 100644 --- a/sdk/eventhub/azure-eventhub/azure/eventhub/aio/_consumer_async.py +++ b/sdk/eventhub/azure-eventhub/azure/eventhub/aio/_consumer_async.py @@ -180,6 +180,13 @@ def _next_message_in_buffer(self): # pylint:disable=protected-access message = self._message_buffer.popleft() event_data = EventData._from_message(message) + if self._client._config.network_tracing and _LOGGER.isEnabledFor(logging.DEBUG): + _LOGGER.debug( + "Received message: seq-num: %r, offset: %r, partition-key: %r", + event_data.sequence_number, + event_data.offset, + event_data.partition_key, + ) if self._amqp_transport.KIND == "uamqp": event_data._uamqp_message = message self._last_received_event = event_data From 03b11a80a0e8b51d8ec37751184b5c3e22c6edc2 Mon Sep 17 00:00:00 2001 From: swathipil Date: Thu, 31 Oct 2024 17:42:08 -0700 Subject: [PATCH 6/7] remove sb pyamqp recvr logs --- .../servicebus/_pyamqp/aio/_receiver_async.py | 31 ------------------- .../azure/servicebus/_pyamqp/receiver.py | 1 - 2 files changed, 32 deletions(-) diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/_pyamqp/aio/_receiver_async.py b/sdk/servicebus/azure-servicebus/azure/servicebus/_pyamqp/aio/_receiver_async.py index 8a1991222937..aeb587aea2dd 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/_pyamqp/aio/_receiver_async.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/_pyamqp/aio/_receiver_async.py @@ -7,7 +7,6 @@ import uuid import logging from typing import Optional, Union -from datetime import datetime, timezone from .._decode import decode_payload from ._link_async import Link @@ -79,36 +78,6 @@ async def _incoming_transfer(self, frame): self._received_payload = bytearray() else: message = decode_payload(frame[11]) - if self.network_trace and _LOGGER.isEnabledFor(logging.DEBUG): - seq_num = None - enqueued_time = None - locked_until = None - ttl = None - delivery_count = None - if message.message_annotations: - seq_num = message.message_annotations.get(b"x-opt-sequence-number") - try: - enqueued_time = message.message_annotations.get(b"x-opt-enqueued-time") - enqueued_time = datetime.fromtimestamp(enqueued_time/1000, tz=timezone.utc) - except TypeError: - enqueued_time = None - try: - locked_until = message.message_annotations.get(b"x-opt-locked-until") - enqueued_time = datetime.fromtimestamp(locked_until/1000, tz=timezone.utc) - except TypeError: - locked_until = None - if message.header: - ttl = message.header.ttl - delivery_count = message.header.delivery_count - _LOGGER.debug( - "Received message: seq-num: %r, enqd-utc: %r, lockd-til-utc: %r, ttl: %r, dlvry-cnt: %r", - seq_num, - enqueued_time, - locked_until, - ttl, - delivery_count, - extra=self.network_trace_params - ) delivery_state = await self._process_incoming_message(self._first_frame, message) if not frame[4] and delivery_state: # settled await self._outgoing_disposition( diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/_pyamqp/receiver.py b/sdk/servicebus/azure-servicebus/azure/servicebus/_pyamqp/receiver.py index 454d0bc6e6b7..e3cb8c77f54a 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/_pyamqp/receiver.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/_pyamqp/receiver.py @@ -7,7 +7,6 @@ import uuid import logging from typing import Optional, Union -from datetime import datetime, timezone from ._decode import decode_payload from .link import Link From 231349a100007d5e069a198271cee386b105bced Mon Sep 17 00:00:00 2001 From: swathipil Date: Thu, 31 Oct 2024 20:47:26 -0700 Subject: [PATCH 7/7] black --- .../azure/servicebus/_transport/_pyamqp_transport.py | 1 + 1 file changed, 1 insertion(+) diff --git a/sdk/servicebus/azure-servicebus/azure/servicebus/_transport/_pyamqp_transport.py b/sdk/servicebus/azure-servicebus/azure/servicebus/_transport/_pyamqp_transport.py index 1530d02f8773..8d694b8d7998 100644 --- a/sdk/servicebus/azure-servicebus/azure/servicebus/_transport/_pyamqp_transport.py +++ b/sdk/servicebus/azure-servicebus/azure/servicebus/_transport/_pyamqp_transport.py @@ -105,6 +105,7 @@ _LOGGER = logging.getLogger(__name__) + class _ServiceBusErrorPolicy(RetryPolicy): no_retry = RetryPolicy.no_retry + cast(