From 5c8546121fe5162e8e854fefc00c86d02007766f Mon Sep 17 00:00:00 2001 From: Rakshith Bhyravabhotla Date: Thu, 22 Jul 2021 17:25:52 -0700 Subject: [PATCH 01/21] Draft for magic events --- sdk/core/azure-core/azure/core/messaging.py | 26 +++++++++++++++++-- .../azure/eventgrid/_models.py | 9 ++++++- .../consume_cloud_events_from_eventhub.py | 6 ++--- ...consume_cloud_events_from_storage_queue.py | 4 +-- ...eventgrid_events_from_service_bus_queue.py | 2 +- 5 files changed, 37 insertions(+), 10 deletions(-) diff --git a/sdk/core/azure-core/azure/core/messaging.py b/sdk/core/azure-core/azure/core/messaging.py index a677ae462dd1..fc8bfbd6435d 100644 --- a/sdk/core/azure-core/azure/core/messaging.py +++ b/sdk/core/azure-core/azure/core/messaging.py @@ -5,6 +5,7 @@ # license information. # -------------------------------------------------------------------------- import uuid +import json from base64 import b64decode from datetime import datetime from .utils._utils import _convert_to_isoformat, TZ_UTC @@ -21,8 +22,28 @@ __all__ = ["CloudEvent"] - -class CloudEvent(object): # pylint:disable=too-many-instance-attributes +class _EventMixin(object): + """Event mixin to have methods that are common to different Event types + like CloudEvent, EventGridEvent etc. + """ + @staticmethod + def _get_bytes(obj): + # type: (Any) -> Dict + try: + # storage queue + return json.loads(obj.content) + except AttributeError: + # eventhubs + try: + return json.loads(next(obj.body))[0] + except KeyError: + # servicebus + return json.loads(next(obj.body)) + except: + return obj + + +class CloudEvent(_EventMixin): # pylint:disable=too-many-instance-attributes """Properties of the CloudEvent 1.0 Schema. All required parameters must be populated in order to send to Azure. @@ -122,6 +143,7 @@ def from_dict(cls, event): :type event: dict :rtype: CloudEvent """ + event = CloudEvent._get_bytes(event) kwargs = {} # type: Dict[Any, Any] reserved_attr = [ "data", diff --git a/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py b/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py index 02f79cb2dcc4..5ce6663707db 100644 --- a/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py +++ b/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py @@ -3,16 +3,18 @@ # Licensed under the MIT License. See License.txt in the project root for license information. # -------------------------------------------------------------------------------------------- # pylint:disable=protected-access +from _typeshed import Self from typing import Any import datetime as dt import uuid +from azure.core.messaging import _EventMixin from msrest.serialization import UTC from ._generated.models import ( EventGridEvent as InternalEventGridEvent, ) -class EventGridEvent(InternalEventGridEvent): +class EventGridEvent(InternalEventGridEvent, _EventMixin): """Properties of an event published to an Event Grid topic using the EventGrid Schema. Variables are only populated by the server, and will be ignored when sending a request. @@ -96,3 +98,8 @@ def __repr__(self): return "EventGridEvent(subject={}, event_type={}, id={}, event_time={})".format( self.subject, self.event_type, self.id, self.event_time )[:1024] + + @classmethod + def from_dict(cls, event): + event = _EventMixin._get_bytes(event) + super(EventGridEvent).from_dict(event) diff --git a/sdk/eventgrid/azure-eventgrid/samples/consume_samples/consume_cloud_events_from_eventhub.py b/sdk/eventgrid/azure-eventgrid/samples/consume_samples/consume_cloud_events_from_eventhub.py index b372d4ea6752..9ce79e5d587e 100644 --- a/sdk/eventgrid/azure-eventgrid/samples/consume_samples/consume_cloud_events_from_eventhub.py +++ b/sdk/eventgrid/azure-eventgrid/samples/consume_samples/consume_cloud_events_from_eventhub.py @@ -24,11 +24,9 @@ CONNECTION_STR = os.environ["EVENT_HUB_CONN_STR"] EVENTHUB_NAME = os.environ["EVENT_HUB_NAME"] - - def on_event(partition_context, event): - dict_event = CloudEvent.from_dict(json.loads(event)[0]) - print("data: {}\n".format(deserialized_event.data)) + dict_event = CloudEvent.from_dict(event) + print("data: {}\n".format(dict_event.data)) consumer_client = EventHubConsumerClient.from_connection_string( conn_str=CONNECTION_STR, diff --git a/sdk/eventgrid/azure-eventgrid/samples/consume_samples/consume_cloud_events_from_storage_queue.py b/sdk/eventgrid/azure-eventgrid/samples/consume_samples/consume_cloud_events_from_storage_queue.py index 2b94f5816c2b..148906586809 100644 --- a/sdk/eventgrid/azure-eventgrid/samples/consume_samples/consume_cloud_events_from_storage_queue.py +++ b/sdk/eventgrid/azure-eventgrid/samples/consume_samples/consume_cloud_events_from_storage_queue.py @@ -27,10 +27,10 @@ payload = qsc.get_queue_client( queue=queue_name, message_decode_policy=BinaryBase64DecodePolicy() - ).peek_messages() + ).peek_messages(max_messages=32) ## deserialize payload into a list of typed Events - events = [CloudEvent.from_dict(json.loads(msg.content)) for msg in payload] + events = [CloudEvent.from_dict(msg) for msg in payload] for event in events: print(type(event)) ## CloudEvent diff --git a/sdk/eventgrid/azure-eventgrid/samples/consume_samples/consume_eventgrid_events_from_service_bus_queue.py b/sdk/eventgrid/azure-eventgrid/samples/consume_samples/consume_eventgrid_events_from_service_bus_queue.py index 0756043a9ac7..2f13f22ba905 100644 --- a/sdk/eventgrid/azure-eventgrid/samples/consume_samples/consume_eventgrid_events_from_service_bus_queue.py +++ b/sdk/eventgrid/azure-eventgrid/samples/consume_samples/consume_eventgrid_events_from_service_bus_queue.py @@ -30,7 +30,7 @@ payload = sb_client.get_queue_receiver(queue_name).receive_messages() ## deserialize payload into a list of typed Events - events = [EventGridEvent.from_dict(json.loads(next(msg.body).decode('utf-8'))) for msg in payload] + events = [EventGridEvent.from_dict(msg) for msg in payload] for event in events: print(type(event)) ## EventGridEvent From 944ed4e4654c483b3205811166f6f9c2d40c2914 Mon Sep 17 00:00:00 2001 From: Rakshith Bhyravabhotla Date: Thu, 22 Jul 2021 17:27:31 -0700 Subject: [PATCH 02/21] Update sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py --- sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py b/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py index 5ce6663707db..48b8ea9b4b8d 100644 --- a/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py +++ b/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py @@ -101,5 +101,5 @@ def __repr__(self): @classmethod def from_dict(cls, event): - event = _EventMixin._get_bytes(event) + event = EventGridEvent._get_bytes(event) super(EventGridEvent).from_dict(event) From 7b0bb8fab86287b9fe28c09c122bae5174c52363 Mon Sep 17 00:00:00 2001 From: Rakshith Bhyravabhotla Date: Thu, 22 Jul 2021 17:28:16 -0700 Subject: [PATCH 03/21] Update sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py --- sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py | 1 - 1 file changed, 1 deletion(-) diff --git a/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py b/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py index 48b8ea9b4b8d..cb968569ec55 100644 --- a/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py +++ b/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py @@ -3,7 +3,6 @@ # Licensed under the MIT License. See License.txt in the project root for license information. # -------------------------------------------------------------------------------------------- # pylint:disable=protected-access -from _typeshed import Self from typing import Any import datetime as dt import uuid From 56bb5d7b2be9f59a54a47623290b97f30852ae04 Mon Sep 17 00:00:00 2001 From: Rakshith Bhyravabhotla Date: Thu, 22 Jul 2021 17:29:35 -0700 Subject: [PATCH 04/21] Update sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py --- sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py b/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py index cb968569ec55..11f7cec8ee97 100644 --- a/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py +++ b/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py @@ -101,4 +101,4 @@ def __repr__(self): @classmethod def from_dict(cls, event): event = EventGridEvent._get_bytes(event) - super(EventGridEvent).from_dict(event) + super(EventGridEvent, cls).from_dict(event) From b2a90e29ebddd485b9b539551171488d2191d52c Mon Sep 17 00:00:00 2001 From: Rakshith Bhyravabhotla Date: Mon, 26 Jul 2021 16:19:26 -0700 Subject: [PATCH 05/21] duplicate code --- sdk/core/azure-core/azure/core/messaging.py | 41 +++++++++---------- .../azure/eventgrid/_models.py | 23 ++++++++++- 2 files changed, 41 insertions(+), 23 deletions(-) diff --git a/sdk/core/azure-core/azure/core/messaging.py b/sdk/core/azure-core/azure/core/messaging.py index fc8bfbd6435d..6748c18c49eb 100644 --- a/sdk/core/azure-core/azure/core/messaging.py +++ b/sdk/core/azure-core/azure/core/messaging.py @@ -22,28 +22,8 @@ __all__ = ["CloudEvent"] -class _EventMixin(object): - """Event mixin to have methods that are common to different Event types - like CloudEvent, EventGridEvent etc. - """ - @staticmethod - def _get_bytes(obj): - # type: (Any) -> Dict - try: - # storage queue - return json.loads(obj.content) - except AttributeError: - # eventhubs - try: - return json.loads(next(obj.body))[0] - except KeyError: - # servicebus - return json.loads(next(obj.body)) - except: - return obj - -class CloudEvent(_EventMixin): # pylint:disable=too-many-instance-attributes +class CloudEvent(object): # pylint:disable=too-many-instance-attributes """Properties of the CloudEvent 1.0 Schema. All required parameters must be populated in order to send to Azure. @@ -203,3 +183,22 @@ def from_dict(cls, event): " The `source` and `type` params are required." ) return event_obj + + @staticmethod + def _get_bytes(obj): + """Event mixin to have methods that are common to different Event types + like CloudEvent, EventGridEvent etc. + """ + # type: (Any) -> Dict + try: + # storage queue + return json.loads(obj.content) + except AttributeError: + # eventhubs + try: + return json.loads(next(obj.body))[0] + except KeyError: + # servicebus + return json.loads(next(obj.body)) + except: + return obj \ No newline at end of file diff --git a/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py b/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py index 11f7cec8ee97..02e0c4e06c45 100644 --- a/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py +++ b/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py @@ -6,14 +6,14 @@ from typing import Any import datetime as dt import uuid -from azure.core.messaging import _EventMixin +import json from msrest.serialization import UTC from ._generated.models import ( EventGridEvent as InternalEventGridEvent, ) -class EventGridEvent(InternalEventGridEvent, _EventMixin): +class EventGridEvent(InternalEventGridEvent): """Properties of an event published to an Event Grid topic using the EventGrid Schema. Variables are only populated by the server, and will be ignored when sending a request. @@ -102,3 +102,22 @@ def __repr__(self): def from_dict(cls, event): event = EventGridEvent._get_bytes(event) super(EventGridEvent, cls).from_dict(event) + + @staticmethod + def _get_bytes(obj): + """Event mixin to have methods that are common to different Event types + like EventGridEvent etc. + """ + # type: (Any) -> Dict + try: + # storage queue + return json.loads(obj.content) + except AttributeError: + # eventhubs + try: + return json.loads(next(obj.body))[0] + except KeyError: + # servicebus + return json.loads(next(obj.body)) + except: + return obj From 393ccb08482678fc72d5f9cf9eeaf9ef009bd99e Mon Sep 17 00:00:00 2001 From: Rakshith Bhyravabhotla Date: Mon, 26 Jul 2021 22:44:05 -0700 Subject: [PATCH 06/21] lint --- sdk/core/azure-core/azure/core/messaging.py | 6 +++--- .../azure-eventgrid/azure/eventgrid/_models.py | 14 +++++++------- 2 files changed, 10 insertions(+), 10 deletions(-) diff --git a/sdk/core/azure-core/azure/core/messaging.py b/sdk/core/azure-core/azure/core/messaging.py index 6748c18c49eb..67b8b2e7ca0f 100644 --- a/sdk/core/azure-core/azure/core/messaging.py +++ b/sdk/core/azure-core/azure/core/messaging.py @@ -186,10 +186,10 @@ def from_dict(cls, event): @staticmethod def _get_bytes(obj): + # type: (Any) -> Dict """Event mixin to have methods that are common to different Event types like CloudEvent, EventGridEvent etc. """ - # type: (Any) -> Dict try: # storage queue return json.loads(obj.content) @@ -200,5 +200,5 @@ def _get_bytes(obj): except KeyError: # servicebus return json.loads(next(obj.body)) - except: - return obj \ No newline at end of file + except: # pylint: disable=bare-except + return obj diff --git a/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py b/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py index 02e0c4e06c45..c894f93ace53 100644 --- a/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py +++ b/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py @@ -3,7 +3,7 @@ # Licensed under the MIT License. See License.txt in the project root for license information. # -------------------------------------------------------------------------------------------- # pylint:disable=protected-access -from typing import Any +from typing import Any, Dict import datetime as dt import uuid import json @@ -97,18 +97,18 @@ def __repr__(self): return "EventGridEvent(subject={}, event_type={}, id={}, event_time={})".format( self.subject, self.event_type, self.id, self.event_time )[:1024] - + @classmethod - def from_dict(cls, event): - event = EventGridEvent._get_bytes(event) - super(EventGridEvent, cls).from_dict(event) + def from_dict(cls, data, key_extractors=None, content_type=None): + event = EventGridEvent._get_bytes(data) + super(EventGridEvent, cls).from_dict(event, key_extractors, content_type) @staticmethod def _get_bytes(obj): + # type: (Any) -> Dict """Event mixin to have methods that are common to different Event types like EventGridEvent etc. """ - # type: (Any) -> Dict try: # storage queue return json.loads(obj.content) @@ -119,5 +119,5 @@ def _get_bytes(obj): except KeyError: # servicebus return json.loads(next(obj.body)) - except: + except: # pylint: disable=bare-except return obj From d9c5481784073323810ab08905fd724c908bbd7f Mon Sep 17 00:00:00 2001 From: Rakshith Bhyravabhotla Date: Mon, 26 Jul 2021 23:12:59 -0700 Subject: [PATCH 07/21] mypy --- sdk/core/azure-core/azure/core/messaging.py | 1 - sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py | 1 - 2 files changed, 2 deletions(-) diff --git a/sdk/core/azure-core/azure/core/messaging.py b/sdk/core/azure-core/azure/core/messaging.py index 67b8b2e7ca0f..6e0fee8ab6ab 100644 --- a/sdk/core/azure-core/azure/core/messaging.py +++ b/sdk/core/azure-core/azure/core/messaging.py @@ -186,7 +186,6 @@ def from_dict(cls, event): @staticmethod def _get_bytes(obj): - # type: (Any) -> Dict """Event mixin to have methods that are common to different Event types like CloudEvent, EventGridEvent etc. """ diff --git a/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py b/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py index c894f93ace53..53b204d334a9 100644 --- a/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py +++ b/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py @@ -105,7 +105,6 @@ def from_dict(cls, data, key_extractors=None, content_type=None): @staticmethod def _get_bytes(obj): - # type: (Any) -> Dict """Event mixin to have methods that are common to different Event types like EventGridEvent etc. """ From e38a4d5e2817379f51a03580cba6b487dba3d7f8 Mon Sep 17 00:00:00 2001 From: Rakshith Bhyravabhotla Date: Tue, 27 Jul 2021 09:28:54 -0700 Subject: [PATCH 08/21] changelog --- sdk/core/azure-core/CHANGELOG.md | 4 +++- sdk/core/azure-core/azure/core/_version.py | 2 +- sdk/eventgrid/azure-eventgrid/CHANGELOG.md | 6 ++++++ sdk/eventgrid/azure-eventgrid/azure/eventgrid/_version.py | 2 +- 4 files changed, 11 insertions(+), 3 deletions(-) diff --git a/sdk/core/azure-core/CHANGELOG.md b/sdk/core/azure-core/CHANGELOG.md index 30e247c25865..9f206f680052 100644 --- a/sdk/core/azure-core/CHANGELOG.md +++ b/sdk/core/azure-core/CHANGELOG.md @@ -1,9 +1,11 @@ # Release History -## 1.16.1 (Unreleased) +## 1.17.0 (Unreleased) ### Features Added +- `CloudEvent`'s `from_dict` method now accepts objects from servicebus, eventhubs and storage directly. + ### Breaking Changes ### Key Bugs Fixed diff --git a/sdk/core/azure-core/azure/core/_version.py b/sdk/core/azure-core/azure/core/_version.py index b9b4f5d4b07b..986690699f98 100644 --- a/sdk/core/azure-core/azure/core/_version.py +++ b/sdk/core/azure-core/azure/core/_version.py @@ -9,4 +9,4 @@ # regenerated. # -------------------------------------------------------------------------- -VERSION = "1.16.1" +VERSION = "1.17.0" diff --git a/sdk/eventgrid/azure-eventgrid/CHANGELOG.md b/sdk/eventgrid/azure-eventgrid/CHANGELOG.md index dae285744e73..04b08d54d94c 100644 --- a/sdk/eventgrid/azure-eventgrid/CHANGELOG.md +++ b/sdk/eventgrid/azure-eventgrid/CHANGELOG.md @@ -1,5 +1,11 @@ # Release History +## 4.5.0 (Unreleased) + +### Features Added + +- `EventGridEvent`'s `from_dict` method now accepts objects from servicebus, eventhubs and storage directly. + ## 4.4.0 (2021-07-19) - Bumped `msrest` dependency to `0.6.21` to align with mgmt package. diff --git a/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_version.py b/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_version.py index b5234b1c4677..5638bcf9e648 100644 --- a/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_version.py +++ b/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_version.py @@ -9,4 +9,4 @@ # regenerated. # -------------------------------------------------------------------------- -VERSION = "4.4.0" +VERSION = "4.5.0" From ce32b6694df36220d84a3e43b8d77b2d6192aabc Mon Sep 17 00:00:00 2001 From: Rakshith Bhyravabhotla Date: Tue, 27 Jul 2021 16:41:04 -0700 Subject: [PATCH 09/21] tests --- sdk/core/azure-core/azure/core/messaging.py | 10 + .../tests/test_messaging_cloud_event.py | 187 ++++++++++++++++- .../azure/eventgrid/_models.py | 12 +- .../tests/test_eg_event_get_bytes.py | 190 ++++++++++++++++++ 4 files changed, 397 insertions(+), 2 deletions(-) create mode 100644 sdk/eventgrid/azure-eventgrid/tests/test_eg_event_get_bytes.py diff --git a/sdk/core/azure-core/azure/core/messaging.py b/sdk/core/azure-core/azure/core/messaging.py index 6e0fee8ab6ab..d3bf0faf4e35 100644 --- a/sdk/core/azure-core/azure/core/messaging.py +++ b/sdk/core/azure-core/azure/core/messaging.py @@ -192,6 +192,11 @@ def _get_bytes(obj): try: # storage queue return json.loads(obj.content) + except ValueError: + raise ValueError( + "Failed to retrieve content from the object. Make sure the " + + "content follows the CloudEvent schema." + ) except AttributeError: # eventhubs try: @@ -199,5 +204,10 @@ def _get_bytes(obj): except KeyError: # servicebus return json.loads(next(obj.body)) + except ValueError: + raise ValueError( + "Failed to retrieve body from the object. Make sure the " + + "body follows the CloudEvent schema." + ) except: # pylint: disable=bare-except return obj diff --git a/sdk/core/azure-core/tests/test_messaging_cloud_event.py b/sdk/core/azure-core/tests/test_messaging_cloud_event.py index c2d76b5a002c..c1042b02ea4f 100644 --- a/sdk/core/azure-core/tests/test_messaging_cloud_event.py +++ b/sdk/core/azure-core/tests/test_messaging_cloud_event.py @@ -4,12 +4,89 @@ # ------------------------------------ import pytest import json +import uuid import datetime from azure.core.messaging import CloudEvent from azure.core.utils._utils import _convert_to_isoformat from azure.core.serialization import NULL +class MockQueueMessage(object): + def __init__(self, content=None): + self.id = uuid.uuid4() + self.inserted_on = datetime.datetime.now() + self.expires_on = datetime.datetime.now() + datetime.timedelta(days=100) + self.dequeue_count = 1 + self.content = content + self.pop_receipt = None + self.next_visible_on = None + +class MockServiceBusReceivedMessage(object): + def __init__(self, body=None, **kwargs): + self.body=body + self.application_properties=None + self.session_id=None + self.message_id='3f6c5441-5be5-4f33-80c3-3ffeb6a090ce' + self.content_type='application/cloudevents+json; charset=utf-8' + self.correlation_id=None + self.to=None + self.reply_to=None + self.reply_to_session_id=None + self.subject=None + self.time_to_live=datetime.timedelta(days=14) + self.partition_key=None + self.scheduled_enqueue_time_utc=None + self.auto_renew_error=None, + self.dead_letter_error_description=None + self.dead_letter_reason=None + self.dead_letter_source=None + self.delivery_count=13 + self.enqueued_sequence_number=0 + self.enqueued_time_utc=datetime.datetime(2021, 7, 22, 22, 27, 41, 236000, tzinfo=datetime.timezone.utc) + self.expires_at_utc=datetime.datetime(2021, 8, 5, 22, 27, 41, 236000, tzinfo=datetime.timezone.utc) + self.sequence_number=11219 + self.lock_token='233146e3-d5a6-45eb-826f-691d82fb8b13' + +class MockEventhubData(object): + def __init__(self, body=None): + self._last_enqueued_event_properties = {} + self._sys_properties = None + if body is None: + raise ValueError("EventData cannot be None.") + + # Internal usage only for transforming AmqpAnnotatedMessage to outgoing EventData + self.body=body + self._raw_amqp_message = "some amqp data" + self.message_id = None + self.content_type = None + self.correlation_id = None + + +class MockBody(object): + def __init__(self, data=None): + self.data = data + + def __iter__(self): + return self + + def __next__(self): + if not self.data: + return """{"id":"f208feff-099b-4bda-a341-4afd0fa02fef","source":"https://egsample.dev/sampleevent","data":"ServiceBus","type":"Azure.Sdk.Sample","time":"2021-07-22T22:27:38.960209Z","specversion":"1.0"}""" + return self.data + +class MockEhBody(object): + def __init__(self, data=None): + self.data = data + + def __iter__(self): + return self + + def __next__(self): + if not self.data: + return b'[{"id":"f208feff-099b-4bda-a341-4afd0fa02fef","source":"https://egsample.dev/sampleevent","data":"Eventhub","type":"Azure.Sdk.Sample","time":"2021-07-22T22:27:38.960209Z","specversion":"1.0"}]' + return self.data + + # Cloud Event tests def test_cloud_event_constructor(): event = CloudEvent( @@ -469,4 +546,112 @@ def test_wrong_schema_raises_no_type(): "specversion":"1.0", } with pytest.raises(ValueError, match="The event does not conform to the cloud event spec https://github.com/cloudevents/spec. The `source` and `type` params are required."): - CloudEvent.from_dict(cloud_custom_dict) \ No newline at end of file + CloudEvent.from_dict(cloud_custom_dict) + +def test_get_bytes_storage_queue(): + cloud_storage_dict = """{ + "id":"a0517898-9fa4-4e70-b4a3-afda1dd68672", + "source":"/subscriptions/{subscription-id}/resourceGroups/{resource-group}/providers/Microsoft.Storage/storageAccounts/{storage-account}", + "data":{ + "api":"PutBlockList", + "client_request_id":"6d79dbfb-0e37-4fc4-981f-442c9ca65760", + "request_id":"831e1650-001e-001b-66ab-eeb76e000000", + "e_tag":"0x8D4BCC2E4835CD0", + "content_type":"application/octet-stream", + "content_length":524288, + "blob_type":"BlockBlob", + "url":"https://oc2d2817345i60006.blob.core.windows.net/oc2d2817345i200097container/oc2d2817345i20002296blob", + "sequencer":"00000000000004420000000000028963", + "storage_diagnostics":{"batchId":"b68529f3-68cd-4744-baa4-3c0498ec19f0"} + }, + "type":"Microsoft.Storage.BlobCreated", + "time":"2021-02-18T20:18:10.581147898Z", + "specversion":"1.0" + }""" + obj = MockQueueMessage(content=cloud_storage_dict) + + dict = CloudEvent._get_bytes(obj) + assert dict.get('data') == { + "api":"PutBlockList", + "client_request_id":"6d79dbfb-0e37-4fc4-981f-442c9ca65760", + "request_id":"831e1650-001e-001b-66ab-eeb76e000000", + "e_tag":"0x8D4BCC2E4835CD0", + "content_type":"application/octet-stream", + "content_length":524288, + "blob_type":"BlockBlob", + "url":"https://oc2d2817345i60006.blob.core.windows.net/oc2d2817345i200097container/oc2d2817345i20002296blob", + "sequencer":"00000000000004420000000000028963", + "storage_diagnostics":{"batchId":"b68529f3-68cd-4744-baa4-3c0498ec19f0"} + } + assert dict.get('specversion') == "1.0" + +def test_get_bytes_storage_queue_wrong_content(): + cloud_storage_string = u'This is a random string which must fail' + obj = MockQueueMessage(content=cloud_storage_string) + + with pytest.raises(ValueError, match="Failed to retrieve content from the object. Make sure the content follows the CloudEvent schema."): + CloudEvent.from_dict(obj) + +def test_get_bytes_servicebus(): + obj = MockServiceBusReceivedMessage( + body=MockBody(), + message_id='3f6c5441-5be5-4f33-80c3-3ffeb6a090ce', + content_type='application/cloudevents+json; charset=utf-8', + time_to_live=datetime.timedelta(days=14), + delivery_count=13, + enqueued_sequence_number=0, + enqueued_time_utc=datetime.datetime(2021, 7, 22, 22, 27, 41, 236000, tzinfo=datetime.timezone.utc), + expires_at_utc=datetime.datetime(2021, 8, 5, 22, 27, 41, 236000, tzinfo=datetime.timezone.utc), + sequence_number=11219, + lock_token='233146e3-d5a6-45eb-826f-691d82fb8b13' + ) + dict = CloudEvent._get_bytes(obj) + assert dict.get('data') == "ServiceBus" + assert dict.get('specversion') == '1.0' + +def test_get_bytes_servicebus_wrong_content(): + obj = MockServiceBusReceivedMessage( + body=MockBody(data="random string"), + message_id='3f6c5441-5be5-4f33-80c3-3ffeb6a090ce', + content_type='application/json; charset=utf-8', + time_to_live=datetime.timedelta(days=14), + delivery_count=13, + enqueued_sequence_number=0, + enqueued_time_utc=datetime.datetime(2021, 7, 22, 22, 27, 41, 236000, tzinfo=datetime.timezone.utc), + expires_at_utc=datetime.datetime(2021, 8, 5, 22, 27, 41, 236000, tzinfo=datetime.timezone.utc), + sequence_number=11219, + lock_token='233146e3-d5a6-45eb-826f-691d82fb8b13' + ) + + with pytest.raises(ValueError, match="Failed to retrieve body from the object. Make sure the body follows the CloudEvent schema."): + CloudEvent._get_bytes(obj) + +def test_get_bytes_eventhubs(): + obj = MockEventhubData( + body=MockEhBody() + ) + dict = CloudEvent._get_bytes(obj) + assert dict.get('data') == 'Eventhub' + assert dict.get('specversion') == '1.0' + +def test_get_bytes_eventhubs_wrong_content(): + obj = MockEventhubData( + body=MockEhBody(data='random string') + ) + + with pytest.raises(ValueError, match="Failed to retrieve body from the object. Make sure the body follows the CloudEvent schema."): + dict = CloudEvent._get_bytes(obj) + +def test_get_bytes_random_obj(): + random_obj = { + "id":"de0fd76c-4ef4-4dfb-ab3a-8f24a307e033", + "source":"https://egtest.dev/cloudcustomevent", + "data":{"team": "event grid squad"}, + "type":"Azure.Sdk.Sample", + "time":"2020-08-07T02:06:08.11969Z", + "specversion":"1.0", + "ext1": "example", + "BADext2": "example2" + } + + assert CloudEvent._get_bytes(random_obj) is random_obj diff --git a/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py b/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py index 53b204d334a9..1d2f58a08007 100644 --- a/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py +++ b/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py @@ -106,11 +106,16 @@ def from_dict(cls, data, key_extractors=None, content_type=None): @staticmethod def _get_bytes(obj): """Event mixin to have methods that are common to different Event types - like EventGridEvent etc. + like CloudEvent, EventGridEvent etc. """ try: # storage queue return json.loads(obj.content) + except ValueError: + raise ValueError( + "Failed to retrieve content from the object. Make sure the " + + "content follows the EventGridEvent schema." + ) except AttributeError: # eventhubs try: @@ -118,5 +123,10 @@ def _get_bytes(obj): except KeyError: # servicebus return json.loads(next(obj.body)) + except ValueError: + raise ValueError( + "Failed to retrieve body from the object. Make sure the " + + "body follows the EventGridEvent schema." + ) except: # pylint: disable=bare-except return obj diff --git a/sdk/eventgrid/azure-eventgrid/tests/test_eg_event_get_bytes.py b/sdk/eventgrid/azure-eventgrid/tests/test_eg_event_get_bytes.py new file mode 100644 index 000000000000..366e6bb01f0c --- /dev/null +++ b/sdk/eventgrid/azure-eventgrid/tests/test_eg_event_get_bytes.py @@ -0,0 +1,190 @@ +# ------------------------------------ +# Copyright (c) Microsoft Corporation. +# Licensed under the MIT License. +# ------------------------------------ +import datetime +import pytest +import uuid +from azure.eventgrid import EventGridEvent + +class MockQueueMessage(object): + def __init__(self, content=None): + self.id = uuid.uuid4() + self.inserted_on = datetime.datetime.now() + self.expires_on = datetime.datetime.now() + datetime.timedelta(days=100) + self.dequeue_count = 1 + self.content = content + self.pop_receipt = None + self.next_visible_on = None + +class MockServiceBusReceivedMessage(object): + def __init__(self, body=None, **kwargs): + self.body=body + self.application_properties=None + self.session_id=None + self.message_id='3f6c5441-5be5-4f33-80c3-3ffeb6a090ce' + self.content_type='application/cloudevents+json; charset=utf-8' + self.correlation_id=None + self.to=None + self.reply_to=None + self.reply_to_session_id=None + self.subject=None + self.time_to_live=datetime.timedelta(days=14) + self.partition_key=None + self.scheduled_enqueue_time_utc=None + self.auto_renew_error=None, + self.dead_letter_error_description=None + self.dead_letter_reason=None + self.dead_letter_source=None + self.delivery_count=13 + self.enqueued_sequence_number=0 + self.enqueued_time_utc=datetime.datetime(2021, 7, 22, 22, 27, 41, 236000, tzinfo=datetime.timezone.utc) + self.expires_at_utc=datetime.datetime(2021, 8, 5, 22, 27, 41, 236000, tzinfo=datetime.timezone.utc) + self.sequence_number=11219 + self.lock_token='233146e3-d5a6-45eb-826f-691d82fb8b13' + +class MockEventhubData(object): + def __init__(self, body=None): + self._last_enqueued_event_properties = {} + self._sys_properties = None + if body is None: + raise ValueError("EventData cannot be None.") + + # Internal usage only for transforming AmqpAnnotatedMessage to outgoing EventData + self.body=body + self._raw_amqp_message = "some amqp data" + self.message_id = None + self.content_type = None + self.correlation_id = None + + +class MockBody(object): + def __init__(self, data=None): + self.data = data + + def __iter__(self): + return self + + def __next__(self): + if not self.data: + return """{"id":"f208feff-099b-4bda-a341-4afd0fa02fef","subject":"https://egsample.dev/sampleevent","data":"ServiceBus","event_type":"Azure.Sdk.Sample","event_time":"2021-07-22T22:27:38.960209Z","data_version":"1.0"}""" + return self.data + +class MockEhBody(object): + def __init__(self, data=None): + self.data = data + + def __iter__(self): + return self + + def __next__(self): + if not self.data: + return b'[{"id":"f208feff-099b-4bda-a341-4afd0fa02fef","subject":"https://egsample.dev/sampleevent","data":"Eventhub","event_type":"Azure.Sdk.Sample","event_time":"2021-07-22T22:27:38.960209Z","data_version":"1.0"}]' + return self.data + +def test_get_bytes_storage_queue(): + cloud_storage_dict = """{ + "id":"a0517898-9fa4-4e70-b4a3-afda1dd68672", + "subject":"/subscriptions/{subscription-id}/resourceGroups/{resource-group}/providers/Microsoft.Storage/storageAccounts/{storage-account}", + "data":{ + "api":"PutBlockList", + "client_request_id":"6d79dbfb-0e37-4fc4-981f-442c9ca65760", + "request_id":"831e1650-001e-001b-66ab-eeb76e000000", + "e_tag":"0x8D4BCC2E4835CD0", + "content_type":"application/octet-stream", + "content_length":524288, + "blob_type":"BlockBlob", + "url":"https://oc2d2817345i60006.blob.core.windows.net/oc2d2817345i200097container/oc2d2817345i20002296blob", + "sequencer":"00000000000004420000000000028963", + "storage_diagnostics":{"batchId":"b68529f3-68cd-4744-baa4-3c0498ec19f0"} + }, + "event_type":"Microsoft.Storage.BlobCreated", + "event_time":"2021-02-18T20:18:10.581147898Z", + "data_version":"1.0" + }""" + obj = MockQueueMessage(content=cloud_storage_dict) + + dict = EventGridEvent._get_bytes(obj) + assert dict.get('data') == { + "api":"PutBlockList", + "client_request_id":"6d79dbfb-0e37-4fc4-981f-442c9ca65760", + "request_id":"831e1650-001e-001b-66ab-eeb76e000000", + "e_tag":"0x8D4BCC2E4835CD0", + "content_type":"application/octet-stream", + "content_length":524288, + "blob_type":"BlockBlob", + "url":"https://oc2d2817345i60006.blob.core.windows.net/oc2d2817345i200097container/oc2d2817345i20002296blob", + "sequencer":"00000000000004420000000000028963", + "storage_diagnostics":{"batchId":"b68529f3-68cd-4744-baa4-3c0498ec19f0"} + } + assert dict.get('data_version') == "1.0" + +def test_get_bytes_storage_queue_wrong_content(): + string = u'This is a random string which must fail' + obj = MockQueueMessage(content=string) + + with pytest.raises(ValueError, match="Failed to retrieve content from the object. Make sure the content follows the EventGridEvent schema."): + EventGridEvent.from_dict(obj) + +def test_get_bytes_servicebus(): + obj = MockServiceBusReceivedMessage( + body=MockBody(), + message_id='3f6c5441-5be5-4f33-80c3-3ffeb6a090ce', + content_type='application/cloudevents+json; charset=utf-8', + time_to_live=datetime.timedelta(days=14), + delivery_count=13, + enqueued_sequence_number=0, + enqueued_time_utc=datetime.datetime(2021, 7, 22, 22, 27, 41, 236000, tzinfo=datetime.timezone.utc), + expires_at_utc=datetime.datetime(2021, 8, 5, 22, 27, 41, 236000, tzinfo=datetime.timezone.utc), + sequence_number=11219, + lock_token='233146e3-d5a6-45eb-826f-691d82fb8b13' + ) + dict = EventGridEvent._get_bytes(obj) + assert dict.get('data') == "ServiceBus" + assert dict.get('data_version') == '1.0' + +def test_get_bytes_servicebus_wrong_content(): + obj = MockServiceBusReceivedMessage( + body=MockBody(data='random'), + message_id='3f6c5441-5be5-4f33-80c3-3ffeb6a090ce', + content_type='application/json; charset=utf-8', + time_to_live=datetime.timedelta(days=14), + delivery_count=13, + enqueued_sequence_number=0, + enqueued_time_utc=datetime.datetime(2021, 7, 22, 22, 27, 41, 236000, tzinfo=datetime.timezone.utc), + expires_at_utc=datetime.datetime(2021, 8, 5, 22, 27, 41, 236000, tzinfo=datetime.timezone.utc), + sequence_number=11219, + lock_token='233146e3-d5a6-45eb-826f-691d82fb8b13' + ) + with pytest.raises(ValueError, match="Failed to retrieve body from the object. Make sure the body follows the EventGridEvent schema."): + dict = EventGridEvent._get_bytes(obj) + + +def test_get_bytes_eventhubs(): + obj = MockEventhubData( + body=MockEhBody() + ) + dict = EventGridEvent._get_bytes(obj) + assert dict.get('data') == 'Eventhub' + assert dict.get('data_version') == '1.0' + +def test_get_bytes_eventhubs_wrong_content(): + obj = MockEventhubData( + body=MockEhBody(data='random string') + ) + + with pytest.raises(ValueError, match="Failed to retrieve body from the object. Make sure the body follows the EventGridEvent schema."): + dict = EventGridEvent._get_bytes(obj) + + +def test_get_bytes_random_obj(): + random_obj = { + "id":"de0fd76c-4ef4-4dfb-ab3a-8f24a307e033", + "subject":"https://egtest.dev/cloudcustomevent", + "data":{"team": "event grid squad"}, + "event_type":"Azure.Sdk.Sample", + "event_time":"2020-08-07T02:06:08.11969Z", + "data_version":"1.0", + } + + assert EventGridEvent._get_bytes(random_obj) is random_obj From b19d37089c0e03b739d254a85736ff9234e4a5ae Mon Sep 17 00:00:00 2001 From: Rakshith Bhyravabhotla Date: Thu, 29 Jul 2021 10:49:27 -0700 Subject: [PATCH 10/21] ci --- sdk/core/azure-core/azure/core/messaging.py | 2 +- .../azure-core/tests/test_messaging_cloud_event.py | 12 ++++++------ .../azure-eventgrid/azure/eventgrid/_models.py | 4 ++-- .../azure-eventgrid/tests/test_eg_event_get_bytes.py | 12 ++++++------ 4 files changed, 15 insertions(+), 15 deletions(-) diff --git a/sdk/core/azure-core/azure/core/messaging.py b/sdk/core/azure-core/azure/core/messaging.py index d3bf0faf4e35..5201333499c8 100644 --- a/sdk/core/azure-core/azure/core/messaging.py +++ b/sdk/core/azure-core/azure/core/messaging.py @@ -208,6 +208,6 @@ def _get_bytes(obj): raise ValueError( "Failed to retrieve body from the object. Make sure the " + "body follows the CloudEvent schema." - ) + ) except: # pylint: disable=bare-except return obj diff --git a/sdk/core/azure-core/tests/test_messaging_cloud_event.py b/sdk/core/azure-core/tests/test_messaging_cloud_event.py index c1042b02ea4f..1adc34b5b073 100644 --- a/sdk/core/azure-core/tests/test_messaging_cloud_event.py +++ b/sdk/core/azure-core/tests/test_messaging_cloud_event.py @@ -42,8 +42,8 @@ def __init__(self, body=None, **kwargs): self.dead_letter_source=None self.delivery_count=13 self.enqueued_sequence_number=0 - self.enqueued_time_utc=datetime.datetime(2021, 7, 22, 22, 27, 41, 236000, tzinfo=datetime.timezone.utc) - self.expires_at_utc=datetime.datetime(2021, 8, 5, 22, 27, 41, 236000, tzinfo=datetime.timezone.utc) + self.enqueued_time_utc=datetime.datetime(2021, 7, 22, 22, 27, 41, 236000) + self.expires_at_utc=datetime.datetime(2021, 8, 5, 22, 27, 41, 236000) self.sequence_number=11219 self.lock_token='233146e3-d5a6-45eb-826f-691d82fb8b13' @@ -600,8 +600,8 @@ def test_get_bytes_servicebus(): time_to_live=datetime.timedelta(days=14), delivery_count=13, enqueued_sequence_number=0, - enqueued_time_utc=datetime.datetime(2021, 7, 22, 22, 27, 41, 236000, tzinfo=datetime.timezone.utc), - expires_at_utc=datetime.datetime(2021, 8, 5, 22, 27, 41, 236000, tzinfo=datetime.timezone.utc), + enqueued_time_utc=datetime.datetime(2021, 7, 22, 22, 27, 41, 236000), + expires_at_utc=datetime.datetime(2021, 8, 5, 22, 27, 41, 236000), sequence_number=11219, lock_token='233146e3-d5a6-45eb-826f-691d82fb8b13' ) @@ -617,8 +617,8 @@ def test_get_bytes_servicebus_wrong_content(): time_to_live=datetime.timedelta(days=14), delivery_count=13, enqueued_sequence_number=0, - enqueued_time_utc=datetime.datetime(2021, 7, 22, 22, 27, 41, 236000, tzinfo=datetime.timezone.utc), - expires_at_utc=datetime.datetime(2021, 8, 5, 22, 27, 41, 236000, tzinfo=datetime.timezone.utc), + enqueued_time_utc=datetime.datetime(2021, 7, 22, 22, 27, 41, 236000), + expires_at_utc=datetime.datetime(2021, 8, 5, 22, 27, 41, 236000), sequence_number=11219, lock_token='233146e3-d5a6-45eb-826f-691d82fb8b13' ) diff --git a/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py b/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py index 1d2f58a08007..c31c15b38fe5 100644 --- a/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py +++ b/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py @@ -3,7 +3,7 @@ # Licensed under the MIT License. See License.txt in the project root for license information. # -------------------------------------------------------------------------------------------- # pylint:disable=protected-access -from typing import Any, Dict +from typing import Any import datetime as dt import uuid import json @@ -127,6 +127,6 @@ def _get_bytes(obj): raise ValueError( "Failed to retrieve body from the object. Make sure the " + "body follows the EventGridEvent schema." - ) + ) except: # pylint: disable=bare-except return obj diff --git a/sdk/eventgrid/azure-eventgrid/tests/test_eg_event_get_bytes.py b/sdk/eventgrid/azure-eventgrid/tests/test_eg_event_get_bytes.py index 366e6bb01f0c..06eb44271233 100644 --- a/sdk/eventgrid/azure-eventgrid/tests/test_eg_event_get_bytes.py +++ b/sdk/eventgrid/azure-eventgrid/tests/test_eg_event_get_bytes.py @@ -38,8 +38,8 @@ def __init__(self, body=None, **kwargs): self.dead_letter_source=None self.delivery_count=13 self.enqueued_sequence_number=0 - self.enqueued_time_utc=datetime.datetime(2021, 7, 22, 22, 27, 41, 236000, tzinfo=datetime.timezone.utc) - self.expires_at_utc=datetime.datetime(2021, 8, 5, 22, 27, 41, 236000, tzinfo=datetime.timezone.utc) + self.enqueued_time_utc=datetime.datetime(2021, 7, 22, 22, 27, 41, 236000) + self.expires_at_utc=datetime.datetime(2021, 8, 5, 22, 27, 41, 236000) self.sequence_number=11219 self.lock_token='233146e3-d5a6-45eb-826f-691d82fb8b13' @@ -134,8 +134,8 @@ def test_get_bytes_servicebus(): time_to_live=datetime.timedelta(days=14), delivery_count=13, enqueued_sequence_number=0, - enqueued_time_utc=datetime.datetime(2021, 7, 22, 22, 27, 41, 236000, tzinfo=datetime.timezone.utc), - expires_at_utc=datetime.datetime(2021, 8, 5, 22, 27, 41, 236000, tzinfo=datetime.timezone.utc), + enqueued_time_utc=datetime.datetime(2021, 7, 22, 22, 27, 41, 236000), + expires_at_utc=datetime.datetime(2021, 8, 5, 22, 27, 41, 236000), sequence_number=11219, lock_token='233146e3-d5a6-45eb-826f-691d82fb8b13' ) @@ -151,8 +151,8 @@ def test_get_bytes_servicebus_wrong_content(): time_to_live=datetime.timedelta(days=14), delivery_count=13, enqueued_sequence_number=0, - enqueued_time_utc=datetime.datetime(2021, 7, 22, 22, 27, 41, 236000, tzinfo=datetime.timezone.utc), - expires_at_utc=datetime.datetime(2021, 8, 5, 22, 27, 41, 236000, tzinfo=datetime.timezone.utc), + enqueued_time_utc=datetime.datetime(2021, 7, 22, 22, 27, 41, 236000), + expires_at_utc=datetime.datetime(2021, 8, 5, 22, 27, 41, 236000), sequence_number=11219, lock_token='233146e3-d5a6-45eb-826f-691d82fb8b13' ) From 2376e4f0441ddff1befc07f6c660c315a11f8510 Mon Sep 17 00:00:00 2001 From: Rakshith Bhyravabhotla Date: Fri, 30 Jul 2021 11:18:26 -0700 Subject: [PATCH 11/21] lint --- sdk/core/azure-core/azure/core/messaging.py | 2 +- sdk/core/azure-core/tests/test_messaging_cloud_event.py | 4 ++++ .../azure-eventgrid/tests/test_eg_event_get_bytes.py | 5 +++++ 3 files changed, 10 insertions(+), 1 deletion(-) diff --git a/sdk/core/azure-core/azure/core/messaging.py b/sdk/core/azure-core/azure/core/messaging.py index 5201333499c8..d41cae166f31 100644 --- a/sdk/core/azure-core/azure/core/messaging.py +++ b/sdk/core/azure-core/azure/core/messaging.py @@ -208,6 +208,6 @@ def _get_bytes(obj): raise ValueError( "Failed to retrieve body from the object. Make sure the " + "body follows the CloudEvent schema." - ) + ) except: # pylint: disable=bare-except return obj diff --git a/sdk/core/azure-core/tests/test_messaging_cloud_event.py b/sdk/core/azure-core/tests/test_messaging_cloud_event.py index 1adc34b5b073..c871eaaceb06 100644 --- a/sdk/core/azure-core/tests/test_messaging_cloud_event.py +++ b/sdk/core/azure-core/tests/test_messaging_cloud_event.py @@ -73,6 +73,8 @@ def __next__(self): if not self.data: return """{"id":"f208feff-099b-4bda-a341-4afd0fa02fef","source":"https://egsample.dev/sampleevent","data":"ServiceBus","type":"Azure.Sdk.Sample","time":"2021-07-22T22:27:38.960209Z","specversion":"1.0"}""" return self.data + + next = __next__ class MockEhBody(object): def __init__(self, data=None): @@ -86,6 +88,8 @@ def __next__(self): return b'[{"id":"f208feff-099b-4bda-a341-4afd0fa02fef","source":"https://egsample.dev/sampleevent","data":"Eventhub","type":"Azure.Sdk.Sample","time":"2021-07-22T22:27:38.960209Z","specversion":"1.0"}]' return self.data + next = __next__ + # Cloud Event tests def test_cloud_event_constructor(): diff --git a/sdk/eventgrid/azure-eventgrid/tests/test_eg_event_get_bytes.py b/sdk/eventgrid/azure-eventgrid/tests/test_eg_event_get_bytes.py index 06eb44271233..3d993516083b 100644 --- a/sdk/eventgrid/azure-eventgrid/tests/test_eg_event_get_bytes.py +++ b/sdk/eventgrid/azure-eventgrid/tests/test_eg_event_get_bytes.py @@ -70,6 +70,9 @@ def __next__(self): return """{"id":"f208feff-099b-4bda-a341-4afd0fa02fef","subject":"https://egsample.dev/sampleevent","data":"ServiceBus","event_type":"Azure.Sdk.Sample","event_time":"2021-07-22T22:27:38.960209Z","data_version":"1.0"}""" return self.data + next = __next__ + + class MockEhBody(object): def __init__(self, data=None): self.data = data @@ -81,6 +84,8 @@ def __next__(self): if not self.data: return b'[{"id":"f208feff-099b-4bda-a341-4afd0fa02fef","subject":"https://egsample.dev/sampleevent","data":"Eventhub","event_type":"Azure.Sdk.Sample","event_time":"2021-07-22T22:27:38.960209Z","data_version":"1.0"}]' return self.data + + next = __next__ def test_get_bytes_storage_queue(): cloud_storage_dict = """{ From d0fa5538270d8b6c13a4f7718f92b710fbd6ff61 Mon Sep 17 00:00:00 2001 From: Rakshith Bhyravabhotla Date: Sun, 1 Aug 2021 17:45:11 -0700 Subject: [PATCH 12/21] rename + static method --- sdk/core/azure-core/azure/core/messaging.py | 32 ++----------------- .../azure-core/azure/core/utils/_utils.py | 28 ++++++++++++++++ .../tests/test_messaging_cloud_event.py | 14 ++++---- .../azure/eventgrid/_helpers.py | 27 ++++++++++++++++ .../azure/eventgrid/_models.py | 3 +- .../tests/test_eg_event_get_bytes.py | 13 ++++---- 6 files changed, 73 insertions(+), 44 deletions(-) diff --git a/sdk/core/azure-core/azure/core/messaging.py b/sdk/core/azure-core/azure/core/messaging.py index d41cae166f31..d3ec35b4b55e 100644 --- a/sdk/core/azure-core/azure/core/messaging.py +++ b/sdk/core/azure-core/azure/core/messaging.py @@ -8,7 +8,7 @@ import json from base64 import b64decode from datetime import datetime -from .utils._utils import _convert_to_isoformat, TZ_UTC +from .utils._utils import _convert_to_isoformat, TZ_UTC, _get_json_content from .serialization import NULL try: @@ -123,7 +123,7 @@ def from_dict(cls, event): :type event: dict :rtype: CloudEvent """ - event = CloudEvent._get_bytes(event) + event = _get_json_content(event) kwargs = {} # type: Dict[Any, Any] reserved_attr = [ "data", @@ -183,31 +183,3 @@ def from_dict(cls, event): " The `source` and `type` params are required." ) return event_obj - - @staticmethod - def _get_bytes(obj): - """Event mixin to have methods that are common to different Event types - like CloudEvent, EventGridEvent etc. - """ - try: - # storage queue - return json.loads(obj.content) - except ValueError: - raise ValueError( - "Failed to retrieve content from the object. Make sure the " - + "content follows the CloudEvent schema." - ) - except AttributeError: - # eventhubs - try: - return json.loads(next(obj.body))[0] - except KeyError: - # servicebus - return json.loads(next(obj.body)) - except ValueError: - raise ValueError( - "Failed to retrieve body from the object. Make sure the " - + "body follows the CloudEvent schema." - ) - except: # pylint: disable=bare-except - return obj diff --git a/sdk/core/azure-core/azure/core/utils/_utils.py b/sdk/core/azure-core/azure/core/utils/_utils.py index 5f3133e02955..bd656d3f3209 100644 --- a/sdk/core/azure-core/azure/core/utils/_utils.py +++ b/sdk/core/azure-core/azure/core/utils/_utils.py @@ -5,6 +5,7 @@ # license information. # -------------------------------------------------------------------------- import datetime +import json class _FixedOffset(datetime.tzinfo): @@ -104,3 +105,30 @@ def _case_insensitive_dict(*args, **kwargs): raise ValueError( "Neither 'requests' or 'multidict' are installed and no case-insensitive dict impl have been found" ) + +def _get_json_content(obj): + """Event mixin to have methods that are common to different Event types + like CloudEvent, EventGridEvent etc. + """ + try: + # storage queue + return json.loads(obj.content) + except ValueError: + raise ValueError( + "Failed to retrieve content from the object. Make sure the " + + "content follows the CloudEvent schema." + ) + except AttributeError: + # eventhubs + try: + return json.loads(next(obj.body))[0] + except KeyError: + # servicebus + return json.loads(next(obj.body)) + except ValueError: + raise ValueError( + "Failed to retrieve body from the object. Make sure the " + + "body follows the CloudEvent schema." + ) + except: # pylint: disable=bare-except + return obj diff --git a/sdk/core/azure-core/tests/test_messaging_cloud_event.py b/sdk/core/azure-core/tests/test_messaging_cloud_event.py index c871eaaceb06..e9a4a94cfaf2 100644 --- a/sdk/core/azure-core/tests/test_messaging_cloud_event.py +++ b/sdk/core/azure-core/tests/test_messaging_cloud_event.py @@ -8,7 +8,7 @@ import datetime from azure.core.messaging import CloudEvent -from azure.core.utils._utils import _convert_to_isoformat +from azure.core.utils._utils import _convert_to_isoformat, _get_json_content from azure.core.serialization import NULL class MockQueueMessage(object): @@ -574,7 +574,7 @@ def test_get_bytes_storage_queue(): }""" obj = MockQueueMessage(content=cloud_storage_dict) - dict = CloudEvent._get_bytes(obj) + dict = _get_json_content(obj) assert dict.get('data') == { "api":"PutBlockList", "client_request_id":"6d79dbfb-0e37-4fc4-981f-442c9ca65760", @@ -609,7 +609,7 @@ def test_get_bytes_servicebus(): sequence_number=11219, lock_token='233146e3-d5a6-45eb-826f-691d82fb8b13' ) - dict = CloudEvent._get_bytes(obj) + dict = _get_json_content(obj) assert dict.get('data') == "ServiceBus" assert dict.get('specversion') == '1.0' @@ -628,13 +628,13 @@ def test_get_bytes_servicebus_wrong_content(): ) with pytest.raises(ValueError, match="Failed to retrieve body from the object. Make sure the body follows the CloudEvent schema."): - CloudEvent._get_bytes(obj) + _get_json_content(obj) def test_get_bytes_eventhubs(): obj = MockEventhubData( body=MockEhBody() ) - dict = CloudEvent._get_bytes(obj) + dict = _get_json_content(obj) assert dict.get('data') == 'Eventhub' assert dict.get('specversion') == '1.0' @@ -644,7 +644,7 @@ def test_get_bytes_eventhubs_wrong_content(): ) with pytest.raises(ValueError, match="Failed to retrieve body from the object. Make sure the body follows the CloudEvent schema."): - dict = CloudEvent._get_bytes(obj) + dict = _get_json_content(obj) def test_get_bytes_random_obj(): random_obj = { @@ -658,4 +658,4 @@ def test_get_bytes_random_obj(): "BADext2": "example2" } - assert CloudEvent._get_bytes(random_obj) is random_obj + assert _get_json_content(random_obj) is random_obj diff --git a/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_helpers.py b/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_helpers.py index 8c51aafffafc..923301bebd6e 100644 --- a/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_helpers.py +++ b/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_helpers.py @@ -165,3 +165,30 @@ def _build_request(endpoint, content_type, events): ) request.format_parameters(query_parameters) return request + +def _get_json_content(obj): + """Event mixin to have methods that are common to different Event types + like CloudEvent, EventGridEvent etc. + """ + try: + # storage queue + return json.loads(obj.content) + except ValueError: + raise ValueError( + "Failed to retrieve content from the object. Make sure the " + + "content follows the EventGridEvent schema." + ) + except AttributeError: + # eventhubs + try: + return json.loads(next(obj.body))[0] + except KeyError: + # servicebus + return json.loads(next(obj.body)) + except ValueError: + raise ValueError( + "Failed to retrieve body from the object. Make sure the " + + "body follows the EventGridEvent schema." + ) + except: # pylint: disable=bare-except + return obj \ No newline at end of file diff --git a/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py b/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py index c31c15b38fe5..6c803d0d1b0c 100644 --- a/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py +++ b/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py @@ -8,6 +8,7 @@ import uuid import json from msrest.serialization import UTC +from ._helpers import _get_json_content from ._generated.models import ( EventGridEvent as InternalEventGridEvent, ) @@ -100,7 +101,7 @@ def __repr__(self): @classmethod def from_dict(cls, data, key_extractors=None, content_type=None): - event = EventGridEvent._get_bytes(data) + event = _get_json_content(data) super(EventGridEvent, cls).from_dict(event, key_extractors, content_type) @staticmethod diff --git a/sdk/eventgrid/azure-eventgrid/tests/test_eg_event_get_bytes.py b/sdk/eventgrid/azure-eventgrid/tests/test_eg_event_get_bytes.py index 3d993516083b..5a04a0acf816 100644 --- a/sdk/eventgrid/azure-eventgrid/tests/test_eg_event_get_bytes.py +++ b/sdk/eventgrid/azure-eventgrid/tests/test_eg_event_get_bytes.py @@ -5,6 +5,7 @@ import datetime import pytest import uuid +from azure.eventgrid._helpers import _get_json_content from azure.eventgrid import EventGridEvent class MockQueueMessage(object): @@ -109,7 +110,7 @@ def test_get_bytes_storage_queue(): }""" obj = MockQueueMessage(content=cloud_storage_dict) - dict = EventGridEvent._get_bytes(obj) + dict = _get_json_content(obj) assert dict.get('data') == { "api":"PutBlockList", "client_request_id":"6d79dbfb-0e37-4fc4-981f-442c9ca65760", @@ -144,7 +145,7 @@ def test_get_bytes_servicebus(): sequence_number=11219, lock_token='233146e3-d5a6-45eb-826f-691d82fb8b13' ) - dict = EventGridEvent._get_bytes(obj) + dict = _get_json_content(obj) assert dict.get('data') == "ServiceBus" assert dict.get('data_version') == '1.0' @@ -162,14 +163,14 @@ def test_get_bytes_servicebus_wrong_content(): lock_token='233146e3-d5a6-45eb-826f-691d82fb8b13' ) with pytest.raises(ValueError, match="Failed to retrieve body from the object. Make sure the body follows the EventGridEvent schema."): - dict = EventGridEvent._get_bytes(obj) + dict = _get_json_content(obj) def test_get_bytes_eventhubs(): obj = MockEventhubData( body=MockEhBody() ) - dict = EventGridEvent._get_bytes(obj) + dict = _get_json_content(obj) assert dict.get('data') == 'Eventhub' assert dict.get('data_version') == '1.0' @@ -179,7 +180,7 @@ def test_get_bytes_eventhubs_wrong_content(): ) with pytest.raises(ValueError, match="Failed to retrieve body from the object. Make sure the body follows the EventGridEvent schema."): - dict = EventGridEvent._get_bytes(obj) + dict = _get_json_content(obj) def test_get_bytes_random_obj(): @@ -192,4 +193,4 @@ def test_get_bytes_random_obj(): "data_version":"1.0", } - assert EventGridEvent._get_bytes(random_obj) is random_obj + assert _get_json_content(random_obj) is random_obj From d2dbd4a9f105bbc6e5c24a37eb82f66cb5f34add Mon Sep 17 00:00:00 2001 From: Rakshith Bhyravabhotla Date: Mon, 2 Aug 2021 10:02:02 -0700 Subject: [PATCH 13/21] lint --- sdk/core/azure-core/azure/core/messaging.py | 1 - sdk/eventgrid/azure-eventgrid/azure/eventgrid/_helpers.py | 2 +- 2 files changed, 1 insertion(+), 2 deletions(-) diff --git a/sdk/core/azure-core/azure/core/messaging.py b/sdk/core/azure-core/azure/core/messaging.py index d3ec35b4b55e..76672ce4d1b3 100644 --- a/sdk/core/azure-core/azure/core/messaging.py +++ b/sdk/core/azure-core/azure/core/messaging.py @@ -5,7 +5,6 @@ # license information. # -------------------------------------------------------------------------- import uuid -import json from base64 import b64decode from datetime import datetime from .utils._utils import _convert_to_isoformat, TZ_UTC, _get_json_content diff --git a/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_helpers.py b/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_helpers.py index 923301bebd6e..2be04b3be457 100644 --- a/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_helpers.py +++ b/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_helpers.py @@ -191,4 +191,4 @@ def _get_json_content(obj): + "body follows the EventGridEvent schema." ) except: # pylint: disable=bare-except - return obj \ No newline at end of file + return obj From 33163dfdaecf82e7bd1130987f5d5f0907c9c9b3 Mon Sep 17 00:00:00 2001 From: Rakshith Bhyravabhotla Date: Mon, 2 Aug 2021 13:05:26 -0700 Subject: [PATCH 14/21] remove method --- .../azure/eventgrid/_models.py | 28 ------------------- 1 file changed, 28 deletions(-) diff --git a/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py b/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py index 6c803d0d1b0c..7bc18b962cb7 100644 --- a/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py +++ b/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py @@ -103,31 +103,3 @@ def __repr__(self): def from_dict(cls, data, key_extractors=None, content_type=None): event = _get_json_content(data) super(EventGridEvent, cls).from_dict(event, key_extractors, content_type) - - @staticmethod - def _get_bytes(obj): - """Event mixin to have methods that are common to different Event types - like CloudEvent, EventGridEvent etc. - """ - try: - # storage queue - return json.loads(obj.content) - except ValueError: - raise ValueError( - "Failed to retrieve content from the object. Make sure the " - + "content follows the EventGridEvent schema." - ) - except AttributeError: - # eventhubs - try: - return json.loads(next(obj.body))[0] - except KeyError: - # servicebus - return json.loads(next(obj.body)) - except ValueError: - raise ValueError( - "Failed to retrieve body from the object. Make sure the " - + "body follows the EventGridEvent schema." - ) - except: # pylint: disable=bare-except - return obj From 70a0168d6edd08abf8c6c3498a913fd0723549ea Mon Sep 17 00:00:00 2001 From: Rakshith Bhyravabhotla Date: Mon, 2 Aug 2021 14:00:02 -0700 Subject: [PATCH 15/21] Update sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py --- sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py | 1 - 1 file changed, 1 deletion(-) diff --git a/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py b/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py index 7bc18b962cb7..b5d82424f8bc 100644 --- a/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py +++ b/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py @@ -6,7 +6,6 @@ from typing import Any import datetime as dt import uuid -import json from msrest.serialization import UTC from ._helpers import _get_json_content from ._generated.models import ( From 40ca6d69dc59b4637be6bead872ece0abd55d965 Mon Sep 17 00:00:00 2001 From: Rakshith Bhyravabhotla Date: Mon, 2 Aug 2021 19:33:41 -0700 Subject: [PATCH 16/21] add a from_json method --- sdk/core/azure-core/CHANGELOG.md | 2 +- sdk/core/azure-core/azure/core/messaging.py | 14 +++- .../azure-core/azure/core/utils/_utils.py | 2 +- .../tests/test_messaging_cloud_event.py | 77 +++++++++++++++++-- sdk/eventgrid/azure-eventgrid/CHANGELOG.md | 2 + .../azure/eventgrid/_helpers.py | 2 +- .../azure/eventgrid/_models.py | 6 +- .../tests/test_eg_event_get_bytes.py | 72 ++++++++++++++++- 8 files changed, 163 insertions(+), 14 deletions(-) diff --git a/sdk/core/azure-core/CHANGELOG.md b/sdk/core/azure-core/CHANGELOG.md index 7b1f836a3196..8ee413fc3350 100644 --- a/sdk/core/azure-core/CHANGELOG.md +++ b/sdk/core/azure-core/CHANGELOG.md @@ -5,7 +5,7 @@ ### Features Added - Cut hard dependency on requests library -- `CloudEvent`'s `from_dict` method now accepts objects from servicebus, eventhubs and storage directly. +- Added a `from_json` method which now accepts storage QueueMessage, eventhub's EventData or ServiceBusMessage or simply json bytes to return a `CloudEvent` ### Breaking Changes diff --git a/sdk/core/azure-core/azure/core/messaging.py b/sdk/core/azure-core/azure/core/messaging.py index 76672ce4d1b3..9bc7684b6929 100644 --- a/sdk/core/azure-core/azure/core/messaging.py +++ b/sdk/core/azure-core/azure/core/messaging.py @@ -122,7 +122,6 @@ def from_dict(cls, event): :type event: dict :rtype: CloudEvent """ - event = _get_json_content(event) kwargs = {} # type: Dict[Any, Any] reserved_attr = [ "data", @@ -182,3 +181,16 @@ def from_dict(cls, event): " The `source` and `type` params are required." ) return event_obj + + @classmethod + def from_json(cls, json): + # type: (Any) -> CloudEvent + """ + Returns the deserialized CloudEvent object when a json is provided. + :param json: The json string that should be converted into a CloudEvent. This can also be + a storage QueueMessage, eventhub's EventData or ServiceBusMessage + :type json: object + :rtype: CloudEvent + """ + event = _get_json_content(json) + return CloudEvent.from_dict(event) diff --git a/sdk/core/azure-core/azure/core/utils/_utils.py b/sdk/core/azure-core/azure/core/utils/_utils.py index 1b2375508b92..d8ae29080b7b 100644 --- a/sdk/core/azure-core/azure/core/utils/_utils.py +++ b/sdk/core/azure-core/azure/core/utils/_utils.py @@ -133,4 +133,4 @@ def _get_json_content(obj): + "body follows the CloudEvent schema." ) except: # pylint: disable=bare-except - return obj + return json.loads(obj) diff --git a/sdk/core/azure-core/tests/test_messaging_cloud_event.py b/sdk/core/azure-core/tests/test_messaging_cloud_event.py index e9a4a94cfaf2..98b064776957 100644 --- a/sdk/core/azure-core/tests/test_messaging_cloud_event.py +++ b/sdk/core/azure-core/tests/test_messaging_cloud_event.py @@ -594,7 +594,7 @@ def test_get_bytes_storage_queue_wrong_content(): obj = MockQueueMessage(content=cloud_storage_string) with pytest.raises(ValueError, match="Failed to retrieve content from the object. Make sure the content follows the CloudEvent schema."): - CloudEvent.from_dict(obj) + _get_json_content(obj) def test_get_bytes_servicebus(): obj = MockServiceBusReceivedMessage( @@ -647,15 +647,82 @@ def test_get_bytes_eventhubs_wrong_content(): dict = _get_json_content(obj) def test_get_bytes_random_obj(): + json_str = '{"id": "de0fd76c-4ef4-4dfb-ab3a-8f24a307e033", "source": "https://egtest.dev/cloudcustomevent", "data": {"team": "event grid squad"}, "type": "Azure.Sdk.Sample", "time": "2020-08-07T02:06:08.11969Z", "specversion": "1.0"}' random_obj = { "id":"de0fd76c-4ef4-4dfb-ab3a-8f24a307e033", "source":"https://egtest.dev/cloudcustomevent", "data":{"team": "event grid squad"}, "type":"Azure.Sdk.Sample", "time":"2020-08-07T02:06:08.11969Z", - "specversion":"1.0", - "ext1": "example", - "BADext2": "example2" + "specversion":"1.0" } - assert _get_json_content(random_obj) is random_obj + assert _get_json_content(json_str) == random_obj + +def test_from_json_sb(): + obj = MockServiceBusReceivedMessage( + body=MockBody(), + message_id='3f6c5441-5be5-4f33-80c3-3ffeb6a090ce', + content_type='application/cloudevents+json; charset=utf-8', + time_to_live=datetime.timedelta(days=14), + delivery_count=13, + enqueued_sequence_number=0, + enqueued_time_utc=datetime.datetime(2021, 7, 22, 22, 27, 41, 236000), + expires_at_utc=datetime.datetime(2021, 8, 5, 22, 27, 41, 236000), + sequence_number=11219, + lock_token='233146e3-d5a6-45eb-826f-691d82fb8b13' + ) + event = CloudEvent.from_json(obj) + + assert event.id == "f208feff-099b-4bda-a341-4afd0fa02fef" + assert event.data == "ServiceBus" + +def test_from_json_eh(): + obj = MockEventhubData( + body=MockEhBody() + ) + event = CloudEvent.from_json(obj) + assert event.id == "f208feff-099b-4bda-a341-4afd0fa02fef" + assert event.data == "Eventhub" + +def test_from_json_storage(): + cloud_storage_dict = """{ + "id":"a0517898-9fa4-4e70-b4a3-afda1dd68672", + "source":"/subscriptions/{subscription-id}/resourceGroups/{resource-group}/providers/Microsoft.Storage/storageAccounts/{storage-account}", + "data":{ + "api":"PutBlockList", + "client_request_id":"6d79dbfb-0e37-4fc4-981f-442c9ca65760", + "request_id":"831e1650-001e-001b-66ab-eeb76e000000", + "e_tag":"0x8D4BCC2E4835CD0", + "content_type":"application/octet-stream", + "content_length":524288, + "blob_type":"BlockBlob", + "url":"https://oc2d2817345i60006.blob.core.windows.net/oc2d2817345i200097container/oc2d2817345i20002296blob", + "sequencer":"00000000000004420000000000028963", + "storage_diagnostics":{"batchId":"b68529f3-68cd-4744-baa4-3c0498ec19f0"} + }, + "type":"Microsoft.Storage.BlobCreated", + "time":"2021-02-18T20:18:10.581147898Z", + "specversion":"1.0" + }""" + obj = MockQueueMessage(content=cloud_storage_dict) + event = CloudEvent.from_json(obj) + assert event.data == { + "api":"PutBlockList", + "client_request_id":"6d79dbfb-0e37-4fc4-981f-442c9ca65760", + "request_id":"831e1650-001e-001b-66ab-eeb76e000000", + "e_tag":"0x8D4BCC2E4835CD0", + "content_type":"application/octet-stream", + "content_length":524288, + "blob_type":"BlockBlob", + "url":"https://oc2d2817345i60006.blob.core.windows.net/oc2d2817345i200097container/oc2d2817345i20002296blob", + "sequencer":"00000000000004420000000000028963", + "storage_diagnostics":{"batchId":"b68529f3-68cd-4744-baa4-3c0498ec19f0"} + } + + +def test_from_json(): + json_str = '{"id": "de0fd76c-4ef4-4dfb-ab3a-8f24a307e033", "source": "https://egtest.dev/cloudcustomevent", "data": {"team": "event grid squad"}, "type": "Azure.Sdk.Sample", "time": "2020-08-07T02:06:08.11969Z", "specversion": "1.0"}' + event = CloudEvent.from_json(json_str) + + assert event.data == {"team": "event grid squad"} diff --git a/sdk/eventgrid/azure-eventgrid/CHANGELOG.md b/sdk/eventgrid/azure-eventgrid/CHANGELOG.md index 04b08d54d94c..81753dfc36ba 100644 --- a/sdk/eventgrid/azure-eventgrid/CHANGELOG.md +++ b/sdk/eventgrid/azure-eventgrid/CHANGELOG.md @@ -5,6 +5,8 @@ ### Features Added - `EventGridEvent`'s `from_dict` method now accepts objects from servicebus, eventhubs and storage directly. +- Added a `from_json` method which now accepts storage QueueMessage, eventhub's EventData or ServiceBusMessage or simply json bytes to return an `EventGridEvent` + ## 4.4.0 (2021-07-19) diff --git a/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_helpers.py b/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_helpers.py index 2be04b3be457..61a2994f7583 100644 --- a/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_helpers.py +++ b/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_helpers.py @@ -191,4 +191,4 @@ def _get_json_content(obj): + "body follows the EventGridEvent schema." ) except: # pylint: disable=bare-except - return obj + return json.loads(obj) diff --git a/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py b/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py index b5d82424f8bc..500ad7fb06a0 100644 --- a/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py +++ b/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py @@ -99,6 +99,6 @@ def __repr__(self): )[:1024] @classmethod - def from_dict(cls, data, key_extractors=None, content_type=None): - event = _get_json_content(data) - super(EventGridEvent, cls).from_dict(event, key_extractors, content_type) + def from_json(cls, json): + event = _get_json_content(json) + return EventGridEvent.from_dict(event) diff --git a/sdk/eventgrid/azure-eventgrid/tests/test_eg_event_get_bytes.py b/sdk/eventgrid/azure-eventgrid/tests/test_eg_event_get_bytes.py index 5a04a0acf816..03d350598213 100644 --- a/sdk/eventgrid/azure-eventgrid/tests/test_eg_event_get_bytes.py +++ b/sdk/eventgrid/azure-eventgrid/tests/test_eg_event_get_bytes.py @@ -130,7 +130,7 @@ def test_get_bytes_storage_queue_wrong_content(): obj = MockQueueMessage(content=string) with pytest.raises(ValueError, match="Failed to retrieve content from the object. Make sure the content follows the EventGridEvent schema."): - EventGridEvent.from_dict(obj) + _get_json_content(obj) def test_get_bytes_servicebus(): obj = MockServiceBusReceivedMessage( @@ -184,6 +184,7 @@ def test_get_bytes_eventhubs_wrong_content(): def test_get_bytes_random_obj(): + json_str = '{"id": "de0fd76c-4ef4-4dfb-ab3a-8f24a307e033", "subject": "https://egtest.dev/cloudcustomevent", "data": {"team": "event grid squad"}, "event_type": "Azure.Sdk.Sample", "event_time": "2020-08-07T02:06:08.11969Z", "data_version": "1.0"}' random_obj = { "id":"de0fd76c-4ef4-4dfb-ab3a-8f24a307e033", "subject":"https://egtest.dev/cloudcustomevent", @@ -193,4 +194,71 @@ def test_get_bytes_random_obj(): "data_version":"1.0", } - assert _get_json_content(random_obj) is random_obj + assert _get_json_content(json_str) == random_obj + +def test_from_json_sb(): + obj = MockServiceBusReceivedMessage( + body=MockBody(), + message_id='3f6c5441-5be5-4f33-80c3-3ffeb6a090ce', + content_type='application/cloudevents+json; charset=utf-8', + time_to_live=datetime.timedelta(days=14), + delivery_count=13, + enqueued_sequence_number=0, + enqueued_time_utc=datetime.datetime(2021, 7, 22, 22, 27, 41, 236000), + expires_at_utc=datetime.datetime(2021, 8, 5, 22, 27, 41, 236000), + sequence_number=11219, + lock_token='233146e3-d5a6-45eb-826f-691d82fb8b13' + ) + event = EventGridEvent.from_json(obj) + + assert event.id == "f208feff-099b-4bda-a341-4afd0fa02fef" + assert event.data == "ServiceBus" + +def test_from_json_eh(): + obj = MockEventhubData( + body=MockEhBody() + ) + event = EventGridEvent.from_json(obj) + assert event.id == "f208feff-099b-4bda-a341-4afd0fa02fef" + assert event.data == "Eventhub" + +def test_from_json_storage(): + eg_storage_dict = """{ + "id":"a0517898-9fa4-4e70-b4a3-afda1dd68672", + "subject":"/subscriptions/{subscription-id}/resourceGroups/{resource-group}/providers/Microsoft.Storage/storageAccounts/{storage-account}", + "data":{ + "api":"PutBlockList", + "client_request_id":"6d79dbfb-0e37-4fc4-981f-442c9ca65760", + "request_id":"831e1650-001e-001b-66ab-eeb76e000000", + "e_tag":"0x8D4BCC2E4835CD0", + "content_type":"application/octet-stream", + "content_length":524288, + "blob_type":"BlockBlob", + "url":"https://oc2d2817345i60006.blob.core.windows.net/oc2d2817345i200097container/oc2d2817345i20002296blob", + "sequencer":"00000000000004420000000000028963", + "storage_diagnostics":{"batchId":"b68529f3-68cd-4744-baa4-3c0498ec19f0"} + }, + "event_type":"Microsoft.Storage.BlobCreated", + "event_time":"2021-02-18T20:18:10.581147898Z", + "data_version":"1.0" + }""" + obj = MockQueueMessage(content=eg_storage_dict) + event = EventGridEvent.from_json(obj) + assert event.data == { + "api":"PutBlockList", + "client_request_id":"6d79dbfb-0e37-4fc4-981f-442c9ca65760", + "request_id":"831e1650-001e-001b-66ab-eeb76e000000", + "e_tag":"0x8D4BCC2E4835CD0", + "content_type":"application/octet-stream", + "content_length":524288, + "blob_type":"BlockBlob", + "url":"https://oc2d2817345i60006.blob.core.windows.net/oc2d2817345i200097container/oc2d2817345i20002296blob", + "sequencer":"00000000000004420000000000028963", + "storage_diagnostics":{"batchId":"b68529f3-68cd-4744-baa4-3c0498ec19f0"} + } + + +def test_from_json(): + json_str = '{"id": "de0fd76c-4ef4-4dfb-ab3a-8f24a307e033", "subject": "https://egtest.dev/cloudcustomevent", "data": {"team": "event grid squad"}, "event_type": "Azure.Sdk.Sample", "event_time": "2020-08-07T02:06:08.11969Z", "data_version": "1.0"}' + event = EventGridEvent.from_json(json_str) + assert event.data == {"team": "event grid squad"} From 1926b64949786cc265a8f5a95671435fcc5b4128 Mon Sep 17 00:00:00 2001 From: Rakshith Bhyravabhotla Date: Tue, 3 Aug 2021 11:40:54 -0700 Subject: [PATCH 17/21] comments --- sdk/core/azure-core/azure/core/messaging.py | 6 ++- .../azure/core/utils/_messaging_shared.py | 41 +++++++++++++++++++ .../azure-core/azure/core/utils/_utils.py | 28 ------------- .../tests/test_messaging_cloud_event.py | 14 ++++--- .../azure/eventgrid/_helpers.py | 27 ------------ .../azure/eventgrid/_messaging_shared.py | 41 +++++++++++++++++++ .../azure/eventgrid/_models.py | 11 ++++- .../consume_cloud_events_from_eventhub.py | 2 +- ...consume_cloud_events_from_storage_queue.py | 2 +- ...eventgrid_events_from_service_bus_queue.py | 2 +- .../tests/test_eg_event_get_bytes.py | 10 +++-- 11 files changed, 114 insertions(+), 70 deletions(-) create mode 100644 sdk/core/azure-core/azure/core/utils/_messaging_shared.py create mode 100644 sdk/eventgrid/azure-eventgrid/azure/eventgrid/_messaging_shared.py diff --git a/sdk/core/azure-core/azure/core/messaging.py b/sdk/core/azure-core/azure/core/messaging.py index 9bc7684b6929..26d05df8b96a 100644 --- a/sdk/core/azure-core/azure/core/messaging.py +++ b/sdk/core/azure-core/azure/core/messaging.py @@ -7,7 +7,8 @@ import uuid from base64 import b64decode from datetime import datetime -from .utils._utils import _convert_to_isoformat, TZ_UTC, _get_json_content +from .utils._utils import _convert_to_isoformat, TZ_UTC +from .utils._messaging_shared import _get_json_content from .serialization import NULL try: @@ -186,11 +187,12 @@ def from_dict(cls, event): def from_json(cls, json): # type: (Any) -> CloudEvent """ - Returns the deserialized CloudEvent object when a json is provided. + Returns the deserialized CloudEvent object when a json payload is provided. :param json: The json string that should be converted into a CloudEvent. This can also be a storage QueueMessage, eventhub's EventData or ServiceBusMessage :type json: object :rtype: CloudEvent + :raises ValueError: If the provided JSON is invalid. """ event = _get_json_content(json) return CloudEvent.from_dict(event) diff --git a/sdk/core/azure-core/azure/core/utils/_messaging_shared.py b/sdk/core/azure-core/azure/core/utils/_messaging_shared.py new file mode 100644 index 000000000000..bc9307d00580 --- /dev/null +++ b/sdk/core/azure-core/azure/core/utils/_messaging_shared.py @@ -0,0 +1,41 @@ +# coding=utf-8 +# -------------------------------------------------------------------------- +# Copyright (c) Microsoft Corporation. All rights reserved. +# Licensed under the MIT License. See License.txt in the project root for +# license information. +# -------------------------------------------------------------------------- + +# ========================================================================== +# This file contains duplicate code that is shared with azure-eventgrid. +# Both the files should always be identical. +# ========================================================================== + + + +import json +from azure.core.exceptions import raise_with_traceback + +def _get_json_content(obj): + """Event mixin to have methods that are common to different Event types + like CloudEvent, EventGridEvent etc. + """ + msg = "Failed to load JSON content from the object." + try: + # storage queue + return json.loads(obj.content) + except ValueError as err: + raise_with_traceback(ValueError, msg, err) + except AttributeError: + # eventhubs + try: + return json.loads(next(obj.body))[0] + except KeyError: + # servicebus + return json.loads(next(obj.body)) + except ValueError as err: + raise_with_traceback(ValueError, msg, err) + except: # pylint: disable=bare-except + try: + return json.loads(obj) + except ValueError as err: + raise_with_traceback(ValueError, msg, err) diff --git a/sdk/core/azure-core/azure/core/utils/_utils.py b/sdk/core/azure-core/azure/core/utils/_utils.py index d8ae29080b7b..08ae1ea21da7 100644 --- a/sdk/core/azure-core/azure/core/utils/_utils.py +++ b/sdk/core/azure-core/azure/core/utils/_utils.py @@ -5,7 +5,6 @@ # license information. # -------------------------------------------------------------------------- import datetime -import json class _FixedOffset(datetime.tzinfo): @@ -107,30 +106,3 @@ def _case_insensitive_dict(*args, **kwargs): raise ValueError( "Neither 'requests' or 'multidict' are installed and no case-insensitive dict impl have been found" ) - -def _get_json_content(obj): - """Event mixin to have methods that are common to different Event types - like CloudEvent, EventGridEvent etc. - """ - try: - # storage queue - return json.loads(obj.content) - except ValueError: - raise ValueError( - "Failed to retrieve content from the object. Make sure the " - + "content follows the CloudEvent schema." - ) - except AttributeError: - # eventhubs - try: - return json.loads(next(obj.body))[0] - except KeyError: - # servicebus - return json.loads(next(obj.body)) - except ValueError: - raise ValueError( - "Failed to retrieve body from the object. Make sure the " - + "body follows the CloudEvent schema." - ) - except: # pylint: disable=bare-except - return json.loads(obj) diff --git a/sdk/core/azure-core/tests/test_messaging_cloud_event.py b/sdk/core/azure-core/tests/test_messaging_cloud_event.py index 98b064776957..c6773747a86d 100644 --- a/sdk/core/azure-core/tests/test_messaging_cloud_event.py +++ b/sdk/core/azure-core/tests/test_messaging_cloud_event.py @@ -3,12 +3,12 @@ # Licensed under the MIT License. # ------------------------------------ import pytest -import json import uuid import datetime from azure.core.messaging import CloudEvent -from azure.core.utils._utils import _convert_to_isoformat, _get_json_content +from azure.core.utils._utils import _convert_to_isoformat +from azure.core.utils._messaging_shared import _get_json_content from azure.core.serialization import NULL class MockQueueMessage(object): @@ -593,7 +593,7 @@ def test_get_bytes_storage_queue_wrong_content(): cloud_storage_string = u'This is a random string which must fail' obj = MockQueueMessage(content=cloud_storage_string) - with pytest.raises(ValueError, match="Failed to retrieve content from the object. Make sure the content follows the CloudEvent schema."): + with pytest.raises(ValueError, match="Failed to load JSON content from the object."): _get_json_content(obj) def test_get_bytes_servicebus(): @@ -627,7 +627,7 @@ def test_get_bytes_servicebus_wrong_content(): lock_token='233146e3-d5a6-45eb-826f-691d82fb8b13' ) - with pytest.raises(ValueError, match="Failed to retrieve body from the object. Make sure the body follows the CloudEvent schema."): + with pytest.raises(ValueError, match="Failed to load JSON content from the object."): _get_json_content(obj) def test_get_bytes_eventhubs(): @@ -643,7 +643,7 @@ def test_get_bytes_eventhubs_wrong_content(): body=MockEhBody(data='random string') ) - with pytest.raises(ValueError, match="Failed to retrieve body from the object. Make sure the body follows the CloudEvent schema."): + with pytest.raises(ValueError, match="Failed to load JSON content from the object."): dict = _get_json_content(obj) def test_get_bytes_random_obj(): @@ -726,3 +726,7 @@ def test_from_json(): event = CloudEvent.from_json(json_str) assert event.data == {"team": "event grid squad"} + assert event.time.year == 2020 + assert event.time.month == 8 + assert event.time.day == 7 + assert event.time.hour == 2 diff --git a/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_helpers.py b/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_helpers.py index 61a2994f7583..8c51aafffafc 100644 --- a/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_helpers.py +++ b/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_helpers.py @@ -165,30 +165,3 @@ def _build_request(endpoint, content_type, events): ) request.format_parameters(query_parameters) return request - -def _get_json_content(obj): - """Event mixin to have methods that are common to different Event types - like CloudEvent, EventGridEvent etc. - """ - try: - # storage queue - return json.loads(obj.content) - except ValueError: - raise ValueError( - "Failed to retrieve content from the object. Make sure the " - + "content follows the EventGridEvent schema." - ) - except AttributeError: - # eventhubs - try: - return json.loads(next(obj.body))[0] - except KeyError: - # servicebus - return json.loads(next(obj.body)) - except ValueError: - raise ValueError( - "Failed to retrieve body from the object. Make sure the " - + "body follows the EventGridEvent schema." - ) - except: # pylint: disable=bare-except - return json.loads(obj) diff --git a/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_messaging_shared.py b/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_messaging_shared.py new file mode 100644 index 000000000000..bc9307d00580 --- /dev/null +++ b/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_messaging_shared.py @@ -0,0 +1,41 @@ +# coding=utf-8 +# -------------------------------------------------------------------------- +# Copyright (c) Microsoft Corporation. All rights reserved. +# Licensed under the MIT License. See License.txt in the project root for +# license information. +# -------------------------------------------------------------------------- + +# ========================================================================== +# This file contains duplicate code that is shared with azure-eventgrid. +# Both the files should always be identical. +# ========================================================================== + + + +import json +from azure.core.exceptions import raise_with_traceback + +def _get_json_content(obj): + """Event mixin to have methods that are common to different Event types + like CloudEvent, EventGridEvent etc. + """ + msg = "Failed to load JSON content from the object." + try: + # storage queue + return json.loads(obj.content) + except ValueError as err: + raise_with_traceback(ValueError, msg, err) + except AttributeError: + # eventhubs + try: + return json.loads(next(obj.body))[0] + except KeyError: + # servicebus + return json.loads(next(obj.body)) + except ValueError as err: + raise_with_traceback(ValueError, msg, err) + except: # pylint: disable=bare-except + try: + return json.loads(obj) + except ValueError as err: + raise_with_traceback(ValueError, msg, err) diff --git a/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py b/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py index 500ad7fb06a0..8dc073994a88 100644 --- a/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py +++ b/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py @@ -7,7 +7,7 @@ import datetime as dt import uuid from msrest.serialization import UTC -from ._helpers import _get_json_content +from ._messaging_shared import _get_json_content from ._generated.models import ( EventGridEvent as InternalEventGridEvent, ) @@ -100,5 +100,14 @@ def __repr__(self): @classmethod def from_json(cls, json): + # type: (Any) -> EventGridEvent + """ + Returns the deserialized EventGridEvent object when a json payload is provided. + :param json: The json string that should be converted into a EventGridEvent. This can also be + a storage QueueMessage, eventhub's EventData or ServiceBusMessage + :type json: object + :rtype: EventGridEvent + :raises ValueError: If the provided JSON is invalid. + """ event = _get_json_content(json) return EventGridEvent.from_dict(event) diff --git a/sdk/eventgrid/azure-eventgrid/samples/consume_samples/consume_cloud_events_from_eventhub.py b/sdk/eventgrid/azure-eventgrid/samples/consume_samples/consume_cloud_events_from_eventhub.py index 9ce79e5d587e..9aee835f0bfc 100644 --- a/sdk/eventgrid/azure-eventgrid/samples/consume_samples/consume_cloud_events_from_eventhub.py +++ b/sdk/eventgrid/azure-eventgrid/samples/consume_samples/consume_cloud_events_from_eventhub.py @@ -25,7 +25,7 @@ CONNECTION_STR = os.environ["EVENT_HUB_CONN_STR"] EVENTHUB_NAME = os.environ["EVENT_HUB_NAME"] def on_event(partition_context, event): - dict_event = CloudEvent.from_dict(event) + dict_event = CloudEvent.from_json(event) print("data: {}\n".format(dict_event.data)) consumer_client = EventHubConsumerClient.from_connection_string( diff --git a/sdk/eventgrid/azure-eventgrid/samples/consume_samples/consume_cloud_events_from_storage_queue.py b/sdk/eventgrid/azure-eventgrid/samples/consume_samples/consume_cloud_events_from_storage_queue.py index 148906586809..4db27504690e 100644 --- a/sdk/eventgrid/azure-eventgrid/samples/consume_samples/consume_cloud_events_from_storage_queue.py +++ b/sdk/eventgrid/azure-eventgrid/samples/consume_samples/consume_cloud_events_from_storage_queue.py @@ -30,7 +30,7 @@ ).peek_messages(max_messages=32) ## deserialize payload into a list of typed Events - events = [CloudEvent.from_dict(msg) for msg in payload] + events = [CloudEvent.from_json(msg) for msg in payload] for event in events: print(type(event)) ## CloudEvent diff --git a/sdk/eventgrid/azure-eventgrid/samples/consume_samples/consume_eventgrid_events_from_service_bus_queue.py b/sdk/eventgrid/azure-eventgrid/samples/consume_samples/consume_eventgrid_events_from_service_bus_queue.py index 2f13f22ba905..9889ea87c112 100644 --- a/sdk/eventgrid/azure-eventgrid/samples/consume_samples/consume_eventgrid_events_from_service_bus_queue.py +++ b/sdk/eventgrid/azure-eventgrid/samples/consume_samples/consume_eventgrid_events_from_service_bus_queue.py @@ -30,7 +30,7 @@ payload = sb_client.get_queue_receiver(queue_name).receive_messages() ## deserialize payload into a list of typed Events - events = [EventGridEvent.from_dict(msg) for msg in payload] + events = [EventGridEvent.from_json(msg) for msg in payload] for event in events: print(type(event)) ## EventGridEvent diff --git a/sdk/eventgrid/azure-eventgrid/tests/test_eg_event_get_bytes.py b/sdk/eventgrid/azure-eventgrid/tests/test_eg_event_get_bytes.py index 03d350598213..1ae80ecde2c9 100644 --- a/sdk/eventgrid/azure-eventgrid/tests/test_eg_event_get_bytes.py +++ b/sdk/eventgrid/azure-eventgrid/tests/test_eg_event_get_bytes.py @@ -5,7 +5,8 @@ import datetime import pytest import uuid -from azure.eventgrid._helpers import _get_json_content +from msrest.serialization import UTC +from azure.eventgrid._messaging_shared import _get_json_content from azure.eventgrid import EventGridEvent class MockQueueMessage(object): @@ -129,7 +130,7 @@ def test_get_bytes_storage_queue_wrong_content(): string = u'This is a random string which must fail' obj = MockQueueMessage(content=string) - with pytest.raises(ValueError, match="Failed to retrieve content from the object. Make sure the content follows the EventGridEvent schema."): + with pytest.raises(ValueError, match="Failed to load JSON content from the object."): _get_json_content(obj) def test_get_bytes_servicebus(): @@ -162,7 +163,7 @@ def test_get_bytes_servicebus_wrong_content(): sequence_number=11219, lock_token='233146e3-d5a6-45eb-826f-691d82fb8b13' ) - with pytest.raises(ValueError, match="Failed to retrieve body from the object. Make sure the body follows the EventGridEvent schema."): + with pytest.raises(ValueError, match="Failed to load JSON content from the object."): dict = _get_json_content(obj) @@ -179,7 +180,7 @@ def test_get_bytes_eventhubs_wrong_content(): body=MockEhBody(data='random string') ) - with pytest.raises(ValueError, match="Failed to retrieve body from the object. Make sure the body follows the EventGridEvent schema."): + with pytest.raises(ValueError, match="Failed to load JSON content from the object."): dict = _get_json_content(obj) @@ -262,3 +263,4 @@ def test_from_json(): json_str = '{"id": "de0fd76c-4ef4-4dfb-ab3a-8f24a307e033", "subject": "https://egtest.dev/cloudcustomevent", "data": {"team": "event grid squad"}, "event_type": "Azure.Sdk.Sample", "event_time": "2020-08-07T02:06:08.11969Z", "data_version": "1.0"}' event = EventGridEvent.from_json(json_str) assert event.data == {"team": "event grid squad"} + assert event.event_time == datetime.datetime(2020, 8, 7, 2, 6, 8, 119690, UTC()) From f6b6593daa94ae04882e0df198ab6423e55c8ca7 Mon Sep 17 00:00:00 2001 From: Rakshith Bhyravabhotla Date: Tue, 3 Aug 2021 14:12:16 -0700 Subject: [PATCH 18/21] handle type error --- .../azure/core/utils/_messaging_shared.py | 6 ++-- .../tests/test_messaging_cloud_event.py | 30 +++++++++++++++++++ .../azure/eventgrid/_messaging_shared.py | 6 ++-- .../tests/test_eg_event_get_bytes.py | 30 +++++++++++++++++++ 4 files changed, 66 insertions(+), 6 deletions(-) diff --git a/sdk/core/azure-core/azure/core/utils/_messaging_shared.py b/sdk/core/azure-core/azure/core/utils/_messaging_shared.py index bc9307d00580..a5b172d20d2e 100644 --- a/sdk/core/azure-core/azure/core/utils/_messaging_shared.py +++ b/sdk/core/azure-core/azure/core/utils/_messaging_shared.py @@ -23,7 +23,7 @@ def _get_json_content(obj): try: # storage queue return json.loads(obj.content) - except ValueError as err: + except (ValueError, TypeError) as err: raise_with_traceback(ValueError, msg, err) except AttributeError: # eventhubs @@ -32,10 +32,10 @@ def _get_json_content(obj): except KeyError: # servicebus return json.loads(next(obj.body)) - except ValueError as err: + except (ValueError, TypeError) as err: raise_with_traceback(ValueError, msg, err) except: # pylint: disable=bare-except try: return json.loads(obj) - except ValueError as err: + except (ValueError, TypeError) as err: raise_with_traceback(ValueError, msg, err) diff --git a/sdk/core/azure-core/tests/test_messaging_cloud_event.py b/sdk/core/azure-core/tests/test_messaging_cloud_event.py index c6773747a86d..b87cf6c7cdc1 100644 --- a/sdk/core/azure-core/tests/test_messaging_cloud_event.py +++ b/sdk/core/azure-core/tests/test_messaging_cloud_event.py @@ -76,6 +76,20 @@ def __next__(self): next = __next__ +class MockTypeErrorBody(object): + def __init__(self, data=None): + self.data = data + + def __iter__(self): + return self + + def __next__(self): + if not self.data: + return {"id":"f208feff-099b-4bda-a341-4afd0fa02fef","source":"https://egsample.dev/sampleevent","data":"ServiceBus","type":"Azure.Sdk.Sample","time":"2021-07-22T22:27:38.960209Z","specversion":"1.0"} + return self.data + + next = __next__ + class MockEhBody(object): def __init__(self, data=None): self.data = data @@ -730,3 +744,19 @@ def test_from_json(): assert event.time.month == 8 assert event.time.day == 7 assert event.time.hour == 2 + +def test_from_json_sb_type_error(): + obj = MockServiceBusReceivedMessage( + body=MockTypeErrorBody(), + message_id='3f6c5441-5be5-4f33-80c3-3ffeb6a090ce', + content_type='application/cloudevents+json; charset=utf-8', + time_to_live=datetime.timedelta(days=14), + delivery_count=13, + enqueued_sequence_number=0, + enqueued_time_utc=datetime.datetime(2021, 7, 22, 22, 27, 41, 236000), + expires_at_utc=datetime.datetime(2021, 8, 5, 22, 27, 41, 236000), + sequence_number=11219, + lock_token='233146e3-d5a6-45eb-826f-691d82fb8b13' + ) + with pytest.raises(ValueError): + _get_json_content(obj) \ No newline at end of file diff --git a/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_messaging_shared.py b/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_messaging_shared.py index bc9307d00580..a5b172d20d2e 100644 --- a/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_messaging_shared.py +++ b/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_messaging_shared.py @@ -23,7 +23,7 @@ def _get_json_content(obj): try: # storage queue return json.loads(obj.content) - except ValueError as err: + except (ValueError, TypeError) as err: raise_with_traceback(ValueError, msg, err) except AttributeError: # eventhubs @@ -32,10 +32,10 @@ def _get_json_content(obj): except KeyError: # servicebus return json.loads(next(obj.body)) - except ValueError as err: + except (ValueError, TypeError) as err: raise_with_traceback(ValueError, msg, err) except: # pylint: disable=bare-except try: return json.loads(obj) - except ValueError as err: + except (ValueError, TypeError) as err: raise_with_traceback(ValueError, msg, err) diff --git a/sdk/eventgrid/azure-eventgrid/tests/test_eg_event_get_bytes.py b/sdk/eventgrid/azure-eventgrid/tests/test_eg_event_get_bytes.py index 1ae80ecde2c9..8463953b883d 100644 --- a/sdk/eventgrid/azure-eventgrid/tests/test_eg_event_get_bytes.py +++ b/sdk/eventgrid/azure-eventgrid/tests/test_eg_event_get_bytes.py @@ -19,6 +19,20 @@ def __init__(self, content=None): self.pop_receipt = None self.next_visible_on = None +class MockTypeErrorBody(object): + def __init__(self, data=None): + self.data = data + + def __iter__(self): + return self + + def __next__(self): + if not self.data: + return {"id":"f208feff-099b-4bda-a341-4afd0fa02fef","source":"https://egsample.dev/sampleevent","data":"ServiceBus","type":"Azure.Sdk.Sample","time":"2021-07-22T22:27:38.960209Z","specversion":"1.0"} + return self.data + + next = __next__ + class MockServiceBusReceivedMessage(object): def __init__(self, body=None, **kwargs): self.body=body @@ -264,3 +278,19 @@ def test_from_json(): event = EventGridEvent.from_json(json_str) assert event.data == {"team": "event grid squad"} assert event.event_time == datetime.datetime(2020, 8, 7, 2, 6, 8, 119690, UTC()) + +def test_from_json_sb_type_error(): + obj = MockServiceBusReceivedMessage( + body=MockTypeErrorBody(), + message_id='3f6c5441-5be5-4f33-80c3-3ffeb6a090ce', + content_type='application/cloudevents+json; charset=utf-8', + time_to_live=datetime.timedelta(days=14), + delivery_count=13, + enqueued_sequence_number=0, + enqueued_time_utc=datetime.datetime(2021, 7, 22, 22, 27, 41, 236000), + expires_at_utc=datetime.datetime(2021, 8, 5, 22, 27, 41, 236000), + sequence_number=11219, + lock_token='233146e3-d5a6-45eb-826f-691d82fb8b13' + ) + with pytest.raises(ValueError): + _get_json_content(obj) From ee4b1ab0f8048ca352decd180e0dce3e0eb307e2 Mon Sep 17 00:00:00 2001 From: Rakshith Bhyravabhotla Date: Tue, 3 Aug 2021 14:52:12 -0700 Subject: [PATCH 19/21] Revert "handle type error" This reverts commit f6b6593daa94ae04882e0df198ab6423e55c8ca7. --- .../azure/core/utils/_messaging_shared.py | 6 ++-- .../tests/test_messaging_cloud_event.py | 30 ------------------- .../azure/eventgrid/_messaging_shared.py | 6 ++-- .../tests/test_eg_event_get_bytes.py | 30 ------------------- 4 files changed, 6 insertions(+), 66 deletions(-) diff --git a/sdk/core/azure-core/azure/core/utils/_messaging_shared.py b/sdk/core/azure-core/azure/core/utils/_messaging_shared.py index a5b172d20d2e..bc9307d00580 100644 --- a/sdk/core/azure-core/azure/core/utils/_messaging_shared.py +++ b/sdk/core/azure-core/azure/core/utils/_messaging_shared.py @@ -23,7 +23,7 @@ def _get_json_content(obj): try: # storage queue return json.loads(obj.content) - except (ValueError, TypeError) as err: + except ValueError as err: raise_with_traceback(ValueError, msg, err) except AttributeError: # eventhubs @@ -32,10 +32,10 @@ def _get_json_content(obj): except KeyError: # servicebus return json.loads(next(obj.body)) - except (ValueError, TypeError) as err: + except ValueError as err: raise_with_traceback(ValueError, msg, err) except: # pylint: disable=bare-except try: return json.loads(obj) - except (ValueError, TypeError) as err: + except ValueError as err: raise_with_traceback(ValueError, msg, err) diff --git a/sdk/core/azure-core/tests/test_messaging_cloud_event.py b/sdk/core/azure-core/tests/test_messaging_cloud_event.py index b87cf6c7cdc1..c6773747a86d 100644 --- a/sdk/core/azure-core/tests/test_messaging_cloud_event.py +++ b/sdk/core/azure-core/tests/test_messaging_cloud_event.py @@ -76,20 +76,6 @@ def __next__(self): next = __next__ -class MockTypeErrorBody(object): - def __init__(self, data=None): - self.data = data - - def __iter__(self): - return self - - def __next__(self): - if not self.data: - return {"id":"f208feff-099b-4bda-a341-4afd0fa02fef","source":"https://egsample.dev/sampleevent","data":"ServiceBus","type":"Azure.Sdk.Sample","time":"2021-07-22T22:27:38.960209Z","specversion":"1.0"} - return self.data - - next = __next__ - class MockEhBody(object): def __init__(self, data=None): self.data = data @@ -744,19 +730,3 @@ def test_from_json(): assert event.time.month == 8 assert event.time.day == 7 assert event.time.hour == 2 - -def test_from_json_sb_type_error(): - obj = MockServiceBusReceivedMessage( - body=MockTypeErrorBody(), - message_id='3f6c5441-5be5-4f33-80c3-3ffeb6a090ce', - content_type='application/cloudevents+json; charset=utf-8', - time_to_live=datetime.timedelta(days=14), - delivery_count=13, - enqueued_sequence_number=0, - enqueued_time_utc=datetime.datetime(2021, 7, 22, 22, 27, 41, 236000), - expires_at_utc=datetime.datetime(2021, 8, 5, 22, 27, 41, 236000), - sequence_number=11219, - lock_token='233146e3-d5a6-45eb-826f-691d82fb8b13' - ) - with pytest.raises(ValueError): - _get_json_content(obj) \ No newline at end of file diff --git a/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_messaging_shared.py b/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_messaging_shared.py index a5b172d20d2e..bc9307d00580 100644 --- a/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_messaging_shared.py +++ b/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_messaging_shared.py @@ -23,7 +23,7 @@ def _get_json_content(obj): try: # storage queue return json.loads(obj.content) - except (ValueError, TypeError) as err: + except ValueError as err: raise_with_traceback(ValueError, msg, err) except AttributeError: # eventhubs @@ -32,10 +32,10 @@ def _get_json_content(obj): except KeyError: # servicebus return json.loads(next(obj.body)) - except (ValueError, TypeError) as err: + except ValueError as err: raise_with_traceback(ValueError, msg, err) except: # pylint: disable=bare-except try: return json.loads(obj) - except (ValueError, TypeError) as err: + except ValueError as err: raise_with_traceback(ValueError, msg, err) diff --git a/sdk/eventgrid/azure-eventgrid/tests/test_eg_event_get_bytes.py b/sdk/eventgrid/azure-eventgrid/tests/test_eg_event_get_bytes.py index 8463953b883d..1ae80ecde2c9 100644 --- a/sdk/eventgrid/azure-eventgrid/tests/test_eg_event_get_bytes.py +++ b/sdk/eventgrid/azure-eventgrid/tests/test_eg_event_get_bytes.py @@ -19,20 +19,6 @@ def __init__(self, content=None): self.pop_receipt = None self.next_visible_on = None -class MockTypeErrorBody(object): - def __init__(self, data=None): - self.data = data - - def __iter__(self): - return self - - def __next__(self): - if not self.data: - return {"id":"f208feff-099b-4bda-a341-4afd0fa02fef","source":"https://egsample.dev/sampleevent","data":"ServiceBus","type":"Azure.Sdk.Sample","time":"2021-07-22T22:27:38.960209Z","specversion":"1.0"} - return self.data - - next = __next__ - class MockServiceBusReceivedMessage(object): def __init__(self, body=None, **kwargs): self.body=body @@ -278,19 +264,3 @@ def test_from_json(): event = EventGridEvent.from_json(json_str) assert event.data == {"team": "event grid squad"} assert event.event_time == datetime.datetime(2020, 8, 7, 2, 6, 8, 119690, UTC()) - -def test_from_json_sb_type_error(): - obj = MockServiceBusReceivedMessage( - body=MockTypeErrorBody(), - message_id='3f6c5441-5be5-4f33-80c3-3ffeb6a090ce', - content_type='application/cloudevents+json; charset=utf-8', - time_to_live=datetime.timedelta(days=14), - delivery_count=13, - enqueued_sequence_number=0, - enqueued_time_utc=datetime.datetime(2021, 7, 22, 22, 27, 41, 236000), - expires_at_utc=datetime.datetime(2021, 8, 5, 22, 27, 41, 236000), - sequence_number=11219, - lock_token='233146e3-d5a6-45eb-826f-691d82fb8b13' - ) - with pytest.raises(ValueError): - _get_json_content(obj) From 262b29ba03cb60a7a201015d28691907891e66f1 Mon Sep 17 00:00:00 2001 From: Rakshith Bhyravabhotla Date: Tue, 3 Aug 2021 14:54:42 -0700 Subject: [PATCH 20/21] mypy fix --- sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py b/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py index 8dc073994a88..3fd5a1005633 100644 --- a/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py +++ b/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py @@ -3,7 +3,7 @@ # Licensed under the MIT License. See License.txt in the project root for license information. # -------------------------------------------------------------------------------------------- # pylint:disable=protected-access -from typing import Any +from typing import Any, cast import datetime as dt import uuid from msrest.serialization import UTC @@ -110,4 +110,4 @@ def from_json(cls, json): :raises ValueError: If the provided JSON is invalid. """ event = _get_json_content(json) - return EventGridEvent.from_dict(event) + return cast(EventGridEvent, EventGridEvent.from_dict(event)) From 5a11e221dbf27d3d8cf26e95a98c0b199a2242f5 Mon Sep 17 00:00:00 2001 From: Rakshith Bhyravabhotla Date: Tue, 3 Aug 2021 16:03:17 -0700 Subject: [PATCH 21/21] rename param to event --- sdk/core/azure-core/azure/core/messaging.py | 10 +++++----- .../azure-eventgrid/azure/eventgrid/_models.py | 10 +++++----- 2 files changed, 10 insertions(+), 10 deletions(-) diff --git a/sdk/core/azure-core/azure/core/messaging.py b/sdk/core/azure-core/azure/core/messaging.py index 26d05df8b96a..0692986f388a 100644 --- a/sdk/core/azure-core/azure/core/messaging.py +++ b/sdk/core/azure-core/azure/core/messaging.py @@ -184,15 +184,15 @@ def from_dict(cls, event): return event_obj @classmethod - def from_json(cls, json): + def from_json(cls, event): # type: (Any) -> CloudEvent """ Returns the deserialized CloudEvent object when a json payload is provided. - :param json: The json string that should be converted into a CloudEvent. This can also be + :param event: The json string that should be converted into a CloudEvent. This can also be a storage QueueMessage, eventhub's EventData or ServiceBusMessage - :type json: object + :type event: object :rtype: CloudEvent :raises ValueError: If the provided JSON is invalid. """ - event = _get_json_content(json) - return CloudEvent.from_dict(event) + dict_event = _get_json_content(event) + return CloudEvent.from_dict(dict_event) diff --git a/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py b/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py index 3fd5a1005633..ecc74505a9db 100644 --- a/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py +++ b/sdk/eventgrid/azure-eventgrid/azure/eventgrid/_models.py @@ -99,15 +99,15 @@ def __repr__(self): )[:1024] @classmethod - def from_json(cls, json): + def from_json(cls, event): # type: (Any) -> EventGridEvent """ Returns the deserialized EventGridEvent object when a json payload is provided. - :param json: The json string that should be converted into a EventGridEvent. This can also be + :param event: The json string that should be converted into a EventGridEvent. This can also be a storage QueueMessage, eventhub's EventData or ServiceBusMessage - :type json: object + :type event: object :rtype: EventGridEvent :raises ValueError: If the provided JSON is invalid. """ - event = _get_json_content(json) - return cast(EventGridEvent, EventGridEvent.from_dict(event)) + dict_event = _get_json_content(event) + return cast(EventGridEvent, EventGridEvent.from_dict(dict_event))